diff --git a/docs/content/docs/(configuration)/config.mdx b/docs/content/docs/(configuration)/config.mdx index f0115567d..b8205ff92 100644 --- a/docs/content/docs/(configuration)/config.mdx +++ b/docs/content/docs/(configuration)/config.mdx @@ -266,14 +266,13 @@ Most config values are hot-reloaded when their files change. Spacebot watches `c | Identity files (SOUL.md, etc.) | Yes | Next channel message renders new identity | | Skills (SKILL.md files) | Yes | Next message / worker spawn sees new skills | | Bindings | Yes | Next message routes using new bindings | -| Discord/Slack permissions | Yes | Next message checks new permission rules | +| Messaging adapters (Discord/Slack/Telegram/Twitch/Email/Webhook) | Yes | Adapter runtime is reconciled live; changed instances are restarted or removed without a full process restart | ### What Needs Restart | Setting | Why | |---------|-----| | LLM API keys | Provider clients are initialized once (applies to `secret:`, `env:`, and literal values) | -| Messaging adapters (Discord token, webhook bind/port) | Adapter connections are long-lived | | Agent topology (adding/removing `[[agents]]`) | Databases and event buses are per-agent | | Database paths | Connections are opened once at startup | | System prompts | Compiled into the binary via `include_str!` | diff --git a/src/api/messaging.rs b/src/api/messaging.rs index 6aad16fe3..ba1c39da0 100644 --- a/src/api/messaging.rs +++ b/src/api/messaging.rs @@ -1059,11 +1059,7 @@ pub(super) async fn disconnect_platform( if platform == "twitch" { let instance_dir = state.instance_dir.load(); if let Some(name) = adapter_name { - let safe_name: String = name - .chars() - .map(|ch| if ch.is_ascii_alphanumeric() { ch } else { '_' }) - .collect(); - let token_path = instance_dir.join(format!("twitch_token_{safe_name}.json")); + let token_path = instance_dir.join(crate::config::named_twitch_token_file_name(name)); match tokio::fs::remove_file(&token_path).await { Ok(()) => { tracing::info!(path = %token_path.display(), "twitch token file deleted"); @@ -1438,14 +1434,8 @@ pub(super) async fn toggle_platform( "twitch", Some(instance.name.as_str()), ); - let token_file_name = format!( - "twitch_token_{}.json", - instance - .name - .chars() - .map(|ch| if ch.is_ascii_alphanumeric() { ch } else { '_' }) - .collect::() - ); + let token_file_name = + crate::config::named_twitch_token_file_name(&instance.name); let instance_dir = state.instance_dir.load(); let token_path = instance_dir.join(token_file_name); let perms = std::sync::Arc::new(arc_swap::ArcSwap::from_pointee( @@ -2269,11 +2259,7 @@ pub(super) async fn delete_messaging_instance( if platform == "twitch" { let instance_dir = state.instance_dir.load(); let token_path = if let Some(name) = adapter_name { - let safe_name: String = name - .chars() - .map(|ch| if ch.is_ascii_alphanumeric() { ch } else { '_' }) - .collect(); - instance_dir.join(format!("twitch_token_{safe_name}.json")) + instance_dir.join(crate::config::named_twitch_token_file_name(name)) } else { instance_dir.join("twitch_token.json") }; diff --git a/src/config.rs b/src/config.rs index 7a4c1e3f0..964142f8d 100644 --- a/src/config.rs +++ b/src/config.rs @@ -21,7 +21,7 @@ pub use permissions::{ pub(crate) use providers::default_provider_config; pub use runtime::RuntimeConfig; pub use types::*; -pub use watcher::spawn_file_watcher; +pub use watcher::{FileWatcherHandle, spawn_file_watcher}; // Re-export pub(crate) items that need crate-wide visibility. // (GEMINI_PROVIDER_BASE_URL is only used within config submodules, no re-export needed.) diff --git a/src/config/types.rs b/src/config/types.rs index 65d646c8f..825acd3cd 100644 --- a/src/config/types.rs +++ b/src/config/types.rs @@ -6,6 +6,7 @@ use crate::secrets::store::{InstancePattern, SecretField, SystemSecrets}; use chrono_tz::Tz; use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; use std::collections::HashMap; use std::path::{Path, PathBuf}; @@ -1962,6 +1963,19 @@ pub fn binding_runtime_adapter_key(platform: &str, adapter: Option<&str>) -> Str platform.to_string() } +/// Build the persisted token filename for a named Twitch adapter instance. +pub fn named_twitch_token_file_name(name: &str) -> String { + let safe_name: String = name + .chars() + .map(|ch| if ch.is_ascii_alphanumeric() { ch } else { '_' }) + .take(64) + .collect(); + let hash = Sha256::digest(name.as_bytes()); + let hash_prefix = hex::encode(&hash[..8]); + + format!("twitch_token_{safe_name}_{hash_prefix}.json") +} + /// Match a binding's adapter selector against an inbound message adapter. pub(super) fn binding_adapter_matches(binding: &Binding, message: &crate::InboundMessage) -> bool { match (&binding.adapter, message.adapter_selector()) { diff --git a/src/config/watcher.rs b/src/config/watcher.rs index f197897d5..fa99a26fd 100644 --- a/src/config/watcher.rs +++ b/src/config/watcher.rs @@ -1,10 +1,11 @@ -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use std::sync::Arc; use super::{ Binding, Config, DiscordPermissions, MattermostPermissions, RuntimeConfig, SignalPermissions, SlackPermissions, TelegramPermissions, TwitchPermissions, binding_runtime_adapter_key, }; +use sha2::{Digest, Sha256}; /// Per-agent context needed by the file watcher: (id, prompt_dir, identity_dir, /// runtime_config, mcp_manager). @@ -16,6 +17,17 @@ type WatchedAgent = ( Arc, ); +pub struct FileWatcherHandle { + shutdown_tx: std::sync::mpsc::Sender<()>, + _task: tokio::task::JoinHandle<()>, +} + +impl Drop for FileWatcherHandle { + fn drop(&mut self) { + self.shutdown_tx.send(()).ok(); + } +} + /// Watches config, prompt, identity, and skill files for changes and triggers /// hot reload on the corresponding RuntimeConfig. /// @@ -38,11 +50,12 @@ pub fn spawn_file_watcher( llm_manager: Arc, agent_links: Arc>>, agent_humans: Arc>>, -) -> tokio::task::JoinHandle<()> { +) -> FileWatcherHandle { use notify::{Event, RecursiveMode, Watcher}; use std::time::Duration; - tokio::task::spawn_blocking(move || { + let (shutdown_tx, shutdown_rx) = std::sync::mpsc::channel::<()>(); + let task = tokio::task::spawn_blocking(move || { let (tx, rx) = std::sync::mpsc::channel::(); let mut watcher = match notify::recommended_watcher( @@ -118,10 +131,21 @@ pub fn spawn_file_watcher( // Debounce loop: collect events for 2 seconds, then reload let debounce = Duration::from_secs(2); - while let Ok(first) = rx.recv() { + loop { + if shutdown_rx.try_recv().is_ok() { + break; + } + let first = match rx.recv_timeout(Duration::from_millis(250)) { + Ok(first) => first, + Err(std::sync::mpsc::RecvTimeoutError::Timeout) => continue, + Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => break, + }; // Drain any additional events within the debounce window let mut changed_paths: Vec = first.paths; while let Ok(event) = rx.recv_timeout(debounce) { + if shutdown_rx.try_recv().is_ok() { + return; + } changed_paths.extend(event.paths); } @@ -259,7 +283,8 @@ pub fn spawn_file_watcher( tracing::info!("signal permissions reloaded"); } - // Hot-start adapters that are newly enabled in the config + // Reconcile adapter runtime state with the new config: start + // newly enabled adapters, stop removed ones, restart changed ones. if let Some(ref manager) = messaging_manager { let rt = tokio::runtime::Handle::current(); let manager = manager.clone(); @@ -273,417 +298,25 @@ pub fn spawn_file_watcher( let instance_dir = instance_dir.clone(); rt.spawn(async move { - // Discord: start default + named instances that are enabled and not already running. - if let Some(discord_config) = &config.messaging.discord - && discord_config.enabled { - if !discord_config.token.is_empty() && !manager.has_adapter("discord").await { - let permissions = match discord_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = DiscordPermissions::from_config(discord_config, &config.bindings); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - let adapter = crate::messaging::discord::DiscordAdapter::new( - "discord", - &discord_config.token, - permissions, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start discord adapter from config change"); - } - } - - for instance in discord_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "discord", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - // TODO: named instance permissions are not hot-updated on - // config reload because each instance owns its own - // Arc with no external handle. Fixing this - // requires either a permission-update method on the - // Messaging trait or a shared handle registry. Permissions - // will be correct after a full restart. - continue; - } - - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - DiscordPermissions::from_instance_config(instance, &config.bindings), - )); - let adapter = crate::messaging::discord::DiscordAdapter::new( - runtime_key, - &instance.token, - permissions, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named discord adapter from config change"); - } - } - } - - // Slack: start default + named instances that are enabled and not already running. - if let Some(slack_config) = &config.messaging.slack - && slack_config.enabled { - if !slack_config.bot_token.is_empty() - && !slack_config.app_token.is_empty() - && !manager.has_adapter("slack").await - { - let permissions = match slack_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = SlackPermissions::from_config(slack_config, &config.bindings); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - match crate::messaging::slack::SlackAdapter::new( - "slack", - &slack_config.bot_token, - &slack_config.app_token, - permissions, - slack_config.commands.clone(), - ) { - Ok(adapter) => { - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start slack adapter from config change"); - } - } - Err(error) => { - tracing::error!(%error, "failed to build slack adapter from config change"); - } - } - } - - for instance in slack_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "slack", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - // TODO: named instance permissions not hot-updated (see discord block comment) - continue; - } - - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - SlackPermissions::from_instance_config(instance, &config.bindings), - )); - match crate::messaging::slack::SlackAdapter::new( - runtime_key, - &instance.bot_token, - &instance.app_token, - permissions, - instance.commands.clone(), - ) { - Ok(adapter) => { - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named slack adapter from config change"); - } - } - Err(error) => { - tracing::error!(%error, adapter = %instance.name, "failed to build named slack adapter from config change"); - } - } + match build_desired_configured_adapters( + &config, + &instance_dir, + discord_permissions, + slack_permissions, + telegram_permissions, + twitch_permissions, + mattermost_permissions, + signal_permissions, + ) { + Ok(desired) => { + if let Err(error) = manager.reconcile_configured(desired).await { + tracing::warn!(%error, "messaging adapter reconciliation encountered errors"); } } - - // Telegram: start default + named instances that are enabled and not already running. - if let Some(telegram_config) = &config.messaging.telegram - && telegram_config.enabled { - if !telegram_config.token.is_empty() - && !manager.has_adapter("telegram").await - { - let permissions = match telegram_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = TelegramPermissions::from_config(telegram_config, &config.bindings); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - let adapter = crate::messaging::telegram::TelegramAdapter::new( - "telegram", - &telegram_config.token, - permissions, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start telegram adapter from config change"); - } - } - - for instance in telegram_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "telegram", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - // TODO: named instance permissions not hot-updated (see discord block comment) - continue; - } - - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - TelegramPermissions::from_instance_config(instance, &config.bindings), - )); - let adapter = crate::messaging::telegram::TelegramAdapter::new( - runtime_key, - &instance.token, - permissions, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named telegram adapter from config change"); - } - } - } - - // Email: start default + named instances that are enabled and not already running. - if let Some(email_config) = &config.messaging.email - && email_config.enabled { - if !email_config.imap_host.is_empty() && !manager.has_adapter("email").await { - match crate::messaging::email::EmailAdapter::from_config(email_config) { - Ok(adapter) => { - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start email adapter from config change"); - } - } - Err(error) => { - tracing::error!(%error, "failed to build email adapter from config change"); - } - } - } - - for instance in email_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "email", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - continue; - } - - match crate::messaging::email::EmailAdapter::from_instance_config( - runtime_key.as_str(), - instance, - ) { - Ok(adapter) => { - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named email adapter from config change"); - } - } - Err(error) => { - tracing::error!(%error, adapter = %instance.name, "failed to build named email adapter from config change"); - } - } - } - } - - // Twitch: start default + named instances that are enabled and not already running. - if let Some(twitch_config) = &config.messaging.twitch - && twitch_config.enabled { - if !twitch_config.username.is_empty() - && !twitch_config.oauth_token.is_empty() - && !manager.has_adapter("twitch").await - { - let permissions = match twitch_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = TwitchPermissions::from_config(twitch_config, &config.bindings); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - let token_path = instance_dir.join("twitch_token.json"); - let adapter = crate::messaging::twitch::TwitchAdapter::new( - "twitch", - &twitch_config.username, - &twitch_config.oauth_token, - twitch_config.client_id.clone(), - twitch_config.client_secret.clone(), - twitch_config.refresh_token.clone(), - Some(token_path), - twitch_config.channels.clone(), - twitch_config.trigger_prefix.clone(), - permissions, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start twitch adapter from config change"); - } - } - - for instance in twitch_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "twitch", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - // TODO: named instance permissions not hot-updated (see discord block comment) - continue; - } - - let token_file_name = { - use std::hash::{Hash, Hasher}; - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - instance.name.hash(&mut hasher); - let name_hash = hasher.finish(); - format!( - "twitch_token_{}_{name_hash:016x}.json", - instance - .name - .chars() - .map(|ch| if ch.is_ascii_alphanumeric() { ch } else { '_' }) - .collect::() - ) - }; - let token_path = instance_dir.join(token_file_name); - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - TwitchPermissions::from_instance_config(instance, &config.bindings), - )); - let adapter = crate::messaging::twitch::TwitchAdapter::new( - runtime_key, - &instance.username, - &instance.oauth_token, - instance.client_id.clone(), - instance.client_secret.clone(), - instance.refresh_token.clone(), - Some(token_path), - instance.channels.clone(), - instance.trigger_prefix.clone(), - permissions, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named twitch adapter from config change"); - } - } - } - - // Signal: start default adapter (requires root enabled) and named instances (independent). - // Unlike Discord/Telegram where named instances inherit the root enabled gate, - // Signal named instances start independently when they have valid credentials - // and their own enabled flag is set. This allows running multiple Signal accounts - // without needing a "default" account enabled. - if let Some(signal_config) = &config.messaging.signal { - // Start default adapter only if root is enabled AND has credentials - if signal_config.enabled - && !signal_config.http_url.is_empty() - && !signal_config.account.is_empty() - && !manager.has_adapter("signal").await - { - let permissions = match signal_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = SignalPermissions::from_config(signal_config); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - let tmp_dir = instance_dir.join("tmp"); - let adapter = crate::messaging::signal::SignalAdapter::new( - "signal", - &signal_config.http_url, - &signal_config.account, - signal_config.ignore_stories, - permissions, - tmp_dir, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start signal adapter from config change"); - } - } - - // Start named instances regardless of root enabled flag (as long as config exists) - for instance in signal_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "signal", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - // TODO: named instance permissions not hot-updated (see discord block comment) - continue; - } - // Skip instances with missing credentials (same gating as cold-start) - if instance.http_url.is_empty() || instance.account.is_empty() { - tracing::warn!(adapter = %instance.name, "skipping enabled signal instance with missing credentials"); - continue; - } - - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - SignalPermissions::from_instance_config(instance), - )); - let tmp_dir = instance_dir.join("tmp"); - let adapter = crate::messaging::signal::SignalAdapter::new( - runtime_key, - &instance.http_url, - &instance.account, - instance.ignore_stories, - permissions, - tmp_dir, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named signal adapter from config change"); - } + Err(error) => { + tracing::error!(%error, "failed to build desired messaging adapters from config change"); } } - - // Mattermost: start default + named instances that are enabled and not already running. - if let Some(mattermost_config) = &config.messaging.mattermost - && mattermost_config.enabled { - if !mattermost_config.base_url.is_empty() - && !mattermost_config.token.is_empty() - && !manager.has_adapter("mattermost").await - { - let permissions = match mattermost_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = MattermostPermissions::from_config(mattermost_config, &config.bindings); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - match crate::messaging::mattermost::MattermostAdapter::new( - "mattermost", - &mattermost_config.base_url, - mattermost_config.token.as_str(), - mattermost_config.team_id.as_deref().map(Arc::from), - mattermost_config.max_attachment_bytes, - permissions, - ) { - Ok(adapter) => { - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start mattermost adapter from config change"); - } - } - Err(error) => { - tracing::error!(%error, "failed to build mattermost adapter from config change"); - } - } - } - - for instance in mattermost_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "mattermost", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - continue; - } - - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - MattermostPermissions::from_instance_config(instance, &config.bindings), - )); - match crate::messaging::mattermost::MattermostAdapter::new( - runtime_key, - &instance.base_url, - instance.token.as_str(), - instance.team_id.as_deref().map(Arc::from), - instance.max_attachment_bytes, - permissions, - ) { - Ok(adapter) => { - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named mattermost adapter from config change"); - } - } - Err(error) => { - tracing::error!(%error, adapter = %instance.name, "failed to build named mattermost adapter from config change"); - } - } - } - } }); } } @@ -713,5 +346,627 @@ pub fn spawn_file_watcher( } tracing::info!("file watcher stopped"); - }) + }); + + FileWatcherHandle { + shutdown_tx, + _task: task, + } +} + +#[allow(clippy::too_many_arguments)] +fn build_desired_configured_adapters( + config: &Config, + instance_dir: &Path, + discord_permissions: Option>>, + slack_permissions: Option>>, + telegram_permissions: Option>>, + twitch_permissions: Option>>, + mattermost_permissions: Option>>, + signal_permissions: Option>>, +) -> anyhow::Result> { + let mut desired = Vec::new(); + + if let Some(discord_config) = &config.messaging.discord + && discord_config.enabled + { + if !discord_config.token.is_empty() { + let permissions_snapshot = + DiscordPermissions::from_config(discord_config, &config.bindings); + let permissions = discord_permissions.unwrap_or_else(|| { + Arc::new(arc_swap::ArcSwap::from_pointee( + permissions_snapshot.clone(), + )) + }); + let fingerprint = format!( + "token={}|dm={:?}|allow_bot_messages={}|permissions={}", + secret_fingerprint(&discord_config.token), + sorted_strings(discord_config.dm_allowed_users.clone()), + discord_config.allow_bot_messages, + discord_permissions_fingerprint(&permissions_snapshot) + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + crate::messaging::discord::DiscordAdapter::new( + "discord", + &discord_config.token, + permissions, + ), + fingerprint, + )); + } + + for instance in discord_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.token.is_empty() { + continue; + } + let permissions_snapshot = + DiscordPermissions::from_instance_config(instance, &config.bindings); + let fingerprint = format!( + "token={}|dm={:?}|allow_bot_messages={}|permissions={}", + secret_fingerprint(&instance.token), + sorted_strings(instance.dm_allowed_users.clone()), + instance.allow_bot_messages, + discord_permissions_fingerprint(&permissions_snapshot) + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + crate::messaging::discord::DiscordAdapter::new( + binding_runtime_adapter_key("discord", Some(instance.name.as_str())), + &instance.token, + Arc::new(arc_swap::ArcSwap::from_pointee(permissions_snapshot)), + ), + fingerprint, + )); + } + } + + if let Some(slack_config) = &config.messaging.slack + && slack_config.enabled + { + if !slack_config.bot_token.is_empty() && !slack_config.app_token.is_empty() { + let permissions_snapshot = + SlackPermissions::from_config(slack_config, &config.bindings); + let permissions = slack_permissions.unwrap_or_else(|| { + Arc::new(arc_swap::ArcSwap::from_pointee( + permissions_snapshot.clone(), + )) + }); + let fingerprint = format!( + "bot_token={}|app_token={}|dm={:?}|commands={:?}|permissions={}", + secret_fingerprint(&slack_config.bot_token), + secret_fingerprint(&slack_config.app_token), + sorted_strings(slack_config.dm_allowed_users.clone()), + sorted_slack_commands(slack_config.commands.clone()), + slack_permissions_fingerprint(&permissions_snapshot) + ); + let adapter = crate::messaging::slack::SlackAdapter::new( + "slack", + &slack_config.bot_token, + &slack_config.app_token, + permissions, + slack_config.commands.clone(), + )?; + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + + for instance in slack_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.bot_token.is_empty() || instance.app_token.is_empty() { + continue; + } + let permissions_snapshot = + SlackPermissions::from_instance_config(instance, &config.bindings); + let fingerprint = format!( + "bot_token={}|app_token={}|dm={:?}|commands={:?}|permissions={}", + secret_fingerprint(&instance.bot_token), + secret_fingerprint(&instance.app_token), + sorted_strings(instance.dm_allowed_users.clone()), + sorted_slack_commands(instance.commands.clone()), + slack_permissions_fingerprint(&permissions_snapshot) + ); + let adapter = crate::messaging::slack::SlackAdapter::new( + binding_runtime_adapter_key("slack", Some(instance.name.as_str())), + &instance.bot_token, + &instance.app_token, + Arc::new(arc_swap::ArcSwap::from_pointee(permissions_snapshot)), + instance.commands.clone(), + )?; + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + } + + if let Some(telegram_config) = &config.messaging.telegram + && telegram_config.enabled + { + if !telegram_config.token.is_empty() { + let permissions_snapshot = + TelegramPermissions::from_config(telegram_config, &config.bindings); + let permissions = telegram_permissions.unwrap_or_else(|| { + Arc::new(arc_swap::ArcSwap::from_pointee( + permissions_snapshot.clone(), + )) + }); + let fingerprint = format!( + "token={}|dm={:?}|permissions={}", + secret_fingerprint(&telegram_config.token), + sorted_strings(telegram_config.dm_allowed_users.clone()), + telegram_permissions_fingerprint(&permissions_snapshot) + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + crate::messaging::telegram::TelegramAdapter::new( + "telegram", + &telegram_config.token, + permissions, + ), + fingerprint, + )); + } + + for instance in telegram_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.token.is_empty() { + continue; + } + let permissions_snapshot = + TelegramPermissions::from_instance_config(instance, &config.bindings); + let fingerprint = format!( + "token={}|dm={:?}|permissions={}", + secret_fingerprint(&instance.token), + sorted_strings(instance.dm_allowed_users.clone()), + telegram_permissions_fingerprint(&permissions_snapshot) + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + crate::messaging::telegram::TelegramAdapter::new( + binding_runtime_adapter_key("telegram", Some(instance.name.as_str())), + &instance.token, + Arc::new(arc_swap::ArcSwap::from_pointee(permissions_snapshot)), + ), + fingerprint, + )); + } + } + + if let Some(email_config) = &config.messaging.email + && email_config.enabled + { + if !email_config.imap_host.is_empty() { + let fingerprint = format!( + "imap_host={};imap_port={};imap_username={};imap_password={};imap_use_tls={};smtp_host={};smtp_port={};smtp_username={};smtp_password={};smtp_use_starttls={};from_address={};from_name={:?};poll_interval_secs={};folders={:?};allowed_senders={:?};max_body_bytes={};max_attachment_bytes={}", + email_config.imap_host, + email_config.imap_port, + secret_fingerprint(&email_config.imap_username), + secret_fingerprint(&email_config.imap_password), + email_config.imap_use_tls, + email_config.smtp_host, + email_config.smtp_port, + secret_fingerprint(&email_config.smtp_username), + secret_fingerprint(&email_config.smtp_password), + email_config.smtp_use_starttls, + email_config.from_address, + email_config.from_name, + email_config.poll_interval_secs, + sorted_strings(email_config.folders.clone()), + sorted_strings(email_config.allowed_senders.clone()), + email_config.max_body_bytes, + email_config.max_attachment_bytes + ); + let adapter = crate::messaging::email::EmailAdapter::from_config(email_config)?; + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + + for instance in email_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.imap_host.is_empty() { + continue; + } + let fingerprint = format!( + "name={};imap_host={};imap_port={};imap_username={};imap_password={};imap_use_tls={};smtp_host={};smtp_port={};smtp_username={};smtp_password={};smtp_use_starttls={};from_address={};from_name={:?};poll_interval_secs={};folders={:?};allowed_senders={:?};max_body_bytes={};max_attachment_bytes={}", + instance.name, + instance.imap_host, + instance.imap_port, + secret_fingerprint(&instance.imap_username), + secret_fingerprint(&instance.imap_password), + instance.imap_use_tls, + instance.smtp_host, + instance.smtp_port, + secret_fingerprint(&instance.smtp_username), + secret_fingerprint(&instance.smtp_password), + instance.smtp_use_starttls, + instance.from_address, + instance.from_name, + instance.poll_interval_secs, + sorted_strings(instance.folders.clone()), + sorted_strings(instance.allowed_senders.clone()), + instance.max_body_bytes, + instance.max_attachment_bytes + ); + let adapter = crate::messaging::email::EmailAdapter::from_instance_config( + binding_runtime_adapter_key("email", Some(instance.name.as_str())), + instance, + )?; + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + } + + if let Some(webhook_config) = &config.messaging.webhook + && webhook_config.enabled + { + let fingerprint = format!( + "port={};bind={};auth_token={:?}", + webhook_config.port, + webhook_config.bind, + webhook_config.auth_token.as_deref().map(secret_fingerprint) + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + crate::messaging::webhook::WebhookAdapter::new( + webhook_config.port, + &webhook_config.bind, + webhook_config.auth_token.clone(), + ), + fingerprint, + )); + } + + if let Some(twitch_config) = &config.messaging.twitch + && twitch_config.enabled + { + if !twitch_config.username.is_empty() && !twitch_config.oauth_token.is_empty() { + let permissions_snapshot = + TwitchPermissions::from_config(twitch_config, &config.bindings); + let permissions = twitch_permissions.unwrap_or_else(|| { + Arc::new(arc_swap::ArcSwap::from_pointee( + permissions_snapshot.clone(), + )) + }); + let fingerprint = format!( + "username={};oauth_token={};client_id={:?};client_secret={:?};refresh_token={:?};channels={:?};trigger_prefix={:?};permissions={}", + secret_fingerprint(&twitch_config.username), + secret_fingerprint(&twitch_config.oauth_token), + twitch_config.client_id.as_deref().map(secret_fingerprint), + twitch_config + .client_secret + .as_deref() + .map(secret_fingerprint), + twitch_config + .refresh_token + .as_deref() + .map(secret_fingerprint), + sorted_strings(twitch_config.channels.clone()), + twitch_config.trigger_prefix, + twitch_permissions_fingerprint(&permissions_snapshot) + ); + let adapter = crate::messaging::twitch::TwitchAdapter::new( + "twitch", + &twitch_config.username, + &twitch_config.oauth_token, + twitch_config.client_id.clone(), + twitch_config.client_secret.clone(), + twitch_config.refresh_token.clone(), + Some(instance_dir.join("twitch_token.json")), + twitch_config.channels.clone(), + twitch_config.trigger_prefix.clone(), + permissions, + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + + for instance in twitch_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.username.is_empty() || instance.oauth_token.is_empty() { + continue; + } + let permissions_snapshot = + TwitchPermissions::from_instance_config(instance, &config.bindings); + let fingerprint = format!( + "username={};oauth_token={};client_id={:?};client_secret={:?};refresh_token={:?};channels={:?};trigger_prefix={:?};permissions={}", + secret_fingerprint(&instance.username), + secret_fingerprint(&instance.oauth_token), + instance.client_id.as_deref().map(secret_fingerprint), + instance.client_secret.as_deref().map(secret_fingerprint), + instance.refresh_token.as_deref().map(secret_fingerprint), + sorted_strings(instance.channels.clone()), + instance.trigger_prefix, + twitch_permissions_fingerprint(&permissions_snapshot) + ); + let token_path = + instance_dir.join(crate::config::named_twitch_token_file_name(&instance.name)); + let adapter = crate::messaging::twitch::TwitchAdapter::new( + binding_runtime_adapter_key("twitch", Some(instance.name.as_str())), + &instance.username, + &instance.oauth_token, + instance.client_id.clone(), + instance.client_secret.clone(), + instance.refresh_token.clone(), + Some(token_path), + instance.channels.clone(), + instance.trigger_prefix.clone(), + Arc::new(arc_swap::ArcSwap::from_pointee(permissions_snapshot)), + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + } + + if let Some(mattermost_config) = &config.messaging.mattermost + && mattermost_config.enabled + { + if !mattermost_config.base_url.is_empty() && !mattermost_config.token.is_empty() { + let permissions_snapshot = + MattermostPermissions::from_config(mattermost_config, &config.bindings); + let permissions = mattermost_permissions.unwrap_or_else(|| { + Arc::new(arc_swap::ArcSwap::from_pointee( + permissions_snapshot.clone(), + )) + }); + let fingerprint = format!( + "base_url={};token={};team_id={:?};max_attachment_bytes={};permissions={}", + mattermost_config.base_url, + secret_fingerprint(&mattermost_config.token), + mattermost_config.team_id, + mattermost_config.max_attachment_bytes, + mattermost_permissions_fingerprint(&permissions_snapshot) + ); + match crate::messaging::mattermost::MattermostAdapter::new( + "mattermost", + &mattermost_config.base_url, + mattermost_config.token.as_str(), + mattermost_config.team_id.as_deref().map(Arc::from), + mattermost_config.max_attachment_bytes, + permissions, + ) { + Ok(adapter) => { + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + Err(error) => { + tracing::error!(%error, "failed to build mattermost adapter from config change"); + } + } + } + + for instance in mattermost_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.base_url.is_empty() || instance.token.is_empty() { + tracing::warn!(adapter = %instance.name, "skipping enabled mattermost instance with missing credentials"); + continue; + } + let permissions_snapshot = + MattermostPermissions::from_instance_config(instance, &config.bindings); + let fingerprint = format!( + "base_url={};token={};team_id={:?};max_attachment_bytes={};permissions={}", + instance.base_url, + secret_fingerprint(&instance.token), + instance.team_id, + instance.max_attachment_bytes, + mattermost_permissions_fingerprint(&permissions_snapshot) + ); + match crate::messaging::mattermost::MattermostAdapter::new( + binding_runtime_adapter_key("mattermost", Some(instance.name.as_str())), + &instance.base_url, + instance.token.as_str(), + instance.team_id.as_deref().map(Arc::from), + instance.max_attachment_bytes, + Arc::new(arc_swap::ArcSwap::from_pointee(permissions_snapshot)), + ) { + Ok(adapter) => { + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + Err(error) => { + tracing::error!(%error, adapter = %instance.name, "failed to build named mattermost adapter from config change"); + } + } + } + } + + // Signal named instances start independently of the root enabled flag, + // matching cold-start: multiple Signal accounts can run without a + // "default" account being enabled. + if let Some(signal_config) = &config.messaging.signal { + let tmp_dir = instance_dir.join("tmp"); + if signal_config.enabled + && !signal_config.http_url.is_empty() + && !signal_config.account.is_empty() + { + let permissions_snapshot = SignalPermissions::from_config(signal_config); + let permissions = signal_permissions.unwrap_or_else(|| { + Arc::new(arc_swap::ArcSwap::from_pointee( + permissions_snapshot.clone(), + )) + }); + let fingerprint = format!( + "http_url={};account={};ignore_stories={};permissions={}", + secret_fingerprint(&signal_config.http_url), + secret_fingerprint(&signal_config.account), + signal_config.ignore_stories, + signal_permissions_fingerprint(&permissions_snapshot) + ); + let adapter = crate::messaging::signal::SignalAdapter::new( + "signal", + &signal_config.http_url, + &signal_config.account, + signal_config.ignore_stories, + permissions, + tmp_dir.clone(), + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + + for instance in signal_config + .instances + .iter() + .filter(|instance| instance.enabled) + { + if instance.http_url.is_empty() || instance.account.is_empty() { + tracing::warn!(adapter = %instance.name, "skipping enabled signal instance with missing credentials"); + continue; + } + let permissions_snapshot = SignalPermissions::from_instance_config(instance); + let fingerprint = format!( + "http_url={};account={};ignore_stories={};permissions={}", + secret_fingerprint(&instance.http_url), + secret_fingerprint(&instance.account), + instance.ignore_stories, + signal_permissions_fingerprint(&permissions_snapshot) + ); + let adapter = crate::messaging::signal::SignalAdapter::new( + binding_runtime_adapter_key("signal", Some(instance.name.as_str())), + &instance.http_url, + &instance.account, + instance.ignore_stories, + Arc::new(arc_swap::ArcSwap::from_pointee(permissions_snapshot)), + tmp_dir.clone(), + ); + desired.push(crate::messaging::ConfiguredAdapter::new( + adapter, + fingerprint, + )); + } + } + + Ok(desired) +} + +fn sorted_strings(mut values: Vec) -> Vec { + values.sort(); + values +} + +fn sorted_u64s(mut values: Vec) -> Vec { + values.sort_unstable(); + values +} + +fn sorted_i64s(mut values: Vec) -> Vec { + values.sort_unstable(); + values +} + +fn format_u64_map(map: &std::collections::HashMap>) -> String { + let mut entries = map + .iter() + .map(|(key, values)| (*key, sorted_u64s(values.clone()))) + .collect::>(); + entries.sort_by_key(|(key, _)| *key); + format!("{entries:?}") +} + +fn format_string_map(map: &std::collections::HashMap>) -> String { + let mut entries = map + .iter() + .map(|(key, values)| (key.clone(), sorted_strings(values.clone()))) + .collect::>(); + entries.sort_by(|left, right| left.0.cmp(&right.0)); + format!("{entries:?}") +} + +fn discord_permissions_fingerprint(permissions: &DiscordPermissions) -> String { + format!( + "guild_filter={:?};channel_filter={};dm_allowed_users={:?};allow_bot_messages={}", + permissions.guild_filter.clone().map(sorted_u64s), + format_u64_map(&permissions.channel_filter), + sorted_u64s(permissions.dm_allowed_users.clone()), + permissions.allow_bot_messages + ) +} + +fn slack_permissions_fingerprint(permissions: &SlackPermissions) -> String { + format!( + "workspace_filter={:?};channel_filter={};dm_allowed_users={:?}", + permissions.workspace_filter.clone().map(sorted_strings), + format_string_map(&permissions.channel_filter), + sorted_strings(permissions.dm_allowed_users.clone()) + ) +} + +fn telegram_permissions_fingerprint(permissions: &TelegramPermissions) -> String { + format!( + "chat_filter={:?};dm_allowed_users={:?}", + permissions.chat_filter.clone().map(sorted_i64s), + sorted_i64s(permissions.dm_allowed_users.clone()) + ) +} + +fn twitch_permissions_fingerprint(permissions: &TwitchPermissions) -> String { + format!( + "channel_filter={:?};allowed_users={:?}", + permissions.channel_filter.clone().map(sorted_strings), + sorted_strings(permissions.allowed_users.clone()) + ) +} + +fn mattermost_permissions_fingerprint(permissions: &MattermostPermissions) -> String { + format!( + "team_filter={:?};channel_filter={};dm_allowed_users={:?}", + permissions.team_filter.clone().map(sorted_strings), + format_string_map(&permissions.channel_filter), + sorted_strings(permissions.dm_allowed_users.clone()) + ) +} + +fn signal_permissions_fingerprint(permissions: &SignalPermissions) -> String { + format!( + "group_filter={:?};dm_allowed_users={:?};group_allowed_users={:?}", + permissions.group_filter.clone().map(sorted_strings), + sorted_strings(permissions.dm_allowed_users.clone()), + sorted_strings(permissions.group_allowed_users.clone()) + ) +} + +fn secret_fingerprint(value: &str) -> String { + let digest = Sha256::digest(value.as_bytes()); + hex::encode(&digest[..8]) +} + +fn sorted_slack_commands( + mut commands: Vec, +) -> Vec<(String, String, Option)> { + let mut normalized = commands + .drain(..) + .map(|command| (command.command, command.agent_id, command.description)) + .collect::>(); + normalized.sort_by(|left, right| left.0.cmp(&right.0).then(left.1.cmp(&right.1))); + normalized } diff --git a/src/main.rs b/src/main.rs index da784382e..4cd975f39 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1790,7 +1790,7 @@ async fn run( let mut agents_initialized = false; // File watcher handle — started after agent init (or in setup mode with empty data) - let mut _file_watcher; + let mut _file_watcher: Option; // If providers are available, initialize agents immediately if has_providers { @@ -1833,7 +1833,7 @@ async fn run( agents_initialized = true; // Start file watcher with populated agent data - _file_watcher = spacebot::config::spawn_file_watcher( + _file_watcher = Some(spacebot::config::spawn_file_watcher( config_path.clone(), config.instance_dir.clone(), watcher_agents, @@ -1848,10 +1848,10 @@ async fn run( llm_manager.clone(), agent_links.clone(), agent_humans.clone(), - ); + )); } else { // Start file watcher in setup mode (no agents to watch yet) - _file_watcher = spacebot::config::spawn_file_watcher( + _file_watcher = Some(spacebot::config::spawn_file_watcher( config_path.clone(), config.instance_dir.clone(), Vec::new(), @@ -1866,7 +1866,7 @@ async fn run( llm_manager.clone(), agent_links.clone(), agent_humans.clone(), - ); + )); } if foreground { @@ -2575,6 +2575,7 @@ async fn run( { Ok(new_llm) => { let new_llm_manager = Arc::new(new_llm); + api_state.set_llm_manager(new_llm_manager.clone()).await; // Update agent_humans from the reloaded config // before initialize_agents so agents see the // latest [[humans]] entries. @@ -2617,7 +2618,8 @@ async fn run( Ok(()) => { agents_initialized = true; // Restart file watcher with the new agent data - _file_watcher = spacebot::config::spawn_file_watcher( + let _old_watcher = _file_watcher.take(); + _file_watcher = Some(spacebot::config::spawn_file_watcher( config_path.clone(), new_config.instance_dir.clone(), new_watcher_agents, @@ -2632,7 +2634,7 @@ async fn run( new_llm_manager.clone(), agent_links.clone(), agent_humans.clone(), - ); + )); tracing::info!("agents initialized after provider setup"); } Err(error) => { @@ -3451,20 +3453,7 @@ async fn initialize_agents( "twitch", Some(instance.name.as_str()), ); - let token_file_name = { - use std::hash::{Hash, Hasher}; - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - instance.name.hash(&mut hasher); - let name_hash = hasher.finish(); - format!( - "twitch_token_{}_{name_hash:016x}.json", - instance - .name - .chars() - .map(|ch| if ch.is_ascii_alphanumeric() { ch } else { '_' }) - .collect::() - ) - }; + let token_file_name = spacebot::config::named_twitch_token_file_name(&instance.name); let token_path = config.instance_dir.join(token_file_name); let perms = Arc::new(ArcSwap::from_pointee( spacebot::config::TwitchPermissions::from_instance_config( @@ -3626,6 +3615,10 @@ async fn initialize_agents( } } + new_messaging_manager + .seed_configured_fingerprints_from_registered() + .await; + let portal_agent_pools = agents .iter() .map(|(agent_id, agent)| (agent_id.to_string(), agent.db.sqlite.clone())) diff --git a/src/messaging.rs b/src/messaging.rs index b3c434e22..e92536683 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -13,6 +13,6 @@ pub mod traits; pub mod twitch; pub mod webhook; -pub use manager::MessagingManager; +pub use manager::{ConfiguredAdapter, MessagingManager}; pub use traits::Messaging; pub use traits::apply_runtime_adapter_to_conversation_id; diff --git a/src/messaging/manager.rs b/src/messaging/manager.rs index f38e43d44..79276d658 100644 --- a/src/messaging/manager.rs +++ b/src/messaging/manager.rs @@ -10,7 +10,41 @@ use anyhow::Context as _; use futures::StreamExt as _; use std::collections::HashMap; use std::sync::Arc; -use tokio::sync::{RwLock, mpsc}; +use std::sync::atomic::{AtomicBool, Ordering}; +use tokio::sync::{RwLock, mpsc, oneshot, watch}; +use tokio::task::JoinHandle; + +#[cfg(test)] +const INITIAL_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(10); +#[cfg(not(test))] +const INITIAL_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); + +#[cfg(test)] +const MAX_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(100); +#[cfg(not(test))] +const MAX_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(60); + +struct AdapterRuntime { + shutdown_tx: watch::Sender, + task: JoinHandle<()>, +} + +pub struct ConfiguredAdapter { + pub name: String, + pub fingerprint: String, + adapter: Arc, +} + +impl ConfiguredAdapter { + pub fn new(adapter: impl Messaging, fingerprint: impl Into) -> Self { + let name = adapter.name().to_string(); + Self { + name, + fingerprint: fingerprint.into(), + adapter: Arc::new(adapter), + } + } +} /// Manages all messaging adapters with support for runtime addition. /// @@ -18,10 +52,14 @@ use tokio::sync::{RwLock, mpsc}; /// can be registered after `start()` without replacing the inbound stream. pub struct MessagingManager { adapters: RwLock>>, + runtimes: RwLock>, + configured_fingerprints: RwLock>, + lifecycle_mutex: tokio::sync::Mutex<()>, /// Sender side of the fan-in channel. Cloned for each adapter's forwarding task. fan_in_tx: mpsc::Sender, /// Receiver side, taken once by `start()`. fan_in_rx: RwLock>>, + started: AtomicBool, } impl MessagingManager { @@ -29,27 +67,55 @@ impl MessagingManager { let (fan_in_tx, fan_in_rx) = mpsc::channel(512); Self { adapters: RwLock::new(HashMap::new()), + runtimes: RwLock::new(HashMap::new()), + configured_fingerprints: RwLock::new(HashMap::new()), + lifecycle_mutex: tokio::sync::Mutex::new(()), fan_in_tx, fan_in_rx: RwLock::new(Some(fan_in_rx)), + started: AtomicBool::new(false), } } /// Register an adapter (before start). Use `register_and_start` for runtime addition. pub async fn register(&self, adapter: impl Messaging) { + let _lifecycle_guard = self.lifecycle_mutex.lock().await; let name = adapter.name().to_string(); + let adapter: Arc = Arc::new(adapter); tracing::info!(adapter = %name, "registered messaging adapter"); - self.adapters.write().await.insert(name, Arc::new(adapter)); + let started = self.started.load(Ordering::SeqCst); + let old_adapter = self + .adapters + .write() + .await + .insert(name.clone(), Arc::clone(&adapter)); + if started && let Err(error) = self.stop_runtime(&name, old_adapter, false).await { + tracing::warn!(adapter = %name, %error, "failed to stop old runtime during register"); + } + if started && let Err(error) = self.start_runtime(name.clone(), adapter, false).await { + tracing::warn!(adapter = %name, %error, "failed to start adapter registered after manager start"); + } } /// Register a pre-wrapped adapter that the caller retains a handle to. pub async fn register_shared(&self, adapter: Arc) { + let _lifecycle_guard = self.lifecycle_mutex.lock().await; let name = adapter.name().to_string(); tracing::info!(adapter = %name, "registered messaging adapter (shared)"); - self.adapters.write().await.insert(name, adapter); + let adapter: Arc = adapter; + let started = self.started.load(Ordering::SeqCst); + let old_adapter = self + .adapters + .write() + .await + .insert(name.clone(), Arc::clone(&adapter)); + if started && let Err(error) = self.stop_runtime(&name, old_adapter, false).await { + tracing::warn!(adapter = %name, %error, "failed to stop old runtime during shared register"); + } + if started && let Err(error) = self.start_runtime(name.clone(), adapter, false).await { + tracing::warn!(adapter = %name, %error, "failed to start shared adapter registered after manager start"); + } } - /// Maximum number of retry attempts for failed adapters before giving up. - const MAX_RETRY_ATTEMPTS: u32 = 12; /// Maximum number of proactive-send retry attempts for transient broadcast failures. const MAX_BROADCAST_RETRY_ATTEMPTS: u32 = 3; #[cfg(test)] @@ -68,25 +134,15 @@ impl MessagingManager { /// Adapters that fail to start (e.g. due to network not being ready) are /// retried in the background with exponential backoff. pub async fn start(&self) -> crate::Result { - let adapters = self.adapters.read().await; - for (name, adapter) in adapters.iter() { - match adapter.start().await { - Ok(stream) => Self::spawn_forwarder(name.clone(), stream, self.fan_in_tx.clone()), - Err(error) => { - tracing::warn!( - adapter = %name, - %error, - "adapter failed to start, will retry in background" - ); - Self::spawn_retry_task( - name.clone(), - Arc::clone(adapter), - self.fan_in_tx.clone(), - ); - } - } - } - drop(adapters); + let _lifecycle_guard = self.lifecycle_mutex.lock().await; + self.started.store(true, Ordering::SeqCst); + let adapters = self + .adapters + .read() + .await + .iter() + .map(|(name, adapter)| (name.clone(), Arc::clone(adapter))) + .collect::>(); let receiver = self .fan_in_rx @@ -95,6 +151,10 @@ impl MessagingManager { .take() .context("start() already called")?; + for (name, adapter) in adapters { + self.start_runtime(name, adapter, false).await?; + } + Ok(Box::pin(tokio_stream::wrappers::ReceiverStream::new( receiver, ))) @@ -106,29 +166,18 @@ impl MessagingManager { /// channel, so the main loop's stream receives messages without any /// stream replacement or restart. pub async fn register_and_start(&self, adapter: impl Messaging) -> crate::Result<()> { + let _lifecycle_guard = self.lifecycle_mutex.lock().await; let name = adapter.name().to_string(); - - // Shut down existing adapter with the same name if present - { - let adapters = self.adapters.read().await; - if let Some(existing) = adapters.get(&name) { - tracing::info!(adapter = %name, "shutting down existing adapter before replacement"); - if let Err(error) = existing.shutdown().await { - tracing::warn!(adapter = %name, %error, "failed to shut down existing adapter"); - } - } - } - let adapter: Arc = Arc::new(adapter); - - let stream = adapter - .start() + let old_adapter = self + .adapters + .write() .await - .with_context(|| format!("failed to start adapter '{name}'"))?; - Self::spawn_forwarder(name.clone(), stream, self.fan_in_tx.clone()); - - self.adapters.write().await.insert(name.clone(), adapter); - + .insert(name.clone(), Arc::clone(&adapter)); + if let Err(error) = self.stop_runtime(&name, old_adapter, false).await { + tracing::warn!(adapter = %name, %error, "failed to stop old runtime during register_and_start"); + } + self.start_runtime(name.clone(), adapter, true).await?; tracing::info!(adapter = %name, "adapter registered and started at runtime"); Ok(()) } @@ -153,70 +202,280 @@ impl MessagingManager { self.adapters.read().await.keys().cloned().collect() } - /// Spawn a background task that retries starting a failed adapter with exponential backoff. - /// - /// Once the adapter starts successfully, its stream is forwarded into the - /// existing fan-in channel — the same mechanism used by `register_and_start`. - fn spawn_retry_task( + pub async fn reconcile_configured(&self, desired: Vec) -> crate::Result<()> { + let _lifecycle_guard = self.lifecycle_mutex.lock().await; + let desired_names = desired + .iter() + .map(|adapter| adapter.name.clone()) + .collect::>(); + let current_configured = self + .configured_fingerprints + .read() + .await + .keys() + .cloned() + .collect::>(); + + let mut first_error = None; + for name in current_configured { + if !desired_names.contains(&name) + && let Err(error) = self.remove_adapter_inner(&name, true).await + { + tracing::warn!(adapter = %name, %error, "failed to remove stale adapter during reconciliation"); + if first_error.is_none() { + first_error = Some(error); + } + } + } + + let current_fingerprints = self.configured_fingerprints.read().await.clone(); + for desired_adapter in desired { + let is_unchanged = current_fingerprints + .get(&desired_adapter.name) + .is_some_and(|fingerprint| fingerprint == &desired_adapter.fingerprint); + if is_unchanged { + continue; + } + + let old_adapter = self.adapters.write().await.insert( + desired_adapter.name.clone(), + Arc::clone(&desired_adapter.adapter), + ); + if let Err(error) = self + .stop_runtime(&desired_adapter.name, old_adapter, false) + .await + { + tracing::warn!( + adapter = %desired_adapter.name, + %error, + "failed to stop old runtime during reconciliation" + ); + } + let replace_result = self + .start_runtime( + desired_adapter.name.clone(), + Arc::clone(&desired_adapter.adapter), + true, + ) + .await; + self.configured_fingerprints.write().await.insert( + desired_adapter.name.clone(), + desired_adapter.fingerprint.clone(), + ); + if let Err(error) = replace_result { + tracing::warn!( + adapter = %desired_adapter.name, + %error, + "adapter reconciliation failed; supervisor will keep retrying" + ); + if first_error.is_none() { + first_error = Some(error); + } + } + } + + if let Some(error) = first_error { + Err(error) + } else { + Ok(()) + } + } + + fn spawn_supervisor( name: String, adapter: Arc, fan_in_tx: mpsc::Sender, - ) { + mut shutdown_rx: watch::Receiver, + first_start_result_tx: Option>>, + ) -> JoinHandle<()> { tokio::spawn(async move { - let mut delay = std::time::Duration::from_secs(5); - let max_delay = std::time::Duration::from_secs(60); + let mut next_retry_delay = INITIAL_RETRY_DELAY; + let mut first_start_result_tx = first_start_result_tx; + #[cfg(test)] + let healthy_stream_duration = std::time::Duration::from_millis(20); + #[cfg(not(test))] + let healthy_stream_duration = std::time::Duration::from_secs(2); - for attempt in 1..=Self::MAX_RETRY_ATTEMPTS { - tokio::time::sleep(delay).await; + loop { + let start_result = tokio::select! { + _ = shutdown_rx.changed() => break, + result = adapter.start() => result, + }; - match adapter.start().await { + let mut stream = match start_result { Ok(stream) => { - tracing::info!( - adapter = %name, - attempt, - "adapter started successfully after retry" - ); - Self::spawn_forwarder(name, stream, fan_in_tx); - return; + tracing::info!(adapter = %name, "adapter started successfully"); + if let Some(tx) = first_start_result_tx.take() { + let _ = tx.send(Ok(())); + } + stream } Err(error) => { + if let Some(tx) = first_start_result_tx.take() { + let _ = tx.send(Err(error.to_string())); + } tracing::warn!( adapter = %name, - attempt, - max_attempts = Self::MAX_RETRY_ATTEMPTS, %error, - "adapter retry failed, next attempt in {:?}", - delay.min(max_delay) + retry_delay = ?next_retry_delay, + "adapter start failed, retrying in background" ); + let current_delay = next_retry_delay; + next_retry_delay = (next_retry_delay * 2).min(MAX_RETRY_DELAY); + tokio::select! { + _ = shutdown_rx.changed() => break, + _ = tokio::time::sleep(current_delay) => continue, + } + } + }; + + let started_at = std::time::Instant::now(); + let mut observed_message = false; + loop { + tokio::select! { + _ = shutdown_rx.changed() => { + tracing::info!(adapter = %name, "adapter supervisor shutting down"); + return; + } + maybe_message = stream.next() => { + match maybe_message { + Some(message) => { + observed_message = true; + next_retry_delay = INITIAL_RETRY_DELAY; + if fan_in_tx.send(message).await.is_err() { + tracing::warn!(adapter = %name, "fan-in channel closed, stopping supervisor"); + return; + } + } + None => { + if !observed_message && started_at.elapsed() < healthy_stream_duration { + let current_delay = next_retry_delay; + next_retry_delay = (next_retry_delay * 2).min(MAX_RETRY_DELAY); + tracing::warn!( + adapter = %name, + retry_delay = ?current_delay, + "adapter stream ended before becoming healthy, backing off before restart" + ); + tokio::select! { + _ = shutdown_rx.changed() => return, + _ = tokio::time::sleep(current_delay) => {} + } + } else { + next_retry_delay = INITIAL_RETRY_DELAY; + tracing::warn!(adapter = %name, "adapter stream ended, restarting"); + } + break; + } + } + } } } + } + }) + } - delay = (delay * 2).min(max_delay); + async fn start_runtime( + &self, + name: String, + adapter: Arc, + surface_start_error: bool, + ) -> crate::Result<()> { + if !self.started.load(Ordering::SeqCst) { + return Ok(()); + } + + let first_start_result_rx = if surface_start_error { + let (tx, rx) = oneshot::channel(); + self.install_supervisor(name.clone(), adapter, Some(tx)) + .await; + Some(rx) + } else { + self.install_supervisor(name.clone(), adapter, None).await; + None + }; + + if let Some(rx) = first_start_result_rx { + match rx.await { + Ok(Ok(())) => {} + Ok(Err(error)) => { + return Err(anyhow::anyhow!("failed to start adapter '{name}': {error}").into()); + } + Err(error) => { + return Err(anyhow::anyhow!( + "failed to receive initial start result for '{name}': {error}" + ) + .into()); + } } + } - tracing::error!( - adapter = %name, - "adapter failed to start after {} attempts, giving up", - Self::MAX_RETRY_ATTEMPTS - ); - }); + self.configured_fingerprints + .write() + .await + .entry(name) + .or_insert_with(String::new); + Ok(()) } - /// Spawn a task that forwards messages from an adapter stream into the fan-in channel. - fn spawn_forwarder( + async fn install_supervisor( + &self, name: String, - mut stream: InboundStream, - fan_in_tx: mpsc::Sender, - ) { - tokio::spawn(async move { - while let Some(message) = stream.next().await { - if fan_in_tx.send(message).await.is_err() { - tracing::warn!(adapter = %name, "fan-in channel closed, stopping forwarder"); - break; + adapter: Arc, + first_start_result_tx: Option>>, + ) -> watch::Sender { + let (shutdown_tx, shutdown_rx) = watch::channel(false); + let task = Self::spawn_supervisor( + name.clone(), + adapter, + self.fan_in_tx.clone(), + shutdown_rx, + first_start_result_tx, + ); + self.runtimes.write().await.insert( + name, + AdapterRuntime { + shutdown_tx: shutdown_tx.clone(), + task, + }, + ); + shutdown_tx + } + + async fn stop_runtime( + &self, + name: &str, + adapter: Option>, + propagate_shutdown_error: bool, + ) -> crate::Result<()> { + let runtime = self.runtimes.write().await.remove(name); + if let Some(runtime) = runtime.as_ref() { + runtime.shutdown_tx.send(true).ok(); + } + + let mut shutdown_error = None; + if let Some(adapter) = adapter + && let Err(error) = adapter.shutdown().await + { + tracing::warn!(adapter = %name, %error, "failed to shut down adapter"); + shutdown_error = Some(error); + } + + if let Some(runtime) = runtime { + runtime.task.abort(); + match runtime.task.await { + Ok(()) => {} + Err(error) if error.is_cancelled() => {} + Err(error) => { + tracing::warn!(adapter = %name, %error, "adapter supervisor join failed"); } } - tracing::info!(adapter = %name, "adapter stream ended"); - }); + } + + if propagate_shutdown_error && let Some(error) = shutdown_error { + return Err(error); + } + + Ok(()) } /// Inject a message directly into the fan-in channel, bypassing adapter streams. @@ -363,11 +622,20 @@ impl MessagingManager { /// Remove and shut down a single adapter by name. pub async fn remove_adapter(&self, name: &str) -> crate::Result<()> { + let _lifecycle_guard = self.lifecycle_mutex.lock().await; + self.remove_adapter_inner(name, true).await + } + + async fn remove_adapter_inner( + &self, + name: &str, + propagate_shutdown_error: bool, + ) -> crate::Result<()> { let adapter = self.adapters.write().await.remove(name); - if let Some(adapter) = adapter { - adapter.shutdown().await?; - tracing::info!(adapter = %name, "adapter removed and shut down"); - } + self.configured_fingerprints.write().await.remove(name); + self.stop_runtime(name, adapter, propagate_shutdown_error) + .await?; + tracing::info!(adapter = %name, "adapter removed and shut down"); Ok(()) } @@ -393,13 +661,33 @@ impl MessagingManager { /// Shut down all adapters gracefully. pub async fn shutdown(&self) { - let adapters = self.adapters.read().await; - for (name, adapter) in adapters.iter() { - if let Err(error) = adapter.shutdown().await { + let names = self + .adapters + .read() + .await + .keys() + .cloned() + .collect::>(); + for name in names { + if let Err(error) = self.remove_adapter(&name).await { tracing::warn!(adapter = %name, %error, "failed to shut down adapter"); } } } + + pub async fn seed_configured_fingerprints_from_registered(&self) { + let names = self + .adapters + .read() + .await + .keys() + .cloned() + .collect::>(); + let mut fingerprints = self.configured_fingerprints.write().await; + for name in names { + fingerprints.entry(name).or_insert_with(String::new); + } + } } impl Default for MessagingManager { @@ -410,6 +698,188 @@ impl Default for MessagingManager { #[cfg(test)] mod tests { + use super::MessagingManager; + use crate::messaging::traits::{InboundStream, Messaging}; + use crate::{InboundMessage, MessageContent}; + use futures::StreamExt as _; + use std::collections::HashMap; + use std::sync::Arc; + use std::sync::atomic::{AtomicUsize, Ordering}; + + enum StartPlan { + Error(&'static str), + Stream(Vec<&'static str>), + Pending, + } + + struct TestAdapter { + name: String, + plans: std::sync::Mutex>, + start_calls: AtomicUsize, + shutdown_calls: AtomicUsize, + start_signal_tx: tokio::sync::watch::Sender, + } + + impl TestAdapter { + fn new(name: &str, plans: Vec) -> Self { + let (start_signal_tx, _start_signal_rx) = tokio::sync::watch::channel(0); + Self { + name: name.to_string(), + plans: std::sync::Mutex::new(plans), + start_calls: AtomicUsize::new(0), + shutdown_calls: AtomicUsize::new(0), + start_signal_tx, + } + } + + fn subscribe_start_calls(&self) -> tokio::sync::watch::Receiver { + self.start_signal_tx.subscribe() + } + } + + impl Messaging for TestAdapter { + fn name(&self) -> &str { + &self.name + } + + async fn start(&self) -> crate::Result { + let start_count = self.start_calls.fetch_add(1, Ordering::SeqCst) + 1; + self.start_signal_tx.send_replace(start_count); + let plan = self.plans.lock().expect("plans lock").remove(0); + match plan { + StartPlan::Error(error) => Err(anyhow::anyhow!(error).into()), + StartPlan::Stream(messages) => { + let adapter_name = self.name.clone(); + let stream = + tokio_stream::iter(messages.into_iter().map(move |text| InboundMessage { + id: uuid::Uuid::new_v4().to_string(), + source: adapter_name.clone(), + adapter: Some(adapter_name.clone()), + conversation_id: format!("{adapter_name}:test"), + sender_id: "sender".to_string(), + agent_id: None, + content: MessageContent::Text(text.to_string()), + timestamp: chrono::Utc::now(), + metadata: HashMap::new(), + formatted_author: None, + })); + Ok(Box::pin(stream)) + } + StartPlan::Pending => Ok(Box::pin(futures::stream::pending())), + } + } + + async fn respond( + &self, + _message: &InboundMessage, + _response: crate::OutboundResponse, + ) -> crate::Result<()> { + Ok(()) + } + + async fn health_check(&self) -> crate::Result<()> { + Ok(()) + } + + async fn shutdown(&self) -> crate::Result<()> { + self.shutdown_calls.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + } + + #[tokio::test] + async fn supervisor_retries_after_start_failure_and_stream_end() { + let manager = MessagingManager::new(); + let adapter = Arc::new(TestAdapter::new( + "test", + vec![ + StartPlan::Error("boom"), + StartPlan::Stream(vec!["first"]), + StartPlan::Stream(vec!["second"]), + StartPlan::Pending, + ], + )); + let mut start_calls = adapter.subscribe_start_calls(); + manager.register_shared(adapter.clone()).await; + let mut stream = manager.start().await.expect("start manager"); + + let first = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next()) + .await + .expect("timeout") + .expect("message"); + let second = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next()) + .await + .expect("timeout") + .expect("message"); + + match first.content { + MessageContent::Text(text) => assert_eq!(text, "first"), + other => panic!("unexpected message content: {other:?}"), + } + match second.content { + MessageContent::Text(text) => assert_eq!(text, "second"), + other => panic!("unexpected message content: {other:?}"), + } + wait_for_start_calls(&mut start_calls, 3).await; + } + + #[tokio::test] + async fn removing_adapter_stops_retry_supervisor() { + let manager = MessagingManager::new(); + let adapter = Arc::new(TestAdapter::new( + "retrying", + vec![ + StartPlan::Error("boom"), + StartPlan::Error("boom"), + StartPlan::Error("boom"), + StartPlan::Pending, + ], + )); + let mut start_calls = adapter.subscribe_start_calls(); + manager.register_shared(adapter.clone()).await; + let _stream = manager.start().await.expect("start manager"); + + wait_for_start_calls(&mut start_calls, 2).await; + let before_remove = *start_calls.borrow_and_update(); + manager + .remove_adapter("retrying") + .await + .expect("remove adapter"); + let no_more_starts = + tokio::time::timeout(std::time::Duration::from_millis(200), start_calls.changed()) + .await; + + assert!( + no_more_starts.is_err(), + "unexpected extra start after remove" + ); + assert_eq!(before_remove, adapter.start_calls.load(Ordering::SeqCst)); + assert!(adapter.shutdown_calls.load(Ordering::SeqCst) >= 1); + } + + async fn wait_for_start_calls( + receiver: &mut tokio::sync::watch::Receiver, + expected: usize, + ) { + if *receiver.borrow() >= expected { + return; + } + + tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + receiver.changed().await.expect("start signal"); + if *receiver.borrow() >= expected { + break; + } + } + }) + .await + .expect("timed out waiting for start calls"); + } +} + +#[cfg(test)] +mod broadcast_tests { use super::MessagingManager; use crate::messaging::traits::{ InboundStream, Messaging, mark_permanent_broadcast, mark_retryable_broadcast,