Repository navigation
feat(storage): background work visits every agent's storage scope #7204
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
senamakel
merged 28 commits into
tinyhumansai:main
from
senamakel:storage-scope-propagation
Oct 9, 2026
Merged
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 06b8343
refactor(storage): extract agent storage helpers
senamakel 061df35
feat(storage): register agents module and record live agents on install
senamakel 28d3875
revert: restore the tinyagents pin #7197 moved back
senamakel 144cc45
Merge branch 'restore-tinyagents-pin' into storage-scope-propagation
senamakel ca3d663
feat(cron): run scheduled jobs and task polls per agent scope
senamakel ff3e176
style: apply rustfmt formatting to storage and task source modules
senamakel d7d7d37
fix(flows): sweep and reconcile runs across all agent storage scopes
senamakel fbfa3b4
fix(agent): reap orphaned runs across every agent scope
senamakel f475d87
feat(devices): track pairing agent and add agent-scoped context helpers
senamakel e756dc6
feat(devices): scope tunnel frames to the owning agent
senamakel 35a48fb
style(security): reformat owner device module
senamakel 30a2578
test(cli): register storage scope e2e test binary
senamakel cf8c019
test(storage): reorder imports in storage scope e2e test
senamakel 4266afd
test(storage): drop spawn_blocking from scope e2e test
senamakel 6e6b901
chore(ci): refresh saas ambient baseline and lockfile
senamakel 4c7e9b1
docs(storage): document background work and agent scopes
senamakel 24a7330
chore: 姫The diff is empty — no changes are shown for any of the liste…
senamakel b980fde
chore(flows): use tracing for boot sweep log
senamakel 5ef831c
test: cover per-agent poll scoping and agent context lifecycle
senamakel 34d032b
test(storage): reformat recorded state setup in agents test
senamakel 73ae76f
refactor(storage): make agent backend selection injectable
senamakel 17091db
test(cron): cover agent tick health reporting
senamakel 72260a2
feat(cron): add cron job management and device owner tracking
senamakel 60f87f7
test(cron): cover agent scope config and policy reuse
senamakel f1d2da9
test(storage): wait for agent contexts to be released
senamakel 110195c
chore: files changed crates/openhuman-core/src/storage/agents_tests.rs
senamakel 481497c
Merge upstream/main into storage-scope-propagation
senamakel File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
31 changes: 31 additions & 0 deletions
31
crates/openhuman-core/src/core/runtime/context_for_agent.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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> { | ||
|
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()), | ||
|
senamakel marked this conversation as resolved.
|
||
| agent: Default::default(), | ||
| }) | ||
| } | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -5,6 +5,7 @@ | |
|
|
||
| pub mod bus; | ||
| pub mod crypto; | ||
| mod owner; | ||
| pub mod rpc; | ||
| pub mod schemas; | ||
| pub mod store; | ||
|
|
||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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); | ||
|
senamakel marked this conversation as resolved.
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 { | ||
|
senamakel marked this conversation as resolved.
|
||
| return None; | ||
| }; | ||
| Some( | ||
| super::store::get_device(&config, channel_id) | ||
| .ok() | ||
| .flatten() | ||
| .is_some(), | ||
| ) | ||
|
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; | ||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.