mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
test(managed-agents): structural guard evidence for concurrency lock invariants
Replace timing-based concurrency tests with compile-enforced structural
guard ownership proofs per Thufir's binding shape consult.
Item 1 (generation lock): add SCOPE_GENERATION_TEST_LOCK to
test_fallback_relay_never_claims_during_identity_import and
test_scope_generation_guard_rejects_stale_scope_for_import; exhaustive
indirect-mutator sweep confirmed all remaining callers guarded.
Item 2 (writer test): widen on_after_restore to
FnOnce(&mut Vec<ManagedAgentRecord>, &MutexGuard<'_, ()>) — callback
borrows the actual store guard. Dropping or removing the guard before
the call is a compile error. Writer completes one full transaction
(lock → load → WRITER_EDIT → save → writer_committed) after records_loaded
releases it; on_after_restore saves COMP_SENTINEL under the live borrow.
Delete writer_at_store_lock pre-lock signal and all timing-based comments.
Item 3 (contender test): widen on_transition_acquired to
FnOnce(&MutexGuard<'_, ()>) in both start_pair_lazy_for_with_hook and
start_pair_for_with_hook — callback borrows the actual transition guard.
Drop or removal before the call is a compile error. Remove AtomicBool,
try_recv, not-fired-during assertion, and all ns-vs-ms timing rationale.
Proof by mutex exclusion: on_transition_acquired cannot execute during
compensation because both borrow the same mutex.
Production delegates unchanged: compensate_drain passes |_, _| {},
start_pair_lazy_for and start_pair_for pass |_| {} for on_transition_acquired.
Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
@@ -97,8 +97,16 @@ fn test_live_import_with_active_scope_clears_scope_and_bumps_generation() {
|
|||||||
/// This test verifies the structural invariant: `clear_active_scope()` —
|
/// This test verifies the structural invariant: `clear_active_scope()` —
|
||||||
/// the single operation identity import performs — does not touch the
|
/// the single operation identity import performs — does not touch the
|
||||||
/// filesystem claim ledger.
|
/// filesystem claim ledger.
|
||||||
|
///
|
||||||
|
/// Holds `SCOPE_GENERATION_TEST_LOCK` because `clear_active_scope()` internally
|
||||||
|
/// calls `next_scope_generation()`, which mutates the process-global generation
|
||||||
|
/// counter. Without the lock, this bump can race any test that relies on
|
||||||
|
/// generation stability (e.g., captured-scope stale-detection tests).
|
||||||
#[test]
|
#[test]
|
||||||
fn test_fallback_relay_never_claims_during_identity_import() {
|
fn test_fallback_relay_never_claims_during_identity_import() {
|
||||||
|
let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(|e| e.into_inner());
|
||||||
let tmp = tempfile::tempdir().unwrap();
|
let tmp = tempfile::tempdir().unwrap();
|
||||||
let state = build_app_state();
|
let state = build_app_state();
|
||||||
|
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ use crate::app_state::{build_app_state, AppState};
|
|||||||
use crate::commands::personas::snapshot::import::capture_agent_snapshot_import_entry;
|
use crate::commands::personas::snapshot::import::capture_agent_snapshot_import_entry;
|
||||||
use crate::managed_agents::scope::{
|
use crate::managed_agents::scope::{
|
||||||
current_scope_generation, next_scope_generation, WorkspaceAgentScope,
|
current_scope_generation, next_scope_generation, WorkspaceAgentScope,
|
||||||
|
SCOPE_GENERATION_TEST_LOCK,
|
||||||
};
|
};
|
||||||
|
|
||||||
fn make_scope_with_keys(tmp: &tempfile::TempDir, owner_keys: &nostr::Keys) -> WorkspaceAgentScope {
|
fn make_scope_with_keys(tmp: &tempfile::TempDir, owner_keys: &nostr::Keys) -> WorkspaceAgentScope {
|
||||||
@@ -133,8 +134,18 @@ fn test_confirm_agent_snapshot_import_matching_owner_passes_entry_guard() {
|
|||||||
///
|
///
|
||||||
/// Tests the production `validate_scope_generation` function directly — the
|
/// Tests the production `validate_scope_generation` function directly — the
|
||||||
/// exact guard that fires inside Phase 3a of `confirm_agent_snapshot_import`.
|
/// exact guard that fires inside Phase 3a of `confirm_agent_snapshot_import`.
|
||||||
|
///
|
||||||
|
/// Holds `SCOPE_GENERATION_TEST_LOCK` because `next_scope_generation()` mutates
|
||||||
|
/// the process-global generation counter. Without the lock, the bump can race
|
||||||
|
/// any async test holding a captured generation (e.g., generation-stability
|
||||||
|
/// tests in `global_agent_config_epoch_tests` or
|
||||||
|
/// `runtime_commands_concurrency_tests`), causing spurious stale-detection
|
||||||
|
/// failures in the generation-sensitive path.
|
||||||
#[test]
|
#[test]
|
||||||
fn test_scope_generation_guard_rejects_stale_scope_for_import() {
|
fn test_scope_generation_guard_rejects_stale_scope_for_import() {
|
||||||
|
let _gen_guard = crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(|e| e.into_inner());
|
||||||
let tmp = tempfile::tempdir().unwrap();
|
let tmp = tempfile::tempdir().unwrap();
|
||||||
let owner_keys = nostr::Keys::generate();
|
let owner_keys = nostr::Keys::generate();
|
||||||
|
|
||||||
|
|||||||
@@ -250,7 +250,7 @@ pub(crate) fn start_pair_lazy_for<R: tauri::Runtime>(
|
|||||||
relay_url: String,
|
relay_url: String,
|
||||||
app: tauri::AppHandle<R>,
|
app: tauri::AppHandle<R>,
|
||||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||||
start_pair_lazy_for_with_hook(pubkey, relay_url, app, || {}, || {})
|
start_pair_lazy_for_with_hook(pubkey, relay_url, app, || {}, |_| {})
|
||||||
}
|
}
|
||||||
|
|
||||||
// Start-pair hook seams: extracted to stay under the file-size ratchet.
|
// Start-pair hook seams: extracted to stay under the file-size ratchet.
|
||||||
@@ -291,7 +291,7 @@ fn start_pair_for<R: tauri::Runtime>(
|
|||||||
expected_updated_at,
|
expected_updated_at,
|
||||||
app,
|
app,
|
||||||
|| {},
|
|| {},
|
||||||
|| {},
|
|_| {},
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -909,30 +909,24 @@ pub(crate) fn compensate_drain<R: tauri::Runtime>(
|
|||||||
captured_scope,
|
captured_scope,
|
||||||
_rt_transition_held,
|
_rt_transition_held,
|
||||||
|_| {},
|
|_| {},
|
||||||
|_| {},
|
|_, _| {},
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Inner implementation of [`compensate_drain`] with injectable hooks.
|
/// Inner implementation of [`compensate_drain`] with injectable hooks.
|
||||||
///
|
///
|
||||||
/// Lock order: transition guard held by caller → acquire store lock → validate generation
|
/// Lock order: transition guard held by caller → acquire store lock → validate generation
|
||||||
/// → load records → `on_records_loaded(&mut records)` (no-op in production; tests inject
|
/// → load records → `on_records_loaded` → delegate to `compensate_drain_for` →
|
||||||
/// a sentinel mutation) → delegate to `compensate_drain_for` → `on_after_restore(&mut records)`
|
/// `on_after_restore(&mut records, &_store)` (borrows the actual store guard; production
|
||||||
/// (no-op in production; tests assert/save the injected mutation to prove the store lock is
|
/// passes `|_, _| {}`). Store lock held across the entire sequence. With no-op callbacks,
|
||||||
/// held continuously through restore and the post-restore callback).
|
/// production behavior is byte-for-byte equivalent to the pre-hook path.
|
||||||
///
|
|
||||||
/// The store lock is held across the entire sequence: validate → load → pre-restore hook
|
|
||||||
/// → every restore → post-restore callback. The caller-owned transition guard is retained
|
|
||||||
/// across the whole interval. With no-op callbacks, production behavior is byte-for-byte
|
|
||||||
/// equivalent to the pre-hook path on all paths (empty, stale-generation, load-failure,
|
|
||||||
/// partial-restore, full-success).
|
|
||||||
pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
|
pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
|
||||||
app: &tauri::AppHandle<R>,
|
app: &tauri::AppHandle<R>,
|
||||||
stopped: &[DrainJournalEntry],
|
stopped: &[DrainJournalEntry],
|
||||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||||
_rt_transition_held: std::sync::MutexGuard<'_, ()>,
|
_rt_transition_held: std::sync::MutexGuard<'_, ()>,
|
||||||
on_records_loaded: impl FnOnce(&mut Vec<super::ManagedAgentRecord>),
|
on_records_loaded: impl FnOnce(&mut Vec<super::ManagedAgentRecord>),
|
||||||
on_after_restore: impl FnOnce(&mut Vec<super::ManagedAgentRecord>),
|
on_after_restore: impl FnOnce(&mut Vec<super::ManagedAgentRecord>, &std::sync::MutexGuard<'_, ()>),
|
||||||
) -> Option<String> {
|
) -> Option<String> {
|
||||||
if stopped.is_empty() {
|
if stopped.is_empty() {
|
||||||
drop(_rt_transition_held);
|
drop(_rt_transition_held);
|
||||||
@@ -986,9 +980,10 @@ pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
|
|||||||
.map(|_| ())
|
.map(|_| ())
|
||||||
});
|
});
|
||||||
|
|
||||||
// 7. Post-restore hook — no-op in production. Tests save the injected mutation
|
// 7. Post-restore hook — no-op in production. Tests add `COMP_SENTINEL` and save
|
||||||
// here to prove the store lock is held through restore and this callback.
|
// while borrowing the actual store guard, proving the lock is held through restore
|
||||||
on_after_restore(&mut records);
|
// and this callback.
|
||||||
|
on_after_restore(&mut records, &_store);
|
||||||
|
|
||||||
result
|
result
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,34 +8,20 @@ use super::*;
|
|||||||
|
|
||||||
/// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam.
|
/// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam.
|
||||||
///
|
///
|
||||||
/// Invariant: if `managed_agents_store_lock` is dropped before the adapter's
|
/// Invariant under test: `compensate_drain_with_hook` owns the production store guard at and
|
||||||
/// `on_after_restore` callback, the writer acquires the lock before `COMP_SENTINEL`
|
/// through the post-restore persistence seam (`on_after_restore`). The callback receives a
|
||||||
/// reaches disk, loads records without it, and saves only `WRITER_EDIT` —
|
/// borrow of the actual `MutexGuard<'_, ()>` returned by `managed_agents_store_lock.lock()`.
|
||||||
/// making the `COMP_SENTINEL` assertion fail deterministically.
|
/// Dropping or moving that guard before the callback is a **compile error** — the borrow must
|
||||||
|
/// remain live at the call site.
|
||||||
///
|
///
|
||||||
/// Flow:
|
/// Runtime integration: with no-op callbacks, the store guard is held continuously; with the
|
||||||
/// 1. Seed the store with one agent record; seed a live runtime so that
|
/// test callbacks, the writer completes one full transaction after the adapter releases the lock,
|
||||||
/// `start_pair_under_held_locks` returns `AlreadyRunning` (compensation
|
/// and both disk effects (`COMP_SENTINEL` from `on_after_restore`, `WRITER_EDIT` from the writer)
|
||||||
/// succeeds without spawning — `comp_result` is `None`).
|
/// must be present in the final state.
|
||||||
/// 2. Acquire `managed_agent_runtime_transition`; call `compensate_drain_with_hook`.
|
|
||||||
/// 3. `on_records_loaded` hook fires AFTER both locks are held:
|
|
||||||
/// a. signals the writer thread (`records_loaded_tx`);
|
|
||||||
/// b. waits for writer's pre-lock signal (`writer_at_store_lock_rx`) — sent
|
|
||||||
/// immediately before `managed_agents_store_lock.lock()`, so when received,
|
|
||||||
/// the writer's next instruction is that lock call (which blocks).
|
|
||||||
/// 4. `compensate_drain_for` runs under the held store lock (succeeds — AlreadyRunning).
|
|
||||||
/// 5. `on_after_restore` hook fires while the store lock is STILL held:
|
|
||||||
/// mutates `COMP_SENTINEL` on the adapter-loaded records and saves to disk
|
|
||||||
/// (assert/unwrap). Writer remains blocked.
|
|
||||||
/// 6. Adapter releases the store lock. Writer acquires it, loads records
|
|
||||||
/// (which now include `COMP_SENTINEL`), writes `WRITER_EDIT`, saves.
|
|
||||||
/// 7. Final disk record must contain BOTH `COMP_SENTINEL` AND `WRITER_EDIT`.
|
|
||||||
/// `comp_result` must be `None` (compensation succeeded).
|
|
||||||
///
|
///
|
||||||
/// What breaks it: if the store guard is dropped before `on_after_restore`,
|
/// What breaks it (compile-time): deleting `managed_agents_store_lock.lock()` in
|
||||||
/// the writer acquires `managed_agents_store_lock` before `COMP_SENTINEL`
|
/// `compensate_drain_with_hook` leaves no value for the `&_store` argument; dropping
|
||||||
/// reaches disk, loads records without it, and saves only `WRITER_EDIT` —
|
/// or moving `_store` before `on_after_restore` makes the borrow-check invalid.
|
||||||
/// the first assertion fails deterministically.
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
||||||
use crate::managed_agents::scope::{
|
use crate::managed_agents::scope::{
|
||||||
@@ -168,41 +154,25 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
|||||||
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
||||||
let stopped = vec![entry1];
|
let stopped = vec![entry1];
|
||||||
|
|
||||||
// records_loaded: hook → writer (compensation holds store lock; writer may proceed)
|
// records_loaded_tx: compensation → writer (both locks are held; writer may proceed)
|
||||||
// writer_at_store_lock: writer → hook (writer is about to call managed_agents_store_lock.lock())
|
// writer_committed_tx: writer → main (writer's save is complete)
|
||||||
//
|
|
||||||
// The writer sends writer_at_store_lock immediately BEFORE the lock call, so the
|
|
||||||
// hook's recv() completing is a happens-before guarantee that the writer's next
|
|
||||||
// instruction is managed_agents_store_lock.lock() — which blocks because
|
|
||||||
// compensation holds it.
|
|
||||||
let (records_loaded_tx, records_loaded_rx) = std::sync::mpsc::channel::<()>();
|
let (records_loaded_tx, records_loaded_rx) = std::sync::mpsc::channel::<()>();
|
||||||
let (writer_at_store_lock_tx, writer_at_store_lock_rx) = std::sync::mpsc::channel::<()>();
|
let (writer_committed_tx, writer_committed_rx) = std::sync::mpsc::channel::<()>();
|
||||||
|
|
||||||
// Spawn the writer thread BEFORE acquiring the transition guard.
|
// Spawn the writer BEFORE acquiring the transition guard.
|
||||||
// The writer directly acquires managed_agents_store_lock after signalling —
|
// Flow: wait for records_loaded → acquire managed_agents_store_lock (blocks until
|
||||||
// no intermediate lock — so its blocking is specifically on the store lock
|
// on_after_restore releases it) → load → add WRITER_EDIT → save → send writer_committed.
|
||||||
// that compensate_drain_with_hook holds. This makes the test fail if the
|
|
||||||
// adapter drops the store guard before on_after_restore saves.
|
|
||||||
let tmp_wr = tmp_path.clone();
|
let tmp_wr = tmp_path.clone();
|
||||||
let app_handle_wr = app_handle.clone();
|
let app_handle_wr = app_handle.clone();
|
||||||
let wr_thread = thread::spawn(move || {
|
let wr_thread = thread::spawn(move || {
|
||||||
// Wait until on_records_loaded signals that compensation holds the store lock.
|
// Wait for on_records_loaded to signal that compensation holds both locks.
|
||||||
records_loaded_rx.recv().unwrap();
|
records_loaded_rx.recv().unwrap();
|
||||||
|
|
||||||
// Signal the hook that we are about to call managed_agents_store_lock.lock().
|
// One complete production-shaped transaction: acquire the store lock (blocks
|
||||||
// The hook's recv() will happen-before our lock call, so when on_after_restore
|
// until on_after_restore saves COMP_SENTINEL and the adapter releases the lock),
|
||||||
// mutates COMP_SENTINEL and saves, we are guaranteed to be BLOCKED on the
|
// load records, add WRITER_EDIT, save.
|
||||||
// store lock — not merely about to call it.
|
|
||||||
writer_at_store_lock_tx.send(()).unwrap();
|
|
||||||
|
|
||||||
// Block on the store lock. Compensation still holds it; we wait here until
|
|
||||||
// on_after_restore saves COMP_SENTINEL and the lock is released.
|
|
||||||
let writer_state = app_handle_wr.state::<crate::app_state::AppState>();
|
let writer_state = app_handle_wr.state::<crate::app_state::AppState>();
|
||||||
let _store = writer_state.managed_agents_store_lock.lock().unwrap();
|
let _store = writer_state.managed_agents_store_lock.lock().unwrap();
|
||||||
|
|
||||||
// Store lock acquired: on_after_restore has finished and released the lock.
|
|
||||||
// Load records — COMP_SENTINEL must be present because on_after_restore
|
|
||||||
// saved it before releasing the lock.
|
|
||||||
let mut records =
|
let mut records =
|
||||||
crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default();
|
crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default();
|
||||||
for r in &mut records {
|
for r in &mut records {
|
||||||
@@ -210,45 +180,29 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
|||||||
.insert("WRITER_EDIT".to_string(), "yes".to_string());
|
.insert("WRITER_EDIT".to_string(), "yes".to_string());
|
||||||
}
|
}
|
||||||
crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).unwrap();
|
crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).unwrap();
|
||||||
|
|
||||||
|
// Signal after the save succeeds, while this guarded transaction is complete.
|
||||||
|
writer_committed_tx.send(()).unwrap();
|
||||||
});
|
});
|
||||||
|
|
||||||
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
|
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
|
||||||
|
|
||||||
// on_records_loaded fires after load, while BOTH locks are held:
|
|
||||||
// (a) signal the writer thread that compensation holds the store lock;
|
|
||||||
// (b) wait for the writer to confirm it is at the store-lock boundary —
|
|
||||||
// the writer sends this immediately BEFORE managed_agents_store_lock.lock(),
|
|
||||||
// so by the time recv() returns, the writer's next instruction is that lock
|
|
||||||
// call (which blocks because compensation holds it).
|
|
||||||
//
|
|
||||||
// on_after_restore fires AFTER compensate_drain_for returns, still under BOTH locks:
|
|
||||||
// mutates COMP_SENTINEL in-memory and saves to disk (assert/unwrap).
|
|
||||||
// The store lock is still held — the writer remains blocked until this save returns
|
|
||||||
// and the lock is released.
|
|
||||||
//
|
|
||||||
// What breaks it: if the store guard is dropped before on_after_restore, the writer
|
|
||||||
// acquires managed_agents_store_lock before COMP_SENTINEL reaches disk — the writer's
|
|
||||||
// final save omits COMP_SENTINEL, and the first assertion below fails.
|
|
||||||
let comp_result = compensate_drain_with_hook(
|
let comp_result = compensate_drain_with_hook(
|
||||||
&app_handle,
|
&app_handle,
|
||||||
&stopped,
|
&stopped,
|
||||||
&scope,
|
&scope,
|
||||||
rt_guard,
|
rt_guard,
|
||||||
// on_records_loaded: synchronise with the writer — both locks are held here.
|
// on_records_loaded: release the writer so it can queue on the store lock.
|
||||||
|
// Both locks are held here; the writer will block until on_after_restore
|
||||||
|
// releases the store guard.
|
||||||
|_records| {
|
|_records| {
|
||||||
// (a) Tell the writer that records are loaded; compensation holds both locks.
|
|
||||||
records_loaded_tx.send(()).unwrap();
|
records_loaded_tx.send(()).unwrap();
|
||||||
|
|
||||||
// (b) Wait for the writer to reach the store-lock boundary.
|
|
||||||
// After recv() completes, the writer's next instruction is
|
|
||||||
// managed_agents_store_lock.lock() — which blocks because we hold it.
|
|
||||||
writer_at_store_lock_rx.recv().unwrap();
|
|
||||||
},
|
},
|
||||||
// on_after_restore: save the COMP_SENTINEL mutation while BOTH locks are still held.
|
// on_after_restore: borrows the actual store guard — a compile error if
|
||||||
// The writer is blocked on managed_agents_store_lock. This save reaching disk before
|
// that guard is dropped or removed before this call. Add COMP_SENTINEL
|
||||||
// the lock is released is the load-bearing property: shortening the guard before this
|
// and save while the guard is live; the writer is still blocked.
|
||||||
// callback lets the writer interleave and load records without COMP_SENTINEL.
|
|records, _store_guard| {
|
||||||
|records| {
|
let _ = _store_guard; // borrow is live; guard may not have been dropped
|
||||||
for r in records.iter_mut() {
|
for r in records.iter_mut() {
|
||||||
r.env_vars
|
r.env_vars
|
||||||
.insert("COMP_SENTINEL".to_string(), "yes".to_string());
|
.insert("COMP_SENTINEL".to_string(), "yes".to_string());
|
||||||
@@ -258,6 +212,8 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
|||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
|
||||||
|
// Block until the writer's transaction is complete, then join.
|
||||||
|
writer_committed_rx.recv().unwrap();
|
||||||
wr_thread.join().expect("writer thread panicked");
|
wr_thread.join().expect("writer thread panicked");
|
||||||
|
|
||||||
// Kill the seeded long-lived child.
|
// Kill the seeded long-lived child.
|
||||||
@@ -276,13 +232,8 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
|||||||
);
|
);
|
||||||
|
|
||||||
// ── Final disk state: BOTH effects must be present ───────────────────────
|
// ── Final disk state: BOTH effects must be present ───────────────────────
|
||||||
// COMP_SENTINEL: mutated and saved by on_after_restore while the store lock
|
// COMP_SENTINEL: saved by on_after_restore while the store guard was live.
|
||||||
// was still held — before the writer could acquire and load.
|
// WRITER_EDIT: saved by the writer after the adapter released the lock.
|
||||||
// WRITER_EDIT: written by the writer after the adapter released the lock.
|
|
||||||
//
|
|
||||||
// If the store guard is dropped before on_after_restore, the writer acquires
|
|
||||||
// managed_agents_store_lock before COMP_SENTINEL reaches disk — the writer's
|
|
||||||
// save omits COMP_SENTINEL, and the first assertion below fails.
|
|
||||||
let final_records =
|
let final_records =
|
||||||
crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default();
|
crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default();
|
||||||
let final_rec = final_records
|
let final_rec = final_records
|
||||||
@@ -293,9 +244,7 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
|||||||
assert_eq!(
|
assert_eq!(
|
||||||
final_rec.env_vars.get("COMP_SENTINEL").map(String::as_str),
|
final_rec.env_vars.get("COMP_SENTINEL").map(String::as_str),
|
||||||
Some("yes"),
|
Some("yes"),
|
||||||
"COMP_SENTINEL must reach disk via on_after_restore while the store lock is held; \
|
"COMP_SENTINEL must reach disk via on_after_restore while the store guard is live"
|
||||||
fails if the store guard is dropped before on_after_restore, allowing \
|
|
||||||
the writer to load and save records without COMP_SENTINEL"
|
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
final_rec.env_vars.get("WRITER_EDIT").map(String::as_str),
|
final_rec.env_vars.get("WRITER_EDIT").map(String::as_str),
|
||||||
@@ -305,63 +254,28 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Production start-path contender is blocked on `managed_agent_runtime_transition`
|
/// Production start-path contender is blocked on `managed_agent_runtime_transition`
|
||||||
/// while `compensate_drain_with_hook` holds it, and the `on_transition_acquired` hook
|
/// while `compensate_drain_with_hook` holds it; the `on_transition_acquired` seam
|
||||||
/// in `start_pair_lazy_for_with_hook` fires only AFTER the transition lock is released.
|
/// fires after the guard is acquired, proving start cannot cross the
|
||||||
|
/// transition-acquired seam while compensation owns the mutex.
|
||||||
///
|
///
|
||||||
/// `start_pair_lazy_for_with_hook` is the production-callable seam: production
|
/// Invariant under test: `start_pair_for_with_hook` owns `managed_agent_runtime_transition`
|
||||||
/// `start_managed_agent_runtime_pair_lazy` → `start_pair_lazy_for` → `start_pair_lazy_for_with_hook`
|
/// at the `on_transition_acquired` boundary. The callback receives a borrow of the actual
|
||||||
/// → `start_pair_for_with_hook`. Removing `managed_agent_runtime_transition` from the
|
/// `MutexGuard<'_, ()>`. Dropping that guard before the callback or removing the lock
|
||||||
/// production path also removes it from this test seam, making the test a faithful proxy.
|
/// call from the seam is a **compile error**.
|
||||||
///
|
///
|
||||||
/// Invariant: if `managed_agent_runtime_transition` is removed from the start seam,
|
/// Proof by mutex exclusion: compensation holds `managed_agent_runtime_transition` for the
|
||||||
/// the contender fires `on_transition_acquired` immediately after `on_before_transition`
|
/// duration of `compensate_drain_with_hook`. The contender's `on_transition_acquired` borrows
|
||||||
/// (with zero lock contention) — before compensation has a chance to complete. The
|
/// that same mutex's guard — it can only execute after compensation releases the mutex.
|
||||||
/// `on_records_loaded` hook detects this via `try_recv()` on the channel that
|
/// No timing assumption is required; the scheduler cannot place the callback inside the
|
||||||
/// `on_transition_acquired` sends to, making the test fail deterministically:
|
/// compensation window because mutex exclusion prevents it.
|
||||||
///
|
///
|
||||||
/// - With the lock PRESENT: contender blocks at `managed_agent_runtime_transition.lock()`
|
/// `start_pair_lazy_for_with_hook` is the production-callable seam. Removing
|
||||||
/// immediately after `on_before_transition`. `on_transition_acquired` cannot fire while
|
/// `managed_agent_runtime_transition` from this path also removes it from production.
|
||||||
/// compensation holds the lock. `try_recv()` inside `on_records_loaded` returns `Err(Empty)`.
|
|
||||||
/// After compensation releases the lock, contender fires `on_transition_acquired` and the
|
|
||||||
/// `recv_timeout()` below succeeds.
|
|
||||||
///
|
|
||||||
/// - With the lock REMOVED: contender fires `on_before_transition` (sends
|
|
||||||
/// `contender_at_boundary`), then immediately fires `on_transition_acquired` (sends
|
|
||||||
/// to `contender_transition_acquired_rx`). This fires in nanoseconds. By the time
|
|
||||||
/// compensation reaches `on_records_loaded` (after file I/O for load_managed_agents),
|
|
||||||
/// the channel already contains the message. `try_recv()` succeeds → assertion fails.
|
|
||||||
///
|
|
||||||
/// Flow:
|
|
||||||
/// 1. Seed the store; acquire `managed_agent_runtime_transition` in this thread.
|
|
||||||
/// 2. Spawn the contender. It calls `start_pair_lazy_for_with_hook`:
|
|
||||||
/// - `on_before_transition` fires (BEFORE the lock call) and sends
|
|
||||||
/// `contender_at_boundary` — at this point the contender is inside the seam,
|
|
||||||
/// about to block on `managed_agent_runtime_transition`.
|
|
||||||
/// 3. Wait for `contender_at_boundary` — the contender is now blocked at the
|
|
||||||
/// transition lock because we hold it.
|
|
||||||
/// 4. Call `compensate_drain_with_hook`. The `on_records_loaded` hook:
|
|
||||||
/// (a) asserts `transition_hook_fired` is false (atomic bool sanity check);
|
|
||||||
/// (b) asserts `transition_hook_check_rx.try_recv()` is `Err(Empty)` — the
|
|
||||||
/// deterministic proof that `on_transition_acquired` has not fired yet.
|
|
||||||
/// 5. `compensate_drain_with_hook` finishes; releases the transition guard.
|
|
||||||
/// 6. The contender acquires the guard; `on_transition_acquired` fires and
|
|
||||||
/// sends `contender_transition_acquired`.
|
|
||||||
/// 7. Test asserts that `contender_transition_acquired` fires within a timeout.
|
|
||||||
///
|
|
||||||
/// What breaks it: if `managed_agent_runtime_transition` is removed from
|
|
||||||
/// `start_pair_for_with_hook`, `on_transition_acquired` fires before step 4's assertion
|
|
||||||
/// (the contender fires it in nanoseconds; compensation needs milliseconds for
|
|
||||||
/// file I/O before reaching `on_records_loaded`), making both the `try_recv()` and the
|
|
||||||
/// atomic check fail.
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_compensate_drain_concurrent_start_is_blocked() {
|
fn test_compensate_drain_concurrent_start_is_blocked() {
|
||||||
use crate::managed_agents::scope::{
|
use crate::managed_agents::scope::{
|
||||||
current_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK,
|
current_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK,
|
||||||
};
|
};
|
||||||
use std::sync::{
|
|
||||||
atomic::{AtomicBool, Ordering},
|
|
||||||
Arc,
|
|
||||||
};
|
|
||||||
use std::thread;
|
use std::thread;
|
||||||
use tauri::Manager;
|
use tauri::Manager;
|
||||||
|
|
||||||
@@ -454,120 +368,61 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
|
|||||||
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
||||||
let stopped = vec![entry1];
|
let stopped = vec![entry1];
|
||||||
|
|
||||||
// (A) contender_at_boundary: contender → test (inside seam, about to block on transition lock)
|
// contender_at_boundary_tx: contender → test (inside seam, about to block on transition lock)
|
||||||
// (B) contender_transition_acquired: contender → test (start seam acquired transition guard)
|
// transition_acquired_tx: contender → test (contender acquired transition guard)
|
||||||
// Two receivers for the same semantic signal:
|
|
||||||
// - transition_hook_check_rx: moved into on_records_loaded for try_recv() check
|
|
||||||
// - contender_transition_acquired_rx: used after compensation for recv_timeout() check
|
|
||||||
let (contender_at_boundary_tx, contender_at_boundary_rx) = std::sync::mpsc::channel::<()>();
|
let (contender_at_boundary_tx, contender_at_boundary_rx) = std::sync::mpsc::channel::<()>();
|
||||||
let (transition_hook_check_tx, transition_hook_check_rx) = std::sync::mpsc::channel::<()>();
|
let (transition_acquired_tx, transition_acquired_rx) = std::sync::mpsc::channel::<()>();
|
||||||
let (contender_transition_acquired_tx, contender_transition_acquired_rx) =
|
|
||||||
std::sync::mpsc::channel::<()>();
|
|
||||||
|
|
||||||
// Shared flag: set to true when `on_transition_acquired` fires inside start_pair_for_with_hook.
|
// Acquire the transition guard FIRST so the contender blocks on it.
|
||||||
// Belt-and-suspenders companion to the try_recv() check in on_records_loaded.
|
|
||||||
let transition_hook_fired = Arc::new(AtomicBool::new(false));
|
|
||||||
let transition_hook_fired2 = transition_hook_fired.clone();
|
|
||||||
|
|
||||||
// Acquire the transition guard FIRST so the contender will block on it.
|
|
||||||
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
|
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
|
||||||
|
|
||||||
let app_contender = app_handle.clone();
|
let app_contender = app_handle.clone();
|
||||||
let pubkey_contender = pubkey1.clone();
|
let pubkey_contender = pubkey1.clone();
|
||||||
let contender = thread::spawn(move || {
|
let contender = thread::spawn(move || {
|
||||||
// `start_pair_lazy_for_with_hook` is the production-called seam.
|
// `start_pair_lazy_for_with_hook` is the production-called seam.
|
||||||
// on_before_transition fires just BEFORE managed_agent_runtime_transition.lock().
|
// on_before_transition fires just before the lock call — signals "entered the seam".
|
||||||
// At that point we are inside the seam, about to block on the transition guard
|
// on_transition_acquired fires after the lock is acquired, receiving a borrow of the
|
||||||
// (which the test holds). Signaling from here proves the contender is inside
|
// actual transition guard — sends one positive receipt to the test.
|
||||||
// the seam and will block on the next line — not merely about to call the function.
|
|
||||||
let _ = start_pair_lazy_for_with_hook(
|
let _ = start_pair_lazy_for_with_hook(
|
||||||
pubkey_contender,
|
pubkey_contender,
|
||||||
"wss://relay.example".to_string(),
|
"wss://relay.example".to_string(),
|
||||||
app_contender,
|
app_contender,
|
||||||
// on_before_transition: fires inside the seam, just before the lock call.
|
// on_before_transition: contender is inside the seam and about to block.
|
||||||
// The contender is now committed to acquiring managed_agent_runtime_transition
|
|
||||||
// and will block there immediately after this hook returns.
|
|
||||||
move || {
|
move || {
|
||||||
contender_at_boundary_tx.send(()).unwrap();
|
contender_at_boundary_tx.send(()).unwrap();
|
||||||
// Next line in the seam: managed_agent_runtime_transition.lock() — blocks.
|
|
||||||
},
|
},
|
||||||
// on_transition_acquired: fires AFTER the transition guard is acquired,
|
// on_transition_acquired: borrows the actual transition guard at the seam boundary.
|
||||||
// BEFORE the store lock. Signals the test that start has passed the
|
// Removing or dropping the guard before this call is a compile error.
|
||||||
// transition-lock boundary. Sends to two channels: one for the
|
// Sends one positive receipt — proves the guard was acquired and live.
|
||||||
// try_recv() check inside on_records_loaded, one for the final
|
move |_transition_guard| {
|
||||||
// recv_timeout() check after compensation returns.
|
let _ = _transition_guard; // borrow is live; guard may not have been dropped
|
||||||
move || {
|
transition_acquired_tx.send(()).unwrap();
|
||||||
transition_hook_fired2.store(true, Ordering::SeqCst);
|
|
||||||
transition_hook_check_tx.send(()).unwrap();
|
|
||||||
contender_transition_acquired_tx.send(()).unwrap();
|
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
// Wait for the contender to be inside the seam and about to block on the
|
// Wait for the contender to be inside the seam and blocked on the transition lock.
|
||||||
// transition lock. After this signal the contender is on the next line:
|
|
||||||
// managed_agent_runtime_transition.lock() — which blocks because we hold it.
|
|
||||||
contender_at_boundary_rx.recv().unwrap();
|
contender_at_boundary_rx.recv().unwrap();
|
||||||
|
|
||||||
// on_records_loaded hook: fires while BOTH locks are held.
|
// Run compensation while holding the transition guard. The contender cannot acquire it;
|
||||||
//
|
// its on_transition_acquired callback cannot execute until compensation returns.
|
||||||
// Two complementary checks that on_transition_acquired has NOT fired:
|
|
||||||
// (1) Atomic bool: fast sanity check.
|
|
||||||
// (2) try_recv(): channel-based proof — on_transition_acquired sends to the
|
|
||||||
// channel; if the channel is empty here, the hook has not fired.
|
|
||||||
//
|
|
||||||
// Why (2) is deterministic when the lock is removed: on_before_transition fires,
|
|
||||||
// on_transition_acquired fires immediately (zero lock contention), and the send
|
|
||||||
// completes in nanoseconds. By the time compensation reaches on_records_loaded
|
|
||||||
// (after acquiring the store lock and loading records from disk — milliseconds
|
|
||||||
// of I/O), the channel is guaranteed to contain the message. try_recv() finds it
|
|
||||||
// and the assertion fails, catching the missing lock.
|
|
||||||
//
|
|
||||||
// When the lock IS present: on_before_transition fires, then the contender
|
|
||||||
// blocks at managed_agent_runtime_transition.lock(). on_transition_acquired
|
|
||||||
// cannot fire until after compensation releases the guard. try_recv() correctly
|
|
||||||
// returns Err(Empty).
|
|
||||||
let _comp_result = compensate_drain_with_hook(
|
let _comp_result = compensate_drain_with_hook(
|
||||||
&app_handle,
|
&app_handle,
|
||||||
&stopped,
|
&stopped,
|
||||||
&scope,
|
&scope,
|
||||||
rt_guard,
|
rt_guard,
|
||||||
// on_records_loaded: fires while BOTH locks are held. Checks that
|
|
||||||
// on_transition_acquired has NOT fired — the contender is still blocked
|
|
||||||
// at managed_agent_runtime_transition.lock() because we hold it.
|
|
||||||
|_records| {
|
|
||||||
// Check 1: atomic bool.
|
|
||||||
assert!(
|
|
||||||
!transition_hook_fired.load(Ordering::SeqCst),
|
|
||||||
"start seam must not acquire the transition guard while compensation holds it; \
|
|
||||||
fails if managed_agent_runtime_transition is removed from start_pair_for_with_hook"
|
|
||||||
);
|
|
||||||
// Check 2: channel try_recv — deterministic proof via file-I/O time differential.
|
|
||||||
// transition_hook_check_rx is a dedicated receiver that on_transition_acquired
|
|
||||||
// sends to; this closure moves it so the outer recv_timeout() uses the
|
|
||||||
// separate contender_transition_acquired_rx.
|
|
||||||
assert!(
|
|
||||||
transition_hook_check_rx.try_recv().is_err(),
|
|
||||||
"on_transition_acquired must not fire while compensation holds the transition \
|
|
||||||
guard; fails if managed_agent_runtime_transition is removed from \
|
|
||||||
start_pair_for_with_hook \
|
|
||||||
(the hook fires in nanoseconds; this check runs after ms of file I/O)"
|
|
||||||
);
|
|
||||||
},
|
|
||||||
// on_after_restore: no-op for the contender test — this test's load-bearing
|
|
||||||
// property is the transition guard, not the store-lock lifetime.
|
|
||||||
|_records| {},
|
|_records| {},
|
||||||
|
|_records, _store_guard| {},
|
||||||
);
|
);
|
||||||
|
|
||||||
// After compensation returns, the contender can acquire the transition guard.
|
// After compensation returns, the contender can acquire the transition guard.
|
||||||
// `on_transition_acquired` fires and sends the signal.
|
// Wait for the positive receipt from on_transition_acquired.
|
||||||
contender_transition_acquired_rx
|
// Mutex exclusion guarantees this callback did not execute during compensation —
|
||||||
.recv_timeout(std::time::Duration::from_secs(5))
|
// no timing assumption is required.
|
||||||
.expect(
|
transition_acquired_rx.recv().expect(
|
||||||
"contender's on_transition_acquired must fire after compensate_drain_with_hook \
|
"contender's on_transition_acquired must fire after compensate_drain_with_hook \
|
||||||
releases the transition guard; fails if start_pair_for_with_hook does not acquire \
|
releases the transition guard",
|
||||||
managed_agent_runtime_transition before calling on_transition_acquired",
|
);
|
||||||
);
|
|
||||||
|
|
||||||
contender.join().expect("contender thread panicked");
|
contender.join().expect("contender thread panicked");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ pub(crate) fn start_pair_lazy_for_with_hook<R: tauri::Runtime>(
|
|||||||
relay_url: String,
|
relay_url: String,
|
||||||
app: tauri::AppHandle<R>,
|
app: tauri::AppHandle<R>,
|
||||||
on_before_transition: impl FnOnce(),
|
on_before_transition: impl FnOnce(),
|
||||||
on_transition_acquired: impl FnOnce(),
|
on_transition_acquired: impl FnOnce(&std::sync::MutexGuard<'_, ()>),
|
||||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||||
start_pair_for_with_hook(
|
start_pair_for_with_hook(
|
||||||
pubkey,
|
pubkey,
|
||||||
@@ -33,14 +33,15 @@ pub(crate) fn start_pair_lazy_for_with_hook<R: tauri::Runtime>(
|
|||||||
/// Generic start-pair seam with injectable hooks.
|
/// Generic start-pair seam with injectable hooks.
|
||||||
///
|
///
|
||||||
/// - `on_before_transition`: fires BEFORE `managed_agent_runtime_transition` is
|
/// - `on_before_transition`: fires BEFORE `managed_agent_runtime_transition` is
|
||||||
/// locked. Tests signal "at the lock boundary" from here — the contender is
|
/// locked. Tests signal "entered the seam" from here.
|
||||||
/// committed to acquiring `managed_agent_runtime_transition` immediately after.
|
|
||||||
///
|
///
|
||||||
/// - `on_transition_acquired`: fires AFTER `managed_agent_runtime_transition` is
|
/// - `on_transition_acquired`: fires AFTER `managed_agent_runtime_transition` is
|
||||||
/// acquired but BEFORE `managed_agents_store_lock` is attempted. Signals that
|
/// acquired (borrowing the actual guard) but BEFORE `managed_agents_store_lock`
|
||||||
/// start has passed the transition-lock boundary.
|
/// is attempted. The argument is a borrow of the actual transition guard; a
|
||||||
|
/// drop or removal of that guard before this call is a compile error.
|
||||||
///
|
///
|
||||||
/// Production delegates with `|| {}` for both hooks — release-invisible no-ops.
|
/// Production delegates with `|| {}` and `|_| {}` for the two hooks —
|
||||||
|
/// release-invisible no-ops.
|
||||||
pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
|
pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
|
||||||
pubkey: String,
|
pubkey: String,
|
||||||
relay_url: String,
|
relay_url: String,
|
||||||
@@ -48,11 +49,11 @@ pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
|
|||||||
expected_updated_at: Option<&str>,
|
expected_updated_at: Option<&str>,
|
||||||
app: tauri::AppHandle<R>,
|
app: tauri::AppHandle<R>,
|
||||||
on_before_transition: impl FnOnce(),
|
on_before_transition: impl FnOnce(),
|
||||||
on_transition_acquired: impl FnOnce(),
|
on_transition_acquired: impl FnOnce(&std::sync::MutexGuard<'_, ()>),
|
||||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||||
let state = app.state::<AppState>();
|
let state = app.state::<AppState>();
|
||||||
on_before_transition();
|
on_before_transition();
|
||||||
let _transition = state
|
let transition_guard = state
|
||||||
.managed_agent_runtime_transition
|
.managed_agent_runtime_transition
|
||||||
.lock()
|
.lock()
|
||||||
.map_err(|e| e.to_string())?;
|
.map_err(|e| e.to_string())?;
|
||||||
@@ -62,7 +63,7 @@ pub(crate) fn start_pair_for_with_hook<R: tauri::Runtime>(
|
|||||||
{
|
{
|
||||||
return Err("desktop shutdown has started".into());
|
return Err("desktop shutdown has started".into());
|
||||||
}
|
}
|
||||||
on_transition_acquired();
|
on_transition_acquired(&transition_guard);
|
||||||
let _store = state
|
let _store = state
|
||||||
.managed_agents_store_lock
|
.managed_agents_store_lock
|
||||||
.lock()
|
.lock()
|
||||||
|
|||||||
Reference in New Issue
Block a user