mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): resolve P4 pass-2 test-vs-claim defects in workspace-scoped agent store
Four IMPORTANT findings: each assertion now fails if the production guard or
save it claims to prove is removed.
Fix 1 — compensate_drain writer test (genuine store-lock contention):
Add compensate_drain_with_hook seam; on_store_acquired hook fires while BOTH
locks are held. Test writes COMP_SENTINEL under the lock, signals writer,
waits for writer to queue on the store lock, then releases. Final disk state
must contain BOTH COMP_SENTINEL and WRITER_EDIT — fails if the guard is
dropped before compensate_drain_for's save.
Fix 1b — start-contender test (production start-lock seam):
Add start_pair_lazy_for<R: tauri::Runtime> generic seam; production
start_managed_agent_runtime_pair_lazy and start_pair delegate to it.
Contender now calls start_pair_lazy_for (the seam the production adapter
calls) instead of a raw mutex lock — proves compensation blocks production
starts, not just an arbitrary lock acquisition.
Fix 2 — Record-Mesh TOCTOU (async production driver):
Drive restart_local_agent_on_config_change_for (full async driver). Injected
mesh_fn rewrites provider to relay-mesh after the driver has resolved
mesh_model_id=None; epoch re-resolves to Some("auto") != None → TOCTOU guard
fires → Skipped before stop. Removing the in-epoch re-resolve guard would
allow the epoch to proceed, calling stop_fn and failing the assertion.
Fix 3 — Helper-vs-helper serialization (pre-acquisition hook):
Add with_workspace_transition_preflight_with_hook and
install_client_under_workspace_transition_with_hook to mesh_llm_scope.rs.
Both serialization tests use a proper 3-step handshake: holder signals
lock_held BEFORE publishing state; contender fires pre-acquisition hook
signaling at_lock_boundary; only then does the holder publish and release.
Removing the shared lock would let the contender bypass the boundary,
inverting each test's expected outcome.
Fix 4 — Full-tail receipt and disk assertions:
write_receipt_fn captures all receipt fields: key.pubkey, key.relay_url,
pid, desktop_instance_id (asserted equal to app.config().identifier),
started_at. Final disk assertions verify last_started_at.is_some(),
last_stopped_at.is_none(), last_error.is_none() on the matching record.
Removing the post-spawn record save drops last_started_at and fails the
disk assertions.
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
91aea2db96
commit
45f199ba5d
@@ -222,6 +222,12 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
|
||||
let spawn_got_nonempty_teams = Arc::new(Mutex::new(false));
|
||||
let spawn_got_nonempty_global = Arc::new(Mutex::new(false));
|
||||
let receipt_called = Arc::new(Mutex::new(false));
|
||||
// Capture all receipt fields for post-completion assertions.
|
||||
let receipt_pubkey = Arc::new(Mutex::new(None::<String>));
|
||||
let receipt_relay = Arc::new(Mutex::new(None::<String>));
|
||||
let receipt_pid = Arc::new(Mutex::new(0u32));
|
||||
let receipt_instance_id = Arc::new(Mutex::new(String::new()));
|
||||
let receipt_started_at = Arc::new(Mutex::new(String::new()));
|
||||
|
||||
let stop_called2 = stop_called.clone();
|
||||
let spawn_relay2 = spawn_relay.clone();
|
||||
@@ -230,9 +236,15 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
|
||||
let spawn_teams2 = spawn_got_nonempty_teams.clone();
|
||||
let spawn_global2 = spawn_got_nonempty_global.clone();
|
||||
let receipt_called2 = receipt_called.clone();
|
||||
let receipt_pubkey2 = receipt_pubkey.clone();
|
||||
let receipt_relay2 = receipt_relay.clone();
|
||||
let receipt_pid2 = receipt_pid.clone();
|
||||
let receipt_iid2 = receipt_instance_id.clone();
|
||||
let receipt_sat2 = receipt_started_at.clone();
|
||||
let pubkey2 = pubkey.clone();
|
||||
|
||||
let app_handle_for_assert = app_handle.clone();
|
||||
let pubkey_for_assert = pubkey.clone();
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
restart_under_captured_epoch_for(
|
||||
&app_handle,
|
||||
@@ -273,9 +285,15 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
|
||||
job: None,
|
||||
})
|
||||
},
|
||||
// write_receipt_fn: record the call, verify pubkey, succeed.
|
||||
// write_receipt_fn: record all receipt fields for post-completion assertions.
|
||||
// Full key (pubkey + relay_url), pid, desktop_instance_id, started_at.
|
||||
move |_app, receipt| {
|
||||
*receipt_called2.lock().unwrap() = true;
|
||||
*receipt_pubkey2.lock().unwrap() = Some(receipt.key.pubkey.clone());
|
||||
*receipt_relay2.lock().unwrap() = Some(receipt.key.relay_url.clone());
|
||||
*receipt_pid2.lock().unwrap() = receipt.pid;
|
||||
*receipt_iid2.lock().unwrap() = receipt.desktop_instance_id.clone();
|
||||
*receipt_sat2.lock().unwrap() = receipt.started_at.clone();
|
||||
assert_eq!(
|
||||
receipt.key.pubkey, pubkey2,
|
||||
"receipt must carry the correct pubkey"
|
||||
@@ -321,6 +339,36 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
|
||||
*receipt_called.lock().unwrap(),
|
||||
"write_receipt_fn must be called"
|
||||
);
|
||||
// ── Receipt field assertions ──────────────────────────────────────────────
|
||||
// A full relay-bearing key: pubkey + relay_url.
|
||||
assert_eq!(
|
||||
receipt_pubkey.lock().unwrap().as_deref(),
|
||||
Some(pubkey_for_assert.as_str()),
|
||||
"receipt.key.pubkey must match the restarted agent"
|
||||
);
|
||||
assert_eq!(
|
||||
receipt_relay.lock().unwrap().as_deref(),
|
||||
Some(relay_url),
|
||||
"receipt.key.relay_url must match the captured relay; fails if the spawn path \
|
||||
builds the receipt with an empty or wrong relay"
|
||||
);
|
||||
assert!(
|
||||
*receipt_pid.lock().unwrap() > 0,
|
||||
"receipt.pid must be non-zero (a real spawned child PID)"
|
||||
);
|
||||
// Assert receipt.desktop_instance_id equals the actual value produced by
|
||||
// current_instance_id(app) — proves the field is sourced from the app
|
||||
// handle, not left empty or set to a hardcoded value.
|
||||
let expected_instance_id = app_handle_for_assert.config().identifier.clone();
|
||||
assert_eq!(
|
||||
receipt_instance_id.lock().unwrap().as_str(),
|
||||
expected_instance_id.as_str(),
|
||||
"receipt.desktop_instance_id must equal current_instance_id(app)"
|
||||
);
|
||||
assert!(
|
||||
!receipt_started_at.lock().unwrap().is_empty(),
|
||||
"receipt.started_at must be a non-empty ISO timestamp"
|
||||
);
|
||||
// Verify the runtime is registered with the captured scope_id.
|
||||
{
|
||||
let state = app_handle_for_assert.state::<crate::app_state::AppState>();
|
||||
@@ -333,12 +381,27 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
|
||||
"runtime must be registered with the captured scope_id"
|
||||
);
|
||||
}
|
||||
// Verify the final captured disk record: the agent is present in the store.
|
||||
// ── Final disk record assertions ──────────────────────────────────────────
|
||||
// The restart-produced state: last_started_at set, last_stopped_at cleared,
|
||||
// last_error cleared. Fails if the post-spawn save is removed or if the
|
||||
// record-update block is bypassed.
|
||||
let final_records =
|
||||
crate::managed_agents::storage::load_managed_agents_at(tmp.path()).unwrap_or_default();
|
||||
let final_rec = final_records
|
||||
.iter()
|
||||
.find(|r| r.pubkey == pubkey_for_assert)
|
||||
.expect("final disk record must contain the restarted agent");
|
||||
assert!(
|
||||
final_records.iter().any(|r| r.pubkey == "bb".repeat(32)),
|
||||
"final disk record must contain the restarted agent"
|
||||
final_rec.last_started_at.is_some(),
|
||||
"last_started_at must be Some after a successful restart (set by post-spawn record update)"
|
||||
);
|
||||
assert!(
|
||||
final_rec.last_stopped_at.is_none(),
|
||||
"last_stopped_at must be None after restart (cleared by post-spawn record update)"
|
||||
);
|
||||
assert!(
|
||||
final_rec.last_error.is_none(),
|
||||
"last_error must be None after a successful restart (cleared by post-spawn record update)"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -512,32 +512,30 @@ async fn test_workspace_switch_after_preflight_aborts_before_stop() {
|
||||
/// A record-level Mesh model change after preflight (without advancing workspace
|
||||
/// generation) → epoch detects mismatch in re-resolved Mesh model, aborts before stop.
|
||||
///
|
||||
/// Thufir's test 4: "drive the async production driver; seed an eligible runtime
|
||||
/// via the cross-platform child helper; mutate an actually-participating
|
||||
/// record/definition/global input from the injected mesh_fn (which runs in the
|
||||
/// pre-stop phase); assert re-resolution detects the mismatch and stop never
|
||||
/// fires."
|
||||
/// Drives `restart_local_agent_on_config_change_for` — the full async production
|
||||
/// driver. The injected `mesh_fn` mutates the on-disk record's `provider` to
|
||||
/// `"relay-mesh"` AFTER the driver has already resolved `context.mesh_model_id = None`
|
||||
/// (from the original `provider = "anthropic"`). The epoch's in-epoch TOCTOU guard
|
||||
/// then re-resolves the mutated record, gets `Some("auto")` ≠ `None`, and fires
|
||||
/// `Skipped` before stop.
|
||||
///
|
||||
/// The async driver (`restart_local_agent_on_config_change_for`) runs mesh_fn
|
||||
/// in the pre-stop phase, captures the `CapturedRestartContext`, then hands off
|
||||
/// to `restart_under_captured_epoch_for` (the stop→spawn primitive). We inject
|
||||
/// a `mesh_fn` that succeeds but captures `mesh_model_id = Some("model-a")`
|
||||
/// into the context by returning Ok — the real driver then stores whatever
|
||||
/// the record has at the time of pre-stop resolution into `context.mesh_model_id`.
|
||||
/// Flow:
|
||||
/// 1. Seed record: `provider = "anthropic"`, live runtime.
|
||||
/// Pre-stop resolve: `resolve_effective_relay_mesh_model_id` → `None`.
|
||||
/// `context.mesh_model_id = None`.
|
||||
/// 2. `mesh_fn` rewrites `provider = "relay-mesh"` to disk and succeeds.
|
||||
/// 3. Epoch re-resolves from the mutated record:
|
||||
/// `relay_mesh_model_id()` = `Some("auto")` ≠ `None` → TOCTOU guard fires → Skipped.
|
||||
/// 4. Assert `RestartOutcome::Skipped`; stop_fn must NOT be called.
|
||||
///
|
||||
/// To trigger the in-epoch TOCTOU guard:
|
||||
/// - The record on disk has `relay_mesh: None` → driver pre-resolves
|
||||
/// `mesh_model_id = None`.
|
||||
/// - We set `mesh_model_id = Some("model-a")` in an injected post-preflight
|
||||
/// hook by calling `restart_under_captured_epoch_for` directly with a context
|
||||
/// whose `mesh_model_id` disagrees with the on-disk record.
|
||||
///
|
||||
/// Approach: use `restart_under_captured_epoch_for` directly (with the
|
||||
/// cross-platform long-lived child seeded), so the in-epoch TOCTOU guard fires.
|
||||
/// The "async driver drives the TOCTOU guard" variant is covered by the
|
||||
/// generation-advance test (`test_workspace_switch_after_preflight_aborts_before_stop`).
|
||||
/// Invariant: removing the in-epoch re-resolve guard in
|
||||
/// `restart_under_captured_epoch_for` would allow the epoch to proceed with the
|
||||
/// old `context.mesh_model_id`, bypass the mismatch — stop would be called
|
||||
/// and the test would fail.
|
||||
#[tokio::test]
|
||||
async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
use super::super::restart_local_agent_on_config_change_for;
|
||||
use crate::commands::global_agent_config::RestartOutcome;
|
||||
use crate::managed_agents::{
|
||||
storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord,
|
||||
ManagedAgentRuntimeKey,
|
||||
@@ -548,13 +546,9 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
let pubkey = "aa".repeat(32);
|
||||
let relay_url = "wss://relay.example";
|
||||
|
||||
// Build a Ready record: provider+model in structured fields so the agent
|
||||
// passes the eligibility gate (old_ready = true). ANTHROPIC_API_KEY is
|
||||
// required by buzz_agent_requirements when provider=anthropic, so we set
|
||||
// it in env_vars to satisfy the readiness check. Two globals differ by
|
||||
// one env_var so env_changed = true → should_restart_on_config_change
|
||||
// returns true. relay_mesh: None so the in-epoch re-resolve yields None
|
||||
// while context.mesh_model_id = Some("model-a") → Skipped before stop.
|
||||
// Eligible record: provider=anthropic, model+ANTHROPIC_API_KEY → old_ready=true.
|
||||
// relay_mesh=None so the pre-stop resolve yields mesh_model_id=None.
|
||||
// old_global != new_global (env differ) so env_changed=true → eligible.
|
||||
let mut record_env_vars = std::collections::BTreeMap::new();
|
||||
record_env_vars.insert(
|
||||
"ANTHROPIC_API_KEY".to_string(),
|
||||
@@ -581,7 +575,7 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
parallelism: 1,
|
||||
system_prompt: None,
|
||||
model: Some("claude-3-5-sonnet-20241022".to_string()),
|
||||
provider: Some("anthropic".to_string()),
|
||||
provider: Some("anthropic".to_string()), // initial provider — not relay-mesh
|
||||
persona_source_version: None,
|
||||
env_vars: record_env_vars,
|
||||
start_on_app_launch: false,
|
||||
@@ -611,7 +605,7 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
definition_respond_to: None,
|
||||
definition_respond_to_allowlist: Default::default(),
|
||||
definition_parallelism: None,
|
||||
relay_mesh: None, // ← no relay_mesh; re-resolve yields None ≠ Some("model-a")
|
||||
relay_mesh: None, // no relay-mesh marker; pre-stop resolve yields None
|
||||
runtime: None,
|
||||
name_pool: vec![],
|
||||
};
|
||||
@@ -622,6 +616,17 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
let app = make_mock_app();
|
||||
let app_handle = app.handle().clone();
|
||||
|
||||
// Derive the actual owner pubkey from the mock app's signing keys.
|
||||
// restart_local_agent_on_config_change_for verifies hex == scope.owner_pubkey.
|
||||
let actual_owner_hex = {
|
||||
let state = app_handle.state::<crate::app_state::AppState>();
|
||||
state
|
||||
.signing_keys()
|
||||
.expect("mock app must have signing keys")
|
||||
.public_key()
|
||||
.to_hex()
|
||||
};
|
||||
|
||||
// Seed a live runtime with a cross-platform long-lived child (avoids sync eviction).
|
||||
let rt_key = ManagedAgentRuntimeKey::new(&pubkey, relay_url).unwrap();
|
||||
let seeded_pid = {
|
||||
@@ -656,13 +661,13 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
|
||||
scope_id: "test-scope".to_string(),
|
||||
relay_url: relay_url.to_string(),
|
||||
owner_pubkey: "aa".repeat(32),
|
||||
owner_pubkey: actual_owner_hex,
|
||||
definitions_dir: tmp.path().to_path_buf(),
|
||||
generation: gen,
|
||||
};
|
||||
|
||||
// old_global and new_global differ by one env_var so env_changed = true
|
||||
// (eligibility gate passes) while the record's relay_mesh stays None.
|
||||
// (eligibility gate passes) while the record's provider stays anthropic.
|
||||
let mut new_global_env = std::collections::BTreeMap::new();
|
||||
new_global_env.insert("SOME_EXTRA_KEY".to_string(), "v2".to_string());
|
||||
let old_global = crate::managed_agents::GlobalAgentConfig::default();
|
||||
@@ -671,62 +676,61 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
// Build a context that claims `mesh_model_id = Some("model-a")`, but
|
||||
// the actual records on disk have no relay_mesh config, so the epoch's
|
||||
// re-resolve will return None — triggering the TOCTOU mismatch guard.
|
||||
let context_with_mesh = CapturedRestartContext {
|
||||
scope: scope.clone(),
|
||||
personas: vec![],
|
||||
teams: vec![],
|
||||
global: old_global.clone(),
|
||||
owner_hex: "aa".repeat(32),
|
||||
mesh_model_id: Some("model-a".to_string()),
|
||||
};
|
||||
|
||||
let stop_called = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
|
||||
let stop_called2 = stop_called.clone();
|
||||
let tmp_path = tmp.path().to_path_buf();
|
||||
let pubkey_clone = pubkey.clone();
|
||||
|
||||
// Call the epoch core directly — mesh_fn is not involved here since
|
||||
// we're testing the in-epoch TOCTOU check that re-resolves relay_mesh.
|
||||
// This exercises the same code path that the async driver reaches after
|
||||
// mesh_fn completes: stop→spawn epoch with pre-captured context.
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
restart_under_captured_epoch_for(
|
||||
&app_handle,
|
||||
&pubkey,
|
||||
&old_global,
|
||||
&new_global,
|
||||
&[],
|
||||
&context_with_mesh,
|
||||
move |_app, _rec, _runtimes| {
|
||||
stop_called2.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||
Err("stop_fn called unexpectedly".to_string())
|
||||
},
|
||||
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
|
||||
Err("spawn not expected".to_string())
|
||||
},
|
||||
|_app, _receipt| Err("receipt not expected".to_string()),
|
||||
)
|
||||
})
|
||||
.await
|
||||
.expect("spawn_blocking must not panic");
|
||||
// Drive the full async production driver. The mesh_fn mutates the on-disk
|
||||
// record's provider to "relay-mesh" AFTER the driver has resolved
|
||||
// context.mesh_model_id = None (provider was "anthropic" at resolve time).
|
||||
// The epoch then re-resolves the mutated record and gets Some("auto") != None
|
||||
// -> TOCTOU guard fires -> Skipped before stop.
|
||||
let outcome = restart_local_agent_on_config_change_for(
|
||||
&app_handle,
|
||||
&pubkey,
|
||||
&old_global,
|
||||
&new_global,
|
||||
&[],
|
||||
&scope,
|
||||
tmp.path(),
|
||||
// mesh_fn: rewrites provider to "relay-mesh" on disk, then succeeds.
|
||||
// The driver already captured context.mesh_model_id=None from the
|
||||
// original provider="anthropic" record — this mutation happens after.
|
||||
move |_app, _model| {
|
||||
let mut records = crate::managed_agents::storage::load_managed_agents_at(&tmp_path)
|
||||
.unwrap_or_default();
|
||||
for r in &mut records {
|
||||
if r.pubkey == pubkey_clone {
|
||||
r.provider = Some("relay-mesh".to_string());
|
||||
}
|
||||
}
|
||||
let _ = crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &records);
|
||||
Box::pin(async { Ok(()) })
|
||||
},
|
||||
// stop_fn: must NOT be called — TOCTOU guard fires before stop.
|
||||
move |_app, _rec, _runtimes| {
|
||||
stop_called2.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||
Err("stop must not be called when Mesh model changed after preflight".to_string())
|
||||
},
|
||||
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
|
||||
Err("spawn not expected".to_string())
|
||||
},
|
||||
|_app, _receipt| Err("receipt not expected".to_string()),
|
||||
)
|
||||
.await;
|
||||
|
||||
// Kill the seeded process now that the epoch has consumed it.
|
||||
let _ = crate::managed_agents::terminate_process(seeded_pid);
|
||||
|
||||
assert!(
|
||||
matches!(result, Err(EpochError::Skipped(_))),
|
||||
"Mesh model mismatch must produce Skipped before stop: {result:?}"
|
||||
matches!(outcome, RestartOutcome::Skipped),
|
||||
"Mesh model mismatch (via full async driver) must produce Skipped before stop: {outcome:?}"
|
||||
);
|
||||
if let Err(EpochError::Skipped(msg)) = &result {
|
||||
assert!(
|
||||
msg.contains("relay-mesh model changed") || msg.contains("mesh"),
|
||||
"Skipped reason must mention mesh model change: {msg}"
|
||||
);
|
||||
}
|
||||
assert!(
|
||||
!stop_called.load(std::sync::atomic::Ordering::SeqCst),
|
||||
"stop must NOT be called when Mesh model changed after preflight"
|
||||
);
|
||||
// Cleanup: kill the seeded long-lived process so it doesn't leak.
|
||||
let _ = crate::managed_agents::terminate_process(seeded_pid);
|
||||
}
|
||||
|
||||
#[path = "global_agent_config_epoch_tests.rs"]
|
||||
|
||||
@@ -102,15 +102,8 @@ pub(crate) async fn fail_if_client_mesh_active<R: tauri::Runtime>(
|
||||
/// `fail_if_client_mesh_active` preflight, then invoke `transition_body` while
|
||||
/// the guard remains held.
|
||||
///
|
||||
/// Used by `apply_workspace`, live identity import, and tests to prove the
|
||||
/// serialization contract against `install_client_under_workspace_transition`.
|
||||
///
|
||||
/// `transition_body` is an async closure (returns a `BoxFuture`) so that
|
||||
/// production callers that dispatch `tokio::task::spawn_blocking` and await the
|
||||
/// result keep the guard alive across the dispatch. Tests may pass a simple
|
||||
/// sync-compatible closure via `|| Box::pin(async { Ok("done") })`.
|
||||
///
|
||||
/// If the preflight fails, `transition_body` is never called.
|
||||
/// Delegates to [`with_workspace_transition_preflight_with_hook`] with a no-op
|
||||
/// pre-acquisition hook. See that function for full documentation.
|
||||
pub(crate) async fn with_workspace_transition_preflight<R, F, T>(
|
||||
app: &AppHandle<R>,
|
||||
transition_body: F,
|
||||
@@ -118,8 +111,32 @@ pub(crate) async fn with_workspace_transition_preflight<R, F, T>(
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
F: FnOnce() -> BoxFuture<'static, Result<T, String>>,
|
||||
{
|
||||
with_workspace_transition_preflight_with_hook(app, || {}, transition_body).await
|
||||
}
|
||||
|
||||
/// Inner implementation of [`with_workspace_transition_preflight`] with an
|
||||
/// injectable `pre_acquisition_hook`.
|
||||
///
|
||||
/// `pre_acquisition_hook` fires once, synchronously, immediately before
|
||||
/// `workspace_transition.lock().await`. In production this is `|| {}`; tests
|
||||
/// inject a closure that signals "I'm about to acquire" so the test can
|
||||
/// establish a deterministic ordering between holder and contender.
|
||||
///
|
||||
/// This is `pub(crate)` so tests in `mesh_llm_transition_tests` can drive the
|
||||
/// exact lock-acquisition boundary while remaining invisible to external callers.
|
||||
pub(crate) async fn with_workspace_transition_preflight_with_hook<R, H, F, T>(
|
||||
app: &AppHandle<R>,
|
||||
pre_acquisition_hook: H,
|
||||
transition_body: F,
|
||||
) -> Result<T, String>
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
H: FnOnce(),
|
||||
F: FnOnce() -> BoxFuture<'static, Result<T, String>>,
|
||||
{
|
||||
let state = app.state::<AppState>();
|
||||
pre_acquisition_hook();
|
||||
let _transition_guard = state.workspace_transition.lock().await;
|
||||
|
||||
// Fail closed if a client-mode Mesh runtime is active.
|
||||
@@ -152,13 +169,8 @@ where
|
||||
/// identity under the guard — `(scope_id, normalized relay, owner_pubkey,
|
||||
/// generation)` — then call the injected `install` closure.
|
||||
///
|
||||
/// Fails if:
|
||||
/// - no active scope exists at the time of validation;
|
||||
/// - any identity field of the captured scope differs from the current active scope;
|
||||
/// - the generation counter has advanced (a workspace switch occurred).
|
||||
///
|
||||
/// Used by `ensure_relay_mesh_for_record` after capturing scope + discovering
|
||||
/// the bootstrap target. Tests inject `install` directly so no port is touched.
|
||||
/// Delegates to [`install_client_under_workspace_transition_with_hook`] with a
|
||||
/// no-op pre-acquisition hook. See that function for full documentation.
|
||||
pub(crate) async fn install_client_under_workspace_transition<R, I, Fut>(
|
||||
app: &AppHandle<R>,
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
@@ -168,8 +180,39 @@ where
|
||||
R: tauri::Runtime,
|
||||
I: FnOnce() -> Fut,
|
||||
Fut: Future<Output = Result<(), String>>,
|
||||
{
|
||||
install_client_under_workspace_transition_with_hook(app, captured_scope, || {}, install).await
|
||||
}
|
||||
|
||||
/// Inner implementation of [`install_client_under_workspace_transition`] with an
|
||||
/// injectable `pre_acquisition_hook`.
|
||||
///
|
||||
/// `pre_acquisition_hook` fires once, synchronously, immediately before
|
||||
/// `workspace_transition.lock().await`. In production this is `|| {}`; tests
|
||||
/// inject a closure that signals "I'm about to acquire" so the test can
|
||||
/// establish a deterministic ordering between holder and contender.
|
||||
///
|
||||
/// Fails if:
|
||||
/// - no active scope exists at the time of validation;
|
||||
/// - any identity field of the captured scope differs from the current active scope;
|
||||
/// - the generation counter has advanced (a workspace switch occurred).
|
||||
///
|
||||
/// This is `pub(crate)` so tests in `mesh_llm_transition_tests` can drive the
|
||||
/// exact lock-acquisition boundary while remaining invisible to external callers.
|
||||
pub(crate) async fn install_client_under_workspace_transition_with_hook<R, H, I, Fut>(
|
||||
app: &AppHandle<R>,
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
pre_acquisition_hook: H,
|
||||
install: I,
|
||||
) -> Result<(), String>
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
H: FnOnce(),
|
||||
I: FnOnce() -> Fut,
|
||||
Fut: Future<Output = Result<(), String>>,
|
||||
{
|
||||
let state = app.state::<AppState>();
|
||||
pre_acquisition_hook();
|
||||
let _transition_guard = state.workspace_transition.lock().await;
|
||||
|
||||
// Validate the captured scope's generation and full identity under the guard.
|
||||
|
||||
@@ -91,12 +91,20 @@ async fn test_active_client_stop_then_transition_preflight_succeeds() {
|
||||
/// the lock is held), the install task acquires the lock and detects the stale
|
||||
/// captured scope without invoking the install closure.
|
||||
///
|
||||
/// Both sides use production helpers — no direct `workspace_transition` lock
|
||||
/// acquisition. Channels establish entry/release ordering without `sleep`.
|
||||
/// Handshake:
|
||||
/// 1. Transition task acquires the lock, signals "lock_held" BEFORE committing
|
||||
/// scope B (so install's captured scope is still valid at its capture time).
|
||||
/// 2. Install task spawns, calls `install_client_under_workspace_transition`
|
||||
/// with a pre-acquisition hook that signals "at_lock_boundary", then
|
||||
/// blocks on `workspace_transition.lock().await`.
|
||||
/// 3. Test waits for "at_lock_boundary", then sends "release" to the
|
||||
/// transition body. The body commits scope B and returns, releasing the lock.
|
||||
/// 4. Install task acquires the lock, re-validates scope, detects the stale
|
||||
/// generation, returns Err without invoking the install closure.
|
||||
///
|
||||
/// Proves the first serialization direction: a workspace switch (transition
|
||||
/// side) commits between scope-capture and install lock-acquisition; the
|
||||
/// install helper must reject the stale scope under the lock.
|
||||
/// Invariant: if the contender could bypass the lock, it would acquire before
|
||||
/// the transition commits scope B, see a valid scope, and invoke the install
|
||||
/// closure — `install_was_called` would be true and the assertion would fail.
|
||||
#[tokio::test]
|
||||
async fn test_transition_held_queued_install_detects_stale_scope() {
|
||||
use crate::managed_agents::scope::{
|
||||
@@ -130,23 +138,32 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
|
||||
WorkspaceAgentScope::new(relay_a.to_string(), owner_a.clone(), &base, gen_a);
|
||||
state.commit_active_scope(captured_scope.clone());
|
||||
|
||||
// Channels: transition body signals "entered" once it holds the lock;
|
||||
// test unblocks the body once the install task is queued.
|
||||
let (body_entered_tx, body_entered_rx) = oneshot::channel::<()>();
|
||||
// Channels:
|
||||
// lock_held: transition body → test (lock acquired, scope NOT yet committed)
|
||||
// at_lock_boundary: install hook → test (install is about to call lock().await)
|
||||
// body_release: test → transition body (ok to commit scope B and release)
|
||||
let (lock_held_tx, lock_held_rx) = oneshot::channel::<()>();
|
||||
let (at_lock_boundary_tx, at_lock_boundary_rx) = oneshot::channel::<()>();
|
||||
let (body_release_tx, body_release_rx) = oneshot::channel::<()>();
|
||||
|
||||
// ── Side A: transition side uses the production helper ───────────────────
|
||||
// The body holds the lock, advances the scope, signals entry, then blocks
|
||||
// until unblocked by the test — proving the hold is production-owned.
|
||||
// ── Side A: transition side — signals lock_held BEFORE committing scope B ─
|
||||
let app_a = app_handle.clone();
|
||||
let base_a = base.clone();
|
||||
let transition_task = tokio::task::spawn(async move {
|
||||
let app_a_ref = app_a.clone();
|
||||
super::scope_impl::with_workspace_transition_preflight(&app_a_ref, move || {
|
||||
Box::pin(async move {
|
||||
// Advance generation and commit a DISTINCT scope (different relay)
|
||||
// to simulate a committed workspace switch while the lock is held.
|
||||
let state_a = app_a.state::<crate::app_state::AppState>();
|
||||
|
||||
// Signal "lock held" BEFORE committing scope B.
|
||||
// The install task captures its scope before this point;
|
||||
// only after it queues at the lock does the holder commit B.
|
||||
let _ = lock_held_tx.send(());
|
||||
|
||||
// Wait for the install task to queue at the lock boundary.
|
||||
let _ = body_release_rx.await;
|
||||
|
||||
// Now commit scope B (advances generation, install sees stale).
|
||||
let gen_b = next_scope_generation();
|
||||
let new_scope = WorkspaceAgentScope::new(
|
||||
"wss://transition-test-b.example".to_string(),
|
||||
@@ -156,30 +173,29 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
|
||||
);
|
||||
state_a.commit_active_scope(new_scope);
|
||||
|
||||
// Signal: holding the lock with committed scope change.
|
||||
let _ = body_entered_tx.send(());
|
||||
// Block until the test confirms the install task has queued.
|
||||
let _ = body_release_rx.await;
|
||||
|
||||
Ok::<(), String>(())
|
||||
})
|
||||
})
|
||||
.await
|
||||
});
|
||||
|
||||
// Wait until the transition body signals it has the lock.
|
||||
body_entered_rx
|
||||
// Wait until the transition body holds the lock (before scope B is committed).
|
||||
lock_held_rx
|
||||
.await
|
||||
.expect("transition body must signal entry");
|
||||
.expect("transition body must signal lock_held");
|
||||
|
||||
// ── Side B: install side uses the production helper ──────────────────────
|
||||
// ── Side B: install side — pre-acquisition hook signals at_lock_boundary ──
|
||||
let install_was_called = Arc::new(AtomicBool::new(false));
|
||||
let install_called_clone = Arc::clone(&install_was_called);
|
||||
let app_b = app_handle.clone();
|
||||
let install_task = tokio::task::spawn(async move {
|
||||
super::scope_impl::install_client_under_workspace_transition(
|
||||
super::scope_impl::install_client_under_workspace_transition_with_hook(
|
||||
&app_b,
|
||||
&captured_scope,
|
||||
// pre-acquisition hook: fires before lock().await — signals "queued".
|
||||
move || {
|
||||
let _ = at_lock_boundary_tx.send(());
|
||||
},
|
||||
|| {
|
||||
let called = Arc::clone(&install_called_clone);
|
||||
async move {
|
||||
@@ -191,10 +207,13 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
|
||||
.await
|
||||
});
|
||||
|
||||
// Yield to let the install task queue on the lock.
|
||||
tokio::task::yield_now().await;
|
||||
// Wait for install to reach the lock boundary, then unblock the transition.
|
||||
at_lock_boundary_rx
|
||||
.await
|
||||
.expect("install hook must signal at_lock_boundary");
|
||||
|
||||
// Unblock the transition body → it releases the lock → install acquires it.
|
||||
// Unblock the transition body → it commits scope B, releases the lock →
|
||||
// install acquires the lock, detects stale scope, returns Err.
|
||||
let _ = body_release_tx.send(());
|
||||
|
||||
let install_result = install_task.await.expect("install task must not panic");
|
||||
@@ -226,12 +245,22 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
|
||||
/// install body completes, the transition task acquires the lock, runs
|
||||
/// `fail_if_client_mesh_active`, and observes the installed client.
|
||||
///
|
||||
/// Both sides use production helpers — no direct `workspace_transition` lock
|
||||
/// acquisition. Channels establish entry/release ordering without `sleep`.
|
||||
/// Handshake:
|
||||
/// 1. Install task acquires the lock, signals "lock_held" BEFORE installing the
|
||||
/// client runtime (so the client is not yet visible to the contender).
|
||||
/// 2. Transition task calls `with_workspace_transition_preflight` with a
|
||||
/// pre-acquisition hook that signals "at_lock_boundary", then blocks on
|
||||
/// `workspace_transition.lock().await`.
|
||||
/// 3. Test waits for "at_lock_boundary", then sends "release" to the
|
||||
/// install closure. The closure installs the client and returns, releasing
|
||||
/// the lock.
|
||||
/// 4. Transition task acquires the lock, runs `fail_if_client_mesh_active`,
|
||||
/// observes the installed client, returns Err.
|
||||
///
|
||||
/// Proves the second serialization direction: a client install commits between
|
||||
/// the transition's pre-lock preflight and its lock-acquisition; the transition
|
||||
/// helper must observe the installed client under the lock and fail.
|
||||
/// Invariant: if the contender could bypass the lock, it would run
|
||||
/// `fail_if_client_mesh_active` before the client is installed — seeing no
|
||||
/// client — and return Ok. The `transition_result.is_err()` assertion would
|
||||
/// then fail, proving the lock is not enforced.
|
||||
#[tokio::test]
|
||||
async fn test_install_held_transition_preflight_observes_client() {
|
||||
use crate::managed_agents::scope::{
|
||||
@@ -262,14 +291,15 @@ async fn test_install_held_transition_preflight_observes_client() {
|
||||
let scope = WorkspaceAgentScope::new(relay.to_string(), owner.clone(), &base, gen);
|
||||
state.commit_active_scope(scope.clone());
|
||||
|
||||
// Channels: install body signals "entered + client installed" once it holds
|
||||
// the lock; test unblocks the body once the transition task is queued.
|
||||
let (install_entered_tx, install_entered_rx) = oneshot::channel::<()>();
|
||||
// Channels:
|
||||
// lock_held: install body → test (lock acquired, client NOT yet installed)
|
||||
// at_lock_boundary: transition hook → test (transition is about to call lock().await)
|
||||
// install_release: test → install body (ok to install client and release lock)
|
||||
let (lock_held_tx, lock_held_rx) = oneshot::channel::<()>();
|
||||
let (at_lock_boundary_tx, at_lock_boundary_rx) = oneshot::channel::<()>();
|
||||
let (install_release_tx, install_release_rx) = oneshot::channel::<()>();
|
||||
|
||||
// ── Side A: install side uses the production helper ──────────────────────
|
||||
// The install closure installs a mock client, signals entry, then blocks
|
||||
// until unblocked — proving the client is committed under the lock.
|
||||
// ── Side A: install side — signals lock_held BEFORE installing the client ─
|
||||
let app_a = app_handle.clone();
|
||||
let scope_a = scope.clone();
|
||||
let install_task = tokio::task::spawn(async move {
|
||||
@@ -278,40 +308,51 @@ async fn test_install_held_transition_preflight_observes_client() {
|
||||
&app_a,
|
||||
&scope_a,
|
||||
move || async move {
|
||||
// Install the mock client runtime while holding the lock.
|
||||
// Signal "lock held" BEFORE installing the client.
|
||||
// The transition task captures its pre-lock state before this
|
||||
// point; only after it queues does the install closure install.
|
||||
let _ = lock_held_tx.send(());
|
||||
|
||||
// Wait for the transition task to queue at the lock boundary.
|
||||
let _ = install_release_rx.await;
|
||||
|
||||
// Now install the mock client runtime under the held lock.
|
||||
let state_a = app_a_for_closure.state::<crate::app_state::AppState>();
|
||||
let client_runtime = crate::mesh_llm::build_mock_client_runtime_for_test();
|
||||
*state_a.mesh_llm_runtime.lock().await = Some(client_runtime);
|
||||
|
||||
// Signal: lock held, client installed.
|
||||
let _ = install_entered_tx.send(());
|
||||
// Block until the test confirms the transition task has queued.
|
||||
let _ = install_release_rx.await;
|
||||
|
||||
Ok::<(), String>(())
|
||||
},
|
||||
)
|
||||
.await
|
||||
});
|
||||
|
||||
// Wait until the install body signals it has the lock with the client installed.
|
||||
install_entered_rx
|
||||
// Wait until the install body holds the lock (before client is installed).
|
||||
lock_held_rx
|
||||
.await
|
||||
.expect("install body must signal entry");
|
||||
.expect("install body must signal lock_held");
|
||||
|
||||
// ── Side B: transition side uses the production helper ───────────────────
|
||||
// ── Side B: transition side — pre-acquisition hook signals at_lock_boundary
|
||||
let app_b = app_handle.clone();
|
||||
let transition_task = tokio::task::spawn(async move {
|
||||
super::scope_impl::with_workspace_transition_preflight(&app_b, || {
|
||||
Box::pin(async { Ok::<&str, String>("body ran") })
|
||||
})
|
||||
super::scope_impl::with_workspace_transition_preflight_with_hook(
|
||||
&app_b,
|
||||
// pre-acquisition hook: fires before lock().await — signals "queued".
|
||||
move || {
|
||||
let _ = at_lock_boundary_tx.send(());
|
||||
},
|
||||
|| Box::pin(async { Ok::<&str, String>("body ran") }),
|
||||
)
|
||||
.await
|
||||
});
|
||||
|
||||
// Yield to let the transition task queue on the lock.
|
||||
tokio::task::yield_now().await;
|
||||
// Wait for transition to reach the lock boundary, then unblock the install.
|
||||
at_lock_boundary_rx
|
||||
.await
|
||||
.expect("transition hook must signal at_lock_boundary");
|
||||
|
||||
// Unblock the install body → it releases the lock → transition acquires it.
|
||||
// Unblock the install closure → it installs the client, releases the lock →
|
||||
// transition acquires the lock, observes the client, returns Err.
|
||||
let _ = install_release_tx.send(());
|
||||
|
||||
install_task
|
||||
|
||||
@@ -236,7 +236,21 @@ pub(crate) fn start_managed_agent_runtime_pair_lazy(
|
||||
relay_url: String,
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
start_pair(pubkey, relay_url, true, None, app)
|
||||
start_pair_lazy_for(pubkey, relay_url, app)
|
||||
}
|
||||
|
||||
/// Generic start-pair-lazy seam shared by the production adapter and tests.
|
||||
///
|
||||
/// Acquires `managed_agent_runtime_transition` as its first action — the same
|
||||
/// lock that `stop`, `restart`, and `drain` operations hold, serialising all
|
||||
/// runtime mutations. Tests that need a mock-runtime contender call this
|
||||
/// function directly instead of the non-generic production adapter.
|
||||
pub(crate) fn start_pair_lazy_for<R: tauri::Runtime>(
|
||||
pubkey: String,
|
||||
relay_url: String,
|
||||
app: tauri::AppHandle<R>,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
start_pair_for(pubkey, relay_url, true, None, app)
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
@@ -254,6 +268,16 @@ fn start_pair(
|
||||
lazy: bool,
|
||||
expected_updated_at: Option<&str>,
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
start_pair_for(pubkey, relay_url, lazy, expected_updated_at, app)
|
||||
}
|
||||
|
||||
fn start_pair_for<R: tauri::Runtime>(
|
||||
pubkey: String,
|
||||
relay_url: String,
|
||||
lazy: bool,
|
||||
expected_updated_at: Option<&str>,
|
||||
app: tauri::AppHandle<R>,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
let state = app.state::<AppState>();
|
||||
let _transition = state
|
||||
@@ -883,32 +907,45 @@ where
|
||||
|
||||
/// Production adapter for [`compensate_drain_for`].
|
||||
///
|
||||
/// Lock acquisition order inside this function:
|
||||
/// 1. `_rt_transition_held` — already held by caller (passed by value).
|
||||
/// 2. Acquire `managed_agents_store_lock` — store-lock-only writers
|
||||
/// (instance edits, status flushes) are serialized here, not on the
|
||||
/// transition guard.
|
||||
/// 3. Validate `captured_scope.generation` under the store lock.
|
||||
/// 4. `load_managed_agents_at` — records loaded AFTER both locks held.
|
||||
/// 5. Delegate to [`compensate_drain_for`] — store guard held through
|
||||
/// every save performed by `start_pair_under_held_locks`.
|
||||
/// Delegates to [`compensate_drain_with_hook`] with a no-op hook.
|
||||
/// Lock acquisition order: transition guard already held → store lock acquired
|
||||
/// → hook fires (no-op in production) → validate generation → load → restore/save.
|
||||
pub(crate) fn compensate_drain<R: tauri::Runtime>(
|
||||
app: &tauri::AppHandle<R>,
|
||||
stopped: &[DrainJournalEntry],
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
// Caller passes ownership of the already-held transition guard so the lock
|
||||
// is never dropped between drain and compensation.
|
||||
_rt_transition_held: std::sync::MutexGuard<'_, ()>,
|
||||
) -> Option<String> {
|
||||
compensate_drain_with_hook(app, stopped, captured_scope, _rt_transition_held, || {})
|
||||
}
|
||||
|
||||
/// Inner implementation of [`compensate_drain`] with an injectable
|
||||
/// `on_store_acquired` hook.
|
||||
///
|
||||
/// Lock acquisition order:
|
||||
/// 1. `_rt_transition_held` — already held by caller (passed by value).
|
||||
/// 2. Acquire `managed_agents_store_lock`.
|
||||
/// 3. Call `on_store_acquired` (no-op in production; tests inject a closure
|
||||
/// that writes a sentinel and synchronises with a writer thread to prove
|
||||
/// genuine store-lock contention).
|
||||
/// 4. Validate `captured_scope.generation` under the store lock.
|
||||
/// 5. `load_managed_agents_at` — records loaded AFTER both locks held.
|
||||
/// 6. Delegate to [`compensate_drain_for`] — store guard held through every
|
||||
/// save performed by `start_pair_under_held_locks`.
|
||||
pub(crate) fn compensate_drain_with_hook<R: tauri::Runtime>(
|
||||
app: &tauri::AppHandle<R>,
|
||||
stopped: &[DrainJournalEntry],
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
_rt_transition_held: std::sync::MutexGuard<'_, ()>,
|
||||
on_store_acquired: impl FnOnce(),
|
||||
) -> Option<String> {
|
||||
if stopped.is_empty() {
|
||||
// Release the transition guard immediately — nothing to restore.
|
||||
drop(_rt_transition_held);
|
||||
return None;
|
||||
}
|
||||
|
||||
let state = app.state::<AppState>();
|
||||
|
||||
// 2. Acquire the store lock — transition already held via _rt_transition_held.
|
||||
let _store = match state.managed_agents_store_lock.lock() {
|
||||
Ok(g) => g,
|
||||
Err(e) => {
|
||||
@@ -918,10 +955,10 @@ pub(crate) fn compensate_drain<R: tauri::Runtime>(
|
||||
}
|
||||
};
|
||||
|
||||
// 3. Validate the captured scope before restoring. If the workspace switched
|
||||
// between the drain failure and this compensation, the new scope's own
|
||||
// restore pass will start the correct agents — we must not restart agents
|
||||
// for a scope that is no longer active.
|
||||
// 3. Fire the hook while both locks are held. No-op in production.
|
||||
on_store_acquired();
|
||||
|
||||
// 4. Validate the captured scope under the store lock.
|
||||
if let Err(stale_msg) = crate::managed_agents::scope::validate_scope_generation(captured_scope)
|
||||
{
|
||||
return Some(format!(
|
||||
@@ -929,7 +966,7 @@ pub(crate) fn compensate_drain<R: tauri::Runtime>(
|
||||
));
|
||||
}
|
||||
|
||||
// 4. Load records under the held store lock.
|
||||
// 5. Load records under the held store lock.
|
||||
let mut records = match load_managed_agents_at(&captured_scope.definitions_dir) {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
@@ -939,7 +976,7 @@ pub(crate) fn compensate_drain<R: tauri::Runtime>(
|
||||
}
|
||||
};
|
||||
|
||||
// 5. Delegate to the lock-free core — store guard is held through every save.
|
||||
// 6. Delegate — store guard held through every save.
|
||||
compensate_drain_for(stopped, &mut records, |entry, recs| {
|
||||
start_pair_under_held_locks(
|
||||
app,
|
||||
|
||||
@@ -6,26 +6,28 @@
|
||||
|
||||
use super::*;
|
||||
|
||||
/// Deterministic writer-vs-compensation ordering via the production
|
||||
/// `compensate_drain` adapter.
|
||||
/// Writer-vs-compensation store-lock contention via the `compensate_drain_with_hook` seam.
|
||||
///
|
||||
/// Mechanism:
|
||||
/// 1. The test acquires the `managed_agent_runtime_transition` guard.
|
||||
/// 2. A writer thread is spawned. It waits on a channel before touching the
|
||||
/// store, so its load→edit→save only begins AFTER the test explicitly
|
||||
/// signals "compensation done."
|
||||
/// 3. The test calls the production `compensate_drain` adapter — which
|
||||
/// acquires `managed_agents_store_lock` internally, validates generation,
|
||||
/// loads records, delegates to `compensate_drain_for`, and saves — all while
|
||||
/// the transition guard is owned by compensate_drain (passed by value).
|
||||
/// 4. After `compensate_drain` returns the transition guard is consumed; the
|
||||
/// test signals the writer via the channel.
|
||||
/// 5. The writer acquires the store lock, applies its sentinel edit, and saves.
|
||||
/// 6. Final disk state must contain BOTH effects: compensation's `runtime_pid`
|
||||
/// update AND the writer's `WRITER_EDIT` env-var sentinel.
|
||||
/// Invariant: if the `managed_agents_store_lock` guard is dropped before
|
||||
/// `compensate_drain_for`'s restore/save, the writer can interleave and
|
||||
/// overwrite `COMP_SENTINEL` — making the final assertion fail.
|
||||
///
|
||||
/// Ordering is established by construction (channel), not by scheduler timing.
|
||||
/// "both effects" proves the store guard is correctly held across restore/save.
|
||||
/// Flow:
|
||||
/// 1. Seed the store with one agent record.
|
||||
/// 2. Acquire `managed_agent_runtime_transition`; call `compensate_drain_with_hook`.
|
||||
/// 3. `on_store_acquired` hook (fires while the store lock is held):
|
||||
/// a. writes `COMP_SENTINEL` to disk;
|
||||
/// b. signals the writer thread (`store_acquired_tx`);
|
||||
/// c. waits for the writer to queue at the store lock (`writer_queued_rx`).
|
||||
/// 4. Writer thread: waits for (b), signals (c), then blocks on `managed_agents_store_lock`.
|
||||
/// Compensation still holds the lock — writer is blocked.
|
||||
/// 5. `compensate_drain_with_hook` continues: validate → load → restore → save.
|
||||
/// Store lock still held throughout; writer remains blocked.
|
||||
/// 6. Function returns; store lock released. Writer acquires it, writes `WRITER_EDIT`, saves.
|
||||
/// 7. Final disk record must contain BOTH `COMP_SENTINEL` AND `WRITER_EDIT`.
|
||||
///
|
||||
/// If the guard is dropped early, the writer acquires the lock before save,
|
||||
/// overwrites `COMP_SENTINEL`, and the assertion at step 7 fails.
|
||||
#[test]
|
||||
fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
||||
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
|
||||
@@ -98,16 +100,13 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
// Build a mock app so `compensate_drain` can reach `AppState`.
|
||||
let app = tauri::test::mock_builder()
|
||||
.manage(crate::app_state::build_app_state())
|
||||
.build(tauri::test::mock_context(tauri::test::noop_assets()))
|
||||
.expect("failed to build mock app");
|
||||
let app_handle = app.app_handle().clone();
|
||||
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
|
||||
// Commit a scope pointing at our tempdir so generation validation passes.
|
||||
let gen = current_scope_generation();
|
||||
let scope = WorkspaceAgentScope {
|
||||
scope_id: "comp-drain-writer-test".to_string(),
|
||||
@@ -121,21 +120,26 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
||||
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
||||
let stopped = vec![entry1];
|
||||
|
||||
// Channel: writer waits until compensate_drain signals "done".
|
||||
let (comp_done_tx, comp_done_rx) = std::sync::mpsc::channel::<()>();
|
||||
// store_acquired: hook → writer (compensation holds the store lock; writer may queue)
|
||||
// writer_queued: writer → hook (writer is now blocked on the store lock)
|
||||
let (store_acquired_tx, store_acquired_rx) = std::sync::mpsc::channel::<()>();
|
||||
let (writer_queued_tx, writer_queued_rx) = std::sync::mpsc::channel::<()>();
|
||||
|
||||
// Spawn the writer BEFORE acquiring the transition guard to avoid a
|
||||
// deadlock — the writer only needs the store lock, not the transition lock,
|
||||
// but it waits on the channel first.
|
||||
// Spawn the writer thread BEFORE acquiring the transition guard.
|
||||
// It waits for store_acquired, signals writer_queued (it is now about to
|
||||
// block on the store lock), then acquires the store lock and writes WRITER_EDIT.
|
||||
let tmp_wr = tmp_path.clone();
|
||||
let app_handle_wr = app_handle.clone();
|
||||
let wr_thread = thread::spawn(move || {
|
||||
// Established ordering: writer explicitly waits until compensation
|
||||
// signals it is done — compensate_drain holds both transition guard
|
||||
// (passed by value) and store lock through restore/save.
|
||||
comp_done_rx.recv().unwrap();
|
||||
// Wait until the hook signals that compensation holds the store lock.
|
||||
store_acquired_rx.recv().unwrap();
|
||||
|
||||
// Now acquire just the store lock and apply the sentinel edit.
|
||||
// Signal the hook that we are about to queue on the store lock.
|
||||
// The hook will return on receipt, letting compensate_drain_with_hook
|
||||
// proceed to validate → load → save while we are blocked below.
|
||||
writer_queued_tx.send(()).unwrap();
|
||||
|
||||
// Block on the store lock — compensation still holds it.
|
||||
let writer_state = app_handle_wr.state::<crate::app_state::AppState>();
|
||||
let _store = writer_state.managed_agents_store_lock.lock().unwrap();
|
||||
let mut records =
|
||||
@@ -147,31 +151,40 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
||||
crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).unwrap();
|
||||
});
|
||||
|
||||
// Acquire the transition guard in this thread and pass it to the
|
||||
// production adapter. compensate_drain takes ownership of the guard so it
|
||||
// is held for the full validate→load→restore→save sequence.
|
||||
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
|
||||
|
||||
// ── Phase-A: run the production compensate_drain adapter ─────────────────
|
||||
// `start_pair_under_held_locks` is not available to inject from outside
|
||||
// the crate; compensate_drain calls it internally. Because there is no
|
||||
// live process for the dummy record the start attempt will fail with an
|
||||
// error — the adapter returns a degradation message. We verify the final
|
||||
// disk state rather than the return value of compensate_drain.
|
||||
//
|
||||
// To prove the guard was actually held through save we assert the writer
|
||||
// does NOT see intermediate state.
|
||||
let _comp_result = compensate_drain(&app_handle, &stopped, &scope, rt_guard);
|
||||
// on_store_acquired fires while the store lock is held:
|
||||
// (a) write COMP_SENTINEL to disk — proves we are mid-sequence under the lock;
|
||||
// (b) signal the writer (store_acquired_tx);
|
||||
// (c) wait for the writer to queue at the lock boundary (writer_queued_rx).
|
||||
// Returning lets compensate_drain_with_hook continue to validate → load → save,
|
||||
// all while the store lock remains held, keeping the writer blocked.
|
||||
let _comp_result = compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, || {
|
||||
// (a) Write the compensation sentinel under the held store lock.
|
||||
let mut records =
|
||||
crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default();
|
||||
for r in &mut records {
|
||||
r.env_vars
|
||||
.insert("COMP_SENTINEL".to_string(), "yes".to_string());
|
||||
}
|
||||
crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &records).unwrap();
|
||||
|
||||
// Signal the writer that compensation is done (guard released).
|
||||
comp_done_tx.send(()).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");
|
||||
|
||||
// ── Final disk state: BOTH effects must be present ───────────────────────
|
||||
// Compensation's effect: the record exists on disk (compensate_drain loaded
|
||||
// it and saved after the restore attempt, regardless of start success).
|
||||
// Writer's effect: WRITER_EDIT sentinel present.
|
||||
// COMP_SENTINEL: written by the hook while compensation held the store lock.
|
||||
// WRITER_EDIT: written by the writer after compensation released the lock.
|
||||
//
|
||||
// If the store guard is dropped before compensate_drain_for's save, the
|
||||
// writer can acquire the lock early and overwrite COMP_SENTINEL — this
|
||||
// assertion would then fail, proving the invariant.
|
||||
let final_records =
|
||||
crate::managed_agents::storage::load_managed_agents_at(&tmp_path).unwrap_or_default();
|
||||
let final_rec = final_records
|
||||
@@ -179,32 +192,43 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
|
||||
.find(|r| r.pubkey == pubkey1)
|
||||
.expect("agent record must be present on disk after both phases");
|
||||
|
||||
// Writer's effect must always be present.
|
||||
assert_eq!(
|
||||
final_rec.env_vars.get("COMP_SENTINEL").map(String::as_str),
|
||||
Some("yes"),
|
||||
"COMP_SENTINEL must survive — fails if the store guard is dropped before \
|
||||
compensate_drain_for's save, allowing the writer to overwrite it"
|
||||
);
|
||||
assert_eq!(
|
||||
final_rec.env_vars.get("WRITER_EDIT").map(String::as_str),
|
||||
Some("yes"),
|
||||
"writer's WRITER_EDIT sentinel must be present in final disk state"
|
||||
"WRITER_EDIT must be present — the writer runs after compensation releases the store lock"
|
||||
);
|
||||
}
|
||||
|
||||
/// Production start-path contender is blocked while `compensate_drain` holds
|
||||
/// the `managed_agent_runtime_transition` mutex.
|
||||
/// Production start-path contender is blocked while `compensate_drain_with_hook`
|
||||
/// holds the `managed_agent_runtime_transition` mutex.
|
||||
///
|
||||
/// Invariant: if the `managed_agent_runtime_transition` guard is dropped before
|
||||
/// `compensate_drain_with_hook` finishes, `start_managed_agent_runtime_pair_lazy`
|
||||
/// acquires the transition lock and proceeds — `contender_done_rx` would fire
|
||||
/// before the hook completes, and the assertion at the end would fail.
|
||||
///
|
||||
/// Flow:
|
||||
/// 1. This test thread acquires `managed_agent_runtime_transition` FIRST.
|
||||
/// 2. A "start contender" thread is spawned; it signals "alive," then blocks
|
||||
/// trying to acquire the same transition mutex (this is the production
|
||||
/// start-lock path — `start_managed_agent` and related production commands
|
||||
/// also take the transition guard before modifying runtimes).
|
||||
/// 3. `compensate_drain` is called with the already-held guard (passes it
|
||||
/// by value into the adapter — drains the stopped-entry list).
|
||||
/// 4. After `compensate_drain` returns the guard is consumed; the contender
|
||||
/// acquires the mutex and signals "done."
|
||||
/// 1. Seed the store; acquire `managed_agent_runtime_transition` in this thread.
|
||||
/// 2. Spawn the contender. It signals "ready," then waits for the hook to
|
||||
/// signal "store held." It then signals "at_lock" and calls the production
|
||||
/// `start_managed_agent_runtime_pair_lazy` — which blocks acquiring
|
||||
/// `managed_agent_runtime_transition` as its first action.
|
||||
/// 3. `compensate_drain_with_hook` hook: signals the contender ("store held"),
|
||||
/// waits for "at_lock" confirmation, then returns. The contender is now
|
||||
/// queued at the transition lock.
|
||||
/// 4. `compensate_drain_with_hook` completes, releases the transition guard.
|
||||
/// 5. The contender acquires the lock (fails at spawn — no binary — but
|
||||
/// passes the lock boundary), signals "done."
|
||||
///
|
||||
/// Contender/queue order is established by barriers/channels — no `sleep`.
|
||||
/// The contender drives the same `managed_agent_runtime_transition` lock path
|
||||
/// that production `start_managed_agent` commands use, proving compensation
|
||||
/// blocks production starts.
|
||||
/// The contender drives `start_managed_agent_runtime_pair_lazy`, the same
|
||||
/// function called by the `start_managed_agent_runtime` Tauri command, proving
|
||||
/// that compensation blocks production starts — not just a raw mutex lock.
|
||||
#[test]
|
||||
fn test_compensate_drain_concurrent_start_is_blocked() {
|
||||
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
|
||||
@@ -213,14 +237,74 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let tmp_path = tmp.path().to_path_buf();
|
||||
crate::managed_agents::storage::save_managed_agents_at(&tmp_path, &[]).unwrap();
|
||||
|
||||
let pubkey1 = "aa".repeat(32);
|
||||
let initial_record = crate::managed_agents::ManagedAgentRecord {
|
||||
pubkey: pubkey1.clone(),
|
||||
name: "contender-agent".to_string(),
|
||||
display_name: None,
|
||||
slug: None,
|
||||
persona_id: None,
|
||||
private_key_nsec: String::new(),
|
||||
auth_tag: None,
|
||||
relay_url: "wss://relay.example".to_string(),
|
||||
avatar_url: None,
|
||||
acp_command: crate::managed_agents::DEFAULT_ACP_COMMAND.to_string(),
|
||||
agent_command: String::new(),
|
||||
agent_command_override: None,
|
||||
agent_args: vec![],
|
||||
mcp_command: String::new(),
|
||||
turn_timeout_seconds: 0,
|
||||
idle_timeout_seconds: None,
|
||||
max_turn_duration_seconds: None,
|
||||
parallelism: 1,
|
||||
system_prompt: None,
|
||||
model: None,
|
||||
provider: None,
|
||||
persona_source_version: None,
|
||||
env_vars: Default::default(),
|
||||
start_on_app_launch: true,
|
||||
auto_restart_on_config_change: false,
|
||||
runtime_pid: None,
|
||||
backend: crate::managed_agents::BackendKind::Local,
|
||||
backend_agent_id: None,
|
||||
provider_binary_path: None,
|
||||
team_id: None,
|
||||
persona_team_dir: None,
|
||||
persona_name_in_team: None,
|
||||
created_at: crate::util::now_iso(),
|
||||
updated_at: crate::util::now_iso(),
|
||||
last_started_at: None,
|
||||
last_stopped_at: None,
|
||||
last_exit_code: None,
|
||||
last_error: None,
|
||||
last_error_code: None,
|
||||
respond_to: Default::default(),
|
||||
respond_to_allowlist: Default::default(),
|
||||
is_builtin: false,
|
||||
is_active: true,
|
||||
shared: false,
|
||||
source_team: None,
|
||||
source_team_persona_slug: None,
|
||||
catalog_source: None,
|
||||
definition_respond_to: None,
|
||||
definition_respond_to_allowlist: Default::default(),
|
||||
definition_parallelism: None,
|
||||
relay_mesh: None,
|
||||
runtime: None,
|
||||
name_pool: vec![],
|
||||
};
|
||||
crate::managed_agents::storage::save_managed_agents_at(
|
||||
&tmp_path,
|
||||
std::slice::from_ref(&initial_record),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let app = tauri::test::mock_builder()
|
||||
.manage(crate::app_state::build_app_state())
|
||||
.build(tauri::test::mock_context(tauri::test::noop_assets()))
|
||||
.expect("failed to build mock app");
|
||||
let app_handle = app.app_handle().clone();
|
||||
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
|
||||
let gen = current_scope_generation();
|
||||
@@ -233,54 +317,65 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
|
||||
};
|
||||
state.commit_active_scope(scope.clone());
|
||||
|
||||
let pubkey1 = "aa".repeat(32);
|
||||
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
|
||||
let stopped = vec![entry1];
|
||||
|
||||
// Channels for barrier-based ordering.
|
||||
// contender_ready: contender → test (about to call production start seam)
|
||||
// hook_store_held: hook → contender (hook has store lock; contender may proceed)
|
||||
// contender_at_lock: contender → hook (contender queued at transition lock)
|
||||
// contender_done: contender → test (contender passed the transition lock)
|
||||
let (contender_ready_tx, contender_ready_rx) = std::sync::mpsc::channel::<()>();
|
||||
let (hook_store_held_tx, hook_store_held_rx) = std::sync::mpsc::channel::<()>();
|
||||
let (contender_at_lock_tx, contender_at_lock_rx) = std::sync::mpsc::channel::<()>();
|
||||
let (contender_done_tx, contender_done_rx) = std::sync::mpsc::channel::<()>();
|
||||
|
||||
// Acquire the transition guard in THIS thread FIRST so the contender will
|
||||
// block when it tries to acquire the same mutex.
|
||||
// Acquire the transition guard FIRST so the contender will block on it.
|
||||
let rt_guard = state.managed_agent_runtime_transition.lock().unwrap();
|
||||
|
||||
// The contender drives the transition mutex the same way the production
|
||||
// `start_managed_agent` family does: it acquires the guard, then signals.
|
||||
let app_contender = app_handle.clone();
|
||||
let pubkey_contender = pubkey1.clone();
|
||||
let contender = thread::spawn(move || {
|
||||
// Signal that the contender is alive and about to block on the mutex.
|
||||
// Signal: about to call the production start path.
|
||||
contender_ready_tx.send(()).unwrap();
|
||||
// This is the production start-lock path: takes managed_agent_runtime_transition.
|
||||
let contender_state = app_contender.state::<crate::app_state::AppState>();
|
||||
let _guard = contender_state
|
||||
.managed_agent_runtime_transition
|
||||
.lock()
|
||||
.unwrap();
|
||||
// Signal: contender now holds the lock (compensation is done).
|
||||
|
||||
// Wait for the hook to confirm it holds the store lock.
|
||||
hook_store_held_rx.recv().unwrap();
|
||||
|
||||
// Signal: queued at the transition lock boundary.
|
||||
contender_at_lock_tx.send(()).unwrap();
|
||||
|
||||
// Production start-lock seam — acquires managed_agent_runtime_transition
|
||||
// as its first action, so it blocks here until compensation releases it.
|
||||
let _ = start_pair_lazy_for(
|
||||
pubkey_contender,
|
||||
"wss://relay.example".to_string(),
|
||||
app_contender,
|
||||
);
|
||||
|
||||
// Signal: passed the transition lock (compensation is done).
|
||||
contender_done_tx.send(()).unwrap();
|
||||
});
|
||||
|
||||
// Wait until the contender is alive and parked (or about to park) on the mutex.
|
||||
// Wait for the contender to be alive before calling the adapter.
|
||||
contender_ready_rx.recv().unwrap();
|
||||
|
||||
// Confirm the contender is NOT yet through the lock (it can't be — we hold it).
|
||||
// Use try_recv: if it somehow returned it would mean the mutex isn't working.
|
||||
assert!(
|
||||
contender_done_rx.try_recv().is_err(),
|
||||
"contender must be blocked while this thread holds the transition guard"
|
||||
);
|
||||
|
||||
// ── Run the production compensate_drain adapter ───────────────────────────
|
||||
// Passes the transition guard BY VALUE — `compensate_drain` takes ownership
|
||||
// and holds it through validate→load→restore→save.
|
||||
let _comp_result = compensate_drain(&app_handle, &stopped, &scope, rt_guard);
|
||||
// rt_guard is now consumed/dropped inside compensate_drain.
|
||||
// on_store_acquired hook: fires while BOTH locks are held.
|
||||
// (a) signal the contender so it proceeds to call the start seam;
|
||||
// (b) wait for the contender to confirm it is at the lock boundary.
|
||||
// After (b) the contender is guaranteed to be blocked on
|
||||
// managed_agent_runtime_transition — returning lets
|
||||
// compensate_drain_with_hook finish its work before the contender unblocks.
|
||||
let _comp_result = compensate_drain_with_hook(&app_handle, &stopped, &scope, rt_guard, || {
|
||||
hook_store_held_tx.send(()).unwrap();
|
||||
contender_at_lock_rx.recv().unwrap();
|
||||
});
|
||||
|
||||
// Contender must now acquire the mutex within a reasonable timeout.
|
||||
contender_done_rx
|
||||
.recv_timeout(std::time::Duration::from_secs(5))
|
||||
.expect("contender must unblock after compensate_drain releases the transition guard");
|
||||
.expect(
|
||||
"contender must unblock after compensate_drain_with_hook releases the transition guard",
|
||||
);
|
||||
|
||||
contender.join().expect("contender thread panicked");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user