diff --git a/crates/puffer-cli/src/daemon.rs b/crates/puffer-cli/src/daemon.rs index 5d4d66341..ef7ddbfec 100644 --- a/crates/puffer-cli/src/daemon.rs +++ b/crates/puffer-cli/src/daemon.rs @@ -1484,21 +1484,21 @@ async fn dispatch_request( .await; } - "list_grouped_sessions" => respond!(handle_list_grouped_sessions(&state)), + "list_grouped_sessions" => respond!(detached!(|s| handle_list_grouped_sessions(&s))), "list_grouped_sessions_page" => { respond!(detached!(|s, p| handle_list_grouped_sessions_page(&s, &p))) } - "load_desktop_pins" => respond!(handle_load_desktop_pins(&state)), - "set_desktop_pin" => respond!(handle_set_desktop_pin(&state, ¶ms)), - "load_file_tabs" => respond!(handle_load_file_tabs(&state, ¶ms)), - "save_file_tabs" => respond!(handle_save_file_tabs(&state, ¶ms)), - "load_session_detail" => respond!(handle_load_session_detail(&state, ¶ms)), - "rename_session" => respond!(handle_rename_session(&state, ¶ms)), - "delete_session" => respond!(handle_delete_session(&state, ¶ms)), - "set_session_tags" => respond!(handle_set_session_tags(&state, ¶ms)), - "delete_project" => respond!(handle_delete_project(&state, ¶ms)), - "set_project_tags" => respond!(handle_set_project_tags(&state, ¶ms)), - "refresh_repo_status" => respond!(handle_refresh_repo_status(&state, ¶ms)), + "load_desktop_pins" => respond!(detached!(|s| handle_load_desktop_pins(&s))), + "set_desktop_pin" => respond!(detached!(|s, p| handle_set_desktop_pin(&s, &p))), + "load_file_tabs" => respond!(detached!(|s, p| handle_load_file_tabs(&s, &p))), + "save_file_tabs" => respond!(detached!(|s, p| handle_save_file_tabs(&s, &p))), + "load_session_detail" => respond!(detached!(|s, p| handle_load_session_detail(&s, &p))), + "rename_session" => respond!(detached!(|s, p| handle_rename_session(&s, &p))), + "delete_session" => respond!(detached!(|s, p| handle_delete_session(&s, &p))), + "set_session_tags" => respond!(detached!(|s, p| handle_set_session_tags(&s, &p))), + "delete_project" => respond!(detached!(|s, p| handle_delete_project(&s, &p))), + "set_project_tags" => respond!(detached!(|s, p| handle_set_project_tags(&s, &p))), + "refresh_repo_status" => respond!(detached!(|s, p| handle_refresh_repo_status(&s, &p))), "load_settings_snapshot" => respond!(detached!(|s| handle_load_settings_snapshot(&s))), "login_with_api_key" => { respond!(detached!(|s, p| handle_login_with_api_key(&s, &p))) @@ -1751,42 +1751,51 @@ async fn dispatch_request( "workflow_open_ui" => respond!(detached!(|s| { crate::daemon_workflow_runtime::handle_workflow_open_ui(&s) })), - "workflow_binding_create" => respond!( - crate::daemon_workflows::handle_workflow_binding_create(&state.paths, ¶ms) - ), - "monitor_create" | "task_monitor_create" => respond!( - crate::daemon_workflows::handle_monitor_create(&state.paths, ¶ms) - ), - "monitor_task_ignore" | "task_monitor_ignore" => respond!( - crate::daemon_workflows::handle_monitor_task_ignore(&state.paths, ¶ms) - ), - "monitor_rule_add" | "task_monitor_rule_add" => respond!( - crate::daemon_workflows::handle_monitor_rule_add(&state.paths, ¶ms) - ), - "monitor_rule_delete" | "task_monitor_rule_delete" => respond!( - crate::daemon_workflows::handle_monitor_rule_delete(&state.paths, ¶ms) - ), - "monitor_task_complete" | "task_monitor_complete" => respond!( - crate::daemon_workflows::handle_monitor_task_complete(&state.paths, ¶ms) - ), - "outbound_action_execute" => respond!( - crate::daemon_workflows::handle_outbound_action_execute(&state.paths, ¶ms) - ), - "outbound_action_cancel" => respond!( - crate::daemon_workflows::handle_outbound_action_cancel(&state.paths, ¶ms) - ), - "outbound_action_status" => respond!( - crate::daemon_workflows::handle_outbound_action_status(&state.paths, ¶ms) - ), - "monitor_memory_save" | "task_monitor_memory_save" => respond!( - crate::daemon_workflows::handle_monitor_memory_save(&state.paths, ¶ms) - ), - "monitor_history_list" | "task_monitor_history_list" => respond!( - crate::daemon_workflows::handle_monitor_history_list(&state.paths, ¶ms) - ), - "monitor_trace_list" | "task_monitor_trace_list" => respond!( - crate::daemon_workflows::handle_monitor_trace_list(&state.paths, ¶ms) - ), + "workflow_binding_create" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_workflow_binding_create(s.config_paths(), &p) + })), + "monitor_create" | "task_monitor_create" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_create(s.config_paths(), &p) + })), + "monitor_task_ignore" | "task_monitor_ignore" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_task_ignore(s.config_paths(), &p) + })), + "monitor_rule_add" | "task_monitor_rule_add" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_rule_add(s.config_paths(), &p) + })), + "monitor_rule_delete" | "task_monitor_rule_delete" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_rule_delete(s.config_paths(), &p) + })), + "monitor_task_complete" | "task_monitor_complete" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_task_complete(s.config_paths(), &p) + })), + "monitor_reply_send" | "task_monitor_reply_send" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_outbound_action_execute(s.config_paths(), &p) + })), + "monitor_action_execute" | "task_monitor_action_execute" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_outbound_action_execute(s.config_paths(), &p) + })), + "connector_action_execute" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_outbound_action_execute(s.config_paths(), &p) + })), + "outbound_action_execute" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_outbound_action_execute(s.config_paths(), &p) + })), + "outbound_action_cancel" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_outbound_action_cancel(s.config_paths(), &p) + })), + "outbound_action_status" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_outbound_action_status(s.config_paths(), &p) + })), + "monitor_memory_save" | "task_monitor_memory_save" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_memory_save(s.config_paths(), &p) + })), + "monitor_history_list" | "task_monitor_history_list" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_history_list(s.config_paths(), &p) + })), + "monitor_trace_list" | "task_monitor_trace_list" => respond!(detached!(|s, p| { + crate::daemon_workflows::handle_monitor_trace_list(s.config_paths(), &p) + })), "telegram_diagnostics_export" | "task_telegram_diagnostics_export" => respond!( crate::daemon_workflows::handle_telegram_diagnostics_export(&state.paths, ¶ms) ), diff --git a/crates/puffer-cli/src/daemon_workflows/monitor_memory.rs b/crates/puffer-cli/src/daemon_workflows/monitor_memory.rs index 7f5d2bfc1..c0996941e 100644 --- a/crates/puffer-cli/src/daemon_workflows/monitor_memory.rs +++ b/crates/puffer-cli/src/daemon_workflows/monitor_memory.rs @@ -46,7 +46,7 @@ pub(crate) fn handle_monitor_memory_save(paths: &ConfigPaths, params: &Value) -> } fs::write(&path, params.content) .with_context(|| format!("failed to write {}", path.display()))?; - super::handle_workflow_list(paths) + super::handle_workflow_list_with_runtime(paths, false) } fn load_monitor_memories(paths: &ConfigPaths) -> Result> { diff --git a/crates/puffer-cli/src/daemon_workflows/monitor_rules.rs b/crates/puffer-cli/src/daemon_workflows/monitor_rules.rs index 64669bb7c..d8f9f1673 100644 --- a/crates/puffer-cli/src/daemon_workflows/monitor_rules.rs +++ b/crates/puffer-cli/src/daemon_workflows/monitor_rules.rs @@ -10,6 +10,8 @@ use serde::Deserialize; use serde_json::Value; use std::collections::BTreeSet; use std::path::PathBuf; +use std::sync::Arc; +use std::thread; #[derive(Debug, Deserialize)] struct MonitorRuleAddParams { @@ -63,8 +65,12 @@ pub(crate) fn handle_monitor_rule_add(paths: &ConfigPaths, params: &Value) -> Re } } manager.store().upsert(binding)?; - manager.refresh_connection_consumers()?; - super::handle_workflow_list(paths) + refresh_connection_consumers_background( + Arc::clone(&manager), + connection_slug.to_string(), + "adding monitor rule", + ); + super::handle_workflow_list_with_runtime(paths, false) } /// Deletes one include or exclude monitor rule and returns a refreshed snapshot. @@ -87,8 +93,40 @@ pub(crate) fn handle_monitor_rule_delete(paths: &ConfigPaths, params: &Value) -> } } manager.store().upsert(binding)?; - manager.refresh_connection_consumers()?; - super::handle_workflow_list(paths) + refresh_connection_consumers_background( + Arc::clone(&manager), + connection_slug.to_string(), + "deleting monitor rule", + ); + super::handle_workflow_list_with_runtime(paths, false) +} + +fn refresh_connection_consumers_background( + manager: Arc, + connection_slug: String, + operation: &'static str, +) { + let spawn_connection_slug = connection_slug.clone(); + if let Err(error) = thread::Builder::new() + .name("puffer-monitor-rule-refresh".to_string()) + .spawn(move || { + if let Err(error) = manager.refresh_connection_consumers() { + tracing::warn!( + connection = %spawn_connection_slug, + %error, + operation, + "failed to refresh connection consumers after monitor rule change" + ); + } + }) + { + tracing::warn!( + connection = %connection_slug, + %error, + operation, + "failed to spawn connection consumer refresh after monitor rule change" + ); + } } pub(super) fn include_filters_json(filter: Option<&FilterSpec>) -> Value { diff --git a/crates/puffer-cli/src/daemon_workflows/monitor_rules_tests.rs b/crates/puffer-cli/src/daemon_workflows/monitor_rules_tests.rs index dd1880896..0104f93c0 100644 --- a/crates/puffer-cli/src/daemon_workflows/monitor_rules_tests.rs +++ b/crates/puffer-cli/src/daemon_workflows/monitor_rules_tests.rs @@ -10,6 +10,7 @@ use serde_json::json; use std::fs; use std::path::Path; use std::sync::{Arc, OnceLock}; +use std::time::{Duration, Instant}; struct TestManager { _runtime: tokio::runtime::Runtime, @@ -87,6 +88,21 @@ fn config_paths_with_bundled_resources(root: &Path) -> ConfigPaths { } } +fn wait_for_has_consumer(manager: &SubscriptionManager, connection_slug: &str) -> bool { + let deadline = Instant::now() + Duration::from_secs(2); + while Instant::now() < deadline { + if manager + .connection_store() + .get(connection_slug) + .is_some_and(|connection| connection.has_consumer) + { + return true; + } + std::thread::sleep(Duration::from_millis(20)); + } + false +} + #[test] fn add_exclude_rule_persists_case_insensitive_keyword_filter_and_ignores_legacy_scope() { let tempdir = tempfile::tempdir().unwrap(); @@ -185,6 +201,88 @@ fn add_include_rule_persists_filter_without_touching_other_monitor() { assert!(other.ignore_filters.is_empty()); } +#[test] +fn add_rule_refreshes_connection_consumers_in_background() { + let tempdir = tempfile::tempdir().unwrap(); + let paths = ConfigPaths::discover(tempdir.path()); + let manager = test_manager(); + let connection_slug = "rule-add-background-refresh"; + let _ = manager + .connection_store() + .create(ConnectionRecord::authenticated( + connection_slug, + "telegram-login", + "Telegram", + )); + manager + .store() + .upsert(monitor_binding( + &format!("monitor-{connection_slug}"), + connection_slug, + )) + .unwrap(); + assert!( + !manager + .connection_store() + .get(connection_slug) + .unwrap() + .has_consumer + ); + + handle_monitor_rule_add( + &paths, + &json!({ + "connection_slug": connection_slug, + "mode": "exclude", + "keywords": ["noise"] + }), + ) + .unwrap(); + + assert!(wait_for_has_consumer(manager.as_ref(), connection_slug)); +} + +#[test] +fn delete_rule_refreshes_connection_consumers_in_background() { + let tempdir = tempfile::tempdir().unwrap(); + let paths = ConfigPaths::discover(tempdir.path()); + let manager = test_manager(); + let connection_slug = "rule-delete-background-refresh"; + let rule = FilterSpec::Tagged(TaggedFilterSpec::Regex { + pattern: "noise".to_string(), + case_insensitive: true, + }); + let mut binding = monitor_binding(&format!("monitor-{connection_slug}"), connection_slug); + binding.ignore_filters.push(rule.clone()); + let _ = manager + .connection_store() + .create(ConnectionRecord::authenticated( + connection_slug, + "telegram-login", + "Telegram", + )); + manager.store().upsert(binding).unwrap(); + assert!( + !manager + .connection_store() + .get(connection_slug) + .unwrap() + .has_consumer + ); + + handle_monitor_rule_delete( + &paths, + &json!({ + "connection_slug": connection_slug, + "mode": "exclude", + "rule": serde_json::to_value(rule).unwrap() + }), + ) + .unwrap(); + + assert!(wait_for_has_consumer(manager.as_ref(), connection_slug)); +} + #[test] fn delete_rule_removes_exact_displayed_rule() { let tempdir = tempfile::tempdir().unwrap();