Replace restored pair runtime before respawn

Co-authored-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co>
Signed-off-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co>

Co-authored-by: Tyler Longwell <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf
2026-07-19 11:28:44 -04:00
co-authored by Tyler Longwell
parent 0550db138f
commit 5232c7b3da
4 changed files with 111 additions and 4 deletions
@@ -263,7 +263,7 @@ pub async fn restore_managed_agents_on_launch(
return Ok(());
}
// ── Phase B (no locks): resolve commands and spawn processes in parallel ──
// ── Phase B (transition lock held): resolve commands and spawn in parallel ──
let spawn_results: Vec<AgentSpawnResult> = std::thread::scope(|scope| {
let owner_hex_ref = owner_hex.as_deref();
let handles: Vec<_> = agents_to_start
@@ -281,6 +281,7 @@ pub async fn restore_managed_agents_on_launch(
let result =
super::ManagedAgentRuntimeKey::new(record.pubkey.clone(), &relay_url)
.and_then(|key| {
super::terminate_untracked_pair_runtime(app, &key)?;
spawn_agent_child(app, record, &key.relay_url, false, owner_hex_ref)
.map(|process| (key, process))
});
@@ -420,6 +420,54 @@ pub(crate) fn valid_agent_runtime_receipt(
&& process_has_buzz_marker(receipt.pid, &receipt.desktop_instance_id)
}
fn terminate_runtime_receipt_with(
path: &std::path::Path,
receipt: &super::ManagedAgentRuntimeReceipt,
terminate: impl FnOnce(u32) -> Result<(), String>,
mut is_running: impl FnMut(u32) -> bool,
remove: impl FnOnce(&std::path::Path),
) -> Result<(), String> {
terminate(receipt.pid)?;
for _ in 0..20 {
if !is_running(receipt.pid) {
remove(path);
return Ok(());
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
Err(format!(
"prior runtime {} for pair {} on {} did not exit",
receipt.pid, receipt.key.pubkey, receipt.key.relay_url
))
}
/// Replace a valid prior-session process before registering a new child for
/// the same pair. The caller must hold the runtime transition lock so receipt
/// inspection, termination, spawn, and registration cannot race shutdown or
/// another start.
pub(crate) fn terminate_untracked_pair_runtime(
app: &AppHandle,
key: &ManagedAgentRuntimeKey,
) -> Result<(), String> {
let instance_id = current_instance_id(app);
let Some((path, receipt)) = super::read_all_agent_runtime_receipts(app)
.into_iter()
.find(|(path, receipt)| {
receipt.key == *key && valid_agent_runtime_receipt(path, receipt, &instance_id)
})
else {
return Ok(());
};
terminate_runtime_receipt_with(
&path,
&receipt,
terminate_process,
process_is_running,
super::remove_agent_runtime_receipt_path,
)
}
/// Kill orphaned agent processes using PID file receipts. Reads all files from
/// `agent-pids/`, verifies each PID still belongs to a known agent binary,
/// then resolves each candidate's actual PGID and signals the process group.
@@ -801,3 +801,59 @@ fn receipt_validation_rejects_wrong_pair_filename() {
"test-instance"
));
}
#[test]
fn replacement_removes_receipt_only_after_confirmed_exit() {
use std::cell::{Cell, RefCell};
let receipt = receipt_fixture(
crate::managed_agents::ManagedAgentRuntimeKey::new("aa".repeat(32), "wss://relay.example")
.unwrap(),
);
let path = std::path::Path::new("pair.json");
let terminated = Cell::new(None);
let polls = Cell::new(0);
let removed = RefCell::new(None);
super::terminate_runtime_receipt_with(
path,
&receipt,
|pid| {
terminated.set(Some(pid));
Ok(())
},
|_| {
let poll = polls.get() + 1;
polls.set(poll);
poll < 2
},
|path| *removed.borrow_mut() = Some(path.to_path_buf()),
)
.unwrap();
assert_eq!(terminated.get(), Some(receipt.pid));
assert_eq!(polls.get(), 2);
assert_eq!(removed.into_inner().as_deref(), Some(path));
}
#[test]
fn replacement_failure_keeps_receipt() {
use std::cell::Cell;
let receipt = receipt_fixture(
crate::managed_agents::ManagedAgentRuntimeKey::new("aa".repeat(32), "wss://relay.example")
.unwrap(),
);
let removed = Cell::new(false);
let error = super::terminate_runtime_receipt_with(
std::path::Path::new("pair.json"),
&receipt,
|_| Err("signal failed".into()),
|_| false,
|_| removed.set(true),
)
.unwrap_err();
assert_eq!(error, "signal failed");
assert!(!removed.get());
}
@@ -6,9 +6,10 @@ use super::{
agent_readiness, append_log_marker, current_instance_id, find_managed_agent_mut,
load_global_agent_config, load_managed_agents, load_personas, managed_agent_runtime_log_path,
process_is_running, record_agent_command, resolve_effective_agent_env, save_managed_agents,
spawn_agent_child, terminate_process, write_agent_runtime_receipt, AgentReadiness, BackendKind,
ManagedAgentPairRuntime, ManagedAgentRuntimeKey, ManagedAgentRuntimeLifecycle,
ManagedAgentRuntimeReceipt, ManagedAgentRuntimeStatus,
spawn_agent_child, terminate_process, terminate_untracked_pair_runtime,
write_agent_runtime_receipt, AgentReadiness, BackendKind, ManagedAgentPairRuntime,
ManagedAgentRuntimeKey, ManagedAgentRuntimeLifecycle, ManagedAgentRuntimeReceipt,
ManagedAgentRuntimeStatus,
};
use crate::app_state::AppState;
@@ -168,6 +169,7 @@ fn start_pair(
return Ok(status);
}
runtimes.remove(&key);
terminate_untracked_pair_runtime(&app, &key)?;
let owner = state
.keys