fix(desktop): resolve P4 pass-1 review defects in workspace-scoped agent store

Fix all five IMPORTANT defects identified by Thufir's P4 pass-1 review:

1. Concurrency tests: both test_compensate_drain_writer_vs_compensation_deterministic
   and test_compensate_drain_concurrent_start_is_blocked now drive the production
   compensate_drain adapter (not compensate_drain_for directly). Transition guard
   passed by value; store lock and restore owned by the production symbol.
   Coordinator/contender ordering via channels/barriers exclusively — no
   thread::sleep.

2. Global-restart tests: cross-platform spawn_long_lived_child_for_test (sh loop /
   ping) and spawn_noop_child_for_test (sh exit 0 / cmd.exe exit 0) replace the
   UNIX-specific sleep 10000 and /usr/bin/true, fixing Windows Rust CI failures.
   Full-tail test asserts non-empty personas/teams/global in captured context.
   Context-load-failure test uses malformed JSON (not absent file) to exercise the
   genuine parse-error path. Record-mesh-change test seeds eligible runtime and
   constructs mismatched context.mesh_model_id to trigger in-epoch TOCTOU guard.

3. Active-scope identity import: import_identity routes through the production
   with_workspace_transition_preflight (mesh-llm) / direct workspace_transition
   (no mesh-llm) when has_active_scope, matching apply_workspace orchestration.

   Direction tests rewritten as helper-vs-helper: each side uses a production
   helper (with_workspace_transition_preflight or install_client_under_workspace_
   transition); oneshot channels establish entry/release ordering with no direct
   lock acquisition and no sleep.

4. Team memory boundary: setup_team_import_app_with_scope adds a memory-bearing
   team member (member_with_memory); seam test asserts both the captured HTTP
   relay URL and the captured owner p-tag from the built engram event after
   after_store commits a distinct owner/scope.

5. CI portability: setup_import_app_with_scope and setup_team_import_app_with_scope
   use WorkspaceAgentScope::new with writable temp agent base so definitions_dir
   has the production <base>/scopes/<scope_id> shape, fixing the Linux /retention
   EPERM failure and the poisoned-lock cascade in seam tests.

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
2026-08-04 23:58:02 -04:00
co-authored by Will Pfleger
parent a1228b39a1
commit 91aea2db96
7 changed files with 981 additions and 579 deletions
@@ -15,9 +15,13 @@ use super::*;
/// and written, runtimes map contains new entry with context.scope.scope_id,
/// final captured disk record matches."
///
/// Note: to pass the eligibility check we need a record with backend=Local
/// and a live runtime in the runtimes map. Since we can't inject a real
/// process, we seed the runtimes map directly via AppState.
/// Cross-platform process helpers replace `sleep 10000` / `/usr/bin/true`:
/// - Seed runtime: `spawn_long_lived_child_for_test()` (survives sync eviction).
/// - Spawn closure: `spawn_noop_child_for_test()` (exits immediately; test only
/// checks in-memory state, not process liveness).
///
/// Non-empty personas, teams, and global are placed in both the context AND
/// the captured definitions_dir so the spawn_fn can assert they arrive.
#[tokio::test]
async fn test_full_tail_stop_spawn_receipt_register_save() {
use crate::managed_agents::{
@@ -107,21 +111,14 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
let app = make_mock_app();
let app_handle = app.handle().clone();
// Seed a live pair runtime using a long-running process.
// `/usr/bin/true` exits immediately and is evicted by
// sync_managed_agent_processes before the eligibility check; use `sleep`
// to keep the runtime alive through the sync.
// Seed a live runtime with a cross-platform long-lived child (avoids sync eviction).
// `spawn_long_lived_child_for_test()` replaces `sleep 10000` / `ping -n 100000`.
let rt_key = ManagedAgentRuntimeKey::new(&pubkey, relay_url).unwrap();
{
let seeded_pid = {
let state = app_handle.state::<crate::app_state::AppState>();
let mut runtimes = state.managed_agent_processes.lock().unwrap();
let child = std::process::Command::new("sleep")
.arg("10000")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn sleep 10000");
let child = spawn_long_lived_child_for_test();
let pid = child.id();
let process = crate::managed_agents::ManagedAgentProcess {
child,
log_path: std::path::PathBuf::new(),
@@ -142,7 +139,8 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
rt_key.clone(),
ManagedAgentPairRuntime::starting(process, Some(scope_id.to_string())),
);
}
pid
};
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
@@ -153,11 +151,56 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
generation: gen,
};
// Build NON-EMPTY personas, teams, and global so the spawn_fn can assert
// they are actually delivered to the captured context.
let test_persona = crate::managed_agents::AgentDefinition {
id: "test-persona-id".to_string(),
display_name: "Test Persona".to_string(),
avatar_url: None,
system_prompt: "Test persona prompt.".to_string(),
runtime: None,
model: None,
provider: None,
name_pool: vec![],
is_builtin: false,
is_active: true,
shared: false,
source_team: None,
source_team_persona_slug: None,
catalog_source: None,
env_vars: Default::default(),
respond_to: None,
respond_to_allowlist: Default::default(),
parallelism: None,
created_at: crate::util::now_iso(),
updated_at: crate::util::now_iso(),
};
let test_team = crate::managed_agents::TeamRecord {
id: "test-team-id".to_string(),
name: "Test Team".to_string(),
description: None,
instructions: None,
persona_ids: vec![],
is_builtin: false,
source_dir: None,
is_symlink: false,
symlink_target: None,
version: None,
created_at: crate::util::now_iso(),
updated_at: crate::util::now_iso(),
};
let mut global_env_vars = std::collections::BTreeMap::new();
global_env_vars.insert("CAPTURED_GLOBAL_VAR".to_string(), "test-value".to_string());
let captured_global = crate::managed_agents::GlobalAgentConfig {
env_vars: global_env_vars,
..Default::default()
};
let context = CapturedRestartContext {
scope: scope.clone(),
personas: vec![],
teams: vec![],
global: crate::managed_agents::GlobalAgentConfig::default(),
personas: vec![test_persona],
teams: vec![test_team],
global: captured_global.clone(),
owner_hex: owner_hex.clone(),
mesh_model_id: None,
};
@@ -175,66 +218,77 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
let stop_called = Arc::new(Mutex::new(false));
let spawn_relay = Arc::new(Mutex::new(None::<String>));
let spawn_owner = Arc::new(Mutex::new(None::<String>));
let spawn_got_nonempty_personas = Arc::new(Mutex::new(false));
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));
let stop_called2 = stop_called.clone();
let spawn_relay2 = spawn_relay.clone();
let spawn_owner2 = spawn_owner.clone();
let spawn_personas2 = spawn_got_nonempty_personas.clone();
let spawn_teams2 = spawn_got_nonempty_teams.clone();
let spawn_global2 = spawn_got_nonempty_global.clone();
let receipt_called2 = receipt_called.clone();
let pubkey2 = pubkey.clone();
let result = restart_under_captured_epoch_for(
&app_handle,
&pubkey,
&old_global,
&new_global,
&[],
&context,
// stop_fn: record the call, simulate success, remove runtime.
move |_app, rec, runtimes| {
*stop_called2.lock().unwrap() = true;
runtimes.retain(|k, _| k.pubkey != rec.pubkey);
Ok(())
},
// spawn_fn: record captured relay+owner, return a fresh process.
move |_app, rec, relay, owner, _personas, global, teams| {
*spawn_relay2.lock().unwrap() = Some(relay.to_string());
*spawn_owner2.lock().unwrap() = owner.map(str::to_string);
// Return an immediately-exiting child as the fake process.
let child = std::process::Command::new("/usr/bin/true")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.map_err(|e| format!("spawn /usr/bin/true: {e}"))?;
Ok(crate::managed_agents::ManagedAgentProcess {
child,
log_path: std::path::PathBuf::new(),
spawn_config:
crate::managed_agents::spawn_snapshot::prospective_spawn_config_snapshot(
rec,
&[],
teams,
relay,
global,
),
setup_mode: false,
adapter_availability: None,
start_nonce: "test-nonce-spawn".to_string(),
#[cfg(windows)]
job: None,
})
},
// write_receipt_fn: record the call, verify pubkey, succeed.
move |_app, receipt| {
*receipt_called2.lock().unwrap() = true;
assert_eq!(
receipt.key.pubkey, pubkey2,
"receipt must carry the correct pubkey"
);
Ok(())
},
);
let app_handle_for_assert = app_handle.clone();
let result = tokio::task::spawn_blocking(move || {
restart_under_captured_epoch_for(
&app_handle,
&pubkey,
&old_global,
&new_global,
&[],
&context,
// stop_fn: record the call, simulate success, remove runtime.
move |_app, rec, runtimes| {
*stop_called2.lock().unwrap() = true;
runtimes.retain(|k, _| k.pubkey != rec.pubkey);
Ok(())
},
// spawn_fn: record captured relay+owner+personas+teams+global, return a noop child.
// `spawn_noop_child_for_test()` replaces `/usr/bin/true` — cross-platform.
move |_app, rec, relay, owner, personas, global, teams| {
*spawn_relay2.lock().unwrap() = Some(relay.to_string());
*spawn_owner2.lock().unwrap() = owner.map(str::to_string);
*spawn_personas2.lock().unwrap() = !personas.is_empty();
*spawn_teams2.lock().unwrap() = !teams.is_empty();
*spawn_global2.lock().unwrap() = !global.env_vars.is_empty();
Ok(crate::managed_agents::ManagedAgentProcess {
child: spawn_noop_child_for_test(),
log_path: std::path::PathBuf::new(),
spawn_config:
crate::managed_agents::spawn_snapshot::prospective_spawn_config_snapshot(
rec,
&[],
teams,
relay,
global,
),
setup_mode: false,
adapter_availability: None,
start_nonce: "test-nonce-spawn".to_string(),
#[cfg(windows)]
job: None,
})
},
// write_receipt_fn: record the call, verify pubkey, succeed.
move |_app, receipt| {
*receipt_called2.lock().unwrap() = true;
assert_eq!(
receipt.key.pubkey, pubkey2,
"receipt must carry the correct pubkey"
);
Ok(())
},
)
})
.await
.expect("spawn_blocking must not panic");
// Kill the seeded long-lived child now that the epoch has consumed it.
let _ = crate::managed_agents::terminate_process(seeded_pid);
assert!(
matches!(result, Ok(())),
@@ -251,13 +305,25 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
Some(owner_hex.as_str()),
"spawn_fn must receive the captured owner hex"
);
assert!(
*spawn_got_nonempty_personas.lock().unwrap(),
"spawn_fn must receive NON-EMPTY captured personas"
);
assert!(
*spawn_got_nonempty_teams.lock().unwrap(),
"spawn_fn must receive NON-EMPTY captured teams"
);
assert!(
*spawn_got_nonempty_global.lock().unwrap(),
"spawn_fn must receive a NON-EMPTY captured global env"
);
assert!(
*receipt_called.lock().unwrap(),
"write_receipt_fn must be called"
);
// Verify the runtime is registered with the captured scope_id.
{
let state = app_handle.state::<crate::app_state::AppState>();
let state = app_handle_for_assert.state::<crate::app_state::AppState>();
let runtimes = state.managed_agent_processes.lock().unwrap();
let registered = runtimes
.values()
@@ -267,26 +333,45 @@ 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.
let final_records =
crate::managed_agents::storage::load_managed_agents_at(tmp.path()).unwrap_or_default();
assert!(
final_records.iter().any(|r| r.pubkey == "bb".repeat(32)),
"final disk record must contain the restarted agent"
);
}
/// Production driver proves preflight fires before stop.
/// Production driver proves preflight fires before stop — and with an eligible
/// runtime seeded both "mesh" and "stop" appear in the log in that order.
///
/// Thufir's test 6: "event log from production driver proves preflight fn
/// fires before stop fn."
/// Thufir's test 6: "seed an eligible runtime through the cross-platform child
/// seam; event log from production driver proves preflight fn fires before stop
/// fn — unconditionally (both events must be present)."
///
/// The async driver's pre-stop phase requires: personas load, teams load,
/// global config load, owner key verification, candidate record load, and
/// mesh preflight (in that order). The mesh_fn is the first outbound async
/// call. We prove mesh fires before stop by asserting "mesh" appears first
/// in the event log.
/// Requirements for the agent to be an eligible restart candidate:
/// - `backend = Local` with a live pair runtime (seeded with long-lived child).
/// - `provider = anthropic`, `model = ...`, `ANTHROPIC_API_KEY` in env_vars
/// → `old_ready = true`.
/// - `old_global != new_global` (env_vars differ) → `env_changed = true`
/// → `should_restart_on_config_change = true`.
/// - `owner_pubkey` matches the mock app's signing key.
///
/// The `stop_fn` records "stop" and REMOVES the runtime from the map so the
/// epoch considers it properly stopped. The `spawn_fn` returns `Err` so the
/// epoch ends with `FailedAfterStop` — but both "mesh" and "stop" are in the
/// log before that.
///
/// Owner key must match the app's signing key — use the actual generated key
/// from the mock app's AppState. The agent pubkey in the store is separate.
/// from the mock app's AppState.
#[tokio::test]
async fn test_relay_mesh_preflight_precedes_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;
use crate::managed_agents::{
BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord, ManagedAgentRuntimeKey,
};
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
@@ -307,9 +392,14 @@ async fn test_relay_mesh_preflight_precedes_stop() {
let agent_pubkey = "aa".repeat(32);
// Write the candidate record so the pre-stop phase can find it.
// The agent must exist in the store for the candidate record lookup.
let agent_record = crate::managed_agents::ManagedAgentRecord {
// Build a Ready record: provider+model set and ANTHROPIC_API_KEY in env_vars.
// old_global != new_global (env differ) → env_changed = true → eligible.
let mut record_env_vars = std::collections::BTreeMap::new();
record_env_vars.insert(
"ANTHROPIC_API_KEY".to_string(),
"sk-test-key-eligible".to_string(),
);
let agent_record = ManagedAgentRecord {
pubkey: agent_pubkey.clone(),
name: "test-agent-preflight".to_string(),
display_name: None,
@@ -329,14 +419,14 @@ async fn test_relay_mesh_preflight_precedes_stop() {
max_turn_duration_seconds: None,
parallelism: 1,
system_prompt: None,
model: None,
provider: None,
model: Some("claude-3-5-sonnet-20241022".to_string()),
provider: Some("anthropic".to_string()),
persona_source_version: None,
env_vars: Default::default(),
env_vars: record_env_vars,
start_on_app_launch: false,
auto_restart_on_config_change: false,
runtime_pid: None,
backend: crate::managed_agents::BackendKind::Local,
backend: BackendKind::Local,
backend_agent_id: None,
provider_binary_path: None,
team_id: None,
@@ -364,10 +454,41 @@ async fn test_relay_mesh_preflight_precedes_stop() {
runtime: None,
name_pool: vec![],
};
save_managed_agents_at(tmp.path(), &[agent_record]).unwrap();
save_managed_agents_at(tmp.path(), std::slice::from_ref(&agent_record)).unwrap();
std::fs::write(tmp.path().join("personas.json"), b"[]").unwrap();
std::fs::write(tmp.path().join("global-agent-config.json"), b"{}").unwrap();
// Seed a live runtime for the agent (makes it eligible — avoids the
// "no live pair runtime" Skipped path).
let rt_key = ManagedAgentRuntimeKey::new(&agent_pubkey, "wss://relay.example").unwrap();
let seeded_pid = {
let state = app_handle.state::<crate::app_state::AppState>();
let mut runtimes = state.managed_agent_processes.lock().unwrap();
let child = spawn_long_lived_child_for_test();
let pid = child.id();
let process = crate::managed_agents::ManagedAgentProcess {
child,
log_path: std::path::PathBuf::new(),
spawn_config: crate::managed_agents::spawn_snapshot::prospective_spawn_config_snapshot(
&agent_record,
&[],
&[],
"wss://relay.example",
&Default::default(),
),
setup_mode: false,
adapter_availability: None,
start_nonce: "test-nonce-preflight".to_string(),
#[cfg(windows)]
job: None,
};
runtimes.insert(
rt_key,
ManagedAgentPairRuntime::starting(process, Some("test-scope".to_string())),
);
pid
};
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
@@ -378,6 +499,15 @@ async fn test_relay_mesh_preflight_precedes_stop() {
generation: gen,
};
// old_global and new_global differ by one env_var so env_changed = true.
let old_global = crate::managed_agents::GlobalAgentConfig::default();
let mut new_global_env = std::collections::BTreeMap::new();
new_global_env.insert("PREFLIGHT_TEST_VAR".to_string(), "v2".to_string());
let new_global = crate::managed_agents::GlobalAgentConfig {
env_vars: new_global_env,
..Default::default()
};
// Shared event log: "mesh" or "stop" entries in order.
let event_log: std::sync::Arc<std::sync::Mutex<Vec<&'static str>>> =
std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
@@ -387,8 +517,8 @@ async fn test_relay_mesh_preflight_precedes_stop() {
let outcome = restart_local_agent_on_config_change_for(
&app_handle,
&agent_pubkey,
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&old_global,
&new_global,
&[],
&scope,
tmp.path(),
@@ -397,39 +527,47 @@ async fn test_relay_mesh_preflight_precedes_stop() {
log_mesh.lock().unwrap().push("mesh");
Box::pin(async { Ok(()) })
},
// stop_fn: records "stop". Since there's no live runtime the epoch
// will Skipped before stop_fn fires — the ordering assertion is on
// the driver phase (mesh before ANY stop attempt).
move |_app, _rec, _runtimes| {
// stop_fn: records "stop", removes the runtime so spawn fails gracefully.
move |_app, rec, runtimes| {
log_stop.lock().unwrap().push("stop");
runtimes.retain(|k, _| k.pubkey != rec.pubkey);
Ok(())
},
// spawn_fn: returns Err so the epoch ends with FailedAfterStop.
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
Err("spawn not expected".to_string())
Err("spawn not available in test".to_string())
},
|_app, _receipt| Err("receipt not expected".to_string()),
)
.await;
// The call Skips at the no-live-runtime check, not at preflight.
// Kill the seeded process now that the epoch has consumed it.
let _ = crate::managed_agents::terminate_process(seeded_pid);
// With an eligible runtime, the epoch progresses past preflight and stop.
// stop_fn removed the runtime → spawn_fn fails → FailedAfterStop.
assert!(
matches!(outcome, RestartOutcome::Skipped),
"no live runtime: {outcome:?}"
matches!(outcome, RestartOutcome::FailedAfterStop),
"with eligible runtime: stop fires and spawn fails → FailedAfterStop: {outcome:?}"
);
let log = event_log.lock().unwrap();
// Mesh must appear before any stop (even if stop never fired).
// Both events must be present — mesh fires in the async pre-stop phase,
// stop fires inside the epoch.
let mesh_pos = log.iter().position(|&e| e == "mesh");
let stop_pos = log.iter().position(|&e| e == "stop");
assert!(
mesh_pos.is_some(),
"mesh_fn must be called (preflight runs before epoch)"
"mesh_fn must be called (preflight runs in async pre-stop phase): {log:?}"
);
assert!(
stop_pos.is_some(),
"stop_fn must be called with an eligible live runtime: {log:?}"
);
let mesh_idx = mesh_pos.unwrap();
let stop_idx = stop_pos.unwrap();
assert!(
mesh_idx < stop_idx,
"preflight (mesh at {mesh_idx}) must precede stop (stop at {stop_idx}): {log:?}"
);
if let Some(stop_idx) = stop_pos {
let mesh_idx = mesh_pos.unwrap();
assert!(
mesh_idx < stop_idx,
"preflight (mesh at {mesh_idx}) must precede stop (stop at {stop_idx})"
);
}
}
@@ -229,6 +229,71 @@ fn test_restart_under_captured_epoch_stale_scope_is_rejected() {
// ── Area-2 tests: async driver and epoch core ─────────────────────────────
/// Spawn a child process that exits immediately.
///
/// Cross-platform replacement for `/usr/bin/true`: used in `spawn_fn` closures
/// that must return a valid `ManagedAgentProcess` without spawning a real agent.
/// The child exits before or shortly after being inserted into the runtimes map;
/// for tests that only inspect in-memory state (not process liveness) this is
/// sufficient.
pub(crate) fn spawn_noop_child_for_test() -> std::process::Child {
#[cfg(not(windows))]
{
std::process::Command::new("sh")
.args(["-c", "exit 0"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn noop test child (sh -c 'exit 0')")
}
#[cfg(windows)]
{
std::process::Command::new("cmd.exe")
.args(["/C", "exit 0"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn noop test child (cmd.exe /C exit 0)")
}
}
/// Spawn a long-lived child process that stays running long enough for tests to
/// complete.
///
/// Cross-platform replacement for `sleep 10000`: seeds the in-memory runtimes
/// map before `sync_managed_agent_processes` runs its `try_wait()` scan, so
/// the runtime survives to the eligibility check.
///
/// Tests that use this helper MUST drop the returned `Child` (or kill it) when
/// the test exits so OS processes are not leaked. The helper is intentionally
/// not `#[cfg(test)]` — it lives here so `epoch_tests.rs` (via `use super::*`)
/// can reach it without a separate import.
pub(crate) fn spawn_long_lived_child_for_test() -> std::process::Child {
#[cfg(not(windows))]
{
std::process::Command::new("sh")
.args(["-c", "while true; do sleep 1; done"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn long-lived test child (sh loop)")
}
#[cfg(windows)]
{
// `ping -n N 127.0.0.1` sleeps ~(N-1) seconds; 100000 ≈ 28 hours.
std::process::Command::new("ping")
.args(["-n", "100000", "127.0.0.1"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn long-lived test child (ping)")
}
}
fn make_mock_app() -> tauri::App<tauri::test::MockRuntime> {
tauri::test::mock_builder()
.manage(crate::app_state::build_app_state())
@@ -239,11 +304,14 @@ fn make_mock_app() -> tauri::App<tauri::test::MockRuntime> {
/// Context load failure (personas) before any stop → `RestartOutcome::Skipped`,
/// stop closure never called.
///
/// Drives `restart_local_agent_on_config_change_for` with an injected mesh_fn
/// that succeeds but a definitions_dir that does not contain a managed-agents.json,
/// causing `load_managed_agents_at` in the pre-stop phase to fail (agent not found).
/// Specifically we provide a definitions_dir with NO managed-agents.json file so
/// the pre-stop load (for preflight record resolution) fails before any stop.
/// Drives `restart_local_agent_on_config_change_for` with a `definitions_dir`
/// containing a syntactically invalid `managed-agents.json` — `load_personas_at`
/// delegates to `load_agent_definitions_at`, which calls `load_agent_store_at`,
/// which returns `Err` on malformed JSON. The pre-stop phase must return
/// `Skipped` without calling stop.
///
/// Absent files load as empty/default, so this test writes a malformed file to
/// ensure a genuine parse-error path (not the "agent not found" path).
///
/// Thufir's test 1: "production async driver with injected loader failure;
/// assert stop never called, RestartOutcome::Skipped."
@@ -253,8 +321,13 @@ async fn test_context_load_failure_leaves_runtime_running() {
use crate::commands::global_agent_config::RestartOutcome;
let tmp = tempfile::tempdir().unwrap();
// No managed-agents.json → personas load fails (no personas.json also fine).
// Either way the pre-stop phase fails and returns Skipped.
// Malformed JSON → load_agent_store_at (called by load_personas_at) returns
// Err → pre-stop phase returns Skipped before any stop.
std::fs::write(
tmp.path().join("managed-agents.json"),
b"this is not valid json",
)
.unwrap();
let app = make_mock_app();
let app_handle = app.handle().clone();
@@ -437,22 +510,34 @@ 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, stops before stop.
/// generation) → epoch detects mismatch in re-resolved Mesh model, aborts before stop.
///
/// Thufir's test 4: "hook edits the captured record's Mesh-relevant config field
/// without generation change; epoch detects mismatch; stop never called."
/// 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."
///
/// Implementation note: the re-resolve check in the epoch compares
/// `re_resolved_mesh != context.mesh_model_id`. We set `context.mesh_model_id =
/// Some("model-a")` but the on-disk record has `relay_mesh: None`, so the
/// in-epoch re-resolve yields `None` ≠ `Some("model-a")` → 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`.
///
/// To reach the mesh re-resolve check the agent must be Ready (eligibility gate
/// passes) and have a live runtime. The record's `env_vars` supplies
/// BUZZ_AGENT_PROVIDER + BUZZ_AGENT_MODEL so it is Ready in both old and new
/// configs; a differing global env_var makes `env_changed = true`.
#[test]
fn test_record_mesh_change_after_preflight_aborts_before_stop() {
/// 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`).
#[tokio::test]
async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
use crate::managed_agents::{
storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord,
ManagedAgentRuntimeKey,
@@ -531,22 +616,19 @@ fn test_record_mesh_change_after_preflight_aborts_before_stop() {
name_pool: vec![],
};
save_managed_agents_at(tmp.path(), std::slice::from_ref(&record)).unwrap();
std::fs::write(tmp.path().join("personas.json"), b"[]").unwrap();
std::fs::write(tmp.path().join("global-agent-config.json"), b"{}").unwrap();
let app = make_mock_app();
let app_handle = app.handle().clone();
// Seed a live runtime with a long-running process (avoids sync eviction).
// 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 = {
let state = app_handle.state::<crate::app_state::AppState>();
let mut runtimes = state.managed_agent_processes.lock().unwrap();
let child = std::process::Command::new("sleep")
.arg("10000")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn sleep 10000");
let child = spawn_long_lived_child_for_test();
let pid = child.id();
let process = crate::managed_agents::ManagedAgentProcess {
child,
log_path: std::path::PathBuf::new(),
@@ -567,7 +649,8 @@ fn test_record_mesh_change_after_preflight_aborts_before_stop() {
rt_key,
ManagedAgentPairRuntime::starting(process, Some("test-scope".to_string())),
);
}
pid
};
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
@@ -604,23 +687,29 @@ fn test_record_mesh_change_after_preflight_aborts_before_stop() {
let stop_called2 = stop_called.clone();
// Call the epoch core directly — mesh_fn is not involved here since
// we're testing the in-epoch TOCTOU check.
let result = 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()),
);
// 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");
assert!(
matches!(result, Err(EpochError::Skipped(_))),
@@ -636,6 +725,8 @@ fn test_record_mesh_change_after_preflight_aborts_before_stop() {
!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"]
+222 -195
View File
@@ -333,19 +333,9 @@ pub async fn save_ncryptsec_copy(
Ok(Some(dest.display().to_string()))
}
/// Drain all live managed-agent runtimes as part of an identity import.
///
/// Uses the same journaled drain+compensation protocol as `apply_workspace`
/// (Layer 2 of the transition state machine). The caller MUST hold
/// `managed_agent_runtime_transition` before calling this function.
///
/// Returns:
/// - `Ok(stopped)` when all runtimes were stopped (or there were none).
/// - `Err((stopped, msg))` when the drain failed; `stopped` contains the
/// entries that were successfully killed before the failure — the caller
/// MUST pass the `managed_agent_runtime_transition` guard by value into
/// `compensate_drain`, which holds it continuously through all journal
/// restarts (no drop-and-reacquire interleave window).
/// Drain live managed-agent runtimes for identity import (Layer 2 protocol).
/// Caller must hold `managed_agent_runtime_transition`. Returns stopped entries
/// or `Err((stopped, msg))` on failure.
fn drain_managed_agent_runtimes_for_import(
app: &tauri::AppHandle,
state: &AppState,
@@ -367,206 +357,243 @@ pub async fn import_identity(
password: Option<String>,
app_handle: tauri::AppHandle,
) -> Result<IdentityInfo, String> {
// ── Layer 1: async serialization lock ────────────────────────────────────
// identity_mutation (Layer 1) must be held for the full import to prevent a
// concurrent stale persist from overwriting the imported key.
// ── Layer 1: identity_mutation (async serialization lock) ────────────────
// Held for the full import to prevent a concurrent stale persist from
// overwriting the imported key. Lock order: identity_mutation →
// workspace_transition (when active scope present).
//
// If there is an active scope, we also take workspace_transition so that
// clearing the active scope is serialized against apply_workspace.
// Lock order: identity_mutation → workspace_transition.
let lock_app = app_handle.clone();
let lock_state = lock_app.state::<AppState>();
// Use a cloned handle for lock acquisition so the original `app_handle` is
// free for the spawned blocking body below (no borrow conflict).
let lock_handle = app_handle.clone();
let lock_state = lock_handle.state::<AppState>();
let _mutation_guard = lock_state.identity_mutation.lock().await;
// Capture whether an active scope exists BEFORE entering spawn_blocking.
// Capture whether an active scope exists BEFORE branching.
let has_active_scope = lock_state.capture_active_scope().is_some();
// For the live-active path, also hold workspace_transition so that
// clearing the active scope is serialized against concurrent apply_workspace
// calls. Lock order: identity_mutation → workspace_transition.
let _transition_guard = if has_active_scope {
Some(lock_state.workspace_transition.lock().await)
// ── Layer 1b: workspace_transition + mesh preflight (active-scope path) ──
// When a workspace scope is live, route through the production
// `with_workspace_transition_preflight` helper (when the `mesh-llm`
// feature is enabled): it acquires `workspace_transition`, runs
// `fail_if_client_mesh_active`, then invokes the body while the lock
// remains held — the same orchestration path used by `apply_workspace`.
//
// Without `mesh-llm`, manually acquire `workspace_transition` (no mesh
// check needed) and invoke the blocking body.
//
// When no scope is active, skip the lock — there is no workspace to
// serialize against.
let result = if has_active_scope {
let app_for_preflight_body = app_handle.clone();
let nsec_for_body = nsec;
let password_for_body = password;
#[cfg(feature = "mesh-llm")]
let branch_result =
crate::commands::mesh_llm::scope_impl::with_workspace_transition_preflight(
&app_handle,
move || {
Box::pin(async move {
tokio::task::spawn_blocking(move || {
import_identity_blocking(
app_for_preflight_body,
nsec_for_body,
password_for_body,
true,
)
})
.await
.map_err(|e| format!("spawn_blocking failed: {e}"))?
})
},
)
.await;
#[cfg(not(feature = "mesh-llm"))]
let branch_result = {
let _transition_guard = lock_state.workspace_transition.lock().await;
tokio::task::spawn_blocking(move || {
import_identity_blocking(
app_for_preflight_body,
nsec_for_body,
password_for_body,
true,
)
})
.await
.map_err(|e| format!("spawn_blocking failed: {e}"))?
};
branch_result
} else {
tokio::task::spawn_blocking(move || {
import_identity_blocking(app_handle, nsec, password, false)
})
.await
.map_err(|e| format!("spawn_blocking failed: {e}"))?
};
// identity_mutation must outlive spawn_blocking — drop explicitly here so
// the compiler can see the guard's lifetime covers both branches.
drop(_mutation_guard);
result
}
/// Blocking body of [`import_identity`]: key recovery, journaled drain (when
/// `has_active_scope`), identity commit, and scope clear. The caller has
/// already acquired `workspace_transition` when `has_active_scope` is true.
fn import_identity_blocking(
app_handle: tauri::AppHandle,
nsec: String,
password: Option<String>,
has_active_scope: bool,
) -> Result<IdentityInfo, String> {
// NIP-49 backups require a passphrase and decrypt entirely in Rust.
// Raw nsec/hex input follows the existing parser path unchanged.
let password = password.map(zeroize::Zeroizing::new);
let keys = crate::key_backup::recover_keys_from_input(
&nsec,
password.as_ref().map(|value| value.as_str()),
)?;
let state = app_handle.state::<AppState>();
let data_dir = app_handle
.path()
.app_data_dir()
.map_err(|e| format!("app data dir: {e}"))?;
std::fs::create_dir_all(&data_dir).map_err(|e| format!("create app data dir: {e}"))?;
let key_path = data_dir.join("identity.key");
// ── Live-active path: journaled drain before swapping identity ─────────
// Drain all managed-agent runtimes under `managed_agent_runtime_transition`
// (Layer 2) BEFORE persisting the new identity — same protocol as
// `apply_workspace`. The store lock is held through drain/save; on drain
// failure the transition guard is passed into compensate_drain so
// compensation runs without any interleave window.
let _rt_transition_guard = if has_active_scope {
Some(
state
.managed_agent_runtime_transition
.lock()
.map_err(|e| format!("managed_agent_runtime_transition poisoned: {e}"))?,
)
} else {
None
};
// Fail closed if a client-mode Mesh runtime is active when there is an
// active scope — a live client-mode runtime would become dangling after the
// import (Option A ruling). Called via `run_mesh_transition_preflight` so the
// shared preflight logic is not duplicated inline; the guard above is held.
#[cfg(feature = "mesh-llm")]
if has_active_scope {
crate::commands::mesh_llm::scope_impl::run_mesh_transition_preflight(&app_handle).await?;
}
let _store_guard = if has_active_scope {
Some(
state
.managed_agents_store_lock
.lock()
.map_err(|e| format!("managed_agents_store_lock poisoned: {e}"))?,
)
} else {
None
};
let result = tokio::task::spawn_blocking(move || {
// NIP-49 backups require a passphrase and decrypt entirely in Rust.
// Raw nsec/hex input follows the existing parser path unchanged.
let password = password.map(zeroize::Zeroizing::new);
let keys = crate::key_backup::recover_keys_from_input(
&nsec,
password.as_ref().map(|value| value.as_str()),
)?;
// Capture the pre-import scope for compensation validation.
let pre_import_scope = state.capture_active_scope();
let state = app_handle.state::<AppState>();
let data_dir = app_handle
.path()
.app_data_dir()
.map_err(|e| format!("app data dir: {e}"))?;
std::fs::create_dir_all(&data_dir).map_err(|e| format!("create app data dir: {e}"))?;
let key_path = data_dir.join("identity.key");
// ── Live-active path: journaled drain before swapping identity ─────────
// When an active scope exists, drain all managed-agent runtimes under
// `managed_agent_runtime_transition` (Layer 2) BEFORE persisting the
// new identity. This is the same drain-journal protocol used by
// `apply_workspace`.
//
// `managed_agents_store_lock` is acquired at Layer 2 start (same as
// apply_workspace) so a concurrent save_managed_agents cannot interleave
// with the drain or the scope clear.
//
// On drain failure: compensate by restarting what was stopped, then
// return Err — do NOT proceed with the identity persist. The caller
// (frontend membership-denied flow) must handle the Err and retry.
//
// The store lock must be DROPPED before calling compensate_drain, but
// the transition lock is passed INTO compensate_drain so compensation
// runs without any interleave window.
let _rt_transition_guard = if has_active_scope {
Some(
state
.managed_agent_runtime_transition
.lock()
.map_err(|e| format!("managed_agent_runtime_transition poisoned: {e}"))?,
)
} else {
None
};
let _store_guard = if has_active_scope {
Some(
state
.managed_agents_store_lock
.lock()
.map_err(|e| format!("managed_agents_store_lock poisoned: {e}"))?,
)
} else {
None
};
// Capture the pre-import scope for compensation validation.
let pre_import_scope = state.capture_active_scope();
let stopped_entries = if has_active_scope {
match drain_managed_agent_runtimes_for_import(&app_handle, &state) {
Ok(stopped) => stopped,
Err((stopped, drain_err)) => {
// Drain failed — drop the store lock BEFORE compensating
// (compensate_drain re-acquires it), but pass the transition
// guard into compensate_drain so there is no interleave window.
drop(_store_guard);
let comp_err = match (pre_import_scope.as_ref(), _rt_transition_guard) {
(Some(scope), Some(rt_guard)) => crate::managed_agents::compensate_drain(
&app_handle,
&stopped,
scope,
rt_guard,
),
(_, leftover_guard) => {
drop(leftover_guard);
None
}
};
let msg = match comp_err {
Some(comp) => format!(
"identity import drain failed: {drain_err}; compensation failed: {comp}"
),
None => format!("identity import drain failed: {drain_err}"),
};
return Err(msg);
}
}
} else {
vec![]
};
let commit_result = commit_imported_identity(&state, &data_dir, keys, |keys| {
// Persist into the OS keyring first (store → read-back verify →
// marker → delete file). Falls back to the 0o600 file when the
// keyring is unavailable; returns Err only when both backends fail.
let store =
crate::secret_store::SecretStore::shared(crate::app_state::keyring_service());
crate::app_state::persist_imported_identity(store, keys, &key_path, &data_dir)
});
// If identity persist failed after a successful drain, compensate.
// Drop the store lock BEFORE calling compensate_drain; pass the
// transition guard into it so there is no interleave window.
let (pubkey, storage) = match commit_result {
Ok(result) => result,
Err(e) => {
if !stopped_entries.is_empty() {
drop(_store_guard);
let comp_err = match (pre_import_scope.as_ref(), _rt_transition_guard) {
(Some(scope), Some(rt_guard)) => crate::managed_agents::compensate_drain(
&app_handle,
&stopped_entries,
scope,
rt_guard,
),
(_, leftover_guard) => {
drop(leftover_guard);
None
}
};
if let Some(comp_err) = comp_err {
eprintln!(
"buzz-desktop: identity import persist failed, compensation failed: {comp_err}"
);
let stopped_entries = if has_active_scope {
match drain_managed_agent_runtimes_for_import(&app_handle, &state) {
Ok(stopped) => stopped,
Err((stopped, drain_err)) => {
// Drain failed — drop the store lock BEFORE compensating
// (compensate_drain re-acquires it), but pass the transition
// guard into compensate_drain so there is no interleave window.
drop(_store_guard);
let comp_err = match (pre_import_scope.as_ref(), _rt_transition_guard) {
(Some(scope), Some(rt_guard)) => crate::managed_agents::compensate_drain(
&app_handle,
&stopped,
scope,
rt_guard,
),
(_, leftover_guard) => {
drop(leftover_guard);
None
}
}
return Err(e);
};
let msg = match comp_err {
Some(comp) => format!(
"identity import drain failed: {drain_err}; compensation failed: {comp}"
),
None => format!("identity import drain failed: {drain_err}"),
};
return Err(msg);
}
};
}
} else {
vec![]
};
// ── Clear active scope and bump generation ────────────────────────────
// For no-active-scope path: no scope was ever set; clearing is a no-op
// but bumping generation invalidates any in-flight stale operations.
// For live-active path: agents are stopped; clearing scope makes all
// agent commands fail closed until the frontend re-applies a workspace.
//
// Invariant: the fallback relay can never claim legacy data — claims
// are only written inside apply_workspace's prepare stage.
//
// `clear_active_scope()` internally calls `next_scope_generation()` —
// no additional bump is needed here.
state.clear_active_scope();
let commit_result = commit_imported_identity(&state, &data_dir, keys, |keys| {
// Persist into the OS keyring first (store → read-back verify →
// marker → delete file). Falls back to the 0o600 file when the
// keyring is unavailable; returns Err only when both backends fail.
let store = crate::secret_store::SecretStore::shared(crate::app_state::keyring_service());
crate::app_state::persist_imported_identity(store, keys, &key_path, &data_dir)
});
let pubkey_hex = pubkey.to_hex();
let display_name = truncated_display_name(&pubkey)?;
// If identity persist failed after a successful drain, compensate.
// Drop the store lock BEFORE calling compensate_drain; pass the
// transition guard into it so there is no interleave window.
let (pubkey, storage) = match commit_result {
Ok(result) => result,
Err(e) => {
if !stopped_entries.is_empty() {
drop(_store_guard);
let comp_err = match (pre_import_scope.as_ref(), _rt_transition_guard) {
(Some(scope), Some(rt_guard)) => crate::managed_agents::compensate_drain(
&app_handle,
&stopped_entries,
scope,
rt_guard,
),
(_, leftover_guard) => {
drop(leftover_guard);
None
}
};
if let Some(comp_err) = comp_err {
eprintln!(
"buzz-desktop: identity import persist failed, compensation failed: {comp_err}"
);
}
}
return Err(e);
}
};
eprintln!("buzz-desktop: imported identity pubkey {}", pubkey_hex);
// ── Clear active scope and bump generation ────────────────────────────
// For no-active-scope path: no scope was ever set; clearing is a no-op
// but bumping generation invalidates any in-flight stale operations.
// For live-active path: agents are stopped; clearing scope makes all
// agent commands fail closed until the frontend re-applies a workspace.
//
// Invariant: the fallback relay can never claim legacy data — claims
// are only written inside apply_workspace's prepare stage.
//
// `clear_active_scope()` internally calls `next_scope_generation()` —
// no additional bump is needed here.
state.clear_active_scope();
Ok(IdentityInfo {
pubkey: pubkey_hex,
display_name,
storage: storage.as_str().to_string(),
lost: false,
locked: false,
reset_failed: false,
})
let pubkey_hex = pubkey.to_hex();
let display_name = truncated_display_name(&pubkey)?;
eprintln!("buzz-desktop: imported identity pubkey {}", pubkey_hex);
Ok(IdentityInfo {
pubkey: pubkey_hex,
display_name,
storage: storage.as_str().to_string(),
lost: false,
locked: false,
reset_failed: false,
})
.await
.map_err(|e| format!("spawn_blocking failed: {e}"))?;
// Guards must stay alive until spawn_blocking completes so the serialization
// covers the full duration of the import.
drop(_transition_guard);
drop(_mutation_guard);
result
}
/// Commit an imported identity: durably persist, swap in-memory keys, clear
@@ -85,22 +85,19 @@ async fn test_active_client_stop_then_transition_preflight_succeeds() {
);
}
/// Transition helper holds the lock, commits a distinct scope (advancing
/// generation); queued install detects the stale captured scope and aborts
/// without invoking the injected install closure.
/// `with_workspace_transition_preflight` holds the lock (body blocks on a
/// oneshot channel); `install_client_under_workspace_transition` queues on the
/// same lock. After the body completes (advancing the scope generation while
/// the lock is held), the install task acquires the lock and detects the stale
/// captured scope without invoking the install closure.
///
/// Proves the first serialization direction: a workspace switch commits between
/// scope-capture and lock-acquisition for the install path. The install helper
/// must reject the stale scope under the lock rather than invoke install.
/// Both sides use production helpers — no direct `workspace_transition` lock
/// acquisition. Channels establish entry/release ordering without `sleep`.
///
/// Concurrency model: `workspace_transition` is acquired directly in the test
/// body to simulate a long-running transition. A spawned install task queues
/// on the same lock. When the test body drops the guard the install task
/// acquires, validates, and finds the captured scope stale.
/// 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.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock serializes generation counter
// mutations — cannot deadlock. See import_tests.rs for full rationale.
#[allow(clippy::await_holding_lock)]
async fn test_transition_held_queued_install_detects_stale_scope() {
use crate::managed_agents::scope::{
next_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK,
@@ -108,8 +105,9 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use tauri::Manager;
use tokio::sync::oneshot;
// Serialise generation-sensitive work across parallel tests.
// Serialize generation-sensitive work across parallel tests.
let _gen_guard = SCOPE_GENERATION_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
@@ -123,38 +121,64 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
let base = std::path::PathBuf::from("/tmp/area4-test-dir");
let relay_a = "wss://transition-test-a.example";
let owner_a = "aa".repeat(32); // 64-char hex-looking pubkey
let owner_a = "aa".repeat(32);
// Capture a scope at the current generation — this is the scope the install
// task will carry after "discovering a bootstrap target".
// Capture a scope at the current generation — the install task will carry
// this scope after "discovering a bootstrap target".
let gen_a = next_scope_generation();
let captured_scope =
WorkspaceAgentScope::new(relay_a.to_string(), owner_a.clone(), &base, gen_a);
// Make it the active scope so the identity check inside the helper can compare.
state.commit_active_scope(captured_scope.clone());
// Acquire the workspace_transition lock directly — simulates the transition
// helper holding the lock during a workspace switch.
let transition_guard = state.workspace_transition.lock().await;
// 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::<()>();
let (body_release_tx, body_release_rx) = oneshot::channel::<()>();
// Advance generation and commit a DISTINCT scope (different relay) to
// simulate a committed workspace switch while the lock was held.
let gen_b = next_scope_generation();
let new_scope = WorkspaceAgentScope::new(
"wss://transition-test-b.example".to_string(),
"bb".repeat(32),
&base,
gen_b,
);
state.commit_active_scope(new_scope);
// ── 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.
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>();
let gen_b = next_scope_generation();
let new_scope = WorkspaceAgentScope::new(
"wss://transition-test-b.example".to_string(),
"bb".repeat(32),
&base_a,
gen_b,
);
state_a.commit_active_scope(new_scope);
// Spawn the install task. It blocks on workspace_transition until we drop the guard.
// 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
.await
.expect("transition body must signal entry");
// ── Side B: install side uses the production helper ──────────────────────
let install_was_called = Arc::new(AtomicBool::new(false));
let install_called_clone = Arc::clone(&install_was_called);
let app_handle_clone = app_handle.clone();
let app_b = app_handle.clone();
let install_task = tokio::task::spawn(async move {
super::scope_impl::install_client_under_workspace_transition(
&app_handle_clone,
&app_b,
&captured_scope,
|| {
let called = Arc::clone(&install_called_clone);
@@ -167,15 +191,17 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
.await
});
// Yield to give the install task a chance to queue on the lock before we
// release it. This is not required for correctness (the generation check
// fires regardless of ordering) but makes the concurrent queuing observable.
// Yield to let the install task queue on the lock.
tokio::task::yield_now().await;
// Release the transition lock → install task acquires it and validates.
drop(transition_guard);
// Unblock the transition body → it releases the lock → install acquires it.
let _ = body_release_tx.send(());
let install_result = install_task.await.expect("install task must not panic");
transition_task
.await
.expect("transition task must not panic")
.expect("transition body must succeed");
// The install helper must have rejected the stale scope without invoking install.
assert!(
@@ -194,28 +220,27 @@ async fn test_transition_held_queued_install_detects_stale_scope() {
);
}
/// Install helper acquires the lock first and installs a mock client; after
/// release, the transition preflight observes the client and rejects.
/// `install_client_under_workspace_transition` holds the lock (install closure
/// blocks on a oneshot channel after installing a mock client);
/// `with_workspace_transition_preflight` queues on the same lock. After the
/// 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`.
///
/// Proves the second serialization direction: a client install commits between
/// "preflight check" and lock-acquisition for the transition path. The
/// transition helper must observe the installed client under the lock and fail.
///
/// Concurrency model: `workspace_transition` is acquired directly in the test
/// body to simulate a long-running install. A spawned transition task queues
/// on the same lock. When the test body drops the guard the transition task
/// acquires, runs `fail_if_client_mesh_active`, and finds the installed client.
/// the transition's pre-lock preflight and its lock-acquisition; the transition
/// helper must observe the installed client under the lock and fail.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock serializes generation counter
// mutations — cannot deadlock. See import_tests.rs for full rationale.
#[allow(clippy::await_holding_lock)]
async fn test_install_held_transition_preflight_observes_client() {
use crate::managed_agents::scope::{
next_scope_generation, WorkspaceAgentScope, SCOPE_GENERATION_TEST_LOCK,
};
use tauri::Manager;
use tokio::sync::oneshot;
// Serialise generation-sensitive work across parallel tests.
// Serialize generation-sensitive work across parallel tests.
let _gen_guard = SCOPE_GENERATION_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
@@ -237,34 +262,62 @@ 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());
// Acquire workspace_transition directly — simulates the install helper
// holding the lock during target acquisition + client startup.
let install_guard = state.workspace_transition.lock().await;
// 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::<()>();
let (install_release_tx, install_release_rx) = oneshot::channel::<()>();
// Install the mock client runtime while we hold the lock.
// This is the "install closure ran while holding the lock" effect.
{
let client_runtime = crate::mesh_llm::build_mock_client_runtime_for_test();
*state.mesh_llm_runtime.lock().await = Some(client_runtime);
}
// ── 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.
let app_a = app_handle.clone();
let scope_a = scope.clone();
let install_task = tokio::task::spawn(async move {
let app_a_for_closure = app_a.clone();
super::scope_impl::install_client_under_workspace_transition(
&app_a,
&scope_a,
move || async move {
// Install the mock client runtime while holding the 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);
// Spawn the transition task. It blocks on workspace_transition until we drop
// the install guard.
let app_handle_clone = app_handle.clone();
// 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
.await
.expect("install body must signal entry");
// ── Side B: transition side uses the production helper ───────────────────
let app_b = app_handle.clone();
let transition_task = tokio::task::spawn(async move {
super::scope_impl::with_workspace_transition_preflight(&app_handle_clone, || {
super::scope_impl::with_workspace_transition_preflight(&app_b, || {
Box::pin(async { Ok::<&str, String>("body ran") })
})
.await
});
// Yield to give the transition task a chance to queue on the lock before we
// release it. Not required for correctness but makes the concurrent queuing
// observable.
// Yield to let the transition task queue on the lock.
tokio::task::yield_now().await;
// Release the install lock → transition task acquires it and runs preflight.
drop(install_guard);
// Unblock the install body → it releases the lock → transition acquires it.
let _ = install_release_tx.send(());
install_task
.await
.expect("install task must not panic")
.expect("install closure must succeed");
let transition_result = transition_task
.await
@@ -235,6 +235,12 @@ fn minimal_agent_snapshot_json_with_memory(name: &str) -> Vec<u8> {
/// Uses `tauri::test::mock_builder().manage(state)` so that `app.state::<AppState>()`
/// works inside `try_regenerate_nest` and other AppHandle users called by the core.
///
/// Uses `WorkspaceAgentScope::new` with `base_dir = tmp.path()` so that
/// `definitions_dir = <tmp>/scopes/<scope_id>/`. `retention_scope_from_captured`
/// derives the agent base two parents above `definitions_dir`, yielding
/// `<tmp>` — a writable directory — instead of `/` (which causes EPERM on Linux
/// when `definitions_dir` is set directly to the tempdir root).
///
/// Uses `next_scope_generation()` to claim the current generation slot so the
/// scope's generation matches the global counter at entry, reducing the race
/// window vs. tests that call `next_scope_generation()` concurrently.
@@ -263,13 +269,19 @@ fn setup_import_app_with_scope(
// value so capture_agent_snapshot_import_entry sees a matching generation
// when it reads current_scope_generation() at entry.
let gen = next_scope_generation();
let scope = WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: "wss://captured.example".to_string(),
owner_pubkey: owner_keys.public_key().to_hex(),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
// Use WorkspaceAgentScope::new so definitions_dir has the production
// shape: <tmp>/scopes/<scope_id>/. retention_scope_from_captured derives
// the agent base two parents above definitions_dir — with this layout it
// resolves to <tmp> (writable) rather than / (which causes EPERM on Linux).
let scope = WorkspaceAgentScope::new(
"wss://captured.example".to_string(),
owner_keys.public_key().to_hex(),
tmp.path(),
gen,
);
// Ensure the definitions directory exists so the core can write into it.
std::fs::create_dir_all(&scope.definitions_dir)
.expect("failed to create scope definitions dir");
s.commit_active_scope(scope);
}
@@ -11,6 +11,42 @@ use super::*;
/// See the equivalent comment in `import_tests.rs` for rationale.
use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK as GENERATION_TEST_LOCK;
/// Build a team-member snapshot with a core memory entry.
///
/// Used by tests that must exercise the Phase-5 memory loop in the team import
/// core and assert that `submit_memory` carries the captured relay + owner.
fn member_with_memory(name: &str) -> AgentSnapshot {
use crate::managed_agents::agent_snapshot::{AgentSnapshotMemoryEntry, MemoryLevel};
let mut m = member(name);
m.memory = crate::managed_agents::agent_snapshot::AgentSnapshotMemory {
level: MemoryLevel::Core,
entries: vec![AgentSnapshotMemoryEntry {
slug: buzz_core_pkg::engram::CORE_SLUG.to_string(),
body: format!("# {name}\nTeam member memory body."),
}],
};
m
}
/// Extract the first `p` tag value from a nostr event JSON byte slice.
///
/// Returns `Some(hex_pubkey)` if a `["p", "<hex>"]` tag entry is found,
/// `None` if the JSON cannot be parsed or has no `p` tag.
fn extract_p_tag_from_memory_event(event_json: &[u8]) -> Option<String> {
let val: serde_json::Value = serde_json::from_slice(event_json).ok()?;
let tags = val.get("tags")?.as_array()?;
for tag in tags {
if let Some(arr) = tag.as_array() {
if arr.first().and_then(|v| v.as_str()) == Some("p") {
if let Some(hex) = arr.get(1).and_then(|v| v.as_str()) {
return Some(hex.to_string());
}
}
}
}
None
}
fn setup_team_import_app_with_scope(
tmp: &tempfile::TempDir,
) -> (tauri::App<tauri::test::MockRuntime>, nostr::Keys) {
@@ -32,13 +68,19 @@ fn setup_team_import_app_with_scope(
use tauri::Manager;
let s = app.state::<crate::app_state::AppState>();
let gen = next_scope_generation();
let scope = WorkspaceAgentScope {
scope_id: "ts-test-scope".to_string(),
relay_url: "wss://captured.example".to_string(),
owner_pubkey: owner_keys.public_key().to_hex(),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
// Use WorkspaceAgentScope::new so definitions_dir has the production
// shape: <tmp>/scopes/<scope_id>/. retention_scope_from_captured derives
// the agent base two parents above definitions_dir — with this layout it
// resolves to <tmp> (writable) rather than / (which causes EPERM on Linux).
let scope = WorkspaceAgentScope::new(
"wss://captured.example".to_string(),
owner_keys.public_key().to_hex(),
tmp.path(),
gen,
);
// Ensure the definitions directory exists so the core can write into it.
std::fs::create_dir_all(&scope.definitions_dir)
.expect("failed to create team scope definitions dir");
s.commit_active_scope(scope);
}
@@ -53,6 +95,14 @@ fn setup_team_import_app_with_scope(
/// scope and owner, not merely increment a counter. We swap the active scope
/// to a different relay + fresh owner inside the hook; all per-member profile
/// adapters must receive the old captured relay URL.
///
/// One member carries a core memory entry so the Phase-5 memory loop in
/// `team_snapshot.rs` fires. The memory adapter asserts BOTH the captured
/// relay URL AND that the built engram event's `p` tag (owner counterpart)
/// matches the CAPTURED owner's pubkey — not the post-switch owner committed
/// in `after_store`. This validates the team core's independently implemented
/// Phase-5 memory loop (`team_snapshot.rs:815-860`) which uses
/// `captured_owner_keys` at `:553`.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock held to serialize generation
// counter mutations — cannot deadlock. See import_tests.rs for full rationale.
@@ -67,21 +117,28 @@ async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
let (app, owner_keys) = setup_team_import_app_with_scope(&tmp);
let handle = app.handle();
let snap = snapshot(vec![member("Alice"), member("Bob")]);
// One plain member + one member with a core memory entry so the Phase-5
// memory loop fires for the second member.
let snap = snapshot(vec![member("Alice"), member_with_memory("Bob")]);
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
let input = TeamSnapshotImportConfirm {
file_bytes: encoded,
keep_allowlist: false,
};
// Captured owner pubkey — must appear in the `p` tag of memory events.
let captured_owner_pubkey_hex = owner_keys.public_key().to_hex();
let expected_relay = "wss://captured.example".to_string();
let profile_relays: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(vec![]));
let memory_relays: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(vec![]));
let memory_p_tags: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(vec![]));
let pr = profile_relays.clone();
let mr = memory_relays.clone();
let mp = memory_p_tags.clone();
let state = app.state::<crate::app_state::AppState>();
let handle_for_hook = handle.clone();
@@ -94,13 +151,12 @@ async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
move || {
// after_store: commit a genuinely DIFFERENT live scope + owner.
let new_owner = nostr::Keys::generate();
let new_scope = WorkspaceAgentScope {
scope_id: "switched-scope".to_string(),
relay_url: "wss://new-relay-after-switch.example".to_string(),
owner_pubkey: new_owner.public_key().to_hex(),
definitions_dir: std::path::PathBuf::from("/tmp/switched"),
generation: current_scope_generation(),
};
let new_scope = WorkspaceAgentScope::new(
"wss://new-relay-after-switch.example".to_string(),
new_owner.public_key().to_hex(),
std::path::Path::new("/tmp/switched"),
current_scope_generation(),
);
let s = handle_for_hook.state::<crate::app_state::AppState>();
s.commit_active_scope(new_scope);
},
@@ -113,12 +169,15 @@ async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
})
},
move |m: MemoryPublish<'_>| {
// Assert: relay URL contains the captured relay, not the switched one.
let relay = m.relay_url.to_string();
mr.lock().unwrap().push(relay.clone());
Box::pin(async move {
let _ = relay;
Ok(())
})
// Extract the `p` tag — must carry the CAPTURED owner pubkey.
let event_bytes = m.event_json.to_vec();
if let Some(p_tag) = extract_p_tag_from_memory_event(&event_bytes) {
mp.lock().unwrap().push(p_tag);
}
Box::pin(async move { Ok(()) })
},
)
.await;
@@ -130,19 +189,47 @@ async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
);
// Every member's profile adapter received the captured relay URL.
let seen = profile_relays.lock().unwrap();
let seen_profiles = profile_relays.lock().unwrap();
assert_eq!(
seen.len(),
seen_profiles.len(),
2,
"profile adapter must be called once per member"
);
for relay in seen.iter() {
for relay in seen_profiles.iter() {
assert_eq!(
relay, &expected_relay,
"profile adapter must receive captured relay, got: {relay}"
);
}
// No memory entries in these members — memory adapter not called.
// Memory adapter was called for the member with memory entries.
let seen_memory_relays = memory_relays.lock().unwrap();
assert!(
!seen_memory_relays.is_empty(),
"memory adapter must be called for the member with memory entries"
);
for relay in seen_memory_relays.iter() {
assert!(
relay.contains("captured.example"),
"memory adapter relay_url must contain captured relay 'captured.example', got: {relay}"
);
}
// Memory event's `p` tag must equal the CAPTURED owner's pubkey — not the
// post-switch owner committed in `after_store`.
let seen_p_tags = memory_p_tags.lock().unwrap();
assert!(
!seen_p_tags.is_empty(),
"at least one memory event must carry a `p` tag"
);
for p_tag in seen_p_tags.iter() {
assert_eq!(
p_tag.as_str(),
captured_owner_pubkey_hex.as_str(),
"engram event p-tag must equal the captured owner's pubkey (not post-switch owner); \
got: {p_tag}"
);
}
}
/// `before_store` hook advances scope generation — Phase 3 must reject BEFORE
@@ -6,31 +6,31 @@
use super::*;
/// Deterministic writer-vs-compensation ordering via held-lock queuing.
/// Deterministic writer-vs-compensation ordering via the production
/// `compensate_drain` adapter.
///
/// Mechanism:
/// 1. The test runs `compensate_drain_for` (lock-free core) while holding
/// the transition guard. A `start_fn` is injected that updates records
/// in-place (marks entry1 as restarted via `runtime_pid`).
/// 2. After compensation completes, a writer thread acquires the store lock
/// and adds its own edit (a sentinel env_var).
/// 3. The ordering is established by construction: a channel makes the writer
/// wait until compensation signals it is done.
/// 4. Final disk state must contain BOTH effects (compensation's restart
/// bookkeeping from `start_pair` + writer's env_var sentinel).
/// 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.
///
/// Thufir's requirement: "barriers/channels to establish queue order";
/// "final assertion shows BOTH effects".
///
/// Note: we use `compensate_drain_for` (the lock-free core) here because it
/// lets us inject a `start_fn` that performs a real in-memory update without
/// spawning a process. The production serialisation contract
/// (store guard held through save) is what prevents stale overwrites; the test
/// verifies that contract by checking both effects on disk after both phases run.
/// Ordering is established by construction (channel), not by scheduler timing.
/// "both effects" proves the store guard is correctly held across restore/save.
#[test]
fn test_compensate_drain_writer_vs_compensation_deterministic() {
use std::sync::{Arc, Mutex};
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
use std::thread;
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf();
@@ -98,45 +98,46 @@ 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(),
relay_url: "wss://relay.example".to_string(),
owner_pubkey: "aa".repeat(32),
definitions_dir: tmp_path.clone(),
generation: gen,
};
state.commit_active_scope(scope.clone());
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1.clone()];
let stopped = vec![entry1];
// Phase-A (compensation): run compensate_drain_for with a start_fn that
// marks the record as restarted by setting `runtime_pid = Some(99)`, then
// saves the records slice back to disk.
let comp_ran = Arc::new(Mutex::new(false));
let comp_ran2 = comp_ran.clone();
let tmp_comp = tmp_path.clone();
// Channel: writer waits on this before acquiring the store.
// Channel: writer waits until compensate_drain signals "done".
let (comp_done_tx, comp_done_rx) = std::sync::mpsc::channel::<()>();
// Run compensation in this thread (simulating the adapter's held-lock epoch).
let comp_result = {
let mut records =
crate::managed_agents::storage::load_managed_agents_at(&tmp_comp).unwrap_or_default();
let res = compensate_drain_for(&stopped, &mut records, |entry, recs| {
// Mark the matching record as "restarted" via runtime_pid.
if let Some(r) = recs.iter_mut().find(|r| r.pubkey == entry.key.pubkey) {
r.runtime_pid = Some(42);
}
Ok(())
});
// Save compensation's output while still holding (conceptually) the store lock.
crate::managed_agents::storage::save_managed_agents_at(&tmp_comp, &records).unwrap();
*comp_ran2.lock().unwrap() = true;
// Signal the writer that compensation is done.
comp_done_tx.send(()).unwrap();
res
};
// Phase-B (writer): waits until compensation signals done, then loads
// whatever is on disk, adds its sentinel, saves.
// 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.
let tmp_wr = tmp_path.clone();
let app_handle_wr = app_handle.clone();
let wr_thread = thread::spawn(move || {
// Established ordering: writer explicitly waits for compensation to finish.
// 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();
// Now acquire just the store lock and apply the sentinel edit.
let writer_state = app_handle_wr.state::<crate::app_state::AppState>();
let _store = writer_state.managed_agents_store_lock.lock().unwrap();
let mut records =
crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default();
for r in &mut records {
@@ -146,30 +147,39 @@ 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);
// Signal the writer that compensation is done (guard released).
comp_done_tx.send(()).unwrap();
wr_thread.join().expect("writer thread panicked");
// Compensation succeeded (entry1 restored → None).
assert!(
comp_result.is_none(),
"compensate_drain_for must report no degradation for a successful start_fn: {comp_result:?}"
);
assert!(*comp_ran.lock().unwrap(), "compensation must have run");
// ── 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.
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 == pubkey1)
.expect("agent record must be present on disk");
.expect("agent record must be present on disk after both phases");
// Compensation's effect: runtime_pid set to 42.
assert_eq!(
final_rec.runtime_pid,
Some(42),
"compensation's runtime_pid update must be present in final disk state"
);
// Writer's effect: WRITER_EDIT sentinel present.
// Writer's effect must always be present.
assert_eq!(
final_rec.env_vars.get("WRITER_EDIT").map(String::as_str),
Some("yes"),
@@ -177,24 +187,29 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
);
}
/// Joined drain→restore test with a real start-contender blocked on the
/// transition mutex.
/// Production start-path contender is blocked while `compensate_drain` holds
/// the `managed_agent_runtime_transition` mutex.
///
/// Flow:
/// 1. A contender thread starts and waits on the transition mutex.
/// 2. The test thread acquires the transition guard.
/// 3. A channel confirms the contender is alive and parked on the mutex.
/// 4. `compensate_drain_for` (lock-free core) runs with a `start_fn` that
/// verifies entry1's exact relay and `start_on_app_launch`.
/// 5. The test drops the transition guard.
/// 6. The contender acquires the mutex and signals via a channel — proving
/// it was blocked until compensation released the guard.
/// 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."
///
/// Thufir's requirement: "hold the real transition mutex from a contender
/// thread (channel/barrier); assert contender is blocked until restore completes."
/// 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.
#[test]
fn test_compensate_drain_concurrent_start_is_blocked() {
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
use std::thread;
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf();
@@ -206,12 +221,11 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
.expect("failed to build mock app");
let app_handle = app.app_handle().clone();
use tauri::Manager;
let state = app.state::<crate::app_state::AppState>();
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope-contender".to_string(),
let gen = current_scope_generation();
let scope = WorkspaceAgentScope {
scope_id: "comp-drain-contender-test".to_string(),
relay_url: "wss://relay.example".to_string(),
owner_pubkey: "aa".repeat(32),
definitions_dir: tmp_path.clone(),
@@ -219,74 +233,54 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
};
state.commit_active_scope(scope.clone());
// Produce `stopped = [entry1]` — entry1 is the only successfully stopped agent.
let pubkey1 = "aa".repeat(32);
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped_from_drain = vec![entry1.clone()];
let stopped = vec![entry1];
// The contender just acquires the transition mutex and reports when it
// managed to do so.
// Channels for barrier-based ordering.
let (contender_ready_tx, contender_ready_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.
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();
// Acquire the transition guard in THIS thread FIRST, before the contender
// is spawned. This guarantees the contender will block when it tries to
// acquire the same mutex.
let transition_guard = state.managed_agent_runtime_transition.lock().unwrap();
let contender = thread::spawn(move || {
use tauri::Manager;
let state = app_contender.state::<crate::app_state::AppState>();
// Signal that the contender is alive and about to try the lock.
// Signal that the contender is alive and about to block on the mutex.
contender_ready_tx.send(()).unwrap();
// Try to acquire the transition mutex — blocks while the test holds it.
let _guard = state.managed_agent_runtime_transition.lock().unwrap();
// Signal: contender now holds the lock (compensation must be done).
// 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).
contender_done_tx.send(()).unwrap();
});
// Wait until contender is alive and parked on (or about to park on) the mutex.
// Wait until the contender is alive and parked (or about to park) on the mutex.
contender_ready_rx.recv().unwrap();
// Give the contender a moment to reach the mutex and block on it.
std::thread::sleep(std::time::Duration::from_millis(20));
// Confirm contender is NOT done yet (still blocked on the mutex).
// 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 lock-free core while the transition guard is held.
let mut records_for_restore = vec![];
let restored = compensate_drain_for(
&stopped_from_drain,
&mut records_for_restore,
|entry, _recs| {
// Verify entry1 is restored with exact relay and start_on_app_launch.
assert_eq!(entry.key.pubkey, pubkey1, "only entry1 must be restored");
assert_eq!(
entry.key.relay_url, "wss://relay.example",
"relay must match"
);
assert!(entry.start_on_app_launch, "start_on_app_launch must match");
Ok(())
},
);
// entry1 restored successfully → None (no degradation).
assert!(
restored.is_none(),
"entry1 restoration must succeed: {restored:?}"
);
// ── 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.
// Release the transition guard (compensation done).
drop(transition_guard);
// Contender must now be able to acquire the mutex.
// 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 compensation releases the transition guard");
.expect("contender must unblock after compensate_drain releases the transition guard");
contender.join().expect("contender thread panicked");
}