mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
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:
co-authored by
Will Pfleger
parent
27ea607243
commit
6739f157e2
@@ -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
|
||||||
|
// 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) 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) 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
|
r.env_vars
|
||||||
.insert("COMP_SENTINEL".to_string(), "yes".to_string());
|
.insert("COMP_SENTINEL".to_string(), "yes".to_string());
|
||||||
}
|
}
|
||||||
crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &records).unwrap();
|
|
||||||
|
|
||||||
// (b) Tell the writer the store lock is held.
|
|
||||||
store_acquired_tx.send(()).unwrap();
|
|
||||||
|
|
||||||
// (c) Wait for the writer to queue at the store lock boundary.
|
|
||||||
writer_queued_rx.recv().unwrap();
|
|
||||||
});
|
});
|
||||||
|
|
||||||
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();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user