From 274b50ed1392b09bd2f65ce5669a97d6c2f0cb9e Mon Sep 17 00:00:00 2001 From: sanil-23 Date: Mon, 13 Apr 2026 19:19:42 +0530 Subject: [PATCH] feat(dispatch): route channel messages to welcome/orchestrator by onboarding flag (#525) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wires the per-agent tool scoping plumbing from commit 4a into the channel-message dispatch path. Each incoming channel message now picks the active agent — `welcome` pre-onboarding, `orchestrator` post — based on `Config::onboarding_completed`, loads the matching definition from the global `AgentDefinitionRegistry`, synthesises any `delegate_*` tools the agent declares in its `subagents` field, and passes everything through to `agent.run_turn` on the bus. This is the half of #525 that makes the welcome agent actually run for new users — the welcome definition has existed since upstream PR #522 but had no caller; nothing in dispatch consulted the onboarding flag, so every channel message ran through the same generic tool loop with the full registry exposed. Changes: src/openhuman/channels/runtime/dispatch.rs - New `AgentScoping` struct carrying the three new `AgentTurnRequest` fields (`target_agent_id`, `visible_tool_names`, `extra_tools`) plus an `unscoped()` constructor for safe-fallback paths. - New async `resolve_target_agent(channel)` helper: * fresh `Config::load_or_init().await` per turn (no cache — the loader reads from disk every call, verified at `config/schema/load.rs:409`, so the welcome→orchestrator handoff is observed on the next message after `complete_onboarding(complete)` flips the flag, with no need for an explicit handoff event); * picks `"welcome"` or `"orchestrator"` based on the flag and emits a structured `[dispatch::routing] selected target agent` info trace recording the choice + the flag value, satisfying the #525 acceptance criterion `"agent-selection logs clearly record why each agent was selected at onboarding boundaries"`; * looks up the definition in `AgentDefinitionRegistry::global()`, gracefully falling back to `AgentScoping::unscoped()` (= legacy behaviour, no filter, no extras) if the registry isn't initialised or the definition isn't found, so a routing miss never fails the user message; * for agents with a non-empty `subagents` field, awaits `composio::fetch_connected_integrations(&config)` and runs `orchestrator_tools::collect_orchestrator_tools` to materialise per-turn delegation tools (`research`, `plan`, `delegate_gmail`, …). Agents with empty `subagents` get an empty extras vec. - New `build_visible_tool_set(definition, &extra_tools)` helper that returns `Some(union)` for `ToolScope::Named` agents (their named list ∪ the names of the synthesised delegation tools) and `None` for `ToolScope::Wildcard` agents to preserve the unfiltered semantics — so agents like `skills_agent` and `morning_briefing` that already work via `wildcard + category_filter` keep their existing behaviour without this layer interfering. - `process_channel_message` calls `resolve_target_agent` once per turn, drops the placeholder defaults from commit 4a, and feeds the real `target_agent_id`/`visible_tool_names`/`extra_tools` into `AgentTurnRequest`. - New imports: `AgentDefinition`, `AgentDefinitionRegistry`, `ToolScope`, `Config`, `fetch_connected_integrations`, `orchestrator_tools`, `Tool`, `HashSet`. End-to-end behaviour after this commit: 1. New user, `onboarding_completed=false`: dispatch picks `welcome`, loads its 2-tool TOML scope, builds `visible_tool_names = {complete_onboarding, memory_recall}`, no extras, hands off to the bus. Bus handler applies the filter → welcome's LLM sees exactly 2 tools. 2. Welcome agent guides the user through setup, eventually calls `complete_onboarding(action="complete")` → flag persists to disk via `config.save()`. 3. Next user message: dispatch reads the flag fresh, picks `orchestrator`, fetches connected Composio integrations, expands `subagents = ["researcher", "planner", "code_executor", "critic", "archivist", { skills = "*" }]` into delegate_research / delegate_plan / delegate_run_code / delegate_review_code / delegate_archive_session + one delegate_ per connected integration. visible_tool_names is the union with the 4 direct tools from orchestrator's `[tools] named` list. LLM sees the scoped delegation surface, not the full 1000+ Composio catalog. #526's runtime leak is now fixed end-to-end: the orchestrator's LLM prompt only contains the tools its TOML allows, and the SkillDelegationTool path narrows skills_agent to a single toolkit via the `skill_filter` propagation fix from commit 4a. No agent at any layer sees more than its definition declares. Tests: 599/599 channel module tests pass — including `runtime_dispatch::dispatch_routes_through_agent_run_turn_bus_handler` and the telegram integration variant, which exercise the full bus roundtrip with the new fields populated. No existing assertions were modified. Co-Authored-By: Claude Opus 4.6 (1M context) --- src/openhuman/channels/runtime/dispatch.rs | 198 ++++++++++++++++++++- 1 file changed, 191 insertions(+), 7 deletions(-) diff --git a/src/openhuman/channels/runtime/dispatch.rs b/src/openhuman/channels/runtime/dispatch.rs index 54cceade4..04b6e7944 100644 --- a/src/openhuman/channels/runtime/dispatch.rs +++ b/src/openhuman/channels/runtime/dispatch.rs @@ -4,6 +4,9 @@ use crate::core::event_bus::{ publish_global, request_native_global, DomainEvent, NativeRequestError, }; use crate::openhuman::agent::bus::{AgentTurnRequest, AgentTurnResponse, AGENT_RUN_TURN_METHOD}; +use crate::openhuman::agent::harness::definition::{ + AgentDefinition, AgentDefinitionRegistry, ToolScope, +}; use crate::openhuman::channels::context::{ build_memory_context, compact_sender_history, conversation_history_key, conversation_memory_key, is_context_window_overflow_error, ChannelRuntimeContext, @@ -14,8 +17,12 @@ use crate::openhuman::channels::routes::{ }; use crate::openhuman::channels::traits; use crate::openhuman::channels::{Channel, SendMessage}; +use crate::openhuman::composio::fetch_connected_integrations; +use crate::openhuman::config::Config; use crate::openhuman::providers::{self, ChatMessage}; +use crate::openhuman::tools::{orchestrator_tools, Tool}; use crate::openhuman::util::truncate_with_ellipsis; +use std::collections::HashSet; use std::sync::Arc; use std::time::{Duration, Instant}; use tokio_util::sync::CancellationToken; @@ -176,6 +183,181 @@ fn spawn_scoped_typing_task( handle } +/// Per-turn scoping fields derived from the active agent definition. +/// +/// Carries the three new fields that get spliced into [`AgentTurnRequest`] +/// in [`process_channel_message`]. Constructed by [`resolve_target_agent`] +/// after reading `config.onboarding_completed`, looking up the matching +/// definition in [`AgentDefinitionRegistry`], and synthesising any +/// per-turn delegation tools the agent needs. +struct AgentScoping { + target_agent_id: Option, + visible_tool_names: Option>, + extra_tools: Vec>, +} + +impl AgentScoping { + /// Empty scoping — preserves the legacy "every tool in the global + /// registry is visible" behaviour. Returned when the registry isn't + /// initialised yet (early startup) or when the target agent + /// definition isn't found, so the channel layer never crashes the + /// runtime over a routing miss. + fn unscoped() -> Self { + Self { + target_agent_id: None, + visible_tool_names: None, + extra_tools: Vec::new(), + } + } +} + +/// Decide which agent should run for this channel turn and build the +/// matching tool-scoping payload. +/// +/// The selection is purely a function of `config.onboarding_completed`: +/// +/// * **`false`** → route to the `welcome` agent. Welcome's TOML +/// restricts it to two tools (`complete_onboarding`, `memory_recall`) +/// so the LLM cannot accidentally send messages or write files +/// while guiding the user through setup. The welcome agent decides +/// when the user is ready and calls +/// `complete_onboarding(action="complete")`, which flips the flag. +/// +/// * **`true`** → route to the `orchestrator` agent. Orchestrator +/// delegates real work to specialist subagents via a `subagents` +/// field in its TOML; this function expands that field into a list +/// of `delegate_*` tools spliced alongside the global registry. +/// +/// The next channel message after `complete_onboarding` flips the flag +/// is automatically routed to the orchestrator because +/// `Config::load_or_init()` reads from disk every call (no in-process +/// cache, verified at `config/schema/load.rs:409`), so the new value +/// is observed on the next turn without any explicit handoff event. +/// +/// On any failure path (missing registry, missing definition, missing +/// orchestrator delegation targets) the function logs and returns +/// [`AgentScoping::unscoped`], which lets the turn run with the legacy +/// unfiltered behaviour rather than failing the whole message. +async fn resolve_target_agent(channel: &str) -> AgentScoping { + let config = match Config::load_or_init().await { + Ok(c) => c, + Err(err) => { + tracing::warn!( + channel = %channel, + error = %err, + "[dispatch::routing] failed to load config — falling back to unscoped turn" + ); + return AgentScoping::unscoped(); + } + }; + + let target_id = if config.onboarding_completed { + "orchestrator" + } else { + "welcome" + }; + + tracing::info!( + channel = %channel, + target_agent = target_id, + onboarding_completed = config.onboarding_completed, + "[dispatch::routing] selected target agent" + ); + + let registry = match AgentDefinitionRegistry::global() { + Some(reg) => reg, + None => { + tracing::warn!( + channel = %channel, + target_agent = target_id, + "[dispatch::routing] AgentDefinitionRegistry not initialised — falling back to unscoped turn" + ); + return AgentScoping::unscoped(); + } + }; + + let definition = match registry.get(target_id) { + Some(def) => def, + None => { + tracing::warn!( + channel = %channel, + target_agent = target_id, + "[dispatch::routing] target agent not in registry — falling back to unscoped turn" + ); + return AgentScoping::unscoped(); + } + }; + + // Synthesise per-turn delegation tools when the target agent has a + // `subagents = [...]` field. Today only the orchestrator does, but + // the helper is agent-agnostic so future agents that delegate + // (e.g. a custom workspace-override planner that subdivides work) + // pick this up for free. + let extra_tools = if !definition.subagents.is_empty() { + let connected = fetch_connected_integrations(&config).await; + tracing::debug!( + channel = %channel, + target_agent = target_id, + connected_integration_count = connected.len(), + "[dispatch::routing] fetched connected integrations for delegation expansion" + ); + orchestrator_tools::collect_orchestrator_tools(definition, registry, &connected) + } else { + Vec::new() + }; + + let visible_tool_names = build_visible_tool_set(definition, &extra_tools); + + tracing::debug!( + channel = %channel, + target_agent = target_id, + named_tool_count = match &definition.tools { + ToolScope::Named(names) => names.len(), + ToolScope::Wildcard => 0, + }, + extra_tool_count = extra_tools.len(), + visible_tool_count = visible_tool_names.as_ref().map(|s| s.len()).unwrap_or(0), + "[dispatch::routing] assembled tool scoping for turn" + ); + + AgentScoping { + target_agent_id: Some(target_id.to_string()), + visible_tool_names, + extra_tools, + } +} + +/// Build the visible-tool whitelist for an agent. +/// +/// The set is the union of: +/// * every tool name in the agent's `[tools] named = [...]` list +/// (when the scope is [`ToolScope::Named`]); and +/// * every name produced by the per-turn synthesised delegation tools +/// in `extra_tools` (e.g. `research`, `delegate_gmail`). +/// +/// When the agent's tool scope is [`ToolScope::Wildcard`] **and** there +/// are no `extra_tools`, returns `None` to preserve the legacy +/// "everything visible" semantics — a `Wildcard` agent that delegates +/// nothing should still see the full registry. When `Wildcard` is +/// combined with non-empty extras (an unusual but legal combination), +/// the legacy unfiltered behaviour also wins because the wildcard +/// implicitly covers anything in the registry plus the extras. +fn build_visible_tool_set( + definition: &AgentDefinition, + extra_tools: &[Box], +) -> Option> { + match &definition.tools { + ToolScope::Wildcard => None, + ToolScope::Named(names) => { + let mut set: HashSet = names.iter().cloned().collect(); + for tool in extra_tools { + set.insert(tool.name().to_string()); + } + Some(set) + } + } +} + pub(crate) async fn process_channel_message( ctx: Arc, msg: traits::ChannelMessage, @@ -377,6 +559,12 @@ pub(crate) async fn process_channel_message( // The agent handler owns the history vector — we `mem::take` the // local one to avoid an unnecessary clone; `history` is not read // again below. + // Pick the active agent for this turn (welcome pre-onboarding, + // orchestrator post) and synthesise its delegation tool surface. + // Fresh disk read of `Config::onboarding_completed` happens inside + // `resolve_target_agent` — see the `[dispatch::routing]` traces. + let scoping = resolve_target_agent(&msg.channel).await; + let turn_request = AgentTurnRequest { provider: Arc::clone(&active_provider), history: std::mem::take(&mut history), @@ -389,13 +577,9 @@ pub(crate) async fn process_channel_message( multimodal: ctx.multimodal.clone(), max_tool_iterations: ctx.max_tool_iterations, on_delta: delta_tx, - // Per-agent scoping fields are populated in commit 4b (the - // dispatch routing logic for #525). For now, leave them at - // their defaults so this commit ships zero behaviour change — - // every channel turn still sees the full unfiltered registry. - target_agent_id: None, - visible_tool_names: None, - extra_tools: Vec::new(), + target_agent_id: scoping.target_agent_id, + visible_tool_names: scoping.visible_tool_names, + extra_tools: scoping.extra_tools, }; tracing::debug!( channel = %msg.channel,