From b450d469330a4ee2ffa069e2c2cee10caf54d786 Mon Sep 17 00:00:00 2001 From: Duncan Date: Wed, 5 Aug 2026 16:56:31 -0400 Subject: [PATCH] test(managed-agents): structural guard evidence for concurrency lock invariants MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace timing-based concurrency tests with compile-enforced structural guard ownership proofs per Thufir's binding shape consult. Item 1 (generation lock): add SCOPE_GENERATION_TEST_LOCK to test_fallback_relay_never_claims_during_identity_import and test_scope_generation_guard_rejects_stale_scope_for_import; exhaustive indirect-mutator sweep confirmed all remaining callers guarded. Item 2 (writer test): widen on_after_restore to FnOnce(&mut Vec, &MutexGuard<'_, ()>) — callback borrows the actual store guard. Dropping or removing the guard before the call is a compile error. Writer completes one full transaction (lock → load → WRITER_EDIT → save → writer_committed) after records_loaded releases it; on_after_restore saves COMP_SENTINEL under the live borrow. Delete writer_at_store_lock pre-lock signal and all timing-based comments. Item 3 (contender test): widen on_transition_acquired to FnOnce(&MutexGuard<'_, ()>) in both start_pair_lazy_for_with_hook and start_pair_for_with_hook — callback borrows the actual transition guard. Drop or removal before the call is a compile error. Remove AtomicBool, try_recv, not-fired-during assertion, and all ns-vs-ms timing rationale. Proof by mutex exclusion: on_transition_acquired cannot execute during compensation because both borrow the same mutex. Production delegates unchanged: compensate_drain passes |_, _| {}, start_pair_lazy_for and start_pair_for pass |_| {} for on_transition_acquired. Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- .../src-tauri/src/app_state_scope_tests.rs | 8 + .../personas/snapshot/tests_captured_scope.rs | 11 + .../src/managed_agents/runtime_commands.rs | 29 +- .../runtime_commands_concurrency_tests.rs | 299 +++++------------- .../managed_agents/runtime_commands_seams.rs | 19 +- 5 files changed, 118 insertions(+), 248 deletions(-) diff --git a/desktop/src-tauri/src/app_state_scope_tests.rs b/desktop/src-tauri/src/app_state_scope_tests.rs index 2a88321ad..5d6d7d46e 100644 --- a/desktop/src-tauri/src/app_state_scope_tests.rs +++ b/desktop/src-tauri/src/app_state_scope_tests.rs @@ -97,8 +97,16 @@ fn test_live_import_with_active_scope_clears_scope_and_bumps_generation() { /// This test verifies the structural invariant: `clear_active_scope()` — /// the single operation identity import performs — does not touch the /// filesystem claim ledger. +/// +/// Holds `SCOPE_GENERATION_TEST_LOCK` because `clear_active_scope()` internally +/// calls `next_scope_generation()`, which mutates the process-global generation +/// counter. Without the lock, this bump can race any test that relies on +/// generation stability (e.g., captured-scope stale-detection tests). #[test] fn test_fallback_relay_never_claims_during_identity_import() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); let state = build_app_state(); diff --git a/desktop/src-tauri/src/commands/personas/snapshot/tests_captured_scope.rs b/desktop/src-tauri/src/commands/personas/snapshot/tests_captured_scope.rs index a47e20f43..a15e81ab8 100644 --- a/desktop/src-tauri/src/commands/personas/snapshot/tests_captured_scope.rs +++ b/desktop/src-tauri/src/commands/personas/snapshot/tests_captured_scope.rs @@ -16,6 +16,7 @@ use crate::app_state::{build_app_state, AppState}; use crate::commands::personas::snapshot::import::capture_agent_snapshot_import_entry; use crate::managed_agents::scope::{ current_scope_generation, next_scope_generation, WorkspaceAgentScope, + SCOPE_GENERATION_TEST_LOCK, }; fn make_scope_with_keys(tmp: &tempfile::TempDir, owner_keys: &nostr::Keys) -> WorkspaceAgentScope { @@ -133,8 +134,18 @@ fn test_confirm_agent_snapshot_import_matching_owner_passes_entry_guard() { /// /// Tests the production `validate_scope_generation` function directly — the /// exact guard that fires inside Phase 3a of `confirm_agent_snapshot_import`. +/// +/// Holds `SCOPE_GENERATION_TEST_LOCK` because `next_scope_generation()` mutates +/// the process-global generation counter. Without the lock, the bump can race +/// any async test holding a captured generation (e.g., generation-stability +/// tests in `global_agent_config_epoch_tests` or +/// `runtime_commands_concurrency_tests`), causing spurious stale-detection +/// failures in the generation-sensitive path. #[test] fn test_scope_generation_guard_rejects_stale_scope_for_import() { + let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); let tmp = tempfile::tempdir().unwrap(); let owner_keys = nostr::Keys::generate(); diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands.rs b/desktop/src-tauri/src/managed_agents/runtime_commands.rs index 428fb36d3..ed3ae2503 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands.rs @@ -250,7 +250,7 @@ pub(crate) fn start_pair_lazy_for( relay_url: String, app: tauri::AppHandle, ) -> Result { - start_pair_lazy_for_with_hook(pubkey, relay_url, app, || {}, || {}) + start_pair_lazy_for_with_hook(pubkey, relay_url, app, || {}, |_| {}) } // Start-pair hook seams: extracted to stay under the file-size ratchet. @@ -291,7 +291,7 @@ fn start_pair_for( expected_updated_at, app, || {}, - || {}, + |_| {}, ) } @@ -909,30 +909,24 @@ pub(crate) fn compensate_drain( captured_scope, _rt_transition_held, |_| {}, - |_| {}, + |_, _| {}, ) } /// Inner implementation of [`compensate_drain`] with injectable hooks. /// /// 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` → `on_after_restore(&mut records)` -/// (no-op in production; tests assert/save the injected mutation to prove the store lock is -/// held continuously through restore and the post-restore callback). -/// -/// The store lock is held across the entire sequence: validate → load → pre-restore hook -/// → every restore → post-restore callback. The caller-owned transition guard is retained -/// across the whole interval. With no-op callbacks, production behavior is byte-for-byte -/// equivalent to the pre-hook path on all paths (empty, stale-generation, load-failure, -/// partial-restore, full-success). +/// → load records → `on_records_loaded` → delegate to `compensate_drain_for` → +/// `on_after_restore(&mut records, &_store)` (borrows the actual store guard; production +/// passes `|_, _| {}`). Store lock held across the entire sequence. With no-op callbacks, +/// production behavior is byte-for-byte equivalent to the pre-hook path. 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_records_loaded: impl FnOnce(&mut Vec), - on_after_restore: impl FnOnce(&mut Vec), + on_after_restore: impl FnOnce(&mut Vec, &std::sync::MutexGuard<'_, ()>), ) -> Option { if stopped.is_empty() { drop(_rt_transition_held); @@ -986,9 +980,10 @@ pub(crate) fn compensate_drain_with_hook( .map(|_| ()) }); - // 7. Post-restore hook — no-op in production. Tests save the injected mutation - // here to prove the store lock is held through restore and this callback. - on_after_restore(&mut records); + // 7. Post-restore hook — no-op in production. Tests add `COMP_SENTINEL` and save + // while borrowing the actual store guard, proving the lock is held through restore + // and this callback. + on_after_restore(&mut records, &_store); result } 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 a9ac70e80..634ecafb5 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,34 +8,20 @@ use super::*; /// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam. /// -/// Invariant: if `managed_agents_store_lock` is dropped before the adapter's -/// `on_after_restore` callback, 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. +/// Invariant under test: `compensate_drain_with_hook` owns the production store guard at and +/// through the post-restore persistence seam (`on_after_restore`). The callback receives a +/// borrow of the actual `MutexGuard<'_, ()>` returned by `managed_agents_store_lock.lock()`. +/// Dropping or moving that guard before the callback is a **compile error** — the borrow must +/// remain live at the call site. /// -/// Flow: -/// 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_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). -/// 4. `compensate_drain_for` runs under the held store lock (succeeds — AlreadyRunning). -/// 5. `on_after_restore` hook fires while the store lock is STILL held: -/// mutates `COMP_SENTINEL` on the adapter-loaded records and saves to disk -/// (assert/unwrap). Writer remains blocked. -/// 6. Adapter releases the store lock. Writer acquires it, loads records -/// (which now include `COMP_SENTINEL`), writes `WRITER_EDIT`, saves. -/// 7. Final disk record must contain BOTH `COMP_SENTINEL` AND `WRITER_EDIT`. -/// `comp_result` must be `None` (compensation succeeded). +/// Runtime integration: with no-op callbacks, the store guard is held continuously; with the +/// test callbacks, the writer completes one full transaction after the adapter releases the lock, +/// and both disk effects (`COMP_SENTINEL` from `on_after_restore`, `WRITER_EDIT` from the writer) +/// must be present in the final state. /// -/// What breaks it: if the store guard is dropped before `on_after_restore`, -/// the writer acquires `managed_agents_store_lock` before `COMP_SENTINEL` -/// reaches disk, loads records without it, and saves only `WRITER_EDIT` — -/// the first assertion fails deterministically. +/// What breaks it (compile-time): deleting `managed_agents_store_lock.lock()` in +/// `compensate_drain_with_hook` leaves no value for the `&_store` argument; dropping +/// or moving `_store` before `on_after_restore` makes the borrow-check invalid. #[test] fn test_compensate_drain_writer_vs_compensation_deterministic() { use crate::managed_agents::scope::{ @@ -168,41 +154,25 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let stopped = vec![entry1]; - // 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. + // records_loaded_tx: compensation → writer (both locks are held; writer may proceed) + // writer_committed_tx: writer → main (writer's save is complete) 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::<()>(); + let (writer_committed_tx, writer_committed_rx) = std::sync::mpsc::channel::<()>(); - // Spawn the writer thread BEFORE acquiring the transition guard. - // 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 on_after_restore saves. + // Spawn the writer BEFORE acquiring the transition guard. + // Flow: wait for records_loaded → acquire managed_agents_store_lock (blocks until + // on_after_restore releases it) → load → add WRITER_EDIT → save → send writer_committed. let tmp_wr = tmp_path.clone(); let app_handle_wr = app_handle.clone(); let wr_thread = thread::spawn(move || { - // Wait until on_records_loaded signals that compensation holds the store lock. + // Wait for on_records_loaded to signal that compensation holds both locks. records_loaded_rx.recv().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_after_restore - // mutates COMP_SENTINEL and saves, 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; we wait here until - // on_after_restore saves COMP_SENTINEL and the lock is released. + // One complete production-shaped transaction: acquire the store lock (blocks + // until on_after_restore saves COMP_SENTINEL and the adapter releases the lock), + // load records, add WRITER_EDIT, save. let writer_state = app_handle_wr.state::(); let _store = writer_state.managed_agents_store_lock.lock().unwrap(); - - // Store lock acquired: on_after_restore has finished and released the lock. - // Load records — COMP_SENTINEL must be present because on_after_restore - // saved 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 { @@ -210,45 +180,29 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { .insert("WRITER_EDIT".to_string(), "yes".to_string()); } crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).unwrap(); + + // Signal after the save succeeds, while this guarded transaction is complete. + writer_committed_tx.send(()).unwrap(); }); let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); - // on_records_loaded 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 immediately BEFORE managed_agents_store_lock.lock(), - // so by the time recv() returns, the writer's next instruction is that lock - // call (which blocks because compensation holds it). - // - // on_after_restore fires AFTER compensate_drain_for returns, still under BOTH locks: - // mutates COMP_SENTINEL in-memory and saves to disk (assert/unwrap). - // The store lock is still held — the writer remains blocked until this save returns - // and the lock is released. - // - // What breaks it: if the store guard is dropped before on_after_restore, the writer - // acquires managed_agents_store_lock before COMP_SENTINEL reaches disk — the writer's - // final save omits COMP_SENTINEL, and the first assertion below fails. let comp_result = compensate_drain_with_hook( &app_handle, &stopped, &scope, rt_guard, - // on_records_loaded: synchronise with the writer — both locks are held here. + // on_records_loaded: release the writer so it can queue on the store lock. + // Both locks are held here; the writer will block until on_after_restore + // releases the store guard. |_records| { - // (a) Tell the writer that records are loaded; compensation holds both locks. records_loaded_tx.send(()).unwrap(); - - // (b) Wait for the writer to reach the store-lock boundary. - // After recv() completes, the writer's next instruction is - // managed_agents_store_lock.lock() — which blocks because we hold it. - writer_at_store_lock_rx.recv().unwrap(); }, - // on_after_restore: save the COMP_SENTINEL mutation while BOTH locks are still held. - // The writer is blocked on managed_agents_store_lock. This save reaching disk before - // the lock is released is the load-bearing property: shortening the guard before this - // callback lets the writer interleave and load records without COMP_SENTINEL. - |records| { + // on_after_restore: borrows the actual store guard — a compile error if + // that guard is dropped or removed before this call. Add COMP_SENTINEL + // and save while the guard is live; the writer is still blocked. + |records, _store_guard| { + let _ = _store_guard; // borrow is live; guard may not have been dropped for r in records.iter_mut() { r.env_vars .insert("COMP_SENTINEL".to_string(), "yes".to_string()); @@ -258,6 +212,8 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { }, ); + // Block until the writer's transaction is complete, then join. + writer_committed_rx.recv().unwrap(); wr_thread.join().expect("writer thread panicked"); // Kill the seeded long-lived child. @@ -276,13 +232,8 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { ); // ── Final disk state: BOTH effects must be present ─────────────────────── - // COMP_SENTINEL: mutated and saved by on_after_restore while the store lock - // was still held — before the writer could acquire and load. - // WRITER_EDIT: written by the writer after the adapter released the lock. - // - // If the store guard is dropped before on_after_restore, the writer acquires - // managed_agents_store_lock before COMP_SENTINEL reaches disk — the writer's - // save omits COMP_SENTINEL, and the first assertion below fails. + // COMP_SENTINEL: saved by on_after_restore while the store guard was live. + // WRITER_EDIT: saved by the writer after the adapter released the lock. let final_records = crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default(); let final_rec = final_records @@ -293,9 +244,7 @@ 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 reach disk via on_after_restore while the store lock is held; \ - fails if the store guard is dropped before on_after_restore, allowing \ - the writer to load and save records without COMP_SENTINEL" + "COMP_SENTINEL must reach disk via on_after_restore while the store guard is live" ); assert_eq!( final_rec.env_vars.get("WRITER_EDIT").map(String::as_str), @@ -305,63 +254,28 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() { } /// 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. +/// while `compensate_drain_with_hook` holds it; the `on_transition_acquired` seam +/// fires after the guard is acquired, proving start cannot cross the +/// transition-acquired seam while compensation owns the mutex. /// -/// `start_pair_lazy_for_with_hook` is the production-callable seam: production -/// `start_managed_agent_runtime_pair_lazy` → `start_pair_lazy_for` → `start_pair_lazy_for_with_hook` -/// → `start_pair_for_with_hook`. Removing `managed_agent_runtime_transition` from the -/// production path also removes it from this test seam, making the test a faithful proxy. +/// Invariant under test: `start_pair_for_with_hook` owns `managed_agent_runtime_transition` +/// at the `on_transition_acquired` boundary. The callback receives a borrow of the actual +/// `MutexGuard<'_, ()>`. Dropping that guard before the callback or removing the lock +/// call from the seam is a **compile error**. /// -/// 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: +/// Proof by mutex exclusion: compensation holds `managed_agent_runtime_transition` for the +/// duration of `compensate_drain_with_hook`. The contender's `on_transition_acquired` borrows +/// that same mutex's guard — it can only execute after compensation releases the mutex. +/// No timing assumption is required; the scheduler cannot place the callback inside the +/// compensation window because mutex exclusion prevents it. /// -/// - 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 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. -/// -/// What breaks it: if `managed_agent_runtime_transition` is removed from -/// `start_pair_for_with_hook`, `on_transition_acquired` fires before step 4's assertion -/// (the contender fires it in nanoseconds; compensation needs milliseconds for -/// file I/O before reaching `on_records_loaded`), making both the `try_recv()` and the -/// atomic check fail. +/// `start_pair_lazy_for_with_hook` is the production-callable seam. Removing +/// `managed_agent_runtime_transition` from this path also removes it from production. #[test] fn test_compensate_drain_concurrent_start_is_blocked() { 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; @@ -454,120 +368,61 @@ fn test_compensate_drain_concurrent_start_is_blocked() { let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let stopped = vec![entry1]; - // (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 + // contender_at_boundary_tx: contender → test (inside seam, about to block on transition lock) + // transition_acquired_tx: contender → test (contender acquired transition guard) 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::<()>(); + let (transition_acquired_tx, transition_acquired_rx) = std::sync::mpsc::channel::<()>(); - // Shared flag: set to true when `on_transition_acquired` fires inside start_pair_for_with_hook. - // 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. + // Acquire the transition guard FIRST so the contender blocks on it. let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); let app_contender = app_handle.clone(); let pubkey_contender = pubkey1.clone(); let contender = thread::spawn(move || { // `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. + // on_before_transition fires just before the lock call — signals "entered the seam". + // on_transition_acquired fires after the lock is acquired, receiving a borrow of the + // actual transition guard — sends one positive receipt to the test. 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. + // on_before_transition: contender is inside the seam and about to block. 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(); + // on_transition_acquired: borrows the actual transition guard at the seam boundary. + // Removing or dropping the guard before this call is a compile error. + // Sends one positive receipt — proves the guard was acquired and live. + move |_transition_guard| { + let _ = _transition_guard; // borrow is live; guard may not have been dropped + transition_acquired_tx.send(()).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. + // Wait for the contender to be inside the seam and blocked on the transition lock. contender_at_boundary_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). + // Run compensation while holding the transition guard. The contender cannot acquire it; + // its on_transition_acquired callback cannot execute until compensation returns. let _comp_result = compensate_drain_with_hook( &app_handle, &stopped, &scope, rt_guard, - // on_records_loaded: fires while BOTH locks are held. Checks that - // on_transition_acquired has NOT fired — the contender is still blocked - // at managed_agent_runtime_transition.lock() because we hold it. - |_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_with_hook" - ); - // 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_with_hook \ - (the hook fires in nanoseconds; this check runs after ms of file I/O)" - ); - }, - // on_after_restore: no-op for the contender test — this test's load-bearing - // property is the transition guard, not the store-lock lifetime. |_records| {}, + |_records, _store_guard| {}, ); // 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'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 on_transition_acquired", - ); + // Wait for the positive receipt from on_transition_acquired. + // Mutex exclusion guarantees this callback did not execute during compensation — + // no timing assumption is required. + transition_acquired_rx.recv().expect( + "contender's on_transition_acquired must fire after compensate_drain_with_hook \ + releases the transition guard", + ); contender.join().expect("contender thread panicked"); } diff --git a/desktop/src-tauri/src/managed_agents/runtime_commands_seams.rs b/desktop/src-tauri/src/managed_agents/runtime_commands_seams.rs index 1790a6b61..9c9b7ab0d 100644 --- a/desktop/src-tauri/src/managed_agents/runtime_commands_seams.rs +++ b/desktop/src-tauri/src/managed_agents/runtime_commands_seams.rs @@ -17,7 +17,7 @@ pub(crate) fn start_pair_lazy_for_with_hook( relay_url: String, app: tauri::AppHandle, on_before_transition: impl FnOnce(), - on_transition_acquired: impl FnOnce(), + on_transition_acquired: impl FnOnce(&std::sync::MutexGuard<'_, ()>), ) -> Result { start_pair_for_with_hook( pubkey, @@ -33,14 +33,15 @@ pub(crate) fn start_pair_lazy_for_with_hook( /// Generic start-pair seam with injectable hooks. /// /// - `on_before_transition`: fires BEFORE `managed_agent_runtime_transition` is -/// locked. Tests signal "at the lock boundary" from here — the contender is -/// committed to acquiring `managed_agent_runtime_transition` immediately after. +/// locked. Tests signal "entered the seam" from here. /// /// - `on_transition_acquired`: fires AFTER `managed_agent_runtime_transition` is -/// acquired but BEFORE `managed_agents_store_lock` is attempted. Signals that -/// start has passed the transition-lock boundary. +/// acquired (borrowing the actual guard) but BEFORE `managed_agents_store_lock` +/// is attempted. The argument is a borrow of the actual transition guard; a +/// drop or removal of that guard before this call is a compile error. /// -/// Production delegates with `|| {}` for both hooks — release-invisible no-ops. +/// Production delegates with `|| {}` and `|_| {}` for the two hooks — +/// release-invisible no-ops. pub(crate) fn start_pair_for_with_hook( pubkey: String, relay_url: String, @@ -48,11 +49,11 @@ pub(crate) fn start_pair_for_with_hook( expected_updated_at: Option<&str>, app: tauri::AppHandle, on_before_transition: impl FnOnce(), - on_transition_acquired: impl FnOnce(), + on_transition_acquired: impl FnOnce(&std::sync::MutexGuard<'_, ()>), ) -> Result { let state = app.state::(); on_before_transition(); - let _transition = state + let transition_guard = state .managed_agent_runtime_transition .lock() .map_err(|e| e.to_string())?; @@ -62,7 +63,7 @@ pub(crate) fn start_pair_for_with_hook( { return Err("desktop shutdown has started".into()); } - on_transition_acquired(); + on_transition_acquired(&transition_guard); let _store = state .managed_agents_store_lock .lock()