diff --git a/desktop/src-tauri/src/commands/global_agent_config_epoch_tests.rs b/desktop/src-tauri/src/commands/global_agent_config_epoch_tests.rs index 3f9b42d86..2e8365bfd 100644 --- a/desktop/src-tauri/src/commands/global_agent_config_epoch_tests.rs +++ b/desktop/src-tauri/src/commands/global_agent_config_epoch_tests.rs @@ -222,6 +222,12 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { let spawn_got_nonempty_teams = Arc::new(Mutex::new(false)); let spawn_got_nonempty_global = Arc::new(Mutex::new(false)); let receipt_called = Arc::new(Mutex::new(false)); + // Capture all receipt fields for post-completion assertions. + let receipt_pubkey = Arc::new(Mutex::new(None::)); + let receipt_relay = Arc::new(Mutex::new(None::)); + let receipt_pid = Arc::new(Mutex::new(0u32)); + let receipt_instance_id = Arc::new(Mutex::new(String::new())); + let receipt_started_at = Arc::new(Mutex::new(String::new())); let stop_called2 = stop_called.clone(); let spawn_relay2 = spawn_relay.clone(); @@ -230,9 +236,15 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { let spawn_teams2 = spawn_got_nonempty_teams.clone(); let spawn_global2 = spawn_got_nonempty_global.clone(); let receipt_called2 = receipt_called.clone(); + let receipt_pubkey2 = receipt_pubkey.clone(); + let receipt_relay2 = receipt_relay.clone(); + let receipt_pid2 = receipt_pid.clone(); + let receipt_iid2 = receipt_instance_id.clone(); + let receipt_sat2 = receipt_started_at.clone(); let pubkey2 = pubkey.clone(); let app_handle_for_assert = app_handle.clone(); + let pubkey_for_assert = pubkey.clone(); let result = tokio::task::spawn_blocking(move || { restart_under_captured_epoch_for( &app_handle, @@ -273,9 +285,15 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { job: None, }) }, - // write_receipt_fn: record the call, verify pubkey, succeed. + // write_receipt_fn: record all receipt fields for post-completion assertions. + // Full key (pubkey + relay_url), pid, desktop_instance_id, started_at. move |_app, receipt| { *receipt_called2.lock().unwrap() = true; + *receipt_pubkey2.lock().unwrap() = Some(receipt.key.pubkey.clone()); + *receipt_relay2.lock().unwrap() = Some(receipt.key.relay_url.clone()); + *receipt_pid2.lock().unwrap() = receipt.pid; + *receipt_iid2.lock().unwrap() = receipt.desktop_instance_id.clone(); + *receipt_sat2.lock().unwrap() = receipt.started_at.clone(); assert_eq!( receipt.key.pubkey, pubkey2, "receipt must carry the correct pubkey" @@ -321,6 +339,36 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { *receipt_called.lock().unwrap(), "write_receipt_fn must be called" ); + // ── Receipt field assertions ────────────────────────────────────────────── + // A full relay-bearing key: pubkey + relay_url. + assert_eq!( + receipt_pubkey.lock().unwrap().as_deref(), + Some(pubkey_for_assert.as_str()), + "receipt.key.pubkey must match the restarted agent" + ); + assert_eq!( + receipt_relay.lock().unwrap().as_deref(), + Some(relay_url), + "receipt.key.relay_url must match the captured relay; fails if the spawn path \ + builds the receipt with an empty or wrong relay" + ); + assert!( + *receipt_pid.lock().unwrap() > 0, + "receipt.pid must be non-zero (a real spawned child PID)" + ); + // Assert receipt.desktop_instance_id equals the actual value produced by + // current_instance_id(app) — proves the field is sourced from the app + // handle, not left empty or set to a hardcoded value. + let expected_instance_id = app_handle_for_assert.config().identifier.clone(); + assert_eq!( + receipt_instance_id.lock().unwrap().as_str(), + expected_instance_id.as_str(), + "receipt.desktop_instance_id must equal current_instance_id(app)" + ); + assert!( + !receipt_started_at.lock().unwrap().is_empty(), + "receipt.started_at must be a non-empty ISO timestamp" + ); // Verify the runtime is registered with the captured scope_id. { let state = app_handle_for_assert.state::(); @@ -333,12 +381,27 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { "runtime must be registered with the captured scope_id" ); } - // Verify the final captured disk record: the agent is present in the store. + // ── Final disk record assertions ────────────────────────────────────────── + // The restart-produced state: last_started_at set, last_stopped_at cleared, + // last_error cleared. Fails if the post-spawn save is removed or if the + // record-update block is bypassed. let final_records = crate::managed_agents::storage::load_managed_agents_at(tmp.path()).unwrap_or_default(); + let final_rec = final_records + .iter() + .find(|r| r.pubkey == pubkey_for_assert) + .expect("final disk record must contain the restarted agent"); assert!( - final_records.iter().any(|r| r.pubkey == "bb".repeat(32)), - "final disk record must contain the restarted agent" + final_rec.last_started_at.is_some(), + "last_started_at must be Some after a successful restart (set by post-spawn record update)" + ); + assert!( + final_rec.last_stopped_at.is_none(), + "last_stopped_at must be None after restart (cleared by post-spawn record update)" + ); + assert!( + final_rec.last_error.is_none(), + "last_error must be None after a successful restart (cleared by post-spawn record update)" ); } diff --git a/desktop/src-tauri/src/commands/global_agent_config_tests.rs b/desktop/src-tauri/src/commands/global_agent_config_tests.rs index 27c1db959..36eb4388b 100644 --- a/desktop/src-tauri/src/commands/global_agent_config_tests.rs +++ b/desktop/src-tauri/src/commands/global_agent_config_tests.rs @@ -512,32 +512,30 @@ async fn test_workspace_switch_after_preflight_aborts_before_stop() { /// A record-level Mesh model change after preflight (without advancing workspace /// generation) → epoch detects mismatch in re-resolved Mesh model, aborts before stop. /// -/// Thufir's test 4: "drive the async production driver; seed an eligible runtime -/// via the cross-platform child helper; mutate an actually-participating -/// record/definition/global input from the injected mesh_fn (which runs in the -/// pre-stop phase); assert re-resolution detects the mismatch and stop never -/// fires." +/// Drives `restart_local_agent_on_config_change_for` — the full async production +/// driver. The injected `mesh_fn` mutates the on-disk record's `provider` to +/// `"relay-mesh"` AFTER the driver has already resolved `context.mesh_model_id = None` +/// (from the original `provider = "anthropic"`). The epoch's in-epoch TOCTOU guard +/// then re-resolves the mutated record, gets `Some("auto")` ≠ `None`, and fires +/// `Skipped` before stop. /// -/// The async driver (`restart_local_agent_on_config_change_for`) runs mesh_fn -/// in the pre-stop phase, captures the `CapturedRestartContext`, then hands off -/// to `restart_under_captured_epoch_for` (the stop→spawn primitive). We inject -/// a `mesh_fn` that succeeds but captures `mesh_model_id = Some("model-a")` -/// into the context by returning Ok — the real driver then stores whatever -/// the record has at the time of pre-stop resolution into `context.mesh_model_id`. +/// Flow: +/// 1. Seed record: `provider = "anthropic"`, live runtime. +/// Pre-stop resolve: `resolve_effective_relay_mesh_model_id` → `None`. +/// `context.mesh_model_id = None`. +/// 2. `mesh_fn` rewrites `provider = "relay-mesh"` to disk and succeeds. +/// 3. Epoch re-resolves from the mutated record: +/// `relay_mesh_model_id()` = `Some("auto")` ≠ `None` → TOCTOU guard fires → Skipped. +/// 4. Assert `RestartOutcome::Skipped`; stop_fn must NOT be called. /// -/// To trigger the in-epoch TOCTOU guard: -/// - The record on disk has `relay_mesh: None` → driver pre-resolves -/// `mesh_model_id = None`. -/// - We set `mesh_model_id = Some("model-a")` in an injected post-preflight -/// hook by calling `restart_under_captured_epoch_for` directly with a context -/// whose `mesh_model_id` disagrees with the on-disk record. -/// -/// Approach: use `restart_under_captured_epoch_for` directly (with the -/// cross-platform long-lived child seeded), so the in-epoch TOCTOU guard fires. -/// The "async driver drives the TOCTOU guard" variant is covered by the -/// generation-advance test (`test_workspace_switch_after_preflight_aborts_before_stop`). +/// Invariant: removing the in-epoch re-resolve guard in +/// `restart_under_captured_epoch_for` would allow the epoch to proceed with the +/// old `context.mesh_model_id`, bypass the mismatch — stop would be called +/// and the test would fail. #[tokio::test] async fn test_record_mesh_change_after_preflight_aborts_before_stop() { + use super::super::restart_local_agent_on_config_change_for; + use crate::commands::global_agent_config::RestartOutcome; use crate::managed_agents::{ storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord, ManagedAgentRuntimeKey, @@ -548,13 +546,9 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() { let pubkey = "aa".repeat(32); let relay_url = "wss://relay.example"; - // Build a Ready record: provider+model in structured fields so the agent - // passes the eligibility gate (old_ready = true). ANTHROPIC_API_KEY is - // required by buzz_agent_requirements when provider=anthropic, so we set - // it in env_vars to satisfy the readiness check. Two globals differ by - // one env_var so env_changed = true → should_restart_on_config_change - // returns true. relay_mesh: None so the in-epoch re-resolve yields None - // while context.mesh_model_id = Some("model-a") → Skipped before stop. + // Eligible record: provider=anthropic, model+ANTHROPIC_API_KEY → old_ready=true. + // relay_mesh=None so the pre-stop resolve yields mesh_model_id=None. + // old_global != new_global (env differ) so env_changed=true → eligible. let mut record_env_vars = std::collections::BTreeMap::new(); record_env_vars.insert( "ANTHROPIC_API_KEY".to_string(), @@ -581,7 +575,7 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() { parallelism: 1, system_prompt: None, model: Some("claude-3-5-sonnet-20241022".to_string()), - provider: Some("anthropic".to_string()), + provider: Some("anthropic".to_string()), // initial provider — not relay-mesh persona_source_version: None, env_vars: record_env_vars, start_on_app_launch: false, @@ -611,7 +605,7 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() { definition_respond_to: None, definition_respond_to_allowlist: Default::default(), definition_parallelism: None, - relay_mesh: None, // ← no relay_mesh; re-resolve yields None ≠ Some("model-a") + relay_mesh: None, // no relay-mesh marker; pre-stop resolve yields None runtime: None, name_pool: vec![], }; @@ -622,6 +616,17 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() { let app = make_mock_app(); let app_handle = app.handle().clone(); + // Derive the actual owner pubkey from the mock app's signing keys. + // restart_local_agent_on_config_change_for verifies hex == scope.owner_pubkey. + let actual_owner_hex = { + let state = app_handle.state::(); + state + .signing_keys() + .expect("mock app must have signing keys") + .public_key() + .to_hex() + }; + // Seed a live runtime with a cross-platform long-lived child (avoids sync eviction). let rt_key = ManagedAgentRuntimeKey::new(&pubkey, relay_url).unwrap(); let seeded_pid = { @@ -656,13 +661,13 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() { let scope = crate::managed_agents::scope::WorkspaceAgentScope { scope_id: "test-scope".to_string(), relay_url: relay_url.to_string(), - owner_pubkey: "aa".repeat(32), + owner_pubkey: actual_owner_hex, definitions_dir: tmp.path().to_path_buf(), generation: gen, }; // old_global and new_global differ by one env_var so env_changed = true - // (eligibility gate passes) while the record's relay_mesh stays None. + // (eligibility gate passes) while the record's provider stays anthropic. let mut new_global_env = std::collections::BTreeMap::new(); new_global_env.insert("SOME_EXTRA_KEY".to_string(), "v2".to_string()); let old_global = crate::managed_agents::GlobalAgentConfig::default(); @@ -671,62 +676,61 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() { ..Default::default() }; - // Build a context that claims `mesh_model_id = Some("model-a")`, but - // the actual records on disk have no relay_mesh config, so the epoch's - // re-resolve will return None — triggering the TOCTOU mismatch guard. - let context_with_mesh = CapturedRestartContext { - scope: scope.clone(), - personas: vec![], - teams: vec![], - global: old_global.clone(), - owner_hex: "aa".repeat(32), - mesh_model_id: Some("model-a".to_string()), - }; - let stop_called = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); let stop_called2 = stop_called.clone(); + let tmp_path = tmp.path().to_path_buf(); + let pubkey_clone = pubkey.clone(); - // Call the epoch core directly — mesh_fn is not involved here since - // we're testing the in-epoch TOCTOU check that re-resolves relay_mesh. - // This exercises the same code path that the async driver reaches after - // mesh_fn completes: stop→spawn epoch with pre-captured context. - let result = tokio::task::spawn_blocking(move || { - restart_under_captured_epoch_for( - &app_handle, - &pubkey, - &old_global, - &new_global, - &[], - &context_with_mesh, - move |_app, _rec, _runtimes| { - stop_called2.store(true, std::sync::atomic::Ordering::SeqCst); - Err("stop_fn called unexpectedly".to_string()) - }, - |_app, _rec, _relay, _owner, _personas, _global, _teams| { - Err("spawn not expected".to_string()) - }, - |_app, _receipt| Err("receipt not expected".to_string()), - ) - }) - .await - .expect("spawn_blocking must not panic"); + // Drive the full async production driver. The mesh_fn mutates the on-disk + // record's provider to "relay-mesh" AFTER the driver has resolved + // context.mesh_model_id = None (provider was "anthropic" at resolve time). + // The epoch then re-resolves the mutated record and gets Some("auto") != None + // -> TOCTOU guard fires -> Skipped before stop. + let outcome = restart_local_agent_on_config_change_for( + &app_handle, + &pubkey, + &old_global, + &new_global, + &[], + &scope, + tmp.path(), + // mesh_fn: rewrites provider to "relay-mesh" on disk, then succeeds. + // The driver already captured context.mesh_model_id=None from the + // original provider="anthropic" record — this mutation happens after. + move |_app, _model| { + let mut records = crate::managed_agents::storage::load_managed_agents_at(&tmp_path) + .unwrap_or_default(); + for r in &mut records { + if r.pubkey == pubkey_clone { + r.provider = Some("relay-mesh".to_string()); + } + } + let _ = crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &records); + Box::pin(async { Ok(()) }) + }, + // stop_fn: must NOT be called — TOCTOU guard fires before stop. + move |_app, _rec, _runtimes| { + stop_called2.store(true, std::sync::atomic::Ordering::SeqCst); + Err("stop must not be called when Mesh model changed after preflight".to_string()) + }, + |_app, _rec, _relay, _owner, _personas, _global, _teams| { + Err("spawn not expected".to_string()) + }, + |_app, _receipt| Err("receipt not expected".to_string()), + ) + .await; + + // Kill the seeded process now that the epoch has consumed it. + let _ = crate::managed_agents::terminate_process(seeded_pid); assert!( - matches!(result, Err(EpochError::Skipped(_))), - "Mesh model mismatch must produce Skipped before stop: {result:?}" + matches!(outcome, RestartOutcome::Skipped), + "Mesh model mismatch (via full async driver) must produce Skipped before stop: {outcome:?}" ); - if let Err(EpochError::Skipped(msg)) = &result { - assert!( - msg.contains("relay-mesh model changed") || msg.contains("mesh"), - "Skipped reason must mention mesh model change: {msg}" - ); - } assert!( !stop_called.load(std::sync::atomic::Ordering::SeqCst), "stop must NOT be called when Mesh model changed after preflight" ); - // Cleanup: kill the seeded long-lived process so it doesn't leak. - let _ = crate::managed_agents::terminate_process(seeded_pid); } #[path = "global_agent_config_epoch_tests.rs"] diff --git a/desktop/src-tauri/src/commands/mesh_llm_scope.rs b/desktop/src-tauri/src/commands/mesh_llm_scope.rs index 48c875269..6e777cc29 100644 --- a/desktop/src-tauri/src/commands/mesh_llm_scope.rs +++ b/desktop/src-tauri/src/commands/mesh_llm_scope.rs @@ -102,15 +102,8 @@ pub(crate) async fn fail_if_client_mesh_active( /// `fail_if_client_mesh_active` preflight, then invoke `transition_body` while /// the guard remains held. /// -/// Used by `apply_workspace`, live identity import, and tests to prove the -/// serialization contract against `install_client_under_workspace_transition`. -/// -/// `transition_body` is an async closure (returns a `BoxFuture`) so that -/// production callers that dispatch `tokio::task::spawn_blocking` and await the -/// result keep the guard alive across the dispatch. Tests may pass a simple -/// sync-compatible closure via `|| Box::pin(async { Ok("done") })`. -/// -/// If the preflight fails, `transition_body` is never called. +/// Delegates to [`with_workspace_transition_preflight_with_hook`] with a no-op +/// pre-acquisition hook. See that function for full documentation. pub(crate) async fn with_workspace_transition_preflight( app: &AppHandle, transition_body: F, @@ -118,8 +111,32 @@ pub(crate) async fn with_workspace_transition_preflight( where R: tauri::Runtime, F: FnOnce() -> BoxFuture<'static, Result>, +{ + with_workspace_transition_preflight_with_hook(app, || {}, transition_body).await +} + +/// Inner implementation of [`with_workspace_transition_preflight`] with an +/// injectable `pre_acquisition_hook`. +/// +/// `pre_acquisition_hook` fires once, synchronously, immediately before +/// `workspace_transition.lock().await`. In production this is `|| {}`; tests +/// inject a closure that signals "I'm about to acquire" so the test can +/// establish a deterministic ordering between holder and contender. +/// +/// This is `pub(crate)` so tests in `mesh_llm_transition_tests` can drive the +/// exact lock-acquisition boundary while remaining invisible to external callers. +pub(crate) async fn with_workspace_transition_preflight_with_hook( + app: &AppHandle, + pre_acquisition_hook: H, + transition_body: F, +) -> Result +where + R: tauri::Runtime, + H: FnOnce(), + F: FnOnce() -> BoxFuture<'static, Result>, { let state = app.state::(); + pre_acquisition_hook(); let _transition_guard = state.workspace_transition.lock().await; // Fail closed if a client-mode Mesh runtime is active. @@ -152,13 +169,8 @@ where /// identity under the guard — `(scope_id, normalized relay, owner_pubkey, /// generation)` — then call the injected `install` closure. /// -/// Fails if: -/// - no active scope exists at the time of validation; -/// - any identity field of the captured scope differs from the current active scope; -/// - the generation counter has advanced (a workspace switch occurred). -/// -/// Used by `ensure_relay_mesh_for_record` after capturing scope + discovering -/// the bootstrap target. Tests inject `install` directly so no port is touched. +/// Delegates to [`install_client_under_workspace_transition_with_hook`] with a +/// no-op pre-acquisition hook. See that function for full documentation. pub(crate) async fn install_client_under_workspace_transition( app: &AppHandle, captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope, @@ -168,8 +180,39 @@ where R: tauri::Runtime, I: FnOnce() -> Fut, Fut: Future>, +{ + install_client_under_workspace_transition_with_hook(app, captured_scope, || {}, install).await +} + +/// Inner implementation of [`install_client_under_workspace_transition`] with an +/// injectable `pre_acquisition_hook`. +/// +/// `pre_acquisition_hook` fires once, synchronously, immediately before +/// `workspace_transition.lock().await`. In production this is `|| {}`; tests +/// inject a closure that signals "I'm about to acquire" so the test can +/// establish a deterministic ordering between holder and contender. +/// +/// Fails if: +/// - no active scope exists at the time of validation; +/// - any identity field of the captured scope differs from the current active scope; +/// - the generation counter has advanced (a workspace switch occurred). +/// +/// This is `pub(crate)` so tests in `mesh_llm_transition_tests` can drive the +/// exact lock-acquisition boundary while remaining invisible to external callers. +pub(crate) async fn install_client_under_workspace_transition_with_hook( + app: &AppHandle, + captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope, + pre_acquisition_hook: H, + install: I, +) -> Result<(), String> +where + R: tauri::Runtime, + H: FnOnce(), + I: FnOnce() -> Fut, + Fut: Future>, { let state = app.state::(); + pre_acquisition_hook(); let _transition_guard = state.workspace_transition.lock().await; // Validate the captured scope's generation and full identity under the guard. diff --git a/desktop/src-tauri/src/commands/mesh_llm_transition_tests.rs b/desktop/src-tauri/src/commands/mesh_llm_transition_tests.rs index 58f1841ca..009c28ae1 100644 --- a/desktop/src-tauri/src/commands/mesh_llm_transition_tests.rs +++ b/desktop/src-tauri/src/commands/mesh_llm_transition_tests.rs @@ -91,12 +91,20 @@ async fn test_active_client_stop_then_transition_preflight_succeeds() { /// the lock is held), the install task acquires the lock and detects the stale /// captured scope without invoking the install closure. /// -/// Both sides use production helpers — no direct `workspace_transition` lock -/// acquisition. Channels establish entry/release ordering without `sleep`. +/// Handshake: +/// 1. Transition task acquires the lock, signals "lock_held" BEFORE committing +/// scope B (so install's captured scope is still valid at its capture time). +/// 2. Install task spawns, calls `install_client_under_workspace_transition` +/// with a pre-acquisition hook that signals "at_lock_boundary", then +/// blocks on `workspace_transition.lock().await`. +/// 3. Test waits for "at_lock_boundary", then sends "release" to the +/// transition body. The body commits scope B and returns, releasing the lock. +/// 4. Install task acquires the lock, re-validates scope, detects the stale +/// generation, returns Err without invoking the install closure. /// -/// Proves the first serialization direction: a workspace switch (transition -/// side) commits between scope-capture and install lock-acquisition; the -/// install helper must reject the stale scope under the lock. +/// Invariant: if the contender could bypass the lock, it would acquire before +/// the transition commits scope B, see a valid scope, and invoke the install +/// closure — `install_was_called` would be true and the assertion would fail. #[tokio::test] async fn test_transition_held_queued_install_detects_stale_scope() { use crate::managed_agents::scope::{ @@ -130,23 +138,32 @@ async fn test_transition_held_queued_install_detects_stale_scope() { WorkspaceAgentScope::new(relay_a.to_string(), owner_a.clone(), &base, gen_a); state.commit_active_scope(captured_scope.clone()); - // Channels: transition body signals "entered" once it holds the lock; - // test unblocks the body once the install task is queued. - let (body_entered_tx, body_entered_rx) = oneshot::channel::<()>(); + // Channels: + // lock_held: transition body → test (lock acquired, scope NOT yet committed) + // at_lock_boundary: install hook → test (install is about to call lock().await) + // body_release: test → transition body (ok to commit scope B and release) + let (lock_held_tx, lock_held_rx) = oneshot::channel::<()>(); + let (at_lock_boundary_tx, at_lock_boundary_rx) = oneshot::channel::<()>(); let (body_release_tx, body_release_rx) = oneshot::channel::<()>(); - // ── Side A: transition side uses the production helper ─────────────────── - // The body holds the lock, advances the scope, signals entry, then blocks - // until unblocked by the test — proving the hold is production-owned. + // ── Side A: transition side — signals lock_held BEFORE committing scope B ─ let app_a = app_handle.clone(); let base_a = base.clone(); let transition_task = tokio::task::spawn(async move { let app_a_ref = app_a.clone(); super::scope_impl::with_workspace_transition_preflight(&app_a_ref, move || { Box::pin(async move { - // Advance generation and commit a DISTINCT scope (different relay) - // to simulate a committed workspace switch while the lock is held. let state_a = app_a.state::(); + + // Signal "lock held" BEFORE committing scope B. + // The install task captures its scope before this point; + // only after it queues at the lock does the holder commit B. + let _ = lock_held_tx.send(()); + + // Wait for the install task to queue at the lock boundary. + let _ = body_release_rx.await; + + // Now commit scope B (advances generation, install sees stale). let gen_b = next_scope_generation(); let new_scope = WorkspaceAgentScope::new( "wss://transition-test-b.example".to_string(), @@ -156,30 +173,29 @@ async fn test_transition_held_queued_install_detects_stale_scope() { ); state_a.commit_active_scope(new_scope); - // Signal: holding the lock with committed scope change. - let _ = body_entered_tx.send(()); - // Block until the test confirms the install task has queued. - let _ = body_release_rx.await; - Ok::<(), String>(()) }) }) .await }); - // Wait until the transition body signals it has the lock. - body_entered_rx + // Wait until the transition body holds the lock (before scope B is committed). + lock_held_rx .await - .expect("transition body must signal entry"); + .expect("transition body must signal lock_held"); - // ── Side B: install side uses the production helper ────────────────────── + // ── Side B: install side — pre-acquisition hook signals at_lock_boundary ── let install_was_called = Arc::new(AtomicBool::new(false)); let install_called_clone = Arc::clone(&install_was_called); let app_b = app_handle.clone(); let install_task = tokio::task::spawn(async move { - super::scope_impl::install_client_under_workspace_transition( + super::scope_impl::install_client_under_workspace_transition_with_hook( &app_b, &captured_scope, + // pre-acquisition hook: fires before lock().await — signals "queued". + move || { + let _ = at_lock_boundary_tx.send(()); + }, || { let called = Arc::clone(&install_called_clone); async move { @@ -191,10 +207,13 @@ async fn test_transition_held_queued_install_detects_stale_scope() { .await }); - // Yield to let the install task queue on the lock. - tokio::task::yield_now().await; + // Wait for install to reach the lock boundary, then unblock the transition. + at_lock_boundary_rx + .await + .expect("install hook must signal at_lock_boundary"); - // Unblock the transition body → it releases the lock → install acquires it. + // Unblock the transition body → it commits scope B, releases the lock → + // install acquires the lock, detects stale scope, returns Err. let _ = body_release_tx.send(()); let install_result = install_task.await.expect("install task must not panic"); @@ -226,12 +245,22 @@ async fn test_transition_held_queued_install_detects_stale_scope() { /// install body completes, the transition task acquires the lock, runs /// `fail_if_client_mesh_active`, and observes the installed client. /// -/// Both sides use production helpers — no direct `workspace_transition` lock -/// acquisition. Channels establish entry/release ordering without `sleep`. +/// Handshake: +/// 1. Install task acquires the lock, signals "lock_held" BEFORE installing the +/// client runtime (so the client is not yet visible to the contender). +/// 2. Transition task calls `with_workspace_transition_preflight` with a +/// pre-acquisition hook that signals "at_lock_boundary", then blocks on +/// `workspace_transition.lock().await`. +/// 3. Test waits for "at_lock_boundary", then sends "release" to the +/// install closure. The closure installs the client and returns, releasing +/// the lock. +/// 4. Transition task acquires the lock, runs `fail_if_client_mesh_active`, +/// observes the installed client, returns Err. /// -/// Proves the second serialization direction: a client install commits between -/// the transition's pre-lock preflight and its lock-acquisition; the transition -/// helper must observe the installed client under the lock and fail. +/// Invariant: if the contender could bypass the lock, it would run +/// `fail_if_client_mesh_active` before the client is installed — seeing no +/// client — and return Ok. The `transition_result.is_err()` assertion would +/// then fail, proving the lock is not enforced. #[tokio::test] async fn test_install_held_transition_preflight_observes_client() { use crate::managed_agents::scope::{ @@ -262,14 +291,15 @@ async fn test_install_held_transition_preflight_observes_client() { let scope = WorkspaceAgentScope::new(relay.to_string(), owner.clone(), &base, gen); state.commit_active_scope(scope.clone()); - // Channels: install body signals "entered + client installed" once it holds - // the lock; test unblocks the body once the transition task is queued. - let (install_entered_tx, install_entered_rx) = oneshot::channel::<()>(); + // Channels: + // lock_held: install body → test (lock acquired, client NOT yet installed) + // at_lock_boundary: transition hook → test (transition is about to call lock().await) + // install_release: test → install body (ok to install client and release lock) + let (lock_held_tx, lock_held_rx) = oneshot::channel::<()>(); + let (at_lock_boundary_tx, at_lock_boundary_rx) = oneshot::channel::<()>(); let (install_release_tx, install_release_rx) = oneshot::channel::<()>(); - // ── Side A: install side uses the production helper ────────────────────── - // The install closure installs a mock client, signals entry, then blocks - // until unblocked — proving the client is committed under the lock. + // ── Side A: install side — signals lock_held BEFORE installing the client ─ let app_a = app_handle.clone(); let scope_a = scope.clone(); let install_task = tokio::task::spawn(async move { @@ -278,40 +308,51 @@ async fn test_install_held_transition_preflight_observes_client() { &app_a, &scope_a, move || async move { - // Install the mock client runtime while holding the lock. + // Signal "lock held" BEFORE installing the client. + // The transition task captures its pre-lock state before this + // point; only after it queues does the install closure install. + let _ = lock_held_tx.send(()); + + // Wait for the transition task to queue at the lock boundary. + let _ = install_release_rx.await; + + // Now install the mock client runtime under the held lock. let state_a = app_a_for_closure.state::(); let client_runtime = crate::mesh_llm::build_mock_client_runtime_for_test(); *state_a.mesh_llm_runtime.lock().await = Some(client_runtime); - // Signal: lock held, client installed. - let _ = install_entered_tx.send(()); - // Block until the test confirms the transition task has queued. - let _ = install_release_rx.await; - Ok::<(), String>(()) }, ) .await }); - // Wait until the install body signals it has the lock with the client installed. - install_entered_rx + // Wait until the install body holds the lock (before client is installed). + lock_held_rx .await - .expect("install body must signal entry"); + .expect("install body must signal lock_held"); - // ── Side B: transition side uses the production helper ─────────────────── + // ── Side B: transition side — pre-acquisition hook signals at_lock_boundary let app_b = app_handle.clone(); let transition_task = tokio::task::spawn(async move { - super::scope_impl::with_workspace_transition_preflight(&app_b, || { - Box::pin(async { Ok::<&str, String>("body ran") }) - }) + super::scope_impl::with_workspace_transition_preflight_with_hook( + &app_b, + // pre-acquisition hook: fires before lock().await — signals "queued". + move || { + let _ = at_lock_boundary_tx.send(()); + }, + || Box::pin(async { Ok::<&str, String>("body ran") }), + ) .await }); - // Yield to let the transition task queue on the lock. - tokio::task::yield_now().await; + // Wait for transition to reach the lock boundary, then unblock the install. + at_lock_boundary_rx + .await + .expect("transition hook must signal at_lock_boundary"); - // Unblock the install body → it releases the lock → transition acquires it. + // Unblock the install closure → it installs the client, releases the lock → + // transition acquires the lock, observes the client, returns Err. let _ = install_release_tx.send(()); install_task diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands.rs b/desktop/src-tauri/src/managed_agents/runtime_commands.rs index 74d96203f..dac8392a8 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands.rs @@ -236,7 +236,21 @@ pub(crate) fn start_managed_agent_runtime_pair_lazy( relay_url: String, app: AppHandle, ) -> Result { - start_pair(pubkey, relay_url, true, None, app) + start_pair_lazy_for(pubkey, relay_url, app) +} + +/// Generic start-pair-lazy seam shared by the production adapter and tests. +/// +/// Acquires `managed_agent_runtime_transition` as its first action — the same +/// lock that `stop`, `restart`, and `drain` operations hold, serialising all +/// runtime mutations. Tests that need a mock-runtime contender call this +/// function directly instead of the non-generic production adapter. +pub(crate) fn start_pair_lazy_for( + pubkey: String, + relay_url: String, + app: tauri::AppHandle, +) -> Result { + start_pair_for(pubkey, relay_url, true, None, app) } #[tauri::command] @@ -254,6 +268,16 @@ fn start_pair( lazy: bool, expected_updated_at: Option<&str>, app: AppHandle, +) -> Result { + start_pair_for(pubkey, relay_url, lazy, expected_updated_at, app) +} + +fn start_pair_for( + pubkey: String, + relay_url: String, + lazy: bool, + expected_updated_at: Option<&str>, + app: tauri::AppHandle, ) -> Result { let state = app.state::(); let _transition = state @@ -883,32 +907,45 @@ where /// Production adapter for [`compensate_drain_for`]. /// -/// Lock acquisition order inside this function: -/// 1. `_rt_transition_held` — already held by caller (passed by value). -/// 2. Acquire `managed_agents_store_lock` — store-lock-only writers -/// (instance edits, status flushes) are serialized here, not on the -/// transition guard. -/// 3. Validate `captured_scope.generation` under the store lock. -/// 4. `load_managed_agents_at` — records loaded AFTER both locks held. -/// 5. Delegate to [`compensate_drain_for`] — store guard held through -/// every save performed by `start_pair_under_held_locks`. +/// Delegates to [`compensate_drain_with_hook`] with a no-op hook. +/// Lock acquisition order: transition guard already held → store lock acquired +/// → hook fires (no-op in production) → validate generation → load → restore/save. pub(crate) fn compensate_drain( app: &tauri::AppHandle, stopped: &[DrainJournalEntry], captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope, - // Caller passes ownership of the already-held transition guard so the lock - // is never dropped between drain and compensation. _rt_transition_held: std::sync::MutexGuard<'_, ()>, +) -> Option { + compensate_drain_with_hook(app, stopped, captured_scope, _rt_transition_held, || {}) +} + +/// Inner implementation of [`compensate_drain`] with an injectable +/// `on_store_acquired` hook. +/// +/// Lock acquisition order: +/// 1. `_rt_transition_held` — already held by caller (passed by value). +/// 2. Acquire `managed_agents_store_lock`. +/// 3. Call `on_store_acquired` (no-op in production; tests inject a closure +/// that writes a sentinel and synchronises with a writer thread to prove +/// genuine store-lock contention). +/// 4. Validate `captured_scope.generation` under the store lock. +/// 5. `load_managed_agents_at` — records loaded AFTER both locks held. +/// 6. Delegate to [`compensate_drain_for`] — store guard held through every +/// save performed by `start_pair_under_held_locks`. +pub(crate) fn compensate_drain_with_hook( + app: &tauri::AppHandle, + stopped: &[DrainJournalEntry], + captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope, + _rt_transition_held: std::sync::MutexGuard<'_, ()>, + on_store_acquired: impl FnOnce(), ) -> Option { if stopped.is_empty() { - // Release the transition guard immediately — nothing to restore. drop(_rt_transition_held); return None; } let state = app.state::(); - // 2. Acquire the store lock — transition already held via _rt_transition_held. let _store = match state.managed_agents_store_lock.lock() { Ok(g) => g, Err(e) => { @@ -918,10 +955,10 @@ pub(crate) fn compensate_drain( } }; - // 3. Validate the captured scope before restoring. If the workspace switched - // between the drain failure and this compensation, the new scope's own - // restore pass will start the correct agents — we must not restart agents - // for a scope that is no longer active. + // 3. Fire the hook while both locks are held. No-op in production. + on_store_acquired(); + + // 4. Validate the captured scope under the store lock. if let Err(stale_msg) = crate::managed_agents::scope::validate_scope_generation(captured_scope) { return Some(format!( @@ -929,7 +966,7 @@ pub(crate) fn compensate_drain( )); } - // 4. Load records under the held store lock. + // 5. Load records under the held store lock. let mut records = match load_managed_agents_at(&captured_scope.definitions_dir) { Ok(r) => r, Err(e) => { @@ -939,7 +976,7 @@ pub(crate) fn compensate_drain( } }; - // 5. Delegate to the lock-free core — store guard is held through every save. + // 6. Delegate — store guard held through every save. compensate_drain_for(stopped, &mut records, |entry, recs| { start_pair_under_held_locks( app, diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands_concurrency_tests.rs b/desktop/src-tauri/src/managed_agents/runtime_commands_concurrency_tests.rs index 1b01e64a6..960acdc64 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands_concurrency_tests.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands_concurrency_tests.rs @@ -6,26 +6,28 @@ use super::*; -/// Deterministic writer-vs-compensation ordering via the production -/// `compensate_drain` adapter. +/// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam. /// -/// Mechanism: -/// 1. The test acquires the `managed_agent_runtime_transition` guard. -/// 2. A writer thread is spawned. It waits on a channel before touching the -/// store, so its load→edit→save only begins AFTER the test explicitly -/// signals "compensation done." -/// 3. The test calls the production `compensate_drain` adapter — which -/// acquires `managed_agents_store_lock` internally, validates generation, -/// loads records, delegates to `compensate_drain_for`, and saves — all while -/// the transition guard is owned by compensate_drain (passed by value). -/// 4. After `compensate_drain` returns the transition guard is consumed; the -/// test signals the writer via the channel. -/// 5. The writer acquires the store lock, applies its sentinel edit, and saves. -/// 6. Final disk state must contain BOTH effects: compensation's `runtime_pid` -/// update AND the writer's `WRITER_EDIT` env-var sentinel. +/// Invariant: if the `managed_agents_store_lock` guard is dropped before +/// `compensate_drain_for`'s restore/save, the writer can interleave and +/// overwrite `COMP_SENTINEL` — making the final assertion fail. /// -/// Ordering is established by construction (channel), not by scheduler timing. -/// "both effects" proves the store guard is correctly held across restore/save. +/// Flow: +/// 1. Seed the store with one agent record. +/// 2. Acquire `managed_agent_runtime_transition`; call `compensate_drain_with_hook`. +/// 3. `on_store_acquired` hook (fires while the store lock is held): +/// a. writes `COMP_SENTINEL` to disk; +/// b. signals the writer thread (`store_acquired_tx`); +/// c. waits for the writer to queue at the store lock (`writer_queued_rx`). +/// 4. Writer thread: waits for (b), signals (c), then blocks on `managed_agents_store_lock`. +/// Compensation still holds the lock — writer is blocked. +/// 5. `compensate_drain_with_hook` continues: validate → load → restore → save. +/// Store lock still held throughout; writer remains blocked. +/// 6. Function returns; store lock released. Writer acquires it, writes `WRITER_EDIT`, saves. +/// 7. Final disk record must contain BOTH `COMP_SENTINEL` AND `WRITER_EDIT`. +/// +/// If the guard is dropped early, the writer acquires the lock before save, +/// overwrites `COMP_SENTINEL`, and the assertion at step 7 fails. #[test] fn test_compensate_drain_writer_vs_compensation_deterministic() { use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope}; @@ -98,16 +100,13 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { ) .unwrap(); - // Build a mock app so `compensate_drain` can reach `AppState`. let app = tauri::test::mock_builder() .manage(crate::app_state::build_app_state()) .build(tauri::test::mock_context(tauri::test::noop_assets())) .expect("failed to build mock app"); let app_handle = app.app_handle().clone(); - let state = app.state::(); - // Commit a scope pointing at our tempdir so generation validation passes. let gen = current_scope_generation(); let scope = WorkspaceAgentScope { scope_id: "comp-drain-writer-test".to_string(), @@ -121,21 +120,26 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let stopped = vec![entry1]; - // Channel: writer waits until compensate_drain signals "done". - let (comp_done_tx, comp_done_rx) = std::sync::mpsc::channel::<()>(); + // store_acquired: hook → writer (compensation holds the store lock; writer may queue) + // writer_queued: writer → hook (writer is now blocked on the store lock) + let (store_acquired_tx, store_acquired_rx) = std::sync::mpsc::channel::<()>(); + let (writer_queued_tx, writer_queued_rx) = std::sync::mpsc::channel::<()>(); - // Spawn the writer BEFORE acquiring the transition guard to avoid a - // deadlock — the writer only needs the store lock, not the transition lock, - // but it waits on the channel first. + // Spawn the writer thread BEFORE acquiring the transition guard. + // It waits for store_acquired, signals writer_queued (it is now about to + // block on the store lock), then acquires the store lock and writes WRITER_EDIT. let tmp_wr = tmp_path.clone(); let app_handle_wr = app_handle.clone(); let wr_thread = thread::spawn(move || { - // Established ordering: writer explicitly waits until compensation - // signals it is done — compensate_drain holds both transition guard - // (passed by value) and store lock through restore/save. - comp_done_rx.recv().unwrap(); + // Wait until the hook signals that compensation holds the store lock. + store_acquired_rx.recv().unwrap(); - // Now acquire just the store lock and apply the sentinel edit. + // Signal the hook that we are about to queue on the store lock. + // The hook will return on receipt, letting compensate_drain_with_hook + // proceed to validate → load → save while we are blocked below. + writer_queued_tx.send(()).unwrap(); + + // Block on the store lock — compensation still holds it. let writer_state = app_handle_wr.state::(); let _store = writer_state.managed_agents_store_lock.lock().unwrap(); let mut records = @@ -147,31 +151,40 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).unwrap(); }); - // Acquire the transition guard in this thread and pass it to the - // production adapter. compensate_drain takes ownership of the guard so it - // is held for the full validate→load→restore→save sequence. let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); - // ── Phase-A: run the production compensate_drain adapter ───────────────── - // `start_pair_under_held_locks` is not available to inject from outside - // the crate; compensate_drain calls it internally. Because there is no - // live process for the dummy record the start attempt will fail with an - // error — the adapter returns a degradation message. We verify the final - // disk state rather than the return value of compensate_drain. - // - // To prove the guard was actually held through save we assert the writer - // does NOT see intermediate state. - let _comp_result = compensate_drain(&app_handle, &stopped, &scope, rt_guard); + // on_store_acquired fires while the store lock is held: + // (a) write COMP_SENTINEL to disk — proves we are mid-sequence under the lock; + // (b) signal the writer (store_acquired_tx); + // (c) wait for the writer to queue at the lock boundary (writer_queued_rx). + // Returning lets compensate_drain_with_hook continue to validate → load → save, + // all while the store lock remains held, keeping the writer blocked. + let _comp_result = compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, || { + // (a) Write the compensation sentinel under the held store lock. + let mut records = + crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default(); + for r in &mut records { + r.env_vars + .insert("COMP_SENTINEL".to_string(), "yes".to_string()); + } + crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &records).unwrap(); - // Signal the writer that compensation is done (guard released). - comp_done_tx.send(()).unwrap(); + // (b) Tell the writer the store lock is held. + store_acquired_tx.send(()).unwrap(); + + // (c) Wait for the writer to queue at the store lock boundary. + writer_queued_rx.recv().unwrap(); + }); wr_thread.join().expect("writer thread panicked"); // ── Final disk state: BOTH effects must be present ─────────────────────── - // Compensation's effect: the record exists on disk (compensate_drain loaded - // it and saved after the restore attempt, regardless of start success). - // Writer's effect: WRITER_EDIT sentinel present. + // COMP_SENTINEL: written by the hook while compensation held the store lock. + // WRITER_EDIT: written by the writer after compensation released the lock. + // + // If the store guard is dropped before compensate_drain_for's save, the + // writer can acquire the lock early and overwrite COMP_SENTINEL — this + // assertion would then fail, proving the invariant. let final_records = crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default(); let final_rec = final_records @@ -179,32 +192,43 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { .find(|r| r.pubkey == pubkey1) .expect("agent record must be present on disk after both phases"); - // Writer's effect must always be present. + assert_eq!( + final_rec.env_vars.get("COMP_SENTINEL").map(String::as_str), + Some("yes"), + "COMP_SENTINEL must survive — fails if the store guard is dropped before \ + compensate_drain_for's save, allowing the writer to overwrite it" + ); assert_eq!( final_rec.env_vars.get("WRITER_EDIT").map(String::as_str), Some("yes"), - "writer's WRITER_EDIT sentinel must be present in final disk state" + "WRITER_EDIT must be present — the writer runs after compensation releases the store lock" ); } -/// Production start-path contender is blocked while `compensate_drain` holds -/// the `managed_agent_runtime_transition` mutex. +/// Production start-path contender is blocked while `compensate_drain_with_hook` +/// holds the `managed_agent_runtime_transition` mutex. +/// +/// Invariant: if the `managed_agent_runtime_transition` guard is dropped before +/// `compensate_drain_with_hook` finishes, `start_managed_agent_runtime_pair_lazy` +/// acquires the transition lock and proceeds — `contender_done_rx` would fire +/// before the hook completes, and the assertion at the end would fail. /// /// Flow: -/// 1. This test thread acquires `managed_agent_runtime_transition` FIRST. -/// 2. A "start contender" thread is spawned; it signals "alive," then blocks -/// trying to acquire the same transition mutex (this is the production -/// start-lock path — `start_managed_agent` and related production commands -/// also take the transition guard before modifying runtimes). -/// 3. `compensate_drain` is called with the already-held guard (passes it -/// by value into the adapter — drains the stopped-entry list). -/// 4. After `compensate_drain` returns the guard is consumed; the contender -/// acquires the mutex and signals "done." +/// 1. Seed the store; acquire `managed_agent_runtime_transition` in this thread. +/// 2. Spawn the contender. It signals "ready," then waits for the hook to +/// signal "store held." It then signals "at_lock" and calls the production +/// `start_managed_agent_runtime_pair_lazy` — which blocks acquiring +/// `managed_agent_runtime_transition` as its first action. +/// 3. `compensate_drain_with_hook` hook: signals the contender ("store held"), +/// waits for "at_lock" confirmation, then returns. The contender is now +/// queued at the transition lock. +/// 4. `compensate_drain_with_hook` completes, releases the transition guard. +/// 5. The contender acquires the lock (fails at spawn — no binary — but +/// passes the lock boundary), signals "done." /// -/// Contender/queue order is established by barriers/channels — no `sleep`. -/// The contender drives the same `managed_agent_runtime_transition` lock path -/// that production `start_managed_agent` commands use, proving compensation -/// blocks production starts. +/// The contender drives `start_managed_agent_runtime_pair_lazy`, the same +/// function called by the `start_managed_agent_runtime` Tauri command, proving +/// that compensation blocks production starts — not just a raw mutex lock. #[test] fn test_compensate_drain_concurrent_start_is_blocked() { use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope}; @@ -213,14 +237,74 @@ fn test_compensate_drain_concurrent_start_is_blocked() { let tmp = tempfile::tempdir().unwrap(); let tmp_path = tmp.path().to_path_buf(); - crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &[]).unwrap(); + + let pubkey1 = "aa".repeat(32); + let initial_record = crate::managed_agents::ManagedAgentRecord { + pubkey: pubkey1.clone(), + name: "contender-agent".to_string(), + display_name: None, + slug: None, + persona_id: None, + private_key_nsec: String::new(), + auth_tag: None, + relay_url: "wss://relay.example".to_string(), + avatar_url: None, + acp_command: crate::managed_agents::DEFAULT_ACP_COMMAND.to_string(), + agent_command: String::new(), + agent_command_override: None, + agent_args: vec![], + mcp_command: String::new(), + turn_timeout_seconds: 0, + idle_timeout_seconds: None, + max_turn_duration_seconds: None, + parallelism: 1, + system_prompt: None, + model: None, + provider: None, + persona_source_version: None, + env_vars: Default::default(), + start_on_app_launch: true, + auto_restart_on_config_change: false, + runtime_pid: None, + backend: crate::managed_agents::BackendKind::Local, + backend_agent_id: None, + provider_binary_path: None, + team_id: None, + persona_team_dir: None, + persona_name_in_team: None, + created_at: crate::util::now_iso(), + updated_at: crate::util::now_iso(), + last_started_at: None, + last_stopped_at: None, + last_exit_code: None, + last_error: None, + last_error_code: None, + respond_to: Default::default(), + respond_to_allowlist: Default::default(), + is_builtin: false, + is_active: true, + shared: false, + source_team: None, + source_team_persona_slug: None, + catalog_source: None, + definition_respond_to: None, + definition_respond_to_allowlist: Default::default(), + definition_parallelism: None, + relay_mesh: None, + runtime: None, + name_pool: vec![], + }; + crate::managed_agents::storage::save_managed_agents_at( + &tmp_path, + std::slice::from_ref(&initial_record), + ) + .unwrap(); let app = tauri::test::mock_builder() .manage(crate::app_state::build_app_state()) .build(tauri::test::mock_context(tauri::test::noop_assets())) .expect("failed to build mock app"); let app_handle = app.app_handle().clone(); - let state = app.state::(); let gen = current_scope_generation(); @@ -233,54 +317,65 @@ fn test_compensate_drain_concurrent_start_is_blocked() { }; state.commit_active_scope(scope.clone()); - let pubkey1 = "aa".repeat(32); let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let stopped = vec![entry1]; - // Channels for barrier-based ordering. + // contender_ready: contender → test (about to call production start seam) + // hook_store_held: hook → contender (hook has store lock; contender may proceed) + // contender_at_lock: contender → hook (contender queued at transition lock) + // contender_done: contender → test (contender passed the transition lock) let (contender_ready_tx, contender_ready_rx) = std::sync::mpsc::channel::<()>(); + let (hook_store_held_tx, hook_store_held_rx) = std::sync::mpsc::channel::<()>(); + let (contender_at_lock_tx, contender_at_lock_rx) = std::sync::mpsc::channel::<()>(); let (contender_done_tx, contender_done_rx) = std::sync::mpsc::channel::<()>(); - // Acquire the transition guard in THIS thread FIRST so the contender will - // block when it tries to acquire the same mutex. + // Acquire the transition guard FIRST so the contender will block on it. let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); - // The contender drives the transition mutex the same way the production - // `start_managed_agent` family does: it acquires the guard, then signals. let app_contender = app_handle.clone(); + let pubkey_contender = pubkey1.clone(); let contender = thread::spawn(move || { - // Signal that the contender is alive and about to block on the mutex. + // Signal: about to call the production start path. contender_ready_tx.send(()).unwrap(); - // This is the production start-lock path: takes managed_agent_runtime_transition. - let contender_state = app_contender.state::(); - let _guard = contender_state - .managed_agent_runtime_transition - .lock() - .unwrap(); - // Signal: contender now holds the lock (compensation is done). + + // Wait for the hook to confirm it holds the store lock. + hook_store_held_rx.recv().unwrap(); + + // Signal: queued at the transition lock boundary. + contender_at_lock_tx.send(()).unwrap(); + + // Production start-lock seam — acquires managed_agent_runtime_transition + // as its first action, so it blocks here until compensation releases it. + let _ = start_pair_lazy_for( + pubkey_contender, + "wss://relay.example".to_string(), + app_contender, + ); + + // Signal: passed the transition lock (compensation is done). contender_done_tx.send(()).unwrap(); }); - // Wait until the contender is alive and parked (or about to park) on the mutex. + // Wait for the contender to be alive before calling the adapter. contender_ready_rx.recv().unwrap(); - // Confirm the contender is NOT yet through the lock (it can't be — we hold it). - // Use try_recv: if it somehow returned it would mean the mutex isn't working. - assert!( - contender_done_rx.try_recv().is_err(), - "contender must be blocked while this thread holds the transition guard" - ); - - // ── Run the production compensate_drain adapter ─────────────────────────── - // Passes the transition guard BY VALUE — `compensate_drain` takes ownership - // and holds it through validate→load→restore→save. - let _comp_result = compensate_drain(&app_handle, &stopped, &scope, rt_guard); - // rt_guard is now consumed/dropped inside compensate_drain. + // on_store_acquired hook: fires while BOTH locks are held. + // (a) signal the contender so it proceeds to call the start seam; + // (b) wait for the contender to confirm it is at the lock boundary. + // After (b) the contender is guaranteed to be blocked on + // managed_agent_runtime_transition — returning lets + // compensate_drain_with_hook finish its work before the contender unblocks. + let _comp_result = compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, || { + hook_store_held_tx.send(()).unwrap(); + contender_at_lock_rx.recv().unwrap(); + }); // Contender must now acquire the mutex within a reasonable timeout. contender_done_rx .recv_timeout(std::time::Duration::from_secs(5)) - .expect("contender must unblock after compensate_drain releases the transition guard"); + .expect( + "contender must unblock after compensate_drain_with_hook releases the transition guard", + ); contender.join().expect("contender thread panicked"); }