diff --git a/crates/openhuman-cli/Cargo.toml b/crates/openhuman-cli/Cargo.toml index 490ee136a6a..4a67e02a8f6 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/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/core/runtime/context.rs b/crates/openhuman-core/src/core/runtime/context.rs index 4bf3711273c..b29a51dc69d 100644 --- a/crates/openhuman-core/src/core/runtime/context.rs +++ b/crates/openhuman-core/src/core/runtime/context.rs @@ -372,7 +372,7 @@ impl CoreContext { } }; let agent = self.agent.derive(&mut overlay); - Arc::new(CoreContext { + crate::storage::agents::registered(Arc::new(CoreContext { host_kind: self.host_kind, workspace_binding: RwLock::new(shared_binding), domains, @@ -383,7 +383,7 @@ impl CoreContext { turn_origin: self.turn_origin.clone(), session_agent: overlay.session_agent.or_else(|| self.session_agent.clone()), agent, - }) + })) } /// The agent a host session store scopes work under this context to, if @@ -689,6 +689,8 @@ pub async fn init_stores(cfg: &crate::config::Config, domains: crate::core::runt } } +#[path = "context_for_agent.rs"] +mod for_agent; #[path = "context_turn_origin.rs"] mod turn_origin_scope; diff --git a/crates/openhuman-core/src/core/runtime/context_for_agent.rs b/crates/openhuman-core/src/core/runtime/context_for_agent.rs new file mode 100644 index 00000000000..8f07ab2f7db --- /dev/null +++ b/crates/openhuman-core/src/core/runtime/context_for_agent.rs @@ -0,0 +1,31 @@ +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()), + agent: Default::default(), + }) + } +} diff --git a/crates/openhuman-core/src/flows/ops/run_management.rs b/crates/openhuman-core/src/flows/ops/run_management.rs index 10be3a787ca..fc75f018e4b 100644 --- a/crates/openhuman-core/src/flows/ops/run_management.rs +++ b/crates/openhuman-core/src/flows/ops/run_management.rs @@ -270,6 +270,32 @@ 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`). + // + // 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() { + 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", || { + 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..2fc64bdc2ec 100644 --- a/crates/openhuman-core/src/flows/ops/triggers.rs +++ b/crates/openhuman-core/src/flows/ops/triggers.rs @@ -244,6 +244,29 @@ 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 diff --git a/crates/openhuman-core/src/integrations/task_sources/periodic.rs b/crates/openhuman-core/src/integrations/task_sources/periodic.rs index faac8d2e004..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 } @@ -70,7 +81,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 +98,14 @@ 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/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/security/devices/bus.rs b/crates/openhuman-core/src/security/devices/bus.rs index 979c0747c95..076d426478d 100644 --- a/crates/openhuman-core/src/security/devices/bus.rs +++ b/crates/openhuman-core/src/security/devices/bus.rs @@ -84,7 +84,25 @@ 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 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), + ) + .await; } _ => {} } @@ -361,6 +379,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..bcb1ad59660 --- /dev/null +++ b/crates/openhuman-core/src/security/devices/owner.rs @@ -0,0 +1,99 @@ +//! 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`. +/// +/// # 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 Ok(session.agent.clone()); + } + if let Some(owner) = cached(channel_id) { + 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 None; + }; + Some( + super::store::get_device(&config, channel_id) + .ok() + .flatten() + .is_some(), + ) + }) + .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/security/devices/owner_tests.rs b/crates/openhuman-core/src/security/devices/owner_tests.rs new file mode 100644 index 00000000000..5c5df21dc61 --- /dev/null +++ b/crates/openhuman-core/src/security/devices/owner_tests.rs @@ -0,0 +1,85 @@ +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 + .unwrap() + .as_deref(), + Some("agent-7") + ); + let local = session(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 + .unwrap() + .as_deref(), + Some("agent-9") + ); + remember("owner-test-local", None); + assert_eq!(owner_of("owner-test-local", None).await, Ok(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"); +} + +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, Ok(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, Ok(None)); + assert_eq!(cached("owner-test-missing"), None); +} 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/README.md b/crates/openhuman-core/src/storage/README.md index 14e34a43d48..d07737d2564 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_live_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'; diff --git a/crates/openhuman-core/src/storage/agents.rs b/crates/openhuman-core/src/storage/agents.rs new file mode 100644 index 00000000000..30ec2e50fce --- /dev/null +++ b/crates/openhuman-core/src/storage/agents.rs @@ -0,0 +1,263 @@ +//! 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. +/// +/// 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); + +/// 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() { + let mut live = LIVE + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + // 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); + } + 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() { + 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; + }; + 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) + .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 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 + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + 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); + } + !entries.is_empty() + }); + } + if let (Some(backend), Some(fallback)) = (backend, fallback) { + for agent in recorded(backend) { + contexts + .entry(agent.clone()) + .or_insert_with(|| fallback.for_agent(&agent)); + } + } + contexts.into_iter().collect() +} + +/// The context to act for `agent` under: its live context when one exists, +/// else the current context acting for it (`CoreContext::for_agent`). +pub fn context_for(agent: &str) -> Option> { + let live = LIVE + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .get(agent) + .and_then(|entries| entries.iter().rev().find_map(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 +/// 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(); + 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 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((agent, value)); + } + results +} + +#[cfg(test)] +#[path = "agents_tests.rs"] +mod tests; 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..7875a3ec084 --- /dev/null +++ b/crates/openhuman-core/src/storage/agents_tests.rs @@ -0,0 +1,158 @@ +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), + ) +} + +/// 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() +} + +#[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!(eventually_gone("agents-test-live")); +} + +#[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()); +} + +#[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!(eventually_gone("agents-test-siblings")); +} + +#[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")); +} + +fn memory_backend() -> 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")); +} + +#[test] +fn registering_forgets_agents_whose_contexts_are_all_gone() { + drop(agent_context("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); +} diff --git a/crates/openhuman-core/src/storage/mod.rs b/crates/openhuman-core/src/storage/mod.rs index 1ff9cf94e89..5d67a8f2e15 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,12 @@ 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); + // What was recorded described the previous backend. + agents::reset_recorded(); + // Agents derived before the backend existed still need recording. + agents::record_live(); + previous } /// The installed backend, when the host configured one. @@ -112,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() } diff --git a/scripts/ci/saas-ambient-baseline.json b/scripts/ci/saas-ambient-baseline.json index fe8978e3a3b..3ef10d6d6a8 100644 --- a/scripts/ci/saas-ambient-baseline.json +++ b/scripts/ci/saas-ambient-baseline.json @@ -657,13 +657,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", diff --git a/tests/storage_scope_e2e.rs b/tests/storage_scope_e2e.rs new file mode 100644 index 00000000000..a2219a0e4fa --- /dev/null +++ b/tests/storage_scope_e2e.rs @@ -0,0 +1,83 @@ +//! 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::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) + .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"), + ); + // 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. + assert!(job_names(&config).is_empty()); + + // Visited through the live agent context … + 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:?}" + ); + + // … and, once the agent is gone, through the id the backend recorded. + drop(agent); + 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:?}" + ); + assert!(recorded.contains(&(None, Vec::new())), "{recorded:?}"); +}