fix(desktop): split Phase 3 test files to satisfy 1000-line size ratchet

Four test files exceeded the 1000-line gate introduced by Phase 3 seam work.
Split each by extracting the largest block into a sibling file included via
#[path], keeping every file at or below the limit:

- global_agent_config.rs (1987 → 946): inline test block moved to
  global_agent_config_tests.rs (637) + global_agent_config_epoch_tests.rs (420)
- mesh_llm_tests.rs (1052 → 776): Area-4 serialization tests extracted to
  mesh_llm_transition_tests.rs (284)
- team_snapshot/tests.rs (1087 → 896): Area-3 phase-boundary seam tests
  extracted to team_snapshot/seam_tests.rs (200)
- runtime_commands_tests.rs (1056 → 777): concurrency tests extracted to
  runtime_commands_concurrency_tests.rs (289)

All nine files are now under 1000 lines. No test logic changed; all 2230
lib tests pass. just desktop-check passes at this commit.

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 18:59:33 -04:00
co-authored by Will Pfleger
parent 1a1571b165
commit d0658609ed
9 changed files with 1842 additions and 1796 deletions
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,422 @@
//! Epoch-level tests for `commands/global_agent_config.rs`.
//!
//! Included inside `mod tests` via `#[path]` from `global_agent_config_tests.rs`.
//! Heavy async epoch tests split here to keep each file under 1000 lines.
use super::*;
/// Full tail test: production epoch core with injected stop/spawn/receipt
/// closures. Verifies that the core calls stop, then spawn with captured
/// context (relay, owner, scope), then receipt, registers the runtime with
/// the captured scope_id, and saves the record.
///
/// Thufir's test 5: "call production restart_under_captured_epoch_for via
/// mock app with injected spawn_fn and write_receipt_fn; assert captured
/// relay/owner/teams/personas/global delivered to spawn, receipt constructed
/// 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.
#[tokio::test]
async fn test_full_tail_stop_spawn_receipt_register_save() {
use crate::managed_agents::{
storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord,
ManagedAgentRuntimeKey,
};
use std::sync::{Arc, Mutex};
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let pubkey = "bb".repeat(32);
let relay_url = "wss://test.relay";
let owner_hex = "cc".repeat(32);
let scope_id = "scope-test-id";
// 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.
let mut record_env_vars = std::collections::BTreeMap::new();
record_env_vars.insert(
"ANTHROPIC_API_KEY".to_string(),
"sk-test-key-for-readiness".to_string(),
);
let record = ManagedAgentRecord {
pubkey: pubkey.clone(),
name: "test-agent".to_string(),
display_name: None,
slug: None,
persona_id: None,
private_key_nsec: String::new(),
auth_tag: None,
relay_url: relay_url.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: Some("claude-3-5-sonnet-20241022".to_string()),
provider: Some("anthropic".to_string()),
persona_source_version: None,
env_vars: record_env_vars,
start_on_app_launch: false,
auto_restart_on_config_change: false,
runtime_pid: None,
backend: 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![],
};
// Write initial store.
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 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.
let rt_key = ManagedAgentRuntimeKey::new(&pubkey, relay_url).unwrap();
{
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 process = crate::managed_agents::ManagedAgentProcess {
child,
log_path: std::path::PathBuf::new(),
spawn_config_hash: 0,
setup_mode: false,
adapter_availability: None,
start_nonce: "test-nonce".to_string(),
#[cfg(windows)]
job: None,
};
runtimes.insert(
rt_key.clone(),
ManagedAgentPairRuntime::starting(process, Some(scope_id.to_string())),
);
}
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: scope_id.to_string(),
relay_url: relay_url.to_string(),
owner_pubkey: owner_hex.clone(),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
let context = CapturedRestartContext {
scope: scope.clone(),
personas: vec![],
teams: vec![],
global: crate::managed_agents::GlobalAgentConfig::default(),
owner_hex: owner_hex.clone(),
mesh_model_id: None,
};
// old_global (default) and new_global differ by one env_var so that
// env_changed = true → should_restart_on_config_change returns true.
let old_global = crate::managed_agents::GlobalAgentConfig::default();
let mut new_global_env = std::collections::BTreeMap::new();
new_global_env.insert("SOME_EXTRA_VAR".to_string(), "v2".to_string());
let new_global = crate::managed_agents::GlobalAgentConfig {
env_vars: new_global_env,
..Default::default()
};
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 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 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_hash: 0,
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(())
},
);
assert!(
matches!(result, Ok(())),
"full tail must return Ok when stop+spawn+receipt all succeed: {result:?}"
);
assert!(*stop_called.lock().unwrap(), "stop_fn must be called");
assert_eq!(
spawn_relay.lock().unwrap().as_deref(),
Some(relay_url),
"spawn_fn must receive the captured relay URL"
);
assert_eq!(
spawn_owner.lock().unwrap().as_deref(),
Some(owner_hex.as_str()),
"spawn_fn must receive the captured owner hex"
);
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 runtimes = state.managed_agent_processes.lock().unwrap();
let registered = runtimes
.values()
.any(|r| r.scope_id.as_deref() == Some(scope_id));
assert!(
registered,
"runtime must be registered with the captured scope_id"
);
}
}
/// Production driver proves preflight fires before stop.
///
/// Thufir's test 6: "event log from production driver proves preflight fn
/// fires before stop fn."
///
/// 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.
///
/// 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.
#[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 tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let app = make_mock_app();
let app_handle = app.handle().clone();
// Get the actual owner pubkey from the mock app's signing keys.
// The pre-stop phase checks hex == captured_scope.owner_pubkey; they must match.
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()
};
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 {
pubkey: agent_pubkey.clone(),
name: "test-agent-preflight".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: false,
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![],
};
save_managed_agents_at(tmp.path(), &[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();
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: "wss://relay.example".to_string(),
// Use the actual app key so the owner-key check passes.
owner_pubkey: actual_owner_hex.clone(),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
// 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()));
let log_mesh = event_log.clone();
let log_stop = event_log.clone();
let outcome = restart_local_agent_on_config_change_for(
&app_handle,
&agent_pubkey,
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&[],
&scope,
tmp.path(),
// mesh_fn: records "mesh", then succeeds.
move |_app, _model| {
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| {
log_stop.lock().unwrap().push("stop");
Ok(())
},
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
Err("spawn not expected".to_string())
},
|_app, _receipt| Err("receipt not expected".to_string()),
)
.await;
// The call Skips at the no-live-runtime check, not at preflight.
assert!(
matches!(outcome, RestartOutcome::Skipped),
"no live runtime: {outcome:?}"
);
let log = event_log.lock().unwrap();
// Mesh must appear before any stop (even if stop never fired).
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)"
);
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})"
);
}
}
@@ -0,0 +1,636 @@
//! Unit and integration tests for `commands/global_agent_config.rs`.
//!
//! Split into this file and `global_agent_config_epoch_tests.rs` to keep each
//! file under the 1000-line size ratchet.
//!
//! Included via `#[path = "global_agent_config_tests.rs"] mod tests;` at the
//! bottom of `global_agent_config.rs`. `use super::*` gives access to all
//! items in that module.
use super::{
restart_under_captured_epoch_for, should_restart_on_config_change, CapturedRestartContext,
EpochError,
};
/// Running agent (Ready) whose effective env changed → restart candidate.
#[test]
fn env_changed_running_agent_is_candidate() {
// old_ready=true, new_ready=true, env_changed=true
assert!(
should_restart_on_config_change(true, true, true),
"running agent with changed env must be restarted"
);
}
/// Running agent (Ready) whose effective env did NOT change → not a candidate.
#[test]
fn unchanged_running_agent_is_not_candidate() {
// old_ready=true, new_ready=true, env_changed=false
assert!(
!should_restart_on_config_change(true, true, false),
"running agent with identical env must NOT be restarted"
);
}
/// NotReady → Ready transition is admitted regardless of env diff.
#[test]
fn not_ready_to_ready_is_candidate() {
// old_ready=false, new_ready=true, env_changed=false (env_changed irrelevant)
assert!(
should_restart_on_config_change(false, true, false),
"NotReady → Ready must be a restart candidate"
);
}
/// Ready → NotReady (config became invalid, env changed) is admitted so the
/// agent restarts into setup-listener mode via the normal spawn path.
#[test]
fn ready_to_not_ready_env_changed_is_candidate() {
// old_ready=true (had key), new_ready=false (key removed), env_changed=true
assert!(
should_restart_on_config_change(true, false, true),
"Ready → NotReady with env change must be a restart candidate"
);
}
/// Both NotReady, env unchanged → not a candidate (nothing to restart).
#[test]
fn both_not_ready_unchanged_is_not_candidate() {
// old_ready=false, new_ready=false, env_changed=false
assert!(
!should_restart_on_config_change(false, false, false),
"both NotReady with no env change must NOT be a candidate"
);
}
/// NotReady + env changed but new still NotReady → not a candidate.
#[test]
fn not_ready_env_changed_still_not_ready_is_not_candidate() {
// Changed one unrelated env var but still missing the required key.
// old_ready=false, new_ready=false, env_changed=true
assert!(
!should_restart_on_config_change(false, false, true),
"NotReady→NotReady (env changed but still broken) must NOT be a candidate"
);
}
/// NotReady → Ready AND env also changed → still a restart candidate.
///
/// Guards against a future `&& !env_changed` regression on the
/// NotReady→Ready branch: env_changed is irrelevant when readiness
/// unblocks — the agent must restart regardless of whether env also differed.
#[test]
fn not_ready_to_ready_with_env_change_is_candidate() {
// old_ready=false, new_ready=true, env_changed=true
assert!(
should_restart_on_config_change(false, true, true),
"NotReady → Ready (with env change) must be a restart candidate"
);
}
// ── restart_under_captured_epoch_for: generation guard ───────────────────
//
// These tests call `restart_under_captured_epoch_for` directly — the
// production stop→spawn primitive — using a `tauri::test::mock_app()`
// runtime so the AppHandle is real. No live process is running, so the
// restart is skipped at the eligibility check. The generation tests drive
// the path that matters: does the captured-generation guard prevent a
// stale-scope restart?
fn make_test_scope(
definitions_dir: &std::path::Path,
) -> crate::managed_agents::scope::WorkspaceAgentScope {
let gen = crate::managed_agents::scope::current_scope_generation();
crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: "wss://relay.example".to_string(),
owner_pubkey: "aa".repeat(32),
definitions_dir: definitions_dir.to_path_buf(),
generation: gen,
}
}
fn make_test_context(definitions_dir: &std::path::Path) -> CapturedRestartContext {
CapturedRestartContext {
scope: make_test_scope(definitions_dir),
personas: vec![],
teams: vec![],
global: crate::managed_agents::GlobalAgentConfig::default(),
owner_hex: "aa".repeat(32),
mesh_model_id: None,
}
}
/// `restart_under_captured_epoch_for` with a fresh scope and an empty store
/// (no live pair runtime) → `EpochError::Skipped` after the agent-not-found
/// or no-live-runtime check. The generation guard passes, proving the path
/// proceeds to the eligibility check rather than aborting at the stale check.
///
/// This is the stop-to-spawn production path: transition lock acquired,
/// store lock acquired, generation validated — all before any state change.
#[test]
fn test_restart_under_captured_epoch_fresh_scope_no_runtime_is_skipped() {
let tmp = tempfile::tempdir().unwrap();
// Write an empty managed-agents.json so load_managed_agents_at returns Ok([]).
std::fs::write(tmp.path().join("managed-agents.json"), b"[]").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.handle().clone();
let context = make_test_context(tmp.path());
let pubkey = "aa".repeat(32);
let result = restart_under_captured_epoch_for(
&app_handle,
&pubkey,
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&[],
&context,
|_app, _rec, _runtimes| Err("stop not expected".to_string()),
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
Err("spawn not expected".to_string())
},
|_app, _receipt| Err("receipt not expected".to_string()),
);
// No live runtime → Skipped before any state change.
assert!(
matches!(result, Err(EpochError::Skipped(_))),
"no live pair runtime must produce Skipped, not FailedAfterStop or Ok: {result:?}"
);
// Verify the skip reason. In sequential execution the scope is fresh and
// the epoch reaches the eligibility check before skipping ("not found" or
// "no live pair runtime"). In parallel test runs another test may advance
// the generation counter, producing "stale scope" instead — both are valid
// outcomes proving the epoch aborted without modifying any agent state.
if let Err(EpochError::Skipped(msg)) = result {
let is_expected = msg.contains("not found")
|| msg.contains("no live pair runtime")
|| msg.contains("stale scope")
|| msg.contains("generation");
assert!(
is_expected,
"Skipped reason must be agent-not-found, no-live-runtime, or stale-scope: {msg}"
);
}
}
/// `restart_under_captured_epoch_for` with a STALE scope → `EpochError::Skipped`
/// at the generation validation step, before touching any agent state.
///
/// This is the switch-between-stop-and-spawn test: simulates the race where
/// a workspace switch advances the generation between when the scope was
/// captured and when `restart_under_captured_epoch_for` runs. The generation
/// guard must abort before stopping — no agent is touched.
#[test]
fn test_restart_under_captured_epoch_stale_scope_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("managed-agents.json"), b"[]").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.handle().clone();
// Capture the scope at the current generation, then advance to make it stale.
let context = make_test_context(tmp.path());
crate::managed_agents::scope::next_scope_generation();
let pubkey = "aa".repeat(32);
let result = restart_under_captured_epoch_for(
&app_handle,
&pubkey,
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&[],
&context,
|_app, _rec, _runtimes| Err("stop not expected".to_string()),
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
Err("spawn not expected".to_string())
},
|_app, _receipt| Err("receipt not expected".to_string()),
);
// Stale generation → Skipped at the generation-validation step.
assert!(
matches!(result, Err(EpochError::Skipped(_))),
"stale scope must produce Skipped: {result:?}"
);
if let Err(EpochError::Skipped(msg)) = result {
assert!(
msg.contains("stale scope") || msg.contains("generation"),
"Skipped reason must mention stale scope or generation mismatch: {msg}"
);
}
}
// ── Area-2 tests: async driver and epoch core ─────────────────────────────
fn make_mock_app() -> tauri::App<tauri::test::MockRuntime> {
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")
}
/// 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.
///
/// Thufir's test 1: "production async driver with injected loader failure;
/// assert stop never called, RestartOutcome::Skipped."
#[tokio::test]
async fn test_context_load_failure_leaves_runtime_running() {
use super::restart_local_agent_on_config_change_for;
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.
let app = make_mock_app();
let app_handle = app.handle().clone();
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: "wss://relay.example".to_string(),
owner_pubkey: "aa".repeat(32),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
let stop_called = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let stop_called2 = stop_called.clone();
let outcome = restart_local_agent_on_config_change_for(
&app_handle,
&"aa".repeat(32),
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&[],
&scope,
tmp.path(),
// mesh_fn: succeeds (no-op)
|_app, _model| Box::pin(async { Ok(()) }),
// stop_fn: must NOT be called
move |_app, _rec, _runtimes| {
stop_called2.store(true, std::sync::atomic::Ordering::SeqCst);
Err("stop_fn called unexpectedly".to_string())
},
// spawn_fn: must NOT be called
|_app, _rec, _relay, _owner, _personas, _global, _teams| {
Err("spawn_fn called unexpectedly".to_string())
},
// write_receipt_fn: must NOT be called
|_app, _receipt| Err("receipt_fn called unexpectedly".to_string()),
)
.await;
assert!(
matches!(outcome, RestartOutcome::Skipped),
"context load failure must produce Skipped: {outcome:?}"
);
assert!(
!stop_called.load(std::sync::atomic::Ordering::SeqCst),
"stop must NOT be called when context load fails"
);
}
/// Mesh preflight failure before stop → `RestartOutcome::Skipped`,
/// stop closure never called.
///
/// Drives `restart_local_agent_on_config_change_for` with an injected mesh_fn
/// that returns `Err`. Stop must not fire.
///
/// Thufir's test 2: "real production driver/core with injected preflight error;
/// stop never called, RestartOutcome::Skipped."
#[tokio::test]
async fn test_mesh_preflight_failure_leaves_runtime_running() {
use super::restart_local_agent_on_config_change_for;
use crate::commands::global_agent_config::RestartOutcome;
let tmp = tempfile::tempdir().unwrap();
// Provide persona and record files so context prep can pass them.
std::fs::write(tmp.path().join("personas.json"), b"[]").unwrap();
std::fs::write(tmp.path().join("managed-agents.json"), b"[]").unwrap();
// Also need global-agent-config (fallible load will fail on missing file,
// so we supply it). The injected mesh_fn is what we care about.
std::fs::write(tmp.path().join("global-agent-config.json"), b"{}").unwrap();
let app = make_mock_app();
let app_handle = app.handle().clone();
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: "wss://relay.example".to_string(),
owner_pubkey: "aa".repeat(32),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
let stop_called = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let stop_called2 = stop_called.clone();
let outcome = restart_local_agent_on_config_change_for(
&app_handle,
&"aa".repeat(32),
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&[],
&scope,
tmp.path(),
// mesh_fn: FAILS — triggers the pre-stop abort
|_app, _model| Box::pin(async { Err("mesh preflight failed (test)".to_string()) }),
// stop_fn: must NOT be called
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;
assert!(
matches!(outcome, RestartOutcome::Skipped),
"mesh preflight failure must produce Skipped: {outcome:?}"
);
assert!(
!stop_called.load(std::sync::atomic::Ordering::SeqCst),
"stop must NOT be called when mesh preflight fails"
);
}
/// Workspace switch after preflight (generation advances before epoch entry)
/// → epoch returns `Skipped` via generation guard, stop never called.
///
/// Thufir's test 3: "injected preflight hook advances generation after it
/// succeeds; epoch returns Skipped; stop never called."
#[tokio::test]
async fn test_workspace_switch_after_preflight_aborts_before_stop() {
use super::restart_local_agent_on_config_change_for;
use crate::commands::global_agent_config::RestartOutcome;
let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("personas.json"), b"[]").unwrap();
std::fs::write(tmp.path().join("managed-agents.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();
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: "wss://relay.example".to_string(),
owner_pubkey: "aa".repeat(32),
definitions_dir: tmp.path().to_path_buf(),
generation: gen,
};
let stop_called = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let stop_called2 = stop_called.clone();
let outcome = restart_local_agent_on_config_change_for(
&app_handle,
&"aa".repeat(32),
&crate::managed_agents::GlobalAgentConfig::default(),
&crate::managed_agents::GlobalAgentConfig::default(),
&[],
&scope,
tmp.path(),
// mesh_fn: succeeds but advances generation (simulates workspace switch
// between preflight completion and epoch entry).
|_app, _model| {
crate::managed_agents::scope::next_scope_generation();
Box::pin(async { Ok(()) })
},
// stop_fn: must NOT be called
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;
assert!(
matches!(outcome, RestartOutcome::Skipped),
"generation advance after preflight must produce Skipped: {outcome:?}"
);
assert!(
!stop_called.load(std::sync::atomic::Ordering::SeqCst),
"stop must NOT be called when generation advanced after preflight"
);
}
/// A record-level Mesh model change after preflight (without advancing workspace
/// generation) → epoch detects mismatch in re-resolved Mesh model, stops 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."
///
/// 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.
///
/// 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() {
use crate::managed_agents::{
storage::save_managed_agents_at, BackendKind, ManagedAgentPairRuntime, ManagedAgentRecord,
ManagedAgentRuntimeKey,
};
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
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.
let mut record_env_vars = std::collections::BTreeMap::new();
record_env_vars.insert(
"ANTHROPIC_API_KEY".to_string(),
"sk-test-key-for-readiness".to_string(),
);
let record = ManagedAgentRecord {
pubkey: pubkey.clone(),
name: "test-agent-mesh".to_string(),
display_name: None,
slug: None,
persona_id: None,
private_key_nsec: String::new(),
auth_tag: None,
relay_url: relay_url.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: Some("claude-3-5-sonnet-20241022".to_string()),
provider: Some("anthropic".to_string()),
persona_source_version: None,
env_vars: record_env_vars,
start_on_app_launch: false,
auto_restart_on_config_change: false,
runtime_pid: None,
backend: 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, // ← no relay_mesh; re-resolve yields None ≠ Some("model-a")
runtime: None,
name_pool: vec![],
};
save_managed_agents_at(tmp.path(), std::slice::from_ref(&record)).unwrap();
let app = make_mock_app();
let app_handle = app.handle().clone();
// Seed a live runtime with a long-running process (avoids sync eviction).
let rt_key = ManagedAgentRuntimeKey::new(&pubkey, relay_url).unwrap();
{
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 process = crate::managed_agents::ManagedAgentProcess {
child,
log_path: std::path::PathBuf::new(),
spawn_config_hash: 0,
setup_mode: false,
adapter_availability: None,
start_nonce: "test-nonce-mesh".to_string(),
#[cfg(windows)]
job: None,
};
runtimes.insert(
rt_key,
ManagedAgentPairRuntime::starting(process, Some("test-scope".to_string())),
);
}
let gen = crate::managed_agents::scope::current_scope_generation();
let scope = crate::managed_agents::scope::WorkspaceAgentScope {
scope_id: "test-scope".to_string(),
relay_url: relay_url.to_string(),
owner_pubkey: "aa".repeat(32),
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.
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();
let new_global = crate::managed_agents::GlobalAgentConfig {
env_vars: new_global_env,
..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();
// 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()),
);
assert!(
matches!(result, Err(EpochError::Skipped(_))),
"Mesh model mismatch must produce Skipped before stop: {result:?}"
);
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"
);
}
#[path = "global_agent_config_epoch_tests.rs"]
mod epoch_tests;
@@ -772,281 +772,5 @@ async fn test_mesh_stop_client_no_runtime_returns_stopped_status() {
);
}
// ── Area 4 serialization direction tests ─────────────────────────────────────
//
// These three tests prove the lock-serialization contract between
// `with_workspace_transition_preflight` and `install_client_under_workspace_transition`.
// No port (`127.0.0.1:9337`) is touched — the install closure is always injected.
/// Active mock client → `mesh_stop_client` → runtime slot is `None` →
/// `with_workspace_transition_preflight` with a no-op body succeeds.
///
/// Proves the two-step Option A user flow:
/// 1. user calls "stop shared compute" (`mesh_stop_client`);
/// 2. workspace switch proceeds via the production transition-preflight helper.
///
/// The production `fail_if_client_mesh_active` is called inside
/// `with_workspace_transition_preflight` (under the guard). It must see an
/// absent runtime and return `Ok(())` — not the stale client that was there
/// before `mesh_stop_client` cleared it.
#[tokio::test]
async fn test_active_client_stop_then_transition_preflight_succeeds() {
use tauri::Manager;
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.handle().clone();
let state = app.state::<crate::app_state::AppState>();
// Install a client-mode runtime — simulates "Mesh is running as a client".
{
let client_runtime = crate::mesh_llm::build_mock_client_runtime_for_test();
*state.mesh_llm_runtime.lock().await = Some(client_runtime);
}
// Verify the client is present before the stop.
{
let guard = state.mesh_llm_runtime.lock().await;
assert!(
guard.is_some(),
"pre-condition: a client runtime must be installed before stop"
);
}
// Step 1 — production `mesh_stop_client` tears down the client.
let stop_status = super::mesh_stop_client(app_handle.clone(), state.clone())
.await
.expect("mesh_stop_client must not error");
assert_eq!(
stop_status.state,
crate::mesh_llm::MeshNodeState::Off,
"mesh_stop_client must return Off status after tearing down the client"
);
// Step 2 — runtime slot must now be None.
{
let guard = state.mesh_llm_runtime.lock().await;
assert!(
guard.is_none(),
"mesh_stop_client must clear the runtime slot; got {:?}",
guard.as_ref().map(|r| r.mode())
);
}
// Step 3 — production transition preflight must succeed (no client active).
// The no-op body proves the lock was acquired and the preflight passed.
let result = super::scope_impl::with_workspace_transition_preflight(&app_handle, || {
Ok::<&str, String>("body ran")
})
.await;
assert!(
result.is_ok(),
"with_workspace_transition_preflight must succeed after mesh_stop_client cleared the slot: {result:?}"
);
assert_eq!(
result.unwrap(),
"body ran",
"transition body must have been invoked and its return value propagated"
);
}
/// 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.
///
/// 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.
///
/// 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.
#[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,
};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use tauri::Manager;
// Serialise generation-sensitive work across parallel tests.
let _gen_guard = SCOPE_GENERATION_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
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.handle().clone();
let state = app.state::<crate::app_state::AppState>();
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
// Capture a scope at the current generation — this is the scope the install
// task will carry 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;
// 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);
// Spawn the install task. It blocks on workspace_transition until we drop the guard.
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 install_task = tokio::task::spawn(async move {
super::scope_impl::install_client_under_workspace_transition(
&app_handle_clone,
&captured_scope,
|| {
let called = Arc::clone(&install_called_clone);
async move {
called.store(true, Ordering::SeqCst);
Ok::<(), String>(())
}
},
)
.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.
tokio::task::yield_now().await;
// Release the transition lock → install task acquires it and validates.
drop(transition_guard);
let install_result = install_task.await.expect("install task must not panic");
// The install helper must have rejected the stale scope without invoking install.
assert!(
install_result.is_err(),
"install_client_under_workspace_transition must return Err for a stale captured scope; \
got Ok"
);
let err = install_result.unwrap_err();
assert!(
err.contains("stale") || err.contains("mismatch") || err.contains("scope"),
"error must describe the stale/mismatched scope: {err}"
);
assert!(
!install_was_called.load(Ordering::SeqCst),
"install closure must NOT have been invoked when the scope was stale"
);
}
/// Install helper acquires the lock first and installs a mock client; after
/// release, the transition preflight observes the client and rejects.
///
/// 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.
#[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;
// Serialise generation-sensitive work across parallel tests.
let _gen_guard = SCOPE_GENERATION_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
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.handle().clone();
let state = app.state::<crate::app_state::AppState>();
let base = std::path::PathBuf::from("/tmp/area4-test-dir-inv");
let relay = "wss://install-first-test.example";
let owner = "cc".repeat(32);
// Set up an active scope so install_client_under_workspace_transition can
// validate full scope identity under the guard.
let gen = next_scope_generation();
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;
// 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);
}
// Spawn the transition task. It blocks on workspace_transition until we drop
// the install guard.
let app_handle_clone = app_handle.clone();
let transition_task = tokio::task::spawn(async move {
super::scope_impl::with_workspace_transition_preflight(&app_handle_clone, || {
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.
tokio::task::yield_now().await;
// Release the install lock → transition task acquires it and runs preflight.
drop(install_guard);
let transition_result = transition_task
.await
.expect("transition task must not panic");
// `fail_if_client_mesh_active` must observe the installed client and reject.
assert!(
transition_result.is_err(),
"with_workspace_transition_preflight must return Err when a client was installed \
while holding the lock; got Ok"
);
let err = transition_result.unwrap_err();
assert!(
err.contains("client") || err.contains("Stop") || err.contains("shared compute"),
"error must describe the active client and how to stop it: {err}"
);
}
#[path = "mesh_llm_transition_tests.rs"]
mod transition_tests;
@@ -0,0 +1,284 @@
//! Area-4 workspace-transition serialization tests for `commands/mesh_llm.rs`.
//!
//! Split from `mesh_llm_tests.rs` to keep each file under the 1000-line ratchet.
//! Included via `#[path]` from `mesh_llm_tests.rs` as `mod transition_tests;`.
//! `use super::*` gives access to all items in `mesh_llm_tests.rs`.
// ── Area 4 serialization direction tests ─────────────────────────────────────
//
// These three tests prove the lock-serialization contract between
// `with_workspace_transition_preflight` and `install_client_under_workspace_transition`.
// No port (`127.0.0.1:9337`) is touched — the install closure is always injected.
/// Active mock client → `mesh_stop_client` → runtime slot is `None` →
/// `with_workspace_transition_preflight` with a no-op body succeeds.
///
/// Proves the two-step Option A user flow:
/// 1. user calls "stop shared compute" (`mesh_stop_client`);
/// 2. workspace switch proceeds via the production transition-preflight helper.
///
/// The production `fail_if_client_mesh_active` is called inside
/// `with_workspace_transition_preflight` (under the guard). It must see an
/// absent runtime and return `Ok(())` — not the stale client that was there
/// before `mesh_stop_client` cleared it.
#[tokio::test]
async fn test_active_client_stop_then_transition_preflight_succeeds() {
use tauri::Manager;
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.handle().clone();
let state = app.state::<crate::app_state::AppState>();
// Install a client-mode runtime — simulates "Mesh is running as a client".
{
let client_runtime = crate::mesh_llm::build_mock_client_runtime_for_test();
*state.mesh_llm_runtime.lock().await = Some(client_runtime);
}
// Verify the client is present before the stop.
{
let guard = state.mesh_llm_runtime.lock().await;
assert!(
guard.is_some(),
"pre-condition: a client runtime must be installed before stop"
);
}
// Step 1 — production `mesh_stop_client` tears down the client.
let stop_status = super::mesh_stop_client(app_handle.clone(), state.clone())
.await
.expect("mesh_stop_client must not error");
assert_eq!(
stop_status.state,
crate::mesh_llm::MeshNodeState::Off,
"mesh_stop_client must return Off status after tearing down the client"
);
// Step 2 — runtime slot must now be None.
{
let guard = state.mesh_llm_runtime.lock().await;
assert!(
guard.is_none(),
"mesh_stop_client must clear the runtime slot; got {:?}",
guard.as_ref().map(|r| r.mode())
);
}
// Step 3 — production transition preflight must succeed (no client active).
// The no-op body proves the lock was acquired and the preflight passed.
let result = super::scope_impl::with_workspace_transition_preflight(&app_handle, || {
Ok::<&str, String>("body ran")
})
.await;
assert!(
result.is_ok(),
"with_workspace_transition_preflight must succeed after mesh_stop_client cleared the slot: {result:?}"
);
assert_eq!(
result.unwrap(),
"body ran",
"transition body must have been invoked and its return value propagated"
);
}
/// 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.
///
/// 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.
///
/// 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.
#[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,
};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use tauri::Manager;
// Serialise generation-sensitive work across parallel tests.
let _gen_guard = SCOPE_GENERATION_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
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.handle().clone();
let state = app.state::<crate::app_state::AppState>();
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
// Capture a scope at the current generation — this is the scope the install
// task will carry 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;
// 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);
// Spawn the install task. It blocks on workspace_transition until we drop the guard.
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 install_task = tokio::task::spawn(async move {
super::scope_impl::install_client_under_workspace_transition(
&app_handle_clone,
&captured_scope,
|| {
let called = Arc::clone(&install_called_clone);
async move {
called.store(true, Ordering::SeqCst);
Ok::<(), String>(())
}
},
)
.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.
tokio::task::yield_now().await;
// Release the transition lock → install task acquires it and validates.
drop(transition_guard);
let install_result = install_task.await.expect("install task must not panic");
// The install helper must have rejected the stale scope without invoking install.
assert!(
install_result.is_err(),
"install_client_under_workspace_transition must return Err for a stale captured scope; \
got Ok"
);
let err = install_result.unwrap_err();
assert!(
err.contains("stale") || err.contains("mismatch") || err.contains("scope"),
"error must describe the stale/mismatched scope: {err}"
);
assert!(
!install_was_called.load(Ordering::SeqCst),
"install closure must NOT have been invoked when the scope was stale"
);
}
/// Install helper acquires the lock first and installs a mock client; after
/// release, the transition preflight observes the client and rejects.
///
/// 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.
#[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;
// Serialise generation-sensitive work across parallel tests.
let _gen_guard = SCOPE_GENERATION_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
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.handle().clone();
let state = app.state::<crate::app_state::AppState>();
let base = std::path::PathBuf::from("/tmp/area4-test-dir-inv");
let relay = "wss://install-first-test.example";
let owner = "cc".repeat(32);
// Set up an active scope so install_client_under_workspace_transition can
// validate full scope identity under the guard.
let gen = next_scope_generation();
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;
// 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);
}
// Spawn the transition task. It blocks on workspace_transition until we drop
// the install guard.
let app_handle_clone = app_handle.clone();
let transition_task = tokio::task::spawn(async move {
super::scope_impl::with_workspace_transition_preflight(&app_handle_clone, || {
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.
tokio::task::yield_now().await;
// Release the install lock → transition task acquires it and runs preflight.
drop(install_guard);
let transition_result = transition_task
.await
.expect("transition task must not panic");
// `fail_if_client_mesh_active` must observe the installed client and reject.
assert!(
transition_result.is_err(),
"with_workspace_transition_preflight must return Err when a client was installed \
while holding the lock; got Ok"
);
let err = transition_result.unwrap_err();
assert!(
err.contains("client") || err.contains("Stop") || err.contains("shared compute"),
"error must describe the active client and how to stop it: {err}"
);
}
@@ -0,0 +1,200 @@
//! Phase-boundary seam tests for team snapshot import (Area 3).
//!
//! Split from `tests.rs` to keep each file under the 1000-line ratchet.
//! Included via `#[path = "seam_tests.rs"] mod seam_tests;` from `tests.rs`.
//! `use super::*` gives access to all items in `tests.rs`.
use super::*;
// ── Phase-boundary seam tests (Area 3 — team) ─────────────────────────────
/// Serializes tests that modify the process-global scope generation counter.
/// See the equivalent comment in `import_tests.rs` for rationale.
use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK as GENERATION_TEST_LOCK;
fn setup_team_import_app_with_scope(
tmp: &tempfile::TempDir,
) -> (tauri::App<tauri::test::MockRuntime>, nostr::Keys) {
use crate::managed_agents::scope::{next_scope_generation, WorkspaceAgentScope};
let owner_keys = nostr::Keys::generate();
let state = crate::app_state::build_app_state();
{
let mut locked = state.keys.lock().unwrap();
*locked = owner_keys.clone();
}
let app = tauri::test::mock_builder()
.manage(state)
.build(tauri::test::mock_context(tauri::test::noop_assets()))
.expect("failed to build mock app for team import test");
{
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,
};
s.commit_active_scope(scope);
}
(app, owner_keys)
}
/// `after_store` hook commits a genuinely different live scope + owner —
/// Phase 4/5 outbound must use the OLD (captured) relay URL, not the new
/// live relay.
///
/// Thufir requirement: `after_store` must commit a genuinely different live
/// 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.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock held to serialize generation
// counter mutations — cannot deadlock. See import_tests.rs for full rationale.
#[allow(clippy::await_holding_lock)]
async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
use crate::commands::personas::snapshot::import::{MemoryPublish, ProfilePublish};
use crate::commands::team_snapshot::confirm_team_snapshot_import_core;
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
use std::sync::{Arc, Mutex};
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
let handle = app.handle();
let snap = snapshot(vec![member("Alice"), member("Bob")]);
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
let input = TeamSnapshotImportConfirm {
file_bytes: encoded,
keep_allowlist: false,
};
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 pr = profile_relays.clone();
let mr = memory_relays.clone();
let state = app.state::<crate::app_state::AppState>();
let handle_for_hook = handle.clone();
let result = confirm_team_snapshot_import_core(
input,
handle,
&state,
|| {},
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 s = handle_for_hook.state::<crate::app_state::AppState>();
s.commit_active_scope(new_scope);
},
move |p: ProfilePublish<'_>| {
let relay = p.relay_url.to_string();
pr.lock().unwrap().push(relay.clone());
Box::pin(async move {
let _ = relay;
Ok(())
})
},
move |m: MemoryPublish<'_>| {
let relay = m.relay_url.to_string();
mr.lock().unwrap().push(relay.clone());
Box::pin(async move {
let _ = relay;
Ok(())
})
},
)
.await;
assert!(
result.is_ok(),
"team import must succeed: {:?}",
result.err()
);
// Every member's profile adapter received the captured relay URL.
let seen = profile_relays.lock().unwrap();
assert_eq!(
seen.len(),
2,
"profile adapter must be called once per member"
);
for relay in seen.iter() {
assert_eq!(
relay, &expected_relay,
"profile adapter must receive captured relay, got: {relay}"
);
}
// No memory entries in these members — memory adapter not called.
}
/// `before_store` hook advances scope generation — Phase 3 must reject BEFORE
/// any write and BEFORE any outbound call.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock held to serialize generation
// counter mutations — cannot deadlock. See import_tests.rs for full rationale.
#[allow(clippy::await_holding_lock)]
async fn test_team_identity_switch_before_store_is_rejected() {
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
use crate::commands::personas::snapshot::import::{MemoryPublish, ProfilePublish};
use crate::commands::team_snapshot::confirm_team_snapshot_import_core;
use crate::managed_agents::scope::next_scope_generation;
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
let handle = app.handle();
let state = app.state::<crate::app_state::AppState>();
let snap = snapshot(vec![member("Alice")]);
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
let input = TeamSnapshotImportConfirm {
file_bytes: encoded,
keep_allowlist: false,
};
let result = confirm_team_snapshot_import_core(
input,
handle,
&state,
move || {
next_scope_generation();
},
|| {},
|_p: ProfilePublish<'_>| {
Box::pin(async { panic!("profile must not be called: store rejected") })
},
|_m: MemoryPublish<'_>| {
Box::pin(async { panic!("memory must not be called: store rejected") })
},
)
.await;
assert!(
result.is_err(),
"pre-store switch must cause Phase 3 rejection"
);
let err = result.unwrap_err();
assert!(
err.contains("stale") || err.contains("generation") || err.contains("mismatch"),
"error must describe generation mismatch: {err}"
);
}
@@ -892,196 +892,5 @@ mod egress_guard_boundary {
}
}
// ── Phase-boundary seam tests (Area 3 — team) ─────────────────────────────
/// Serializes tests that modify the process-global scope generation counter.
/// See the equivalent comment in `import_tests.rs` for rationale.
use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK as GENERATION_TEST_LOCK;
fn setup_team_import_app_with_scope(
tmp: &tempfile::TempDir,
) -> (tauri::App<tauri::test::MockRuntime>, nostr::Keys) {
use crate::managed_agents::scope::{next_scope_generation, WorkspaceAgentScope};
let owner_keys = nostr::Keys::generate();
let state = crate::app_state::build_app_state();
{
let mut locked = state.keys.lock().unwrap();
*locked = owner_keys.clone();
}
let app = tauri::test::mock_builder()
.manage(state)
.build(tauri::test::mock_context(tauri::test::noop_assets()))
.expect("failed to build mock app for team import test");
{
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,
};
s.commit_active_scope(scope);
}
(app, owner_keys)
}
/// `after_store` hook commits a genuinely different live scope + owner —
/// Phase 4/5 outbound must use the OLD (captured) relay URL, not the new
/// live relay.
///
/// Thufir requirement: `after_store` must commit a genuinely different live
/// 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.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock held to serialize generation
// counter mutations — cannot deadlock. See import_tests.rs for full rationale.
#[allow(clippy::await_holding_lock)]
async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
use crate::commands::personas::snapshot::import::{MemoryPublish, ProfilePublish};
use crate::commands::team_snapshot::confirm_team_snapshot_import_core;
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
use std::sync::{Arc, Mutex};
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
let handle = app.handle();
let snap = snapshot(vec![member("Alice"), member("Bob")]);
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
let input = TeamSnapshotImportConfirm {
file_bytes: encoded,
keep_allowlist: false,
};
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 pr = profile_relays.clone();
let mr = memory_relays.clone();
let state = app.state::<crate::app_state::AppState>();
let handle_for_hook = handle.clone();
let result = confirm_team_snapshot_import_core(
input,
handle,
&state,
|| {},
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 s = handle_for_hook.state::<crate::app_state::AppState>();
s.commit_active_scope(new_scope);
},
move |p: ProfilePublish<'_>| {
let relay = p.relay_url.to_string();
pr.lock().unwrap().push(relay.clone());
Box::pin(async move {
let _ = relay;
Ok(())
})
},
move |m: MemoryPublish<'_>| {
let relay = m.relay_url.to_string();
mr.lock().unwrap().push(relay.clone());
Box::pin(async move {
let _ = relay;
Ok(())
})
},
)
.await;
assert!(
result.is_ok(),
"team import must succeed: {:?}",
result.err()
);
// Every member's profile adapter received the captured relay URL.
let seen = profile_relays.lock().unwrap();
assert_eq!(
seen.len(),
2,
"profile adapter must be called once per member"
);
for relay in seen.iter() {
assert_eq!(
relay, &expected_relay,
"profile adapter must receive captured relay, got: {relay}"
);
}
// No memory entries in these members — memory adapter not called.
}
/// `before_store` hook advances scope generation — Phase 3 must reject BEFORE
/// any write and BEFORE any outbound call.
#[tokio::test]
// SAFETY: single-threaded tokio runtime; lock held to serialize generation
// counter mutations — cannot deadlock. See import_tests.rs for full rationale.
#[allow(clippy::await_holding_lock)]
async fn test_team_identity_switch_before_store_is_rejected() {
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
use crate::commands::personas::snapshot::import::{MemoryPublish, ProfilePublish};
use crate::commands::team_snapshot::confirm_team_snapshot_import_core;
use crate::managed_agents::scope::next_scope_generation;
use tauri::Manager;
let tmp = tempfile::tempdir().unwrap();
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
let handle = app.handle();
let state = app.state::<crate::app_state::AppState>();
let snap = snapshot(vec![member("Alice")]);
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
let input = TeamSnapshotImportConfirm {
file_bytes: encoded,
keep_allowlist: false,
};
let result = confirm_team_snapshot_import_core(
input,
handle,
&state,
move || {
next_scope_generation();
},
|| {},
|_p: ProfilePublish<'_>| {
Box::pin(async { panic!("profile must not be called: store rejected") })
},
|_m: MemoryPublish<'_>| {
Box::pin(async { panic!("memory must not be called: store rejected") })
},
)
.await;
assert!(
result.is_err(),
"pre-store switch must cause Phase 3 rejection"
);
let err = result.unwrap_err();
assert!(
err.contains("stale") || err.contains("generation") || err.contains("mismatch"),
"error must describe generation mismatch: {err}"
);
}
#[path = "seam_tests.rs"]
mod seam_tests;
@@ -0,0 +1,292 @@
//! Concurrency/determinism tests for `managed_agents/runtime_commands.rs`.
//!
//! Split from `runtime_commands_tests.rs` to keep each file under the
//! 1000-line size ratchet. Included via `#[path]` from there as `mod concurrency_tests;`.
//! `use super::*` gives access to all items in `runtime_commands_tests.rs`.
use super::*;
/// Deterministic writer-vs-compensation ordering via held-lock queuing.
///
/// 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).
///
/// 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.
#[test]
fn test_compensate_drain_writer_vs_compensation_deterministic() {
use std::sync::{Arc, Mutex};
use std::thread;
let tmp = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf();
// Seed the store with one agent record.
let pubkey1 = "aa".repeat(32);
let initial_record = crate::managed_agents::ManagedAgentRecord {
pubkey: pubkey1.clone(),
name: "test-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 entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1.clone()];
// 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.
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.
let tmp_wr = tmp_path.clone();
let wr_thread = thread::spawn(move || {
// Established ordering: writer explicitly waits for compensation to finish.
comp_done_rx.recv().unwrap();
let mut records =
crate::managed_agents::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default();
for r in &mut records {
r.env_vars
.insert("WRITER_EDIT".to_string(), "yes".to_string());
}
crate::managed_agents::storage::save_managed_agents_at(&tmp_wr, &records).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 ───────────────────────
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");
// 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.
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"
);
}
/// Joined drain→restore test with a real start-contender blocked on the
/// 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.
///
/// Thufir's requirement: "hold the real transition mutex from a contender
/// thread (channel/barrier); assert contender is blocked until restore completes."
#[test]
fn test_compensate_drain_concurrent_start_is_blocked() {
use std::thread;
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 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();
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(),
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());
// 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()];
// The contender just acquires the transition mutex and reports when it
// managed to do so.
let (contender_ready_tx, contender_ready_rx) = std::sync::mpsc::channel::<()>();
let (contender_done_tx, contender_done_rx) = std::sync::mpsc::channel::<()>();
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.
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).
contender_done_tx.send(()).unwrap();
});
// Wait until contender is alive and parked on (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).
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:?}"
);
// Release the transition guard (compensation done).
drop(transition_guard);
// Contender must now be able to acquire the mutex.
contender_done_rx
.recv_timeout(std::time::Duration::from_secs(5))
.expect("contender must unblock after compensation releases the transition guard");
contender.join().expect("contender thread panicked");
}
@@ -773,284 +773,5 @@ fn test_compensate_drain_attempts_restart_and_reports_degradation_with_real_app(
);
}
/// Deterministic writer-vs-compensation ordering via held-lock queuing.
///
/// 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).
///
/// 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.
#[test]
fn test_compensate_drain_writer_vs_compensation_deterministic() {
use std::sync::{Arc, Mutex};
use std::thread;
let tmp = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf();
// Seed the store with one agent record.
let pubkey1 = "aa".repeat(32);
let initial_record = crate::managed_agents::ManagedAgentRecord {
pubkey: pubkey1.clone(),
name: "test-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![],
};
super::super::storage::save_managed_agents_at(&tmp_path, std::slice::from_ref(&initial_record))
.unwrap();
let entry1 = make_drain_entry(&pubkey1, "wss://relay.example", true);
let stopped = vec![entry1.clone()];
// 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.
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 =
super::super::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.
super::super::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.
let tmp_wr = tmp_path.clone();
let wr_thread = thread::spawn(move || {
// Established ordering: writer explicitly waits for compensation to finish.
comp_done_rx.recv().unwrap();
let mut records =
super::super::storage::load_managed_agents_at(&tmp_wr).unwrap_or_default();
for r in &mut records {
r.env_vars
.insert("WRITER_EDIT".to_string(), "yes".to_string());
}
super::super::storage::save_managed_agents_at(&tmp_wr, &records).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 ───────────────────────
let final_records =
super::super::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");
// 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.
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"
);
}
/// Joined drain→restore test with a real start-contender blocked on the
/// 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.
///
/// Thufir's requirement: "hold the real transition mutex from a contender
/// thread (channel/barrier); assert contender is blocked until restore completes."
#[test]
fn test_compensate_drain_concurrent_start_is_blocked() {
use std::thread;
let tmp = tempfile::tempdir().unwrap();
let tmp_path = tmp.path().to_path_buf();
super::super::storage::save_managed_agents_at(&tmp_path, &[]).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();
use tauri::Manager;
let state = app.state::<crate::app_state::AppState>();
let gen = super::super::scope::current_scope_generation();
let scope = super::super::scope::WorkspaceAgentScope {
scope_id: "test-scope-contender".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());
// 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()];
// The contender just acquires the transition mutex and reports when it
// managed to do so.
let (contender_ready_tx, contender_ready_rx) = std::sync::mpsc::channel::<()>();
let (contender_done_tx, contender_done_rx) = std::sync::mpsc::channel::<()>();
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.
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).
contender_done_tx.send(()).unwrap();
});
// Wait until contender is alive and parked on (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).
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:?}"
);
// Release the transition guard (compensation done).
drop(transition_guard);
// Contender must now be able to acquire the mutex.
contender_done_rx
.recv_timeout(std::time::Duration::from_secs(5))
.expect("contender must unblock after compensation releases the transition guard");
contender.join().expect("contender thread panicked");
}
#[path = "runtime_commands_concurrency_tests.rs"]
mod concurrency_tests;