mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
Add pair runtime lifecycle commands
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:
co-authored by
Tyler Longwell
parent
901c2c1308
commit
0e4921cb79
@@ -248,6 +248,12 @@ impl AppState {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn clear_agent_session_cache(&self, key: &ManagedAgentRuntimeKey) {
|
||||
if let Ok(mut map) = self.session_config_cache.lock() {
|
||||
map.remove(key);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn clear_agent_session_caches(&self, pubkey: &str) {
|
||||
if let Ok(mut map) = self.session_config_cache.lock() {
|
||||
map.retain(|key, _| key.pubkey != pubkey);
|
||||
|
||||
@@ -46,7 +46,11 @@ use huddle::{
|
||||
join_huddle, leave_huddle, push_audio_pcm, set_huddle_transcription_enabled, set_tts_enabled,
|
||||
set_voice_input_mode, speak_agent_message, start_huddle, start_stt_pipeline,
|
||||
};
|
||||
use managed_agents::{backfill_persona_snapshots, ensure_nest, try_regenerate_nest};
|
||||
use managed_agents::{
|
||||
backfill_persona_snapshots, ensure_nest, list_managed_agent_runtimes,
|
||||
reconcile_managed_agent_runtimes, restart_managed_agent_runtime, start_managed_agent_runtime,
|
||||
stop_managed_agent_runtime, try_regenerate_nest,
|
||||
};
|
||||
#[cfg(not(feature = "mesh-llm"))]
|
||||
use mesh_llm_stubs::*;
|
||||
#[cfg(all(feature = "mesh-llm", target_os = "macos"))]
|
||||
@@ -857,6 +861,11 @@ pub fn run() {
|
||||
resolve_oa_owner,
|
||||
list_relay_agents,
|
||||
list_managed_agents,
|
||||
list_managed_agent_runtimes,
|
||||
start_managed_agent_runtime,
|
||||
stop_managed_agent_runtime,
|
||||
restart_managed_agent_runtime,
|
||||
reconcile_managed_agent_runtimes,
|
||||
create_managed_agent,
|
||||
start_managed_agent,
|
||||
stop_managed_agent,
|
||||
|
||||
@@ -26,6 +26,7 @@ mod restore;
|
||||
pub mod retention;
|
||||
mod runtime;
|
||||
mod runtime_types;
|
||||
mod runtime_commands;
|
||||
pub(crate) mod spawn_hash;
|
||||
pub(crate) mod storage;
|
||||
pub(crate) mod team_events;
|
||||
@@ -68,6 +69,7 @@ pub use repos::{
|
||||
pub use restore::*;
|
||||
pub use runtime::*;
|
||||
pub use runtime_types::*;
|
||||
pub use runtime_commands::*;
|
||||
pub use storage::*;
|
||||
pub use teams::*;
|
||||
pub use types::*;
|
||||
|
||||
@@ -114,10 +114,25 @@ pub async fn restore_managed_agents_on_launch(
|
||||
changed |=
|
||||
kill_stale_tracked_processes(&mut records, &runtimes, &super::current_instance_id(app));
|
||||
|
||||
let tracked_pids: Vec<u32> = records
|
||||
.iter()
|
||||
.filter_map(|r| r.runtime_pid)
|
||||
.chain(runtimes.values().map(|rt| rt.child.id()))
|
||||
let tracked_pids: Vec<u32> = runtimes
|
||||
.values()
|
||||
.map(|runtime| runtime.child.id())
|
||||
.chain(
|
||||
super::read_all_agent_runtime_receipts(app)
|
||||
.into_iter()
|
||||
.filter_map(|(_, receipt)| {
|
||||
let canonical = super::ManagedAgentRuntimeKey::new(
|
||||
receipt.key.pubkey.clone(),
|
||||
&receipt.key.relay_url,
|
||||
)
|
||||
.ok()?;
|
||||
(canonical == receipt.key
|
||||
&& receipt.desktop_instance_id == super::current_instance_id(app)
|
||||
&& super::process_is_running(receipt.pid)
|
||||
&& super::process_belongs_to_us(receipt.pid))
|
||||
.then_some(receipt.pid)
|
||||
}),
|
||||
)
|
||||
.collect();
|
||||
super::sweep_orphaned_agent_processes(app, &tracked_pids);
|
||||
|
||||
|
||||
@@ -409,20 +409,46 @@ fn resolve_pgids_and_kill(candidate_pids: &[i32]) {
|
||||
/// `skip_pids` are PIDs already handled by the tracked-agent path.
|
||||
#[cfg(unix)]
|
||||
pub(crate) fn sweep_orphaned_agent_processes(app: &AppHandle, skip_pids: &[u32]) {
|
||||
let entries = super::read_all_agent_pid_files(app);
|
||||
let legacy_entries = super::read_all_agent_pid_files(app);
|
||||
let instance_id = current_instance_id(app);
|
||||
let receipt_entries: Vec<_> = super::read_all_agent_runtime_receipts(app)
|
||||
.into_iter()
|
||||
.filter_map(|(path, receipt)| {
|
||||
let canonical = ManagedAgentRuntimeKey::new(
|
||||
receipt.key.pubkey.clone(),
|
||||
&receipt.key.relay_url,
|
||||
)
|
||||
.ok();
|
||||
let valid = canonical.as_ref() == Some(&receipt.key)
|
||||
&& path.file_name().and_then(|name| name.to_str())
|
||||
== Some(&format!("{}.json", receipt.key.runtime_id()))
|
||||
&& receipt.desktop_instance_id == instance_id
|
||||
&& process_is_running(receipt.pid)
|
||||
&& process_belongs_to_us(receipt.pid)
|
||||
&& process_has_buzz_marker(receipt.pid, &receipt.desktop_instance_id);
|
||||
if valid {
|
||||
Some((path, receipt))
|
||||
} else {
|
||||
super::remove_agent_runtime_receipt_path(&path);
|
||||
None
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
// Collect live orphans AND dead-leader groups into a single kill batch.
|
||||
// Dead leaders: PGID may have been recycled, but the window is narrow
|
||||
// (PID files are from this session) and the cost of missing surviving
|
||||
// group members outweighs the recycling risk.
|
||||
let targets: Vec<i32> = entries
|
||||
let targets: Vec<i32> = legacy_entries
|
||||
.iter()
|
||||
.filter(|(_, pid)| {
|
||||
.map(|(_, pid)| *pid)
|
||||
.chain(receipt_entries.iter().map(|(_, receipt)| receipt.pid))
|
||||
.filter(|pid| {
|
||||
if skip_pids.contains(pid) {
|
||||
return false;
|
||||
}
|
||||
(process_is_running(*pid) && process_belongs_to_us(*pid)) || !process_is_running(*pid)
|
||||
})
|
||||
.map(|(_, pid)| *pid as i32)
|
||||
.map(|pid| pid as i32)
|
||||
.collect();
|
||||
|
||||
if !targets.is_empty() {
|
||||
@@ -430,7 +456,7 @@ pub(crate) fn sweep_orphaned_agent_processes(app: &AppHandle, skip_pids: &[u32])
|
||||
}
|
||||
|
||||
// Clean up PID files for processes we just killed or that are already gone.
|
||||
for (pubkey, pid) in &entries {
|
||||
for (pubkey, pid) in &legacy_entries {
|
||||
if skip_pids.contains(pid) {
|
||||
continue;
|
||||
}
|
||||
@@ -438,6 +464,14 @@ pub(crate) fn sweep_orphaned_agent_processes(app: &AppHandle, skip_pids: &[u32])
|
||||
super::remove_agent_pid_file(app, pubkey);
|
||||
}
|
||||
}
|
||||
for (_, receipt) in &receipt_entries {
|
||||
if skip_pids.contains(&receipt.pid) {
|
||||
continue;
|
||||
}
|
||||
if !process_is_running(receipt.pid) || !process_belongs_to_us(receipt.pid) {
|
||||
super::remove_agent_runtime_receipt(app, &receipt.key);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
|
||||
@@ -0,0 +1,201 @@
|
||||
use std::sync::atomic::Ordering;
|
||||
|
||||
use tauri::{AppHandle, Emitter, Manager};
|
||||
|
||||
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,
|
||||
};
|
||||
use crate::app_state::AppState;
|
||||
|
||||
const STATUS_EVENT: &str = "managed-agent-runtime-status";
|
||||
|
||||
fn status_for(
|
||||
app: &AppHandle,
|
||||
record: &super::ManagedAgentRecord,
|
||||
key: &ManagedAgentRuntimeKey,
|
||||
runtime: Option<&ManagedAgentPairRuntime>,
|
||||
requested_relay_url: Option<String>,
|
||||
) -> ManagedAgentRuntimeStatus {
|
||||
let personas = load_personas(app).unwrap_or_default();
|
||||
let global = load_global_agent_config(app).unwrap_or_default();
|
||||
let command = record_agent_command(record, &personas);
|
||||
let metadata = super::known_acp_runtime(&command);
|
||||
let effective = resolve_effective_agent_env(record, &personas, metadata, &global);
|
||||
let local_setup = matches!(agent_readiness(&effective), AgentReadiness::Ready);
|
||||
ManagedAgentRuntimeStatus {
|
||||
pubkey: key.pubkey.clone(),
|
||||
relay_url: key.relay_url.clone(),
|
||||
requested_relay_url,
|
||||
local_setup,
|
||||
lifecycle: runtime
|
||||
.map(|runtime| runtime.lifecycle.clone())
|
||||
.unwrap_or(ManagedAgentRuntimeLifecycle::Stopped),
|
||||
pid: runtime.map(|runtime| runtime.child.id()),
|
||||
error: runtime.and_then(|runtime| runtime.error.clone()),
|
||||
log_path: managed_agent_runtime_log_path(app, key)
|
||||
.ok()
|
||||
.map(|path| path.display().to_string()),
|
||||
}
|
||||
}
|
||||
|
||||
fn emit_status(app: &AppHandle, status: &ManagedAgentRuntimeStatus) {
|
||||
let _ = app.emit(STATUS_EVENT, status);
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn list_managed_agent_runtimes(app: AppHandle) -> Result<Vec<ManagedAgentRuntimeStatus>, String> {
|
||||
let state = app.state::<AppState>();
|
||||
let records = load_managed_agents(&app)?;
|
||||
let runtimes = state.managed_agent_processes.lock().map_err(|e| e.to_string())?;
|
||||
Ok(runtimes
|
||||
.iter()
|
||||
.filter_map(|(key, runtime)| {
|
||||
let record = records.iter().find(|record| record.pubkey == key.pubkey)?;
|
||||
Some(status_for(&app, record, key, Some(runtime), None))
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn start_managed_agent_runtime(
|
||||
pubkey: String,
|
||||
relay_url: String,
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
start_pair(pubkey, relay_url, true, app)
|
||||
}
|
||||
|
||||
fn start_pair(
|
||||
pubkey: String,
|
||||
relay_url: String,
|
||||
lazy: bool,
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
let state = app.state::<AppState>();
|
||||
let _transition = state
|
||||
.managed_agent_runtime_transition
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
if state.shutdown_started.load(Ordering::Acquire) {
|
||||
return Err("desktop shutdown has started".into());
|
||||
}
|
||||
let _store = state.managed_agents_store_lock.lock().map_err(|e| e.to_string())?;
|
||||
let mut records = load_managed_agents(&app)?;
|
||||
let record = find_managed_agent_mut(&mut records, &pubkey)?;
|
||||
if record.backend != BackendKind::Local {
|
||||
return Err("managed runtime pairs require a local agent".into());
|
||||
}
|
||||
let key = ManagedAgentRuntimeKey::new(pubkey, &relay_url)?;
|
||||
let mut runtimes = state.managed_agent_processes.lock().map_err(|e| e.to_string())?;
|
||||
if runtimes.get_mut(&key).is_some_and(|runtime| {
|
||||
runtime.child.try_wait().ok().flatten().is_none()
|
||||
}) {
|
||||
let status = status_for(&app, record, &key, runtimes.get(&key), None);
|
||||
return Ok(status);
|
||||
}
|
||||
runtimes.remove(&key);
|
||||
|
||||
let owner = state.keys.lock().ok().map(|keys| keys.public_key().to_hex());
|
||||
let mut process = spawn_agent_child(&app, record, &key.relay_url, lazy, owner.as_deref())?;
|
||||
let now = crate::util::now_iso();
|
||||
let receipt = ManagedAgentRuntimeReceipt {
|
||||
key: key.clone(),
|
||||
pid: process.child.id(),
|
||||
desktop_instance_id: current_instance_id(&app),
|
||||
started_at: now.clone(),
|
||||
};
|
||||
if let Err(error) = write_agent_runtime_receipt(&app, &receipt) {
|
||||
let _ = terminate_process(process.child.id());
|
||||
let _ = process.child.wait();
|
||||
return Err(error);
|
||||
}
|
||||
record.runtime_pid = None;
|
||||
record.updated_at = now.clone();
|
||||
record.last_started_at = Some(now);
|
||||
record.last_stopped_at = None;
|
||||
record.last_error = None;
|
||||
runtimes.insert(key.clone(), ManagedAgentPairRuntime::starting(process));
|
||||
let status = status_for(&app, record, &key, runtimes.get(&key), None);
|
||||
drop(runtimes);
|
||||
save_managed_agents(&app, &records)?;
|
||||
emit_status(&app, &status);
|
||||
Ok(status)
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn stop_managed_agent_runtime(
|
||||
pubkey: String,
|
||||
relay_url: String,
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
let state = app.state::<AppState>();
|
||||
let _transition = state
|
||||
.managed_agent_runtime_transition
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
let _store = state.managed_agents_store_lock.lock().map_err(|e| e.to_string())?;
|
||||
let mut records = load_managed_agents(&app)?;
|
||||
let record = find_managed_agent_mut(&mut records, &pubkey)?;
|
||||
let key = ManagedAgentRuntimeKey::new(pubkey, &relay_url)?;
|
||||
let mut runtimes = state.managed_agent_processes.lock().map_err(|e| e.to_string())?;
|
||||
if let Some(mut runtime) = runtimes.remove(&key) {
|
||||
if process_is_running(runtime.child.id()) {
|
||||
terminate_process(runtime.child.id())?;
|
||||
}
|
||||
let status = runtime.child.wait().map_err(|e| e.to_string())?;
|
||||
record.last_exit_code = status.code();
|
||||
let _ = append_log_marker(&runtime.log_path, "=== stopped pair runtime ===");
|
||||
}
|
||||
super::remove_agent_runtime_receipt(&app, &key);
|
||||
state.clear_agent_session_cache(&key);
|
||||
record.runtime_pid = None;
|
||||
record.updated_at = crate::util::now_iso();
|
||||
record.last_stopped_at = Some(record.updated_at.clone());
|
||||
let status = status_for(&app, record, &key, None, None);
|
||||
drop(runtimes);
|
||||
save_managed_agents(&app, &records)?;
|
||||
emit_status(&app, &status);
|
||||
Ok(status)
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn restart_managed_agent_runtime(
|
||||
pubkey: String,
|
||||
relay_url: String,
|
||||
app: AppHandle,
|
||||
) -> Result<ManagedAgentRuntimeStatus, String> {
|
||||
stop_managed_agent_runtime(pubkey.clone(), relay_url.clone(), app.clone())?;
|
||||
start_pair(pubkey, relay_url, true, app)
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn reconcile_managed_agent_runtimes(
|
||||
communities: Vec<super::ManagedAgentCommunityTarget>,
|
||||
app: AppHandle,
|
||||
) -> Result<Vec<ManagedAgentRuntimeStatus>, String> {
|
||||
// Discovery is Rust-owned. Until membership has been authenticated, never
|
||||
// infer an empty membership set as success: return explicit failed rows.
|
||||
let records = load_managed_agents(&app)?;
|
||||
let mut rows = Vec::new();
|
||||
for community in communities {
|
||||
for record in records.iter().filter(|record| record.backend == BackendKind::Local) {
|
||||
let key = ManagedAgentRuntimeKey::new(record.pubkey.clone(), &community.relay_url)?;
|
||||
let mut row = status_for(
|
||||
&app,
|
||||
record,
|
||||
&key,
|
||||
None,
|
||||
Some(community.relay_url.clone()),
|
||||
);
|
||||
row.lifecycle = ManagedAgentRuntimeLifecycle::Failed;
|
||||
row.error = Some("authenticated membership discovery is not yet available".into());
|
||||
rows.push(row);
|
||||
}
|
||||
}
|
||||
Ok(rows)
|
||||
}
|
||||
@@ -628,7 +628,13 @@ pub fn remove_agent_runtime_receipt(app: &AppHandle, key: &ManagedAgentRuntimeKe
|
||||
}
|
||||
}
|
||||
|
||||
pub fn read_all_agent_runtime_receipts(app: &AppHandle) -> Vec<ManagedAgentRuntimeReceipt> {
|
||||
pub fn remove_agent_runtime_receipt_path(path: &Path) {
|
||||
let _ = fs::remove_file(path);
|
||||
}
|
||||
|
||||
pub fn read_all_agent_runtime_receipts(
|
||||
app: &AppHandle,
|
||||
) -> Vec<(PathBuf, ManagedAgentRuntimeReceipt)> {
|
||||
let Ok(dir) = agent_pids_dir(app) else {
|
||||
return Vec::new();
|
||||
};
|
||||
@@ -638,8 +644,13 @@ pub fn read_all_agent_runtime_receipts(app: &AppHandle) -> Vec<ManagedAgentRunti
|
||||
entries
|
||||
.flatten()
|
||||
.filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
|
||||
.filter_map(|entry| fs::read(entry.path()).ok())
|
||||
.filter_map(|bytes| serde_json::from_slice(&bytes).ok())
|
||||
.filter_map(|entry| {
|
||||
let path = entry.path();
|
||||
let bytes = fs::read(&path).ok()?;
|
||||
serde_json::from_slice(&bytes)
|
||||
.ok()
|
||||
.map(|receipt| (path, receipt))
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user