Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
5c18808
refactor(runtime): register derived contexts with the agent registry
senamakel Oct 9, 2026
06b8343
refactor(storage): extract agent storage helpers
senamakel Oct 9, 2026
061df35
feat(storage): register agents module and record live agents on install
senamakel Oct 9, 2026
28d3875
revert: restore the tinyagents pin #7197 moved back
senamakel Oct 9, 2026
144cc45
Merge branch 'restore-tinyagents-pin' into storage-scope-propagation
senamakel Oct 9, 2026
ca3d663
feat(cron): run scheduled jobs and task polls per agent scope
senamakel Oct 9, 2026
ff3e176
style: apply rustfmt formatting to storage and task source modules
senamakel Oct 9, 2026
d7d7d37
fix(flows): sweep and reconcile runs across all agent storage scopes
senamakel Oct 9, 2026
fbfa3b4
fix(agent): reap orphaned runs across every agent scope
senamakel Oct 9, 2026
f475d87
feat(devices): track pairing agent and add agent-scoped context helpers
senamakel Oct 9, 2026
e756dc6
feat(devices): scope tunnel frames to the owning agent
senamakel Oct 9, 2026
35a48fb
style(security): reformat owner device module
senamakel Oct 9, 2026
30a2578
test(cli): register storage scope e2e test binary
senamakel Oct 9, 2026
cf8c019
test(storage): reorder imports in storage scope e2e test
senamakel Oct 9, 2026
4266afd
test(storage): drop spawn_blocking from scope e2e test
senamakel Oct 9, 2026
6e6b901
chore(ci): refresh saas ambient baseline and lockfile
senamakel Oct 9, 2026
4c7e9b1
docs(storage): document background work and agent scopes
senamakel Oct 9, 2026
24a7330
chore: 姫The diff is empty — no changes are shown for any of the liste…
senamakel Oct 9, 2026
b980fde
chore(flows): use tracing for boot sweep log
senamakel Oct 9, 2026
5ef831c
test: cover per-agent poll scoping and agent context lifecycle
senamakel Oct 9, 2026
34d032b
test(storage): reformat recorded state setup in agents test
senamakel Oct 9, 2026
73ae76f
refactor(storage): make agent backend selection injectable
senamakel Oct 9, 2026
17091db
test(cron): cover agent tick health reporting
senamakel Oct 9, 2026
72260a2
feat(cron): add cron job management and device owner tracking
senamakel Oct 9, 2026
60f87f7
test(cron): cover agent scope config and policy reuse
senamakel Oct 9, 2026
f1d2da9
test(storage): wait for agent contexts to be released
senamakel Oct 9, 2026
110195c
chore: files changed crates/openhuman-core/src/storage/agents_tests.rs
senamakel Oct 9, 2026
481497c
Merge upstream/main into storage-scope-propagation
senamakel Oct 9, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions crates/openhuman-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
senamakel marked this conversation as resolved.
# 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"
Expand Down
10 changes: 9 additions & 1 deletion crates/openhuman-core/src/agent/tinyagents/reaper.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
6 changes: 4 additions & 2 deletions crates/openhuman-core/src/core/runtime/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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;

Expand Down
31 changes: 31 additions & 0 deletions crates/openhuman-core/src/core/runtime/context_for_agent.rs
Original file line number Diff line number Diff line change
@@ -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<Self>, agent: &str) -> Arc<CoreContext> {
Comment thread
senamakel marked this conversation as resolved.
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()),
Comment thread
senamakel marked this conversation as resolved.
agent: Default::default(),
})
}
}
26 changes: 26 additions & 0 deletions crates/openhuman-core/src/flows/ops/run_management.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Comment thread
senamakel marked this conversation as resolved.
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", || {
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.
sweep_orphaned_running_runs_in_scope(config)
Comment thread
senamakel marked this conversation as resolved.
})
.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.";
Expand Down
23 changes: 23 additions & 0 deletions crates/openhuman-core/src/flows/ops/triggers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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", || {
Comment thread
senamakel marked this conversation as resolved.
reconcile_schedule_triggers_in_scope(config)
})
Comment thread
senamakel marked this conversation as resolved.
.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
Expand Down
27 changes: 22 additions & 5 deletions crates/openhuman-core/src/integrations/task_sources/periodic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Comment thread
senamakel marked this conversation as resolved.
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());
}
}

Expand All @@ -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
}
Expand All @@ -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 {
Comment thread
senamakel marked this conversation as resolved.
tracing::info!(
tick_seconds = TICK_SECONDS,
"[task_sources:periodic] scheduler starting"
Expand All @@ -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
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.
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)");
}
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Comment thread
senamakel marked this conversation as resolved.
CoreContext::scope(agent("ts-agent-a"), async { record_poll(&s.id) }).await;
Comment thread
senamakel marked this conversation as resolved.
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");
}
25 changes: 24 additions & 1 deletion crates/openhuman-core/src/security/devices/bus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,25 @@ impl EventHandler<DomainEvent> 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(
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.
owner.as_deref(),
handle_tunnel_frame(channel_id, payload_b64),
)
.await;
}
_ => {}
}
Expand Down Expand Up @@ -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,
Expand Down
1 change: 1 addition & 0 deletions crates/openhuman-core/src/security/devices/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

pub mod bus;
pub mod crypto;
mod owner;
pub mod rpc;
pub mod schemas;
pub mod store;
Expand Down
99 changes: 99 additions & 0 deletions crates/openhuman-core/src/security/devices/owner.rs
Original file line number Diff line number Diff line change
@@ -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<Mutex<HashMap<String, Option<String>>>> = LazyLock::new(Default::default);
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.

/// Records that `channel_id` belongs to `agent` (`None` = `local`).
pub(super) fn remember(channel_id: &str, agent: Option<String>) {
OWNERS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(channel_id.to_string(), agent);
}

fn cached(channel_id: &str) -> Option<Option<String>> {
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<Option<String>, 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 {
Comment thread
senamakel marked this conversation as resolved.
return None;
};
Some(
super::store::get_device(&config, channel_id)
.ok()
.flatten()
.is_some(),
)
Comment thread
senamakel marked this conversation as resolved.
})
.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;
Loading
Loading