From cb1b48a15d3c26be609d554753161df3c3eac550 Mon Sep 17 00:00:00 2001 From: xavix-yo Date: Mon, 20 Jul 2026 00:07:33 +0100 Subject: [PATCH] Honcho: route per-user chat turns onto shared bus, expose to plugins The honcho memory sink needs to observe every user's completed chat turns from one subscription, keyed by ChatEvent.user_id. But each UserContext minted its own per-user ChatEventBus, so a single global subscription saw nothing. - UserContext now publishes onto the shared Runtime.event_bus (the one Skald::subscribe_chat_events reads) instead of a fresh per-user bus. - Expose that bus to plugins as PluginContext.chat_bus (distinct from system_bus, which carries only infra lifecycle events). - plugin-honcho subscribes via ctx.chat_bus. Co-Authored-By: Claude Opus 4.8 --- crates/core-api/src/plugin.rs | 6 ++++++ crates/plugin-honcho/src/lib.rs | 2 +- crates/skald-core/src/plugin/mod.rs | 1 + crates/skald-core/src/skald/user_context.rs | 11 ++++++++++- 4 files changed, 18 insertions(+), 2 deletions(-) diff --git a/crates/core-api/src/plugin.rs b/crates/core-api/src/plugin.rs index 8f6fb7f..cb5d871 100644 --- a/crates/core-api/src/plugin.rs +++ b/crates/core-api/src/plugin.rs @@ -5,6 +5,7 @@ use async_trait::async_trait; use serde_json::Value; use tokio::sync::RwLock; +use crate::bus::ChatEventBus; use crate::command::CommandApi; use crate::config_api::ConfigApi; use crate::i18n::I18nApi; @@ -84,6 +85,11 @@ pub struct PluginContext { pub api_provider_registry: Arc, pub location: Arc, pub system_bus: Arc, + /// The single shared chat-turn bus. Every user's completed turns are published + /// here, tagged with `ChatEvent.user_id`. A plugin that builds long-term memory + /// (Honcho) subscribes once and demuxes per user. Distinct from `system_bus`, + /// which carries only infra lifecycle events. + pub chat_bus: Arc, /// Channel-to-session resolver (blueprint §13). Lets channel plugins /// (Telegram, mobile, …) look up an unlocked user's chat hub, approval /// manager and event stream by user id. diff --git a/crates/plugin-honcho/src/lib.rs b/crates/plugin-honcho/src/lib.rs index 9e1177b..bd32991 100644 --- a/crates/plugin-honcho/src/lib.rs +++ b/crates/plugin-honcho/src/lib.rs @@ -892,7 +892,7 @@ impl core_api::plugin::Plugin for HonchoPlugin { self.honcho_memory.activate(Arc::clone(&client), workspace_id.clone(), Arc::clone(&user_config)); let session_map = Arc::clone(&self.honcho_memory.session_map); - let mut rx = ctx.event_bus.subscribe(); + let mut rx = ctx.chat_bus.subscribe(); let cancel = CancellationToken::new(); let cancel_clone = cancel.clone(); let running = Arc::clone(&self.running); diff --git a/crates/skald-core/src/plugin/mod.rs b/crates/skald-core/src/plugin/mod.rs index 423fc8c..b8901a9 100644 --- a/crates/skald-core/src/plugin/mod.rs +++ b/crates/skald-core/src/plugin/mod.rs @@ -187,6 +187,7 @@ impl PluginManager { api_provider_registry: Arc::clone(skald.provider_registry()) as _, location: Arc::clone(skald.location_manager()) as _, system_bus: Arc::clone(skald.system_bus()), + chat_bus: Arc::clone(skald.event_bus()), user_channel: self.skald()? as Arc, user_config: Arc::clone(&self.user_config) as _, i18n: self.i18n(), diff --git a/crates/skald-core/src/skald/user_context.rs b/crates/skald-core/src/skald/user_context.rs index 96d8e31..b3818d1 100644 --- a/crates/skald-core/src/skald/user_context.rs +++ b/crates/skald-core/src/skald/user_context.rs @@ -105,6 +105,11 @@ pub(super) struct UserContextFactory { image_generator_manager: Arc, run_context_manager: Arc, system_bus: Arc, + /// The single shared chat-turn bus. Every per-user `UserContext` publishes its + /// completed turns here (tagged with `user_id`) so a global consumer — the + /// Honcho memory sink — can observe every user's turns from one subscription + /// (`Skald::subscribe_chat_events`) and demux by `ChatEvent.user_id`. + event_bus: Arc, supervisor: Arc, shutdown_token: CancellationToken, max_history_messages: usize, @@ -138,6 +143,7 @@ impl UserContextFactory { image_generator_manager: Arc::clone(&media.image_generator_manager), run_context_manager: Arc::clone(&conversation.run_context_manager), system_bus: Arc::clone(&rt.system_bus), + event_bus: Arc::clone(&rt.event_bus), supervisor: Arc::clone(&rt.supervisor), shutdown_token: rt.shutdown_token.clone(), max_history_messages: config.llm.max_history_messages, @@ -156,7 +162,10 @@ impl UserContextFactory { // A shared swappable cell — a shared-folder membership change is applied in // place while the user is live (§6 remount), not deferred to next login. let fs = SharedFs::new(crate::container::build_user_fs(&self.registry_pool, user_id).await?); - let event_bus = Arc::new(ChatEventBus::new()); + // Shared, not per-user: publish this user's turns onto the one global bus so + // the Honcho sink sees every user from a single subscription (demux by + // `ChatEvent.user_id`). See the field doc on `UserContextFactory::event_bus`. + let event_bus = Arc::clone(&self.event_bus); let (global_tx, _) = broadcast::channel::(512); // Interaction stack, per-user. Approval reads the shared registry rules but