From bde1f5414c0cdd191ce3703f14bfd2d9f3e18a3d Mon Sep 17 00:00:00 2001 From: psumotek <30962811+psumo@users.noreply.github.com> Date: Sun, 15 Mar 2026 13:55:08 +0000 Subject: [PATCH] channel agent reresolution MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When an agent is restarted, its UUID changes but the channel bridge still holds the old UUID from startup. This causes "Agent not found" errors. This fix stores the agent *name* alongside the cached UUID at bridge startup and, on "Agent not found" errors, re-resolves the name to a fresh UUID via find_agent_by_name(), updates the cache, and retries the message — all transparently to the end user. Changes: - router.rs: add channel_default_names DashMap, set_channel_default_with_name(), channel_default_name(), update_channel_default() - channel_bridge.rs: use set_channel_default_with_name() at startup - bridge.rs: add try_reresolution() helper, integrate retry logic into dispatch_message() and dispatch_with_blocks() error paths with proper lifecycle_reactions guards and sanitize_agent_error() usage --- crates/openfang-api/src/channel_bridge.rs | 2 +- crates/openfang-channels/src/bridge.rs | 201 +++++++++++++++++++++- crates/openfang-channels/src/router.rs | 29 ++++ 3 files changed, 230 insertions(+), 2 deletions(-) diff --git a/crates/openfang-api/src/channel_bridge.rs b/crates/openfang-api/src/channel_bridge.rs index a989ec5b..2715fcc9 100644 --- a/crates/openfang-api/src/channel_bridge.rs +++ b/crates/openfang-api/src/channel_bridge.rs @@ -1644,7 +1644,7 @@ pub async fn start_channel_bridge_with_config( "{} default agent: {name} ({agent_id}) [channel: {channel_key}]", adapter.name() ); - router.set_channel_default(channel_key, agent_id); + router.set_channel_default_with_name(channel_key, agent_id, name.clone()); // First configured default also becomes system-wide fallback if !system_default_set { router.set_default(agent_id); diff --git a/crates/openfang-channels/src/bridge.rs b/crates/openfang-channels/src/bridge.rs index c60b4936..2810d991 100644 --- a/crates/openfang-channels/src/bridge.rs +++ b/crates/openfang-channels/src/bridge.rs @@ -471,6 +471,45 @@ fn sender_user_id(message: &ChannelMessage) -> &str { .unwrap_or(&message.sender.platform_id) } +/// If an error contains "Agent not found", try to re-resolve the channel's default agent +/// by name (the name stored at bridge startup). Returns `Some(new_id)` on success. +async fn try_reresolution( + err: &str, + channel_key: &str, + handle: &Arc, + router: &Arc, +) -> Option { + if !err.to_lowercase().contains("agent not found") { + return None; + } + let name = router.channel_default_name(channel_key)?; + info!( + channel = channel_key, + agent_name = %name, + "Agent not found — attempting re-resolution by name" + ); + match handle.find_agent_by_name(&name).await { + Ok(Some(new_id)) => { + router.update_channel_default(channel_key, new_id); + info!( + channel = channel_key, + agent_name = %name, + new_id = %new_id, + "Re-resolved agent successfully" + ); + Some(new_id) + } + _ => { + warn!( + channel = channel_key, + agent_name = %name, + "Re-resolution failed — agent not found by name" + ); + None + } + } +} + /// Dispatch a single incoming message — handles bot commands or routes to an agent. /// /// Applies per-channel policies (DM/group filtering, rate limiting, formatting, threading). @@ -787,6 +826,9 @@ async fn dispatch_message( return; } + // Build channel key for re-resolution lookups + let channel_key = format!("{:?}", message.channel); + // Auto-reply check — if enabled, the engine decides whether to process this message. // If auto-reply is enabled but suppressed for this message, skip agent call entirely. if let Some(reply) = handle.check_auto_reply(agent_id, &text).await { @@ -828,6 +870,81 @@ async fn dispatch_message( .await; } Err(e) => { + // Try re-resolution before reporting error + if let Some(new_id) = + try_reresolution(&e, &channel_key, handle, router).await + { + let typing_task2 = + spawn_typing_loop(adapter_arc.clone(), message.sender.clone()); + let retry = handle.send_message(new_id, &text).await; + typing_task2.abort(); + match retry { + Ok(response) => { + if lifecycle_reactions { + send_lifecycle_reaction( + adapter, + &message.sender, + msg_id, + AgentPhase::Done, + ) + .await; + } + send_response( + adapter, + &message.sender, + response, + thread_id, + output_format, + ) + .await; + handle + .record_delivery( + new_id, + ct_str, + &message.sender.platform_id, + true, + None, + thread_id, + ) + .await; + } + Err(e2) => { + if lifecycle_reactions { + send_lifecycle_reaction( + adapter, + &message.sender, + msg_id, + AgentPhase::Error, + ) + .await; + } + warn!("Agent error after re-resolution for {new_id}: {e2}"); + let err_msg = sanitize_agent_error(&e2.to_string()); + if !adapter.suppress_error_responses() { + send_response( + adapter, + &message.sender, + err_msg.clone(), + thread_id, + output_format, + ) + .await; + } + handle + .record_delivery( + new_id, + ct_str, + &message.sender.platform_id, + false, + Some(&err_msg), + thread_id, + ) + .await; + } + } + return; + } + if lifecycle_reactions { send_lifecycle_reaction(adapter, &message.sender, msg_id, AgentPhase::Error).await; } @@ -1097,6 +1214,9 @@ async fn dispatch_with_blocks( } }; + // Build channel key for re-resolution lookups + let channel_key = format!("{:?}", message.channel); + // RBAC check if let Err(denied) = handle .authorize_channel_user(ct_str, sender_user_id(message), "chat") @@ -1125,7 +1245,9 @@ async fn dispatch_with_blocks( // Continuous typing indicator (see spawn_typing_loop doc) let typing_task = spawn_typing_loop(adapter_arc.clone(), message.sender.clone()); - let result = handle.send_message_with_blocks(agent_id, blocks).await; + let result = handle + .send_message_with_blocks(agent_id, blocks.clone()) + .await; typing_task.abort(); @@ -1140,6 +1262,83 @@ async fn dispatch_with_blocks( .await; } Err(e) => { + // Try re-resolution before reporting error + if let Some(new_id) = + try_reresolution(&e, &channel_key, handle, router).await + { + let typing_task2 = + spawn_typing_loop(adapter_arc.clone(), message.sender.clone()); + let retry = handle + .send_message_with_blocks(new_id, blocks) + .await; + typing_task2.abort(); + match retry { + Ok(response) => { + if lifecycle_reactions { + send_lifecycle_reaction( + adapter, + &message.sender, + msg_id, + AgentPhase::Done, + ) + .await; + } + send_response( + adapter, + &message.sender, + response, + thread_id, + output_format, + ) + .await; + handle + .record_delivery( + new_id, + ct_str, + &message.sender.platform_id, + true, + None, + thread_id, + ) + .await; + } + Err(e2) => { + if lifecycle_reactions { + send_lifecycle_reaction( + adapter, + &message.sender, + msg_id, + AgentPhase::Error, + ) + .await; + } + warn!("Agent error after re-resolution for {new_id}: {e2}"); + let err_msg = sanitize_agent_error(&e2.to_string()); + if !adapter.suppress_error_responses() { + send_response( + adapter, + &message.sender, + err_msg.clone(), + thread_id, + output_format, + ) + .await; + } + handle + .record_delivery( + new_id, + ct_str, + &message.sender.platform_id, + false, + Some(&err_msg), + thread_id, + ) + .await; + } + } + return; + } + if lifecycle_reactions { send_lifecycle_reaction(adapter, &message.sender, msg_id, AgentPhase::Error).await; } diff --git a/crates/openfang-channels/src/router.rs b/crates/openfang-channels/src/router.rs index 985f5d48..941e8173 100644 --- a/crates/openfang-channels/src/router.rs +++ b/crates/openfang-channels/src/router.rs @@ -34,6 +34,8 @@ pub struct AgentRouter { default_agent: Option, /// Per-channel-type default agent (e.g., Telegram -> agent_a, Discord -> agent_b). channel_defaults: DashMap, + /// Per-channel-type default agent *name* (for re-resolution when UUID becomes stale). + channel_default_names: DashMap, /// Sorted bindings (most specific first). Uses Mutex for runtime updates via Arc. bindings: Mutex>, /// Broadcast configuration. Uses Mutex for runtime updates via Arc. @@ -50,6 +52,7 @@ impl AgentRouter { direct_routes: DashMap::new(), default_agent: None, channel_defaults: DashMap::new(), + channel_default_names: DashMap::new(), bindings: Mutex::new(Vec::new()), broadcast: Mutex::new(BroadcastConfig::default()), agent_name_cache: DashMap::new(), @@ -66,6 +69,32 @@ impl AgentRouter { self.channel_defaults.insert(channel_key, agent_id); } + /// Set a per-channel-type default agent AND remember the agent name for + /// re-resolution when the cached UUID becomes stale (e.g. after agent restart). + pub fn set_channel_default_with_name( + &self, + channel_key: String, + agent_id: AgentId, + agent_name: String, + ) { + self.channel_defaults.insert(channel_key.clone(), agent_id); + self.channel_default_names + .insert(channel_key, agent_name); + } + + /// Retrieve the stored agent name for a channel default (if any). + pub fn channel_default_name(&self, channel_key: &str) -> Option { + self.channel_default_names + .get(channel_key) + .map(|r| r.clone()) + } + + /// Update the cached agent ID for a channel default (after re-resolution). + pub fn update_channel_default(&self, channel_key: &str, new_agent_id: AgentId) { + self.channel_defaults + .insert(channel_key.to_string(), new_agent_id); + } + /// Set a user's default agent. pub fn set_user_default(&self, user_key: String, agent_id: AgentId) { self.user_defaults.insert(user_key, agent_id);