test(managed-agents): structural guard evidence for concurrency lock invariants

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<ManagedAgentRecord>, &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 <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
Duncan
2026-08-05 16:56:31 -04:00
co-authored by Will Pfleger
parent d71cd62c7f
commit b450d46933
5 changed files with 118 additions and 248 deletions
@@ -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()` — /// This test verifies the structural invariant: `clear_active_scope()` —
/// the single operation identity import performs — does not touch the /// the single operation identity import performs — does not touch the
/// filesystem claim ledger. /// 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] #[test]
fn test_fallback_relay_never_claims_during_identity_import() { 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 tmp = tempfile::tempdir().unwrap();
let state = build_app_state(); let state = build_app_state();
@@ -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::commands::personas::snapshot::import::capture_agent_snapshot_import_entry;
use crate::managed_agents::scope::{ use crate::managed_agents::scope::{
current_scope_generation, next_scope_generation, WorkspaceAgentScope, current_scope_generation, next_scope_generation, WorkspaceAgentScope,
SCOPE_GENERATION_TEST_LOCK,
}; };
fn make_scope_with_keys(tmp: &tempfile::TempDir, owner_keys: &nostr::Keys) -> WorkspaceAgentScope { 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 /// Tests the production `validate_scope_generation` function directly — the
/// exact guard that fires inside Phase 3a of `confirm_agent_snapshot_import`. /// 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] #[test]
fn test_scope_generation_guard_rejects_stale_scope_for_import() { 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 tmp = tempfile::tempdir().unwrap();
let owner_keys = nostr::Keys::generate(); let owner_keys = nostr::Keys::generate();
@@ -250,7 +250,7 @@ pub(crate) fn start_pair_lazy_for<R: tauri::Runtime>(
relay_url: String, relay_url: String,
app: tauri::AppHandle<R>, app: tauri::AppHandle<R>,
) -> Result<ManagedAgentRuntimeStatus, String> { ) -> Result<ManagedAgentRuntimeStatus, String> {
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. // Start-pair hook seams: extracted to stay under the file-size ratchet.
@@ -291,7 +291,7 @@ fn start_pair_for<R: tauri::Runtime>(
expected_updated_at, expected_updated_at,
app, app,
|| {}, || {},
|| {}, |_| {},
) )
} }
@@ -909,30 +909,24 @@ pub(crate) fn compensate_drain<R: tauri::Runtime>(
captured_scope, captured_scope,
_rt_transition_held, _rt_transition_held,
|_| {}, |_| {},
|_| {}, |_, _| {},
) )
} }
/// Inner implementation of [`compensate_drain`] with injectable hooks. /// Inner implementation of [`compensate_drain`] with injectable hooks.
/// ///
/// Lock order: transition guard held by caller → acquire store lock → validate generation /// 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 /// → load records → `on_records_loaded` → delegate to `compensate_drain_for` →
/// a sentinel mutation) → delegate to `compensate_drain_for` → `on_after_restore(&mut records)` /// `on_after_restore(&mut records, &_store)` (borrows the actual store guard; production
/// (no-op in production; tests assert/save the injected mutation to prove the store lock is /// passes `|_, _| {}`). Store lock held across the entire sequence. With no-op callbacks,
/// held continuously through restore and the post-restore callback). /// production behavior is byte-for-byte equivalent to the pre-hook path.
///
/// 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).
pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>( pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
app: &tauri::AppHandle<R>, app: &tauri::AppHandle<R>,
stopped: &[DrainJournalEntry], stopped: &[DrainJournalEntry],
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope, captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
_rt_transition_held: std::sync::MutexGuard<'_, ()>, _rt_transition_held: std::sync::MutexGuard<'_, ()>,
on_records_loaded: impl FnOnce(&mut Vec<super::ManagedAgentRecord>), on_records_loaded: impl FnOnce(&mut Vec<super::ManagedAgentRecord>),
on_after_restore: impl FnOnce(&mut Vec<super::ManagedAgentRecord>), on_after_restore: impl FnOnce(&mut Vec<super::ManagedAgentRecord>, &std::sync::MutexGuard<'_, ()>),
) -> Option<String> { ) -> Option<String> {
if stopped.is_empty() { if stopped.is_empty() {
drop(_rt_transition_held); drop(_rt_transition_held);
@@ -986,9 +980,10 @@ pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
.map(|_| ()) .map(|_| ())
}); });
// 7. Post-restore hook — no-op in production. Tests save the injected mutation // 7. Post-restore hook — no-op in production. Tests add `COMP_SENTINEL` and save
// here to prove the store lock is held through restore and this callback. // while borrowing the actual store guard, proving the lock is held through restore
on_after_restore(&mut records); // and this callback.
on_after_restore(&mut records, &_store);
result result
} }
@@ -8,34 +8,20 @@ use super::*;
/// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam. /// 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 /// Invariant under test: `compensate_drain_with_hook` owns the production store guard at and
/// `on_after_restore` callback, the writer acquires the lock before `COMP_SENTINEL` /// through the post-restore persistence seam (`on_after_restore`). The callback receives a
/// reaches disk, loads records without it, and saves only `WRITER_EDIT` — /// borrow of the actual `MutexGuard<'_, ()>` returned by `managed_agents_store_lock.lock()`.
/// making the `COMP_SENTINEL` assertion fail deterministically. /// Dropping or moving that guard before the callback is a **compile error** — the borrow must
/// remain live at the call site.
/// ///
/// Flow: /// Runtime integration: with no-op callbacks, the store guard is held continuously; with the
/// 1. Seed the store with one agent record; seed a live runtime so that /// test callbacks, the writer completes one full transaction after the adapter releases the lock,
/// `start_pair_under_held_locks` returns `AlreadyRunning` (compensation /// and both disk effects (`COMP_SENTINEL` from `on_after_restore`, `WRITER_EDIT` from the writer)
/// succeeds without spawning — `comp_result` is `None`). /// must be present in the final state.
/// 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).
/// ///
/// What breaks it: if the store guard is dropped before `on_after_restore`, /// What breaks it (compile-time): deleting `managed_agents_store_lock.lock()` in
/// the writer acquires `managed_agents_store_lock` before `COMP_SENTINEL` /// `compensate_drain_with_hook` leaves no value for the `&_store` argument; dropping
/// reaches disk, loads records without it, and saves only `WRITER_EDIT` — /// or moving `_store` before `on_after_restore` makes the borrow-check invalid.
/// the first assertion fails deterministically.
#[test] #[test]
fn test_compensate_drain_writer_vs_compensation_deterministic() { fn test_compensate_drain_writer_vs_compensation_deterministic() {
use crate::managed_agents::scope::{ 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 entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1]; let stopped = vec![entry1];
// records_loaded: hook → writer (compensation holds store lock; writer may proceed) // records_loaded_tx: compensation → writer (both locks are held; writer may proceed)
// writer_at_store_lock: writer → hook (writer is about to call managed_agents_store_lock.lock()) // writer_committed_tx: writer → main (writer's save is complete)
//
// 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 (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. // Spawn the writer BEFORE acquiring the transition guard.
// The writer directly acquires managed_agents_store_lock after signalling — // Flow: wait for records_loaded → acquire managed_agents_store_lock (blocks until
// no intermediate lock — so its blocking is specifically on the store lock // on_after_restore releases it) → load → add WRITER_EDIT → save → send writer_committed.
// that compensate_drain_with_hook holds. This makes the test fail if the
// adapter drops the store guard before on_after_restore saves.
let tmp_wr = tmp_path.clone(); let tmp_wr = tmp_path.clone();
let app_handle_wr = app_handle.clone(); let app_handle_wr = app_handle.clone();
let wr_thread = thread::spawn(move || { 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(); records_loaded_rx.recv().unwrap();
// Signal the hook that we are about to call managed_agents_store_lock.lock(). // One complete production-shaped transaction: acquire the store lock (blocks
// The hook's recv() will happen-before our lock call, so when on_after_restore // until on_after_restore saves COMP_SENTINEL and the adapter releases the lock),
// mutates COMP_SENTINEL and saves, we are guaranteed to be BLOCKED on the // load records, add WRITER_EDIT, save.
// 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.
let writer_state = app_handle_wr.state::<crate::app_state::AppState>(); let writer_state = app_handle_wr.state::<crate::app_state::AppState>();
let _store = writer_state.managed_agents_store_lock.lock().unwrap(); 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 = let mut records =
crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default(); crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default();
for r in &mut records { 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()); .insert("WRITER_EDIT".to_string(), "yes".to_string());
} }
crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).unwrap(); 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(); 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( let comp_result = compensate_drain_with_hook(
&app_handle, &app_handle,
&stopped, &stopped,
&scope, &scope,
rt_guard, 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| { |_records| {
// (a) Tell the writer that records are loaded; compensation holds both locks.
records_loaded_tx.send(()).unwrap(); 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. // on_after_restore: borrows the actual store guard — a compile error if
// The writer is blocked on managed_agents_store_lock. This save reaching disk before // that guard is dropped or removed before this call. Add COMP_SENTINEL
// the lock is released is the load-bearing property: shortening the guard before this // and save while the guard is live; the writer is still blocked.
// callback lets the writer interleave and load records without COMP_SENTINEL. |records, _store_guard| {
|records| { let _ = _store_guard; // borrow is live; guard may not have been dropped
for r in records.iter_mut() { for r in records.iter_mut() {
r.env_vars r.env_vars
.insert("COMP_SENTINEL".to_string(), "yes".to_string()); .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"); wr_thread.join().expect("writer thread panicked");
// Kill the seeded long-lived child. // 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 ─────────────────────── // ── Final disk state: BOTH effects must be present ───────────────────────
// COMP_SENTINEL: mutated and saved by on_after_restore while the store lock // COMP_SENTINEL: saved by on_after_restore while the store guard was live.
// was still held — before the writer could acquire and load. // WRITER_EDIT: saved by the writer after the adapter released the lock.
// 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.
let final_records = let final_records =
crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default(); crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default();
let final_rec = final_records let final_rec = final_records
@@ -293,9 +244,7 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
assert_eq!( assert_eq!(
final_rec.env_vars.get("COMP_SENTINEL").map(String::as_str), final_rec.env_vars.get("COMP_SENTINEL").map(String::as_str),
Some("yes"), Some("yes"),
"COMP_SENTINEL must reach disk via on_after_restore while the store lock is held; \ "COMP_SENTINEL must reach disk via on_after_restore while the store guard is live"
fails if the store guard is dropped before on_after_restore, allowing \
the writer to load and save records without COMP_SENTINEL"
); );
assert_eq!( assert_eq!(
final_rec.env_vars.get("WRITER_EDIT").map(String::as_str), 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` /// Production start-path contender is blocked on `managed_agent_runtime_transition`
/// while `compensate_drain_with_hook` holds it, and the `on_transition_acquired` hook /// while `compensate_drain_with_hook` holds it; the `on_transition_acquired` seam
/// in `start_pair_lazy_for_with_hook` fires only AFTER the transition lock is released. /// 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 /// Invariant under test: `start_pair_for_with_hook` owns `managed_agent_runtime_transition`
/// `start_managed_agent_runtime_pair_lazy` → `start_pair_lazy_for` → `start_pair_lazy_for_with_hook` /// at the `on_transition_acquired` boundary. The callback receives a borrow of the actual
/// → `start_pair_for_with_hook`. Removing `managed_agent_runtime_transition` from the /// `MutexGuard<'_, ()>`. Dropping that guard before the callback or removing the lock
/// production path also removes it from this test seam, making the test a faithful proxy. /// call from the seam is a **compile error**.
/// ///
/// Invariant: if `managed_agent_runtime_transition` is removed from the start seam, /// Proof by mutex exclusion: compensation holds `managed_agent_runtime_transition` for the
/// the contender fires `on_transition_acquired` immediately after `on_before_transition` /// duration of `compensate_drain_with_hook`. The contender's `on_transition_acquired` borrows
/// (with zero lock contention) — before compensation has a chance to complete. The /// that same mutex's guard — it can only execute after compensation releases the mutex.
/// `on_records_loaded` hook detects this via `try_recv()` on the channel that /// No timing assumption is required; the scheduler cannot place the callback inside the
/// `on_transition_acquired` sends to, making the test fail deterministically: /// compensation window because mutex exclusion prevents it.
/// ///
/// - With the lock PRESENT: contender blocks at `managed_agent_runtime_transition.lock()` /// `start_pair_lazy_for_with_hook` is the production-callable seam. Removing
/// immediately after `on_before_transition`. `on_transition_acquired` cannot fire while /// `managed_agent_runtime_transition` from this path also removes it from production.
/// 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.
#[test] #[test]
fn test_compensate_drain_concurrent_start_is_blocked() { fn test_compensate_drain_concurrent_start_is_blocked() {
use crate::managed_agents::scope::{ use crate::managed_agents::scope::{
current_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK, current_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK,
}; };
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc,
};
use std::thread; use std::thread;
use tauri::Manager; use tauri::Manager;
@@ -454,119 +368,60 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true); let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1]; let stopped = vec![entry1];
// (A) contender_at_boundary: contender → test (inside seam, about to block on transition lock) // contender_at_boundary_tx: contender → test (inside seam, about to block on transition lock)
// (B) contender_transition_acquired: contender → test (start seam acquired transition guard) // transition_acquired_tx: contender → test (contender 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 (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 (transition_acquired_tx, transition_acquired_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_with_hook. // Acquire the transition guard FIRST so the contender blocks on it.
// 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(); let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
let app_contender = app_handle.clone(); let app_contender = app_handle.clone();
let pubkey_contender = pubkey1.clone(); let pubkey_contender = pubkey1.clone();
let contender = thread::spawn(move || { let contender = thread::spawn(move || {
// `start_pair_lazy_for_with_hook` is the production-called seam. // `start_pair_lazy_for_with_hook` is the production-called seam.
// on_before_transition fires just BEFORE managed_agent_runtime_transition.lock(). // on_before_transition fires just before the lock call — signals "entered the seam".
// At that point we are inside the seam, about to block on the transition guard // on_transition_acquired fires after the lock is acquired, receiving a borrow of the
// (which the test holds). Signaling from here proves the contender is inside // actual transition guard — sends one positive receipt to the test.
// the seam and will block on the next line — not merely about to call the function.
let _ = start_pair_lazy_for_with_hook( let _ = start_pair_lazy_for_with_hook(
pubkey_contender, pubkey_contender,
"wss://relay.example".to_string(), "wss://relay.example".to_string(),
app_contender, app_contender,
// on_before_transition: fires inside the seam, just before the lock call. // on_before_transition: contender is inside the seam and about to block.
// The contender is now committed to acquiring managed_agent_runtime_transition
// and will block there immediately after this hook returns.
move || { move || {
contender_at_boundary_tx.send(()).unwrap(); 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, // on_transition_acquired: borrows the actual transition guard at the seam boundary.
// BEFORE the store lock. Signals the test that start has passed the // Removing or dropping the guard before this call is a compile error.
// transition-lock boundary. Sends to two channels: one for the // Sends one positive receipt — proves the guard was acquired and live.
// try_recv() check inside on_records_loaded, one for the final move |_transition_guard| {
// recv_timeout() check after compensation returns. let _ = _transition_guard; // borrow is live; guard may not have been dropped
move || { transition_acquired_tx.send(()).unwrap();
transition_hook_fired2.store(true, Ordering::SeqCst);
transition_hook_check_tx.send(()).unwrap();
contender_transition_acquired_tx.send(()).unwrap();
}, },
); );
}); });
// Wait for the contender to be inside the seam and about to block on the // Wait for the contender to be inside the seam and blocked on the transition lock.
// 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(); contender_at_boundary_rx.recv().unwrap();
// on_records_loaded hook: fires while BOTH locks are held. // Run compensation while holding the transition guard. The contender cannot acquire it;
// // its on_transition_acquired callback cannot execute until compensation returns.
// 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( let _comp_result = compensate_drain_with_hook(
&app_handle, &app_handle,
&stopped, &stopped,
&scope, &scope,
rt_guard, 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| {},
|_records, _store_guard| {},
); );
// After compensation returns, the contender can acquire the transition guard. // After compensation returns, the contender can acquire the transition guard.
// `on_transition_acquired` fires and sends the signal. // Wait for the positive receipt from on_transition_acquired.
contender_transition_acquired_rx // Mutex exclusion guarantees this callback did not execute during compensation —
.recv_timeout(std::time::Duration::from_secs(5)) // no timing assumption is required.
.expect( transition_acquired_rx.recv().expect(
"contender's on_transition_acquired must fire after compensate_drain_with_hook \ "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 \ releases the transition guard",
managed_agent_runtime_transition before calling on_transition_acquired",
); );
contender.join().expect("contender thread panicked"); contender.join().expect("contender thread panicked");
@@ -17,7 +17,7 @@ pub(crate) fn start_pair_lazy_for_with_hook<R: tauri::Runtime>(
relay_url: String, relay_url: String,
app: tauri::AppHandle<R>, app: tauri::AppHandle<R>,
on_before_transition: impl FnOnce(), on_before_transition: impl FnOnce(),
on_transition_acquired: impl FnOnce(), on_transition_acquired: impl FnOnce(&std::sync::MutexGuard<'_, ()>),
) -> Result<ManagedAgentRuntimeStatus, String> { ) -> Result<ManagedAgentRuntimeStatus, String> {
start_pair_for_with_hook( start_pair_for_with_hook(
pubkey, pubkey,
@@ -33,14 +33,15 @@ pub(crate) fn start_pair_lazy_for_with_hook<R: tauri::Runtime>(
/// Generic start-pair seam with injectable hooks. /// Generic start-pair seam with injectable hooks.
/// ///
/// - `on_before_transition`: fires BEFORE `managed_agent_runtime_transition` is /// - `on_before_transition`: fires BEFORE `managed_agent_runtime_transition` is
/// locked. Tests signal "at the lock boundary" from here — the contender is /// locked. Tests signal "entered the seam" from here.
/// committed to acquiring `managed_agent_runtime_transition` immediately after.
/// ///
/// - `on_transition_acquired`: fires AFTER `managed_agent_runtime_transition` is /// - `on_transition_acquired`: fires AFTER `managed_agent_runtime_transition` is
/// acquired but BEFORE `managed_agents_store_lock` is attempted. Signals that /// acquired (borrowing the actual guard) but BEFORE `managed_agents_store_lock`
/// start has passed the transition-lock boundary. /// 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<R: tauri::Runtime>( pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
pubkey: String, pubkey: String,
relay_url: String, relay_url: String,
@@ -48,11 +49,11 @@ pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
expected_updated_at: Option<&str>, expected_updated_at: Option<&str>,
app: tauri::AppHandle<R>, app: tauri::AppHandle<R>,
on_before_transition: impl FnOnce(), on_before_transition: impl FnOnce(),
on_transition_acquired: impl FnOnce(), on_transition_acquired: impl FnOnce(&std::sync::MutexGuard<'_, ()>),
) -> Result<ManagedAgentRuntimeStatus, String> { ) -> Result<ManagedAgentRuntimeStatus, String> {
let state = app.state::<AppState>(); let state = app.state::<AppState>();
on_before_transition(); on_before_transition();
let _transition = state let transition_guard = state
.managed_agent_runtime_transition .managed_agent_runtime_transition
.lock() .lock()
.map_err(|e| e.to_string())?; .map_err(|e| e.to_string())?;
@@ -62,7 +63,7 @@ pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
{ {
return Err("desktop shutdown has started".into()); return Err("desktop shutdown has started".into());
} }
on_transition_acquired(); on_transition_acquired(&transition_guard);
let _store = state let _store = state
.managed_agents_store_lock .managed_agents_store_lock
.lock() .lock()