From 0db0b114e96a133c27f955b9ba17fb692aa28fc9 Mon Sep 17 00:00:00 2001 From: Logan Johnson Date: Fri, 4 Sep 2026 15:59:52 -0400 Subject: [PATCH] feat(desktop): authenticate and durably consume remote Stop Bind owner delegation, installation, agent and community before ordinary Stop. Persist admission before effects and retain exact signed outcomes for retries across reopening and bounded eviction. Signed-off-by: Logan Johnson --- .../src-tauri/src/commands/desktop_stop.rs | 145 +++++++++ desktop/src-tauri/src/commands/mod.rs | 2 + desktop/src-tauri/src/lib.rs | 3 + .../src/managed_agents/remote_stop.rs | 299 +++++++++++++++++- desktop/src/features/agents/AGENTS.md | 11 +- 5 files changed, 453 insertions(+), 7 deletions(-) create mode 100644 desktop/src-tauri/src/commands/desktop_stop.rs diff --git a/desktop/src-tauri/src/commands/desktop_stop.rs b/desktop/src-tauri/src/commands/desktop_stop.rs new file mode 100644 index 00000000000..11ab7fe4e14 --- /dev/null +++ b/desktop/src-tauri/src/commands/desktop_stop.rs @@ -0,0 +1,145 @@ +//! Native owner/host validation and ordinary Stop; no keys cross IPC. +use super::desktop_profiles::{prepare, scope}; +use crate::{ + app_state::AppState, + managed_agents::{self, remote_stop, retention::open_retention_db}, +}; +use buzz_core_pkg::{ + desktop_profile::DesktopProfile, + desktop_stop::{StopOutcome, StopResult, StopTarget}, +}; +use nostr::{Event, JsonUtil, PublicKey}; +use serde_json::{json, Value}; +use tauri::{AppHandle, Manager}; + +fn local_id( + conn: &mut rusqlite::Connection, + scope: &managed_agents::retention::RetentionScope, +) -> Result { + let saved = prepare(conn, scope)?; + let event: Event = serde_json::from_value(saved["event"].clone()).map_err(|e| e.to_string())?; + Ok(DesktopProfile::read( + &event, + &scope.owner_keys, + scope.relay_url.trim_end_matches('/'), + )? + .id) +} + +/// Persist exact signed bytes before the UI sends a new Stop. No boot replay. +#[tauri::command] +pub fn prepare_desktop_stop( + app: AppHandle, + owner: String, + community: String, + desktop: String, + agent: String, +) -> Result { + let state = app.state::(); + let scope = scope(&app, &state, &owner, &community)?; + let event = StopTarget { + v: 1, + community, + desktop, + agent, + } + .sign(&scope.owner_keys)?; + let conn = open_retention_db(&scope.db_path)?; + // Only the current UI operation needs retry bytes; retained receiver fences + // and results are separate. Nothing automatically drains this slot. + conn.execute_batch("CREATE TABLE IF NOT EXISTS desktop_stop_outgoing (slot INTEGER PRIMARY KEY CHECK(slot=1), raw TEXT NOT NULL)") + .map_err(|e| e.to_string())?; + conn.execute("INSERT INTO desktop_stop_outgoing VALUES (1, ?1) ON CONFLICT(slot) DO UPDATE SET raw=excluded.raw", [event.as_json()]) + .map_err(|e| e.to_string())?; + Ok(event) +} + +/// Called only for live owner-private delivery. Reopening never fetches commands. +#[tauri::command] +pub async fn receive_desktop_stop( + app: AppHandle, + owner: String, + community: String, + event: Event, +) -> Result, String> { + tokio::task::spawn_blocking(move || { + let state = app.state::(); + let _transition = state + .managed_agent_runtime_transition + .lock() + .map_err(|e| e.to_string())?; + let scope = scope(&app, &state, &owner, &community)?; + let target = StopTarget::read(&event, &scope.owner_keys, &community)?; + let mut conn = open_retention_db(&scope.db_path)?; + let desktop = local_id(&mut conn, &scope)?; + if desktop != target.desktop { + return Ok(None); + } + // Local possession alone is insufficient after an account switch: + // verify the stored agent's owner delegation against the request author. + let owned = { + let _store = state + .managed_agents_store_lock + .lock() + .map_err(|e| e.to_string())?; + let records = managed_agents::load_managed_agents(&app)?; + records + .iter() + .find(|r| r.pubkey == target.agent) + .is_some_and(|r| { + r.backend == managed_agents::BackendKind::Local + && r.auth_tag + .as_deref() + .and_then(|tag| { + let key = PublicKey::from_hex(&target.agent).ok()?; + buzz_sdk_pkg::nip_oa::verify_auth_tag(tag, &key).ok() + }) + .is_some_and(|key| key.to_hex() == owner) + }) + }; + remote_stop::receive( + &mut conn, + &event, + &scope.owner_keys, + &community, + &desktop, + owned, + |target| { + managed_agents::stop_pair_locked( + target.agent.clone(), + community.clone(), + app.clone(), + ) + .map(|_| ()) + }, + ) + }) + .await + .map_err(|e| format!("Desktop Stop task failed: {e}"))? +} + +/// Result queries never dispatch/replay a request. Missing means Unknown. +#[tauri::command] +pub fn read_desktop_stop_results( + app: AppHandle, + owner: String, + community: String, + request: Event, + events: Vec, +) -> Result { + let state = app.state::(); + let scope = scope(&app, &state, &owner, &community)?; + StopTarget::read(&request, &scope.owner_keys, &community)?; + if events.len() > 16 { + return Err("too many Desktop Stop results".into()); + } + let mut outcome = StopOutcome::Unknown; + for event in events { + let result = StopResult::read(&event, &scope.owner_keys, &request, &community)?; + // Persisted terminal result beats a later Unknown after bounded eviction. + if result.outcome != StopOutcome::Unknown { + outcome = result.outcome; + } + } + Ok(json!(outcome)) +} diff --git a/desktop/src-tauri/src/commands/mod.rs b/desktop/src-tauri/src/commands/mod.rs index 0cb37bf7e02..5a5c2566ac1 100644 --- a/desktop/src-tauri/src/commands/mod.rs +++ b/desktop/src-tauri/src/commands/mod.rs @@ -20,6 +20,7 @@ mod channels; mod clipboard; mod desktop_capabilities; mod desktop_profiles; +mod desktop_stop; mod dms; mod engrams; mod export_util; @@ -96,6 +97,7 @@ pub use channels::*; pub use clipboard::*; pub use desktop_capabilities::*; pub use desktop_profiles::*; +pub use desktop_stop::*; pub use dms::*; pub use engrams::*; pub use global_agent_config::*; diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index f794243354d..2eb43d10670 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -553,6 +553,9 @@ pub fn run() { title_bar_double_click, get_identity, prepare_desktop_profile, + prepare_desktop_stop, + receive_desktop_stop, + read_desktop_stop_results, read_desktop_profiles, prepare_desktop_observation, read_desktop_observations, diff --git a/desktop/src-tauri/src/managed_agents/remote_stop.rs b/desktop/src-tauri/src/managed_agents/remote_stop.rs index 818c7072064..93e3e9b5adc 100644 --- a/desktop/src-tauri/src/managed_agents/remote_stop.rs +++ b/desktop/src-tauri/src/managed_agents/remote_stop.rs @@ -1,15 +1,120 @@ -//! Durable no-auto-start fence shared by ordinary Desktop launch paths. +//! Durable Stop admission. Compact outcomes never compact the per-agent fence. +use buzz_core_pkg::desktop_stop::{StopOutcome, StopResult, StopTarget}; +use nostr::{Event, JsonUtil, Keys}; +use rusqlite::{params, Connection, OptionalExtension, TransactionBehavior}; +use tauri::{AppHandle, Manager}; + use super::retention::{open_retention_db, scoped_retention_db_path}; use super::ManagedAgentRuntimeKey; -use rusqlite::{Connection, OptionalExtension}; -use tauri::{AppHandle, Manager}; + +const HISTORY_LIMIT: i64 = 256; fn schema(conn: &Connection) -> Result<(), String> { conn.execute_batch("CREATE TABLE IF NOT EXISTS desktop_stop_fence ( - agent TEXT PRIMARY KEY, stamp INTEGER NOT NULL, event_id TEXT NOT NULL, blocked INTEGER NOT NULL);") + agent TEXT PRIMARY KEY, stamp INTEGER NOT NULL, event_id TEXT NOT NULL, blocked INTEGER NOT NULL); + CREATE TABLE IF NOT EXISTS desktop_stop_results ( + id TEXT PRIMARY KEY, raw TEXT NOT NULL);") .map_err(|e| e.to_string()) } +/// Persist admission before effect; duplicates/interruption never repeat Stop. +pub(crate) fn admit( + conn: &mut Connection, + request: &Event, + target: &StopTarget, +) -> Result { + schema(conn)?; + let tx = conn + .transaction_with_behavior(TransactionBehavior::Immediate) + .map_err(|e| e.to_string())?; + let previous: Option<(u64, String)> = tx + .query_row( + "SELECT stamp, event_id FROM desktop_stop_fence WHERE agent = ?1", + [&target.agent], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .optional() + .map_err(|e| e.to_string())?; + let id = request.id.to_hex(); + let stamp = request.created_at.as_secs(); + if previous.is_some_and(|(time, key)| time > stamp || (time == stamp && key <= id)) { + return Ok(false); + } + tx.execute("INSERT INTO desktop_stop_fence VALUES (?1, ?2, ?3, 1) + ON CONFLICT(agent) DO UPDATE SET stamp=excluded.stamp, event_id=excluded.event_id, blocked=1", + params![target.agent, stamp, id]).map_err(|e| e.to_string())?; + tx.commit().map_err(|e| e.to_string())?; + Ok(true) +} + +pub(crate) fn saved_result(conn: &Connection, id: &str) -> Result, String> { + schema(conn)?; + conn.query_row( + "SELECT raw FROM desktop_stop_results WHERE id=?1", + [id], + |r| r.get(0), + ) + .optional() + .map_err(|e| e.to_string()) +} + +pub(crate) fn save_result(conn: &mut Connection, id: &str, raw: &str) -> Result<(), String> { + schema(conn)?; + let tx = conn + .transaction_with_behavior(TransactionBehavior::Immediate) + .map_err(|e| e.to_string())?; + tx.execute( + "INSERT OR IGNORE INTO desktop_stop_results VALUES (?1, ?2)", + params![id, raw], + ) + .map_err(|e| e.to_string())?; + tx.execute( + "DELETE FROM desktop_stop_results WHERE rowid NOT IN + (SELECT rowid FROM desktop_stop_results ORDER BY rowid DESC LIMIT ?1)", + [HISTORY_LIMIT], + ) + .map_err(|e| e.to_string())?; + tx.commit().map_err(|e| e.to_string()) +} + +/// Authenticate and durably consume a live request before invoking ordinary Stop. +/// The caller holds the runtime transition lock across this entire operation. +pub(crate) fn receive( + conn: &mut Connection, + request: &Event, + keys: &Keys, + community: &str, + desktop: &str, + owned: bool, + stop: impl FnOnce(&StopTarget) -> Result<(), String>, +) -> Result, String> { + let target = StopTarget::read(request, keys, community)?; + if target.desktop != desktop { + return Ok(None); + } + let id = request.id.to_hex(); + if let Some(raw) = saved_result(conn, &id)? { + let result = Event::from_json(raw).map_err(|e| e.to_string())?; + StopResult::read(&result, keys, request, community)?; + return Ok(Some(result)); + } + let outcome = if !owned { + StopOutcome::Failed + } else if admit(conn, request, &target)? { + outcome(stop(&target)) + } else { + StopOutcome::Unknown + }; + let result = StopResult { + target, + request: id.clone(), + outcome, + } + .sign(keys)?; + save_result(conn, &id, &result.as_json())?; + Ok(Some(result)) +} + /// Explicit local Start captures the Stop fence before its asynchronous preflight. /// Automatic starts and Restart continuations never receive this permission. pub(crate) struct ResumeTicket { @@ -108,9 +213,27 @@ pub(crate) fn finish_resume( Ok(()) } +/// Expose outcomes without mistaking a missing record for success. +pub(crate) fn outcome(stopped: Result<(), String>) -> StopOutcome { + if stopped.is_ok() { + StopOutcome::Stopped + } else { + StopOutcome::Failed + } +} + #[cfg(test)] mod tests { use super::*; + use nostr::{EventBuilder, Keys, Timestamp}; + fn request(keys: &Keys, target: &StopTarget, time: u64) -> Event { + let e = target.sign(keys).unwrap(); + EventBuilder::new(e.kind, e.content) + .tags(e.tags.to_vec()) + .custom_created_at(Timestamp::from(time)) + .sign_with_keys(keys) + .unwrap() + } #[test] fn launch_fence_requires_explicit_start_and_rejects_delayed_preflight() { let stopped = ("stop-a".to_owned(), true); @@ -128,4 +251,172 @@ mod tests { assert!(allow_launch(Some(&stopped), Some(&before_any_stop)).is_err()); assert!(allow_launch(None, Some(&before_any_stop)).is_ok()); } + + #[test] + fn saved_result_is_immutable_and_duplicate_after_interruption_is_unknown() { + let mut conn = Connection::open_in_memory().unwrap(); + let keys = Keys::generate(); + let target = StopTarget { + v: 1, + community: "wss://one.example".into(), + desktop: "a".repeat(32), + agent: Keys::generate().public_key().to_hex(), + }; + let event = request(&keys, &target, 100); + assert!(admit(&mut conn, &event, &target).unwrap()); + // A crash between admission and recording an outcome never reexecutes. + assert!(saved_result(&conn, &event.id.to_hex()).unwrap().is_none()); + assert!(!admit(&mut conn, &event, &target).unwrap()); + save_result(&mut conn, &event.id.to_hex(), "original bytes").unwrap(); + save_result(&mut conn, &event.id.to_hex(), "replacement").unwrap(); + assert_eq!( + saved_result(&conn, &event.id.to_hex()).unwrap().as_deref(), + Some("original bytes") + ); + assert_eq!( + outcome(Err("ordinary Stop failed".into())), + StopOutcome::Failed + ); + assert_eq!(outcome(Ok(())), StopOutcome::Stopped); + } + + #[test] + fn receiver_authenticates_routes_and_returns_exact_saved_result_without_effect() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("receiver.db"); + let mut conn = open_retention_db(&path).unwrap(); + let keys = Keys::generate(); + let target = StopTarget { + v: 1, + community: "wss://one.example".into(), + desktop: "a".repeat(32), + agent: Keys::generate().public_key().to_hex(), + }; + let event = request(&keys, &target, 100); + let no_effect = |_: &StopTarget| panic!("must not invoke ordinary Stop"); + let foreign = Keys::generate(); + for (signer, community, host, rejected) in [ + ( + &foreign, + target.community.as_str(), + target.desktop.as_str(), + true, + ), + (&keys, "wss://other.example", target.desktop.as_str(), true), + ( + &keys, + target.community.as_str(), + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + false, + ), + ] { + let result = receive(&mut conn, &event, signer, community, host, true, no_effect); + if rejected { + assert!(result.is_err()); + } else { + assert!(result.unwrap().is_none()); + } + } + let receive_owned = |conn: &mut Connection, event: &Event, owned, stop| { + receive( + conn, + event, + &keys, + &target.community, + &target.desktop, + owned, + stop, + ) + .unwrap() + .unwrap() + }; + let fail: fn(&StopTarget) -> Result<(), String> = |_| Err("ordinary Stop error".into()); + let no_effect: fn(&StopTarget) -> Result<(), String> = no_effect; + // The first effect succeeds. The reopened retry must return its exact + // signed bytes without invoking the callback at all. + let mut effects = 0; + let result = receive( + &mut conn, + &event, + &keys, + &target.community, + &target.desktop, + true, + |actual| { + assert_eq!(actual, &target); + effects += 1; + Ok(()) + }, + ) + .unwrap() + .unwrap(); + assert_eq!(effects, 1); + let assert_outcome = |result: &Event, request: &Event, expected| { + assert_eq!( + StopResult::read(result, &keys, request, &target.community) + .unwrap() + .outcome, + expected + ); + }; + assert_outcome(&result, &event, StopOutcome::Stopped); + drop(conn); + let mut conn = open_retention_db(&path).unwrap(); + let duplicate = receive_owned(&mut conn, &event, true, no_effect); + assert_eq!(result.as_json(), duplicate.as_json()); + let next = request(&keys, &target, 101); + let failed = receive_owned(&mut conn, &next, true, fail); + assert_outcome(&failed, &next, StopOutcome::Failed); + assert_eq!( + failed.as_json(), + receive_owned(&mut conn, &next, true, no_effect).as_json() + ); + let unowned = request(&keys, &target, 102); + let denied = receive_owned(&mut conn, &unowned, false, no_effect); + assert_outcome(&denied, &unowned, StopOutcome::Failed); + assert_eq!( + conn.query_row( + "SELECT event_id FROM desktop_stop_fence WHERE agent=?1", + [&target.agent], + |r| r.get::<_, String>(0) + ) + .unwrap(), + next.id.to_hex() + ); + let interrupted = request(&keys, &target, 103); + assert!(admit(&mut conn, &interrupted, &target).unwrap()); + let unknown = receive_owned(&mut conn, &interrupted, true, no_effect); + assert_outcome(&unknown, &interrupted, StopOutcome::Unknown); + } + + #[test] + fn durable_fence_survives_outcome_eviction_and_accepts_fresh_stop() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("stop.db"); + let mut conn = open_retention_db(&path).unwrap(); + let keys = Keys::generate(); + let target = StopTarget { + v: 1, + community: "wss://one.example".into(), + desktop: "a".repeat(32), + agent: Keys::generate().public_key().to_hex(), + }; + let first = request(&keys, &target, 100); + assert!(admit(&mut conn, &first, &target).unwrap()); + assert!(!admit(&mut conn, &first, &target).unwrap()); + for i in 0..HISTORY_LIMIT + 2 { + save_result(&mut conn, &format!("{i}"), "result").unwrap(); + } + drop(conn); + let mut conn = open_retention_db(&path).unwrap(); + assert!(!admit(&mut conn, &first, &target).unwrap()); + assert!(admit(&mut conn, &request(&keys, &target, 101), &target).unwrap()); + assert!(!admit(&mut conn, &request(&keys, &target, 99), &target).unwrap()); + let a = request(&keys, &target, 102); + let b = request(&keys, &target, 102); + let (low, high) = if a.id < b.id { (a, b) } else { (b, a) }; + assert!(admit(&mut conn, &high, &target).unwrap()); + assert!(admit(&mut conn, &low, &target).unwrap()); + assert!(!admit(&mut conn, &high, &target).unwrap()); + } } diff --git a/desktop/src/features/agents/AGENTS.md b/desktop/src/features/agents/AGENTS.md index cf0cb2c6524..63ad3ee07a9 100644 --- a/desktop/src/features/agents/AGENTS.md +++ b/desktop/src/features/agents/AGENTS.md @@ -304,13 +304,18 @@ with a TypeScript lookup table or an id comparison in a component. 17. **Databricks model discovery has one shared catalog authority.** Desktop and ACP call the shared `buzz-agent` discovery library; Desktop passes the effective merged `DATABRICKS_MODEL_FILTER` explicitly, and the library applies it to raw workspace endpoint IDs and Unity Catalog model-service FQNs after the additive union. A successful filtered-empty catalog is authoritative: it stays empty, disables switching, and never falls through to configured or known-model fallback. UC FQNs are catalog data and always use the MLflow Chat Completions route, regardless of family-looking text in their components. Global Defaults preserves the discovered model ID as the selected value while its closed trigger renders the provider-scoped display label; do not force the raw persisted ID over that label. -## Desktop Stop launch fence +## Remote Desktop Stop +Native IPC accepts an owner-private, explicitly selected agent+Desktop Stop, +not inferred agent location. The relay redelivers stored Stop duplicates without +repeating relay side effects. +The receiver returns saved results or Unknown, never repeats a consumed Stop. +Native owner-delegation and community checks +precede durable admission and ordinary pair Stop. A delivery ACK is not success. All local spawn paths consume the durable Stop fence at the shared native spawn boundary. Only a deliberate **Start agent** action can supersede that fence; config/restore/reconcile and Restart continuations cannot. Explicit Start -captures its fence before preflight and fails if a newer Stop arrives. Fence -release happens only after the new child has its receipt and tracked handle. +captures its fence before preflight and fails if a newer Stop arrives. ## Channel-only runtime controls