Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
0f430ae
fix(session): hydrate integrations on every session start, not only c…
senamakel Sep 23, 2026
ad05dab
fix(test): add missing `tools` field to test transcript constructors
senamakel Sep 23, 2026
fd85718
chore(deps): add tinyagents vendor dependency
senamakel Sep 23, 2026
94fbb14
chore(deps): update tinyagents submodule
senamakel Sep 23, 2026
39283e2
chore(deps): add tinyagents vendor dependency
senamakel Sep 23, 2026
62e3306
chore(deps): update tinyagents vendor dependency
senamakel Sep 23, 2026
04c9b2d
chore: add recorded tools module
senamakel Sep 23, 2026
34ec573
chore: add recorded tools tests
senamakel Sep 23, 2026
5aab52f
fix(agent): handle session host runtime errors gracefully
senamakel Sep 23, 2026
83ee6c0
chore(openhuman-core): remove unused prefix recovery helper
senamakel Sep 23, 2026
bc2a4f8
fix(runtime_session): restore session state after host restart
senamakel Sep 23, 2026
6752e42
fix(scope): reduce visibility of internal methods and add test module
senamakel Sep 23, 2026
9a636ae
fix(runtime_session): fix formatting of rehydrate_integration_actions…
senamakel Sep 23, 2026
51b5cd9
chore(openhuman-core): update session host runtime imports
senamakel Sep 23, 2026
4f0cb7b
fix(session_host): correct module path for recorded_integration_actions
senamakel Sep 23, 2026
de9108d
refactor(session-host): move delegation tool surface doc to its imple…
senamakel Sep 23, 2026
2ae204b
chore: add blank line and Arc import in session host
senamakel Sep 23, 2026
9f68947
docs(AGENTS.md): document tool list persistence in session transcripts
senamakel Sep 24, 2026
20030cf
chore(AGENTS.md): remove duplicate bullet point about tool list recor…
senamakel Sep 24, 2026
6f8bacf
chore: resolve tinyagents submodule merge
senamakel Sep 24, 2026
dcc4582
feat(session): filter rehydrated actions by current integration policy
senamakel Sep 24, 2026
2db8081
feat(session): track whether connected integrations are authoritative
senamakel Sep 24, 2026
527e21c
fix(recorded_tools): fix formatting of chained condition in rehydrate…
senamakel Sep 24, 2026
6664e79
test(recorded_tools): add assertion that rebuilt declaration matches …
senamakel Sep 24, 2026
bca79b5
Merge remote-tracking branch 'upstream/main' into pr/6589
senamakel Sep 24, 2026
8f4963d
fix(session): prune stale integration and MCP announcements on refresh
senamakel Sep 24, 2026
dc5b75b
fix(test): ensure rehydration test sets connected integrations
senamakel Sep 24, 2026
02f10bf
fix(agent): reformat integration prelude and test code
senamakel Sep 24, 2026
8ad1c03
chore(deps): update tinyagents submodule
senamakel Sep 24, 2026
094e35e
chore(deps): update tinyagents submodule
senamakel Sep 24, 2026
ef23298
chore(deps): update tinyagents submodule
senamakel Sep 24, 2026
f0105e8
fix(session_host): skip non-integration tools during rehydration
senamakel Sep 24, 2026
d2349f1
test(recorded-tools): format test vector for readability
senamakel Sep 24, 2026
0f4473b
fix(runtime_session): restore recorded tools on session resume
senamakel Sep 24, 2026
1b9788f
fix(builder): remove authoritative flag from session host builder
senamakel Sep 24, 2026
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
7 changes: 7 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -389,6 +389,13 @@ seed history by hand, or pick a transcript by recency.
messages verbatim, which is what keeps the provider's prefix cache warm
across a restart. The accepted consequence is that prompt edits, new skills
and newly connected integrations do not reach an existing thread.
- **So is the tool list it was sent.** Every turn records its tool
declarations in the transcript (a `{"kind":"tools"}` record, written only
when they change); resume restores them, the prefix, and the committed-turn
count. The host never shrinks a thread's tools because a cache went cold:
`session_host/recorded_tools.rs` rebuilds recorded Composio actions as
Comment thread
senamakel marked this conversation as resolved.
deferred executors, and the prelude fetches integrations on the first turn
of every session instance, not only on a brand-new thread.
- **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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,7 @@ fn fake_meta(thread_id: Option<&str>) -> TranscriptMeta {
async fn ingest_extracts_high_importance_preference_with_provenance() {
let mem = InMemory::new();
let transcript = SessionTranscript {
tools: None,
meta: fake_meta(Some("thr_alpha")),
messages: durable_messages([
ChatMessage::user("hi"),
Expand Down Expand Up @@ -236,6 +237,7 @@ async fn ingest_extracts_high_importance_preference_with_provenance() {
async fn re_ingest_is_idempotent() {
let mem = InMemory::new();
let transcript = SessionTranscript {
tools: None,
meta: fake_meta(Some("thr_beta")),
messages: durable_messages([ChatMessage::user(
"I prefer Postgres for everything new — please default to it.",
Expand All @@ -260,6 +262,7 @@ async fn re_ingest_is_idempotent() {
async fn ingest_captures_user_reflection_and_recurring_pattern() {
let mem = InMemory::new();
let transcript = SessionTranscript {
tools: None,
meta: fake_meta(Some("thr_gamma")),
messages: durable_messages([
ChatMessage::user("I prefer terse responses with no preamble."),
Expand Down Expand Up @@ -296,6 +299,7 @@ async fn ingest_captures_user_reflection_and_recurring_pattern() {
async fn ingest_filters_low_signal_chatter() {
let mem = InMemory::new();
let transcript = SessionTranscript {
tools: None,
meta: fake_meta(None),
messages: durable_messages([
ChatMessage::user("ok"),
Expand Down Expand Up @@ -323,6 +327,7 @@ async fn ingest_persists_candidates_with_bounded_concurrency() {
// PERSIST_CONCURRENCY (8), so an unbounded fan-out would push more than 8
// stores in flight at once — the bound assertion below would then fail.
let transcript = SessionTranscript {
tools: None,
meta: fake_meta(Some("thr_bound")),
messages: durable_messages([
ChatMessage::user("I prefer Postgres over MySQL for new metadata services."),
Expand Down
1 change: 1 addition & 0 deletions crates/openhuman-core/src/agent/session_host/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ mod factory;
mod hooks;
mod policy;
mod prefix_snapshot;
mod recorded_tools;
mod runtime;
mod runtime_session;
mod session_api;
Expand Down
18 changes: 4 additions & 14 deletions crates/openhuman-core/src/agent/session_host/prefix_snapshot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,10 @@
//!
//! The prompt is rendered once per session as cache tiers (see
//! `agent::prompts::TieredPrompt::system_messages`) and sent as one leading
//! system message per tier. These two helpers turn a rendered prompt into
//! that prefix and recover it from a resumed transcript, so `runtime_session`
//! never has to know how many messages a prefix is.
//! system message per tier. This helper turns a rendered prompt into that
//! prefix, so `runtime_session` never has to know how many messages a prefix
//! is. A resumed thread's prefix is restored by the tinyagents session from
//! the transcript's leading system rows; nothing here re-derives it.

use tinyagents_runtime::PrefixSnapshot;
use tinyinference_llm::message::Message;
Expand All @@ -24,14 +25,3 @@ pub(super) fn tiered_prefix_snapshot(tiered: &TieredPrompt) -> PrefixSnapshot {
);
PrefixSnapshot::new(messages.into_iter().map(Message::system).collect())
}

/// The frozen prefix of a resumed transcript: every leading system message,
/// not only the first, because the prompt is sent as one message per tier.
pub(super) fn leading_system_prefix(history: &[Message]) -> Option<PrefixSnapshot> {
let leading: Vec<Message> = history
.iter()
.take_while(|message| matches!(message, Message::System(_)))
.cloned()
.collect();
(!leading.is_empty()).then(|| PrefixSnapshot::new(leading))
}
223 changes: 223 additions & 0 deletions crates/openhuman-core/src/agent/session_host/prelude_integrations.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,223 @@
//! Turn-boundary refresh of the prelude's integration state: hydrating the
//! connected integrations a session starts without, tracking connects and
//! revokes between turns, and adopting the integration actions the tinyagents
//! session restored for a resumed thread.

use std::sync::Arc;

use tinyagents_runtime::ToolSnapshot;

use super::OpenHumanTurnPrelude;

impl OpenHumanTurnPrelude {
/// Takes the declarations the tinyagents session restored for this
/// thread. Called before the boundary refresh so the rebuilt surface can
/// include them.
pub(super) fn adopt_recorded_tools(&self, recorded: Option<&ToolSnapshot>) {
let Some(recorded) = recorded else {
return;
};
let actions = super::super::recorded_tools::recorded_integration_actions(recorded.specs());
Comment thread
senamakel marked this conversation as resolved.
log::debug!(
"[session] adopting {} recorded integration action declaration(s) agent={}",
actions.len(),
self.agent_definition_id
);
self.mutable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.recorded_integration_actions = actions;
}

pub(super) async fn refresh_turn_boundary(&self, cold: bool) {
// Hydrate on the first turn of *this session instance*, not only on a
// brand-new thread. A resumed thread is never `cold`, and a session
// rebuilt after a restart (or any rebuild past the 60 s integrations
// cache TTL) is seeded from an empty cache — gating the fetch on
// `cold` left it with zero integrations, no deferred Composio
// actions, and no `tool_search` bridge for the whole thread.
// `refresh_cold_integrations` is a no-op once hydrated.
self.refresh_cold_integrations().await;
if !cold {
self.refresh_dynamic_announcements().await;
}
// Integration changes are authority changes, not only display
// announcements. Refresh the delegation executable set and rebuild
// its schema/policy in the same hook pass before the driver sees it.
self.refresh_delegation_tool_surface();
}

pub(super) async fn refresh_cold_integrations(&self) {
let should_fetch = !self
.mutable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.connected_integrations_initialized;
if !should_fetch {
return;
}
let config = match self.runtime_config.clone() {
Some(config) => Some(config),
None => crate::config::Config::load_or_init()
.await
.ok()
.map(Arc::new),
};
let Some(config) = config else {
return;
};
let Some((connected, authoritative)) = load_connected_integrations(&config).await else {
// Backend unreachable and nothing cached: stay un-hydrated so the
// next turn retries rather than pinning an empty surface.
log::warn!(
"[session] integrations unavailable and no cached snapshot; will retry next turn agent={}",
self.agent_definition_id
);
return;
};
log::info!(
"[session] hydrated connected integrations count={} agent={}",
connected.len(),
self.agent_definition_id
);
let mcp_servers = crate::mcp::registry::connections::connected_overview()
.await
.into_iter()
.map(|server| server.qualified_name)
.collect::<std::collections::HashSet<_>>();
let mut mutable = self
.mutable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
mutable.connected_integrations = connected;
// A stale fallback is useful for announcements but cannot authorize
// restored executors. Leave hydration pending so a later turn retries
// the live lookup rather than pinning this session to the snapshot.
mutable.connected_integrations_initialized = authoritative;
mutable.connected_integrations_authoritative = authoritative;
Comment thread
senamakel marked this conversation as resolved.
mutable.announced_integrations = mutable
.connected_integrations
.iter()
.map(|item| item.toolkit.clone())
.collect();
mutable.announced_mcp_servers = mcp_servers;
}

pub(super) async fn refresh_dynamic_announcements(&self) {
let skills_changed = self.drain_host_events();
let config = match self.runtime_config.clone() {
Some(config) => Some(config),
None => crate::config::Config::load_or_init()
.await
.ok()
.map(Arc::new),
};
if let Some(config) = config.as_deref() {
// An expired cache is refetched rather than skipped, so a
// long-lived session keeps tracking connects/revokes.
let current = match crate::integrations::composio::cached_active_integrations(config) {
Comment thread
senamakel marked this conversation as resolved.
Some(current) => Some((current, true)),
None => load_connected_integrations(config).await,
};
if let Some((current, authoritative)) = current {
let mut mutable = self
.mutable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let current_slugs: std::collections::HashSet<_> =
current.iter().map(|item| item.toolkit.clone()).collect();
mutable
.announced_integrations
.retain(|slug| current_slugs.contains(slug));
mutable
.pending_integration_announcement
.retain(|slug| current_slugs.contains(slug));
for slug in &current_slugs {
if mutable.announced_integrations.insert(slug.clone())
&& !mutable.pending_integration_announcement.contains(slug)
{
mutable.pending_integration_announcement.push(slug.clone());
}
}
mutable.connected_integrations = current;
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.
mutable.connected_integrations_authoritative = authoritative;
}
}
let connected_mcp = crate::mcp::registry::connections::connected_overview()
.await
.into_iter()
.map(|server| server.qualified_name)
.collect::<Vec<_>>();
let mut mutable = self
.mutable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let connected_mcp: std::collections::HashSet<_> = connected_mcp.into_iter().collect();
mutable
.announced_mcp_servers
.retain(|server| connected_mcp.contains(server));
mutable
.pending_mcp_announcement
.retain(|server| connected_mcp.contains(server));
for server in connected_mcp {
if mutable.announced_mcp_servers.insert(server.clone())
&& !mutable.pending_mcp_announcement.contains(&server)
{
mutable.pending_mcp_announcement.push(server);
}
}
if !skills_changed {
return;
}
// Event-driven metadata refresh keeps the steady-state hot path free
// of the old per-turn filesystem scan.
let latest = crate::skills::load_workflow_metadata(&self.workspace_dir);
let id = |workflow: &crate::skills::Workflow| {
if workflow.dir_name.is_empty() {
workflow.name.clone()
} else {
workflow.dir_name.clone()
}
};
let previous: std::collections::HashSet<_> = mutable.workflows.iter().map(&id).collect();
let current: std::collections::HashSet<_> = latest.iter().map(&id).collect();
for id in current.difference(&previous) {
if mutable.announced_skills.insert((*id).clone())
&& !mutable.pending_skill_announcement.contains(id)
{
mutable.pending_skill_announcement.push((*id).clone());
}
}
for id in previous.difference(&current) {
mutable.announced_skills.remove(id);
mutable
.pending_skill_announcement
.retain(|pending| pending != id);
if !mutable.pending_skill_retraction.contains(id) {
mutable.pending_skill_retraction.push((*id).clone());
}
}
mutable.workflows = latest;
}
}

/// Live connected integrations, falling back to the last cached snapshot
/// (even past its TTL) when the backend is unreachable. `None` only when
/// there is neither a live answer nor any snapshot to fall back to.
async fn load_connected_integrations(
config: &crate::config::Config,
) -> Option<(Vec<crate::agent::prompts::ConnectedIntegration>, bool)> {
use crate::integrations::composio::FetchConnectedIntegrationsStatus;
match crate::integrations::composio::fetch_connected_integrations_status(config).await {
FetchConnectedIntegrationsStatus::Authoritative(connected) => Some((connected, true)),
FetchConnectedIntegrationsStatus::Unavailable => {
let stale =
Comment thread
senamakel marked this conversation as resolved.
crate::integrations::composio::cached_active_integrations_including_expired(config);
log::warn!(
Comment thread
senamakel marked this conversation as resolved.
"[session] integrations fetch unavailable; using stale snapshot={}",
stale.as_ref().map_or(0, Vec::len)
);
stale.map(|connected| (connected, false))
Comment thread
senamakel marked this conversation as resolved.
}
}
}
Loading
Loading