|
|
@@ -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");
|
|
|
|
}
|
|
|
|
}
|
|
|
|