From 0b2b102133a2a2d472aff23660f1b8a9aef6466d Mon Sep 17 00:00:00 2001 From: Olen Latham Date: Sun, 8 Mar 2026 15:54:07 -0500 Subject: [PATCH 1/5] fix(messaging): reconcile adapter runtime state --- docs/content/docs/(configuration)/config.mdx | 3 +- src/config/types.rs | 13 + src/config/watcher.rs | 729 ++++++++++++------- src/main.rs | 32 +- src/messaging.rs | 2 +- src/messaging/manager.rs | 561 +++++++++++--- 6 files changed, 937 insertions(+), 403 deletions(-) diff --git a/docs/content/docs/(configuration)/config.mdx b/docs/content/docs/(configuration)/config.mdx index e85ab5a7d..ba6dda684 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/config/types.rs b/src/config/types.rs index 6c134457f..a905de799 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}; @@ -1554,6 +1555,18 @@ 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 { '_' }) + .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 e39ec1533..95d5e44ca 100644 --- a/src/config/watcher.rs +++ b/src/config/watcher.rs @@ -1,4 +1,4 @@ -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use std::sync::Arc; use super::{ @@ -240,7 +240,6 @@ pub fn spawn_file_watcher( tracing::info!("twitch permissions reloaded"); } - // Hot-start adapters that are newly enabled in the config if let Some(ref manager) = messaging_manager { let rt = tokio::runtime::Handle::current(); let manager = manager.clone(); @@ -252,285 +251,23 @@ 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"); - } - } - } - } - - // 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"); - } + match build_desired_configured_adapters( + &config, + &instance_dir, + discord_permissions, + slack_permissions, + telegram_permissions, + twitch_permissions, + ) { + Ok(desired) => { + if let Err(error) = manager.reconcile_configured(desired).await { + tracing::warn!(%error, "messaging adapter reconciliation encountered errors"); } } - - // 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"); - } - } + Err(error) => { + tracing::error!(%error, "failed to build desired messaging adapters from config change"); } + } }); } } @@ -562,3 +299,439 @@ pub fn spawn_file_watcher( tracing::info!("file watcher stopped"); }) } + +fn build_desired_configured_adapters( + config: &Config, + instance_dir: &Path, + discord_permissions: Option>>, + slack_permissions: Option>>, + telegram_permissions: Option>>, + twitch_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={}", + 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={}", + 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={}", + slack_config.bot_token, + 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={}", + instance.bot_token, + 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={}", + 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={}", + 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, + email_config.imap_username, + email_config.imap_password, + email_config.imap_use_tls, + email_config.smtp_host, + email_config.smtp_port, + email_config.smtp_username, + 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, + instance.imap_username, + instance.imap_password, + instance.imap_use_tls, + instance.smtp_host, + instance.smtp_port, + instance.smtp_username, + 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 + ); + 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={}", + twitch_config.username, + twitch_config.oauth_token, + twitch_config.client_id, + twitch_config.client_secret, + twitch_config.refresh_token, + 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={}", + instance.username, + instance.oauth_token, + instance.client_id, + instance.client_secret, + instance.refresh_token, + 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, + )); + } + } + + 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 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 3f032db9c..37fbcfbba 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1495,7 +1495,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 { @@ -1531,7 +1531,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, @@ -1544,10 +1544,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(), @@ -1560,7 +1560,7 @@ async fn run( llm_manager.clone(), agent_links.clone(), agent_humans.clone(), - ); + )); } if foreground { @@ -2215,7 +2215,10 @@ async fn run( Ok(()) => { agents_initialized = true; // Restart file watcher with the new agent data - _file_watcher = spacebot::config::spawn_file_watcher( + if let Some(old_watcher) = _file_watcher.take() { + old_watcher.abort(); + } + _file_watcher = Some(spacebot::config::spawn_file_watcher( config_path.clone(), new_config.instance_dir.clone(), new_watcher_agents, @@ -2228,7 +2231,7 @@ async fn run( new_llm_manager.clone(), agent_links.clone(), agent_humans.clone(), - ); + )); tracing::info!("agents initialized after provider setup"); } Err(error) => { @@ -2988,20 +2991,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( diff --git a/src/messaging.rs b/src/messaging.rs index f57d5d68a..69e50962b 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -11,6 +11,6 @@ pub mod twitch; pub mod webchat; 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 a55a6e293..f851fdd88 100644 --- a/src/messaging/manager.rs +++ b/src/messaging/manager.rs @@ -7,7 +7,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, 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. /// @@ -15,10 +49,13 @@ 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>, /// 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 { @@ -26,28 +63,54 @@ 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()), 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 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 old_adapter = self + .adapters + .write() + .await + .insert(name.clone(), Arc::clone(&adapter)); + if self.started.load(Ordering::SeqCst) { + self.stop_runtime(&name, old_adapter, false).await.ok(); + } + if self.started.load(Ordering::SeqCst) + && 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 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 old_adapter = self + .adapters + .write() + .await + .insert(name.clone(), Arc::clone(&adapter)); + if self.started.load(Ordering::SeqCst) { + self.stop_runtime(&name, old_adapter, false).await.ok(); + } + if self.started.load(Ordering::SeqCst) + && 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; - /// Start all registered adapters and return the merged inbound stream. /// /// Each adapter's stream is forwarded into a shared channel, so adapters @@ -55,25 +118,13 @@ 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 adapters = self + .adapters + .read() + .await + .iter() + .map(|(name, adapter)| (name.clone(), Arc::clone(adapter))) + .collect::>(); let receiver = self .fan_in_rx @@ -81,6 +132,11 @@ impl MessagingManager { .await .take() .context("start() already called")?; + self.started.store(true, Ordering::SeqCst); + + for (name, adapter) in adapters { + self.start_runtime(name, adapter, false).await?; + } Ok(Box::pin(tokio_stream::wrappers::ReceiverStream::new( receiver, @@ -94,28 +150,14 @@ impl MessagingManager { /// stream replacement or restart. pub async fn register_and_start(&self, adapter: impl Messaging) -> crate::Result<()> { 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)); + self.stop_runtime(&name, old_adapter, false).await.ok(); + self.start_runtime(name.clone(), adapter, true).await?; tracing::info!(adapter = %name, "adapter registered and started at runtime"); Ok(()) } @@ -140,70 +182,232 @@ 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 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(&name).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), + ); + self.stop_runtime(&desired_adapter.name, old_adapter, false) + .await + .ok(); + 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, + initial_stream: Option, + ) -> JoinHandle<()> { tokio::spawn(async move { - let mut delay = std::time::Duration::from_secs(5); - let max_delay = std::time::Duration::from_secs(60); - - for attempt in 1..=Self::MAX_RETRY_ATTEMPTS { - tokio::time::sleep(delay).await; - - match adapter.start().await { - Ok(stream) => { - tracing::info!( - adapter = %name, - attempt, - "adapter started successfully after retry" - ); - Self::spawn_forwarder(name, stream, fan_in_tx); - return; + let mut next_retry_delay = INITIAL_RETRY_DELAY; + let mut initial_stream = initial_stream; + + loop { + let mut stream = if let Some(stream) = initial_stream.take() { + next_retry_delay = INITIAL_RETRY_DELAY; + stream + } else { + let start_result = tokio::select! { + _ = shutdown_rx.changed() => break, + result = adapter.start() => result, + }; + + match start_result { + Ok(stream) => { + tracing::info!(adapter = %name, "adapter started successfully"); + next_retry_delay = INITIAL_RETRY_DELAY; + stream + } + Err(error) => { + tracing::warn!( + adapter = %name, + %error, + 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, + } + } } - Err(error) => { - tracing::warn!( - adapter = %name, - attempt, - max_attempts = Self::MAX_RETRY_ATTEMPTS, - %error, - "adapter retry failed, next attempt in {:?}", - delay.min(max_delay) - ); + }; + + loop { + tokio::select! { + _ = shutdown_rx.changed() => { + tracing::info!(adapter = %name, "adapter supervisor shutting down"); + return; + } + maybe_message = stream.next() => { + match maybe_message { + Some(message) => { + if fan_in_tx.send(message).await.is_err() { + tracing::warn!(adapter = %name, "fan-in channel closed, stopping supervisor"); + return; + } + } + None => { + 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 start_result = adapter.start().await; + let initial_stream = match start_result { + Ok(stream) => stream, + Err(error) => { + self.install_supervisor(name.clone(), Arc::clone(&adapter), None) + .await; + tracing::warn!(adapter = %name, "scheduled background retry for failed adapter start"); + if surface_start_error { + return Err(anyhow::anyhow!("failed to start adapter '{name}': {error}").into()); + } + return Ok(()); } + }; - tracing::error!( - adapter = %name, - "adapter failed to start after {} attempts, giving up", - Self::MAX_RETRY_ATTEMPTS - ); - }); + self.install_supervisor(name, adapter, Some(initial_stream)) + .await; + 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, + initial_stream: 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, + initial_stream, + ); + 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(); + } + + if let Some(adapter) = adapter + && let Err(error) = adapter.shutdown().await + { + if propagate_shutdown_error { + return Err(error); + } + tracing::warn!(adapter = %name, %error, "failed to shut down adapter"); + } + + 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"); - }); + } + + Ok(()) } /// Inject a message directly into the fan-in channel, bypassing adapter streams. @@ -273,10 +477,9 @@ impl MessagingManager { /// Remove and shut down a single adapter by name. pub async fn remove_adapter(&self, name: &str) -> 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, true).await?; + tracing::info!(adapter = %name, "adapter removed and shut down"); Ok(()) } @@ -302,9 +505,15 @@ 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"); } } @@ -316,3 +525,153 @@ impl Default for MessagingManager { Self::new() } } + +#[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, + } + + impl TestAdapter { + fn new(name: &str, plans: Vec) -> Self { + Self { + name: name.to_string(), + plans: std::sync::Mutex::new(plans), + start_calls: AtomicUsize::new(0), + shutdown_calls: AtomicUsize::new(0), + } + } + } + + impl Messaging for TestAdapter { + fn name(&self) -> &str { + &self.name + } + + async fn start(&self) -> crate::Result { + self.start_calls.fetch_add(1, Ordering::SeqCst); + 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, + ], + )); + 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:?}"), + } + assert!( + adapter.start_calls.load(Ordering::SeqCst) >= 3, + "expected retries after failure and stream end" + ); + } + + #[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, + ], + )); + manager.register_shared(adapter.clone()).await; + let _stream = manager.start().await.expect("start manager"); + + tokio::time::sleep(std::time::Duration::from_millis(40)).await; + let before_remove = adapter.start_calls.load(Ordering::SeqCst); + manager + .remove_adapter("retrying") + .await + .expect("remove adapter"); + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + let after_remove = adapter.start_calls.load(Ordering::SeqCst); + + assert_eq!(before_remove, after_remove); + assert!(adapter.shutdown_calls.load(Ordering::SeqCst) >= 1); + } +} From f7ea53a874ed6973caf527dc45d94aa394da66c4 Mon Sep 17 00:00:00 2001 From: Olen Latham Date: Sun, 8 Mar 2026 16:06:12 -0500 Subject: [PATCH 2/5] fix(messaging): harden reconciler lifecycle details --- src/config/types.rs | 1 + src/config/watcher.rs | 68 ++++++++++++++++++++++++---------------- src/messaging/manager.rs | 66 +++++++++++++++++++++++++++----------- 3 files changed, 90 insertions(+), 45 deletions(-) diff --git a/src/config/types.rs b/src/config/types.rs index a905de799..36ae98cde 100644 --- a/src/config/types.rs +++ b/src/config/types.rs @@ -1560,6 +1560,7 @@ 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]); diff --git a/src/config/watcher.rs b/src/config/watcher.rs index 95d5e44ca..d300af707 100644 --- a/src/config/watcher.rs +++ b/src/config/watcher.rs @@ -5,6 +5,7 @@ use super::{ Binding, Config, DiscordPermissions, RuntimeConfig, 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). @@ -323,7 +324,7 @@ fn build_desired_configured_adapters( }); let fingerprint = format!( "token={}|dm={:?}|allow_bot_messages={}|permissions={}", - discord_config.token, + secret_fingerprint(&discord_config.token), sorted_strings(discord_config.dm_allowed_users.clone()), discord_config.allow_bot_messages, discord_permissions_fingerprint(&permissions_snapshot) @@ -350,7 +351,7 @@ fn build_desired_configured_adapters( DiscordPermissions::from_instance_config(instance, &config.bindings); let fingerprint = format!( "token={}|dm={:?}|allow_bot_messages={}|permissions={}", - instance.token, + secret_fingerprint(&instance.token), sorted_strings(instance.dm_allowed_users.clone()), instance.allow_bot_messages, discord_permissions_fingerprint(&permissions_snapshot) @@ -379,8 +380,8 @@ fn build_desired_configured_adapters( }); let fingerprint = format!( "bot_token={}|app_token={}|dm={:?}|commands={:?}|permissions={}", - slack_config.bot_token, - slack_config.app_token, + 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) @@ -410,8 +411,8 @@ fn build_desired_configured_adapters( SlackPermissions::from_instance_config(instance, &config.bindings); let fingerprint = format!( "bot_token={}|app_token={}|dm={:?}|commands={:?}|permissions={}", - instance.bot_token, - instance.app_token, + 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) @@ -443,7 +444,7 @@ fn build_desired_configured_adapters( }); let fingerprint = format!( "token={}|dm={:?}|permissions={}", - telegram_config.token, + secret_fingerprint(&telegram_config.token), sorted_strings(telegram_config.dm_allowed_users.clone()), telegram_permissions_fingerprint(&permissions_snapshot) ); @@ -469,7 +470,7 @@ fn build_desired_configured_adapters( TelegramPermissions::from_instance_config(instance, &config.bindings); let fingerprint = format!( "token={}|dm={:?}|permissions={}", - instance.token, + secret_fingerprint(&instance.token), sorted_strings(instance.dm_allowed_users.clone()), telegram_permissions_fingerprint(&permissions_snapshot) ); @@ -492,13 +493,13 @@ fn build_desired_configured_adapters( "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, - email_config.imap_username, - email_config.imap_password, + 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, - email_config.smtp_username, - email_config.smtp_password, + 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, @@ -528,13 +529,13 @@ fn build_desired_configured_adapters( instance.name, instance.imap_host, instance.imap_port, - instance.imap_username, - instance.imap_password, + secret_fingerprint(&instance.imap_username), + secret_fingerprint(&instance.imap_password), instance.imap_use_tls, instance.smtp_host, instance.smtp_port, - instance.smtp_username, - instance.smtp_password, + secret_fingerprint(&instance.smtp_username), + secret_fingerprint(&instance.smtp_password), instance.smtp_use_starttls, instance.from_address, instance.from_name, @@ -560,7 +561,9 @@ fn build_desired_configured_adapters( { let fingerprint = format!( "port={};bind={};auth_token={:?}", - webhook_config.port, webhook_config.bind, webhook_config.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( @@ -585,11 +588,17 @@ fn build_desired_configured_adapters( }); let fingerprint = format!( "username={};oauth_token={};client_id={:?};client_secret={:?};refresh_token={:?};channels={:?};trigger_prefix={:?};permissions={}", - twitch_config.username, - twitch_config.oauth_token, - twitch_config.client_id, - twitch_config.client_secret, - twitch_config.refresh_token, + 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) @@ -624,11 +633,11 @@ fn build_desired_configured_adapters( TwitchPermissions::from_instance_config(instance, &config.bindings); let fingerprint = format!( "username={};oauth_token={};client_id={:?};client_secret={:?};refresh_token={:?};channels={:?};trigger_prefix={:?};permissions={}", - instance.username, - instance.oauth_token, - instance.client_id, - instance.client_secret, - instance.refresh_token, + 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) @@ -725,6 +734,11 @@ fn twitch_permissions_fingerprint(permissions: &TwitchPermissions) -> String { ) } +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)> { diff --git a/src/messaging/manager.rs b/src/messaging/manager.rs index f851fdd88..aece02c2f 100644 --- a/src/messaging/manager.rs +++ b/src/messaging/manager.rs @@ -76,17 +76,16 @@ impl MessagingManager { let name = adapter.name().to_string(); let adapter: Arc = Arc::new(adapter); tracing::info!(adapter = %name, "registered messaging adapter"); + let started = self.started.load(Ordering::SeqCst); let old_adapter = self .adapters .write() .await .insert(name.clone(), Arc::clone(&adapter)); - if self.started.load(Ordering::SeqCst) { + if started { self.stop_runtime(&name, old_adapter, false).await.ok(); } - if self.started.load(Ordering::SeqCst) - && let Err(error) = self.start_runtime(name.clone(), adapter, false).await - { + 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"); } } @@ -96,17 +95,16 @@ impl MessagingManager { let name = adapter.name().to_string(); tracing::info!(adapter = %name, "registered messaging adapter (shared)"); 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 self.started.load(Ordering::SeqCst) { + if started { self.stop_runtime(&name, old_adapter, false).await.ok(); } - if self.started.load(Ordering::SeqCst) - && let Err(error) = self.start_runtime(name.clone(), adapter, false).await - { + 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"); } } @@ -547,17 +545,24 @@ mod tests { 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 { @@ -566,7 +571,8 @@ mod tests { } async fn start(&self) -> crate::Result { - self.start_calls.fetch_add(1, Ordering::SeqCst); + 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()), @@ -621,6 +627,7 @@ mod tests { 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"); @@ -641,10 +648,7 @@ mod tests { MessageContent::Text(text) => assert_eq!(text, "second"), other => panic!("unexpected message content: {other:?}"), } - assert!( - adapter.start_calls.load(Ordering::SeqCst) >= 3, - "expected retries after failure and stream end" - ); + wait_for_start_calls(&mut start_calls, 3).await; } #[tokio::test] @@ -659,19 +663,45 @@ mod tests { StartPlan::Pending, ], )); + let mut start_calls = adapter.subscribe_start_calls(); manager.register_shared(adapter.clone()).await; let _stream = manager.start().await.expect("start manager"); - tokio::time::sleep(std::time::Duration::from_millis(40)).await; - let before_remove = adapter.start_calls.load(Ordering::SeqCst); + 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"); - tokio::time::sleep(std::time::Duration::from_millis(50)).await; - let after_remove = adapter.start_calls.load(Ordering::SeqCst); + let no_more_starts = + tokio::time::timeout(std::time::Duration::from_millis(200), start_calls.changed()) + .await; - assert_eq!(before_remove, after_remove); + 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"); + } } From bbc3d4b22998bf46371fe0b6c532eef27bb46492 Mon Sep 17 00:00:00 2001 From: Olen Latham Date: Sun, 8 Mar 2026 16:19:26 -0500 Subject: [PATCH 3/5] fix(messaging): close reconciler lifecycle gaps --- src/config.rs | 2 +- src/config/watcher.rs | 36 +++++++++++++++++++++--- src/main.rs | 10 ++++--- src/messaging/manager.rs | 61 +++++++++++++++++++++++++++++++++++----- 4 files changed, 93 insertions(+), 16 deletions(-) diff --git a/src/config.rs b/src/config.rs index dea4b597c..2148eceb2 100644 --- a/src/config.rs +++ b/src/config.rs @@ -20,7 +20,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/watcher.rs b/src/config/watcher.rs index d300af707..b3923a4a8 100644 --- a/src/config/watcher.rs +++ b/src/config/watcher.rs @@ -17,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. /// @@ -37,11 +48,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( @@ -117,10 +129,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); } @@ -298,7 +321,12 @@ pub fn spawn_file_watcher( } tracing::info!("file watcher stopped"); - }) + }); + + FileWatcherHandle { + shutdown_tx, + _task: task, + } } fn build_desired_configured_adapters( diff --git a/src/main.rs b/src/main.rs index 37fbcfbba..7d1a03e2d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1495,7 +1495,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: Option>; + let mut _file_watcher: Option; // If providers are available, initialize agents immediately if has_providers { @@ -2215,9 +2215,7 @@ async fn run( Ok(()) => { agents_initialized = true; // Restart file watcher with the new agent data - if let Some(old_watcher) = _file_watcher.take() { - old_watcher.abort(); - } + let _old_watcher = _file_watcher.take(); _file_watcher = Some(spacebot::config::spawn_file_watcher( config_path.clone(), new_config.instance_dir.clone(), @@ -3015,6 +3013,10 @@ async fn initialize_agents( } } + new_messaging_manager + .seed_configured_fingerprints_from_registered() + .await; + let webchat_agent_pools = agents .iter() .map(|(agent_id, agent)| (agent_id.to_string(), agent.db.sqlite.clone())) diff --git a/src/messaging/manager.rs b/src/messaging/manager.rs index aece02c2f..c38c48716 100644 --- a/src/messaging/manager.rs +++ b/src/messaging/manager.rs @@ -51,6 +51,7 @@ 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()`. @@ -65,6 +66,7 @@ impl MessagingManager { 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), @@ -73,6 +75,7 @@ impl MessagingManager { /// 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"); @@ -92,6 +95,7 @@ impl MessagingManager { /// 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)"); let adapter: Arc = adapter; @@ -116,6 +120,8 @@ 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 _lifecycle_guard = self.lifecycle_mutex.lock().await; + self.started.store(true, Ordering::SeqCst); let adapters = self .adapters .read() @@ -130,7 +136,6 @@ impl MessagingManager { .await .take() .context("start() already called")?; - self.started.store(true, Ordering::SeqCst); for (name, adapter) in adapters { self.start_runtime(name, adapter, false).await?; @@ -147,6 +152,7 @@ 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(); let adapter: Arc = Arc::new(adapter); let old_adapter = self @@ -181,12 +187,13 @@ impl MessagingManager { } 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 + .adapters .read() .await .keys() @@ -261,10 +268,13 @@ impl MessagingManager { tokio::spawn(async move { let mut next_retry_delay = INITIAL_RETRY_DELAY; let mut initial_stream = initial_stream; + #[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); loop { let mut stream = if let Some(stream) = initial_stream.take() { - next_retry_delay = INITIAL_RETRY_DELAY; stream } else { let start_result = tokio::select! { @@ -295,6 +305,8 @@ impl MessagingManager { } }; + let started_at = std::time::Instant::now(); + let mut observed_message = false; loop { tokio::select! { _ = shutdown_rx.changed() => { @@ -304,13 +316,30 @@ impl MessagingManager { 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 => { - tracing::warn!(adapter = %name, "adapter stream ended, restarting"); + 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; } } @@ -385,13 +414,12 @@ impl MessagingManager { runtime.shutdown_tx.send(true).ok(); } + let mut shutdown_error = None; if let Some(adapter) = adapter && let Err(error) = adapter.shutdown().await { - if propagate_shutdown_error { - return Err(error); - } tracing::warn!(adapter = %name, %error, "failed to shut down adapter"); + shutdown_error = Some(error); } if let Some(runtime) = runtime { @@ -405,6 +433,10 @@ impl MessagingManager { } } + if propagate_shutdown_error && let Some(error) = shutdown_error { + return Err(error); + } + Ok(()) } @@ -474,6 +506,7 @@ 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; let adapter = self.adapters.write().await.remove(name); self.configured_fingerprints.write().await.remove(name); self.stop_runtime(name, adapter, true).await?; @@ -516,6 +549,20 @@ impl MessagingManager { } } } + + 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 { From c3bae3e7fc2865c9794d880fa6b8876cc7656ad1 Mon Sep 17 00:00:00 2001 From: Olen Latham Date: Sun, 8 Mar 2026 16:29:58 -0500 Subject: [PATCH 4/5] fix(messaging): address latest review feedback --- src/api/messaging.rs | 22 ++----- src/main.rs | 1 + src/messaging/manager.rs | 124 +++++++++++++++++++++++---------------- 3 files changed, 79 insertions(+), 68 deletions(-) diff --git a/src/api/messaging.rs b/src/api/messaging.rs index 692dea009..28622e7fa 100644 --- a/src/api/messaging.rs +++ b/src/api/messaging.rs @@ -683,11 +683,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"); @@ -1052,14 +1048,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( @@ -1603,11 +1593,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/main.rs b/src/main.rs index 7d1a03e2d..0776836bb 100644 --- a/src/main.rs +++ b/src/main.rs @@ -2180,6 +2180,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. diff --git a/src/messaging/manager.rs b/src/messaging/manager.rs index c38c48716..d1a64e11a 100644 --- a/src/messaging/manager.rs +++ b/src/messaging/manager.rs @@ -8,7 +8,7 @@ use futures::StreamExt as _; use std::collections::HashMap; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; -use tokio::sync::{RwLock, mpsc, watch}; +use tokio::sync::{RwLock, mpsc, oneshot, watch}; use tokio::task::JoinHandle; #[cfg(test)] @@ -85,8 +85,8 @@ impl MessagingManager { .write() .await .insert(name.clone(), Arc::clone(&adapter)); - if started { - self.stop_runtime(&name, old_adapter, false).await.ok(); + 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"); @@ -105,8 +105,8 @@ impl MessagingManager { .write() .await .insert(name.clone(), Arc::clone(&adapter)); - if started { - self.stop_runtime(&name, old_adapter, false).await.ok(); + 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"); @@ -160,7 +160,9 @@ impl MessagingManager { .write() .await .insert(name.clone(), Arc::clone(&adapter)); - self.stop_runtime(&name, old_adapter, false).await.ok(); + 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(()) @@ -225,9 +227,16 @@ impl MessagingManager { desired_adapter.name.clone(), Arc::clone(&desired_adapter.adapter), ); - self.stop_runtime(&desired_adapter.name, old_adapter, false) + if let Err(error) = self + .stop_runtime(&desired_adapter.name, old_adapter, false) .await - .ok(); + { + tracing::warn!( + adapter = %desired_adapter.name, + %error, + "failed to stop old runtime during reconciliation" + ); + } let replace_result = self .start_runtime( desired_adapter.name.clone(), @@ -263,44 +272,45 @@ impl MessagingManager { adapter: Arc, fan_in_tx: mpsc::Sender, mut shutdown_rx: watch::Receiver, - initial_stream: Option, + first_start_result_tx: Option>>, ) -> JoinHandle<()> { tokio::spawn(async move { let mut next_retry_delay = INITIAL_RETRY_DELAY; - let mut initial_stream = initial_stream; + 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); loop { - let mut stream = if let Some(stream) = initial_stream.take() { - stream - } else { - let start_result = tokio::select! { - _ = shutdown_rx.changed() => break, - result = adapter.start() => result, - }; - - match start_result { - Ok(stream) => { - tracing::info!(adapter = %name, "adapter started successfully"); - next_retry_delay = INITIAL_RETRY_DELAY; - stream + let start_result = tokio::select! { + _ = shutdown_rx.changed() => break, + result = adapter.start() => result, + }; + + let mut stream = match start_result { + Ok(stream) => { + tracing::info!(adapter = %name, "adapter started successfully"); + if let Some(tx) = first_start_result_tx.take() { + let _ = tx.send(Ok(())); } - Err(error) => { - tracing::warn!( - adapter = %name, - %error, - 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, - } + stream + } + Err(error) => { + if let Some(tx) = first_start_result_tx.take() { + let _ = tx.send(Err(error.to_string())); + } + tracing::warn!( + adapter = %name, + %error, + 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, } } }; @@ -360,22 +370,36 @@ impl MessagingManager { return Ok(()); } - let start_result = adapter.start().await; - let initial_stream = match start_result { - Ok(stream) => stream, - Err(error) => { - self.install_supervisor(name.clone(), Arc::clone(&adapter), None) - .await; - tracing::warn!(adapter = %name, "scheduled background retry for failed adapter start"); - if surface_start_error { + 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()); } - return Ok(()); + Err(error) => { + return Err(anyhow::anyhow!( + "failed to receive initial start result for '{name}': {error}" + ) + .into()); + } } - }; + } - self.install_supervisor(name, adapter, Some(initial_stream)) - .await; + self.configured_fingerprints + .write() + .await + .entry(name) + .or_insert_with(String::new); Ok(()) } @@ -383,7 +407,7 @@ impl MessagingManager { &self, name: String, adapter: Arc, - initial_stream: Option, + first_start_result_tx: Option>>, ) -> watch::Sender { let (shutdown_tx, shutdown_rx) = watch::channel(false); let task = Self::spawn_supervisor( @@ -391,7 +415,7 @@ impl MessagingManager { adapter, self.fan_in_tx.clone(), shutdown_rx, - initial_stream, + first_start_result_tx, ); self.runtimes.write().await.insert( name, From 433ee3fa3566509f41dc2946a8b1362313679f62 Mon Sep 17 00:00:00 2001 From: Olen Latham Date: Sun, 8 Mar 2026 16:43:04 -0500 Subject: [PATCH 5/5] Fix messaging reconciliation lockup --- src/messaging/manager.rs | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/src/messaging/manager.rs b/src/messaging/manager.rs index d1a64e11a..10a2c8af7 100644 --- a/src/messaging/manager.rs +++ b/src/messaging/manager.rs @@ -195,7 +195,7 @@ impl MessagingManager { .map(|adapter| adapter.name.clone()) .collect::>(); let current_configured = self - .adapters + .configured_fingerprints .read() .await .keys() @@ -205,7 +205,7 @@ impl MessagingManager { let mut first_error = None; for name in current_configured { if !desired_names.contains(&name) - && let Err(error) = self.remove_adapter(&name).await + && 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() { @@ -531,9 +531,18 @@ 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); self.configured_fingerprints.write().await.remove(name); - self.stop_runtime(name, adapter, true).await?; + self.stop_runtime(name, adapter, propagate_shutdown_error) + .await?; tracing::info!(adapter = %name, "adapter removed and shut down"); Ok(()) }