test(desktop): resolve P4 pass-3 test-vs-claim defects in workspace-scoped agent store

Writer test: change writer to block directly on managed_agents_store_lock
(not via start_pair_lazy_for_with_hook) so the test fails deterministically
if the store guard is dropped before the step-7 save. The hook mutates
records in-memory; the step-7 save writes the mutation to disk while the
lock is still held; the writer loads after the lock is released.

Contender test: add try_recv() check on a dedicated channel inside
on_records_loaded alongside the existing atomic bool. on_transition_acquired
sends to both channels. try_recv() is deterministic: without the transition
lock, on_transition_acquired fires in nanoseconds; by the time on_records_loaded
runs (after file I/O), the channel contains the message. Documents the
time-differential argument explicitly in the test comment.

Full-tail PID: capture child.id() before wrapping into ManagedAgentProcess;
assert exact equality with receipt.pid (assert_eq! not assert! pid > 0).

SCOPE_GENERATION_TEST_LOCK: add to all remaining tests that mutate or
capture scope generation: app_state_scope_tests.rs (4 tests),
global_agent_config_tests.rs (4 tests), runtime_commands_tests.rs (2 tests),
scope.rs (5 generation tests). All tests that advance SCOPE_GENERATION
now participate in the process-global serialization mutex.

Trim runtime_commands.rs doc comments to stay at 1000 lines (gate limit).

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
2026-08-05 04:03:26 -04:00
co-authored by Will Pfleger
parent 27ea607243
commit 6739f157e2
7 changed files with 495 additions and 143 deletions
@@ -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();
@@ -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
@@ -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<tauri::test::MockRuntime> {
/// 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";
@@ -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<R: tauri::Runtime>(
app: &tauri::AppHandle<R>,
stopped: &[DrainJournalEntry],
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
_rt_transition_held: std::sync::MutexGuard<'_, ()>,
) -> Option<String> {
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<R: tauri::Runtime>(
app: &tauri::AppHandle<R>,
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<super::ManagedAgentRecord>),
) -> Option<String> {
if stopped.is_empty() {
drop(_rt_transition_held);
@@ -955,10 +946,7 @@ pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
}
};
// 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<R: tauri::Runtime>(
));
}
// 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<R: tauri::Runtime>(
}
};
// 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<R: tauri::Runtime>(
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)]
@@ -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::<crate::app_state::AppState>();
// 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::<crate::app_state::AppState>();
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");
@@ -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<R: tauri::Runtime>(
pubkey: String,
relay_url: String,
lazy: bool,
expected_updated_at: Option<&str>,
app: tauri::AppHandle<R>,
on_before_transition: impl FnOnce(),
on_transition_acquired: impl FnOnce(),
) -> Result<ManagedAgentRuntimeStatus, String> {
let state = app.state::<AppState>();
// 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<R: tauri::Runtime>(
pubkey: String,
relay_url: String,
app: tauri::AppHandle<R>,
on_before_transition: impl FnOnce(),
on_transition_acquired: impl FnOnce(),
) -> Result<ManagedAgentRuntimeStatus, String> {
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()
@@ -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();