feat(memory-sync): wake Composio periodic loop on resume + gate on per-source toggle (#2831) (#3417)

Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
sanil-23
2026-06-07 05:41:06 -04:00
committed by GitHub
co-authored by Claude
parent d7c4f35155
commit 4c6fa39afd
2 changed files with 374 additions and 18 deletions
+238 -16
View File
@@ -52,7 +52,10 @@ use tokio::time::interval;
use crate::openhuman::config::rpc as config_rpc;
use crate::openhuman::config::DEFAULT_MEMORY_SYNC_INTERVAL_SECS;
use crate::openhuman::scheduler_gate::gate::current_policy;
use crate::openhuman::memory_sources::{
memory_sync_defaults_for_toolkit, MemorySourceEntry, SourceKind,
};
use crate::openhuman::scheduler_gate::gate::{current_policy, resume_notify};
use crate::openhuman::scheduler_gate::policy::PauseReason;
use super::providers::{get_provider, ComposioUsage, ProviderContext, SyncReason};
@@ -192,6 +195,57 @@ fn persisted_since_last_sync(
})
}
/// Outcome of consulting the per-source registry for one Composio connection
/// during a periodic tick (#2831).
#[derive(Debug, PartialEq, Eq)]
enum PeriodicSourceDecision {
/// The user toggled this source **off** — skip background sync entirely.
/// Manual `memory_sources_sync` still works (it has its own `enabled`
/// guard); only the automatic loop honours this here.
Skip,
/// Sync this connection with the given caps (`None` = uncapped for that
/// dimension).
Sync {
max_items: Option<u32>,
sync_depth_days: Option<u32>,
},
}
/// Decide whether — and with what caps — to periodically sync one connection,
/// honouring the per-source `enabled` toggle (#2831). Pure so the three
/// branches can be unit-tested without async/registry I/O.
///
/// - **disabled row** → [`PeriodicSourceDecision::Skip`]: the background loop
/// must not sync a source the user switched off (this was the row-2 leak —
/// previously the loop only read the registry for caps and synced disabled
/// sources anyway, uncapped).
/// - **enabled row** → `Sync` with the row's caps.
/// - **no row yet** → `Sync` with conservative per-toolkit defaults
/// ([`memory_sync_defaults_for_toolkit`]). `reconcile` normally backfills an
/// enabled, capped row for every connection; this covers the brief
/// pre-reconcile window. The no-match path defaults to *sync-bounded*, never
/// *skip* and never *uncapped*, so a missing/mismatched row degrades safely
/// (data keeps flowing, just capped) instead of silently going dark.
fn decide_periodic_source(
source: Option<&MemorySourceEntry>,
toolkit: &str,
) -> PeriodicSourceDecision {
match source {
Some(s) if !s.enabled => PeriodicSourceDecision::Skip,
Some(s) => PeriodicSourceDecision::Sync {
max_items: s.max_items,
sync_depth_days: s.sync_depth_days,
},
None => {
let (max_items, sync_depth_days) = memory_sync_defaults_for_toolkit(toolkit);
PeriodicSourceDecision::Sync {
max_items,
sync_depth_days,
}
}
}
}
/// Spawn the periodic sync background task. Idempotent: only the
/// first call actually spawns the loop, every subsequent call is a
/// cheap no-op (logged at `debug` so it's visible during startup
@@ -221,14 +275,35 @@ pub fn start_periodic_sync() {
/// Inner loop, broken out so it's easy to mock-replace in tests if we
/// ever want to drive ticks deterministically.
///
/// Each iteration waits on whichever comes first (#2831):
/// * the 20-min `ticker` — the steady-state cadence, or
/// * the scheduler-gate **resume** notify — fired when the user toggles
/// Memory Tree back on or signs back in.
///
/// On a resume wake we run a tick **immediately** (so sync restarts within
/// seconds, not at the next ≤20-min boundary) and `reset()` the ticker so the
/// *next* scheduled tick is a full `TICK_SECONDS` out. The reset is what stops
/// rapid off-on toggling from bunch-firing: many wakes collapse into at most
/// one extra tick (the `Notify` stores a single permit), and the cadence
/// re-bases from the last actual tick.
async fn run_loop() {
let mut ticker = interval(Duration::from_secs(TICK_SECONDS));
let resume = resume_notify();
// Skip the immediate-fire tick so startup isn't slammed before the
// user even has time to sign in.
ticker.tick().await;
loop {
ticker.tick().await;
tokio::select! {
_ = ticker.tick() => {}
_ = resume.notified() => {
// Woke early on a resume transition. Re-base the cadence so the
// next scheduled tick is TICK_SECONDS from now, then fall
// through and run the tick immediately.
ticker.reset();
}
}
if let Err(e) = run_one_tick().await {
tracing::warn!(
error = %e,
@@ -378,6 +453,28 @@ pub(crate) async fn run_one_tick() -> Result<(), String> {
let audit_index = index_last_success_by_connection(&read_audit_log(&config));
let now = Utc::now();
// Per-source registry snapshot (#2831). The periodic loop gates on the
// per-source `enabled` toggle so a source the user switched off stops
// syncing in the background — matching the manual paths
// (`memory_sources::sync_source`, `memory_sources_sync_all`), which already
// early-return on `!enabled`. Index every Composio source (enabled and
// disabled) by connection id; the per-connection branch below resolves
// skip/caps via `decide_periodic_source`.
//
// Built from the **already-loaded** `config` snapshot (Step 1), not a second
// `list_sources()` read. A separate read whose error we swallowed to an
// empty map would make every disabled source fall through to the
// `decide_periodic_source(None, ..)` default-caps path — silently
// re-enabling background sync for sources the user switched off on a
// transient config-read failure. Reusing the tick's snapshot is fail-closed
// (a disabled row stays disabled) and avoids the extra read entirely.
let composio_sources: HashMap<String, MemorySourceEntry> = config
.memory_sources
.iter()
.filter(|s| s.kind == SourceKind::Composio)
.filter_map(|s| s.connection_id.clone().map(|id| (id, s.clone())))
.collect();
let mut considered = 0usize;
let mut fired = 0usize;
for conn in resp.connections {
@@ -425,20 +522,25 @@ pub(crate) async fn run_one_tick() -> Result<(), String> {
continue;
}
// Look up per-source caps from the memory_sources registry.
// Non-fatal: if the lookup fails we proceed without caps.
let (src_max_items, src_sync_depth_days) = {
let registry_sources = crate::openhuman::memory_sources::list_enabled_by_kind(
crate::openhuman::memory_sources::SourceKind::Composio,
)
.await
.unwrap_or_default();
registry_sources
.iter()
.find(|s| s.connection_id.as_deref() == Some(conn.id.as_str()))
.map(|s| (s.max_items, s.sync_depth_days))
.unwrap_or((None, None))
};
// Per-source gate + caps from the memory_sources registry (#2831).
// A disabled source is skipped here (the background-sync half of the
// toggle); enabled sources sync with their caps; a connection with no
// registry row yet syncs with conservative per-toolkit defaults.
let (src_max_items, src_sync_depth_days) =
match decide_periodic_source(composio_sources.get(&conn.id), &toolkit) {
PeriodicSourceDecision::Skip => {
tracing::debug!(
toolkit = %toolkit,
connection_id = %conn.id,
"[composio:periodic] source disabled — skipping periodic sync"
);
continue;
}
PeriodicSourceDecision::Sync {
max_items,
sync_depth_days,
} => (max_items, sync_depth_days),
};
tracing::debug!(
toolkit = %toolkit,
@@ -633,6 +735,126 @@ mod tests {
));
}
/// Build a minimal Composio `MemorySourceEntry` for the per-source gate
/// tests — only the fields `decide_periodic_source` reads are meaningful.
fn composio_source(
enabled: bool,
max_items: Option<u32>,
sync_depth_days: Option<u32>,
) -> MemorySourceEntry {
MemorySourceEntry {
id: "src_test".to_string(),
kind: SourceKind::Composio,
label: "test".to_string(),
enabled,
toolkit: Some("gmail".to_string()),
connection_id: Some("cmp-1".to_string()),
path: None,
glob: None,
url: None,
branch: None,
paths: Vec::new(),
max_commits: None,
max_issues: None,
max_prs: None,
query: None,
since_days: None,
max_items,
selector: None,
max_tokens_per_sync: None,
max_cost_per_sync_usd: None,
sync_depth_days,
}
}
/// #2831 row 2: a source explicitly toggled **off** must be skipped by the
/// background loop — this is the leak the gate closes.
#[test]
fn decide_periodic_source_skips_disabled_source() {
let src = composio_source(false, Some(100), Some(30));
assert_eq!(
decide_periodic_source(Some(&src), "gmail"),
PeriodicSourceDecision::Skip
);
}
/// An enabled source syncs with exactly its configured caps (no defaulting).
#[test]
fn decide_periodic_source_uses_enabled_source_caps() {
let src = composio_source(true, Some(42), Some(7));
assert_eq!(
decide_periodic_source(Some(&src), "gmail"),
PeriodicSourceDecision::Sync {
max_items: Some(42),
sync_depth_days: Some(7),
}
);
}
/// A connection with no registry row yet (pre-reconcile window) syncs with
/// the conservative per-toolkit defaults — **bounded**, never uncapped, and
/// never skipped. This is the safe-direction fallback for a missing match.
#[test]
fn decide_periodic_source_defaults_caps_when_no_row() {
let (want_items, want_depth) = memory_sync_defaults_for_toolkit("gmail");
assert_eq!(
decide_periodic_source(None, "gmail"),
PeriodicSourceDecision::Sync {
max_items: want_items,
sync_depth_days: want_depth,
}
);
// The defaults are bounded for a known toolkit (regression guard against
// an accidental return to uncapped background fetches).
assert!(want_items.is_some());
}
/// Multi-account regression (#3443 added multiple account connections per
/// toolkit): two live connections of the *same* toolkit must be gated
/// **independently** by their own per-`connection_id` source rows. This
/// pins the loop's connection-id keying — re-keying the lookup by toolkit
/// would collapse the two accounts and is the regression this guards.
#[test]
fn per_connection_gate_is_independent_across_accounts_of_same_toolkit() {
// gmail account A: enabled with caps; gmail account B: disabled.
let mut a = composio_source(true, Some(10), Some(5));
a.connection_id = Some("conn-A".to_string());
a.toolkit = Some("gmail".to_string());
let mut b = composio_source(false, Some(99), Some(99));
b.connection_id = Some("conn-B".to_string());
b.toolkit = Some("gmail".to_string());
// Build the same connection_id → entry index the live tick builds.
let index: HashMap<String, MemorySourceEntry> = [a, b]
.into_iter()
.filter_map(|s| s.connection_id.clone().map(|id| (id, s)))
.collect();
// Account A (enabled) syncs with its own caps...
assert_eq!(
decide_periodic_source(index.get("conn-A"), "gmail"),
PeriodicSourceDecision::Sync {
max_items: Some(10),
sync_depth_days: Some(5),
}
);
// ...account B (disabled) is skipped, even though it shares the toolkit.
assert_eq!(
decide_periodic_source(index.get("conn-B"), "gmail"),
PeriodicSourceDecision::Skip
);
// A third, not-yet-registered account of the same toolkit falls back to
// bounded defaults (never skipped, never uncapped).
let (def_items, def_depth) = memory_sync_defaults_for_toolkit("gmail");
assert_eq!(
decide_periodic_source(index.get("conn-C"), "gmail"),
PeriodicSourceDecision::Sync {
max_items: def_items,
sync_depth_days: def_depth,
}
);
}
/// End-to-end simulation of the scheduler's per-connection decision: prove
/// that **changing the global setting changes when the next sync fires**
/// (issue #3302 acceptance criterion). We drive the same two pure helpers
+136 -2
View File
@@ -10,7 +10,7 @@ use std::sync::{Arc, OnceLock};
use std::time::Duration;
use parking_lot::RwLock;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore};
use crate::openhuman::config::{Config, SchedulerGateConfig};
use crate::openhuman::scheduler_gate::policy::{decide, PauseReason, Policy};
@@ -240,12 +240,50 @@ pub fn init_global(config: &Config) {
});
}
/// Process-wide resume signal (#2831). Fired whenever the gate transitions
/// **out of** a paused state — the user toggles Memory Tree back on
/// ([`update_config`]) or signs back in ([`set_signed_out`]). Background loops
/// (e.g. the Composio periodic scheduler) park on [`resume_notify`] so they can
/// resume work within seconds instead of waiting out their next tick boundary.
static RESUME_NOTIFY: OnceLock<Arc<Notify>> = OnceLock::new();
/// Handle to the process-wide resume [`Notify`] (#2831).
///
/// Both the firing side (`update_config` / `set_signed_out`) and the waiting
/// side (background loops) call this, so they share one instance. Use
/// `notify_one()` to fire: if a loop is parked it wakes immediately; if it's
/// mid-tick, a single permit is stored so the *next* `notified()` returns at
/// once — a resume that arrives during a tick is never lost.
///
/// **Over-notifying is safe by design.** A spurious wake (e.g. Memory Tree
/// toggled on while still signed out, so the effective policy is still paused)
/// just causes one cheap gate-checked tick that re-reads [`current_policy`] and
/// no-ops. We therefore fire on each individual un-pause transition rather than
/// computing the precise combined (config × signed-out) edge.
pub fn resume_notify() -> Arc<Notify> {
RESUME_NOTIFY
.get_or_init(|| Arc::new(Notify::new()))
.clone()
}
/// Update the gate's view of user config (e.g. after a settings change).
///
/// Fires [`resume_notify`] when this update moves the policy out of a paused
/// state (e.g. Memory Tree toggled back on), so parked background loops resume
/// promptly (#2831).
pub fn update_config(cfg: SchedulerGateConfig) {
if let Some(state) = STATE.get() {
let Some(state) = STATE.get() else {
return;
};
let resumed = {
let mut guard = state.write();
let was_paused = matches!(guard.policy, Policy::Paused { .. });
guard.cfg = cfg;
guard.policy = decide(&guard.signals, &guard.cfg);
was_paused && !matches!(guard.policy, Policy::Paused { .. })
};
if resumed {
resume_notify().notify_one();
}
}
@@ -311,6 +349,12 @@ pub fn set_signed_out(signed_out: bool) {
let prev = SIGNED_OUT.swap(signed_out, Ordering::AcqRel);
if prev != signed_out {
log::info!("[scheduler_gate] signed_out {} -> {}", prev, signed_out);
// #2831: signing back in (true -> false) is a transition out of
// `Policy::Paused { SignedOut }`. Wake any periodic loop so background
// sync restarts immediately rather than at the next tick boundary.
if prev && !signed_out {
resume_notify().notify_one();
}
}
}
@@ -325,6 +369,11 @@ pub fn set_signed_out(signed_out: bool) {
let prev = test_state::set_signed_out_for(id, signed_out);
if prev != signed_out {
log::info!("[scheduler_gate] signed_out {} -> {}", prev, signed_out);
// #2831: mirror the production sign-in wake so tests exercise the
// same resume-notify path (true -> false fires the loop wake).
if prev && !signed_out {
resume_notify().notify_one();
}
}
}
@@ -700,4 +749,89 @@ mod tests {
let p2 = try_acquire_llm_permit().expect("permit free after drop");
drop(p2);
}
/// #2831: both the firing side (`update_config` / `set_signed_out`) and the
/// waiting side (the periodic loop) must observe the *same* `Notify`, so
/// `resume_notify` must hand back one process-wide instance.
#[tokio::test]
async fn resume_notify_is_a_stable_singleton() {
let _g = lock();
assert!(
Arc::ptr_eq(&resume_notify(), &resume_notify()),
"resume_notify must return one shared instance"
);
}
/// A `notify_one()` wakes a task parked on `notified()` — the mechanism the
/// periodic loop relies on to resume early. Proves the singleton wiring
/// end-to-end (fire on one handle, wake on another).
#[tokio::test]
async fn resume_notify_wakes_a_parked_waiter() {
let _g = lock();
let waiter = resume_notify();
let parked = tokio::spawn(async move { waiter.notified().await });
// Yield so the spawned task reaches `.notified()` before we fire.
tokio::task::yield_now().await;
resume_notify().notify_one();
timeout(TokioDuration::from_secs(1), parked)
.await
.expect("parked waiter must wake promptly after notify_one")
.expect("waiter task must not panic");
}
/// #2831 wiring: a paused→running `update_config` transition and a
/// sign-in (`set_signed_out` true→false) each fire the resume notify.
///
/// Seeds `STATE` directly (the test module can reach it) so `update_config`
/// / `set_signed_out` are live without spawning the real sampler. We drive
/// `Off → AlwaysOn`, which is deterministically `Paused → Aggressive`
/// regardless of any policy a prior test left behind, so the transition —
/// and thus the `notify_one()` — is guaranteed.
#[tokio::test]
async fn resume_transitions_fire_the_notify() {
use crate::openhuman::config::SchedulerGateMode;
let _g = lock();
// Ensure STATE is initialised. `set` is a no-op if an earlier test
// already promoted it — that's fine, we re-drive the transition below.
let cfg = SchedulerGateConfig {
mode: SchedulerGateMode::Off,
..Default::default()
};
let signals = Signals::sample();
let policy = decide(&signals, &cfg);
let _ = STATE.set(Arc::new(RwLock::new(State {
cfg,
signals,
policy,
})));
// --- update_config: Paused -> running fires the notify ---
let waiter = resume_notify();
let parked = tokio::spawn(async move { waiter.notified().await });
tokio::task::yield_now().await;
update_config(SchedulerGateConfig {
mode: SchedulerGateMode::Off,
..Default::default()
}); // -> Paused { UserDisabled }
update_config(SchedulerGateConfig {
mode: SchedulerGateMode::AlwaysOn,
..Default::default()
}); // Paused -> Aggressive => resume fires
timeout(TokioDuration::from_secs(1), parked)
.await
.expect("update_config un-pause must wake the resume waiter")
.expect("waiter task must not panic");
// --- set_signed_out true -> false fires the notify ---
let waiter2 = resume_notify();
let parked2 = tokio::spawn(async move { waiter2.notified().await });
tokio::task::yield_now().await;
set_signed_out(true);
set_signed_out(false); // true -> false => resume fires
timeout(TokioDuration::from_secs(1), parked2)
.await
.expect("sign-in (signed_out true->false) must wake the resume waiter")
.expect("waiter task must not panic");
}
}