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. /// only written inside `apply_workspace`'s prepare stage.
#[test] #[test]
fn test_import_before_first_apply_leaves_scope_none() { 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(); let state = build_app_state();
// Boot state: no active scope. // 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. /// commands fail closed until the frontend re-applies a workspace.
#[test] #[test]
fn test_live_import_with_active_scope_clears_scope_and_bumps_generation() { 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 state = build_app_state();
let base = std::env::temp_dir(); let base = std::env::temp_dir();
@@ -129,6 +135,9 @@ fn test_fallback_relay_never_claims_during_identity_import() {
/// original scope. /// original scope.
#[test] #[test]
fn test_prepare_failure_leaves_old_scope_intact() { 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 state = build_app_state();
let base = std::env::temp_dir(); 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. /// returning `None` is handled gracefully by callers that check the scope.
#[test] #[test]
fn test_inactive_runtime_exit_after_scope_cleared_is_safe() { 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 state = build_app_state();
let base = std::env::temp_dir(); 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_pid = Arc::new(Mutex::new(0u32));
let receipt_instance_id = Arc::new(Mutex::new(String::new())); let receipt_instance_id = Arc::new(Mutex::new(String::new()));
let receipt_started_at = 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 stop_called2 = stop_called.clone();
let spawn_relay2 = spawn_relay.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_pid2 = receipt_pid.clone();
let receipt_iid2 = receipt_instance_id.clone(); let receipt_iid2 = receipt_instance_id.clone();
let receipt_sat2 = receipt_started_at.clone(); let receipt_sat2 = receipt_started_at.clone();
let spawned_pid2 = spawned_child_pid.clone();
let pubkey2 = pubkey.clone(); let pubkey2 = pubkey.clone();
let app_handle_for_assert = app_handle.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_fn: record captured relay+owner+personas+teams+global, return a noop child.
// `spawn_noop_child_for_test()` replaces `/usr/bin/true` — cross-platform. // `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| { move |_app, rec, relay, owner, personas, global, teams| {
*spawn_relay2.lock().unwrap() = Some(relay.to_string()); *spawn_relay2.lock().unwrap() = Some(relay.to_string());
*spawn_owner2.lock().unwrap() = owner.map(str::to_string); *spawn_owner2.lock().unwrap() = owner.map(str::to_string);
*spawn_personas2.lock().unwrap() = !personas.is_empty(); *spawn_personas2.lock().unwrap() = !personas.is_empty();
*spawn_teams2.lock().unwrap() = !teams.is_empty(); *spawn_teams2.lock().unwrap() = !teams.is_empty();
*spawn_global2.lock().unwrap() = !global.env_vars.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 { Ok(crate::managed_agents::ManagedAgentProcess {
child: spawn_noop_child_for_test(), child,
log_path: std::path::PathBuf::new(), log_path: std::path::PathBuf::new(),
spawn_config: spawn_config:
crate::managed_agents::spawn_snapshot::prospective_spawn_config_snapshot( 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 \ "receipt.key.relay_url must match the captured relay; fails if the spawn path \
builds the receipt with an empty or wrong relay" builds the receipt with an empty or wrong relay"
); );
assert!( // Assert receipt.pid exactly matches the PID captured from the spawned child.
*receipt_pid.lock().unwrap() > 0, // Fails if production writes any PID other than `child.id()` into the receipt.
"receipt.pid must be non-zero (a real spawned child PID)" 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 // Assert receipt.desktop_instance_id equals the actual value produced by
// current_instance_id(app) — proves the field is sourced from the app // 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. /// guard must abort before stopping — no agent is touched.
#[test] #[test]
fn test_restart_under_captured_epoch_stale_scope_is_rejected() { 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(); let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("managed-agents.json"), b"[]").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; /// Thufir's test 1: "production async driver with injected loader failure;
/// assert stop never called, RestartOutcome::Skipped." /// assert stop never called, RestartOutcome::Skipped."
#[tokio::test] #[tokio::test]
#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests
async fn test_context_load_failure_leaves_runtime_running() { async fn test_context_load_failure_leaves_runtime_running() {
use super::restart_local_agent_on_config_change_for; use super::restart_local_agent_on_config_change_for;
use crate::commands::global_agent_config::RestartOutcome; 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(); let tmp = tempfile::tempdir().unwrap();
// Malformed JSON → load_agent_store_at (called by load_personas_at) returns // 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; /// Thufir's test 2: "real production driver/core with injected preflight error;
/// stop never called, RestartOutcome::Skipped." /// stop never called, RestartOutcome::Skipped."
#[tokio::test] #[tokio::test]
#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests
async fn test_mesh_preflight_failure_leaves_runtime_running() { async fn test_mesh_preflight_failure_leaves_runtime_running() {
use super::restart_local_agent_on_config_change_for; use super::restart_local_agent_on_config_change_for;
use crate::commands::global_agent_config::RestartOutcome; 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(); let tmp = tempfile::tempdir().unwrap();
// Provide persona and record files so context prep can pass them. // 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 /// Thufir's test 3: "injected preflight hook advances generation after it
/// succeeds; epoch returns Skipped; stop never called." /// succeeds; epoch returns Skipped; stop never called."
#[tokio::test] #[tokio::test]
#[allow(clippy::await_holding_lock)] // SCOPE_GENERATION_TEST_LOCK serialises parallel tests
async fn test_workspace_switch_after_preflight_aborts_before_stop() { async fn test_workspace_switch_after_preflight_aborts_before_stop() {
use super::restart_local_agent_on_config_change_for; use super::restart_local_agent_on_config_change_for;
use crate::commands::global_agent_config::RestartOutcome; 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(); let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("personas.json"), b"[]").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 /// old `context.mesh_model_id`, bypass the mismatch — stop would be called
/// and the test would fail. /// and the test would fail.
#[tokio::test] #[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() { async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
use super::super::restart_local_agent_on_config_change_for; use super::super::restart_local_agent_on_config_change_for;
use crate::commands::global_agent_config::RestartOutcome; use crate::commands::global_agent_config::RestartOutcome;
use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK;
use crate::managed_agents::{ use crate::managed_agents::{
storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord, storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord,
ManagedAgentRuntimeKey, ManagedAgentRuntimeKey,
}; };
use tauri::Manager; 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 = tempfile::tempdir().unwrap();
let pubkey = "aa".repeat(32); let pubkey = "aa".repeat(32);
let relay_url = "wss://relay.example"; let relay_url = "wss://relay.example";
@@ -908,36 +908,27 @@ where
/// Production adapter for [`compensate_drain_for`]. /// Production adapter for [`compensate_drain_for`].
/// ///
/// Delegates to [`compensate_drain_with_hook`] with a no-op hook. /// 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>( pub(crate) fn compensate_drain<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<'_, ()>,
) -> Option<String> { ) -> 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 /// Inner implementation of [`compensate_drain`] with an injectable `on_records_loaded` hook.
/// `on_store_acquired` hook.
/// ///
/// Lock acquisition order: /// Lock order: transition guard held by caller → acquire store lock → validate generation
/// 1. `_rt_transition_held` — already held by caller (passed by value). /// → load records → `on_records_loaded(&mut records)` (no-op in production; tests inject
/// 2. Acquire `managed_agents_store_lock`. /// a sentinel mutation) → delegate to `compensate_drain_for` → save records (step 7, still
/// 3. Call `on_store_acquired` (no-op in production; tests inject a closure /// under the store lock so the writer cannot interleave before the mutation hits disk).
/// that writes a sentinel and synchronises with a writer thread to prove
/// genuine store-lock contention).
/// 4. Validate `captured_scope.generation` under the store lock.
/// 5. `load_managed_agents_at` — records loaded AFTER both locks held.
/// 6. Delegate to [`compensate_drain_for`] — store guard held through every
/// save performed by `start_pair_under_held_locks`.
pub(crate) fn compensate_drain_with_hook<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_store_acquired: impl FnOnce(), on_records_loaded: impl FnOnce(&mut Vec<super::ManagedAgentRecord>),
) -> Option<String> { ) -> Option<String> {
if stopped.is_empty() { if stopped.is_empty() {
drop(_rt_transition_held); 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. // 3. Validate the captured scope under the store lock.
on_store_acquired();
// 4. Validate the captured scope under the store lock.
if let Err(stale_msg) = crate::managed_agents::scope::validate_scope_generation(captured_scope) if let Err(stale_msg) = crate::managed_agents::scope::validate_scope_generation(captured_scope)
{ {
return Some(format!( 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) { let mut records = match load_managed_agents_at(&captured_scope.definitions_dir) {
Ok(r) => r, Ok(r) => r,
Err(e) => { Err(e) => {
@@ -976,8 +964,13 @@ pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
} }
}; };
// 6. Delegate — store guard held through every save. // 5. Fire the hook with the loaded records. No-op in production.
compensate_drain_for(stopped, &mut records, |entry, recs| { // 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( start_pair_under_held_locks(
app, app,
&state, &state,
@@ -988,7 +981,17 @@ pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
recs, recs,
) )
.map(|_| ()) .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)] #[cfg(test)]
@@ -8,32 +8,47 @@ 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 the `managed_agents_store_lock` guard is dropped before /// Invariant: if `managed_agents_store_lock` is dropped before the adapter's
/// `compensate_drain_for`'s restore/save, the writer can interleave and /// step-7 save in `compensate_drain_with_hook`, the writer acquires the lock
/// overwrite `COMP_SENTINEL` — making the final assertion fail. /// before `COMP_SENTINEL` reaches disk, loads records without it, and saves
/// only `WRITER_EDIT` — making the `COMP_SENTINEL` assertion fail deterministically.
/// ///
/// Flow: /// 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`. /// 2. Acquire `managed_agent_runtime_transition`; call `compensate_drain_with_hook`.
/// 3. `on_store_acquired` hook (fires while the store lock is held): /// 3. `on_records_loaded` hook fires AFTER both locks are held:
/// a. writes `COMP_SENTINEL` to disk; /// a. signals the writer thread (`records_loaded_tx`);
/// b. signals the writer thread (`store_acquired_tx`); /// b. waits for writer's pre-lock signal (`writer_at_store_lock_rx`) — sent
/// c. waits for the writer to queue at the store lock (`writer_queued_rx`). /// immediately before `managed_agents_store_lock.lock()`, so when received,
/// 4. Writer thread: waits for (b), signals (c), then blocks on `managed_agents_store_lock`. /// the writer's next instruction is that lock call (which blocks);
/// Compensation still holds the lock — writer is blocked. /// c. mutates `COMP_SENTINEL` on the in-memory adapter-loaded records.
/// 5. `compensate_drain_with_hook` continues: validate → load → restore → save. /// 4. `compensate_drain_with_hook` continues: delegate to `compensate_drain_for`
/// Store lock still held throughout; writer remains blocked. /// (succeeds — AlreadyRunning), then saves the hook-mutated records at step 7.
/// 6. Function returns; store lock released. Writer acquires it, writes `WRITER_EDIT`, saves. /// Store lock held throughout; writer remains blocked on the store lock.
/// 7. Final disk record must contain BOTH `COMP_SENTINEL` AND `WRITER_EDIT`. /// 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, /// What breaks it: if the store guard is dropped before the adapter's step-7 save,
/// overwrites `COMP_SENTINEL`, and the assertion at step 7 fails. /// 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] #[test]
fn test_compensate_drain_writer_vs_compensation_deterministic() { 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 std::thread;
use tauri::Manager; 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 = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf(); 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 app_handle = app.app_handle().clone();
let state = app.state::<crate::app_state::AppState>(); 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 gen = current_scope_generation();
let scope = WorkspaceAgentScope { let scope = WorkspaceAgentScope {
scope_id: "comp-drain-writer-test".to_string(), 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 entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1]; let stopped = vec![entry1];
// store_acquired: hook → writer (compensation holds the store lock; writer may queue) // records_loaded: hook → writer (compensation holds store lock; writer may proceed)
// writer_queued: writer → hook (writer is now blocked on the store lock) // writer_at_store_lock: writer → hook (writer is about to call managed_agents_store_lock.lock())
let (store_acquired_tx, store_acquired_rx) = std::sync::mpsc::channel::<()>(); //
let (writer_queued_tx, writer_queued_rx) = std::sync::mpsc::channel::<()>(); // 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. // Spawn the writer thread BEFORE acquiring the transition guard.
// It waits for store_acquired, signals writer_queued (it is now about to // The writer directly acquires managed_agents_store_lock after signalling —
// block on the store lock), then acquires the store lock and writes WRITER_EDIT. // 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 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 the hook signals that compensation holds the store lock. // Wait until on_records_loaded signals that compensation holds the store lock.
store_acquired_rx.recv().unwrap(); records_loaded_rx.recv().unwrap();
// Signal the hook that we are about to queue on the store lock. // Signal the hook that we are about to call managed_agents_store_lock.lock().
// The hook will return on receipt, letting compensate_drain_with_hook // The hook's recv() will happen-before our lock call, so when on_records_loaded
// proceed to validate → load → save while we are blocked below. // mutates COMP_SENTINEL and returns, we are guaranteed to be BLOCKED on the
writer_queued_tx.send(()).unwrap(); // 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 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: 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 = 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 {
@@ -153,38 +215,66 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap(); let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
// on_store_acquired fires while the store lock is held: // on_records_loaded hook fires after load, while BOTH locks are held:
// (a) write COMP_SENTINEL to disk — proves we are mid-sequence under the lock; // (a) signal the writer thread that compensation holds the store lock;
// (b) signal the writer (store_acquired_tx); // (b) wait for the writer to confirm it is at the store-lock boundary —
// (c) wait for the writer to queue at the lock boundary (writer_queued_rx). // the writer sends this signal immediately BEFORE calling
// Returning lets compensate_drain_with_hook continue to validate → load → save, // managed_agents_store_lock.lock(), so by the time recv() returns
// all while the store lock remains held, keeping the writer blocked. // here, the writer is queued (blocked) on the store lock;
let _comp_result = compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, || { // (c) mutate COMP_SENTINEL in-memory on the adapter-loaded records.
// (a) Write the compensation sentinel under the held store lock. // The adapter then saves these hook-mutated records (step 7) before releasing
let mut records = // the store lock. Only then does the writer acquire the lock and read records.
crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default(); //
for r in &mut records { // What breaks it: if the store guard is dropped before the adapter's step-7
r.env_vars // save, the writer acquires managed_agents_store_lock before COMP_SENTINEL
.insert("COMP_SENTINEL".to_string(), "yes".to_string()); // reaches disk — the writer's final save omits COMP_SENTINEL, failing the
} // assertion below.
crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &records).unwrap(); 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. // (b) Wait for the writer to reach the store-lock boundary.
store_acquired_tx.send(()).unwrap(); // 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. // (c) Mutate COMP_SENTINEL in-memory. The adapter's step-7 save writes this
writer_queued_rx.recv().unwrap(); // 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"); 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 ─────────────────────── // ── Final disk state: BOTH effects must be present ───────────────────────
// COMP_SENTINEL: written by the hook while compensation held the store lock. // COMP_SENTINEL: set in-memory by the hook, saved by the adapter's step-7 save
// WRITER_EDIT: written by the writer after compensation released the lock. // 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 // If the store guard is dropped before the adapter's step-7 save, the writer
// writer can acquire the lock early and overwrite COMP_SENTINEL — this // acquires managed_agents_store_lock and loads records before COMP_SENTINEL
// assertion would then fail, proving the invariant. // 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
@@ -195,8 +285,9 @@ 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 survive — fails if the store guard is dropped before \ "COMP_SENTINEL must reach disk via the adapter's step-7 save while the store \
compensate_drain_for's save, allowing the writer to overwrite it" 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!( assert_eq!(
final_rec.env_vars.get("WRITER_EDIT").map(String::as_str), 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` /// Production start-path contender is blocked on `managed_agent_runtime_transition`
/// holds the `managed_agent_runtime_transition` mutex. /// 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 /// Invariant: if `managed_agent_runtime_transition` is removed from the start seam,
/// `compensate_drain_with_hook` finishes, `start_managed_agent_runtime_pair_lazy` /// the contender fires `on_transition_acquired` immediately after `on_before_transition`
/// acquires the transition lock and proceeds — `contender_done_rx` would fire /// (with zero lock contention) — before compensation has a chance to complete. The
/// before the hook completes, and the assertion at the end would fail. /// `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: /// Flow:
/// 1. Seed the store; acquire `managed_agent_runtime_transition` in this thread. /// 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 /// 2. Spawn the contender. It calls `start_pair_lazy_for_with_hook`:
/// signal "store held." It then signals "at_lock" and calls the production /// - `on_before_transition` fires (BEFORE the lock call) and sends
/// `start_managed_agent_runtime_pair_lazy` — which blocks acquiring /// `contender_at_boundary` — at this point the contender is inside the seam,
/// `managed_agent_runtime_transition` as its first action. /// about to block on `managed_agent_runtime_transition`.
/// 3. `compensate_drain_with_hook` hook: signals the contender ("store held"), /// 3. Wait for `contender_at_boundary` — the contender is now blocked at the
/// waits for "at_lock" confirmation, then returns. The contender is now /// transition lock because we hold it.
/// queued at the transition lock. /// 4. Call `compensate_drain_with_hook`. The `on_records_loaded` hook:
/// 4. `compensate_drain_with_hook` completes, releases the transition guard. /// (a) asserts `transition_hook_fired` is false (atomic bool sanity check);
/// 5. The contender acquires the lock (fails at spawn — no binary — but /// (b) asserts `transition_hook_check_rx.try_recv()` is `Err(Empty)` — the
/// passes the lock boundary), signals "done." /// 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 /// What breaks it: if `managed_agent_runtime_transition` is removed from
/// function called by the `start_managed_agent_runtime` Tauri command, proving /// `start_pair_for`, `on_transition_acquired` fires before step 4's assertion (because
/// that compensation blocks production starts — not just a raw mutex lock. /// 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] #[test]
fn test_compensate_drain_concurrent_start_is_blocked() { 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 std::thread;
use tauri::Manager; 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 = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf(); 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 entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1]; let stopped = vec![entry1];
// contender_ready: contender → test (about to call production start seam) // (A) contender_at_boundary: contender → test (inside seam, about to block on transition lock)
// hook_store_held: hook → contender (hook has store lock; contender may proceed) // (B) contender_transition_acquired: contender → test (start seam acquired transition guard)
// contender_at_lock: contender → hook (contender queued at transition lock) // Two receivers for the same semantic signal:
// contender_done: contender → test (contender passed the transition lock) // - transition_hook_check_rx: moved into on_records_loaded for try_recv() check
let (contender_ready_tx, contender_ready_rx) = std::sync::mpsc::channel::<()>(); // - contender_transition_acquired_rx: used after compensation for recv_timeout() check
let (hook_store_held_tx, hook_store_held_rx) = std::sync::mpsc::channel::<()>(); let (contender_at_boundary_tx, contender_at_boundary_rx) = std::sync::mpsc::channel::<()>();
let (contender_at_lock_tx, contender_at_lock_rx) = std::sync::mpsc::channel::<()>(); let (transition_hook_check_tx, transition_hook_check_rx) = std::sync::mpsc::channel::<()>();
let (contender_done_tx, contender_done_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. // 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();
@@ -335,46 +462,86 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
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 || {
// Signal: about to call the production start path. // `start_pair_lazy_for_with_hook` is the production-called seam.
contender_ready_tx.send(()).unwrap(); // 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
// Wait for the hook to confirm it holds the store lock. // (which the test holds). Signaling from here proves the contender is inside
hook_store_held_rx.recv().unwrap(); // the seam and will block on the next line — not merely about to call the function.
let _ = start_pair_lazy_for_with_hook(
// Signal: queued at the transition lock boundary.
contender_at_lock_tx.send(()).unwrap();
// Production start-lock seam — acquires managed_agent_runtime_transition
// as its first action, so it blocks here until compensation releases it.
let _ = start_pair_lazy_for(
pubkey_contender, 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.
// 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. // Wait for the contender to be inside the seam and about to block on the
contender_ready_rx.recv().unwrap(); // 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. // on_records_loaded 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. // Two complementary checks that on_transition_acquired has NOT fired:
// After (b) the contender is guaranteed to be blocked on // (1) Atomic bool: fast sanity check.
// managed_agent_runtime_transition — returning lets // (2) try_recv(): channel-based proof — on_transition_acquired sends to the
// compensate_drain_with_hook finish its work before the contender unblocks. // channel; if the channel is empty here, the hook has not fired.
let _comp_result = compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, || { //
hook_store_held_tx.send(()).unwrap(); // Why (2) is deterministic when the lock is removed: on_before_transition fires,
contender_at_lock_rx.recv().unwrap(); // 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. // After compensation returns, the contender can acquire the transition guard.
contender_done_rx // `on_transition_acquired` fires and sends the signal.
contender_transition_acquired_rx
.recv_timeout(std::time::Duration::from_secs(5)) .recv_timeout(std::time::Duration::from_secs(5))
.expect( .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"); contender.join().expect("contender thread panicked");
@@ -5,6 +5,112 @@
use super::*; 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( fn payload(
relay_url: &str, relay_url: &str,
lifecycle: ManagedAgentRuntimeLifecycle, lifecycle: ManagedAgentRuntimeLifecycle,
@@ -689,6 +795,9 @@ fn test_compensate_drain_empty_stopped_returns_none_with_real_app() {
/// generation guard fires before any spawn attempt. /// generation guard fires before any spawn attempt.
#[test] #[test]
fn test_compensate_drain_stale_scope_skips_all_with_real_app() { 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 tmp = tempfile::tempdir().unwrap();
let (app, scope) = build_mock_app_with_scope(&tmp); let (app, scope) = build_mock_app_with_scope(&tmp);
let app_handle = app.app_handle().clone(); 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. /// with the stale-scope test that also needs it, making the window near-zero.
#[test] #[test]
fn test_compensate_drain_attempts_restart_and_reports_degradation_with_real_app() { 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(); let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("managed-agents.json"), b"[]").unwrap(); std::fs::write(tmp.path().join("managed-agents.json"), b"[]").unwrap();
let app = tauri::test::mock_builder() let app = tauri::test::mock_builder()
@@ -341,6 +341,9 @@ mod tests {
#[test] #[test]
fn test_generation_increments_monotonically() { 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 before = current_scope_generation();
let next = next_scope_generation(); let next = next_scope_generation();
assert_eq!(next, before + 1); assert_eq!(next, before + 1);
@@ -356,6 +359,9 @@ mod tests {
/// check `captured != current` is reliable. /// check `captured != current` is reliable.
#[test] #[test]
fn test_generation_staleness_detected_after_scope_change() { 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 let captured = next_scope_generation(); // capture at operation entry
// Simulate a concurrent workspace switch bumping the generation. // Simulate a concurrent workspace switch bumping the generation.
let after_switch = next_scope_generation(); let after_switch = next_scope_generation();
@@ -377,6 +383,9 @@ mod tests {
/// fields correctly reflect the active workspace at each step. /// fields correctly reflect the active workspace at each step.
#[test] #[test]
fn test_scope_switch_a_to_b_to_a_advances_generation() { 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 base = std::env::temp_dir();
let owner = "aa".repeat(32); let owner = "aa".repeat(32);
@@ -411,6 +420,9 @@ mod tests {
/// staleness after C is committed. /// staleness after C is committed.
#[test] #[test]
fn test_rapid_scope_switch_a_b_c_all_stale_after_c() { 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 base = std::env::temp_dir();
let owner = "bb".repeat(32); let owner = "bb".repeat(32);
@@ -438,6 +450,9 @@ mod tests {
/// source of truth. /// source of truth.
#[test] #[test]
fn test_switch_during_restore_detected_by_generation_check() { 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. // Restore captures the generation at its entry.
let captured_at_restore_entry = next_scope_generation(); let captured_at_restore_entry = next_scope_generation();