Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
7 changes: 7 additions & 0 deletions crates/openhuman-core/src/agent/hooks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<String>,
/// Working root of this tool dispatch, including a per-turn cwd override.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cwd: Option<std::path::PathBuf>,
}

/// What a pre-tool hook decided about a call.
Expand Down
67 changes: 67 additions & 0 deletions crates/openhuman-core/src/agent/hooks_scope.rs
Original file line number Diff line number Diff line change
@@ -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<Arc<dyn ToolHook>>,
post_turn: Vec<Arc<dyn PostTurnHook>>,
stop: Vec<Arc<dyn crate::agent::stop_hooks::StopHook>>,
}

impl HookScope {
/// Append a tool callback, preserving registration order.
pub fn push_tool(&mut self, hook: Arc<dyn ToolHook>) {
self.tools.push(hook);
}

/// Append a completed-turn callback, preserving registration order.
pub fn push_post_turn(&mut self, hook: Arc<dyn PostTurnHook>) {
self.post_turn.push(hook);
}

/// Append a usage observer/budget policy checked after each model call.
pub fn push_stop(&mut self, hook: Arc<dyn crate::agent::stop_hooks::StopHook>) {
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<T>(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<F: Future>(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<Arc<dyn ToolHook>> {
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<Arc<dyn PostTurnHook>> {
let mut hooks = super::embedder_post_turn_hooks();
let _ = ACTIVE.try_with(|scope| hooks.extend(scope.post_turn.iter().cloned()));
hooks
}
25 changes: 16 additions & 9 deletions crates/openhuman-core/src/agent/host_agents.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ pub struct HostAgent {
pub host_tools: Option<HostTools>,
/// The context every read during the agent's turn must see.
pub context: Arc<CoreContext>,
/// Agent-owned callbacks inherited by core-scheduled turns.
pub hooks: crate::agent::hooks::HookScope,
}

impl HostAgent {
Expand All @@ -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<F: std::future::Future>(&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
}
}

Expand Down
3 changes: 3 additions & 0 deletions crates/openhuman-core/src/agent/host_agents_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
})
}
Expand Down Expand Up @@ -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:?}");
Expand All @@ -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
Expand Down
1 change: 1 addition & 0 deletions crates/openhuman-core/src/agent/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -393,7 +393,7 @@ impl OpenHumanSessionHost {
None => SystemPromptBuilder::with_defaults(),
};
let post_turn_hooks: Vec<Arc<dyn crate::agent::hooks::PostTurnHook>> =
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
Expand Down
64 changes: 4 additions & 60 deletions crates/openhuman-core/src/agent/session_host/runtime_session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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::{
Expand Down Expand Up @@ -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());
}
Expand Down Expand Up @@ -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<dyn tinyagents_session::transcript::TranscriptLocator> {
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"]
Expand Down
Original file line number Diff line number Diff line change
@@ -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<dyn tinyagents_session::transcript::TranscriptLocator> {
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()),
}
}
}
Loading
Loading