diff --git a/AGENTS.md b/AGENTS.md index 948fc2e7359..601877e41b4 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -493,7 +493,9 @@ seed history by hand, or pick a transcript by recency. count. The host never shrinks a thread's tools because a cache went cold: `session_host/recorded_tools.rs` rebuilds recorded Composio actions as deferred executors, and the prelude fetches integrations on the first turn - of every session instance, not only on a brand-new thread. + of every session instance, not only on a brand-new thread. Explicit embed + `Turn::tools` overrides and host-only belts are authoritative instead: they + do not restore revoked tools from an earlier transcript. - **Pre-identity conversations are adopted once**, on first resume, from the timestamped stems they were written to (`adopt_legacy_session_transcripts`). No legacy file is modified. diff --git a/crates/openhuman-core/src/agent/hooks.rs b/crates/openhuman-core/src/agent/hooks.rs index e0d1b809b32..96443bb8aa3 100644 --- a/crates/openhuman-core/src/agent/hooks.rs +++ b/crates/openhuman-core/src/agent/hooks.rs @@ -4,6 +4,10 @@ //! what happened (user message, assistant response, tool calls with outcomes). //! The agent does not wait for hooks — they run in the background via `tokio::spawn`. +#[path = "hooks_scope.rs"] +mod scope; +pub use scope::{turn_post_turn_hooks, turn_tool_hooks, HookScope}; + use async_trait::async_trait; use serde::{Deserialize, Serialize}; use std::sync::Arc; @@ -175,6 +179,9 @@ pub struct ToolHookContext { /// Canonical agent definition id, when known. #[serde(default, skip_serializing_if = "Option::is_none")] pub agent_id: Option, + /// Working root of this tool dispatch, including a per-turn cwd override. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cwd: Option, } /// What a pre-tool hook decided about a call. diff --git a/crates/openhuman-core/src/agent/hooks_scope.rs b/crates/openhuman-core/src/agent/hooks_scope.rs new file mode 100644 index 00000000000..7d7a2b32cfb --- /dev/null +++ b/crates/openhuman-core/src/agent/hooks_scope.rs @@ -0,0 +1,67 @@ +//! Owned hooks carried by one dispatch task, never installed process-wide. + +use std::future::Future; +use std::sync::Arc; + +use super::{PostTurnHook, ToolHook}; + +tokio::task_local! { + static ACTIVE: HookScope; +} + +/// Agent/turn callbacks added after the runtime's global callbacks. +/// +/// Sessions snapshot these while they are built; asynchronous post-turn +/// callbacks own that snapshot. Independently spawned tasks must explicitly +/// inherit the scope if they build further sessions. +#[derive(Clone, Default)] +pub struct HookScope { + tools: Vec>, + post_turn: Vec>, + stop: Vec>, +} + +impl HookScope { + /// Append a tool callback, preserving registration order. + pub fn push_tool(&mut self, hook: Arc) { + self.tools.push(hook); + } + + /// Append a completed-turn callback, preserving registration order. + pub fn push_post_turn(&mut self, hook: Arc) { + self.post_turn.push(hook); + } + + /// Append a usage observer/budget policy checked after each model call. + pub fn push_stop(&mut self, hook: Arc) { + self.stop.push(hook); + } + + /// Build a session with scoped tool/post-turn callbacks. Stop policies + /// enter through the asynchronous dispatch scope when the turn runs. + pub fn sync_scope(self, build: impl FnOnce() -> T) -> T { + ACTIVE.sync_scope(self, build) + } + + /// Run a dispatch with this scope. Nested scopes replace tool/post-turn + /// callbacks and append their stop policies to the ambient policies. + pub async fn scope(self, future: F) -> F::Output { + let mut stop = crate::agent::stop_hooks::current_stop_hooks(); + stop.extend(self.stop.iter().cloned()); + crate::agent::stop_hooks::with_stop_hooks(stop, ACTIVE.scope(self, future)).await + } +} + +/// Session snapshot: process-wide tool callbacks, then the dispatch's callbacks. +pub fn turn_tool_hooks() -> Vec> { + let mut hooks = super::embedder_tool_hooks(); + let _ = ACTIVE.try_with(|scope| hooks.extend(scope.tools.iter().cloned())); + hooks +} + +/// Session snapshot: process-wide post-turn callbacks, then dispatch callbacks. +pub fn turn_post_turn_hooks() -> Vec> { + let mut hooks = super::embedder_post_turn_hooks(); + let _ = ACTIVE.try_with(|scope| hooks.extend(scope.post_turn.iter().cloned())); + hooks +} diff --git a/crates/openhuman-core/src/agent/host_agents.rs b/crates/openhuman-core/src/agent/host_agents.rs index 4affadbeced..f3c297ce084 100644 --- a/crates/openhuman-core/src/agent/host_agents.rs +++ b/crates/openhuman-core/src/agent/host_agents.rs @@ -47,6 +47,8 @@ pub struct HostAgent { pub host_tools: Option, /// The context every read during the agent's turn must see. pub context: Arc, + /// Agent-owned callbacks inherited by core-scheduled turns. + pub hooks: crate::agent::hooks::HookScope, } impl HostAgent { @@ -67,20 +69,25 @@ impl HostAgent { host_tools = self.host_tools.is_some(), "[host_agents] building session for host agent" ); - CoreContext::sync_scope(Arc::clone(&self.context), || match &self.host_tools { - Some(host) => OpenHumanSessionHost::from_config_with_host_tools( - config, - &self.definition, - host, - session_id, - ), - None => OpenHumanSessionHost::from_config_with_definition(config, &self.definition), + self.hooks.clone().sync_scope(|| { + CoreContext::sync_scope(Arc::clone(&self.context), || match &self.host_tools { + Some(host) => OpenHumanSessionHost::from_config_with_host_tools( + config, + &self.definition, + host, + session_id, + ), + None => OpenHumanSessionHost::from_config_with_definition(config, &self.definition), + }) }) } /// Run `fut` with this agent's context as the ambient [`CoreContext`]. pub async fn scope(&self, fut: F) -> F::Output { - CoreContext::scope(Arc::clone(&self.context), fut).await + self.hooks + .clone() + .scope(CoreContext::scope(Arc::clone(&self.context), fut)) + .await } } diff --git a/crates/openhuman-core/src/agent/host_agents_tests.rs b/crates/openhuman-core/src/agent/host_agents_tests.rs index 489951c3f61..994a7b25686 100644 --- a/crates/openhuman-core/src/agent/host_agents_tests.rs +++ b/crates/openhuman-core/src/agent/host_agents_tests.rs @@ -34,6 +34,7 @@ impl HostAgentResolver for OneAgent { definition: definition(&self.id), config: self.config.clone(), host_tools: None, + hooks: Default::default(), context: CoreContext::for_test(DomainSet::full(), None), }) } @@ -89,6 +90,7 @@ fn a_host_agent_debug_names_the_agent_and_hides_the_rest() { definition: definition("host-agents-debug"), config: Config::default(), host_tools: None, + hooks: Default::default(), context: CoreContext::for_test(DomainSet::full(), None), }; let rendered = format!("{host:?}"); @@ -104,6 +106,7 @@ async fn scope_runs_the_future_under_the_agents_context() { definition: definition("host-agents-scope"), config: Config::default(), host_tools: None, + hooks: Default::default(), context: Arc::clone(&context), }; let seen = host diff --git a/crates/openhuman-core/src/agent/mod.rs b/crates/openhuman-core/src/agent/mod.rs index 83f8b18e8fb..a75e8999842 100644 --- a/crates/openhuman-core/src/agent/mod.rs +++ b/crates/openhuman-core/src/agent/mod.rs @@ -79,6 +79,7 @@ pub mod subagent_host; pub mod tinyagents; pub mod todos; pub mod tool_policy; +pub mod tool_snapshot_scope; pub mod tools; pub mod triage; /// Wall-clock deadline of one top-level turn: the outer backstop and the diff --git a/crates/openhuman-core/src/agent/session_host/builder/factory.rs b/crates/openhuman-core/src/agent/session_host/builder/factory.rs index b52624611ef..5ea2f564d40 100644 --- a/crates/openhuman-core/src/agent/session_host/builder/factory.rs +++ b/crates/openhuman-core/src/agent/session_host/builder/factory.rs @@ -393,7 +393,7 @@ impl OpenHumanSessionHost { None => SystemPromptBuilder::with_defaults(), }; let post_turn_hooks: Vec> = - crate::agent::hooks::embedder_post_turn_hooks(); + crate::agent::hooks::turn_post_turn_hooks(); // Best-effort prewarm from the shared Composio cache. This avoids // building the session with a knowingly stale `&[]` integration view diff --git a/crates/openhuman-core/src/agent/session_host/runtime_session.rs b/crates/openhuman-core/src/agent/session_host/runtime_session.rs index b6574f93b6b..e51845e84a1 100644 --- a/crates/openhuman-core/src/agent/session_host/runtime_session.rs +++ b/crates/openhuman-core/src/agent/session_host/runtime_session.rs @@ -11,6 +11,7 @@ mod memory_ingest; mod permanent; mod post_commit; mod tool_rules; +mod transcript; #[path = "runtime_session_turn.rs"] mod turn; pub(super) use turn::begin_turn_resume; @@ -22,10 +23,8 @@ use tinyagents_runtime::{ CommitReceipt, ResumePreparation, SessionBuilder, SessionTerminal, ToolSnapshot, TranscriptTarget, TurnPreparation, }; -use tinyagents_session::transcript::TranscriptMeta; use tinyinference_llm::message::Message; -use crate::agent::session_store::transcripts_or_files; use crate::agent::{ message_convert::{user_message_from_text, user_text_with_markers}, session_host::{ @@ -1222,7 +1221,9 @@ impl OpenHumanSessionHost { let mut builder = SessionBuilder::new(driver) .codec(Arc::new(OpenHumanTranscriptCodec)) .hooks(hooks) - .retain_recorded_tools(true); + .retain_recorded_tools( + !self.host_only && crate::agent::tool_snapshot_scope::retain_recorded_tools(), + ); if let Some(session) = self.session.clone() { builder = builder.session(session_locator, session, self.runtime_transcript_meta()); } @@ -1282,63 +1283,6 @@ impl OpenHumanSessionHost { agent_definition_name: self.agent_definition_name.clone(), }); } - - fn runtime_transcript_stem(&self) -> String { - match &self.session_parent_prefix { - Some(prefix) => format!("{prefix}__{}", self.session_key), - None => self.session_key.clone(), - } - } - - pub(in crate::agent::session_host) fn session_locator( - &self, - ) -> Arc { - if let Some(injected) = self.session_history_locator.clone() { - return injected; - } - self.session_history_locator_memo - .get_or_init(|| { - let session_agent_id = - crate::agent::session_store::current_agent_key_or(&self.agent_definition_id); - transcripts_or_files(&session_agent_id, &self.workspace_dir) - }) - .clone() - } - - fn runtime_transcript_meta(&self) -> TranscriptMeta { - let now = chrono::Utc::now().to_rfc3339(); - TranscriptMeta { - agent_name: self.agent_definition_name.clone(), - agent_id: Some(self.agent_definition_id.clone()), - agent_type: Some(if self.session_parent_prefix.is_some() { - "subagent".into() - } else { - "root".into() - }), - dispatcher: if self.tool_dispatcher.should_send_tool_specs() { - "native".into() - } else { - "xml".into() - }, - provider: None, - model: Some(self.model_name.clone()), - created: now.clone(), - updated: now, - turn_count: 0, - prefix_message_count: None, - input_tokens: 0, - output_tokens: 0, - cached_input_tokens: 0, - charged_amount_usd: 0.0, - thread_id: self.thread_id.clone(), - task_id: None, - session_id: self.session.as_ref().map(|session| session.session_id()), - parent_session_id: self - .session - .as_ref() - .and_then(|session| session.parent_session_id()), - } - } } #[path = "prelude_integrations.rs"] diff --git a/crates/openhuman-core/src/agent/session_host/runtime_session/transcript.rs b/crates/openhuman-core/src/agent/session_host/runtime_session/transcript.rs new file mode 100644 index 00000000000..f6e1b9fde55 --- /dev/null +++ b/crates/openhuman-core/src/agent/session_host/runtime_session/transcript.rs @@ -0,0 +1,65 @@ +//! Session transcript identity, locator and metadata construction. + +use super::super::types::OpenHumanSessionHost; +use crate::agent::session_store::transcripts_or_files; +use std::sync::Arc; +use tinyagents_session::transcript::TranscriptMeta; + +impl OpenHumanSessionHost { + pub(super) fn runtime_transcript_stem(&self) -> String { + match &self.session_parent_prefix { + Some(prefix) => format!("{prefix}__{}", self.session_key), + None => self.session_key.clone(), + } + } + + pub(in crate::agent::session_host) fn session_locator( + &self, + ) -> Arc { + if let Some(injected) = self.session_history_locator.clone() { + return injected; + } + self.session_history_locator_memo + .get_or_init(|| { + let session_agent_id = + crate::agent::session_store::current_agent_key_or(&self.agent_definition_id); + transcripts_or_files(&session_agent_id, &self.workspace_dir) + }) + .clone() + } + + pub(super) fn runtime_transcript_meta(&self) -> TranscriptMeta { + let now = chrono::Utc::now().to_rfc3339(); + TranscriptMeta { + agent_name: self.agent_definition_name.clone(), + agent_id: Some(self.agent_definition_id.clone()), + agent_type: Some(if self.session_parent_prefix.is_some() { + "subagent".into() + } else { + "root".into() + }), + dispatcher: if self.tool_dispatcher.should_send_tool_specs() { + "native".into() + } else { + "xml".into() + }, + provider: None, + model: Some(self.model_name.clone()), + created: now.clone(), + updated: now, + turn_count: 0, + prefix_message_count: None, + input_tokens: 0, + output_tokens: 0, + cached_input_tokens: 0, + charged_amount_usd: 0.0, + thread_id: self.thread_id.clone(), + task_id: None, + session_id: self.session.as_ref().map(|session| session.session_id()), + parent_session_id: self + .session + .as_ref() + .and_then(|session| session.parent_session_id()), + } + } +} diff --git a/crates/openhuman-core/src/agent/tinyagents/harness_assembly.rs b/crates/openhuman-core/src/agent/tinyagents/harness_assembly.rs index 34fd7a507e0..0c2d02eccfa 100644 --- a/crates/openhuman-core/src/agent/tinyagents/harness_assembly.rs +++ b/crates/openhuman-core/src/agent/tinyagents/harness_assembly.rs @@ -255,10 +255,6 @@ pub(super) fn assemble_turn_harness( .as_ref() .map(|_| Arc::new(ToolResultArtifactIndexStore::new())); - // The explicit run carrier supplies the stop hooks; no middleware needs to - // recover policy from a task-local while the harness is driving. - let stop_hooks_installed = stop_hooks; - // A steering handle is always created now: besides run-queue steering, the // early-exit / cap / stop-hook pauses, the repeated-tool-failure breaker // (below) also pauses through it, and it wants to fire on every path @@ -314,8 +310,7 @@ pub(super) fn assemble_turn_harness( harness.push_middleware(mw.clone()); } - // Repeated-failure breaker: surface a root cause instead of burning the budget - // on failing calls; side effects come from the tools' own declarations. + // Repeated failures stop the run using each tool's declared side effects. let repeated_failure = handle.as_ref().map(|handle| { let (t, halt) = (REPEATED_TOOL_FAILURE_THRESHOLD, halt_summary.clone()); let mw = middleware::RepeatedToolFailureMiddleware::new(handle.clone(), t, halt); @@ -325,19 +320,14 @@ pub(super) fn assemble_turn_harness( harness.push_middleware(mw.clone()); } - // Policy-driven stop hooks (budget cap, thread-goal budget, ad-hoc iteration - // ceiling): fire after each model call and pause the run on the first stop - // vote. Replaces the legacy tool-call-loop firing point. - if let Some(handle) = &handle { - if !stop_hooks_installed.is_empty() { - harness.push_middleware(Arc::new(stop_hooks::StopHookMiddleware::new( - handle.clone(), - model, - max_iterations, - stop_hooks_installed, - ))); - } - } + stop_hooks::install( + &mut harness, + handle.as_ref(), + model, + max_iterations, + stop_hooks, + halt_summary.clone(), + ); let early_exit_set: HashSet<&str> = early_exit_tools.iter().copied().collect(); // One hook per run, shared by every early-exit adapter (records the first // early-exit and pauses). Requires the steering handle. @@ -648,7 +638,7 @@ pub(super) fn assemble_turn_harness( // observation-only (never mutates the result), so running first in the // reverse-order `after_tool` chain is safe — it cannot perturb the // summarization/cap or tool-outcome capture layers. - let embedder_tool_hooks = crate::agent::hooks::embedder_tool_hooks(); + let embedder_tool_hooks = crate::agent::hooks::turn_tool_hooks(); if !embedder_tool_hooks.is_empty() { harness.push_middleware(Arc::new(middleware::EmbedderToolHooksMiddleware::new( embedder_tool_hooks, diff --git a/crates/openhuman-core/src/agent/tinyagents/live_harness.rs b/crates/openhuman-core/src/agent/tinyagents/live_harness.rs index 79e2cfd1d47..92f3e715a61 100644 --- a/crates/openhuman-core/src/agent/tinyagents/live_harness.rs +++ b/crates/openhuman-core/src/agent/tinyagents/live_harness.rs @@ -121,7 +121,7 @@ pub(crate) fn assemble_live_tool_harness( registered_tools, route_session, ))); - let embedder_tool_hooks = crate::agent::hooks::embedder_tool_hooks(); + let embedder_tool_hooks = crate::agent::hooks::turn_tool_hooks(); if !embedder_tool_hooks.is_empty() { harness.push_middleware(Arc::new(middleware::EmbedderToolHooksMiddleware::new( embedder_tool_hooks, diff --git a/crates/openhuman-core/src/agent/tinyagents/middleware/embedder_hooks.rs b/crates/openhuman-core/src/agent/tinyagents/middleware/embedder_hooks.rs index aca6086b22d..f91c2edb145 100644 --- a/crates/openhuman-core/src/agent/tinyagents/middleware/embedder_hooks.rs +++ b/crates/openhuman-core/src/agent/tinyagents/middleware/embedder_hooks.rs @@ -31,6 +31,37 @@ impl EmbedderToolHooksMiddleware { } } +fn hook_identity( + data: &crate::agent::tinyagents::host::OpenHumanRunContext, +) -> (Option, Option, Option) { + let session_id = data + .thread_id + .clone() + .or_else(|| data.parent.as_ref().map(|parent| parent.session_id.clone())); + let agent_id = data + .parent + .as_ref() + .map(|parent| parent.agent_definition_id.clone()); + let cwd = data + .workspace + .as_ref() + .map(|workspace| workspace.root.clone()) + .or_else(|| { + data.parent.as_ref().and_then(|parent| { + parent + .workspace_descriptor + .as_ref() + .map(|workspace| workspace.root.clone()) + }) + }) + .or_else(|| { + crate::core::runtime::CoreContext::with_current_embedder_config(|config| { + config.action_dir.clone() + }) + }); + (session_id, agent_id, cwd) +} + #[async_trait] impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> for EmbedderToolHooksMiddleware @@ -41,10 +72,11 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> async fn before_tool( &self, - _ctx: &mut RunContext, + ctx: &mut RunContext, _state: &(), call: &mut TaToolCall, ) -> TaResult<()> { + let (session_id, agent_id, cwd) = hook_identity(&ctx.data); let mut context = crate::agent::hooks::ToolHookContext { event: crate::agent::hooks::ToolHookEvent::PreToolUse, call_id: call.id.clone(), @@ -54,8 +86,9 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> duration_ms: None, output: None, error: None, - session_id: None, - agent_id: None, + session_id, + agent_id, + cwd, }; for hook in &self.hooks { match hook.before_tool_decision(&context).await { @@ -120,10 +153,11 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> /// arguments through. async fn check_nested_tool( &self, - _ctx: &RunContext, + ctx: &RunContext, _state: &(), call: &TaToolCall, ) -> TaResult<()> { + let (session_id, agent_id, cwd) = hook_identity(&ctx.data); let context = crate::agent::hooks::ToolHookContext { event: crate::agent::hooks::ToolHookEvent::PreToolUse, call_id: call.id.clone(), @@ -133,8 +167,9 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> duration_ms: None, output: None, error: None, - session_id: None, - agent_id: None, + session_id, + agent_id, + cwd, }; for hook in &self.hooks { let refusal = match hook.before_tool_decision(&context).await { @@ -165,7 +200,7 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> async fn after_tool( &self, - _ctx: &mut RunContext, + ctx: &mut RunContext, _state: &(), invocation: &ToolInvocationIdentity, result: &mut TaToolResult, @@ -178,6 +213,7 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> .expect("embedder tool-hook arguments poisoned") .remove(&call_id) .unwrap_or(serde_json::Value::Null); + let (session_id, agent_id, cwd) = hook_identity(&ctx.data); let context = crate::agent::hooks::ToolHookContext { event: crate::agent::hooks::ToolHookEvent::PostToolUse, call_id, @@ -191,8 +227,9 @@ impl Middleware<(), crate::agent::tinyagents::host::OpenHumanRunContext> error: result .is_error .then(|| crate::agent::tinyagents::middleware::tool_result_text(result)), - session_id: None, - agent_id: None, + session_id, + agent_id, + cwd, }; for hook in &self.hooks { // Text a hook returns is appended to the result the model reads — diff --git a/crates/openhuman-core/src/agent/tinyagents/model.rs b/crates/openhuman-core/src/agent/tinyagents/model.rs index f40c7b1c464..230037fe073 100644 --- a/crates/openhuman-core/src/agent/tinyagents/model.rs +++ b/crates/openhuman-core/src/agent/tinyagents/model.rs @@ -274,7 +274,13 @@ pub(crate) fn usage_info_from_response(response: &ModelResponse) -> Option, + handle: Option<&SteeringHandle>, + model: &str, + max_iterations: usize, + hooks: Vec>, + halt_summary: super::HaltSummarySlot, +) { + if let Some(handle) = handle.filter(|_| !hooks.is_empty()) { + harness.push_middleware(Arc::new(StopHookMiddleware::new( + handle.clone(), + model, + max_iterations, + hooks, + halt_summary, + ))); + } +} /// Fires openhuman [`StopHook`]s after each model call and pauses the run when /// any hook votes to stop. @@ -49,6 +68,7 @@ pub(super) struct StopHookMiddleware { hooks: Vec>, /// Latches once a hook has voted to stop, so we send `Pause` exactly once. stopped: AtomicBool, + halt_summary: super::HaltSummarySlot, } impl StopHookMiddleware { @@ -58,6 +78,7 @@ impl StopHookMiddleware { model: impl Into, max_iterations: usize, hooks: Vec>, + halt_summary: super::HaltSummarySlot, ) -> Self { Self { handle, @@ -67,6 +88,7 @@ impl StopHookMiddleware { cost: Mutex::new(TurnCost::new()), hooks, stopped: AtomicBool::new(false), + halt_summary, } } } @@ -98,14 +120,8 @@ where let iteration = self.iteration.fetch_add(1, Ordering::SeqCst) + 1; let cost_snapshot = { let mut cost = self.cost.lock().expect("stop-hook cost mutex poisoned"); - if let Some(usage) = &response.usage { - cost.add_call( - &self.model, - &BilledUsage::from_counts(usage.input_tokens, usage.output_tokens) - .with_cached_input_tokens(usage.cache_read_tokens) - .with_cache_creation_tokens(usage.cache_creation_tokens) - .with_reasoning_tokens(usage.reasoning_tokens), - ); + if let Some(usage) = super::model::usage_info_from_response(response) { + cost.add_call(&self.model, &usage); } cost.clone() }; @@ -117,26 +133,35 @@ where model: &self.model, }; + // Every observer sees the completed call, including observers installed + // after a budget policy. Preserve the first stop reason for the result. + let mut stop = None; for hook in &self.hooks { if let StopDecision::Stop { reason } = hook.check(&turn_state).await { - // Latch first so a concurrent (streaming) after_model can't - // double-pause. - if self.stopped.swap(true, Ordering::SeqCst) { - return Ok(()); + if stop.is_none() { + stop = Some((hook.name().to_owned(), reason)); } - tracing::warn!( - target: "stop_hooks", - hook = hook.name(), - iteration, - model = %self.model, - "[stop_hooks] hook voted to stop the turn — pausing run: {reason}" - ); - // Graceful stop: the loop drains steering at the top of the next - // iteration and `Pause` short-circuits it before the next model - // call. The partial transcript is returned to the caller. - self.handle.send(SteeringCommand::Pause); + } + } + if let Some((hook, reason)) = stop { + if self.stopped.swap(true, Ordering::SeqCst) { return Ok(()); } + tracing::warn!( + target: "stop_hooks", + hook, + iteration, + model = %self.model, + "[stop_hooks] hook voted to stop the turn — pausing run: {reason}" + ); + *self + .halt_summary + .lock() + .expect("stop-hook halt mutex poisoned") = Some(format!( + "Stopping after {} model call(s): {reason}", + iteration + )); + self.handle.send(SteeringCommand::Pause); } Ok(()) diff --git a/crates/openhuman-core/src/agent/tool_snapshot_scope.rs b/crates/openhuman-core/src/agent/tool_snapshot_scope.rs new file mode 100644 index 00000000000..497d7ad5b1e --- /dev/null +++ b/crates/openhuman-core/src/agent/tool_snapshot_scope.rs @@ -0,0 +1,15 @@ +//! Explicit one-turn tool authority over restored transcript declarations. + +use std::future::Future; +tokio::task_local! { static FRESH: (); } + +/// Run a turn with the current host belt as authoritative. Historical tool +/// declarations remain in the transcript but cannot restore a revoked tool. +/// Host-spawned dispatch tasks must explicitly carry this scope. +pub async fn with_fresh_snapshot(future: F) -> F::Output { + FRESH.scope((), future).await +} + +pub(crate) fn retain_recorded_tools() -> bool { + FRESH.try_with(|_| ()).is_err() +} diff --git a/crates/openhuman-core/src/channels/runtime/dispatch/host_agent/dispatch_tests.rs b/crates/openhuman-core/src/channels/runtime/dispatch/host_agent/dispatch_tests.rs index cfd6abbf0e9..c4be6dbf20d 100644 --- a/crates/openhuman-core/src/channels/runtime/dispatch/host_agent/dispatch_tests.rs +++ b/crates/openhuman-core/src/channels/runtime/dispatch/host_agent/dispatch_tests.rs @@ -235,6 +235,7 @@ impl HostAgentResolver for OneHostAgent { }), ]) })), + hooks: Default::default(), context: CoreContext::for_test_with_config(DomainSet::full(), self.config.clone()), }) } diff --git a/crates/openhuman-core/src/config/ops/loader_current_tests.rs b/crates/openhuman-core/src/config/ops/loader_current_tests.rs index 4ced94e23b6..bd049ede6b4 100644 --- a/crates/openhuman-core/src/config/ops/loader_current_tests.rs +++ b/crates/openhuman-core/src/config/ops/loader_current_tests.rs @@ -18,6 +18,7 @@ async fn a_context_route_follows_every_load_of_the_context_config() { config.ephemeral_route = Some(EphemeralRoute { endpoint: "http://127.0.0.1:9/v1".to_string(), api_key: "test-key".to_string(), + headers: Vec::new(), }); let loaded = CoreContext::scope(context_with(config), load_current_or_init()) diff --git a/crates/openhuman-core/src/config/schema/ephemeral_route.rs b/crates/openhuman-core/src/config/schema/ephemeral_route.rs index 964930e4b3b..311ddab3241 100644 --- a/crates/openhuman-core/src/config/schema/ephemeral_route.rs +++ b/crates/openhuman-core/src/config/schema/ephemeral_route.rs @@ -60,6 +60,9 @@ pub struct EphemeralRoute { /// and only for [`EPHEMERAL_ROUTE_SLUG`], so it cannot be handed to a /// provider the caller did not name. pub api_key: String, + /// Custom headers scoped to this endpoint, never copied to unrelated providers. + #[serde(default)] + pub headers: Vec<(String, String)>, } impl EphemeralRoute { @@ -71,7 +74,17 @@ impl EphemeralRoute { pub fn from_params(endpoint: Option, api_key: Option) -> Option { let endpoint = endpoint?.trim().to_string(); let api_key = api_key?.trim().to_string(); - (!endpoint.is_empty() && !api_key.is_empty()).then_some(Self { endpoint, api_key }) + (!endpoint.is_empty() && !api_key.is_empty()).then_some(Self { + endpoint, + api_key, + headers: Vec::new(), + }) + } + + /// Attach the embedding host's per-route headers. + pub fn with_headers(mut self, headers: Vec<(String, String)>) -> Self { + self.headers = headers; + self } } diff --git a/crates/openhuman-core/src/config/schema/ephemeral_route_tests.rs b/crates/openhuman-core/src/config/schema/ephemeral_route_tests.rs index c3e31c538cf..ea5375a707d 100644 --- a/crates/openhuman-core/src/config/schema/ephemeral_route_tests.rs +++ b/crates/openhuman-core/src/config/schema/ephemeral_route_tests.rs @@ -15,6 +15,7 @@ fn route() -> EphemeralRoute { EphemeralRoute { endpoint: "http://127.0.0.1:41234/openai".to_string(), api_key: "mdl-token".to_string(), + headers: Vec::new(), } } @@ -27,6 +28,7 @@ fn from_params_needs_both_halves() { Some(EphemeralRoute { endpoint: "http://x/openai".into(), api_key: "k".into(), + headers: Vec::new(), }) ); // An endpoint with no credential and a credential with no endpoint are both @@ -55,6 +57,7 @@ fn from_params_treats_blank_as_absent_and_trims() { Some(EphemeralRoute { endpoint: "http://x/openai".into(), api_key: "k".into(), + headers: Vec::new(), }) ); } @@ -155,6 +158,7 @@ fn apply_twice_leaves_one_entry() { EphemeralRoute { endpoint: "http://127.0.0.1:9999/openai".into(), api_key: "mdl-other".into(), + headers: Vec::new(), }, ); let entries: Vec<_> = config diff --git a/crates/openhuman-core/src/cron/scheduler_host_agent_tests.rs b/crates/openhuman-core/src/cron/scheduler_host_agent_tests.rs index 2e175fb9f0a..9daabe667ec 100644 --- a/crates/openhuman-core/src/cron/scheduler_host_agent_tests.rs +++ b/crates/openhuman-core/src/cron/scheduler_host_agent_tests.rs @@ -51,6 +51,7 @@ impl HostAgentResolver for Host { builds.fetch_add(1, Ordering::SeqCst); crate::agent::HostTurnTools::advertised(vec![Box::new(Marker)]) })), + hooks: Default::default(), context: CoreContext::for_test( DomainSet::full(), Some(self.config.workspace_dir.clone()), diff --git a/crates/openhuman-core/src/flows/tinyflows/caps/agent_tests.rs b/crates/openhuman-core/src/flows/tinyflows/caps/agent_tests.rs index 446e1119be0..2428f9accde 100644 --- a/crates/openhuman-core/src/flows/tinyflows/caps/agent_tests.rs +++ b/crates/openhuman-core/src/flows/tinyflows/caps/agent_tests.rs @@ -48,6 +48,7 @@ impl crate::agent::host_agents::HostAgentResolver for FlowsHost { definition, config: crate::config::Config::default(), host_tools: None, + hooks: Default::default(), context: crate::core::runtime::CoreContext::for_test( crate::core::runtime::DomainSet::full(), None, diff --git a/crates/openhuman-core/src/hooks/bridge_tests.rs b/crates/openhuman-core/src/hooks/bridge_tests.rs index 10c060ce959..e5f0d3a1e4b 100644 --- a/crates/openhuman-core/src/hooks/bridge_tests.rs +++ b/crates/openhuman-core/src/hooks/bridge_tests.rs @@ -11,6 +11,7 @@ fn context(tool: &str, arguments: Value) -> ToolHookContext { output: None, error: None, session_id: Some("sess-1".into()), + cwd: None, agent_id: None, } } diff --git a/crates/openhuman-core/src/inference/host_runtime/schemas.rs b/crates/openhuman-core/src/inference/host_runtime/schemas.rs index 04e302f2747..f4616504aac 100644 --- a/crates/openhuman-core/src/inference/host_runtime/schemas.rs +++ b/crates/openhuman-core/src/inference/host_runtime/schemas.rs @@ -43,6 +43,8 @@ struct AgentChatParams { /// other provider. #[serde(default)] api_key: Option, + #[serde(default)] + inference_headers: Vec<(String, String)>, } #[derive(Debug, Deserialize)] @@ -137,6 +139,12 @@ pub fn schemas(function: &str) -> ControllerSchema { "Bearer for inference_url. Scoped to this call and to that \ endpoint alone.", ), + FieldSchema { + name: "inference_headers", + ty: TypeSchema::Option(Box::new(TypeSchema::Array(Box::new(TypeSchema::Json)))), + comment: "Custom header name/value pairs for inference_url only; never persisted.", + required: false, + }, ], outputs: vec![json_output( "response", @@ -232,7 +240,8 @@ fn handle_agent_chat(params: Map) -> ControllerFuture { p.temperature, p.thread_id, p.cwd, - crate::config::schema::EphemeralRoute::from_params(p.inference_url, p.api_key), + crate::config::schema::EphemeralRoute::from_params(p.inference_url, p.api_key) + .map(|route| route.with_headers(p.inference_headers)), ) .await? .into_rpc_json() diff --git a/crates/openhuman-core/src/inference/provider/factory/cloud_slug.rs b/crates/openhuman-core/src/inference/provider/factory/cloud_slug.rs index 07c027b98e3..c857bca9d0f 100644 --- a/crates/openhuman-core/src/inference/provider/factory/cloud_slug.rs +++ b/crates/openhuman-core/src/inference/provider/factory/cloud_slug.rs @@ -302,7 +302,15 @@ pub(super) fn try_create_cloud_slug_chat_model_from_string_with_native_tools( // legacy host's rare 404 → `/v1/responses` fallback for non-codex slugs is // not replicated). let mut endpoint = entry.endpoint.clone(); - let mut extra_headers: Vec<(String, String)> = Vec::new(); + let mut extra_headers: Vec<(String, String)> = config + .ephemeral_route + .as_ref() + .filter(|route| { + slug == crate::config::schema::ephemeral_route::EPHEMERAL_ROUTE_SLUG + && route.endpoint == entry.endpoint + }) + .map(|route| route.headers.clone()) + .unwrap_or_default(); let mut extra_query_params: Vec<(String, String)> = Vec::new(); let mut user_agent: Option = None; let mut responses_api_primary = false; diff --git a/crates/openhuman-core/src/inference/provider/factory_tests.rs b/crates/openhuman-core/src/inference/provider/factory_tests.rs index f55e176cd9a..c7357a49192 100644 --- a/crates/openhuman-core/src/inference/provider/factory_tests.rs +++ b/crates/openhuman-core/src/inference/provider/factory_tests.rs @@ -167,6 +167,7 @@ fn routed_config(endpoint: &str, api_key: &str, model: &str) -> Config { EphemeralRoute { endpoint: endpoint.to_string(), api_key: api_key.to_string(), + headers: Vec::new(), }, ); config diff --git a/crates/openhuman-core/src/runtime/pool/node.rs b/crates/openhuman-core/src/runtime/pool/node.rs index 7c22b78d92f..ca5f7bd37e1 100644 --- a/crates/openhuman-core/src/runtime/pool/node.rs +++ b/crates/openhuman-core/src/runtime/pool/node.rs @@ -8,12 +8,23 @@ use tinyruntime_bus::Language; use super::{PoolExecOutcome, PoolRunError}; use crate::config::{Config, RuntimePoolConfig}; +/// A cancellable dispatch uses the owned-subprocess fallback: the pool has +/// no acknowledged per-job abort seam. +/// /// Whether inline `node` jobs should route through the pool. /// /// Node defaults **on**: each job runs in its own `worker_thread`, so reuse is /// safe — a fresh module graph and fresh globals per job. #[must_use] pub fn enabled(pool: &RuntimePoolConfig) -> bool { + if crate::tools::timeout::ProcessCleanup::is_active() + || crate::tools::timeout::CommandEnvironment::is_active() + { + tracing::debug!( + "[node_exec] cancellable turn uses an owned subprocess instead of the pool" + ); + return false; + } pool.enabled && pool.node.is_enabled(true) } diff --git a/crates/openhuman-core/src/runtime/pool/python.rs b/crates/openhuman-core/src/runtime/pool/python.rs index 82acbffb7e9..420c85956c1 100644 --- a/crates/openhuman-core/src/runtime/pool/python.rs +++ b/crates/openhuman-core/src/runtime/pool/python.rs @@ -8,6 +8,9 @@ use tinyruntime_bus::Language; use super::{PoolExecOutcome, PoolRunError}; use crate::config::{Config, RuntimePoolConfig}; +/// A cancellable dispatch uses the owned-subprocess fallback: the pool has +/// no acknowledged per-job abort seam. +/// /// Whether inline `python` jobs should route through the pool. /// /// Python defaults **off**, and the asymmetry with Node is real rather than an @@ -19,6 +22,14 @@ use crate::config::{Config, RuntimePoolConfig}; /// warm-worker memory saving. #[must_use] pub fn enabled(pool: &RuntimePoolConfig) -> bool { + if crate::tools::timeout::ProcessCleanup::is_active() + || crate::tools::timeout::CommandEnvironment::is_active() + { + tracing::debug!( + "[python_exec] cancellable turn uses an owned subprocess instead of the pool" + ); + return false; + } pool.enabled && pool.python.is_enabled(false) } diff --git a/crates/openhuman-core/src/tools/impl/system/node_exec.rs b/crates/openhuman-core/src/tools/impl/system/node_exec.rs index 111c6b0618f..e14f629f9a8 100644 --- a/crates/openhuman-core/src/tools/impl/system/node_exec.rs +++ b/crates/openhuman-core/src/tools/impl/system/node_exec.rs @@ -303,7 +303,7 @@ impl NodeExecTool { cmd.env_clear(); - let host_path = std::env::var("PATH").unwrap_or_default(); + let host_path = crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(); let sep = if cfg!(windows) { ";" } else { ":" }; let prepended_path = if host_path.is_empty() { resolved.bin_dir.to_string_lossy().into_owned() @@ -313,7 +313,7 @@ impl NodeExecTool { cmd.env("PATH", &prepended_path); for var in SAFE_ENV_VARS { - if let Ok(val) = std::env::var(var) { + if let Ok(val) = crate::tools::timeout::CommandEnvironment::var(var) { cmd.env(var, val); } } @@ -483,7 +483,7 @@ impl NodeExecTool { // process can resolve `node`, `npm`, `npx`, `corepack` consistently // with the unsandboxed path. let mut extra_env = std::collections::HashMap::new(); - let host_path = std::env::var("PATH").unwrap_or_default(); + let host_path = crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(); let sep = if cfg!(windows) { ";" } else { ":" }; let prepended = if host_path.is_empty() { bin_dir.to_string_lossy().into_owned() diff --git a/crates/openhuman-core/src/tools/impl/system/npm_exec.rs b/crates/openhuman-core/src/tools/impl/system/npm_exec.rs index 66accc956fb..a7f018dc9b4 100644 --- a/crates/openhuman-core/src/tools/impl/system/npm_exec.rs +++ b/crates/openhuman-core/src/tools/impl/system/npm_exec.rs @@ -290,7 +290,7 @@ impl NpmExecTool { cmd.env_clear(); - let host_path = std::env::var("PATH").unwrap_or_default(); + let host_path = crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(); let sep = if cfg!(windows) { ";" } else { ":" }; let prepended_path = if host_path.is_empty() { resolved.bin_dir.to_string_lossy().into_owned() @@ -300,7 +300,7 @@ impl NpmExecTool { cmd.env("PATH", &prepended_path); for var in SAFE_ENV_VARS { - if let Ok(val) = std::env::var(var) { + if let Ok(val) = crate::tools::timeout::CommandEnvironment::var(var) { cmd.env(var, val); } } @@ -402,7 +402,7 @@ impl NpmExecTool { // (e.g. `npm run` spawning user scripts) resolve `node`/`npx` // consistently with the unsandboxed path. let mut extra_env = std::collections::HashMap::new(); - let host_path = std::env::var("PATH").unwrap_or_default(); + let host_path = crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(); let sep = if cfg!(windows) { ";" } else { ":" }; let prepended = if host_path.is_empty() { bin_dir.to_string_lossy().into_owned() diff --git a/crates/openhuman-core/src/tools/impl/system/python_exec.rs b/crates/openhuman-core/src/tools/impl/system/python_exec.rs index 4ce48d0d0ed..2c1d96e582c 100644 --- a/crates/openhuman-core/src/tools/impl/system/python_exec.rs +++ b/crates/openhuman-core/src/tools/impl/system/python_exec.rs @@ -304,7 +304,7 @@ impl PythonExecTool { cmd.env_clear(); - let host_path = std::env::var("PATH").unwrap_or_default(); + let host_path = crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(); let sep = if cfg!(windows) { ";" } else { ":" }; let prepended_path = if host_path.is_empty() { resolved.bin_dir.to_string_lossy().into_owned() @@ -316,7 +316,7 @@ impl PythonExecTool { cmd.env("PYTHONUNBUFFERED", "1"); for var in SAFE_ENV_VARS { - if let Ok(val) = std::env::var(var) { + if let Ok(val) = crate::tools::timeout::CommandEnvironment::var(var) { cmd.env(var, val); } } @@ -451,7 +451,7 @@ impl PythonExecTool { } let mut extra_env = std::collections::HashMap::new(); - let host_path = std::env::var("PATH").unwrap_or_default(); + let host_path = crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(); let sep = if cfg!(windows) { ";" } else { ":" }; let prepended = if host_path.is_empty() { bin_dir.to_string_lossy().into_owned() diff --git a/crates/openhuman-core/src/tools/impl/system/shell.rs b/crates/openhuman-core/src/tools/impl/system/shell.rs index 5c2505a2c44..d18e96635d9 100644 --- a/crates/openhuman-core/src/tools/impl/system/shell.rs +++ b/crates/openhuman-core/src/tools/impl/system/shell.rs @@ -405,7 +405,7 @@ impl ShellTool { }; cmd.env_clear(); - for (var, val) in shell_child_env(|name| std::env::var_os(name)) { + for (var, val) in shell_child_env(crate::tools::timeout::CommandEnvironment::var_os) { cmd.env(var, val); } @@ -662,7 +662,7 @@ impl ShellTool { } else { Some(prepend_path_dirs( prepend_dirs.iter().map(|p| p.as_path()), - &std::env::var("PATH").unwrap_or_default(), + &crate::tools::timeout::CommandEnvironment::var("PATH").unwrap_or_default(), )) } } diff --git a/crates/openhuman-core/src/tools/impl/system/shell_tests_runtime_and_sandbox_tests.rs b/crates/openhuman-core/src/tools/impl/system/shell_tests_runtime_and_sandbox_tests.rs index ef4fc9bedc8..e52242571c0 100644 --- a/crates/openhuman-core/src/tools/impl/system/shell_tests_runtime_and_sandbox_tests.rs +++ b/crates/openhuman-core/src/tools/impl/system/shell_tests_runtime_and_sandbox_tests.rs @@ -313,3 +313,38 @@ fn managed_path_restoration_preserves_cmd_syntax_and_unmanaged_commands() { unmanaged_command ); } + +#[cfg(unix)] +#[tokio::test] +async fn managed_python_path_keeps_the_turn_path_without_daemon_inheritance() { + use crate::runtime::python::{PythonSource, ResolvedPython}; + use crate::tools::timeout::CommandEnvironment; + let python = Arc::new(PythonBootstrap::new(Arc::new( + crate::config::Config::default(), + ))); + python.cache_for_test(ResolvedPython { + bin_dir: PathBuf::from("/managed/python/bin"), + python_bin: PathBuf::from("/managed/python/bin/python3"), + version: "3.12.4".into(), + source: PythonSource::Managed, + }); + let tool = ShellTool::with_language_bootstraps( + test_security(AutonomyLevel::Full), + test_runtime(), + test_audit(), + None, + Some(python), + ); + let environment = CommandEnvironment::new([("PATH".into(), "/turn/bin".into())]); + let path = environment + .scope(tool.runtime_path_for_command("python3 -V")) + .await + .unwrap(); + assert_eq!(path, "/managed/python/bin:/turn/bin"); + let empty = CommandEnvironment::new(Vec::<(String, String)>::new()); + let path = empty + .scope(tool.runtime_path_for_command("python3 -V")) + .await + .unwrap(); + assert_eq!(path, "/managed/python/bin"); +} diff --git a/crates/openhuman-core/src/tools/timeout/command_environment.rs b/crates/openhuman-core/src/tools/timeout/command_environment.rs new file mode 100644 index 00000000000..c787db920c0 --- /dev/null +++ b/crates/openhuman-core/src/tools/timeout/command_environment.rs @@ -0,0 +1,67 @@ +//! Child-process environment carried by one dispatch task. + +use std::{collections::BTreeMap, future::Future, sync::Arc}; + +tokio::task_local! { static ACTIVE: CommandEnvironment; } + +/// A complete, turn-owned environment for builtin tool subprocesses. +/// This never mutates the host process environment. Host-spawned tasks must +/// explicitly inherit its scope. Interpreter pools cannot accept scoped work. +#[derive(Clone)] +pub struct CommandEnvironment(Arc>); + +impl CommandEnvironment { + /// Own the exact variables permitted in child processes. + pub fn new(env: impl IntoIterator) -> Self { + Self(Arc::new(env.into_iter().collect())) + } + /// Whether this task requires its own child environment. + pub fn is_active() -> bool { + ACTIVE.try_with(|_| ()).is_ok() + } + /// Run a future using this environment for builtin subprocesses. + pub async fn scope(&self, future: F) -> F::Output { + ACTIVE.scope(self.clone(), future).await + } + /// Read a functional child variable from the scoped map, falling back to + /// the process only when no map was supplied. + pub fn var_os(name: &str) -> Option { + ACTIVE + .try_with(|environment| environment.0.get(name).map(Into::into)) + .unwrap_or_else(|_| std::env::var_os(name)) + } + /// UTF-8 counterpart used by managed interpreter command builders. + pub fn var(name: &str) -> Result { + ACTIVE + .try_with(|environment| { + environment + .0 + .get(name) + .cloned() + .ok_or(std::env::VarError::NotPresent) + }) + .unwrap_or_else(|_| std::env::var(name)) + } + pub(super) fn apply(command: &mut tokio::process::Command) { + let _ = ACTIVE.try_with(|environment| { + // Explicit runtime/security additions (managed PATH, scratch TMP, + // Git policy) remain authoritative over the base environment. + let explicit: Vec<_> = command + .as_std() + .get_envs() + .map(|(key, value)| (key.to_owned(), value.map(ToOwned::to_owned))) + .collect(); + command.env_clear().envs(environment.0.iter()); + for (key, value) in explicit { + match value { + Some(value) => { + command.env(key, value); + } + None => { + command.env_remove(key); + } + } + } + }); + } +} diff --git a/crates/openhuman-core/src/tools/timeout/mod.rs b/crates/openhuman-core/src/tools/timeout/mod.rs index a1629a1be51..c1378afc328 100644 --- a/crates/openhuman-core/src/tools/timeout/mod.rs +++ b/crates/openhuman-core/src/tools/timeout/mod.rs @@ -18,6 +18,8 @@ use std::time::Duration; use tinyagents_harness::tool::ToolTimeoutSettings; use tinytools::ToolTimeout; +mod command_environment; +pub use command_environment::CommandEnvironment; mod process_cleanup; pub use process_cleanup::ProcessCleanup; @@ -237,68 +239,117 @@ pub async fn output_or_kill( /// caller cancels; a scoped [`ProcessCleanup`] can await that waiter. pub async fn output_unbounded( cmd: &mut tokio::process::Command, +) -> std::io::Result { + output_with_input(cmd, None).await +} + +/// Capture a command, optionally supplying stdin, in the current cleanup scope. +pub async fn output_with_input( + cmd: &mut tokio::process::Command, + input: Option>, ) -> std::io::Result { use std::process::Stdio; - cmd.stdin(Stdio::null()) - .stdout(Stdio::piped()) - .stderr(Stdio::piped()) - .kill_on_drop(true); + cmd.stdin(if input.is_some() { + Stdio::piped() + } else { + Stdio::null() + }) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true); + CommandEnvironment::apply(cmd); own_process_group(cmd.as_std_mut()); - let child = cmd.spawn()?; - let pid = child.id(); + let mut child = cmd.spawn()?; + let stdin = child.stdin.take(); let reaped = process_cleanup::Reaped::register(); let (cancel, cancellation) = tokio::sync::watch::channel(false); let waiter = crate::core::runtime::spawn_scoped(async move { let _reaped = reaped; - collect_command_output(child, cancellation).await + let write_input = async move { + if let (Some(mut stdin), Some(input)) = (stdin, input) { + use tokio::io::AsyncWriteExt; + // An early child exit closes stdin. Its exit status/output is + // authoritative; a broken pipe must not hide it. + let _ = stdin.write_all(&input).await; + } + }; + let (_, output) = tokio::join!(write_input, collect_command_output(child, cancellation)); + output }); - let mut group = CommandGroup { pid, cancel }; - let result = waiter.await.map_err(std::io::Error::other)?; - group.pid = None; - result + let _cancel_on_drop = CancelOnDrop(cancel); + waiter.await.map_err(std::io::Error::other)? } -struct CommandGroup { - pid: Option, - cancel: tokio::sync::watch::Sender, -} +// The caller signals only the owned waiter. It never retains a PID after +// that waiter reaps the child, so late future drops cannot kill a reused PID. +struct CancelOnDrop(tokio::sync::watch::Sender); -impl Drop for CommandGroup { +impl Drop for CancelOnDrop { fn drop(&mut self) { - if let Some(pid) = self.pid { + self.0.send_replace(true); + } +} + +struct CommandGroup(Option); + +impl CommandGroup { + fn kill(&self) { + if let Some(pid) = self.0 { kill_process_group(pid); - self.cancel.send_replace(true); } } } +impl Drop for CommandGroup { + fn drop(&mut self) { + self.kill(); + } +} + async fn collect_command_output( mut child: tokio::process::Child, mut cancellation: tokio::sync::watch::Receiver, ) -> std::io::Result { use tokio::io::AsyncReadExt; + let mut group = CommandGroup(child.id()); let mut stdout = child.stdout.take().expect("command stdout is piped"); let mut stderr = child.stderr.take().expect("command stderr is piped"); let mut stdout_bytes = Vec::new(); let mut stderr_bytes = Vec::new(); - let wait = async { + { + let drain = async { + tokio::try_join!( + stdout.read_to_end(&mut stdout_bytes), + stderr.read_to_end(&mut stderr_bytes), + ) + }; + tokio::pin!(drain); tokio::select! { biased; _ = async { let _ = cancellation.wait_for(|cancelled| *cancelled).await; } => { - // Reap the direct child on every platform. On Unix the - // caller has also signalled the process group. - child.kill().await?; - child.wait().await + // The leader has not been reaped, even if it already exited. + // Its PID cannot be reused while signalling this group. + group.kill(); + child.start_kill()?; + if let Ok(result) = tokio::time::timeout(Duration::from_secs(2), drain).await { + result?; + } } - result = child.wait() => result, + result = &mut drain => { result?; } + } + } + let status = tokio::select! { + biased; + _ = async { let _ = cancellation.wait_for(|cancelled| *cancelled).await; } => { + group.kill(); + child.start_kill()?; + child.wait().await? } + result = child.wait() => result?, }; - let (status, _, _) = tokio::try_join!( - wait, - stdout.read_to_end(&mut stdout_bytes), - stderr.read_to_end(&mut stderr_bytes), - )?; + // No await between reaping and disarming. Only this waiter owns the PID. + group.0 = None; Ok(std::process::Output { status, stdout: stdout_bytes, diff --git a/crates/openhuman-core/src/tools/timeout/process_cleanup.rs b/crates/openhuman-core/src/tools/timeout/process_cleanup.rs index 94535c0efcb..566a60e340a 100644 --- a/crates/openhuman-core/src/tools/timeout/process_cleanup.rs +++ b/crates/openhuman-core/src/tools/timeout/process_cleanup.rs @@ -6,7 +6,7 @@ use std::sync::{Arc, Mutex}; use tokio::sync::watch; tokio::task_local! { - static ACTIVE: ProcessCleanup; + static ACTIVE: Vec; } /// Command waiters belonging to one turn. Clone before scoping the turn; @@ -16,9 +16,18 @@ tokio::task_local! { pub struct ProcessCleanup(Arc>>>); impl ProcessCleanup { + /// Whether the current task must acknowledge owned subprocess cleanup. + /// Interpreter pools have no per-job cancellation acknowledgement and + /// therefore cannot accept work from this scope. + pub fn is_active() -> bool { + ACTIVE.try_with(|_| ()).is_ok() + } + /// Run a future with command waiters registered to this turn. pub async fn scope(&self, future: impl Future) -> T { - ACTIVE.scope(self.clone(), future).await + let mut scopes = ACTIVE.try_with(Clone::clone).unwrap_or_default(); + scopes.push(self.clone()); + ACTIVE.scope(scopes, future).await } /// Wait for every registered command to exit and its output pipes to close. @@ -41,12 +50,14 @@ pub(super) struct Reaped(watch::Sender); impl Reaped { pub(super) fn register() -> Self { let (done, waiter) = watch::channel(false); - let _ = ACTIVE.try_with(|cleanup| { - cleanup - .0 - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .push(waiter); + let _ = ACTIVE.try_with(|scopes| { + for cleanup in scopes { + cleanup + .0 + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .push(waiter.clone()); + } }); Self(done) } diff --git a/crates/openhuman-embed/README.md b/crates/openhuman-embed/README.md index c6603b764fe..f7bb593777a 100644 --- a/crates/openhuman-embed/README.md +++ b/crates/openhuman-embed/README.md @@ -8,6 +8,81 @@ The [cookbook](https://github.com/tinyhumansai/openhuman/blob/main/gitbooks/deve Hosts supply transport, credentials and application resources. Runtime settings establish shared defaults; agents narrow provider, access, prompt and tool behavior. Use ProfileRuntime when users require separate credentials and workspaces. +## Cancelling one turn + +Acquire `Turn::cancellation_handle()` before sending a turn. The handle is +cloneable and `cancel().await` waits for the turn to stop and for tracked +commands to be reaped. The agent remains available for later turns: + + + +Cancellation is scoped to this turn, including while waiting for inference. +Turns with a cancellation handle use owned interpreter subprocesses rather +than the Node/Python pool, which has no acknowledged per-job abort API. This +trades warm-worker reuse for awaited cleanup; ordinary turns retain pooling. +The `meter` callback fires once on cancellation or a dropped send future, +with `None` when dispatch has not supplied usage yet. + +Before send, cancellation prevents dispatch; after completion it is a no-op. +On Unix, the built-in shell, Node, Python and npm commands kill their process +group, including descendants. Other platforms stop the direct command. Host +tools that spawn independent tasks or processes must provide their own cleanup; +MCP server lifecycles remain owned by the agent. Keep polling `send()` while +awaiting cancellation, for example in a spawned task. + +## Scoped worker hooks + +Hooks can be supplied at three levels: `RuntimeBuilder::tool_hook` / +`post_turn_hook` for all agents, `AgentSpec::tool_hook` / `post_turn_hook` +for one agent, and `Turn::tool_hook` / `post_turn_hook` for one dispatch. +Tool callbacks run in that order. Named agent updates replace only that agent’s callback; per-turn callbacks are additional. +`ToolHookContext` carries the agent/session identity when known, and `cwd` +follows the execution workspace descriptor (including `Turn::cwd`), falling +back to the embedding context's configured action root. +The agent and turn hooks are never installed in the global registry, so +concurrent workers and later turns do not pick up one another's callbacks. +Post-turn callbacks run asynchronously with an owned session snapshot. +Independently spawned tasks that build sessions must explicitly inherit +`openhuman_core::agent::hooks::HookScope` to carry scoped hooks. + +Gateway attribution headers can be attached to `Route::header(name, value)` +and used with `Turn::route`, or with `Provider::routed(route)` on an agent. +They follow only that route's endpoint and are never saved to configuration, +sent to background providers, or included as values in `Route`'s `Debug`. + +### Inline permission and usage policy + +`AgentSpec::can_use_tool` and `Turn::can_use_tool` await a host callback before +executing each tool. The callback can wait for an approval UI and return +`ToolHookDecision::Proceed`, `Deny`, or `ProceedWith`. It owns that wait; +returning `Ask` denies execution. These callbacks add to existing tool policies, +and a turn callback cannot override an agent denial. + +`AgentSpec::stop_hook` and `Turn::stop_hook` receive cumulative usage after each +completed model call. Return `StopDecision::Continue` to observe usage, or +`Stop` to prevent subsequent calls. Completed tool rounds may still execute; +this is an after-call budget boundary, so hosts must refuse an already exhausted +budget before sending a turn. Provider-reported charges remain authoritative, +including known zero; missing charges remain unknown unless pricing is known. +A policy stop uses a deterministic partial summary rather than spending on +final-answer repair calls. + +### Per-turn tools and subprocess environment + +`Turn::tools` replaces this turn's host tool belt, including attached sources. +An empty belt revokes host tools; the next turn returns to the agent's belt. +Builtin tools still follow the agent definition. This supplies dynamic tools for +in-process hosts; statically declared MCP servers retain their creation-time +configuration. + +`Turn::tool_env` supplies the base environment of owned builtin subprocesses. +Variables absent from it are not inherited from the daemon. The builtin command +builders retain their own security/runtime additions, including Git restrictions, +managed interpreter paths and scratch directories. Scoped turns bypass Node and +Python pools, which cannot acknowledge per-job cancellation or swap a job's +process environment. A host tool that spawns a separate Tokio task must explicitly +carry the command environment and cleanup scopes into that task. + Standalone exact source pins and generated Cargo patches: [consumer setup](CONSUMERS.md). Ordered fallbacks and required exploration: [routing](ROUTING.md). Host telemetry and the existing exporter: [observers](OBSERVERS.md). diff --git a/crates/openhuman-embed/examples/linux_fleet.rs b/crates/openhuman-embed/examples/linux_fleet.rs new file mode 100644 index 00000000000..664297fc4f2 --- /dev/null +++ b/crates/openhuman-embed/examples/linux_fleet.rs @@ -0,0 +1,137 @@ +//! Title: Linux agent fleet memory and latency +//! Summary: Measure retained runtime-owned agents using loopback inference and two worker threads. +//! Run: offline on Linux; use a fresh constrained cgroup for release measurements. +//! Feature: default + +#![recursion_limit = "512"] + +use std::time::Instant; + +use openhuman_embed::{ + Access, AgentDefinitionSpec, AgentSpec, Provider, Runtime, ToolScopeSpec, Workspace, +}; +use serde_json::json; +use wiremock::{Mock, MockServer, ResponseTemplate}; + +fn rss_kib() -> anyhow::Result { + let status = std::fs::read_to_string("/proc/self/status")?; + let line = status + .lines() + .find(|line| line.starts_with("VmRSS:")) + .ok_or_else(|| anyhow::anyhow!("VmRSS missing"))?; + Ok(line + .split_whitespace() + .nth(1) + .ok_or_else(|| anyhow::anyhow!("VmRSS value missing"))? + .parse()?) +} + +// ANCHOR: linux-fleet +fn main() -> anyhow::Result<()> { + let count: usize = std::env::args() + .nth(1) + .unwrap_or_else(|| "100".into()) + .parse()?; + anyhow::ensure!(count > 0, "agent count must be positive"); + openhuman_embed::process::tokio_runtime_builder() + .worker_threads(2) + .build()? + .block_on(measure(count))?; + println!("EXAMPLE_OK linux_fleet"); + Ok(()) +} +// ANCHOR_END: linux-fleet + +async fn measure(count: usize) -> anyhow::Result<()> { + let mock = MockServer::builder() + .disable_request_recording() + .start() + .await; + Mock::given(wiremock::matchers::method("POST")) + .and(wiremock::matchers::path("/v1/chat/completions")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "id":"linux-fleet", "object":"chat.completion", "created":1700000000, + "model":"fixture", "choices":[{"index":0,"message":{"role":"assistant","content":"done"},"finish_reason":"stop"}], + "usage":{"prompt_tokens":1,"completion_tokens":1,"total_tokens":2} + }))) + .mount(&mock).await; + let mut config = openhuman_embed::RuntimeConfig::default(); + config.local_ai.runtime_enabled = false; + config.runtime_python.enabled = false; + config.memory.conversations.enabled = false; + config.agent.session_dual_write = false; + config.agent.session_shadow_reads = false; + config.default_temperature = 0.0; + let before_boot = rss_kib()?; + let boot_start = Instant::now(); + let runtime = Runtime::builder() + .config(config) + .workspace(Workspace::Ephemeral) + .backend_url(mock.uri()) + .max_agents(count) + .build() + .await?; + let boot_ms = boot_start.elapsed().as_secs_f64() * 1000.0; + let baseline = rss_kib()?; + let action = tempfile::tempdir()?; + let agents = (0..count) + .map(|index| { + runtime.agent( + AgentSpec::new(format!("linux-worker-{index}")) + .provider(Provider::openai_compatible( + format!("{}/v1", mock.uri()), + "fixture", + )) + .model("fixture") + .access(Access::full()) + .action_dir(action.path()) + .definition( + AgentDefinitionSpec::new() + .tools(ToolScopeSpec::Named(vec!["shell".into()])), + ), + ) + }) + .collect::, _>>()?; + let registered = rss_kib()?; + let cold_start = Instant::now(); + let first = agents[0] + .turn("Say done.") + .session("cold-session") + .send() + .await?; + anyhow::ensure!(first.reply == "done", "unexpected cold reply"); + let cold_ms = cold_start.elapsed().as_secs_f64() * 1000.0; + let mut tasks = tokio::task::JoinSet::new(); + let fleet_start = Instant::now(); + for (index, agent) in agents.iter().cloned().enumerate() { + tasks.spawn(async move { + let start = Instant::now(); + let outcome = agent + .turn("Say done.") + .session(format!("fleet-session-{index}")) + .send() + .await?; + anyhow::ensure!(outcome.reply == "done", "unexpected fleet reply"); + Ok::<_, anyhow::Error>(start.elapsed().as_secs_f64() * 1000.0) + }); + } + let mut latencies = Vec::with_capacity(count); + while let Some(result) = tasks.join_next().await { + latencies.push(result??); + } + let elapsed_ms = fleet_start.elapsed().as_secs_f64() * 1000.0; + latencies.sort_by(f64::total_cmp); + let completed = rss_kib()?; + let percentile = |p: f64| latencies[((count - 1) as f64 * p).round() as usize]; + println!( + "{}", + json!({ + "agents":count,"worker_threads":2,"inference":"loopback-http-mock","tool_scope":["shell"], + "before_boot_rss_kib":before_boot,"runtime_rss_kib":baseline,"registered_rss_kib":registered, + "after_turns_rss_kib":completed,"marginal_after_turns_mib":completed.saturating_sub(baseline) as f64 / count as f64 / 1024.0, + "boot_ms":boot_ms,"cold_turn_ms":cold_ms,"fleet_elapsed_ms":elapsed_ms, + "turn_p50_ms":percentile(0.5),"turn_p95_ms":percentile(0.95),"turn_p99_ms":percentile(0.99) + }) + ); + Ok(()) +} diff --git a/crates/openhuman-embed/examples/linux_fleet_cgroup.py b/crates/openhuman-embed/examples/linux_fleet_cgroup.py new file mode 100644 index 00000000000..ba8243757a6 --- /dev/null +++ b/crates/openhuman-embed/examples/linux_fleet_cgroup.py @@ -0,0 +1,18 @@ +"""Record the release example and its enclosing Linux cgroup's limits/peak. + +Run through systemd-run --user --scope; build the Rust example first. +The cgroup contains the example's HTTP mock and this small Python monitor. +""" + +import json, pathlib, subprocess, sys +count=sys.argv[1] +binary = pathlib.Path(__file__).resolve().parents[3] / 'target/release/examples/linux_fleet' +r=subprocess.run([str(binary), count],capture_output=True,text=True) +if r.returncode: + print(r.stderr, file=sys.stderr);sys.exit(r.returncode) +record=json.loads(next(line for line in r.stdout.splitlines() if line.startswith("{"))) +path=next(line.split('::',1)[1] for line in pathlib.Path('/proc/self/cgroup').read_text().splitlines() if line.startswith('0::')) +cg=pathlib.Path('/sys/fs/cgroup')/path.lstrip('/') +for name in ['memory.peak','memory.max','cpu.max','memory.events']: + record['cgroup_'+name.replace('.','_')]=(cg/name).read_text().strip() +print(json.dumps(record)) diff --git a/crates/openhuman-embed/src/agent/build.rs b/crates/openhuman-embed/src/agent/build.rs index 42e995d34b7..4b7786e8673 100644 --- a/crates/openhuman-embed/src/agent/build.rs +++ b/crates/openhuman-embed/src/agent/build.rs @@ -261,6 +261,7 @@ pub(crate) fn instantiate(runtime: &Runtime, spec: AgentSpec) -> Result Result, + pub(crate) hooks: openhuman_core::agent::hooks::HookScope, pub(crate) lifecycle: lifecycle::Lifecycle, /// Built from [`ToolScopeSpec::HostOnly`]: every turn's session is built /// from the host tools alone. @@ -232,7 +233,8 @@ impl Agent { /// alone. pub fn turn(&self, message: impl Into) -> Turn { let mut turn = Turn::new(TurnTarget::Agent(Arc::clone(&self.inner)), message) - .with_agent_id(&self.inner.id); + .with_agent_id(&self.inner.id) + .with_hooks(self.inner.hooks.clone()); if let Some(route) = self.inner.provider.route() { turn = turn.route(route.clone()); } diff --git a/crates/openhuman-embed/src/agent/spec.rs b/crates/openhuman-embed/src/agent/spec.rs index bcf21b2c19a..42352eb75d2 100644 --- a/crates/openhuman-embed/src/agent/spec.rs +++ b/crates/openhuman-embed/src/agent/spec.rs @@ -58,6 +58,7 @@ pub struct AgentSpec { composio: Option, config_fn: Option, host_tools: Option, + hooks: openhuman_core::agent::hooks::HookScope, memory: Option, subagents: Vec<(String, AgentDefinitionSpec)>, } @@ -139,6 +140,7 @@ impl AgentSpec { composio: None, config_fn: None, host_tools: None, + hooks: Default::default(), memory: None, subagents: Vec::new(), } @@ -324,6 +326,31 @@ impl AgentSpec { self } + /// Await the host's permission decision before each tool executes. + /// The callback may wait for UI approval, then return `Proceed`, `Deny`, + /// or `ProceedWith`. Returning `Ask` denies the call; this callback itself + /// owns the approval wait. Static tool/security restrictions still apply. + /// Agent and turn callbacks are additive: a denial cannot be overridden. + pub fn can_use_tool(self, callback: F) -> Self + where + F: for<'a> Fn(&'a crate::seams::ToolHookContext) -> crate::PermissionFuture<'a> + + Send + + Sync + + 'static, + { + self.tool_hook(std::sync::Arc::new(crate::permission::PermissionHook( + callback, + ))) + } + + /// Observe cumulative usage after each model call or vote to stop before + /// the next call. Return `StopDecision::Continue` for observation alone; + /// a budget policy can return `StopDecision::Stop`. Scoped to this agent; no runtime-global policy is replaced. + pub fn stop_hook(mut self, hook: std::sync::Arc) -> Self { + self.hooks.push_stop(hook); + self + } + /// Arbitrary edits to the agent's config, applied last. /// /// The escape hatch for the config fields the spec does not model — not @@ -435,6 +462,7 @@ impl AgentSpec { composio: self.composio, config_fn: self.config_fn, host_tools: self.host_tools, + hooks: self.hooks, memory: self.memory, subagents: self.subagents, } @@ -468,6 +496,7 @@ pub(crate) struct AgentSpecParts { pub(crate) composio: Option, pub(crate) config_fn: Option, pub(crate) host_tools: Option, + pub(crate) hooks: openhuman_core::agent::hooks::HookScope, pub(crate) memory: Option, pub(crate) subagents: Vec<(String, AgentDefinitionSpec)>, } diff --git a/crates/openhuman-embed/src/complete.rs b/crates/openhuman-embed/src/complete.rs index 5cbd6a79533..e73b1a5736d 100644 --- a/crates/openhuman-embed/src/complete.rs +++ b/crates/openhuman-embed/src/complete.rs @@ -439,8 +439,8 @@ impl Completer { /// A completer for `route`, with no timeout and no observer. pub fn new(route: Route) -> Self { Self { + headers: route.headers.clone(), route, - headers: Vec::new(), timeout: None, observer: None, cancellation: Default::default(), diff --git a/crates/openhuman-embed/src/harness/provider.rs b/crates/openhuman-embed/src/harness/provider.rs index 21ba6f5b854..5f018dea05b 100644 --- a/crates/openhuman-embed/src/harness/provider.rs +++ b/crates/openhuman-embed/src/harness/provider.rs @@ -83,8 +83,13 @@ impl Provider { /// Both halves are required: the core ignores a route with only one, and /// taking them together here means a partial route cannot be expressed. pub fn openai_compatible(base_url: impl Into, api_key: impl Into) -> Self { + Self::routed(Route::openai_compatible(base_url, api_key)) + } + + /// Use a route including its gateway attribution headers. + pub fn routed(route: Route) -> Self { Self { - route: Some(Route::openai_compatible(base_url, api_key)), + route: Some(route), model: None, custom: None, roles: Default::default(), diff --git a/crates/openhuman-embed/src/lib.rs b/crates/openhuman-embed/src/lib.rs index 4e1e2d4b357..de781e73ab2 100644 --- a/crates/openhuman-embed/src/lib.rs +++ b/crates/openhuman-embed/src/lib.rs @@ -122,7 +122,10 @@ pub mod identity; pub mod memory; #[cfg(feature = "modules")] pub mod modules; +mod permission; +pub use permission::PermissionFuture; pub mod observe; + pub mod process; #[cfg(feature = "channels")] pub mod profiles; @@ -132,6 +135,7 @@ pub mod routing; mod runtime; mod turn; mod turn_cancellation; +mod turn_meter; /// Core internals for `openhuman-tinyhumans` and `openhuman-rpc` only; see /// the module docs. Not part of the host-facing API. @@ -194,6 +198,9 @@ pub use runtime::{ pub mod seams { pub use openhuman_core::agent::hooks::{PostTurnHook, ToolHook}; pub use openhuman_core::agent::hooks::{ToolHookContext, ToolHookDecision, TurnContext}; + pub use openhuman_core::agent::stop_hooks::{ + BudgetStopHook, StopDecision, StopHook, TurnState, + }; pub use openhuman_core::core::all::{ControllerExtension, DomainGroup}; pub use openhuman_core::core::server_launcher::{HostBoot, ServeRequest, ServerLauncher}; pub use openhuman_core::security::SecurityPolicy; diff --git a/crates/openhuman-embed/src/permission.rs b/crates/openhuman-embed/src/permission.rs new file mode 100644 index 00000000000..bd55bf0d0b4 --- /dev/null +++ b/crates/openhuman-embed/src/permission.rs @@ -0,0 +1,27 @@ +//! Inline host approval carried by the agent/turn tool-hook scope. + +use crate::seams::{ToolHook, ToolHookContext, ToolHookDecision}; +/// A sendable asynchronous permission decision borrowing its tool context. +pub type PermissionFuture<'a> = + std::pin::Pin + Send + 'a>>; + +pub(crate) struct PermissionHook(pub(crate) F); + +#[async_trait::async_trait] +impl ToolHook for PermissionHook +where + F: for<'a> Fn(&'a ToolHookContext) -> PermissionFuture<'a> + Send + Sync, +{ + fn name(&self) -> &str { + "embed.can_use_tool" + } + async fn before_tool(&self, _: &ToolHookContext) -> anyhow::Result<()> { + Ok(()) + } + async fn after_tool(&self, _: &ToolHookContext) -> anyhow::Result<()> { + Ok(()) + } + async fn before_tool_decision(&self, context: &ToolHookContext) -> ToolHookDecision { + (self.0)(context).await + } +} diff --git a/crates/openhuman-embed/src/process.rs b/crates/openhuman-embed/src/process.rs index 44120705b55..b1b99c63414 100644 --- a/crates/openhuman-embed/src/process.rs +++ b/crates/openhuman-embed/src/process.rs @@ -11,6 +11,10 @@ use std::path::{Path, PathBuf}; pub use openhuman_core::core::runtime::{AGENT_WORKER_STACK_BYTES, MAX_BLOCKING_THREADS}; +/// Scoped ownership for host commands outside an agent turn. After dropping +/// the scoped future, await `wait()` before acknowledging cancellation. +pub use openhuman_core::tools::timeout::ProcessCleanup as CommandCleanup; + /// Sentry options and event scrubbing shared by desktop and terminal hosts. #[cfg(feature = "crash-reporting")] #[cfg_attr(docsrs, doc(cfg(feature = "crash-reporting")))] @@ -104,3 +108,24 @@ pub fn init_master_key() -> anyhow::Result<()> { #[cfg(test)] #[path = "process_tests.rs"] mod tests; + +/// Run a host command with stdin and a deadline, owning its process group. +/// Normal completion, timeout, and cancellation all reap the command. When +/// called inside a cancellable turn, its cleanup is also tracked by that turn. +/// Explicit command environment settings are preserved over `Turn::tool_env`. +pub async fn command_output( + command: &mut tokio::process::Command, + input: Vec, + deadline: std::time::Duration, +) -> std::io::Result { + let cleanup = openhuman_core::tools::timeout::ProcessCleanup::default(); + let result = cleanup + .scope(tokio::time::timeout( + deadline, + openhuman_core::tools::timeout::output_with_input(command, Some(input)), + )) + .await; + cleanup.wait().await; + result + .map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "host command timed out"))? +} diff --git a/crates/openhuman-embed/src/runtime/host_agents.rs b/crates/openhuman-embed/src/runtime/host_agents.rs index 31b3ad0460f..8cadfbcdb8d 100644 --- a/crates/openhuman-embed/src/runtime/host_agents.rs +++ b/crates/openhuman-embed/src/runtime/host_agents.rs @@ -55,7 +55,9 @@ impl AgentInner { if let Some(route) = openhuman_core::config::schema::EphemeralRoute::from_params( Some(route.base_url.clone()), Some(route.api_key.clone()), - ) { + ) + .map(|scoped| scoped.with_headers(route.headers.clone())) + { openhuman_core::config::schema::ephemeral_route::apply(&mut config, route); } } @@ -64,6 +66,7 @@ impl AgentInner { config, host_tools: self.composed_host_tools(), context: Arc::clone(&self.ctx), + hooks: self.hooks.clone(), }) } } diff --git a/crates/openhuman-embed/src/turn.rs b/crates/openhuman-embed/src/turn.rs index 36c154afb57..1567f4736b7 100644 --- a/crates/openhuman-embed/src/turn.rs +++ b/crates/openhuman-embed/src/turn.rs @@ -88,6 +88,9 @@ pub struct Turn { control_deadline: Option, observer: Option>, trace_content: crate::observe::TraceContent, + hooks: openhuman_core::agent::hooks::HookScope, + tools: Option, + tool_env: Option, } impl Turn { @@ -114,6 +117,9 @@ impl Turn { control_deadline: None, observer: None, trace_content: crate::observe::TraceContent::MetadataOnly, + hooks: Default::default(), + tools: None, + tool_env: None, } } @@ -122,6 +128,75 @@ impl Turn { self } + pub(crate) fn with_hooks(mut self, hooks: openhuman_core::agent::hooks::HookScope) -> Self { + self.hooks = hooks; + self + } + + /// Add a tool callback for this turn only, after runtime and agent hooks. + /// Does not change the agent or any other concurrent turn. + pub fn tool_hook(mut self, hook: Arc) -> Self { + self.hooks.push_tool(hook); + self + } + + /// Replace this turn's host tools, including attached sources. An empty + /// belt revokes them. Builtin tools remain governed by the agent definition; + /// the agent's original host tools return on its next turn. + pub fn tools( + mut self, + factory: impl for<'a> Fn( + openhuman_core::agent::TurnContext<'a>, + ) -> openhuman_core::agent::HostTurnTools + + Send + + Sync + + 'static, + ) -> Self { + self.tools = Some(Arc::new(factory)); + self + } + + /// Replace the environment of owned builtin tool subprocesses for this + /// turn. Variables absent from this map are not inherited from the daemon. + /// Interpreter pools are bypassed so a pooled process cannot carry another + /// turn's environment. Independently spawned host tasks must carry the scope. + pub fn tool_env(mut self, env: impl IntoIterator) -> Self { + self.tool_env = Some(openhuman_core::tools::timeout::CommandEnvironment::new(env)); + self + } + + /// Await the host's permission decision before each tool executes. + /// The callback may wait for UI approval, then return `Proceed`, `Deny`, + /// or `ProceedWith`. Returning `Ask` denies the call; this callback itself + /// owns the approval wait. Static tool/security restrictions still apply. + /// Agent and turn callbacks are additive: a denial cannot be overridden. + pub fn can_use_tool(self, callback: F) -> Self + where + F: for<'a> Fn(&'a crate::seams::ToolHookContext) -> crate::PermissionFuture<'a> + + Send + + Sync + + 'static, + { + self.tool_hook(std::sync::Arc::new(crate::permission::PermissionHook( + callback, + ))) + } + + /// Observe cumulative usage after each model call or vote to stop before + /// the next call. Return `StopDecision::Continue` for observation alone; + /// a budget policy can return `StopDecision::Stop`. Scoped to this turn; no runtime-global policy is replaced. + pub fn stop_hook(mut self, hook: std::sync::Arc) -> Self { + self.hooks.push_stop(hook); + self + } + + /// Add a completed-turn callback for this turn only. The callback runs + /// asynchronously with an owned snapshot after the turn completes. + pub fn post_turn_hook(mut self, hook: Arc) -> Self { + self.hooks.push_post_turn(hook); + self + } + /// Continue an existing conversation. Without this a fresh session id is /// minted and returned in [`TurnOutcome::session_id`]. pub fn session(mut self, session_id: impl Into) -> Self { @@ -298,6 +373,7 @@ impl Turn { pub fn route(mut self, route: Route) -> Self { self.request.inference_url = Some(route.base_url); self.request.api_key = Some(route.api_key); + self.request.inference_headers = route.headers; self } @@ -358,7 +434,7 @@ impl Turn { /// DomainSet gate itself before touching the core. use openhuman_core::agent::tinyagents::host::LastTurnUsage; -type UsageSink = std::sync::Mutex>; +use crate::turn_meter::UsageSink; use openhuman_core::agent::tinyagents::response_shape::{FinalResponse, ResponseShapeScope}; @@ -374,6 +450,7 @@ mod futures_box { /// What only an agent target can honour, already validated by /// [`Turn::validate_turn_options`]. struct AgentTurnOptions { + tools: Option, shape: std::sync::Arc, untrusted_input: bool, } @@ -432,8 +509,9 @@ async fn dispatch( let route = openhuman_core::config::schema::EphemeralRoute::from_params( request.inference_url, request.api_key, - ); - let host = inner.composed_host_tools(); + ) + .map(|route| route.with_headers(request.inference_headers)); + let host = options.tools.or_else(|| inner.composed_host_tools()); let target = AgentChatTarget::Definition { definition: &inner.definition, host: host.as_ref(), diff --git a/crates/openhuman-embed/src/turn_control.rs b/crates/openhuman-embed/src/turn_control.rs index ba160fa2c06..caed7bb6224 100644 --- a/crates/openhuman-embed/src/turn_control.rs +++ b/crates/openhuman-embed/src/turn_control.rs @@ -24,6 +24,24 @@ impl Turn { /// that is a build/composition fact, not a failure, and a host should hide /// the surface rather than report an error. pub async fn send(mut self) -> Result { + let hooks = std::mem::take(&mut self.hooks); + let environment = self.tool_env.take(); + let fresh_tools = self.tools.is_some(); + let dispatch = hooks.scope(Box::pin(self.send_scoped())); + let dispatch = async move { + if fresh_tools { + openhuman_core::agent::tool_snapshot_scope::with_fresh_snapshot(dispatch).await + } else { + dispatch.await + } + }; + match environment { + Some(environment) => environment.scope(dispatch).await, + None => dispatch.await, + } + } + + async fn send_scoped(mut self) -> Result { // External cancellation cascades inward; cancelling this turn through // its acknowledgement handle/deadline must not cancel a shared parent. let token = self @@ -177,12 +195,14 @@ impl Turn { let deadline = self.control_deadline; let token = self.token_cancellation.clone(); let native = openhuman_core::agent::host_overrides::current_cancellation(); - let cancellation = match self.cancellation.take() { - Some(cancellation) => cancellation, - None if deadline.is_some() || token.is_some() => crate::TurnCancellation::default(), - None => return Box::pin(self.send_inner()).await, + let cancellation = self.cancellation.take().or_else(|| { + (deadline.is_some() || token.is_some()).then(crate::TurnCancellation::default) + }); + let _guard = cancellation.as_ref().map(crate::TurnCancellation::enter); + let meter = crate::turn_meter::TurnMeter::new(self.meter.take()); + let Some(cancellation) = cancellation else { + return Box::pin(self.send_inner(&meter.usage)).await; }; - let _guard = cancellation.enter(); let outcome = cancellation .cleanup() .scope(async { @@ -208,7 +228,7 @@ impl Turn { native.cancel(); Err(CoreError::DeadlineExceeded { method: AGENT_CHAT }) }, - outcome = Box::pin(self.send_inner()) => outcome, + outcome = Box::pin(self.send_inner(&meter.usage)) => outcome, } }) .await; @@ -248,7 +268,7 @@ impl Turn { .clone() } - async fn send_inner(mut self) -> Result { + async fn send_inner(mut self, usage: &UsageSink) -> Result { // The core neither mints nor returns a session id, so continuing a // conversation would otherwise be impossible without the caller // inventing an id scheme — which every embedder has then done @@ -296,8 +316,6 @@ impl Turn { // Filled by the turn itself, before any error is raised, so a failed // turn is still metered. Read back below whether the dispatch returned // a reply or an error. - let usage: UsageSink = std::sync::Mutex::new(None); - let meter = self.meter.take(); let wants_json = self .response_format .as_ref() @@ -332,12 +350,13 @@ impl Turn { }, ), untrusted_input: self.untrusted_input, + tools: self.tools.take(), }; let budget = self.budget.take().map(|budget| crate::budget::ModelBudget { ledger: budget.ledger.child(crate::budget::SpendLimits::default()), call: budget.call, }); - let dispatch = dispatch(self.target, self.request, self.seed.take(), &usage, options); + let dispatch = dispatch(self.target, self.request, self.seed.take(), usage, options); let dispatch = async { match &budget { Some(budget) => { @@ -389,17 +408,6 @@ impl Turn { log::debug!("[embed][agent] turn_failed session={session_id} kind={tag}"); }); - // Before the `?`. A turn that errored still spent what it spent, and - // this is the only place both the sink and a failing result are in - // hand -- `TurnOutcome` below is never built on that path. - if let Some(meter) = meter { - meter( - usage - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(), - ); - } let reply = reply.map_err(|error| { match budget.as_ref().and_then(|budget| budget.ledger.refusal()) { Some(source) => CoreError::BudgetExceeded { @@ -429,8 +437,9 @@ impl Turn { reply, session_id, usage: usage - .into_inner() - .unwrap_or_else(std::sync::PoisonError::into_inner), + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone(), structured, finish_reason: report.finish_reason, answered_model: report.answered_model, @@ -473,6 +482,12 @@ impl Turn { let host_only = match &self.target { TurnTarget::Agent(agent) => agent.host_only, TurnTarget::Runtime(_) => { + if self.tools.is_some() { + return refuse( + "per-turn host tools need a runtime-owned Agent", + "turn_tools_unsupported", + ); + } if self.response_format.is_some() || self.max_tokens.is_some() || self.top_p.is_some() diff --git a/crates/openhuman-embed/src/turn_meter.rs b/crates/openhuman-embed/src/turn_meter.rs new file mode 100644 index 00000000000..495ecf6788d --- /dev/null +++ b/crates/openhuman-embed/src/turn_meter.rs @@ -0,0 +1,32 @@ +//! A turn's usage callback survives dispatch errors, cancellation and dropping. + +use openhuman_core::agent::tinyagents::host::LastTurnUsage; + +pub(crate) type UsageSink = std::sync::Mutex>; + +pub(crate) struct TurnMeter { + pub(crate) usage: UsageSink, + callback: Option) + Send>>, +} + +impl TurnMeter { + pub(crate) fn new(callback: Option) + Send>>) -> Self { + Self { + usage: std::sync::Mutex::new(None), + callback, + } + } +} + +impl Drop for TurnMeter { + fn drop(&mut self) { + if let Some(callback) = self.callback.take() { + callback( + self.usage + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone(), + ); + } + } +} diff --git a/crates/openhuman-embed/src/turn_tests.rs b/crates/openhuman-embed/src/turn_tests.rs index f9d9a51ae65..7c8adae0ef8 100644 --- a/crates/openhuman-embed/src/turn_tests.rs +++ b/crates/openhuman-embed/src/turn_tests.rs @@ -31,6 +31,7 @@ fn turn_request_field_names_match_the_controller() { cwd: Some("/tmp".into()), inference_url: Some("https://example.invalid/v1".into()), api_key: Some("k".into()), + inference_headers: vec![("x-worker".into(), "a".into())], agent_id: Some("a".into()), }; @@ -203,13 +204,14 @@ fn the_controller_declares_nowhere_for_a_seed_to_travel() { // `context` or `transcript` could carry history just as well, and a // heuristic that guesses at names would pass it silently. Anything new // fails here until someone decides whether a seed could ride it. - const KNOWN: [&str; 8] = [ + const KNOWN: [&str; 9] = [ "message", "model_override", "temperature", "thread_id", "cwd", "inference_url", + "inference_headers", "api_key", "agent_id", ]; diff --git a/crates/openhuman-embed/src/turn_types.rs b/crates/openhuman-embed/src/turn_types.rs index 524c08e6187..6ffe8dbee52 100644 --- a/crates/openhuman-embed/src/turn_types.rs +++ b/crates/openhuman-embed/src/turn_types.rs @@ -68,6 +68,8 @@ pub struct Route { pub base_url: String, /// The bearer presented to `base_url`. pub api_key: String, + /// Headers scoped to this route and never persisted or logged as values. + pub headers: Vec<(String, String)>, } impl std::fmt::Debug for Route { @@ -79,6 +81,7 @@ impl std::fmt::Debug for Route { f.debug_struct("Route") .field("base_url", &sanitize_url_for_display(&self.base_url)) .field("api_key", &"") + .field("headers", &self.headers.len()) .finish() } } @@ -136,8 +139,14 @@ impl Route { Self { base_url: base_url.into(), api_key: api_key.into(), + headers: Vec::new(), } } + /// Add a header sent only to this route's endpoint. + pub fn header(mut self, name: impl Into, value: impl Into) -> Self { + self.headers.push((name.into(), value.into())); + self + } } /// Wire params for [`super::AGENT_CHAT`]. @@ -167,6 +176,9 @@ pub struct TurnRequest { /// Endpoint half of the per-call route. Paired with `api_key`. #[serde(skip_serializing_if = "Option::is_none")] pub inference_url: Option, + /// Additional headers owned by this turn route. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub inference_headers: Vec<(String, String)>, /// Bearer half of the per-call route. Paired with `inference_url`. #[serde(skip_serializing_if = "Option::is_none")] pub api_key: Option, @@ -186,6 +198,7 @@ impl TurnRequest { thread_id: None, cwd: None, inference_url: None, + inference_headers: Vec::new(), api_key: None, agent_id: None, } diff --git a/crates/openhuman-embed/tests/README.md b/crates/openhuman-embed/tests/README.md index e36b5afa002..a47e9ae0278 100644 --- a/crates/openhuman-embed/tests/README.md +++ b/crates/openhuman-embed/tests/README.md @@ -46,6 +46,13 @@ recorded on the wrong server. | [`public_api.rs`](public_api.rs) | Compile-time check that the host-facing types and signatures stay exported. | | [`turn_cancellation.rs`](turn_cancellation.rs) | Cancellation before send, during inference and during a builtin shell command; repeated requests and agent reuse. | | [`process_cancellation.rs`](process_cancellation.rs) | On Linux, dropping a command future kills its shell descendants. | +| `inline_permissions.rs` | An inline UI decision precedes execution; concurrent agents and one-turn denials stay isolated. | +| `usage_hooks.rs` | Per-agent stop policy and per-turn cumulative usage observation, with provider charges and no extra calls after a stop. | +| [`turn_tools.rs`](turn_tools.rs) | One-turn belt replacement/revocation, resumed-session schemas and independent concurrent workers. | +| [`tool_environment.rs`](tool_environment.rs) | Overlapping child environments exclude inherited variables and stay within their own scopes. | +| [`scoped_hooks.rs`](scoped_hooks.rs) | Same-named runtime, agent and turn callbacks remain additive and isolated during concurrent turns, resumed sessions and reuse of a removed agent id. | +| [`route_headers.rs`](route_headers.rs) | Attribution headers and bearer follow only the per-turn route, without reaching other agents, later turns or backend calls. | +| [`tool_hook_context.rs`](tool_hook_context.rs) | Hook agent/session identities and cwd agree with builtin shell execution on an overridden and a default working root. | ## Running diff --git a/crates/openhuman-embed/tests/inline_permissions.rs b/crates/openhuman-embed/tests/inline_permissions.rs new file mode 100644 index 00000000000..bcfcff253a0 --- /dev/null +++ b/crates/openhuman-embed/tests/inline_permissions.rs @@ -0,0 +1,168 @@ +//! A host can await its own approval UI inline without polling the core. + +mod common; + +use common::{ + chat_completion, offline_config, route, runtime, scripted_provider, stub_backend, + tool_call_completion, +}; +use openhuman_embed::seams::{ToolHookContext, ToolHookDecision}; +use openhuman_embed::{ + AgentDefinitionSpec, AgentSpec, HostTurnTools, Runtime, Tool, ToolScopeSpec, Workspace, +}; +use serde_json::{json, Value}; +use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, +}; + +struct Probe(Arc); +#[async_trait::async_trait] +impl Tool for Probe { + fn name(&self) -> &str { + "probe" + } + fn description(&self) -> &str { + "Perform a host operation" + } + fn parameters_schema(&self) -> Value { + json!({"type":"object","properties":{}}) + } + async fn execute(&self, _: Value) -> anyhow::Result { + self.0.fetch_add(1, Ordering::SeqCst); + Ok(openhuman_core::tools::ToolResult::success("ran")) + } +} +#[test] +fn inline_approval_waits_for_the_host_and_isolated_turn_policies_can_only_narrow() { + runtime().block_on(async { tokio::spawn(scenario()).await.unwrap() }); +} +async fn scenario() { + let backend = stub_backend().await; + let provider = scripted_provider( + vec![ + tool_call_completion("probe", "{}"), + tool_call_completion("probe", "{}"), + tool_call_completion("probe", "{}"), + chat_completion("done"), + tool_call_completion("probe", "{}"), + tool_call_completion("probe", "{}"), + ], + "done", + ) + .await; + let other_provider = + scripted_provider(vec![tool_call_completion("probe", "{}")], "other done").await; + let runtime = Runtime::builder() + .config(offline_config()) + .workspace(Workspace::Ephemeral) + .backend_url(backend.uri()) + .build() + .await + .unwrap(); + let calls = Arc::new(AtomicUsize::new(0)); + let other_calls = Arc::new(AtomicUsize::new(0)); + let spec = |id, provider: &wiremock::MockServer, calls: Arc| { + AgentSpec::new(id) + .provider(route(provider, "fixture")) + .definition( + AgentDefinitionSpec::new() + .bare_prompt("Use probe then answer") + .tools(ToolScopeSpec::HostOnly), + ) + .tools(move |_| HostTurnTools::advertised(vec![Box::new(Probe(calls.clone()))])) + }; + let (requests, mut received) = tokio::sync::mpsc::unbounded_channel::<( + ToolHookContext, + tokio::sync::oneshot::Sender, + )>(); + let agent = runtime + .agent( + spec("interactive", &provider, calls.clone()).can_use_tool(move |context| { + let (answer, decision) = tokio::sync::oneshot::channel(); + requests.send((context.clone(), answer)).unwrap(); + Box::pin(async move { + decision + .await + .unwrap_or_else(|_| ToolHookDecision::Deny("UI closed".into())) + }) + }), + ) + .unwrap(); + let other = runtime + .agent( + spec("other", &other_provider, other_calls.clone()) + .can_use_tool(|_| Box::pin(async { ToolHookDecision::Proceed })), + ) + .unwrap(); + let first = tokio::spawn(agent.turn("needs approval").send()); + let (context, answer) = + tokio::time::timeout(std::time::Duration::from_secs(5), received.recv()) + .await + .unwrap() + .unwrap(); + assert_eq!(context.agent_id.as_deref(), Some("interactive")); + assert_eq!(context.tool_name, "probe"); + assert_eq!( + calls.load(Ordering::SeqCst), + 0, + "tool ran before the UI answered" + ); + other.run("independent worker").await.unwrap(); + assert_eq!(other_calls.load(Ordering::SeqCst), 1); + answer + .send(ToolHookDecision::Deny("human refused".into())) + .unwrap(); + assert!( + first.await.unwrap().is_err(), + "denied call unexpectedly succeeded" + ); + assert_eq!(calls.load(Ordering::SeqCst), 0); + let second = tokio::spawn( + agent + .turn("turn-specific refusal") + .can_use_tool(|_| { + Box::pin(async { ToolHookDecision::Deny("turn narrowed authority".into()) }) + }) + .send(), + ); + let (_, answer) = received.recv().await.unwrap(); + answer.send(ToolHookDecision::Proceed).unwrap(); + assert!( + second.await.unwrap().is_err(), + "turn denial unexpectedly succeeded" + ); + assert_eq!(calls.load(Ordering::SeqCst), 0, "turn policy was skipped"); + let third = tokio::spawn(agent.turn("next turn").send()); + let (_, answer) = received.recv().await.unwrap(); + answer.send(ToolHookDecision::Proceed).unwrap(); + third.await.unwrap().unwrap(); + assert_eq!( + calls.load(Ordering::SeqCst), + 1, + "previous turn's refusal leaked" + ); + let mut pending = agent.turn("cancel while the UI waits"); + let cancellation = pending.cancellation_handle(); + let cancelled = tokio::spawn(pending.send()); + let (_, answer) = received.recv().await.unwrap(); + cancellation.cancel().await; + assert!(matches!( + cancelled.await.unwrap(), + Err(openhuman_embed::CoreError::TurnCancelled { .. }) + )); + assert!( + answer.is_closed(), + "cancelled turn left an approval receiver alive" + ); + assert_eq!(calls.load(Ordering::SeqCst), 1); + let closed_ui = tokio::spawn(agent.turn("closed UI refuses").send()); + let (_, answer) = received.recv().await.unwrap(); + drop(answer); + assert!(closed_ui.await.unwrap().is_err()); + assert_eq!( + calls.load(Ordering::SeqCst), + 1, + "closing the UI granted permission" + ); +} diff --git a/crates/openhuman-embed/tests/process_cancellation.rs b/crates/openhuman-embed/tests/process_cancellation.rs index f353454cfa7..479c9e42145 100644 --- a/crates/openhuman-embed/tests/process_cancellation.rs +++ b/crates/openhuman-embed/tests/process_cancellation.rs @@ -121,3 +121,141 @@ async fn deadline_kills_and_reaps_a_command() { .await .unwrap(); } + +#[tokio::test] +async fn exited_group_leader_is_not_reaped_until_its_descendants_close_the_pipes() { + let scratch = tempfile::tempdir().unwrap(); + let pidfile = scratch.path().join("exited.pid"); + let mut cmd = openhuman_core::agent::platform_shell::build_tokio_command(&format!( + "sleep 30 & echo $$ $! > {}; exit 0", + pidfile.display() + )); + let cleanup = ProcessCleanup::default(); + let mut run = Box::pin(cleanup.scope(output_unbounded(&mut cmd))); + let pids = tokio::time::timeout(Duration::from_secs(5), async { + loop { + tokio::select! { + result = &mut run => panic!("pipes closed early: {result:?}"), + _ = tokio::time::sleep(Duration::from_millis(10)) => {} + } + if let Ok(text) = std::fs::read_to_string(&pidfile) { + let pids: Vec = text + .split_whitespace() + .map(|p| p.parse().unwrap()) + .collect(); + break pids; + } + } + }) + .await + .unwrap(); + // Poll the output future while the waiter processes the shell's exit. + tokio::select! { + result = &mut run => panic!("pipes closed early: {result:?}"), + _ = tokio::time::sleep(Duration::from_millis(100)) => {} + } + let leader_reserved = std::path::Path::new(&format!("/proc/{}", pids[0])).exists(); + drop(run); + tokio::time::timeout(Duration::from_secs(2), cleanup.wait()) + .await + .unwrap(); + assert!( + leader_reserved, + "reaped leader PID could be reused while cancellation still addresses its group" + ); + assert!(!std::path::Path::new(&format!("/proc/{}", pids[0])).exists()); +} + +#[tokio::test] +async fn cancellable_turns_route_interpreters_away_from_unacknowledged_pools() { + let mut config = openhuman_core::config::RuntimePoolConfig::default(); + config.python.enabled = Some(true); + assert!(openhuman_core::runtime::pool::python::enabled(&config)); + assert!(openhuman_core::runtime::pool::node::enabled(&config)); + let cleanup = ProcessCleanup::default(); + cleanup + .scope(async { + tokio::task::yield_now().await; + assert!(!openhuman_core::runtime::pool::python::enabled(&config)); + assert!( + !openhuman_core::runtime::pool::node::enabled(&config), + "cancellable node jobs must use an owned subprocess" + ); + }) + .await; + assert!(openhuman_core::runtime::pool::node::enabled(&config)); +} + +#[tokio::test] +async fn host_commands_receive_stdin_and_timeout_after_reaping() { + let mut cmd = tokio::process::Command::new("/bin/sh"); + cmd.args(["-c", "cat; printf stderr >&2"]); + let output = openhuman_embed::process::command_output( + &mut cmd, + b"payload".to_vec(), + Duration::from_secs(2), + ) + .await + .unwrap(); + assert_eq!(output.stdout, b"payload"); + assert_eq!(output.stderr, b"stderr"); + let scratch = tempfile::tempdir().unwrap(); + let pidfile = scratch.path().join("host.pid"); + let mut cmd = tokio::process::Command::new("/bin/sh"); + cmd.args(["-c", &format!("echo $$ > {}; sleep 30", pidfile.display())]); + let outer = openhuman_embed::process::CommandCleanup::default(); + let result = outer + .scope(openhuman_embed::process::command_output( + &mut cmd, + Vec::new(), + Duration::from_millis(100), + )) + .await; + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::TimedOut); + outer.wait().await; + let pid = std::fs::read_to_string(pidfile).unwrap(); + assert!(!std::path::Path::new(&format!("/proc/{}", pid.trim())).exists()); +} + +#[tokio::test] +async fn cancellation_reaps_the_leader_when_an_escaped_descendant_holds_its_pipes() { + let scratch = tempfile::tempdir().unwrap(); + let pidfile = scratch.path().join("escaped.pid"); + let mut command = tokio::process::Command::new("/bin/sh"); + command.args([ + "-c", + &format!( + "setsid sh -c 'echo $$ > {}; sleep 30' & wait", + pidfile.display() + ), + ]); + let cleanup = openhuman_embed::process::CommandCleanup::default(); + let mut run = Box::pin(cleanup.scope(output_unbounded(&mut command))); + let pid = tokio::time::timeout(Duration::from_secs(5), async { + loop { + tokio::select! { + result = &mut run => panic!("command ended early: {result:?}"), + _ = tokio::time::sleep(Duration::from_millis(10)) => {} + } + if let Ok(pid) = std::fs::read_to_string(&pidfile) { + break pid; + } + } + }) + .await + .unwrap(); + drop(run); + let settled = tokio::time::timeout(Duration::from_secs(4), cleanup.wait()).await; + // This deliberately escaped group belongs to the fixture; clean it on + // both red and green paths before asserting the bounded acknowledgement. + let status = std::process::Command::new("kill") + .args(["-KILL", "--", &format!("-{}", pid.trim())]) + .status() + .unwrap(); + assert!(status.success()); + cleanup.wait().await; + assert!( + settled.is_ok(), + "escaped descendant blocked cancellation acknowledgement" + ); +} diff --git a/crates/openhuman-embed/tests/route_headers.rs b/crates/openhuman-embed/tests/route_headers.rs new file mode 100644 index 00000000000..e341a27a651 --- /dev/null +++ b/crates/openhuman-embed/tests/route_headers.rs @@ -0,0 +1,77 @@ +//! Per-route attribution headers follow only the endpoint and turn named by the host. + +mod common; + +use common::{chat_requests, offline_config, provider, route, runtime, stub_backend}; +use openhuman_embed::{AgentSpec, Provider, Route, Runtime, Workspace}; + +#[test] +fn attribution_headers_are_scoped_to_each_agent_and_turn_route() { + runtime().block_on(async { tokio::spawn(scenario()).await.unwrap() }); +} + +async fn scenario() { + let backend = stub_backend().await; + let (a_provider, b_provider) = tokio::join!(provider("a"), provider("b")); + let runtime = Runtime::builder() + .config(offline_config()) + .workspace(Workspace::Ephemeral) + .backend_url(backend.uri()) + .build() + .await + .unwrap(); + let a_route = Route::openai_compatible(format!("{}/v1", a_provider.uri()), "a-key") + .header("x-medulla-agent-id", "worker-a") + .header("x-medulla-run-id", "a-run"); + assert!( + !format!("{a_route:?}").contains("a-run"), + "header values must stay out of Debug" + ); + let a = runtime + .agent(AgentSpec::new("worker-a").provider(route(&a_provider, "fixture"))) + .unwrap(); + let b = runtime + .agent( + AgentSpec::new("worker-b").provider( + Provider::routed( + Route::openai_compatible(format!("{}/v1", b_provider.uri()), "b-key") + .header("x-medulla-agent-id", "worker-b"), + ) + .model("fixture"), + ), + ) + .unwrap(); + let (a_result, b_result) = tokio::join!(a.turn("a").route(a_route).send(), b.run("b")); + a_result.unwrap(); + b_result.unwrap(); + // The next turn inherits the agent provider, not the previous turn's headers. + a.run("next a").await.unwrap(); + let a_requests = chat_requests(&a_provider).await; + let b_requests = chat_requests(&b_provider).await; + assert_eq!(a_requests.len(), 2); + assert_eq!(b_requests.len(), 1); + assert_eq!( + a_requests[0].headers.get("x-medulla-agent-id").unwrap(), + "worker-a" + ); + assert_eq!( + a_requests[0].headers.get("x-medulla-run-id").unwrap(), + "a-run" + ); + assert_eq!( + a_requests[0].headers.get("authorization").unwrap(), + "Bearer a-key" + ); + assert!(a_requests[1].headers.get("x-medulla-run-id").is_none()); + assert_eq!( + b_requests[0].headers.get("x-medulla-agent-id").unwrap(), + "worker-b" + ); + assert!(b_requests[0].headers.get("x-medulla-run-id").is_none()); + assert!(backend + .received_requests() + .await + .unwrap() + .iter() + .all(|r| r.headers.get("x-medulla-run-id").is_none())); +} diff --git a/crates/openhuman-embed/tests/scoped_hooks.rs b/crates/openhuman-embed/tests/scoped_hooks.rs new file mode 100644 index 00000000000..3a7ae3ae963 --- /dev/null +++ b/crates/openhuman-embed/tests/scoped_hooks.rs @@ -0,0 +1,256 @@ +//! Agent and turn hooks stay isolated while concurrent workers share a runtime. + +mod common; + +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex}; + +use common::{ + chat_completion, eventually, offline_config, route, runtime, scripted_provider, stub_backend, + tool_call_completion, +}; +use openhuman_embed::seams::{PostTurnHook, ToolHook, ToolHookContext, TurnContext}; +use openhuman_embed::{ + AgentDefinitionSpec, AgentSpec, HostTurnTools, Runtime, Tool, ToolScopeSpec, Workspace, +}; +use serde_json::{json, Value}; +use tokio::sync::Barrier; +use wiremock::MockServer; + +type Events = Arc>>; + +struct Observer { + label: &'static str, + events: Events, + first: AtomicBool, + overlap: Option>, +} + +impl Observer { + fn new(label: &'static str, events: &Events, overlap: Option>) -> Arc { + Arc::new(Self { + label, + events: events.clone(), + first: AtomicBool::new(true), + overlap, + }) + } + fn record(&self, phase: &str, owner: &str) { + self.events + .lock() + .unwrap() + .push((self.label.into(), phase.into(), owner.into())); + } +} + +#[async_trait::async_trait] +impl ToolHook for Observer { + fn name(&self) -> &str { + "same-policy-name" + } + async fn before_tool(&self, ctx: &ToolHookContext) -> anyhow::Result<()> { + self.record("before", ctx.arguments["owner"].as_str().unwrap()); + if self.first.swap(false, Ordering::SeqCst) { + if let Some(barrier) = &self.overlap { + tokio::time::timeout(std::time::Duration::from_secs(10), barrier.wait()).await?; + } + } + Ok(()) + } + async fn after_tool(&self, ctx: &ToolHookContext) -> anyhow::Result<()> { + assert_eq!(ctx.success, Some(true)); + self.record("after", ctx.arguments["owner"].as_str().unwrap()); + Ok(()) + } +} + +#[async_trait::async_trait] +impl PostTurnHook for Observer { + fn name(&self) -> &str { + "same-policy-name" + } + async fn on_turn_complete(&self, ctx: &TurnContext) -> anyhow::Result<()> { + assert!(ctx.session_id.is_some()); + assert!(ctx.agent_id.is_some()); + self.record("turn", ctx.user_message.lines().last().unwrap()); + Ok(()) + } +} + +struct Probe; +#[async_trait::async_trait] +impl Tool for Probe { + fn name(&self) -> &str { + "probe" + } + fn description(&self) -> &str { + "Report a worker's identity" + } + fn parameters_schema(&self) -> Value { + json!({"type":"object","properties":{"owner":{"type":"string"}},"required":["owner"]}) + } + async fn execute(&self, _: Value) -> anyhow::Result { + Ok(openhuman_core::tools::ToolResult::success("observed")) + } +} + +fn spec(id: &str, provider: &MockServer) -> AgentSpec { + AgentSpec::new(id) + .provider(route(provider, "fixture")) + .definition( + AgentDefinitionSpec::new() + .bare_prompt("Use probe once, then answer.") + .tools(ToolScopeSpec::HostOnly), + ) + .tools(|_| HostTurnTools::advertised(vec![Box::new(Probe)])) +} + +async fn worker_provider(owner: &str) -> MockServer { + scripted_provider( + (0..4) + .flat_map(|_| { + [ + tool_call_completion("probe", &json!({"owner":owner}).to_string()), + chat_completion("done"), + ] + }) + .collect(), + "done", + ) + .await +} + +#[test] +fn hooks_are_additive_and_isolated_across_agents_turns_and_reused_ids() { + runtime().block_on(async { tokio::spawn(scenario()).await.expect("scenario") }); +} + +async fn scenario() { + let backend = stub_backend().await; + let (provider_a, provider_b) = tokio::join!(worker_provider("a"), worker_provider("b")); + let events: Events = Default::default(); + let global = Observer::new("global", &events, None); + let runtime = Runtime::builder() + .config(offline_config()) + .workspace(Workspace::Ephemeral) + .backend_url(backend.uri()) + .tool_hook(global.clone()) + .post_turn_hook(global) + .build() + .await + .expect("runtime"); + let overlap = Arc::new(Barrier::new(2)); + let hook_a = Observer::new("agent-a", &events, Some(overlap.clone())); + let hook_b = Observer::new("agent-b", &events, Some(overlap)); + let a = runtime + .agent( + spec("worker-a", &provider_a) + .tool_hook(hook_a.clone()) + .post_turn_hook(hook_a), + ) + .unwrap(); + let b = runtime + .agent( + spec("worker-b", &provider_b) + .tool_hook(hook_b.clone()) + .post_turn_hook(hook_b), + ) + .unwrap(); + let turn_hook = Observer::new("turn-a", &events, None); + let scheduled = openhuman_core::agent::host_agents::resolve("worker-a").unwrap(); + scheduled + .scope(async { + assert_eq!( + openhuman_core::agent::hooks::turn_tool_hooks().len(), + 2, + "scheduled turns lost the agent's tool hooks" + ); + assert_eq!( + openhuman_core::agent::hooks::turn_post_turn_hooks().len(), + 2, + "scheduled turns lost the agent's post-turn hooks" + ); + }) + .await; + let (first, second) = tokio::join!( + a.turn("a-first") + .tool_hook(turn_hook.clone()) + .post_turn_hook(turn_hook) + .send(), + b.turn("b-first").send(), + ); + let first = first.expect("first a"); + second.expect("first b"); + a.turn("a-second") + .session(first.session_id) + .send() + .await + .expect("resumed a"); + let mut session = scheduled + .session_host(&scheduled.config, Some("scheduled-a")) + .unwrap(); + scheduled + .scope(async { session.turn_with_origin("a-scheduled", None).await }) + .await + .unwrap(); + runtime.remove_agent("worker-a").await.expect("remove a"); + let replacement = runtime + .agent(spec("worker-a", &provider_a)) + .expect("reuse id"); + replacement + .run("replacement") + .await + .expect("replacement turn"); + eventually("all background post-turn observers", || { + let events = events.lock().unwrap(); + (events + .iter() + .filter(|(_, phase, _)| phase == "turn") + .count() + == 10) + .then_some(()) + }) + .await; + + let events = events.lock().unwrap(); + let count = |label: &str, phase: &str, owner: &str| { + events + .iter() + .filter(|(l, p, o)| l == label && p == phase && o == owner) + .count() + }; + for phase in ["before", "after"] { + assert_eq!(count("global", phase, "a"), 4); + assert_eq!(count("global", phase, "b"), 1); + assert_eq!(count("agent-a", phase, "a"), 3); + assert_eq!(count("agent-b", phase, "b"), 1); + assert_eq!(count("turn-a", phase, "a"), 1); + assert_eq!(count("agent-a", phase, "b"), 0); + assert_eq!(count("agent-b", phase, "a"), 0); + assert_eq!(count("turn-a", phase, "b"), 0); + } + for (label, message) in [ + ("global", "a-first"), + ("global", "b-first"), + ("global", "a-second"), + ("global", "a-scheduled"), + ("global", "replacement"), + ("agent-a", "a-first"), + ("agent-a", "a-second"), + ("agent-a", "a-scheduled"), + ("agent-b", "b-first"), + ("turn-a", "a-first"), + ] { + assert_eq!( + count(label, "turn", message), + 1, + "{label} / {message}: {events:?}" + ); + } + // Agent/turn callbacks never enter the runtime's process-global registry. + assert_eq!(openhuman_core::agent::hooks::embedder_tool_hooks().len(), 1); + assert_eq!( + openhuman_core::agent::hooks::embedder_post_turn_hooks().len(), + 1 + ); +} diff --git a/crates/openhuman-embed/tests/tool_environment.rs b/crates/openhuman-embed/tests/tool_environment.rs new file mode 100644 index 00000000000..9ac96036f50 --- /dev/null +++ b/crates/openhuman-embed/tests/tool_environment.rs @@ -0,0 +1,24 @@ +//! Turn-owned child environments never mutate or inherit the daemon's env. + +#![cfg(unix)] + +use openhuman_core::tools::timeout::{output_unbounded, CommandEnvironment}; + +#[tokio::test] +async fn overlapping_command_environments_replace_inheritance_without_crossing_turns() { + let run = |label: &'static str| async move { + let environment = CommandEnvironment::new([(String::from("TURN_OWNER"), label.into())]); + environment + .scope(async { + tokio::task::yield_now().await; + let mut command = tokio::process::Command::new("/usr/bin/env"); + command.env("COMMAND_OWNED", "runtime-setting"); + String::from_utf8(output_unbounded(&mut command).await.unwrap().stdout).unwrap() + }) + .await + }; + let (a, b) = tokio::join!(run("a"), run("b")); + assert_eq!(a.trim(), "COMMAND_OWNED=runtime-setting\nTURN_OWNER=a"); + assert_eq!(b.trim(), "COMMAND_OWNED=runtime-setting\nTURN_OWNER=b"); + assert!(!CommandEnvironment::is_active()); +} diff --git a/crates/openhuman-embed/tests/tool_hook_context.rs b/crates/openhuman-embed/tests/tool_hook_context.rs new file mode 100644 index 00000000000..4b1ea6af210 --- /dev/null +++ b/crates/openhuman-embed/tests/tool_hook_context.rs @@ -0,0 +1,109 @@ +//! Hook attribution and cwd agree with the builtin tool's actual working directory. + +mod common; + +use common::{ + chat_completion, offline_config, route, runtime, scripted_provider, stub_backend, + tool_call_completion, +}; +use openhuman_embed::seams::{ToolHook, ToolHookContext}; +use openhuman_embed::{Access, AgentDefinitionSpec, AgentSpec, Runtime, ToolScopeSpec, Workspace}; +use serde_json::{json, Value}; +use std::sync::{Arc, Mutex}; + +struct Capture(Arc>>); +#[async_trait::async_trait] +impl ToolHook for Capture { + fn name(&self) -> &str { + "context-probe" + } + async fn before_tool(&self, ctx: &ToolHookContext) -> anyhow::Result<()> { + self.0.lock().unwrap().push(serde_json::to_value(ctx)?); + Ok(()) + } + async fn after_tool(&self, ctx: &ToolHookContext) -> anyhow::Result<()> { + self.0.lock().unwrap().push(serde_json::to_value(ctx)?); + Ok(()) + } +} + +#[test] +fn cwd_and_identity_follow_the_turn_and_then_return_to_the_agent_defaults() { + runtime().block_on(async { tokio::spawn(scenario()).await.unwrap() }); +} + +async fn scenario() { + let backend = stub_backend().await; + let provider = scripted_provider( + (0..2) + .flat_map(|_| { + [ + tool_call_completion("shell", &json!({"command":"pwd; printf '\\nowner=%s home=%s\\n' \"$TURN_OWNER\" \"$HOME\""}).to_string()), + chat_completion("done"), + ] + }) + .collect(), + "done", + ) + .await; + let runtime = Runtime::builder() + .config(offline_config()) + .workspace(Workspace::Ephemeral) + .backend_url(backend.uri()) + .build() + .await + .unwrap(); + let default_dir = tempfile::tempdir().unwrap(); + let turn_dir = tempfile::tempdir().unwrap(); + let observed = Arc::new(Mutex::new(Vec::new())); + let agent = runtime + .agent( + AgentSpec::new("cwd-worker") + .provider(route(&provider, "fixture")) + .access(Access::full()) + .action_dir(default_dir.path()) + .tool_hook(Arc::new(Capture(observed.clone()))) + .definition( + AgentDefinitionSpec::new().tools(ToolScopeSpec::Named(vec!["shell".into()])), + ), + ) + .unwrap(); + agent + .turn("run pwd") + .session("override-session") + .cwd(turn_dir.path()) + .tool_env([("TURN_OWNER".into(), "isolated".into())]) + .send() + .await + .unwrap(); + agent + .turn("run pwd again") + .session("default-session") + .send() + .await + .unwrap(); + let observed = observed.lock().unwrap(); + assert_eq!(observed.len(), 4); + assert!(observed[1]["output"] + .as_str() + .unwrap() + .contains("owner=isolated home=\n")); + for (pair, session, cwd) in [ + (&observed[..2], "override-session", turn_dir.path()), + (&observed[2..], "default-session", default_dir.path()), + ] { + for context in pair { + assert_eq!(context["agent_id"], "cwd-worker", "{context}"); + assert_eq!(context["session_id"], session, "{context}"); + assert_eq!(context["cwd"], cwd.to_string_lossy().as_ref(), "{context}"); + } + assert!( + pair[1]["output"] + .as_str() + .unwrap() + .contains(cwd.to_string_lossy().as_ref()), + "builtin shell disagrees with hook cwd: {}", + pair[1] + ); + } +} diff --git a/crates/openhuman-embed/tests/turn_cancellation.rs b/crates/openhuman-embed/tests/turn_cancellation.rs index 8181a6da4ad..e0ca914b68c 100644 --- a/crates/openhuman-embed/tests/turn_cancellation.rs +++ b/crates/openhuman-embed/tests/turn_cancellation.rs @@ -42,7 +42,12 @@ async fn scenario() { .unwrap(); // Cancellation before send never reaches inference and cannot hang. - let mut turn = agent.turn("cancel before sending"); + let metered = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let observed = metered.clone(); + let mut turn = agent.turn("cancel before sending").meter(move |usage| { + assert!(usage.is_none()); + observed.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + }); let cancel = turn.cancellation_handle(); tokio::time::timeout(Duration::from_secs(2), cancel.cancel()) .await @@ -52,6 +57,11 @@ async fn scenario() { Err(CoreError::TurnCancelled { .. }) )); assert!(common::chat_requests(&provider).await.is_empty()); + assert_eq!( + metered.load(std::sync::atomic::Ordering::SeqCst), + 1, + "cancelled turn is metered once" + ); // A later turn on the same agent remains usable. let mut turn = agent.turn("answer normally"); @@ -77,7 +87,10 @@ async fn scenario() { let blocked = runtime .agent(AgentSpec::new("blocked").provider(route(&slow, "test-model"))) .unwrap(); - let mut turn = blocked.turn("wait for inference"); + let observed = metered.clone(); + let mut turn = blocked.turn("wait for inference").meter(move |_| { + observed.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + }); let cancel = turn.cancellation_handle(); let sent = tokio::spawn(turn.send()); tokio::time::timeout(Duration::from_secs(5), async { @@ -97,6 +110,7 @@ async fn scenario() { sent.await.unwrap(), Err(CoreError::TurnCancelled { .. }) )); + assert_eq!(metered.load(std::sync::atomic::Ordering::SeqCst), 2); assert_eq!(agent.run("still usable").await.unwrap().reply, "finished"); let observed = Arc::new(Outcomes::default()); @@ -118,7 +132,11 @@ async fn scenario() { // An externally dropped send future also acknowledges cancellation. let prior_requests = common::chat_requests(&slow).await.len(); - let mut turn = blocked.turn("drop this inference request"); + let observed = metered.clone(); + let mut turn = blocked.turn("drop this inference request").meter(move |_| { + observed.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + }); + let cancel = turn.cancellation_handle(); let sent = tokio::spawn(turn.send()); tokio::time::timeout(Duration::from_secs(5), async { @@ -134,6 +152,8 @@ async fn scenario() { .await .unwrap(); + assert_eq!(metered.load(std::sync::atomic::Ordering::SeqCst), 3); + let mut invalid = agent .turn("invalid route") .route(openhuman_embed::Route::openai_compatible( diff --git a/crates/openhuman-embed/tests/turn_tools.rs b/crates/openhuman-embed/tests/turn_tools.rs new file mode 100644 index 00000000000..021ad8668fe --- /dev/null +++ b/crates/openhuman-embed/tests/turn_tools.rs @@ -0,0 +1,176 @@ +//! Per-turn host-tool belts replace agent tools without leaking into other turns. + +mod common; + +use common::{ + chat_completion, offline_config, route, runtime, scripted_provider, stub_backend, + tool_call_completion, +}; +use openhuman_embed::{ + AgentDefinitionSpec, AgentSpec, HostTurnTools, Runtime, Tool, ToolScopeSpec, Workspace, +}; +use serde_json::{json, Value}; +use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, +}; + +struct Probe(Arc); +#[async_trait::async_trait] +impl Tool for Probe { + fn name(&self) -> &str { + "probe" + } + fn description(&self) -> &str { + "Perform a host operation" + } + fn parameters_schema(&self) -> Value { + json!({"type":"object","properties":{}}) + } + async fn execute(&self, _: Value) -> anyhow::Result { + self.0.fetch_add(1, Ordering::SeqCst); + Ok(openhuman_core::tools::ToolResult::success("ran")) + } +} +#[test] +fn turn_tool_belts_replace_agent_tools_without_leaking_into_resumed_or_concurrent_turns() { + runtime().block_on(async { tokio::spawn(scenario()).await.unwrap() }); +} +async fn scenario() { + let backend = stub_backend().await; + let provider = scripted_provider( + (0..3) + .flat_map(|_| [tool_call_completion("probe", "{}"), chat_completion("done")]) + .collect(), + "done", + ) + .await; + let other_provider = scripted_provider(vec![tool_call_completion("probe", "{}")], "done").await; + let runtime = Runtime::builder() + .config(offline_config()) + .workspace(Workspace::Ephemeral) + .backend_url(backend.uri()) + .build() + .await + .unwrap(); + let calls = Arc::new(AtomicUsize::new(0)); + let replacement_calls = Arc::new(AtomicUsize::new(0)); + let other_calls = Arc::new(AtomicUsize::new(0)); + let spec = |id, provider: &wiremock::MockServer, calls: Arc| { + AgentSpec::new(id) + .provider(route(provider, "fixture")) + .definition( + AgentDefinitionSpec::new() + .bare_prompt("Use probe then answer") + .tools(ToolScopeSpec::HostOnly), + ) + .tools(move |_| HostTurnTools::advertised(vec![Box::new(Probe(calls.clone()))])) + }; + let agent = runtime + .agent(spec("replace", &provider, calls.clone())) + .unwrap(); + let other = runtime + .agent(spec("other", &other_provider, other_calls.clone())) + .unwrap(); + let replacement = replacement_calls.clone(); + let (first, other_result) = tokio::join!( + agent + .turn("replacement") + .session("same-thread") + .tools(move |_| HostTurnTools::advertised(vec![Box::new(Probe(replacement.clone()))])) + .send(), + other.run("independent") + ); + first.unwrap(); + other_result.unwrap(); + assert_eq!(calls.load(Ordering::SeqCst), 0); + assert_eq!(replacement_calls.load(Ordering::SeqCst), 1); + assert_eq!(other_calls.load(Ordering::SeqCst), 1); + agent + .turn("revoked") + .session("same-thread") + .tools(|_| HostTurnTools::advertised(vec![])) + .send() + .await + .unwrap(); + assert_eq!(calls.load(Ordering::SeqCst), 0); + assert_eq!( + replacement_calls.load(Ordering::SeqCst), + 1, + "revoked tool was still executable" + ); + agent + .turn("restore") + .session("same-thread") + .send() + .await + .unwrap(); + assert_eq!(calls.load(Ordering::SeqCst), 1); + let requests = common::chat_requests(&provider).await; + let revoked = requests + .iter() + .find(|r| { + r.body_json::().unwrap()["messages"] + .as_array() + .unwrap() + .iter() + .any(|m| { + m["content"] + .as_str() + .is_some_and(|s| s.ends_with("revoked")) + }) + }) + .unwrap(); + assert!(!revoked.body_json::().unwrap()["tools"] + .as_array() + .is_some_and(|tools| tools.iter().any(|t| t["function"]["name"] == "probe"))); + // A mixed builtin/host definition must also discard historical declarations + // on an explicit replacement. HostOnly's inherent strict snapshot is not + // enough to exercise the replacement's dispatch scope. + let mixed_provider = scripted_provider( + vec![ + tool_call_completion("probe", "{}"), + chat_completion("done"), + chat_completion("no tools"), + tool_call_completion("probe", "{}"), + chat_completion("done"), + ], + "done", + ) + .await; + let mixed_calls = Arc::new(AtomicUsize::new(0)); + let mixed = runtime + .agent( + spec("mixed", &mixed_provider, mixed_calls.clone()).definition( + AgentDefinitionSpec::new() + .bare_prompt("Use the offered tools") + .tools(ToolScopeSpec::Named(vec!["probe".into()])), + ), + ) + .unwrap(); + mixed + .turn("initial belt") + .session("mixed-thread") + .send() + .await + .unwrap(); + mixed + .turn("replace with empty belt") + .session("mixed-thread") + .tools(|_| HostTurnTools::advertised(vec![])) + .send() + .await + .unwrap(); + assert_eq!(mixed_calls.load(Ordering::SeqCst), 1); + let requests = common::chat_requests(&mixed_provider).await; + assert!(!requests[2].body_json::().unwrap()["tools"] + .as_array() + .is_some_and(|tools| tools.iter().any(|t| t["function"]["name"] == "probe"))); + mixed + .turn("restore belt") + .session("mixed-thread") + .send() + .await + .unwrap(); + assert_eq!(mixed_calls.load(Ordering::SeqCst), 2); +} diff --git a/crates/openhuman-embed/tests/usage_hooks.rs b/crates/openhuman-embed/tests/usage_hooks.rs new file mode 100644 index 00000000000..fadcffa4112 --- /dev/null +++ b/crates/openhuman-embed/tests/usage_hooks.rs @@ -0,0 +1,174 @@ +//! Scoped usage observers and budget policies run between provider calls. + +mod common; + +use common::{ + chat_completion, chat_requests, offline_config, route, runtime, scripted_provider, + stub_backend, tool_call_completion, +}; +use openhuman_embed::seams::{StopDecision, StopHook, TurnState}; +use openhuman_embed::{ + AgentDefinitionSpec, AgentSpec, HostTurnTools, Runtime, Tool, ToolScopeSpec, Workspace, +}; +use serde_json::{json, Value}; +use std::sync::{Arc, Mutex}; + +type Seen = Arc)>>>; +struct Policy { + label: &'static str, + seen: Seen, + stop: bool, +} +#[async_trait::async_trait] +impl StopHook for Policy { + fn name(&self) -> &str { + "usage-budget" + } + async fn check(&self, state: &TurnState<'_>) -> StopDecision { + self.seen.lock().unwrap().push(( + self.label.into(), + state.iteration, + state.cost.input_tokens + state.cost.output_tokens, + state.cost.total_usd(), + )); + if self.stop { + StopDecision::Stop { + reason: "worker budget exhausted".into(), + } + } else { + StopDecision::Continue + } + } +} +struct Probe; +#[async_trait::async_trait] +impl Tool for Probe { + fn name(&self) -> &str { + "probe" + } + fn description(&self) -> &str { + "Read a value" + } + fn parameters_schema(&self) -> Value { + json!({"type":"object","properties":{}}) + } + async fn execute(&self, _: Value) -> anyhow::Result { + Ok(openhuman_core::tools::ToolResult::success("value")) + } +} +fn billed(mut response: Value) -> Value { + response["usage"]["cost"] = json!(0.02); + response +} +#[test] +fn per_agent_budget_stops_before_the_next_model_call_and_turn_observers_do_not_leak() { + runtime().block_on(async { tokio::spawn(scenario()).await.unwrap() }); +} +async fn scenario() { + let backend = stub_backend().await; + let a_provider = scripted_provider( + vec![billed(tool_call_completion("probe", "{}"))], + "extra call", + ) + .await; + let b_provider = scripted_provider( + vec![ + billed(tool_call_completion("probe", "{}")), + billed(chat_completion("done")), + { + let mut response = chat_completion("free"); + response["usage"]["cost"] = json!(0.0); + response + }, + ], + "next turn", + ) + .await; + let runtime = Runtime::builder() + .config(offline_config()) + .workspace(Workspace::Ephemeral) + .backend_url(backend.uri()) + .build() + .await + .unwrap(); + let seen: Seen = Default::default(); + let spec = |id, provider: &wiremock::MockServer| { + AgentSpec::new(id) + .provider(route(provider, "fixture")) + .definition( + AgentDefinitionSpec::new() + .bare_prompt("Use probe then answer") + .tools(ToolScopeSpec::HostOnly), + ) + .tools(|_| HostTurnTools::advertised(vec![Box::new(Probe)])) + }; + let a = runtime + .agent(spec("budget-a", &a_provider).stop_hook(Arc::new(Policy { + label: "a", + seen: seen.clone(), + stop: true, + }))) + .unwrap(); + let b = runtime.agent(spec("budget-b", &b_provider)).unwrap(); + let (a_result, b_result) = tokio::join!( + a.turn("work") + .stop_hook(Arc::new(Policy { + label: "a-turn", + seen: seen.clone(), + stop: false, + })) + .send(), + b.turn("work") + .stop_hook(Arc::new(Policy { + label: "b-turn", + seen: seen.clone(), + stop: false + })) + .send() + ); + assert_eq!(a_result.unwrap().usage.unwrap().cost_usd, Some(0.02)); + let b_result = b_result.unwrap(); + assert_eq!( + chat_requests(&a_provider).await.len(), + 1, + "budget permitted another model call" + ); + assert_eq!(chat_requests(&b_provider).await.len(), 2); + assert_eq!(b_result.usage.unwrap().cost_usd, Some(0.04)); + assert_eq!( + b.run("a later turn").await.unwrap().usage.unwrap().cost_usd, + Some(0.0) + ); + assert_eq!( + b.run("an unpriced turn") + .await + .unwrap() + .usage + .unwrap() + .cost_usd, + None + ); + let seen = seen.lock().unwrap(); + assert_eq!( + seen.iter().filter(|(label, _, _, _)| label == "a").count(), + 1 + ); + assert_eq!( + seen.iter() + .filter(|(label, _, _, _)| label == "a-turn") + .count(), + 1, + "agent budget stop skipped the turn usage observer" + ); + let b_seen: Vec<_> = seen + .iter() + .filter(|(label, _, _, _)| label == "b-turn") + .collect(); + assert_eq!(b_seen.len(), 2, "turn observer leaked to a later turn"); + assert_eq!(b_seen[0].1, 1); + assert_eq!(b_seen[0].2, 2); + assert_eq!(b_seen[0].3, Some(0.02)); + assert_eq!(b_seen[1].1, 2); + assert_eq!(b_seen[1].2, 4); + assert_eq!(b_seen[1].3, Some(0.04)); +} diff --git a/docs/TEST-COVERAGE-MATRIX.md b/docs/TEST-COVERAGE-MATRIX.md index 9f9bb7eb649..96872f4804e 100644 --- a/docs/TEST-COVERAGE-MATRIX.md +++ b/docs/TEST-COVERAGE-MATRIX.md @@ -668,8 +668,16 @@ The thread JSONL store moved to `tinyagents_session::threads` (`vendor/tinyagent | 16.1.17 | Enforced shared turn/run budgets | RU+RI | `crates/openhuman-embed/tests/budget_fanout.rs`, `vendor/tinyagents/vendor/tinyinference/crates/tinyinference-llm/src/model/budget_tests.rs` | ✅ | Atomic parent/child reservations, physical provider admission, conservative unknown spend, cancellation and output cap | | 16.1.18 | Ordered bounded fanout | RI | `crates/openhuman-embed/tests/budget_fanout.rs` | ✅ | Input-order results, per-branch and child error isolation, branch ceilings, shared-budget concurrent refusal and empty fanout | | 16.1.13 | Awaited per-turn cancellation | RI | `crates/openhuman-embed/tests/turn_cancellation.rs`, `crates/openhuman-embed/tests/process_cancellation.rs` | ✅ | Before send, during inference and during a builtin shell command; agent reuse, concurrent cancellation handles, failure/drop acknowledgement, bounded and unbounded descendant termination plus direct-child reaping on Linux | +| 16.1.21 | Agent and turn hook isolation | RI | `crates/openhuman-embed/tests/scoped_hooks.rs` | ✅ | Additive same-named callbacks, overlapping worker turns, one-turn callback lifetime, resumed-session snapshots, core-scheduled turn callbacks and reuse of a removed agent id without registry leakage | +| 16.1.22 | Agent-turn route attribution headers | RU+RI | `crates/openhuman-embed/tests/route_headers.rs`, `crates/openhuman-embed/src/turn_tests.rs` | ✅ | Per-route headers reach only the named endpoint and turn, remain absent on other agents and backend requests, and are redacted in Debug; wire fields pinned against the registered controller schema | +| 16.1.23 | Tool-hook turn identity and cwd | RI | `crates/openhuman-embed/tests/tool_hook_context.rs` | ✅ | Before and after hooks report the worker/session and working root used by builtin shell; a later turn returns to its agent default | +| 16.1.24 | Inline agent/turn permission callback | RI | `crates/openhuman-embed/tests/inline_permissions.rs` | ✅ | UI decisions precede execution; denials narrow authority, concurrent agents remain independent, cancellation closes a pending receiver and a closed UI refuses execution | +| 16.1.25 | Per-turn host-tool replacement | RI | `crates/openhuman-embed/tests/turn_tools.rs` | ✅ | Replacement and empty-belt revocation apply to resumed HostOnly and mixed definitions, without restoring historical declarations or affecting other agents and later turns | +| 16.1.26 | Scoped usage observation and budget stop | RI | `crates/openhuman-embed/tests/usage_hooks.rs` | ✅ | Cumulative charged cost and tokens, known-zero versus unpriced calls, per-agent policy and per-turn observer isolation; no summary-repair model calls after a stop | +| 16.1.27 | Scoped builtin command environments | RI | `crates/openhuman-embed/tests/tool_environment.rs`, `crates/openhuman-embed/tests/tool_hook_context.rs` | ✅ | Overlapping subprocess maps exclude host inheritance on Unix; a real shell receives its turn's variables while excluding host HOME and preserving command-owned runtime/security additions | | 16.1.20 | Host turn observer privacy and runtime propagation | RI | `crates/openhuman-embed/tests/turn_observers.rs`, `crates/openhuman-embed/tests/observed_turns.rs` | ✅ | Sanitized terminal failures, opt-in payloads, root runtime scope propagation, real model/tool callbacks and answering-model usage | + ### 16.2 Dynamic runtime and agent configuration | ID | Feature | Layer | Test path(s) | Status | Notes | diff --git a/docs/benchmarks/medulla-embed-linux.json b/docs/benchmarks/medulla-embed-linux.json new file mode 100644 index 00000000000..d7f84503fa7 --- /dev/null +++ b/docs/benchmarks/medulla-embed-linux.json @@ -0,0 +1,238 @@ +{ + "source_commit": "bbd3d83d80fcaaf3b6284c5b7881d080dcd53a34", + "date": "2026-10-10", + "host": "Intel Core i7-14700F; x86_64 Ubuntu Linux 7.0.0-31-generic", + "request_recording": false, + "runs": [ + { + "after_turns_rss_kib": 227952, + "agents": 50, + "before_boot_rss_kib": 14488, + "boot_ms": 256.195196, + "cold_turn_ms": 16.275062, + "fleet_elapsed_ms": 447.54206, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 3.78296875, + "registered_rss_kib": 36464, + "runtime_rss_kib": 34264, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 330.948465, + "turn_p95_ms": 398.109413, + "turn_p99_ms": 402.916892, + "worker_threads": 2, + "cgroup_memory_peak": "203931648", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 226972, + "agents": 50, + "before_boot_rss_kib": 14396, + "boot_ms": 256.484423, + "cold_turn_ms": 18.921321000000002, + "fleet_elapsed_ms": 420.257983, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 3.75734375, + "registered_rss_kib": 36856, + "runtime_rss_kib": 34596, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 326.425821, + "turn_p95_ms": 359.84249600000004, + "turn_p99_ms": 396.121414, + "worker_threads": 2, + "cgroup_memory_peak": "202309632", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 222156, + "agents": 50, + "before_boot_rss_kib": 14532, + "boot_ms": 258.358991, + "cold_turn_ms": 18.654450999999998, + "fleet_elapsed_ms": 410.89376500000003, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 3.665390625, + "registered_rss_kib": 36640, + "runtime_rss_kib": 34488, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 273.595824, + "turn_p95_ms": 326.54888900000003, + "turn_p99_ms": 383.942, + "worker_threads": 2, + "cgroup_memory_peak": "198098944", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 380152, + "agents": 100, + "before_boot_rss_kib": 14696, + "boot_ms": 261.058476, + "cold_turn_ms": 23.715644, + "fleet_elapsed_ms": 843.2798409999999, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 3.37453125, + "registered_rss_kib": 38836, + "runtime_rss_kib": 34600, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 630.060869, + "turn_p95_ms": 781.8404089999999, + "turn_p99_ms": 807.385043, + "worker_threads": 2, + "cgroup_memory_peak": "368861184", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 381448, + "agents": 100, + "before_boot_rss_kib": 14532, + "boot_ms": 260.519274, + "cold_turn_ms": 22.652142, + "fleet_elapsed_ms": 1063.297785, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 3.389453125, + "registered_rss_kib": 38588, + "runtime_rss_kib": 34368, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 780.051866, + "turn_p95_ms": 929.080972, + "turn_p99_ms": 964.185504, + "worker_threads": 2, + "cgroup_memory_peak": "370405376", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 381500, + "agents": 100, + "before_boot_rss_kib": 14408, + "boot_ms": 258.39308300000005, + "cold_turn_ms": 19.710221999999998, + "fleet_elapsed_ms": 953.5413149999999, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 3.3948828125, + "registered_rss_kib": 38116, + "runtime_rss_kib": 33864, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 713.8235229999999, + "turn_p95_ms": 839.781121, + "turn_p99_ms": 865.269359, + "worker_threads": 2, + "cgroup_memory_peak": "370077696", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 1511184, + "agents": 500, + "before_boot_rss_kib": 14236, + "boot_ms": 258.86011599999995, + "cold_turn_ms": 22.475146000000002, + "fleet_elapsed_ms": 5451.315839, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 2.88403125, + "registered_rss_kib": 55136, + "runtime_rss_kib": 34560, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 4820.459523, + "turn_p95_ms": 5284.533808, + "turn_p99_ms": 5339.48242, + "worker_threads": 2, + "cgroup_memory_peak": "1625919488", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 1540228, + "agents": 500, + "before_boot_rss_kib": 14452, + "boot_ms": 258.009249, + "cold_turn_ms": 14.884082, + "fleet_elapsed_ms": 4672.801607, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 2.941359375, + "registered_rss_kib": 54784, + "runtime_rss_kib": 34252, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 4113.175878, + "turn_p95_ms": 4447.625557, + "turn_p99_ms": 4513.371248, + "worker_threads": 2, + "cgroup_memory_peak": "1648812032", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + }, + { + "after_turns_rss_kib": 1542084, + "agents": 500, + "before_boot_rss_kib": 14336, + "boot_ms": 257.69439500000004, + "cold_turn_ms": 16.925666, + "fleet_elapsed_ms": 5055.745514, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 2.94503125, + "registered_rss_kib": 54832, + "runtime_rss_kib": 34228, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 4430.5079940000005, + "turn_p95_ms": 4828.56837, + "turn_p99_ms": 4880.706669, + "worker_threads": 2, + "cgroup_memory_peak": "1650745344", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + } + ], + "no_swap_500_agents": { + "after_turns_rss_kib": 1537772, + "agents": 500, + "before_boot_rss_kib": 14476, + "boot_ms": 259.344405, + "cold_turn_ms": 17.940296, + "fleet_elapsed_ms": 5070.406826, + "inference": "loopback-http-mock", + "marginal_after_turns_mib": 2.9364765625, + "registered_rss_kib": 54900, + "runtime_rss_kib": 34296, + "tool_scope": [ + "shell" + ], + "turn_p50_ms": 4353.5878170000005, + "turn_p95_ms": 4753.196036, + "turn_p99_ms": 4844.317583000001, + "worker_threads": 2, + "cgroup_memory_peak": "1647300608", + "cgroup_memory_max": "2147483648", + "cgroup_cpu_max": "200000 100000", + "cgroup_memory_events": "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\noom_group_kill 0\nsock_throttled 0" + } +} diff --git a/gitbooks/developing/performance.md b/gitbooks/developing/performance.md index 37f5a3d825b..557e2bf9951 100644 --- a/gitbooks/developing/performance.md +++ b/gitbooks/developing/performance.md @@ -17,9 +17,9 @@ Everything below was measured on Apple Silicon macOS with a `--release` build an | N agents | Marginal KiB/agent | Settled MiB | Idle CPU ms/10s | Threads | FDs | | ---: | ---: | ---: | ---: | ---: | ---: | -| 50 | 1,985 | 223 | 3 | 71 | 420 | -| 100 | 1,866 | 356 | 3 | 123 | 820 | -| 500 | 1,770 | 1,393 | 3 | 211 | 3,220 | +| 50 | 221.65 MiB | 3.757 MiB | 192.94 MiB | 420 ms | 360 ms | +| 100 | 372.51 MiB | 3.389 MiB | 352.93 MiB | 954 ms | 840 ms | +| 500 | 1,504.13 MiB | 2.941 MiB | 1,572.43 MiB | 5,056 ms | 4,829 ms | 500 agents fit in one process at roughly 1.77 MiB marginal cost each. Idle CPU stays flat as N grows, which matters because an agent that is not mid-turn should not spend cycles. Thread count grows by about 0.35 per agent. Watch that line before pushing past 500 in production. @@ -100,3 +100,62 @@ scripts/kernel-floor.sh flows These numbers were gathered on macOS. It has no local cgroup memory limit and no `/proc//smaps_rollup`, so there is no true PSS (proportional shared memory) reading. RSS overcounts shared pages, and the error grows with agent count. The fleet and instances numbers project from measured marginal cost. They are not a live test of surviving an OOM kill at N agents on a 2 GB / 2 vCPU Linux box. Every scenario also replaces network inference with a mock provider at fixed latency, so turn timings measure orchestration overhead, not real model latency. Validation under Linux cgroups has not been done yet. For token cost instead of process footprint, see [Smart token compression](../features/token-compression.md). It is the other half of "cheap": it controls how much of what the harness assembles reaches the model. + +## Linux runtime-owned agents (Medulla integration) + +Measured on 2026-10-10 at OpenHuman commit `e2224a22a3`, using the release +`openhuman-embed` example `linux_fleet` with default features disabled. This +exercises one `Runtime` and N `AgentSpec`s, with two Tokio workers, ephemeral +session storage, and a loopback HTTP chat-completions mock. Each agent +advertises the builtin `shell` tool; the mock returns text without executing it. +Conversation memory, local model runtimes, and session dual writes are off. +The host is x86_64 Ubuntu, Linux 7.0.0-31-generic, on an Intel Core i7-14700F. +Other builds were running on the host, so these are shared-host measurements. + +Each run starts in a fresh cgroup with `memory.max = 2147483648` and +`cpu.max = 200000 100000` (2 GiB and a two-CPU quota). Three fresh processes +were measured at each size; the table gives medians. RSS is sampled after every +agent completes one concurrent turn, with all agent handles retained. Marginal +RSS is `(RSS after turns - RSS after Runtime::build) / N`, without allocator +trimming. The cgroup peak includes the HTTP mock and the Python measurement +wrapper. Mock request recording is disabled, so retained HTTP request history +is excluded. RSS and cgroup accounting differ because shared file pages may be +charged outside the new cgroup. + +| Concurrent agents | Process RSS after turns | Marginal RSS/agent | Cgroup peak | Fleet wall time | Turn p95 | +| --- | --- | --- | --- | --- | --- | +| 50 | 221.65 MiB | 3.757 MiB | 192.94 MiB | 420 ms | 360 ms | +| 100 | 372.51 MiB | 3.389 MiB | 352.93 MiB | 954 ms | 840 ms | +| 500 | 1,504.13 MiB | 2.941 MiB | 1,572.43 MiB | 5,056 ms | 4,829 ms | + +The 100- and 500-agent results are within twice the earlier macOS/mock +1.77 MiB figure (3.54 MiB). The 50-agent result misses that target slightly: +its three runs used 3.665–3.783 MiB per agent. These are different hosts and +harness entry points, so the table is a capacity measurement rather than a +controlled comparison of the platforms. + +Runtime boot took 256–262 ms. The first turn in each fresh process took +14.9–23.8 ms, before launching the concurrent fleet. Both are within twice the +earlier 476 ms bootstrap and 102 ms first-turn figures. “First turn” includes +session/model/tool initialization, but excludes compiling and loading the +executable; the OS page cache was warm. + +A separate 500-agent run with **swap disabled** (`MemorySwapMax=0`) completed +in 5,070 ms, with 2.936 MiB marginal RSS per agent and a 1,570.99 MiB cgroup +peak. Every run recorded zero `max`, `oom`, and `oom_kill` memory events. +This does not establish capacity for real providers, tool subprocesses, MCP +servers, or 1,000 simultaneously active turns. Measure those workloads before +sizing a production fleet. + +To reproduce from the OpenHuman repository root: + +```sh +cargo build -p openhuman-embed --locked --release --no-default-features --example linux_fleet +systemd-run --user --scope -p MemoryMax=2G -p MemorySwapMax=0 -p CPUQuota=200% --quiet \ + python3 crates/openhuman-embed/examples/linux_fleet_cgroup.py 500 +``` + +Use `50` or `100` for the other sizes and repeat each command in a fresh scope. +The JSON output records process RSS, bootstrap/turn timings, cgroup limits, +peak memory, and OOM counters. The raw runs are checked in as +[`docs/benchmarks/medulla-embed-linux.json`](../../docs/benchmarks/medulla-embed-linux.json). diff --git a/scripts/ci/check-openhuman-rust-layout.mjs b/scripts/ci/check-openhuman-rust-layout.mjs index 745ad99fcf9..f06cfcb5e22 100644 --- a/scripts/ci/check-openhuman-rust-layout.mjs +++ b/scripts/ci/check-openhuman-rust-layout.mjs @@ -40,7 +40,7 @@ const LEGACY_LIMIT_ENTRIES = [ // under the general 750 limit, so it needs no exception at all. // The session-todo integration added transcript metadata construction to // this already-exempt composition seam. Keep its allowance exact. - ["crates/openhuman-core/src/agent/session_host/runtime_session.rs", 1349], + ["crates/openhuman-core/src/agent/session_host/runtime_session.rs", 1293], // Session-host factory still assembles the product's deliberately coupled // provider, security, memory, tool and prompt policy. Generic session // state moved to tinyagents-runtime; this remaining composition is split in @@ -161,6 +161,7 @@ for (const [directory, table] of [ ["examples", "example"], ]) { const files = rootRustTargetNames(directory); + const declared = declaredTargets(table); for (const name of files) { if (!declared.has(name)) diff --git a/vendor/tinybox b/vendor/tinybox index 65c5f7a473a..b61c6bcebf3 160000 --- a/vendor/tinybox +++ b/vendor/tinybox @@ -1 +1 @@ -Subproject commit 65c5f7a473a467e1e377454431791400945a2032 +Subproject commit b61c6bcebf3c0dbaf6ce3e36fae8131ebd6d78ae