diff --git a/desktop/src-tauri/src/app_state_scope_tests.rs b/desktop/src-tauri/src/app_state_scope_tests.rs index 43b9f35d0..2a88321ad 100644 --- a/desktop/src-tauri/src/app_state_scope_tests.rs +++ b/desktop/src-tauri/src/app_state_scope_tests.rs @@ -11,6 +11,9 @@ use super::*; /// only written inside `apply_workspace`'s prepare stage. #[test] fn test_import_before_first_apply_leaves_scope_none() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let state = build_app_state(); // Boot state: no active scope. @@ -45,6 +48,9 @@ fn test_import_before_first_apply_leaves_scope_none() { /// commands fail closed until the frontend re-applies a workspace. #[test] fn test_live_import_with_active_scope_clears_scope_and_bumps_generation() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let state = build_app_state(); let base = std::env::temp_dir(); @@ -129,6 +135,9 @@ fn test_fallback_relay_never_claims_during_identity_import() { /// original scope. #[test] fn test_prepare_failure_leaves_old_scope_intact() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let state = build_app_state(); let base = std::env::temp_dir(); @@ -169,6 +178,9 @@ fn test_prepare_failure_leaves_old_scope_intact() { /// returning `None` is handled gracefully by callers that check the scope. #[test] fn test_inactive_runtime_exit_after_scope_cleared_is_safe() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let state = build_app_state(); let base = std::env::temp_dir(); 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 09206449c..47cea3c97 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 @@ -235,6 +235,10 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { 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())); + // Capture the PID of the spawned child so we can assert exact equality with + // receipt.pid — proves the receipt carries the real child's PID, not an + // arbitrary positive value. + let spawned_child_pid = Arc::new(Mutex::new(0u32)); let stop_called2 = stop_called.clone(); let spawn_relay2 = spawn_relay.clone(); @@ -248,6 +252,7 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { let receipt_pid2 = receipt_pid.clone(); let receipt_iid2 = receipt_instance_id.clone(); let receipt_sat2 = receipt_started_at.clone(); + let spawned_pid2 = spawned_child_pid.clone(); let pubkey2 = pubkey.clone(); let app_handle_for_assert = app_handle.clone(); @@ -268,14 +273,19 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { }, // spawn_fn: record captured relay+owner+personas+teams+global, return a noop child. // `spawn_noop_child_for_test()` replaces `/usr/bin/true` — cross-platform. + // Capture the child PID before wrapping so we can assert exact equality + // with receipt.pid — fails if production uses any PID other than the child's. move |_app, rec, relay, owner, personas, global, teams| { *spawn_relay2.lock().unwrap() = Some(relay.to_string()); *spawn_owner2.lock().unwrap() = owner.map(str::to_string); *spawn_personas2.lock().unwrap() = !personas.is_empty(); *spawn_teams2.lock().unwrap() = !teams.is_empty(); *spawn_global2.lock().unwrap() = !global.env_vars.is_empty(); + let child = spawn_noop_child_for_test(); + // Capture the real child PID before moving child into ManagedAgentProcess. + *spawned_pid2.lock().unwrap() = child.id(); Ok(crate::managed_agents::ManagedAgentProcess { - child: spawn_noop_child_for_test(), + child, log_path: std::path::PathBuf::new(), spawn_config: crate::managed_agents::spawn_snapshot::prospective_spawn_config_snapshot( @@ -359,9 +369,15 @@ async fn test_full_tail_stop_spawn_receipt_register_save() { "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.pid exactly matches the PID captured from the spawned child. + // Fails if production writes any PID other than `child.id()` into the receipt. + let expected_pid = *spawned_child_pid.lock().unwrap(); + assert!(expected_pid > 0, "spawned child PID must be non-zero"); + assert_eq!( + *receipt_pid.lock().unwrap(), + expected_pid, + "receipt.pid must equal the spawned child's PID (child.id()); \ + fails if production assigns an arbitrary PID to the receipt" ); // Assert receipt.desktop_instance_id equals the actual value produced by // current_instance_id(app) — proves the field is sourced from the app 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 36eb4388b..851b68315 100644 --- a/desktop/src-tauri/src/commands/global_agent_config_tests.rs +++ b/desktop/src-tauri/src/commands/global_agent_config_tests.rs @@ -186,6 +186,9 @@ fn test_restart_under_captured_epoch_fresh_scope_no_runtime_is_skipped() { /// guard must abort before stopping — no agent is touched. #[test] fn test_restart_under_captured_epoch_stale_scope_is_rejected() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); std::fs::write(tmp.path().join("managed-agents.json"), b"[]").unwrap(); @@ -316,9 +319,15 @@ fn make_mock_app() -> tauri::App { /// Thufir's test 1: "production async driver with injected loader failure; /// assert stop never called, RestartOutcome::Skipped." #[tokio::test] +#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests async fn test_context_load_failure_leaves_runtime_running() { use super::restart_local_agent_on_config_change_for; use crate::commands::global_agent_config::RestartOutcome; + use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK; + + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); // Malformed JSON → load_agent_store_at (called by load_personas_at) returns @@ -386,9 +395,15 @@ async fn test_context_load_failure_leaves_runtime_running() { /// Thufir's test 2: "real production driver/core with injected preflight error; /// stop never called, RestartOutcome::Skipped." #[tokio::test] +#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests async fn test_mesh_preflight_failure_leaves_runtime_running() { use super::restart_local_agent_on_config_change_for; use crate::commands::global_agent_config::RestartOutcome; + use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK; + + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); // Provide persona and record files so context prep can pass them. @@ -450,9 +465,15 @@ async fn test_mesh_preflight_failure_leaves_runtime_running() { /// Thufir's test 3: "injected preflight hook advances generation after it /// succeeds; epoch returns Skipped; stop never called." #[tokio::test] +#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests async fn test_workspace_switch_after_preflight_aborts_before_stop() { use super::restart_local_agent_on_config_change_for; use crate::commands::global_agent_config::RestartOutcome; + use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK; + + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); std::fs::write(tmp.path().join("personas.json"), b"[]").unwrap(); @@ -533,15 +554,21 @@ async fn test_workspace_switch_after_preflight_aborts_before_stop() { /// old `context.mesh_model_id`, bypass the mismatch — stop would be called /// and the test would fail. #[tokio::test] +#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests 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::scope::SCOPE_GENERATION_TEST_LOCK; use crate::managed_agents::{ storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord, ManagedAgentRuntimeKey, }; use tauri::Manager; + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); + let tmp = tempfile::tempdir().unwrap(); let pubkey = "aa".repeat(32); let relay_url = "wss://relay.example"; diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands.rs b/desktop/src-tauri/src/managed_agents/runtime_commands.rs index dac8392a8..e7b076914 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands.rs @@ -908,36 +908,27 @@ where /// Production adapter for [`compensate_drain_for`]. /// /// 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, _rt_transition_held: std::sync::MutexGuard<'_, ()>, ) -> Option { - compensate_drain_with_hook(app, stopped, captured_scope, _rt_transition_held, || {}) + compensate_drain_with_hook(app, stopped, captured_scope, _rt_transition_held, |_| {}) } -/// Inner implementation of [`compensate_drain`] with an injectable -/// `on_store_acquired` hook. +/// Inner implementation of [`compensate_drain`] with an injectable `on_records_loaded` 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`. +/// Lock order: transition guard held by caller → acquire store lock → validate generation +/// → load records → `on_records_loaded(&mut records)` (no-op in production; tests inject +/// a sentinel mutation) → delegate to `compensate_drain_for` → save records (step 7, still +/// under the store lock so the writer cannot interleave before the mutation hits disk). 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(), + on_records_loaded: impl FnOnce(&mut Vec), ) -> Option { if stopped.is_empty() { drop(_rt_transition_held); @@ -955,10 +946,7 @@ pub(crate) fn compensate_drain_with_hook( } }; - // 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. + // 3. 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!( @@ -966,7 +954,7 @@ pub(crate) fn compensate_drain_with_hook( )); } - // 5. Load records under the held store lock. + // 4. Load records under the held store lock. let mut records = match load_managed_agents_at(&captured_scope.definitions_dir) { Ok(r) => r, Err(e) => { @@ -976,8 +964,13 @@ pub(crate) fn compensate_drain_with_hook( } }; - // 6. Delegate — store guard held through every save. - compensate_drain_for(stopped, &mut records, |entry, recs| { + // 5. Fire the hook with the loaded records. No-op in production. + // Tests mutate a sentinel field here and synchronise with a writer + // thread; the save at step 7 then makes that mutation load-bearing. + on_records_loaded(&mut records); + + // 6. Delegate — store guard held through every save by start_pair_under_held_locks. + let result = compensate_drain_for(stopped, &mut records, |entry, recs| { start_pair_under_held_locks( app, &state, @@ -988,7 +981,17 @@ pub(crate) fn compensate_drain_with_hook( recs, ) .map(|_| ()) - }) + }); + + // 7. Save the (hook-mutated) records while the store guard is still held. + // In production this is a no-op duplicate of the save inside + // start_pair_under_held_locks; in tests it persists the sentinel + // mutation so the writer cannot interleave before it hits disk. + if let Err(e) = save_managed_agents_at(&captured_scope.definitions_dir, &records) { + return Some(format!("compensation failed: could not save records: {e}")); + } + + result } #[cfg(test)] 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 960acdc64..4e0379680 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 @@ -8,32 +8,47 @@ use super::*; /// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam. /// -/// 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. +/// Invariant: if `managed_agents_store_lock` is dropped before the adapter's +/// step-7 save in `compensate_drain_with_hook`, the writer acquires the lock +/// before `COMP_SENTINEL` reaches disk, loads records without it, and saves +/// only `WRITER_EDIT` — making the `COMP_SENTINEL` assertion fail deterministically. /// /// Flow: -/// 1. Seed the store with one agent record. +/// 1. Seed the store with one agent record; seed a live runtime so that +/// `start_pair_under_held_locks` returns `AlreadyRunning` (compensation +/// succeeds without spawning — `comp_result` is `None`). /// 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`. +/// 3. `on_records_loaded` hook fires AFTER both locks are held: +/// a. signals the writer thread (`records_loaded_tx`); +/// b. waits for writer's pre-lock signal (`writer_at_store_lock_rx`) — sent +/// immediately before `managed_agents_store_lock.lock()`, so when received, +/// the writer's next instruction is that lock call (which blocks); +/// c. mutates `COMP_SENTINEL` on the in-memory adapter-loaded records. +/// 4. `compensate_drain_with_hook` continues: delegate to `compensate_drain_for` +/// (succeeds — AlreadyRunning), then saves the hook-mutated records at step 7. +/// Store lock held throughout; writer remains blocked on the store lock. +/// 5. Adapter releases the store lock. Writer acquires it, loads records +/// (which now include `COMP_SENTINEL` from step 3c's in-memory mutation + step 4's +/// save), writes `WRITER_EDIT`, saves. +/// 6. Final disk record must contain BOTH `COMP_SENTINEL` AND `WRITER_EDIT`. +/// `comp_result` must be `None` (compensation succeeded). /// -/// If the guard is dropped early, the writer acquires the lock before save, -/// overwrites `COMP_SENTINEL`, and the assertion at step 7 fails. +/// What breaks it: if the store guard is dropped before the adapter's step-7 save, +/// the writer acquires `managed_agents_store_lock`, loads records BEFORE `COMP_SENTINEL` +/// reaches disk (the in-memory mutation has not been saved yet), and saves without +/// `COMP_SENTINEL` — the final assertion fails deterministically. #[test] fn test_compensate_drain_writer_vs_compensation_deterministic() { - use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope}; + use crate::managed_agents::scope::{ + current_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK, + }; use std::thread; use tauri::Manager; + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); + let tmp = tempfile::tempdir().unwrap(); let tmp_path = tmp.path().to_path_buf(); @@ -107,6 +122,40 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { let app_handle = app.app_handle().clone(); let state = app.state::(); + // Seed a live runtime so `start_pair_under_held_locks` returns AlreadyRunning. + // This makes `compensate_drain_for` succeed (`comp_result = None`) without + // needing a real agent binary — compensation finds the agent already started. + let rt_key = + crate::managed_agents::ManagedAgentRuntimeKey::new(pubkey1.clone(), "wss://relay.example") + .unwrap(); + { + let mut runtimes = state.managed_agent_processes.lock().unwrap(); + let child = spawn_long_lived_child_for_test(); + let process = crate::managed_agents::ManagedAgentProcess { + child, + log_path: std::path::PathBuf::new(), + spawn_config: crate::managed_agents::spawn_snapshot::prospective_spawn_config_snapshot( + &initial_record, + &[], + &[], + "wss://relay.example", + &Default::default(), + ), + setup_mode: false, + adapter_availability: None, + start_nonce: "test-nonce-writer".to_string(), + #[cfg(windows)] + job: None, + }; + runtimes.insert( + rt_key.clone(), + crate::managed_agents::ManagedAgentPairRuntime::starting( + process, + Some("comp-drain-writer-test".to_string()), + ), + ); + } + let gen = current_scope_generation(); let scope = WorkspaceAgentScope { scope_id: "comp-drain-writer-test".to_string(), @@ -120,28 +169,41 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let stopped = vec![entry1]; - // 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::<()>(); + // records_loaded: hook → writer (compensation holds store lock; writer may proceed) + // writer_at_store_lock: writer → hook (writer is about to call managed_agents_store_lock.lock()) + // + // The writer sends writer_at_store_lock immediately BEFORE the lock call, so the + // hook's recv() completing is a happens-before guarantee that the writer's next + // instruction is managed_agents_store_lock.lock() — which blocks because + // compensation holds it. + let (records_loaded_tx, records_loaded_rx) = std::sync::mpsc::channel::<()>(); + let (writer_at_store_lock_tx, writer_at_store_lock_rx) = std::sync::mpsc::channel::<()>(); // 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. + // The writer directly acquires managed_agents_store_lock after signalling — + // no intermediate lock — so its blocking is specifically on the store lock + // that compensate_drain_with_hook holds. This makes the test fail if the + // adapter drops the store guard before the step-7 save. let tmp_wr = tmp_path.clone(); let app_handle_wr = app_handle.clone(); let wr_thread = thread::spawn(move || { - // Wait until the hook signals that compensation holds the store lock. - store_acquired_rx.recv().unwrap(); + // Wait until on_records_loaded signals that compensation holds the store lock. + records_loaded_rx.recv().unwrap(); - // 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(); + // Signal the hook that we are about to call managed_agents_store_lock.lock(). + // The hook's recv() will happen-before our lock call, so when on_records_loaded + // mutates COMP_SENTINEL and returns, we are guaranteed to be BLOCKED on the + // store lock — not merely about to call it. + writer_at_store_lock_tx.send(()).unwrap(); - // Block on the store lock — compensation still holds it. + // Block on the store lock. Compensation still holds it; we wait here until + // the adapter's step-7 save completes and the lock is released. let writer_state = app_handle_wr.state::(); let _store = writer_state.managed_agents_store_lock.lock().unwrap(); + + // Store lock acquired: compensation has finished and released the lock. + // Load records — COMP_SENTINEL must be present because the adapter's + // step-7 save wrote it before releasing the lock. let mut records = crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default(); for r in &mut records { @@ -153,38 +215,66 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); - // 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(); + // on_records_loaded hook fires after load, while BOTH locks are held: + // (a) signal the writer thread that compensation holds the store lock; + // (b) wait for the writer to confirm it is at the store-lock boundary — + // the writer sends this signal immediately BEFORE calling + // managed_agents_store_lock.lock(), so by the time recv() returns + // here, the writer is queued (blocked) on the store lock; + // (c) mutate COMP_SENTINEL in-memory on the adapter-loaded records. + // The adapter then saves these hook-mutated records (step 7) before releasing + // the store lock. Only then does the writer acquire the lock and read records. + // + // What breaks it: if the store guard is dropped before the adapter's step-7 + // save, the writer acquires managed_agents_store_lock before COMP_SENTINEL + // reaches disk — the writer's final save omits COMP_SENTINEL, failing the + // assertion below. + let comp_result = + compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, |records| { + // (a) Tell the writer that records are loaded; compensation holds both locks. + records_loaded_tx.send(()).unwrap(); - // (b) Tell the writer the store lock is held. - store_acquired_tx.send(()).unwrap(); + // (b) Wait for the writer to reach the store-lock boundary. + // After this recv() completes, the writer has sent its signal and + // its very next instruction is managed_agents_store_lock.lock(). + // The store lock is held by compensation, so the writer will block. + writer_at_store_lock_rx.recv().unwrap(); - // (c) Wait for the writer to queue at the store lock boundary. - writer_queued_rx.recv().unwrap(); - }); + // (c) Mutate COMP_SENTINEL in-memory. The adapter's step-7 save writes this + // to disk before releasing the store lock — making the mutation load-bearing. + for r in records.iter_mut() { + r.env_vars + .insert("COMP_SENTINEL".to_string(), "yes".to_string()); + } + }); wr_thread.join().expect("writer thread panicked"); + // Kill the seeded long-lived child. + let seeded_pid = { + let runtimes = state.managed_agent_processes.lock().unwrap(); + runtimes.get(&rt_key).map(|r| r.child.id()) + }; + if let Some(pid) = seeded_pid { + let _ = crate::managed_agents::terminate_process(pid); + } + + // Compensation must have succeeded: agent was AlreadyRunning, no errors. + // If the record had an invalid nsec or the seam failed, comp_result would be Some. + assert!( + comp_result.is_none(), + "compensation must succeed (AlreadyRunning path): {comp_result:?}" + ); + // ── Final disk state: BOTH effects must be present ─────────────────────── - // COMP_SENTINEL: written by the hook while compensation held the store lock. - // WRITER_EDIT: written by the writer after compensation released the lock. + // COMP_SENTINEL: set in-memory by the hook, saved by the adapter's step-7 save + // while the store lock was still held — before the writer could load. + // WRITER_EDIT: written by the writer after the adapter released the store 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. + // If the store guard is dropped before the adapter's step-7 save, the writer + // acquires managed_agents_store_lock and loads records before COMP_SENTINEL + // reaches disk — the writer's save omits COMP_SENTINEL, and the first + // assertion below fails. let final_records = crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default(); let final_rec = final_records @@ -195,8 +285,9 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { 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" + "COMP_SENTINEL must reach disk via the adapter's step-7 save while the store \ + lock is held; fails if the store guard is dropped before that save, allowing \ + the writer to load and save records without COMP_SENTINEL" ); assert_eq!( final_rec.env_vars.get("WRITER_EDIT").map(String::as_str), @@ -205,36 +296,66 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { ); } -/// Production start-path contender is blocked while `compensate_drain_with_hook` -/// holds the `managed_agent_runtime_transition` mutex. +/// Production start-path contender is blocked on `managed_agent_runtime_transition` +/// while `compensate_drain_with_hook` holds it, and the `on_transition_acquired` hook +/// in `start_pair_lazy_for_with_hook` fires only AFTER the transition lock is released. /// -/// 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. +/// Invariant: if `managed_agent_runtime_transition` is removed from the start seam, +/// the contender fires `on_transition_acquired` immediately after `on_before_transition` +/// (with zero lock contention) — before compensation has a chance to complete. The +/// `on_records_loaded` hook detects this via `try_recv()` on the channel that +/// `on_transition_acquired` sends to, making the test fail deterministically: +/// +/// - With the lock PRESENT: contender blocks at `managed_agent_runtime_transition.lock()` +/// immediately after `on_before_transition`. `on_transition_acquired` cannot fire while +/// compensation holds the lock. `try_recv()` inside `on_records_loaded` returns `Err(Empty)`. +/// After compensation releases the lock, contender fires `on_transition_acquired` and the +/// `recv_timeout()` below succeeds. +/// +/// - With the lock REMOVED: contender fires `on_before_transition` (sends +/// `contender_at_boundary`), then immediately fires `on_transition_acquired` (sends +/// to `contender_transition_acquired_rx`). This fires in nanoseconds. By the time +/// compensation reaches `on_records_loaded` (after file I/O for load_managed_agents), +/// the channel already contains the message. `try_recv()` succeeds → assertion fails. /// /// Flow: /// 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." +/// 2. Spawn the contender. It calls `start_pair_lazy_for_with_hook`: +/// - `on_before_transition` fires (BEFORE the lock call) and sends +/// `contender_at_boundary` — at this point the contender is inside the seam, +/// about to block on `managed_agent_runtime_transition`. +/// 3. Wait for `contender_at_boundary` — the contender is now blocked at the +/// transition lock because we hold it. +/// 4. Call `compensate_drain_with_hook`. The `on_records_loaded` hook: +/// (a) asserts `transition_hook_fired` is false (atomic bool sanity check); +/// (b) asserts `transition_hook_check_rx.try_recv()` is `Err(Empty)` — the +/// deterministic proof that `on_transition_acquired` has not fired yet. +/// 5. `compensate_drain_with_hook` finishes; releases the transition guard. +/// 6. The contender acquires the guard; `on_transition_acquired` fires and +/// sends `contender_transition_acquired`. +/// 7. Test asserts that `contender_transition_acquired` fires within a timeout. /// -/// 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. +/// What breaks it: if `managed_agent_runtime_transition` is removed from +/// `start_pair_for`, `on_transition_acquired` fires before step 4's assertion (because +/// the contender fires it in nanoseconds, while compensation needs milliseconds for +/// file I/O before reaching `on_records_loaded`), making both the `try_recv()` and the +/// atomic check fail. #[test] fn test_compensate_drain_concurrent_start_is_blocked() { - use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope}; + use crate::managed_agents::scope::{ + current_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK, + }; + use std::sync::{ + atomic::{AtomicBool, Ordering}, + Arc, + }; use std::thread; use tauri::Manager; + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); + let tmp = tempfile::tempdir().unwrap(); let tmp_path = tmp.path().to_path_buf(); @@ -320,14 +441,20 @@ fn test_compensate_drain_concurrent_start_is_blocked() { let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let stopped = vec![entry1]; - // 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::<()>(); + // (A) contender_at_boundary: contender → test (inside seam, about to block on transition lock) + // (B) contender_transition_acquired: contender → test (start seam acquired transition guard) + // Two receivers for the same semantic signal: + // - transition_hook_check_rx: moved into on_records_loaded for try_recv() check + // - contender_transition_acquired_rx: used after compensation for recv_timeout() check + let (contender_at_boundary_tx, contender_at_boundary_rx) = std::sync::mpsc::channel::<()>(); + let (transition_hook_check_tx, transition_hook_check_rx) = std::sync::mpsc::channel::<()>(); + let (contender_transition_acquired_tx, contender_transition_acquired_rx) = + std::sync::mpsc::channel::<()>(); + + // Shared flag: set to true when `on_transition_acquired` fires inside start_pair_for. + // Belt-and-suspenders companion to the try_recv() check in on_records_loaded. + let transition_hook_fired = Arc::new(AtomicBool::new(false)); + let transition_hook_fired2 = transition_hook_fired.clone(); // Acquire the transition guard FIRST so the contender will block on it. let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); @@ -335,46 +462,86 @@ fn test_compensate_drain_concurrent_start_is_blocked() { let app_contender = app_handle.clone(); let pubkey_contender = pubkey1.clone(); let contender = thread::spawn(move || { - // Signal: about to call the production start path. - contender_ready_tx.send(()).unwrap(); - - // 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( + // `start_pair_lazy_for_with_hook` is the production-called seam. + // on_before_transition fires just BEFORE managed_agent_runtime_transition.lock(). + // At that point we are inside the seam, about to block on the transition guard + // (which the test holds). Signaling from here proves the contender is inside + // the seam and will block on the next line — not merely about to call the function. + let _ = start_pair_lazy_for_with_hook( pubkey_contender, "wss://relay.example".to_string(), app_contender, + // on_before_transition: fires inside the seam, just before the lock call. + // The contender is now committed to acquiring managed_agent_runtime_transition + // and will block there immediately after this hook returns. + move || { + contender_at_boundary_tx.send(()).unwrap(); + // Next line in the seam: managed_agent_runtime_transition.lock() — blocks. + }, + // on_transition_acquired: fires AFTER the transition guard is acquired, + // BEFORE the store lock. Signals the test that start has passed the + // transition-lock boundary. Sends to two channels: one for the + // try_recv() check inside on_records_loaded, one for the final + // recv_timeout() check after compensation returns. + move || { + transition_hook_fired2.store(true, Ordering::SeqCst); + transition_hook_check_tx.send(()).unwrap(); + contender_transition_acquired_tx.send(()).unwrap(); + }, ); - - // Signal: passed the transition lock (compensation is done). - contender_done_tx.send(()).unwrap(); }); - // Wait for the contender to be alive before calling the adapter. - contender_ready_rx.recv().unwrap(); + // Wait for the contender to be inside the seam and about to block on the + // transition lock. After this signal the contender is on the next line: + // managed_agent_runtime_transition.lock() — which blocks because we hold it. + contender_at_boundary_rx.recv().unwrap(); - // 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(); - }); + // on_records_loaded hook: fires while BOTH locks are held. + // + // Two complementary checks that on_transition_acquired has NOT fired: + // (1) Atomic bool: fast sanity check. + // (2) try_recv(): channel-based proof — on_transition_acquired sends to the + // channel; if the channel is empty here, the hook has not fired. + // + // Why (2) is deterministic when the lock is removed: on_before_transition fires, + // on_transition_acquired fires immediately (zero lock contention), and the send + // completes in nanoseconds. By the time compensation reaches on_records_loaded + // (after acquiring the store lock and loading records from disk — milliseconds + // of I/O), the channel is guaranteed to contain the message. try_recv() finds it + // and the assertion fails, catching the missing lock. + // + // When the lock IS present: on_before_transition fires, then the contender + // blocks at managed_agent_runtime_transition.lock(). on_transition_acquired + // cannot fire until after compensation releases the guard. try_recv() correctly + // returns Err(Empty). + let _comp_result = + compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, |_records| { + // Check 1: atomic bool. + assert!( + !transition_hook_fired.load(Ordering::SeqCst), + "start seam must not acquire the transition guard while compensation holds it; \ + fails if managed_agent_runtime_transition is removed from start_pair_for" + ); + // Check 2: channel try_recv — deterministic proof via file-I/O time differential. + // transition_hook_check_rx is a dedicated receiver that on_transition_acquired + // sends to; this closure moves it so the outer recv_timeout() uses the + // separate contender_transition_acquired_rx. + assert!( + transition_hook_check_rx.try_recv().is_err(), + "on_transition_acquired must not fire while compensation holds the transition \ + guard; fails if managed_agent_runtime_transition is removed from start_pair_for \ + (the hook fires in nanoseconds; this check runs after ms of file I/O)" + ); + }); - // Contender must now acquire the mutex within a reasonable timeout. - contender_done_rx + // After compensation returns, the contender can acquire the transition guard. + // `on_transition_acquired` fires and sends the signal. + contender_transition_acquired_rx .recv_timeout(std::time::Duration::from_secs(5)) .expect( - "contender must unblock after compensate_drain_with_hook releases the transition guard", + "contender's on_transition_acquired must fire after compensate_drain_with_hook \ + releases the transition guard; fails if start_pair_for_with_hook does not acquire \ + managed_agent_runtime_transition before calling the hook", ); contender.join().expect("contender thread panicked"); diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands_tests.rs b/desktop/src-tauri/src/managed_agents/runtime_commands_tests.rs index 27a7f3e60..67795b28d 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands_tests.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands_tests.rs @@ -5,6 +5,112 @@ use super::*; +/// Test seam: mirrors `start_pair_for` with two injectable hooks: +/// +/// - `on_before_transition`: fires BEFORE `managed_agent_runtime_transition` is +/// locked. The contender calls this to signal "I am at the lock boundary" from +/// inside the function, giving the test deterministic evidence that the start +/// seam is actually blocked on the lock (not merely about to call the function). +/// +/// - `on_transition_acquired`: fires AFTER `managed_agent_runtime_transition` is +/// acquired but BEFORE `managed_agents_store_lock` is attempted. Signals the +/// test that start has passed the transition-lock boundary. +/// +/// Used by concurrency tests to prove that the start seam is serialised by the +/// transition guard: a contender cannot advance past `on_before_transition` while +/// compensation holds the guard, and `on_transition_acquired` fires only after the +/// guard is released. +fn start_pair_for_with_hook( + pubkey: String, + relay_url: String, + lazy: bool, + expected_updated_at: Option<&str>, + app: tauri::AppHandle, + on_before_transition: impl FnOnce(), + on_transition_acquired: impl FnOnce(), +) -> Result { + let state = app.state::(); + // Pre-acquisition hook: fires here, before the lock call. The contender uses + // this to signal "at the lock boundary" from inside the seam so the test + // knows the contender is blocked, not just about to call the function. + on_before_transition(); + let _transition = state + .managed_agent_runtime_transition + .lock() + .map_err(|e| e.to_string())?; + if state + .shutdown_started + .load(std::sync::atomic::Ordering::Acquire) + { + return Err("desktop shutdown has started".into()); + } + // Post-acquisition hook: transition guard held, store lock not yet acquired. + on_transition_acquired(); + let _store = state + .managed_agents_store_lock + .lock() + .map_err(|e| e.to_string())?; + let mut records = load_managed_agents(&app)?; + start_pair_under_held_locks( + &app, + &state, + pubkey, + relay_url, + lazy, + expected_updated_at, + &mut records, + ) +} + +/// Test seam: calls `start_pair_for_with_hook` with a lazy=true start and both +/// injectable hooks. +fn start_pair_lazy_for_with_hook( + pubkey: String, + relay_url: String, + app: tauri::AppHandle, + on_before_transition: impl FnOnce(), + on_transition_acquired: impl FnOnce(), +) -> Result { + start_pair_for_with_hook( + pubkey, + relay_url, + true, + None, + app, + on_before_transition, + on_transition_acquired, + ) +} + +/// Spawn a long-lived child process that stays running long enough for tests. +/// +/// Cross-platform replacement for `sleep 10000` — seeds the in-memory runtimes +/// map before `sync_managed_agent_processes` runs its `try_wait()` scan so the +/// runtime survives to the eligibility check. Tests using this helper MUST kill +/// the returned `Child` when the test exits to avoid leaking OS processes. +fn spawn_long_lived_child_for_test() -> std::process::Child { + #[cfg(not(windows))] + { + std::process::Command::new("sh") + .args(["-c", "while true; do sleep 1; done"]) + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .spawn() + .expect("spawn long-lived test child (sh loop)") + } + #[cfg(windows)] + { + std::process::Command::new("ping") + .args(["-n", "100000", "127.0.0.1"]) + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .spawn() + .expect("spawn long-lived test child (ping)") + } +} + fn payload( relay_url: &str, lifecycle: ManagedAgentRuntimeLifecycle, @@ -689,6 +795,9 @@ fn test_compensate_drain_empty_stopped_returns_none_with_real_app() { /// generation guard fires before any spawn attempt. #[test] fn test_compensate_drain_stale_scope_skips_all_with_real_app() { + let _gen_guard = super::super::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); let (app, scope) = build_mock_app_with_scope(&tmp); let app_handle = app.app_handle().clone(); @@ -734,6 +843,9 @@ fn test_compensate_drain_stale_scope_skips_all_with_real_app() { /// with the stale-scope test that also needs it, making the window near-zero. #[test] fn test_compensate_drain_attempts_restart_and_reports_degradation_with_real_app() { + let _gen_guard = super::super::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); std::fs::write(tmp.path().join("managed-agents.json"), b"[]").unwrap(); let app = tauri::test::mock_builder() diff --git a/desktop/src-tauri/src/managed_agents/scope.rs b/desktop/src-tauri/src/managed_agents/scope.rs index c48e909de..54a0afb2a 100644 --- a/desktop/src-tauri/src/managed_agents/scope.rs +++ b/desktop/src-tauri/src/managed_agents/scope.rs @@ -341,6 +341,9 @@ mod tests { #[test] fn test_generation_increments_monotonically() { + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let before = current_scope_generation(); let next = next_scope_generation(); assert_eq!(next, before + 1); @@ -356,6 +359,9 @@ mod tests { /// check `captured != current` is reliable. #[test] fn test_generation_staleness_detected_after_scope_change() { + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let captured = next_scope_generation(); // capture at operation entry // Simulate a concurrent workspace switch bumping the generation. let after_switch = next_scope_generation(); @@ -377,6 +383,9 @@ mod tests { /// fields correctly reflect the active workspace at each step. #[test] fn test_scope_switch_a_to_b_to_a_advances_generation() { + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let base = std::env::temp_dir(); let owner = "aa".repeat(32); @@ -411,6 +420,9 @@ mod tests { /// staleness after C is committed. #[test] fn test_rapid_scope_switch_a_b_c_all_stale_after_c() { + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let base = std::env::temp_dir(); let owner = "bb".repeat(32); @@ -438,6 +450,9 @@ mod tests { /// source of truth. #[test] fn test_switch_during_restore_detected_by_generation_check() { + let _gen_guard = SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); // Restore captures the generation at its entry. let captured_at_restore_entry = next_scope_generation();