From 5c18808659f6f82064fa5382c991807249eb2844 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:39:29 +0300 Subject: [PATCH 01/26] refactor(runtime): register derived contexts with the agent registry Contexts produced by the overlay path are now wrapped through the agent registry instead of being constructed directly, so derived contexts are tracked alongside their originals. A new agent scope module is wired in to support this. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/core/runtime/context.rs | 6 ++-- .../src/core/runtime/context_agent.rs | 30 +++++++++++++++++++ 2 files changed, 34 insertions(+), 2 deletions(-) create mode 100644 crates/openhuman-core/src/core/runtime/context_agent.rs diff --git a/crates/openhuman-core/src/core/runtime/context.rs b/crates/openhuman-core/src/core/runtime/context.rs index ddf07b16129..d3956881709 100644 --- a/crates/openhuman-core/src/core/runtime/context.rs +++ b/crates/openhuman-core/src/core/runtime/context.rs @@ -428,7 +428,7 @@ impl CoreContext { })), } }; - Arc::new(CoreContext { + crate::storage::agents::registered(Arc::new(CoreContext { host_kind: self.host_kind, workspace_binding: RwLock::new(shared_binding), domains, @@ -438,7 +438,7 @@ impl CoreContext { backend_transport: self.backend_transport.clone(), turn_origin: self.turn_origin.clone(), session_agent: overlay.session_agent.or_else(|| self.session_agent.clone()), - }) + })) } /// The agent a host session store scopes work under this context to, if @@ -740,6 +740,8 @@ pub async fn init_stores(cfg: &crate::config::Config, domains: crate::core::runt #[path = "context_turn_origin.rs"] mod turn_origin_scope; +#[path = "context_agent.rs"] +mod agent_scope; #[cfg(test)] #[path = "context_tests.rs"] diff --git a/crates/openhuman-core/src/core/runtime/context_agent.rs b/crates/openhuman-core/src/core/runtime/context_agent.rs new file mode 100644 index 00000000000..90da0a62884 --- /dev/null +++ b/crates/openhuman-core/src/core/runtime/context_agent.rs @@ -0,0 +1,30 @@ +use super::*; + +impl CoreContext { + /// This context, acting for `agent`: the same configuration, workspace + /// binding, domains and transport, with the session agent replaced. + /// + /// For background work that visits an agent's records when no live + /// context of that agent is at hand (`crate::storage::agents`): the agent + /// was registered by an earlier process, or its host dropped it. A live + /// agent context is always preferred, since it carries the agent's own + /// configuration. + pub fn for_agent(self: &Arc, agent: &str) -> Arc { + let binding = self + .workspace_binding + .read() + .map(|handle| Arc::clone(&*handle)) + .unwrap_or_else(|poisoned| Arc::clone(&*poisoned.into_inner())); + Arc::new(CoreContext { + host_kind: self.host_kind, + workspace_binding: RwLock::new(binding), + domains: self.domains, + tool_groups: self.tool_groups.clone(), + embedder_config: self.embedder_config.clone(), + user_skill_roots: self.user_skill_roots, + backend_transport: self.backend_transport.clone(), + turn_origin: self.turn_origin.clone(), + session_agent: Some(agent.to_string()), + }) + } +} From 06b8343a0e7be3d05ffb6e818a53d2cbe9fd45e7 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:39:57 +0300 Subject: [PATCH 02/26] refactor(storage): extract agent storage helpers Split the agent storage module into smaller helper functions to make the persistence logic easier to follow and reuse. Behaviour is unchanged. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/storage/agents.rs | 183 ++++++++++++++++++++ 1 file changed, 183 insertions(+) create mode 100644 crates/openhuman-core/src/storage/agents.rs diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs new file mode 100644 index 00000000000..a5a71952502 --- /dev/null +++ b/crates/openhuman-core/src/storage/agents.rs @@ -0,0 +1,183 @@ +//! The agents that keep records in their own storage scope, and a way for +//! background work to visit each of them. +//! +//! Work done inside an agent's turn runs under that agent's `CoreContext` +//! (`session_agent`), so with a storage backend installed its cron jobs, +//! flows, approvals and the rest land in that agent's scope. Background work +//! — the cron scheduler, pollers, boot sweeps — runs under the process +//! default context, which names no agent, and on its own would only ever see +//! the `local` scope. [`for_each_scope`] closes that gap: it runs a step once +//! for `local` and once under each known agent's context. +//! +//! An agent is known when a context is derived for it in this process +//! ([`registered`], called by `CoreContext::derive_with`) — the live context +//! is used, with the agent's own configuration — or when an earlier process +//! did and recorded its id in the backend's `local` scope +//! (`storage_agents`), so a restart still visits agents the host has not +//! re-created yet; those are visited under the default context with the +//! agent swapped in (`CoreContext::for_agent`). +//! +//! Without a backend nothing here changes behavior: the SQLite stores do not +//! split by agent, so [`for_each_scope`] runs the step once. In SaaS mode the +//! process has no `local` scope and per-user background work is driven by +//! `user_agents::background`, so agent ids are not recorded and +//! [`for_each_scope`] visits only live agent contexts. + +use std::collections::{BTreeMap, HashSet}; +use std::future::Future; +use std::sync::{Arc, LazyLock, Mutex, Weak}; + +use serde_json::json; +use tinystoragedrivers::{CollectionSpec, Precondition, Query, Scope}; + +use super::{block_on, installed, DocumentStoreExt, StorageBackend}; +use crate::core::runtime::CoreContext; + +/// The `local`-scope collection recording which agents have their own scope. +const AGENTS: &str = "storage_agents"; + +/// Live agent contexts, by agent id. +static LIVE: LazyLock>>> = + LazyLock::new(Default::default); + +/// Agent ids this process has already recorded in the backend. +static RECORDED: LazyLock>> = LazyLock::new(Default::default); + +/// Records `context` as its agent's live context (when it names one) and +/// returns it unchanged. `CoreContext::derive_with` passes every derived +/// context through here. +pub fn registered(context: Arc) -> Arc { + if let Some(agent) = context.session_agent() { + LIVE.lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert(agent.to_string(), Arc::downgrade(&context)); + record(agent); + } + context +} + +/// Records every live agent in the backend — for agents derived before the +/// host installed its backend. Called by [`super::install`]. +pub(super) fn record_live() { + let agents: Vec = LIVE + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .keys() + .cloned() + .collect(); + for agent in agents { + record(&agent); + } +} + +/// Writes `agent` to the backend's `storage_agents` collection, once per +/// process. Best effort: a failure is logged and retried on the next +/// registration, and only costs a restarted process its visits to that +/// agent until the agent is derived again. +fn record(agent: &str) { + if crate::core::runtime::mode::is_saas() { + return; + } + let Some(backend) = installed() else { + return; + }; + if RECORDED + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .contains(agent) + { + return; + } + let id = agent.to_string(); + let result = block_on(async move { + let docs = Arc::clone(backend.for_scope(&Scope::local())?.documents()); + docs.ensure_collection(&CollectionSpec::new(AGENTS)).await?; + docs.put(AGENTS, &id, json!({}), Precondition::None).await?; + Ok(()) + }); + match result { + Ok(()) => { + RECORDED + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert(agent.to_string()); + tracing::debug!(%agent, "[storage::agents] recorded agent scope"); + } + Err(error) => { + tracing::warn!(%agent, %error, "[storage::agents] could not record agent scope"); + } + } +} + +/// Agent ids recorded in the backend by this or an earlier process. +fn recorded(backend: Arc) -> Vec { + let result = block_on(async move { + let docs = Arc::clone(backend.for_scope(&Scope::local())?.documents()); + docs.ensure_collection(&CollectionSpec::new(AGENTS)).await?; + let stored = docs.query_all(AGENTS, &Query::all()).await?; + Ok(stored.into_iter().map(|doc| doc.id).collect::>()) + }); + result.unwrap_or_else(|error| { + tracing::warn!(%error, "[storage::agents] could not list recorded agent scopes"); + Vec::new() + }) +} + +/// The agents [`for_each_scope`] visits, each with the context to visit it +/// under: its live context when one exists, else `fallback` acting for it. +fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc)> { + let mut contexts: BTreeMap> = BTreeMap::new(); + { + let mut live = LIVE.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + live.retain(|agent, context| match context.upgrade() { + Some(context) => { + contexts.insert(agent.clone(), context); + true + } + None => false, + }); + } + if !crate::core::runtime::mode::is_saas() { + if let (Some(backend), Some(fallback)) = (installed(), fallback) { + for agent in recorded(backend) { + contexts + .entry(agent.clone()) + .or_insert_with(|| fallback.for_agent(&agent)); + } + } + } + contexts.into_iter().collect() +} + +/// Runs `step` for every storage scope background work must cover: once +/// under the current context (the `local` scope, outside SaaS mode), then — +/// when a storage backend is installed — once under each known agent's +/// context. Steps run one after another; each result is returned with the +/// agent it ran for (`None` for `local`). +/// +/// `label` names the caller in logs. +pub async fn for_each_scope(label: &str, step: F) -> Vec<(Option, T)> +where + F: Fn() -> Fut, + Fut: Future, +{ + let mut results = Vec::new(); + let saas = crate::core::runtime::mode::is_saas(); + if !saas { + results.push((None, step().await)); + } + if installed().is_none() { + return results; + } + let fallback = CoreContext::current(); + for (agent, context) in agent_contexts(fallback.as_ref()) { + tracing::trace!(%agent, label, "[storage::agents] visiting agent scope"); + let value = CoreContext::scope(context, step()).await; + results.push((Some(agent), value)); + } + results +} + +#[cfg(test)] +#[path = "agents_tests.rs"] +mod tests; From 061df3588a4f91ad5d5f5068483949f7ad3f12fd Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:40:43 +0300 Subject: [PATCH 03/26] feat(storage): register agents module and record live agents on install Adds an agents submodule to storage and calls its record_live hook when a backend is installed, so agents derived before the backend existed are still recorded. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/storage/mod.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/crates/openhuman-core/src/storage/mod.rs b/crates/openhuman-core/src/storage/mod.rs index 1ff9cf94e89..490f93423f3 100644 --- a/crates/openhuman-core/src/storage/mod.rs +++ b/crates/openhuman-core/src/storage/mod.rs @@ -18,6 +18,7 @@ //! driver the build does not carry fails at [`open`] naming the feature, so a //! misconfigured deployment stops at boot instead of at its first write. +pub mod agents; pub mod documents; pub mod secrets; @@ -102,7 +103,10 @@ pub async fn open(url: &str) -> Result, StorageError> { /// Makes `backend` the process's storage backend; returns the previous one. pub fn install(backend: Arc) -> Option> { - BACKEND.install(backend) + let previous = BACKEND.install(backend); + // Agents derived before the backend existed still need recording. + agents::record_live(); + previous } /// The installed backend, when the host configured one. From 28d3875973bc7dcbf636a13e71c387f5fe7c795f Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:41:23 +0300 Subject: [PATCH 04/26] revert: restore the tinyagents pin #7197 moved back #7197 (a leftover storage-secrets branch) moved vendor/tinyagents from 33a86887 back to 37185fdc, which predates the tool-rules API (ToolRulePolicy, tinytools::ToolRules/Surface) that main's code uses since #7175, so main stopped compiling. Restores the pin and the Cargo.lock line. This reverts commit b2924d852d, reversing changes made to d670efc2b6. Co-authored-by: Medulla --- Cargo.lock | 1 + vendor/tinyagents | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/Cargo.lock b/Cargo.lock index 04bc3202020..a00d25c649a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7008,6 +7008,7 @@ version = "2.1.3" dependencies = [ "async-trait", "serde", + "tinytools", ] [[package]] diff --git a/vendor/tinyagents b/vendor/tinyagents index 37185fdcfc5..33a86887d52 160000 --- a/vendor/tinyagents +++ b/vendor/tinyagents @@ -1 +1 @@ -Subproject commit 37185fdcfc513ea2d9ea53608c069683975966e0 +Subproject commit 33a86887d52f1dba4e39fdb598b7694c20cfdc84 From ca3d663c92ffea9165815bd9212e1c5906e1865e Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:45:25 +0300 Subject: [PATCH 05/26] feat(cron): run scheduled jobs and task polls per agent scope Background work now runs once per agent storage scope in addition to the local pass, so jobs and task sources an agent scheduled from its own turn execute under that agent's context. The scheduler keeps its process-wide health tracking on the local pass only, while agent passes still report failing jobs. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/cron/scheduler.rs | 16 ++++++++++ .../src/integrations/task_sources/periodic.rs | 10 ++++-- crates/openhuman-core/src/storage/agents.rs | 32 ++++++++++++++----- 3 files changed, 47 insertions(+), 11 deletions(-) diff --git a/crates/openhuman-core/src/cron/scheduler.rs b/crates/openhuman-core/src/cron/scheduler.rs index 8e87f3fa2fc..0b51b21ac12 100644 --- a/crates/openhuman-core/src/cron/scheduler.rs +++ b/crates/openhuman-core/src/cron/scheduler.rs @@ -75,9 +75,25 @@ pub async fn run(config: Config) -> Result<()> { loop { interval.tick().await; tick_once(&config, &security, &mut last_emitted_health).await; + tick_agents(&config, &security).await; } } +/// The same poll for every agent that keeps its jobs in its own storage scope +/// (`crate::storage::agents`): a job an agent scheduled from its own turn +/// lives there, and runs there — under that agent's context and, when the +/// agent is live, its configuration. Only the `local` pass reports the +/// scheduler's health; an agent pass still reports a failing job. +pub(crate) async fn tick_agents(config: &Config, security: &Arc) { + crate::storage::agents::for_each_agent("cron", || async { + let config = crate::core::runtime::CoreContext::current_embedder_config() + .unwrap_or_else(|| config.clone()); + let mut steady = Some(true); + tick_once(&config, security, &mut steady).await; + }) + .await; +} + /// Single poll cycle of the scheduler loop, extracted so tests can drive /// it without owning `tokio::time::interval`. /// diff --git a/crates/openhuman-core/src/integrations/task_sources/periodic.rs b/crates/openhuman-core/src/integrations/task_sources/periodic.rs index faac8d2e004..bd34c8656c8 100644 --- a/crates/openhuman-core/src/integrations/task_sources/periodic.rs +++ b/crates/openhuman-core/src/integrations/task_sources/periodic.rs @@ -70,7 +70,7 @@ pub fn start_periodic_poll() { tracing::debug!("[task_sources:periodic] scheduler already running, skipping start"); return; } - tokio::spawn(async move { + crate::core::runtime::spawn_scoped(async move { tracing::info!( tick_seconds = TICK_SECONDS, "[task_sources:periodic] scheduler starting" @@ -87,8 +87,12 @@ async fn run_loop() { ticker.tick().await; loop { ticker.tick().await; - if let Err(e) = run_one_tick().await { - tracing::warn!(error = %e, "[task_sources:periodic] tick failed (continuing)"); + // `local`, then every agent that keeps its sources in its own storage + // scope (`crate::storage::agents`). + for (agent, result) in crate::storage::agents::for_each_scope("task_sources", run_one_tick).await { + if let Err(e) = result { + tracing::warn!(error = %e, agent = agent.as_deref().unwrap_or("local"), "[task_sources:periodic] tick failed (continuing)"); + } } } } diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs index a5a71952502..493b5dc0b03 100644 --- a/crates/openhuman-core/src/storage/agents.rs +++ b/crates/openhuman-core/src/storage/agents.rs @@ -150,10 +150,9 @@ fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc(label: &str, step: F) -> Vec<(Option, T)> @@ -162,18 +161,35 @@ where Fut: Future, { let mut results = Vec::new(); - let saas = crate::core::runtime::mode::is_saas(); - if !saas { + if !crate::core::runtime::mode::is_saas() { results.push((None, step().await)); } + for (agent, value) in for_each_agent(label, step).await { + results.push((Some(agent), value)); + } + results +} + +/// Runs `step` once under each known agent's context, one after another — +/// only when a storage backend is installed, since without one the stores do +/// not split by agent and the `local` pass already covers everything. +/// +/// For a loop that handles the `local` scope itself (the cron scheduler keeps +/// its process-wide health tracking there) and needs the agents on top. +pub async fn for_each_agent(label: &str, step: F) -> Vec<(String, T)> +where + F: Fn() -> Fut, + Fut: Future, +{ if installed().is_none() { - return results; + return Vec::new(); } let fallback = CoreContext::current(); + let mut results = Vec::new(); for (agent, context) in agent_contexts(fallback.as_ref()) { tracing::trace!(%agent, label, "[storage::agents] visiting agent scope"); let value = CoreContext::scope(context, step()).await; - results.push((Some(agent), value)); + results.push((agent, value)); } results } From ff3e1761c2e69dd7a4973e4c521cd0b46efbf0d9 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:48:24 +0300 Subject: [PATCH 06/26] style: apply rustfmt formatting to storage and task source modules Reformatted long expressions and reordered module declarations to match rustfmt output. No behaviour changes. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/core/runtime/context.rs | 4 +- .../src/integrations/task_sources/periodic.rs | 4 +- crates/openhuman-core/src/storage/agents.rs | 7 ++- .../src/storage/agents_tests.rs | 61 +++++++++++++++++++ 4 files changed, 70 insertions(+), 6 deletions(-) create mode 100644 crates/openhuman-core/src/storage/agents_tests.rs diff --git a/crates/openhuman-core/src/core/runtime/context.rs b/crates/openhuman-core/src/core/runtime/context.rs index d3956881709..5970b2d4982 100644 --- a/crates/openhuman-core/src/core/runtime/context.rs +++ b/crates/openhuman-core/src/core/runtime/context.rs @@ -738,10 +738,10 @@ pub async fn init_stores(cfg: &crate::config::Config, domains: crate::core::runt } } -#[path = "context_turn_origin.rs"] -mod turn_origin_scope; #[path = "context_agent.rs"] mod agent_scope; +#[path = "context_turn_origin.rs"] +mod turn_origin_scope; #[cfg(test)] #[path = "context_tests.rs"] diff --git a/crates/openhuman-core/src/integrations/task_sources/periodic.rs b/crates/openhuman-core/src/integrations/task_sources/periodic.rs index bd34c8656c8..d1c236b003c 100644 --- a/crates/openhuman-core/src/integrations/task_sources/periodic.rs +++ b/crates/openhuman-core/src/integrations/task_sources/periodic.rs @@ -89,7 +89,9 @@ async fn run_loop() { ticker.tick().await; // `local`, then every agent that keeps its sources in its own storage // scope (`crate::storage::agents`). - for (agent, result) in crate::storage::agents::for_each_scope("task_sources", run_one_tick).await { + for (agent, result) in + crate::storage::agents::for_each_scope("task_sources", run_one_tick).await + { if let Err(e) = result { tracing::warn!(error = %e, agent = agent.as_deref().unwrap_or("local"), "[task_sources:periodic] tick failed (continuing)"); } diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs index 493b5dc0b03..aa669310f49 100644 --- a/crates/openhuman-core/src/storage/agents.rs +++ b/crates/openhuman-core/src/storage/agents.rs @@ -37,8 +37,7 @@ use crate::core::runtime::CoreContext; const AGENTS: &str = "storage_agents"; /// Live agent contexts, by agent id. -static LIVE: LazyLock>>> = - LazyLock::new(Default::default); +static LIVE: LazyLock>>> = LazyLock::new(Default::default); /// Agent ids this process has already recorded in the backend. static RECORDED: LazyLock>> = LazyLock::new(Default::default); @@ -128,7 +127,9 @@ fn recorded(backend: Arc) -> Vec { fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc)> { let mut contexts: BTreeMap> = BTreeMap::new(); { - let mut live = LIVE.lock().unwrap_or_else(std::sync::PoisonError::into_inner); + let mut live = LIVE + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); live.retain(|agent, context| match context.upgrade() { Some(context) => { contexts.insert(agent.clone(), context); diff --git a/crates/openhuman-core/src/storage/agents_tests.rs b/crates/openhuman-core/src/storage/agents_tests.rs new file mode 100644 index 00000000000..4e7160e29c1 --- /dev/null +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -0,0 +1,61 @@ +use super::*; +use crate::core::runtime::{ContextOverlay, DomainSet}; + +fn agent_context(agent: &str) -> Arc { + CoreContext::for_test(DomainSet::full(), None).derive_with( + ContextOverlay::new( + crate::config::Config::default(), + DomainSet::full(), + Default::default(), + ) + .session_agent(agent), + ) +} + +fn live_agents() -> Vec { + agent_contexts(None).into_iter().map(|(id, _)| id).collect() +} + +#[test] +fn deriving_an_agent_context_registers_it_until_dropped() { + let context = agent_context("agents-test-live"); + assert!(live_agents().contains(&"agents-test-live".to_string())); + let (_, found) = agent_contexts(None) + .into_iter() + .find(|(id, _)| id == "agents-test-live") + .unwrap(); + assert!(Arc::ptr_eq(&found, &context), "the live context is used"); + drop((found, context)); + assert!(!live_agents().contains(&"agents-test-live".to_string())); +} + +#[test] +fn a_context_without_an_agent_is_not_registered() { + let before = live_agents().len(); + let _plain = CoreContext::for_test(DomainSet::full(), None); + assert_eq!(live_agents().len(), before); +} + +#[test] +fn for_agent_swaps_only_the_agent() { + let parent = agent_context("agents-test-parent"); + let child = parent.for_agent("agents-test-other"); + assert_eq!(child.session_agent(), Some("agents-test-other")); + assert_eq!(parent.session_agent(), Some("agents-test-parent")); +} + +#[tokio::test] +async fn without_a_backend_only_the_local_scope_runs() { + // The lib test binary never installs a backend into the process slot. + if installed().is_some() { + return; + } + let _agent = agent_context("agents-test-unvisited"); + let runs = for_each_scope("test", || async { + CoreContext::current().and_then(|c| c.session_agent().map(str::to_string)) + }) + .await; + assert_eq!(runs.len(), 1); + assert_eq!(runs[0].0, None); + assert!(for_each_agent("test", || async {}).await.is_empty()); +} From d7d7d37e03882c0f0be4640a8ca5eae5503c4675 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:49:55 +0300 Subject: [PATCH 07/26] fix(flows): sweep and reconcile runs across all agent storage scopes Boot-time sweeps for orphaned runs and schedule trigger reconciliation now run once per storage scope, covering the local scope and every agent that keeps its own runs and flows, instead of only the local scope. This ensures agents' interrupted runs are marked resumable and their cron jobs are re-registered on boot. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/flows/ops/run_management.rs | 13 +++++++++++ .../openhuman-core/src/flows/ops/triggers.rs | 22 +++++++++++++++++++ 2 files changed, 35 insertions(+) diff --git a/crates/openhuman-core/src/flows/ops/run_management.rs b/crates/openhuman-core/src/flows/ops/run_management.rs index 10be3a787ca..684da8c6346 100644 --- a/crates/openhuman-core/src/flows/ops/run_management.rs +++ b/crates/openhuman-core/src/flows/ops/run_management.rs @@ -270,6 +270,19 @@ pub async fn sweep_expired_parked_runs(config: &Config) -> usize { /// resumable — only `pending_approval` is). Best-effort by construction: a store /// error is logged and the sweep returns what it managed. pub async fn sweep_orphaned_running_runs_on_boot(config: &Config) -> usize { + // `local`, then every agent that keeps its runs in its own storage scope + // (`crate::storage::agents`). + crate::storage::agents::for_each_scope("flows boot sweep", || { + sweep_orphaned_running_runs_in_scope(config) + }) + .await + .into_iter() + .map(|(_, swept)| swept) + .sum() +} + +/// [`sweep_orphaned_running_runs_on_boot`] for the current storage scope. +async fn sweep_orphaned_running_runs_in_scope(config: &Config) -> usize { let now_str = Utc::now().to_rfc3339(); const REASON: &str = "Run interrupted by an app restart — no live run was executing this row after boot."; diff --git a/crates/openhuman-core/src/flows/ops/triggers.rs b/crates/openhuman-core/src/flows/ops/triggers.rs index 93a1212c779..c7b48941c8d 100644 --- a/crates/openhuman-core/src/flows/ops/triggers.rs +++ b/crates/openhuman-core/src/flows/ops/triggers.rs @@ -244,6 +244,28 @@ fn log_webhook_trigger_deferred(flow: &Flow, enabled: bool) { /// was lost some other way) gets its schedule re-registered on the next /// boot without the user having to toggle it off and on. pub async fn reconcile_schedule_triggers_on_boot(config: &Config) -> Result<(), String> { + // `local`, then every agent that keeps its flows in its own storage scope + // (`crate::storage::agents`); each scope's flows get their cron jobs in + // that scope, where the scheduler visits them. + let mut errors = Vec::new(); + for (agent, result) in crate::storage::agents::for_each_scope("flows schedule reconcile", || { + reconcile_schedule_triggers_in_scope(config) + }) + .await + { + if let Err(error) = result { + errors.push(format!("{}: {error}", agent.as_deref().unwrap_or("local"))); + } + } + if errors.is_empty() { + Ok(()) + } else { + Err(errors.join("; ")) + } +} + +/// [`reconcile_schedule_triggers_on_boot`] for the current storage scope. +async fn reconcile_schedule_triggers_in_scope(config: &Config) -> Result<(), String> { let (flows, skipped) = store::list_enabled_flows(config).map_err(|e| e.to_string())?; if skipped > 0 { // R-M4: a corrupt/unmigratable row must not abort boot reconciliation From fbfa3b449f6a659695dc08906c2117663ec6bbb4 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:51:52 +0300 Subject: [PATCH 08/26] fix(agent): reap orphaned runs across every agent scope The run reaper now sweeps the local scope and each known agent scope via for_each_scope, summing the reaped counts, since every agent keeps its own status store on a storage backend. The schedule trigger reconcile loop was reformatted without behaviour change. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/agent/tinyagents/reaper.rs | 10 +++++++++- crates/openhuman-core/src/flows/ops/triggers.rs | 9 +++++---- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/crates/openhuman-core/src/agent/tinyagents/reaper.rs b/crates/openhuman-core/src/agent/tinyagents/reaper.rs index 266cf6c8c50..5b66fbb5385 100644 --- a/crates/openhuman-core/src/agent/tinyagents/reaper.rs +++ b/crates/openhuman-core/src/agent/tinyagents/reaper.rs @@ -17,8 +17,16 @@ use tinyagents_session::transcript::import::ops::open_session_stores; /// Skipped on a storage backend other processes may share (MongoDB): there a /// non-terminal status can belong to a run another replica is still driving, /// and cancelling it would hide that run from active and late-attach status. +/// +/// With a storage backend every agent keeps its own status store, so the +/// sweep visits `local` and then each known agent (`crate::storage::agents`). pub(crate) async fn reap_orphaned_runs(workspace: &Path) -> usize { - reap_unless_shared(workspace, crate::storage::installed_is_shared()).await + let shared = crate::storage::installed_is_shared(); + crate::storage::agents::for_each_scope("run reaper", || reap_unless_shared(workspace, shared)) + .await + .into_iter() + .map(|(_, reaped)| reaped) + .sum() } async fn reap_unless_shared(workspace: &Path, shared: bool) -> usize { diff --git a/crates/openhuman-core/src/flows/ops/triggers.rs b/crates/openhuman-core/src/flows/ops/triggers.rs index c7b48941c8d..2fc64bdc2ec 100644 --- a/crates/openhuman-core/src/flows/ops/triggers.rs +++ b/crates/openhuman-core/src/flows/ops/triggers.rs @@ -248,10 +248,11 @@ pub async fn reconcile_schedule_triggers_on_boot(config: &Config) -> Result<(), // (`crate::storage::agents`); each scope's flows get their cron jobs in // that scope, where the scheduler visits them. let mut errors = Vec::new(); - for (agent, result) in crate::storage::agents::for_each_scope("flows schedule reconcile", || { - reconcile_schedule_triggers_in_scope(config) - }) - .await + for (agent, result) in + crate::storage::agents::for_each_scope("flows schedule reconcile", || { + reconcile_schedule_triggers_in_scope(config) + }) + .await { if let Err(error) = result { errors.push(format!("{}: {error}", agent.as_deref().unwrap_or("local"))); From f475d87bc16ad285ec723671192758e208c191c2 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:52:38 +0300 Subject: [PATCH 09/26] feat(devices): track pairing agent and add agent-scoped context helpers Pairing sessions now record the agent that started them, so a paired device is stored under and its tunnel frames are handled as that agent. Added storage helpers to resolve an agent's live context and to run background work within that agent's scope. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/security/devices/rpc.rs | 2 ++ .../src/security/devices/types.rs | 5 +++++ crates/openhuman-core/src/storage/agents.rs | 21 +++++++++++++++++++ 3 files changed, 28 insertions(+) diff --git a/crates/openhuman-core/src/security/devices/rpc.rs b/crates/openhuman-core/src/security/devices/rpc.rs index 88546d95260..f4558da1942 100644 --- a/crates/openhuman-core/src/security/devices/rpc.rs +++ b/crates/openhuman-core/src/security/devices/rpc.rs @@ -137,6 +137,8 @@ pub async fn devices_create_pairing( core_pubkey: core_pubkey.clone(), rpc_url: rpc_url.clone(), expires_at: expires_at.clone(), + agent: crate::core::runtime::CoreContext::current() + .and_then(|context| context.session_agent().map(str::to_string)), }, ); diff --git a/crates/openhuman-core/src/security/devices/types.rs b/crates/openhuman-core/src/security/devices/types.rs index 1f2642f2488..135a1bdac35 100644 --- a/crates/openhuman-core/src/security/devices/types.rs +++ b/crates/openhuman-core/src/security/devices/types.rs @@ -39,6 +39,11 @@ pub struct PairingSession { pub rpc_url: Option, /// ISO 8601 timestamp when the pairing token expires. pub expires_at: String, + /// The agent that started the pairing (`CoreContext::session_agent`), + /// when it was acting for one: the paired device is stored in, and its + /// tunnel frames are handled as, that agent (`devices::owner`). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub agent: Option, } /// Response payload for `devices_create_pairing`. diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs index aa669310f49..1fbc6eaf452 100644 --- a/crates/openhuman-core/src/storage/agents.rs +++ b/crates/openhuman-core/src/storage/agents.rs @@ -150,6 +150,27 @@ fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc Option> { + let live = LIVE + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .get(agent) + .and_then(Weak::upgrade); + live.or_else(|| CoreContext::current().map(|current| current.for_agent(agent))) +} + +/// Runs `fut` acting for `agent` when there is one — background work that +/// learned whose record it is handling (a device's pairing agent, an event's +/// publisher) re-enters that agent's scope — and as-is otherwise. +pub async fn within_agent(agent: Option<&str>, fut: F) -> F::Output { + match agent.and_then(context_for) { + Some(context) => CoreContext::scope(context, fut).await, + None => fut.await, + } +} + /// Runs `step` for every storage scope background work must cover: once /// under the current context (the `local` scope, outside SaaS mode), then /// once per known agent ([`for_each_agent`]). Each result is returned with From e756dc6bbb3041066d302769e91a202b0c61a936 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:52:58 +0300 Subject: [PATCH 10/26] feat(devices): scope tunnel frames to the owning agent Tunnel frames are now dispatched inside the agent that owns the device, so pairing records and RPCs land in the correct agent scope. The owner is resolved from the pending session and remembered once the device is persisted. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/security/devices/bus.rs | 20 ++++- .../src/security/devices/mod.rs | 1 + .../src/security/devices/owner.rs | 81 +++++++++++++++++++ 3 files changed, 101 insertions(+), 1 deletion(-) create mode 100644 crates/openhuman-core/src/security/devices/owner.rs diff --git a/crates/openhuman-core/src/security/devices/bus.rs b/crates/openhuman-core/src/security/devices/bus.rs index 979c0747c95..2b1a829b5b8 100644 --- a/crates/openhuman-core/src/security/devices/bus.rs +++ b/crates/openhuman-core/src/security/devices/bus.rs @@ -84,7 +84,20 @@ impl EventHandler for DeviceTunnelSubscriber { channel_id, payload_b64, } => { - handle_tunnel_frame(channel_id, payload_b64).await; + // Handle the frame as the agent the device belongs to, so its + // pairing record and the RPCs it sends land in that agent's + // scope (`super::owner`). + let pending = PENDING_SESSIONS + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .get(channel_id.as_str()) + .cloned(); + let owner = super::owner::owner_of(channel_id, pending.as_ref()).await; + crate::storage::agents::within_agent( + owner.as_deref(), + handle_tunnel_frame(channel_id, payload_b64), + ) + .await; } _ => {} } @@ -361,6 +374,11 @@ async fn handle_tunnel_frame(channel_id: &str, payload_b64: &str) { &session_token_hash, ) { Ok(device) => { + super::owner::remember( + channel_id, + crate::core::runtime::CoreContext::current() + .and_then(|context| context.session_agent().map(str::to_string)), + ); log::info!( "[devices/bus] device persisted channel_id={} label={}", device.channel_id, diff --git a/crates/openhuman-core/src/security/devices/mod.rs b/crates/openhuman-core/src/security/devices/mod.rs index e9127e206a7..0f0f85409bb 100644 --- a/crates/openhuman-core/src/security/devices/mod.rs +++ b/crates/openhuman-core/src/security/devices/mod.rs @@ -5,6 +5,7 @@ pub mod bus; pub mod crypto; +mod owner; pub mod rpc; pub mod schemas; pub mod store; diff --git a/crates/openhuman-core/src/security/devices/owner.rs b/crates/openhuman-core/src/security/devices/owner.rs new file mode 100644 index 00000000000..559753e9b7d --- /dev/null +++ b/crates/openhuman-core/src/security/devices/owner.rs @@ -0,0 +1,81 @@ +//! Which agent a paired device belongs to. +//! +//! A device paired from inside an agent's context is stored in that agent's +//! storage scope (`crate::storage`), and the RPCs it sends through the tunnel +//! must run as that agent — not as the process default, which would read and +//! write somebody else's records. The tunnel subscriber runs outside any +//! agent, so it asks here, per frame: +//! +//! 1. the pending pairing session, which recorded the agent that started it; +//! 2. the owners this process has already resolved; +//! 3. otherwise each storage scope, for a device paired by an earlier process +//! (`crate::storage::agents::for_each_scope`). +//! +//! `None` means the `local` scope (a single-user host, or no backend). + +use std::collections::HashMap; +use std::sync::{LazyLock, Mutex}; + +use super::types::PairingSession; + +/// Channels whose owner this process has resolved: `Some(agent)`, or `None` +/// for `local`. +static OWNERS: LazyLock>>> = + LazyLock::new(Default::default); + +/// Records that `channel_id` belongs to `agent` (`None` = `local`). +pub(super) fn remember(channel_id: &str, agent: Option) { + OWNERS + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert(channel_id.to_string(), agent); +} + +fn cached(channel_id: &str) -> Option> { + OWNERS + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .get(channel_id) + .cloned() +} + +/// The agent `channel_id` belongs to, resolved as described above. A channel +/// no scope knows yet (a handshake still in flight) is `local`. +pub(super) async fn owner_of( + channel_id: &str, + pending: Option<&PairingSession>, +) -> Option { + if let Some(session) = pending { + return session.agent.clone(); + } + if let Some(owner) = cached(channel_id) { + return owner; + } + let Ok(config) = crate::config::rpc::load_config_with_timeout().await else { + return None; + }; + let found = crate::storage::agents::for_each_scope("device owner", || async { + super::store::get_device(&config, channel_id) + .ok() + .flatten() + .is_some() + }) + .await + .into_iter() + .find_map(|(agent, has_device)| has_device.then_some(agent)); + match found { + Some(owner) => { + log::debug!( + "[devices/owner] channel_id={channel_id} belongs to agent={}", + owner.as_deref().unwrap_or("local") + ); + remember(channel_id, owner.clone()); + owner + } + None => None, + } +} + +#[cfg(test)] +#[path = "owner_tests.rs"] +mod tests; From 35a48fbaa18ed0bdff38189b54d5371339398168 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:55:04 +0300 Subject: [PATCH 11/26] style(security): reformat owner device module Collapse the static initializer and the owner_of signature onto single lines to match rustfmt output. No behaviour changes. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/security/devices/owner.rs | 8 +--- .../src/security/devices/owner_tests.rs | 44 +++++++++++++++++++ 2 files changed, 46 insertions(+), 6 deletions(-) create mode 100644 crates/openhuman-core/src/security/devices/owner_tests.rs diff --git a/crates/openhuman-core/src/security/devices/owner.rs b/crates/openhuman-core/src/security/devices/owner.rs index 559753e9b7d..ed5f8952eb1 100644 --- a/crates/openhuman-core/src/security/devices/owner.rs +++ b/crates/openhuman-core/src/security/devices/owner.rs @@ -20,8 +20,7 @@ use super::types::PairingSession; /// Channels whose owner this process has resolved: `Some(agent)`, or `None` /// for `local`. -static OWNERS: LazyLock>>> = - LazyLock::new(Default::default); +static OWNERS: LazyLock>>> = LazyLock::new(Default::default); /// Records that `channel_id` belongs to `agent` (`None` = `local`). pub(super) fn remember(channel_id: &str, agent: Option) { @@ -41,10 +40,7 @@ fn cached(channel_id: &str) -> Option> { /// The agent `channel_id` belongs to, resolved as described above. A channel /// no scope knows yet (a handshake still in flight) is `local`. -pub(super) async fn owner_of( - channel_id: &str, - pending: Option<&PairingSession>, -) -> Option { +pub(super) async fn owner_of(channel_id: &str, pending: Option<&PairingSession>) -> Option { if let Some(session) = pending { return session.agent.clone(); } diff --git a/crates/openhuman-core/src/security/devices/owner_tests.rs b/crates/openhuman-core/src/security/devices/owner_tests.rs new file mode 100644 index 00000000000..c756f226f15 --- /dev/null +++ b/crates/openhuman-core/src/security/devices/owner_tests.rs @@ -0,0 +1,44 @@ +use super::*; + +fn session(agent: Option<&str>) -> PairingSession { + PairingSession { + channel_id: "owner-test".to_string(), + pairing_token: "token".to_string(), + core_pubkey: "pk".to_string(), + rpc_url: None, + expires_at: "2099-01-01T00:00:00Z".to_string(), + agent: agent.map(str::to_string), + } +} + +#[tokio::test] +async fn a_pending_pairing_names_its_agent() { + let pending = session(Some("agent-7")); + assert_eq!( + owner_of("owner-test-pending", Some(&pending)) + .await + .as_deref(), + Some("agent-7") + ); + let local = session(None); + assert_eq!(owner_of("owner-test-pending", Some(&local)).await, None); +} + +#[tokio::test] +async fn a_remembered_owner_is_used_without_a_lookup() { + remember("owner-test-cached", Some("agent-9".to_string())); + assert_eq!( + owner_of("owner-test-cached", None).await.as_deref(), + Some("agent-9") + ); + remember("owner-test-local", None); + assert_eq!(owner_of("owner-test-local", None).await, None); +} + +#[test] +fn the_agent_is_not_sent_over_the_wire_when_absent() { + let json = serde_json::to_value(session(None)).unwrap(); + assert!(json.get("agent").is_none()); + let json = serde_json::to_value(session(Some("a"))).unwrap(); + assert_eq!(json["agent"], "a"); +} From 30a257822ff6ed86b6ad3d08a92237bf220e589e Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:58:16 +0300 Subject: [PATCH 12/26] test(cli): register storage scope e2e test binary Add a dedicated test target for the new storage scope end-to-end test so it runs in its own binary, since it boots a core and installs a storage backend into the process-wide slot. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-cli/Cargo.toml | 6 +++ tests/storage_scope_e2e.rs | 96 +++++++++++++++++++++++++++++++++ 2 files changed, 102 insertions(+) create mode 100644 tests/storage_scope_e2e.rs diff --git a/crates/openhuman-cli/Cargo.toml b/crates/openhuman-cli/Cargo.toml index e15948e9fc5..44b97bf7f39 100644 --- a/crates/openhuman-cli/Cargo.toml +++ b/crates/openhuman-cli/Cargo.toml @@ -104,6 +104,12 @@ path = "../../tests/storage_flows_e2e.rs" name = "storage_secrets_e2e" path = "../../tests/storage_secrets_e2e.rs" +[[test]] +# Its own binary: it boots a core and installs a storage backend into the +# process-wide slot. +name = "storage_scope_e2e" +path = "../../tests/storage_scope_e2e.rs" + [[test]] name = "embedded_server_shutdown_e2e" path = "../../tests/embedded_server_shutdown_e2e.rs" diff --git a/tests/storage_scope_e2e.rs b/tests/storage_scope_e2e.rs new file mode 100644 index 00000000000..a45f1c958cb --- /dev/null +++ b/tests/storage_scope_e2e.rs @@ -0,0 +1,96 @@ +//! Background work reaches the records an agent keeps in its own storage +//! scope. +//! +//! A job scheduled from inside an agent's context is stored in that agent's +//! scope, so the `local` pass a background loop makes on its own never sees +//! it; `storage::agents::for_each_scope` visits the agent too — through its +//! live context while it exists, and through the agent id the backend +//! recorded once it is gone (a restarted process). +//! +//! Its own test binary because it installs a backend into the process-wide +//! storage slot and boots a core. One test, so nothing in it races either. + +use std::sync::Arc; + +use openhuman_core::config::Config; +use openhuman_core::core::runtime::{ + ContextOverlay, CoreBuilder, CoreContext, DomainSet, ServiceSet, +}; +use openhuman_core::core::HostKind; +use openhuman_core::cron::{self, Schedule}; +use openhuman_core::storage::agents::for_each_scope; + +fn job_names(config: &Config) -> Vec { + cron::list_jobs(config) + .unwrap() + .into_iter() + .filter_map(|job| job.name) + .collect() +} + +#[tokio::test(flavor = "multi_thread")] +async fn background_work_visits_every_agent_scope() { + let workspace = tempfile::tempdir().unwrap(); + let config = Config { + workspace_dir: workspace.path().join("workspace"), + config_path: workspace.path().join("config.toml"), + ..Config::default() + }; + let _runtime = CoreBuilder::new(HostKind::Library) + .config(config.clone()) + .services(ServiceSet::none()) + .domains(DomainSet::none()) + .build() + .await + .unwrap(); + openhuman_core::storage::install(Arc::new(openhuman_core::storage::MemoryStorage::new())); + + let agent = CoreContext::current().unwrap().derive_with( + ContextOverlay::new(config.clone(), DomainSet::none(), Default::default()) + .session_agent("agent-e2e"), + ); + let scheduled = config.clone(); + CoreContext::scope(Arc::clone(&agent), async move { + tokio::task::spawn_blocking(move || { + cron::add_shell_job( + &scheduled, + Some("agent-job".to_string()), + Schedule::Every { every_ms: 60_000 }, + "echo hi", + ) + .unwrap(); + }) + .await + .unwrap(); + }) + .await; + + // The agent's job is invisible to the `local` scope. + let local = config.clone(); + assert!(tokio::task::spawn_blocking(move || job_names(&local)) + .await + .unwrap() + .is_empty()); + + let visit = || async { + let config = config.clone(); + tokio::task::spawn_blocking(move || job_names(&config)) + .await + .unwrap() + }; + // Visited through the live agent context … + let live = for_each_scope("e2e", || CoreContext::propagate(visit())).await; + assert!( + live.contains(&(Some("agent-e2e".to_string()), vec!["agent-job".to_string()])), + "{live:?}" + ); + + // … and, once the agent is gone, through the id the backend recorded. + drop(agent); + let recorded = for_each_scope("e2e", || CoreContext::propagate(visit())).await; + assert!( + recorded.contains(&(Some("agent-e2e".to_string()), vec!["agent-job".to_string()])), + "{recorded:?}" + ); + assert!(recorded.contains(&(None, Vec::new())), "{recorded:?}"); +} From cf8c019f8217f216c0811660c69b1bae8b2bffed Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:58:39 +0300 Subject: [PATCH 13/26] test(storage): reorder imports in storage scope e2e test Move the HostKind import alongside the other openhuman_core imports to keep the use statements grouped consistently. Auto-committed-on: dragonfly Co-authored-by: Medulla --- tests/storage_scope_e2e.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/storage_scope_e2e.rs b/tests/storage_scope_e2e.rs index a45f1c958cb..c6df149cec3 100644 --- a/tests/storage_scope_e2e.rs +++ b/tests/storage_scope_e2e.rs @@ -16,9 +16,9 @@ use openhuman_core::config::Config; use openhuman_core::core::runtime::{ ContextOverlay, CoreBuilder, CoreContext, DomainSet, ServiceSet, }; -use openhuman_core::core::HostKind; use openhuman_core::cron::{self, Schedule}; use openhuman_core::storage::agents::for_each_scope; +use openhuman_core::HostKind; fn job_names(config: &Config) -> Vec { cron::list_jobs(config) From 4266afd72febdfa31ca1082534bc2c3c85105a7e Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 17:59:05 +0300 Subject: [PATCH 14/26] test(storage): drop spawn_blocking from scope e2e test The store calls now block on storage's own bridge thread, so the test no longer needs to wrap them in spawn_blocking and can invoke them directly from the agent's task where its context is in scope. Auto-committed-on: dragonfly Co-authored-by: Medulla --- tests/storage_scope_e2e.rs | 37 ++++++++++++------------------------- 1 file changed, 12 insertions(+), 25 deletions(-) diff --git a/tests/storage_scope_e2e.rs b/tests/storage_scope_e2e.rs index c6df149cec3..a2219a0e4fa 100644 --- a/tests/storage_scope_e2e.rs +++ b/tests/storage_scope_e2e.rs @@ -49,37 +49,24 @@ async fn background_work_visits_every_agent_scope() { ContextOverlay::new(config.clone(), DomainSet::none(), Default::default()) .session_agent("agent-e2e"), ); - let scheduled = config.clone(); - CoreContext::scope(Arc::clone(&agent), async move { - tokio::task::spawn_blocking(move || { - cron::add_shell_job( - &scheduled, - Some("agent-job".to_string()), - Schedule::Every { every_ms: 60_000 }, - "echo hi", - ) - .unwrap(); - }) - .await + // The store calls block on storage's own bridge thread, so they are + // made straight from the agent's task, where its context is in scope. + CoreContext::scope(Arc::clone(&agent), async { + cron::add_shell_job( + &config, + Some("agent-job".to_string()), + Schedule::Every { every_ms: 60_000 }, + "echo hi", + ) .unwrap(); }) .await; // The agent's job is invisible to the `local` scope. - let local = config.clone(); - assert!(tokio::task::spawn_blocking(move || job_names(&local)) - .await - .unwrap() - .is_empty()); + assert!(job_names(&config).is_empty()); - let visit = || async { - let config = config.clone(); - tokio::task::spawn_blocking(move || job_names(&config)) - .await - .unwrap() - }; // Visited through the live agent context … - let live = for_each_scope("e2e", || CoreContext::propagate(visit())).await; + let live = for_each_scope("e2e", || async { job_names(&config) }).await; assert!( live.contains(&(Some("agent-e2e".to_string()), vec!["agent-job".to_string()])), "{live:?}" @@ -87,7 +74,7 @@ async fn background_work_visits_every_agent_scope() { // … and, once the agent is gone, through the id the backend recorded. drop(agent); - let recorded = for_each_scope("e2e", || CoreContext::propagate(visit())).await; + let recorded = for_each_scope("e2e", || async { job_names(&config) }).await; assert!( recorded.contains(&(Some("agent-e2e".to_string()), vec!["agent-job".to_string()])), "{recorded:?}" From 6e6b9018e671a9ca4940f85d598f919bacdb8cc0 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 18:07:29 +0300 Subject: [PATCH 15/26] chore(ci): refresh saas ambient baseline and lockfile Update the ambient baseline line numbers to match the current source and drop the periodic task source entry that no longer exists. The app lockfile now resolves hmac to 0.13.0 and adds sha2 0.11.0. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-app/Cargo.lock | 3 ++- scripts/ci/saas-ambient-baseline.json | 21 +++++++-------------- 2 files changed, 9 insertions(+), 15 deletions(-) diff --git a/crates/openhuman-app/Cargo.lock b/crates/openhuman-app/Cargo.lock index 27ee3deead7..8bba1d354fc 100644 --- a/crates/openhuman-app/Cargo.lock +++ b/crates/openhuman-app/Cargo.lock @@ -4426,7 +4426,7 @@ dependencies = [ "futures-util", "hex", "hkdf", - "hmac 0.12.1", + "hmac 0.13.0", "hostname", "iana-time-zone", "keyring", @@ -4447,6 +4447,7 @@ dependencies = [ "serde_json", "serde_path_to_error", "sha2 0.10.9", + "sha2 0.11.0", "starship-battery", "sysinfo", "tar", diff --git a/scripts/ci/saas-ambient-baseline.json b/scripts/ci/saas-ambient-baseline.json index fe123202861..f3669beb607 100644 --- a/scripts/ci/saas-ambient-baseline.json +++ b/scripts/ci/saas-ambient-baseline.json @@ -37,7 +37,7 @@ { "rule": "bare-spawn", "path": "crates/openhuman-core/src/agent/orchestration/background_delivery.rs", - "line": 164, + "line": 170, "text": "tokio::spawn(async move {", "occurrence": 1 }, @@ -121,14 +121,14 @@ { "rule": "bare-spawn", "path": "crates/openhuman-core/src/agent/session_host/runtime_session.rs", - "line": 732, + "line": 724, "text": "tokio::spawn(async move {", "occurrence": 1 }, { "rule": "bare-spawn", "path": "crates/openhuman-core/src/agent/session_host/runtime_session.rs", - "line": 734, + "line": 726, "text": "let transcript = tokio::task::spawn_blocking(move || {", "occurrence": 1 }, @@ -457,7 +457,7 @@ { "rule": "home-dir", "path": "crates/openhuman-core/src/core/runtime/saas.rs", - "line": 182, + "line": 238, "text": "home: dirs::home_dir(),", "occurrence": 1 }, @@ -790,13 +790,6 @@ "text": "tokio::spawn(async move {", "occurrence": 1 }, - { - "rule": "bare-spawn", - "path": "crates/openhuman-core/src/integrations/task_sources/periodic.rs", - "line": 73, - "text": "tokio::spawn(async move {", - "occurrence": 1 - }, { "rule": "bare-spawn", "path": "crates/openhuman-core/src/integrations/test_support_backend.rs", @@ -870,21 +863,21 @@ { "rule": "bare-spawn", "path": "crates/openhuman-core/src/memory/import.rs", - "line": 312, + "line": 313, "text": "let counts = tokio::task::spawn_blocking(move || count_legacy(&workspace_dir))", "occurrence": 1 }, { "rule": "bare-spawn", "path": "crates/openhuman-core/src/memory/import.rs", - "line": 551, + "line": 552, "text": "tokio::spawn(async move {", "occurrence": 1 }, { "rule": "bare-spawn", "path": "crates/openhuman-core/src/memory/import.rs", - "line": 584, + "line": 585, "text": "let reader = tokio::task::spawn_blocking(move || {", "occurrence": 1 }, From 4c7e9b183273698273fcab56f5222986704e4dd1 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 18:07:46 +0300 Subject: [PATCH 16/26] docs(storage): document background work and agent scopes Explain how background work runs under the process default context and would otherwise only see the local scope, and how storage::agents closes that gap through registered, for_each_scope, and within_agent. List the callers that use these helpers and note the behaviour without a backend and in SaaS mode. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/storage/README.md | 25 +++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/crates/openhuman-core/src/storage/README.md b/crates/openhuman-core/src/storage/README.md index 14e34a43d48..f659c2bcf35 100644 --- a/crates/openhuman-core/src/storage/README.md +++ b/crates/openhuman-core/src/storage/README.md @@ -71,6 +71,31 @@ on `storage-mongodb`. master key they fail closed. The config encryption key stays on the process keyring, because `config.toml` is loaded before any agent acts. +## Background work and agent scopes + +Work done inside an agent's turn runs under that agent's `CoreContext` +(`session_agent`, set for embed agents and SaaS user agents), so with a +backend installed its records land in that agent's scope. Background work +runs under the process default context and on its own would only see +`local`. `storage::agents` closes the gap: + +- `registered(context)`: `CoreContext::derive_with` records every agent + context it builds (live, by agent id), and the agent id is written to the + backend's `local` scope (`storage_agents`) so a restarted process still + knows it. +- `for_each_scope(label, step)`: runs `step` for `local`, then under each + known agent's context — its live one, or the default context acting for it + (`CoreContext::for_agent`). `for_each_agent` skips `local`. +- `within_agent(agent, fut)` / `context_for(agent)`: re-enter an agent's + scope when background work learned whose record it is handling. + +Users: the cron scheduler (`cron::scheduler::tick_agents`), the task-source +poller, the flows boot sweep and schedule-trigger reconcile, the run reaper, +and the device tunnel (a paired device's frames run as the agent that paired +it, `security::devices::owner`). Without a backend these run once, as +before. In SaaS mode agent ids are not recorded and `local` is skipped; +per-user background work there is `user_agents::background`. + ## Boundaries The ports, drivers, scopes and conformance suites are tinystoragedrivers'; From 24a7330ae290e16cfda01a983488faf4c5ee31c4 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 18:29:57 +0300 Subject: [PATCH 17/26] =?UTF-8?q?chore:=20=E5=A7=ABThe=20diff=20is=20empty?= =?UTF-8?q?=20=E2=80=94=20no=20changes=20are=20shown=20for=20any=20of=20th?= =?UTF-8?q?e=20listed=20files.=20Without=20any=20added?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit If you can paste the actual diff, I'll write the Conventional Commits message for it. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/cron/scheduler.rs | 27 +++++++++++--- .../src/flows/ops/run_management.rs | 9 +++++ .../src/integrations/task_sources/periodic.rs | 15 +++++++- .../src/security/devices/owner.rs | 8 ++-- crates/openhuman-core/src/storage/agents.rs | 37 ++++++++++++++----- crates/openhuman-core/src/storage/mod.rs | 3 ++ 6 files changed, 80 insertions(+), 19 deletions(-) diff --git a/crates/openhuman-core/src/cron/scheduler.rs b/crates/openhuman-core/src/cron/scheduler.rs index 0b51b21ac12..5eaab974361 100644 --- a/crates/openhuman-core/src/cron/scheduler.rs +++ b/crates/openhuman-core/src/cron/scheduler.rs @@ -75,7 +75,7 @@ pub async fn run(config: Config) -> Result<()> { loop { interval.tick().await; tick_once(&config, &security, &mut last_emitted_health).await; - tick_agents(&config, &security).await; + tick_agents(&config, &mut last_emitted_health).await; } } @@ -83,15 +83,32 @@ pub async fn run(config: Config) -> Result<()> { /// (`crate::storage::agents`): a job an agent scheduled from its own turn /// lives there, and runs there — under that agent's context and, when the /// agent is live, its configuration. Only the `local` pass reports the -/// scheduler's health; an agent pass still reports a failing job. -pub(crate) async fn tick_agents(config: &Config, security: &Arc) { +/// scheduler's health tracker (`last_emitted_health`, shared with the `local` +/// pass so a failure or recovery reported from an agent pass is seen by the +/// next one); each agent pass authorizes its jobs with a policy built from +/// that agent's own configuration. +pub(crate) async fn tick_agents(config: &Config, last_emitted_health: &mut Option) { + let health = std::sync::Mutex::new(*last_emitted_health); crate::storage::agents::for_each_agent("cron", || async { let config = crate::core::runtime::CoreContext::current_embedder_config() .unwrap_or_else(|| config.clone()); - let mut steady = Some(true); - tick_once(&config, security, &mut steady).await; + let security = Arc::new(SecurityPolicy::from_config( + &config.autonomy, + &config.workspace_dir, + &config.action_dir, + )); + let mut steady = *health + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + tick_once(&config, &security, &mut steady).await; + *health + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = steady; }) .await; + *last_emitted_health = health + .into_inner() + .unwrap_or_else(std::sync::PoisonError::into_inner); } /// Single poll cycle of the scheduler loop, extracted so tests can drive diff --git a/crates/openhuman-core/src/flows/ops/run_management.rs b/crates/openhuman-core/src/flows/ops/run_management.rs index 684da8c6346..77fb425d3c8 100644 --- a/crates/openhuman-core/src/flows/ops/run_management.rs +++ b/crates/openhuman-core/src/flows/ops/run_management.rs @@ -272,6 +272,15 @@ pub async fn sweep_expired_parked_runs(config: &Config) -> usize { pub async fn sweep_orphaned_running_runs_on_boot(config: &Config) -> usize { // `local`, then every agent that keeps its runs in its own storage scope // (`crate::storage::agents`). + // + // On a backend other processes share (MongoDB), a `running` row below the + // boot floor can belong to a run another replica is still driving, and + // sweeping it would drop that run's checkpoint — so the agent scopes are + // left alone there, as the agent run reaper does. + if crate::storage::installed_is_shared() && !crate::core::runtime::mode::is_saas() { + log::info!("[flows] boot sweep: agent scopes skipped, the storage backend is shared"); + return sweep_orphaned_running_runs_in_scope(config).await; + } crate::storage::agents::for_each_scope("flows boot sweep", || { sweep_orphaned_running_runs_in_scope(config) }) diff --git a/crates/openhuman-core/src/integrations/task_sources/periodic.rs b/crates/openhuman-core/src/integrations/task_sources/periodic.rs index d1c236b003c..f5464fa82ef 100644 --- a/crates/openhuman-core/src/integrations/task_sources/periodic.rs +++ b/crates/openhuman-core/src/integrations/task_sources/periodic.rs @@ -43,10 +43,21 @@ fn last_poll_map() -> &'static LastPollMap { LAST_POLL_AT.get_or_init(|| Mutex::new(HashMap::new())) } +/// The poll-tracking key for `source_id` in the current storage scope: the +/// same source id in two agents' scopes is two different sources. +fn poll_key(source_id: &str) -> String { + let agent = crate::core::runtime::CoreContext::current() + .and_then(|context| context.session_agent().map(str::to_string)); + match agent { + Some(agent) => format!("{agent}\u{0}{source_id}"), + None => source_id.to_string(), + } +} + /// Record a successful (or attempted) poll for a source id. fn record_poll(source_id: &str) { if let Ok(mut map) = last_poll_map().lock() { - map.insert(source_id.to_string(), Instant::now()); + map.insert(poll_key(source_id), Instant::now()); } } @@ -57,7 +68,7 @@ fn is_due(source: &TaskSource) -> bool { Ok(map) => map, Err(poisoned) => poisoned.into_inner(), }; - match map.get(&source.id) { + match map.get(&poll_key(&source.id)) { Some(when) => when.elapsed() >= Duration::from_secs(interval_secs), None => true, // never polled this run — fire immediately } diff --git a/crates/openhuman-core/src/security/devices/owner.rs b/crates/openhuman-core/src/security/devices/owner.rs index ed5f8952eb1..54fc515b68e 100644 --- a/crates/openhuman-core/src/security/devices/owner.rs +++ b/crates/openhuman-core/src/security/devices/owner.rs @@ -47,10 +47,12 @@ pub(super) async fn owner_of(channel_id: &str, pending: Option<&PairingSession>) if let Some(owner) = cached(channel_id) { return owner; } - let Ok(config) = crate::config::rpc::load_config_with_timeout().await else { - return None; - }; + // The configuration is loaded inside each scope: in SaaS mode loading it + // needs an acting agent, which this tunnel task does not have. let found = crate::storage::agents::for_each_scope("device owner", || async { + let Ok(config) = crate::config::rpc::load_config_with_timeout().await else { + return false; + }; super::store::get_device(&config, channel_id) .ok() .flatten() diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs index 1fbc6eaf452..9d709001b2f 100644 --- a/crates/openhuman-core/src/storage/agents.rs +++ b/crates/openhuman-core/src/storage/agents.rs @@ -37,7 +37,11 @@ use crate::core::runtime::CoreContext; const AGENTS: &str = "storage_agents"; /// Live agent contexts, by agent id. -static LIVE: LazyLock>>> = LazyLock::new(Default::default); +/// +/// Several contexts can name one agent (each `derive_with` makes one), so each +/// agent keeps all of its live contexts, newest last. +static LIVE: LazyLock>>>> = + LazyLock::new(Default::default); /// Agent ids this process has already recorded in the backend. static RECORDED: LazyLock>> = LazyLock::new(Default::default); @@ -47,14 +51,28 @@ static RECORDED: LazyLock>> = LazyLock::new(Default::defau /// context through here. pub fn registered(context: Arc) -> Arc { if let Some(agent) = context.session_agent() { - LIVE.lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .insert(agent.to_string(), Arc::downgrade(&context)); + let mut live = LIVE + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let contexts = live.entry(agent.to_string()).or_default(); + contexts.retain(|existing| existing.strong_count() > 0); + contexts.push(Arc::downgrade(&context)); + drop(live); record(agent); } context } +/// Forgets which agents were recorded, so the next [`record`] writes them to +/// the backend now installed (or removed). Called by [`super::install`] and +/// [`super::clear`]: the record cache describes one backend, not the process. +pub(super) fn reset_recorded() { + RECORDED + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clear(); +} + /// Records every live agent in the backend — for agents derived before the /// host installed its backend. Called by [`super::install`]. pub(super) fn record_live() { @@ -130,12 +148,13 @@ fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc { + live.retain(|agent, entries| { + entries.retain(|entry| entry.strong_count() > 0); + // The newest context that is still alive acts for the agent. + if let Some(context) = entries.iter().rev().find_map(Weak::upgrade) { contexts.insert(agent.clone(), context); - true } - None => false, + !entries.is_empty() }); } if !crate::core::runtime::mode::is_saas() { @@ -157,7 +176,7 @@ pub fn context_for(agent: &str) -> Option> { .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .get(agent) - .and_then(Weak::upgrade); + .and_then(|entries| entries.iter().rev().find_map(Weak::upgrade)); live.or_else(|| CoreContext::current().map(|current| current.for_agent(agent))) } diff --git a/crates/openhuman-core/src/storage/mod.rs b/crates/openhuman-core/src/storage/mod.rs index 490f93423f3..5d67a8f2e15 100644 --- a/crates/openhuman-core/src/storage/mod.rs +++ b/crates/openhuman-core/src/storage/mod.rs @@ -104,6 +104,8 @@ pub async fn open(url: &str) -> Result, StorageError> { /// Makes `backend` the process's storage backend; returns the previous one. pub fn install(backend: Arc) -> Option> { let previous = BACKEND.install(backend); + // What was recorded described the previous backend. + agents::reset_recorded(); // Agents derived before the backend existed still need recording. agents::record_live(); previous @@ -116,6 +118,7 @@ pub fn installed() -> Option> { /// Removes the installed backend; returns whether there was one. pub fn clear() -> bool { + agents::reset_recorded(); BACKEND.clear() } From b980fde9b9bb8dccffcbc5ab480262f985fbcd45 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 18:30:03 +0300 Subject: [PATCH 18/26] chore(flows): use tracing for boot sweep log Switch the boot sweep skip message from log to tracing with the flows target so it is emitted through the same logging pipeline as the rest of the module. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/flows/ops/run_management.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/openhuman-core/src/flows/ops/run_management.rs b/crates/openhuman-core/src/flows/ops/run_management.rs index 77fb425d3c8..9c8a37592bd 100644 --- a/crates/openhuman-core/src/flows/ops/run_management.rs +++ b/crates/openhuman-core/src/flows/ops/run_management.rs @@ -278,7 +278,7 @@ pub async fn sweep_orphaned_running_runs_on_boot(config: &Config) -> usize { // sweeping it would drop that run's checkpoint — so the agent scopes are // left alone there, as the agent run reaper does. if crate::storage::installed_is_shared() && !crate::core::runtime::mode::is_saas() { - log::info!("[flows] boot sweep: agent scopes skipped, the storage backend is shared"); + tracing::info!(target: "flows", "[flows] boot sweep: agent scopes skipped, the storage backend is shared"); return sweep_orphaned_running_runs_in_scope(config).await; } crate::storage::agents::for_each_scope("flows boot sweep", || { From 5ef831cb5714f5191cf38d06b0bd3318a460b30e Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 18:30:20 +0300 Subject: [PATCH 19/26] test: cover per-agent poll scoping and agent context lifecycle Add tests asserting that poll timestamps are tracked per agent scope so one agent's poll does not mark the same source id as due for another, and that dropping a sibling context leaves the live one reachable while reset_recorded clears previously recorded entries. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../task_sources/periodic_tests.rs | 21 ++++++++++++++++++ .../src/storage/agents_tests.rs | 22 +++++++++++++++++++ 2 files changed, 43 insertions(+) diff --git a/crates/openhuman-core/src/integrations/task_sources/periodic_tests.rs b/crates/openhuman-core/src/integrations/task_sources/periodic_tests.rs index 3fa09a3ec8c..ebad54b7611 100644 --- a/crates/openhuman-core/src/integrations/task_sources/periodic_tests.rs +++ b/crates/openhuman-core/src/integrations/task_sources/periodic_tests.rs @@ -206,3 +206,24 @@ async fn manual_fetch_still_returns_the_explanatory_error() { assert_eq!(outcome.fetched, 0); assert_eq!(outcome.routed, 0); } + +#[tokio::test] +async fn poll_timestamps_are_kept_per_agent_scope() { + use crate::core::runtime::{ContextOverlay, CoreContext, DomainSet}; + let agent = |name: &str| { + CoreContext::for_test(DomainSet::full(), None).derive_with( + ContextOverlay::new( + crate::config::Config::default(), + DomainSet::full(), + Default::default(), + ) + .session_agent(name), + ) + }; + let s = source("ts-per-agent-scope-xyz", 1800); + CoreContext::scope(agent("ts-agent-a"), async { record_poll(&s.id) }).await; + let due_for_a = CoreContext::scope(agent("ts-agent-a"), async { is_due(&s) }).await; + let due_for_b = CoreContext::scope(agent("ts-agent-b"), async { is_due(&s) }).await; + assert!(!due_for_a, "agent a just polled it"); + assert!(due_for_b, "agent b's same-id source is a different source"); +} diff --git a/crates/openhuman-core/src/storage/agents_tests.rs b/crates/openhuman-core/src/storage/agents_tests.rs index 4e7160e29c1..5d7d9d6d948 100644 --- a/crates/openhuman-core/src/storage/agents_tests.rs +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -59,3 +59,25 @@ async fn without_a_backend_only_the_local_scope_runs() { assert_eq!(runs[0].0, None); assert!(for_each_agent("test", || async {}).await.is_empty()); } + +#[test] +fn a_dropped_sibling_context_does_not_hide_a_live_one() { + let first = agent_context("agents-test-siblings"); + let second = agent_context("agents-test-siblings"); + drop(second); + let (_, found) = agent_contexts(None) + .into_iter() + .find(|(id, _)| id == "agents-test-siblings") + .expect("the first context is still live"); + assert!(Arc::ptr_eq(&found, &first)); + assert!(context_for("agents-test-siblings").is_some()); + drop((found, first)); + assert!(!live_agents().contains(&"agents-test-siblings".to_string())); +} + +#[test] +fn reset_recorded_forgets_what_was_recorded() { + RECORDED.lock().unwrap().insert("agents-test-recorded".into()); + reset_recorded(); + assert!(!RECORDED.lock().unwrap().contains("agents-test-recorded")); +} From 34d032b769b21c30290ce450117c33140bf35ce7 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 18:33:38 +0300 Subject: [PATCH 20/26] test(storage): reformat recorded state setup in agents test Reformatted the chained lock and insert call in the reset_recorded test to satisfy rustfmt line width limits. No behaviour change. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/storage/agents_tests.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/crates/openhuman-core/src/storage/agents_tests.rs b/crates/openhuman-core/src/storage/agents_tests.rs index 5d7d9d6d948..f5fd254db30 100644 --- a/crates/openhuman-core/src/storage/agents_tests.rs +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -77,7 +77,10 @@ fn a_dropped_sibling_context_does_not_hide_a_live_one() { #[test] fn reset_recorded_forgets_what_was_recorded() { - RECORDED.lock().unwrap().insert("agents-test-recorded".into()); + RECORDED + .lock() + .unwrap() + .insert("agents-test-recorded".into()); reset_recorded(); assert!(!RECORDED.lock().unwrap().contains("agents-test-recorded")); } From 73ae76f4fd9d24efcf1f71560057cc906bf65d2f Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 19:01:39 +0300 Subject: [PATCH 21/26] refactor(storage): make agent backend selection injectable Agent context and recording helpers now take the storage backend as an explicit argument, with the public entry points resolving it from the installed backend and SaaS mode. This lets tests drive the recorded-agent paths without touching global state, and the cron scheduler's per-agent poll is extracted into its own function. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/cron/scheduler.rs | 22 ++++++---- .../src/security/devices/owner_tests.rs | 37 ++++++++++++++++ crates/openhuman-core/src/storage/agents.rs | 31 ++++++++++--- .../src/storage/agents_tests.rs | 44 +++++++++++++++++++ 4 files changed, 119 insertions(+), 15 deletions(-) diff --git a/crates/openhuman-core/src/cron/scheduler.rs b/crates/openhuman-core/src/cron/scheduler.rs index 5eaab974361..ba27d2fa91a 100644 --- a/crates/openhuman-core/src/cron/scheduler.rs +++ b/crates/openhuman-core/src/cron/scheduler.rs @@ -90,17 +90,10 @@ pub async fn run(config: Config) -> Result<()> { pub(crate) async fn tick_agents(config: &Config, last_emitted_health: &mut Option) { let health = std::sync::Mutex::new(*last_emitted_health); crate::storage::agents::for_each_agent("cron", || async { - let config = crate::core::runtime::CoreContext::current_embedder_config() - .unwrap_or_else(|| config.clone()); - let security = Arc::new(SecurityPolicy::from_config( - &config.autonomy, - &config.workspace_dir, - &config.action_dir, - )); let mut steady = *health .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); - tick_once(&config, &security, &mut steady).await; + tick_agent_scope(config, &mut steady).await; *health .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) = steady; @@ -111,6 +104,19 @@ pub(crate) async fn tick_agents(config: &Config, last_emitted_health: &mut Optio .unwrap_or_else(std::sync::PoisonError::into_inner); } +/// One agent's pass: the poll under the current (agent) context, with that +/// agent's configuration when it has one and a policy built from it. +async fn tick_agent_scope(config: &Config, last_emitted_health: &mut Option) { + let config = crate::core::runtime::CoreContext::current_embedder_config() + .unwrap_or_else(|| config.clone()); + let security = Arc::new(SecurityPolicy::from_config( + &config.autonomy, + &config.workspace_dir, + &config.action_dir, + )); + tick_once(&config, &security, last_emitted_health).await; +} + /// Single poll cycle of the scheduler loop, extracted so tests can drive /// it without owning `tokio::time::interval`. /// diff --git a/crates/openhuman-core/src/security/devices/owner_tests.rs b/crates/openhuman-core/src/security/devices/owner_tests.rs index c756f226f15..4476e9adf89 100644 --- a/crates/openhuman-core/src/security/devices/owner_tests.rs +++ b/crates/openhuman-core/src/security/devices/owner_tests.rs @@ -42,3 +42,40 @@ fn the_agent_is_not_sent_over_the_wire_when_absent() { let json = serde_json::to_value(session(Some("a"))).unwrap(); assert_eq!(json["agent"], "a"); } + +fn scoped_config(dir: &std::path::Path) -> crate::config::Config { + crate::config::Config { + workspace_dir: dir.join("workspace"), + config_path: dir.join("config.toml"), + ..Default::default() + } +} + +#[tokio::test] +async fn a_device_found_in_the_local_scope_belongs_to_local_and_is_remembered() { + let tmp = tempfile::tempdir().unwrap(); + let config = scoped_config(tmp.path()); + super::super::store::insert_device(&config, "owner-test-found", "label", "pk", "hash") + .unwrap(); + let context = crate::core::runtime::CoreContext::for_test_with_config( + crate::core::runtime::DomainSet::full(), + config, + ); + let owner = crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-found", None)) + .await; + assert_eq!(owner, None); + assert_eq!(cached("owner-test-found"), Some(None)); +} + +#[tokio::test] +async fn a_channel_no_scope_knows_is_local_and_not_remembered() { + let tmp = tempfile::tempdir().unwrap(); + let context = crate::core::runtime::CoreContext::for_test_with_config( + crate::core::runtime::DomainSet::full(), + scoped_config(tmp.path()), + ); + let owner = crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-missing", None)) + .await; + assert_eq!(owner, None); + assert_eq!(cached("owner-test-missing"), None); +} diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs index 9d709001b2f..53f430b0709 100644 --- a/crates/openhuman-core/src/storage/agents.rs +++ b/crates/openhuman-core/src/storage/agents.rs @@ -98,6 +98,11 @@ fn record(agent: &str) { let Some(backend) = installed() else { return; }; + record_in(backend, agent); +} + +/// [`record`] against an explicit `backend`. +fn record_in(backend: Arc, agent: &str) { if RECORDED .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) @@ -143,6 +148,20 @@ fn recorded(backend: Arc) -> Vec { /// The agents [`for_each_scope`] visits, each with the context to visit it /// under: its live context when one exists, else `fallback` acting for it. fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc)> { + let backend = if crate::core::runtime::mode::is_saas() { + None + } else { + installed() + }; + agent_contexts_in(backend, fallback) +} + +/// [`agent_contexts`] with the backend whose recorded agents are visited +/// made explicit (`None` visits live contexts only). +fn agent_contexts_in( + backend: Option>, + fallback: Option<&Arc>, +) -> Vec<(String, Arc)> { let mut contexts: BTreeMap> = BTreeMap::new(); { let mut live = LIVE @@ -157,13 +176,11 @@ fn agent_contexts(fallback: Option<&Arc>) -> Vec<(String, Arc Arc { + Arc::new(crate::storage::MemoryStorage::new()) +} + +#[test] +fn a_recorded_agent_is_visited_through_the_fallback_context() { + let backend = memory_backend(); + reset_recorded(); + record_in(Arc::clone(&backend), "agents-test-recorded-only"); + assert!(recorded(Arc::clone(&backend)).contains(&"agents-test-recorded-only".to_string())); + + let fallback = CoreContext::for_test(DomainSet::full(), None); + let visited = agent_contexts_in(Some(backend), Some(&fallback)); + let (_, context) = visited + .iter() + .find(|(id, _)| id == "agents-test-recorded-only") + .expect("the recorded agent is visited"); + assert_eq!(context.session_agent(), Some("agents-test-recorded-only")); +} + +#[test] +fn recording_is_skipped_for_an_agent_already_recorded_to_that_backend() { + let first = memory_backend(); + reset_recorded(); + record_in(Arc::clone(&first), "agents-test-once"); + // A second backend is only populated once the cache is reset. + let second = memory_backend(); + record_in(Arc::clone(&second), "agents-test-once"); + assert!(recorded(Arc::clone(&second)).is_empty()); + reset_recorded(); + record_in(Arc::clone(&second), "agents-test-once"); + assert_eq!(recorded(second), vec!["agents-test-once".to_string()]); +} + +#[test] +fn without_a_fallback_only_live_contexts_are_visited() { + let backend = memory_backend(); + reset_recorded(); + record_in(Arc::clone(&backend), "agents-test-no-fallback"); + assert!(!agent_contexts_in(Some(backend), None) + .iter() + .any(|(id, _)| id == "agents-test-no-fallback")); +} From 17091db4f50773d1d35a9f81a146b5703091e529 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 19:02:15 +0300 Subject: [PATCH 22/26] test(cron): cover agent tick health reporting Add tests asserting that an agent pass polls with its own config and reports healthy on success, and that ticking agents without an installed storage backend leaves the health tracker untouched. The device owner tests were also reformatted to satisfy rustfmt. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/cron/scheduler_tests.rs | 22 +++++++++++++++++++ .../src/security/devices/owner_tests.rs | 12 +++++----- 2 files changed, 28 insertions(+), 6 deletions(-) diff --git a/crates/openhuman-core/src/cron/scheduler_tests.rs b/crates/openhuman-core/src/cron/scheduler_tests.rs index dea7b60edd2..80184d947df 100644 --- a/crates/openhuman-core/src/cron/scheduler_tests.rs +++ b/crates/openhuman-core/src/cron/scheduler_tests.rs @@ -313,3 +313,25 @@ async fn a_pipeline_reports_its_last_stage_rather_than_pipefail() { a failure and re-status every existing job with a pipeline in it: {output}" ); } + +#[tokio::test] +async fn an_agent_pass_polls_with_its_own_config_and_reports_health() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp).await; + let mut health = None; + tick_agent_scope(&config, &mut health).await; + assert_eq!(health, Some(true), "a successful poll reports healthy"); +} + +#[tokio::test] +async fn tick_agents_without_a_backend_leaves_the_health_tracker_alone() { + // The lib test binary installs no storage backend, so no agent is visited. + if crate::storage::installed().is_some() { + return; + } + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp).await; + let mut health = Some(false); + tick_agents(&config, &mut health).await; + assert_eq!(health, Some(false)); +} diff --git a/crates/openhuman-core/src/security/devices/owner_tests.rs b/crates/openhuman-core/src/security/devices/owner_tests.rs index 4476e9adf89..4a4575c2415 100644 --- a/crates/openhuman-core/src/security/devices/owner_tests.rs +++ b/crates/openhuman-core/src/security/devices/owner_tests.rs @@ -55,14 +55,13 @@ fn scoped_config(dir: &std::path::Path) -> crate::config::Config { async fn a_device_found_in_the_local_scope_belongs_to_local_and_is_remembered() { let tmp = tempfile::tempdir().unwrap(); let config = scoped_config(tmp.path()); - super::super::store::insert_device(&config, "owner-test-found", "label", "pk", "hash") - .unwrap(); + super::super::store::insert_device(&config, "owner-test-found", "label", "pk", "hash").unwrap(); let context = crate::core::runtime::CoreContext::for_test_with_config( crate::core::runtime::DomainSet::full(), config, ); - let owner = crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-found", None)) - .await; + let owner = + crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-found", None)).await; assert_eq!(owner, None); assert_eq!(cached("owner-test-found"), Some(None)); } @@ -74,8 +73,9 @@ async fn a_channel_no_scope_knows_is_local_and_not_remembered() { crate::core::runtime::DomainSet::full(), scoped_config(tmp.path()), ); - let owner = crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-missing", None)) - .await; + let owner = + crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-missing", None)) + .await; assert_eq!(owner, None); assert_eq!(cached("owner-test-missing"), None); } From 72260a27376fd09cdbc87d8ea4aa65aed9c87ae1 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 19:29:03 +0300 Subject: [PATCH 23/26] feat(cron): add cron job management and device owner tracking Adds cron job scheduling support and device owner tracking so scheduled tasks can be registered and their owning devices recorded. Run management and agent storage were extended to persist and query this new state. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/cron/scheduler.rs | 54 +++++++++++++--- .../src/flows/ops/run_management.rs | 6 +- .../src/security/devices/bus.rs | 7 ++- .../src/security/devices/owner.rs | 62 ++++++++++++------- crates/openhuman-core/src/storage/agents.rs | 12 +++- 5 files changed, 107 insertions(+), 34 deletions(-) diff --git a/crates/openhuman-core/src/cron/scheduler.rs b/crates/openhuman-core/src/cron/scheduler.rs index ba27d2fa91a..ad5f27e2952 100644 --- a/crates/openhuman-core/src/cron/scheduler.rs +++ b/crates/openhuman-core/src/cron/scheduler.rs @@ -75,7 +75,7 @@ pub async fn run(config: Config) -> Result<()> { loop { interval.tick().await; tick_once(&config, &security, &mut last_emitted_health).await; - tick_agents(&config, &mut last_emitted_health).await; + tick_agents(&mut last_emitted_health).await; } } @@ -87,13 +87,13 @@ pub async fn run(config: Config) -> Result<()> { /// pass so a failure or recovery reported from an agent pass is seen by the /// next one); each agent pass authorizes its jobs with a policy built from /// that agent's own configuration. -pub(crate) async fn tick_agents(config: &Config, last_emitted_health: &mut Option) { +pub(crate) async fn tick_agents(last_emitted_health: &mut Option) { let health = std::sync::Mutex::new(*last_emitted_health); crate::storage::agents::for_each_agent("cron", || async { let mut steady = *health .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); - tick_agent_scope(config, &mut steady).await; + tick_agent_scope(&mut steady).await; *health .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) = steady; @@ -106,15 +106,53 @@ pub(crate) async fn tick_agents(config: &Config, last_emitted_health: &mut Optio /// One agent's pass: the poll under the current (agent) context, with that /// agent's configuration when it has one and a policy built from it. -async fn tick_agent_scope(config: &Config, last_emitted_health: &mut Option) { - let config = crate::core::runtime::CoreContext::current_embedder_config() - .unwrap_or_else(|| config.clone()); - let security = Arc::new(SecurityPolicy::from_config( +/// +/// Fails closed: an agent whose context carries no configuration of its own +/// (one known only from the backend's record of it, not yet re-created by the +/// host) is skipped rather than run under another agent's workspace, settings +/// and autonomy policy. Its jobs wait for the host to derive it again. +async fn tick_agent_scope(last_emitted_health: &mut Option) { + use crate::core::runtime::CoreContext; + let (Some(context), Some(config)) = + (CoreContext::current(), CoreContext::current_embedder_config()) + else { + tracing::debug!( + "[cron:scheduler] skipping an agent scope: its context has no configuration of its own" + ); + return; + }; + let security = agent_policy(&context, &config); + tick_once(&config, &security, last_emitted_health).await; +} + +/// The agent's security policy, kept across ticks so its rolling action +/// budget is not reset by every poll. Rebuilt when the host re-derives the +/// agent's context (a new context may carry new settings). +fn agent_policy( + context: &Arc, + config: &Config, +) -> Arc { + use std::collections::HashMap; + use std::sync::{LazyLock, Mutex}; + static POLICIES: LazyLock)>>> = + LazyLock::new(Default::default); + let identity = Arc::as_ptr(context) as usize; + let agent = context.session_agent().unwrap_or_default().to_string(); + let mut policies = POLICIES + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if let Some((built_for, policy)) = policies.get(&agent) { + if *built_for == identity { + return Arc::clone(policy); + } + } + let policy = Arc::new(SecurityPolicy::from_config( &config.autonomy, &config.workspace_dir, &config.action_dir, )); - tick_once(&config, &security, last_emitted_health).await; + policies.insert(agent, (identity, Arc::clone(&policy))); + policy } /// Single poll cycle of the scheduler loop, extracted so tests can drive diff --git a/crates/openhuman-core/src/flows/ops/run_management.rs b/crates/openhuman-core/src/flows/ops/run_management.rs index 9c8a37592bd..fc75f018e4b 100644 --- a/crates/openhuman-core/src/flows/ops/run_management.rs +++ b/crates/openhuman-core/src/flows/ops/run_management.rs @@ -277,8 +277,12 @@ pub async fn sweep_orphaned_running_runs_on_boot(config: &Config) -> usize { // boot floor can belong to a run another replica is still driving, and // sweeping it would drop that run's checkpoint — so the agent scopes are // left alone there, as the agent run reaper does. - if crate::storage::installed_is_shared() && !crate::core::runtime::mode::is_saas() { + if crate::storage::installed_is_shared() { tracing::info!(target: "flows", "[flows] boot sweep: agent scopes skipped, the storage backend is shared"); + if crate::core::runtime::mode::is_saas() { + // No `local` scope to fall back on: every row belongs to an agent. + return 0; + } return sweep_orphaned_running_runs_in_scope(config).await; } crate::storage::agents::for_each_scope("flows boot sweep", || { diff --git a/crates/openhuman-core/src/security/devices/bus.rs b/crates/openhuman-core/src/security/devices/bus.rs index 2b1a829b5b8..076d426478d 100644 --- a/crates/openhuman-core/src/security/devices/bus.rs +++ b/crates/openhuman-core/src/security/devices/bus.rs @@ -92,7 +92,12 @@ impl EventHandler for DeviceTunnelSubscriber { .unwrap_or_else(std::sync::PoisonError::into_inner) .get(channel_id.as_str()) .cloned(); - let owner = super::owner::owner_of(channel_id, pending.as_ref()).await; + let Ok(owner) = super::owner::owner_of(channel_id, pending.as_ref()).await else { + log::warn!( + "[devices/bus] dropping tunnel frame channel_id={channel_id}: owner unresolved" + ); + return; + }; crate::storage::agents::within_agent( owner.as_deref(), handle_tunnel_frame(channel_id, payload_b64), diff --git a/crates/openhuman-core/src/security/devices/owner.rs b/crates/openhuman-core/src/security/devices/owner.rs index 54fc515b68e..bcb1ad59660 100644 --- a/crates/openhuman-core/src/security/devices/owner.rs +++ b/crates/openhuman-core/src/security/devices/owner.rs @@ -40,40 +40,60 @@ fn cached(channel_id: &str) -> Option> { /// The agent `channel_id` belongs to, resolved as described above. A channel /// no scope knows yet (a handshake still in flight) is `local`. -pub(super) async fn owner_of(channel_id: &str, pending: Option<&PairingSession>) -> Option { +/// +/// # Errors +/// +/// [`OwnerLookupFailed`] when no scope claimed the channel and at least one +/// could not be searched (its configuration would not load): the device may +/// well belong to that scope, so the caller must refuse the frame instead of +/// running it as `local`. +pub(super) async fn owner_of( + channel_id: &str, + pending: Option<&PairingSession>, +) -> Result, OwnerLookupFailed> { if let Some(session) = pending { - return session.agent.clone(); + return Ok(session.agent.clone()); } if let Some(owner) = cached(channel_id) { - return owner; + return Ok(owner); } // The configuration is loaded inside each scope: in SaaS mode loading it // needs an acting agent, which this tunnel task does not have. let found = crate::storage::agents::for_each_scope("device owner", || async { let Ok(config) = crate::config::rpc::load_config_with_timeout().await else { - return false; + return None; }; - super::store::get_device(&config, channel_id) - .ok() - .flatten() - .is_some() + Some( + super::store::get_device(&config, channel_id) + .ok() + .flatten() + .is_some(), + ) }) - .await - .into_iter() - .find_map(|(agent, has_device)| has_device.then_some(agent)); - match found { - Some(owner) => { - log::debug!( - "[devices/owner] channel_id={channel_id} belongs to agent={}", - owner.as_deref().unwrap_or("local") - ); - remember(channel_id, owner.clone()); - owner - } - None => None, + .await; + if let Some(owner) = found + .iter() + .find_map(|(agent, has_device)| (*has_device == Some(true)).then(|| agent.clone())) + { + log::debug!( + "[devices/owner] channel_id={channel_id} belongs to agent={}", + owner.as_deref().unwrap_or("local") + ); + remember(channel_id, owner.clone()); + return Ok(owner); } + if found.iter().any(|(_, has_device)| has_device.is_none()) { + log::warn!("[devices/owner] channel_id={channel_id} could not be searched in every scope"); + return Err(OwnerLookupFailed); + } + Ok(None) } +/// A scope's configuration would not load, so the owner of a channel could not +/// be ruled out there. +#[derive(Debug, PartialEq, Eq)] +pub(super) struct OwnerLookupFailed; + #[cfg(test)] #[path = "owner_tests.rs"] mod tests; diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs index 53f430b0709..30ec2e50fce 100644 --- a/crates/openhuman-core/src/storage/agents.rs +++ b/crates/openhuman-core/src/storage/agents.rs @@ -54,9 +54,15 @@ pub fn registered(context: Arc) -> Arc { let mut live = LIVE .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); - let contexts = live.entry(agent.to_string()).or_default(); - contexts.retain(|existing| existing.strong_count() > 0); - contexts.push(Arc::downgrade(&context)); + // Agents come and go; forget the ones whose contexts are all gone so + // the registry stays as small as the set of live agents. + live.retain(|_, entries| { + entries.retain(|entry| entry.strong_count() > 0); + !entries.is_empty() + }); + live.entry(agent.to_string()) + .or_default() + .push(Arc::downgrade(&context)); drop(live); record(agent); } From 60f87f7b7560a2588cd157e0fdca58ebefbc0832 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 19:29:32 +0300 Subject: [PATCH 24/26] test(cron): cover agent scope config and policy reuse Add tests for the agent scheduler scope, covering the case where an agent has no configuration of its own and is skipped, and verifying that a context keeps a single policy across ticks until it is re-derived. The owner and storage tests were updated to match the new fallible owner lookup and to assert that registering drops agents whose contexts are gone. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/cron/scheduler.rs | 7 +-- .../src/cron/scheduler_tests.rs | 43 +++++++++++++++++-- .../src/security/devices/owner_tests.rs | 14 +++--- .../src/storage/agents_tests.rs | 7 +++ 4 files changed, 59 insertions(+), 12 deletions(-) diff --git a/crates/openhuman-core/src/cron/scheduler.rs b/crates/openhuman-core/src/cron/scheduler.rs index ad5f27e2952..77ae699654a 100644 --- a/crates/openhuman-core/src/cron/scheduler.rs +++ b/crates/openhuman-core/src/cron/scheduler.rs @@ -113,9 +113,10 @@ pub(crate) async fn tick_agents(last_emitted_health: &mut Option) { /// and autonomy policy. Its jobs wait for the host to derive it again. async fn tick_agent_scope(last_emitted_health: &mut Option) { use crate::core::runtime::CoreContext; - let (Some(context), Some(config)) = - (CoreContext::current(), CoreContext::current_embedder_config()) - else { + let (Some(context), Some(config)) = ( + CoreContext::current(), + CoreContext::current_embedder_config(), + ) else { tracing::debug!( "[cron:scheduler] skipping an agent scope: its context has no configuration of its own" ); diff --git a/crates/openhuman-core/src/cron/scheduler_tests.rs b/crates/openhuman-core/src/cron/scheduler_tests.rs index 80184d947df..d062c3c9855 100644 --- a/crates/openhuman-core/src/cron/scheduler_tests.rs +++ b/crates/openhuman-core/src/cron/scheduler_tests.rs @@ -314,24 +314,59 @@ async fn a_pipeline_reports_its_last_stage_rather_than_pipefail() { ); } +fn agent_context_with_config( + agent: &str, + config: Config, +) -> Arc { + use crate::core::runtime::{ContextOverlay, CoreContext, DomainSet}; + CoreContext::for_test(DomainSet::full(), None).derive_with( + ContextOverlay::new(config, DomainSet::full(), Default::default()).session_agent(agent), + ) +} + #[tokio::test] async fn an_agent_pass_polls_with_its_own_config_and_reports_health() { let tmp = TempDir::new().unwrap(); let config = test_config(&tmp).await; + let context = agent_context_with_config("cron-test-agent", config); let mut health = None; - tick_agent_scope(&config, &mut health).await; + crate::core::runtime::CoreContext::scope(context, tick_agent_scope(&mut health)).await; assert_eq!(health, Some(true), "a successful poll reports healthy"); } +#[tokio::test] +async fn an_agent_pass_without_its_own_config_is_skipped() { + use crate::core::runtime::{CoreContext, DomainSet}; + // `for_agent` swaps only the identity: no configuration of the agent's own. + let fallback = CoreContext::for_test(DomainSet::full(), None); + let mut health = None; + CoreContext::scope( + fallback.for_agent("cron-test-recorded"), + tick_agent_scope(&mut health), + ) + .await; + assert_eq!(health, None, "nothing ran, so nothing was reported"); +} + +#[tokio::test] +async fn an_agent_keeps_its_policy_across_ticks_until_re_derived() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp).await; + let first = agent_context_with_config("cron-test-policy", config.clone()); + let a = agent_policy(&first, &config); + let b = agent_policy(&first, &config); + assert!(Arc::ptr_eq(&a, &b), "the same context keeps one policy"); + let second = agent_context_with_config("cron-test-policy", config.clone()); + assert!(!Arc::ptr_eq(&a, &agent_policy(&second, &config))); +} + #[tokio::test] async fn tick_agents_without_a_backend_leaves_the_health_tracker_alone() { // The lib test binary installs no storage backend, so no agent is visited. if crate::storage::installed().is_some() { return; } - let tmp = TempDir::new().unwrap(); - let config = test_config(&tmp).await; let mut health = Some(false); - tick_agents(&config, &mut health).await; + tick_agents(&mut health).await; assert_eq!(health, Some(false)); } diff --git a/crates/openhuman-core/src/security/devices/owner_tests.rs b/crates/openhuman-core/src/security/devices/owner_tests.rs index 4a4575c2415..5c5df21dc61 100644 --- a/crates/openhuman-core/src/security/devices/owner_tests.rs +++ b/crates/openhuman-core/src/security/devices/owner_tests.rs @@ -17,22 +17,26 @@ async fn a_pending_pairing_names_its_agent() { assert_eq!( owner_of("owner-test-pending", Some(&pending)) .await + .unwrap() .as_deref(), Some("agent-7") ); let local = session(None); - assert_eq!(owner_of("owner-test-pending", Some(&local)).await, None); + assert_eq!(owner_of("owner-test-pending", Some(&local)).await, Ok(None)); } #[tokio::test] async fn a_remembered_owner_is_used_without_a_lookup() { remember("owner-test-cached", Some("agent-9".to_string())); assert_eq!( - owner_of("owner-test-cached", None).await.as_deref(), + owner_of("owner-test-cached", None) + .await + .unwrap() + .as_deref(), Some("agent-9") ); remember("owner-test-local", None); - assert_eq!(owner_of("owner-test-local", None).await, None); + assert_eq!(owner_of("owner-test-local", None).await, Ok(None)); } #[test] @@ -62,7 +66,7 @@ async fn a_device_found_in_the_local_scope_belongs_to_local_and_is_remembered() ); let owner = crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-found", None)).await; - assert_eq!(owner, None); + assert_eq!(owner, Ok(None)); assert_eq!(cached("owner-test-found"), Some(None)); } @@ -76,6 +80,6 @@ async fn a_channel_no_scope_knows_is_local_and_not_remembered() { let owner = crate::core::runtime::CoreContext::scope(context, owner_of("owner-test-missing", None)) .await; - assert_eq!(owner, None); + assert_eq!(owner, Ok(None)); assert_eq!(cached("owner-test-missing"), None); } diff --git a/crates/openhuman-core/src/storage/agents_tests.rs b/crates/openhuman-core/src/storage/agents_tests.rs index 83c18ffddeb..85effe50789 100644 --- a/crates/openhuman-core/src/storage/agents_tests.rs +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -128,3 +128,10 @@ fn without_a_fallback_only_live_contexts_are_visited() { .iter() .any(|(id, _)| id == "agents-test-no-fallback")); } + +#[test] +fn registering_forgets_agents_whose_contexts_are_all_gone() { + drop(agent_context("agents-test-gone")); + let _other = agent_context("agents-test-other-live"); + assert!(!LIVE.lock().unwrap().contains_key("agents-test-gone")); +} From f1d2da97cda40b3d7df95a9878accd57adb51654 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 19:32:40 +0300 Subject: [PATCH 25/26] test(storage): wait for agent contexts to be released Agent registry tests asserted immediately that a dropped context was gone, which could flake when a concurrent walk of the registry still held a reference for an instant. The assertions now poll briefly for the agent to disappear before failing. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../openhuman-core/src/storage/agents_tests.rs | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/crates/openhuman-core/src/storage/agents_tests.rs b/crates/openhuman-core/src/storage/agents_tests.rs index 85effe50789..52710cce63a 100644 --- a/crates/openhuman-core/src/storage/agents_tests.rs +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -12,6 +12,18 @@ fn agent_context(agent: &str) -> Arc { ) } +/// Whether `agent` stops being live. Another test's concurrent walk of the +/// registry can hold a context for an instant, so allow it to let go. +fn eventually_gone(agent: &str) -> bool { + (0..100).any(|_| { + let gone = !live_agents().contains(&agent.to_string()); + if !gone { + std::thread::sleep(std::time::Duration::from_millis(10)); + } + gone + }) +} + fn live_agents() -> Vec { agent_contexts(None).into_iter().map(|(id, _)| id).collect() } @@ -26,7 +38,7 @@ fn deriving_an_agent_context_registers_it_until_dropped() { .unwrap(); assert!(Arc::ptr_eq(&found, &context), "the live context is used"); drop((found, context)); - assert!(!live_agents().contains(&"agents-test-live".to_string())); + assert!(eventually_gone("agents-test-live")); } #[test] @@ -72,7 +84,7 @@ fn a_dropped_sibling_context_does_not_hide_a_live_one() { assert!(Arc::ptr_eq(&found, &first)); assert!(context_for("agents-test-siblings").is_some()); drop((found, first)); - assert!(!live_agents().contains(&"agents-test-siblings".to_string())); + assert!(eventually_gone("agents-test-siblings")); } #[test] From 110195c5de66b949cbf4255cd70fb8cf94d95b72 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 19:34:20 +0300 Subject: [PATCH 26/26] chore: files changed crates/openhuman-core/src/storage/agents_tests.rs Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/openhuman-core/src/storage/agents_tests.rs | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/crates/openhuman-core/src/storage/agents_tests.rs b/crates/openhuman-core/src/storage/agents_tests.rs index 52710cce63a..7875a3ec084 100644 --- a/crates/openhuman-core/src/storage/agents_tests.rs +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -144,6 +144,15 @@ fn without_a_fallback_only_live_contexts_are_visited() { #[test] fn registering_forgets_agents_whose_contexts_are_all_gone() { drop(agent_context("agents-test-gone")); - let _other = agent_context("agents-test-other-live"); - assert!(!LIVE.lock().unwrap().contains_key("agents-test-gone")); + // Each registration prunes; a concurrent walk of the registry may hold the + // dropped context for an instant, so give it a few registrations. + let forgotten = (0..100).any(|_| { + let _other = agent_context("agents-test-other-live"); + let gone = !LIVE.lock().unwrap().contains_key("agents-test-gone"); + if !gone { + std::thread::sleep(std::time::Duration::from_millis(10)); + } + gone + }); + assert!(forgotten); }