diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 68b2df4d607..5acc815834c 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -114,6 +114,56 @@ fn emit_runtime_lifecycle( } } +fn startup_runtime_lifecycle(lazy_pool: bool) -> &'static str { + if lazy_pool { + "listening" + } else { + "ready" + } +} + +fn emit_startup_runtime_lifecycle( + observer: Option<&observer::ObserverHandle>, + start_nonce: &str, + pubkey: &str, + relay_url: &str, + lazy_pool: bool, +) { + emit_runtime_lifecycle( + observer, + start_nonce, + pubkey, + relay_url, + startup_runtime_lifecycle(lazy_pool), + None, + ); +} + +#[cfg(test)] +mod runtime_lifecycle_tests { + use crate::observer::ObserverHandle; + + #[test] + fn startup_lifecycle_emits_once_for_each_pool_mode() { + for (lazy_pool, expected) in [(false, "ready"), (true, "listening")] { + let observer = ObserverHandle::in_process(); + + super::emit_startup_runtime_lifecycle( + Some(&observer), + "nonce", + "pubkey", + "ws://localhost:3000", + lazy_pool, + ); + + let events = observer.snapshot(); + assert_eq!(events.len(), 1); + assert_eq!(events[0].kind, "managed_agent_runtime_lifecycle"); + assert_eq!(events[0].payload["lifecycle"], expected); + } + } +} + /// Resolve the agent's owner pubkey at startup. /// /// Priority: @@ -2174,16 +2224,13 @@ async fn tokio_main() -> Result<()> { } } - if config.lazy_pool { - emit_runtime_lifecycle( - observer.as_ref(), - &runtime_start_nonce, - &pubkey_hex, - &config.relay_url, - "listening", - None, - ); - } + emit_startup_runtime_lifecycle( + observer.as_ref(), + &runtime_start_nonce, + &pubkey_hex, + &config.relay_url, + config.lazy_pool, + ); let base_prompt_content = config.base_prompt_content.take(); let ctx = Arc::new(PromptContext { diff --git a/desktop/src/features/agents/observerEventOrder.ts b/desktop/src/features/agents/observerEventOrder.ts new file mode 100644 index 00000000000..59460b16be8 --- /dev/null +++ b/desktop/src/features/agents/observerEventOrder.ts @@ -0,0 +1,37 @@ +import type { ObserverEvent } from "./ui/agentSessionTypes"; + +export function compareObserverEvents( + left: ObserverEvent, + right: ObserverEvent, +) { + const leftTime = Date.parse(left.timestamp); + const rightTime = Date.parse(right.timestamp); + if (Number.isFinite(leftTime) && Number.isFinite(rightTime)) { + const timeDiff = leftTime - rightTime; + if (timeDiff !== 0) { + return timeDiff; + } + } + + return left.seq - right.seq; +} + +/** + * Returns true if `candidate` sorts strictly after `stored` using the same + * two-key ordering as `compareObserverEvents`: later timestamp wins; equal + * timestamp falls back to higher seq. Extracted so latest-live advancement + * cannot drift from transcript ordering. + */ +export function isObserverEventAfter( + candidate: { timestamp: string; seq: number }, + stored: { timestamp: string; seq: number }, +): boolean { + const candidateTime = Date.parse(candidate.timestamp); + const storedTime = Date.parse(stored.timestamp); + if (Number.isFinite(candidateTime) && Number.isFinite(storedTime)) { + if (candidateTime !== storedTime) { + return candidateTime > storedTime; + } + } + return candidate.seq > stored.seq; +} diff --git a/desktop/src/features/agents/observerLifecycleSeq.ts b/desktop/src/features/agents/observerLifecycleSeq.ts new file mode 100644 index 00000000000..83d3b630d42 --- /dev/null +++ b/desktop/src/features/agents/observerLifecycleSeq.ts @@ -0,0 +1,42 @@ +import { normalizePubkey } from "@/shared/lib/pubkey"; +import type { ObserverEvent } from "./ui/agentSessionTypes"; + +const appliedLifecycleSeqByPairNonce = new Map(); + +function lifecycleSequenceKey( + agentPubkey: string, + payload: unknown, +): string | null { + if (payload === null || typeof payload !== "object") return null; + const record = payload as { relayUrl?: unknown; startNonce?: unknown }; + if ( + typeof record.relayUrl !== "string" || + typeof record.startNonce !== "string" + ) { + return null; + } + return `${normalizePubkey(agentPubkey)}\0${record.relayUrl}\0${record.startNonce}`; +} + +export function shouldApplyLifecycleFrame( + agentPubkey: string, + event: ObserverEvent, +): boolean { + const key = lifecycleSequenceKey(agentPubkey, event.payload); + if (key === null) return false; + const applied = appliedLifecycleSeqByPairNonce.get(key); + return applied === undefined || event.seq > applied; +} + +export function recordAppliedLifecycleFrame( + agentPubkey: string, + event: ObserverEvent, +): void { + const key = lifecycleSequenceKey(agentPubkey, event.payload); + if (key === null) return; + appliedLifecycleSeqByPairNonce.set(key, event.seq); +} + +export function resetAppliedLifecycleSeq(): void { + appliedLifecycleSeqByPairNonce.clear(); +} diff --git a/desktop/src/features/agents/observerRelayStore.test.mjs b/desktop/src/features/agents/observerRelayStore.test.mjs new file mode 100644 index 00000000000..9a14404acab --- /dev/null +++ b/desktop/src/features/agents/observerRelayStore.test.mjs @@ -0,0 +1,152 @@ +import assert from "node:assert/strict"; +import { afterEach, beforeEach, test } from "node:test"; + +import { + _testProcessLiveObserverEvents, + getAgentObserverSnapshot, + resetAgentObserverStore, +} from "./observerRelayStore.ts"; + +const AGENT = "a".repeat(64); + +function lifecycleEvent(seq, lifecycle, startNonce = "gen-1") { + return { + seq, + timestamp: `2026-08-15T00:00:0${seq}Z`, + kind: "managed_agent_runtime_lifecycle", + agentIndex: 0, + channelId: null, + sessionId: null, + turnId: null, + payload: { + relayUrl: "ws://localhost:3000", + startNonce, + lifecycle, + pid: 123, + }, + }; +} + +beforeEach(() => { + resetAgentObserverStore(); +}); + +afterEach(() => { + delete globalThis.__TAURI_INTERNALS__; + delete globalThis.window; +}); + +test("applies lifecycle frames in observer order", async () => { + const started = []; + const completed = []; + let releaseWaking; + + const tauriInternals = { + invoke: (command, args) => { + assert.equal(command, "put_managed_agent_runtime_lifecycle"); + const lifecycle = args.payload.lifecycle; + started.push(lifecycle); + if (lifecycle === "waking") { + return new Promise((resolve) => { + releaseWaking = () => { + completed.push(lifecycle); + resolve({}); + }; + }); + } + completed.push(lifecycle); + return Promise.resolve({}); + }, + }; + globalThis.__TAURI_INTERNALS__ = tauriInternals; + globalThis.window = { __TAURI_INTERNALS__: tauriInternals }; + + const processing = _testProcessLiveObserverEvents(AGENT, [ + lifecycleEvent(1, "waking"), + lifecycleEvent(2, "ready"), + ]); + await Promise.resolve(); + + assert.deepEqual( + started, + ["waking"], + "ready must wait for the preceding waking write", + ); + releaseWaking(); + await processing; + + assert.deepEqual(started, ["waking", "ready"]); + assert.deepEqual(completed, ["waking", "ready"]); +}); + +test("drops the remainder of a lifecycle batch after a store reset", async () => { + const started = []; + let releaseWaking; + + const tauriInternals = { + invoke: (_command, args) => { + const lifecycle = args.payload.lifecycle; + started.push(lifecycle); + if (lifecycle === "waking") { + return new Promise((resolve) => { + releaseWaking = resolve; + }); + } + return Promise.resolve({}); + }, + }; + globalThis.__TAURI_INTERNALS__ = tauriInternals; + globalThis.window = { __TAURI_INTERNALS__: tauriInternals }; + + const processing = _testProcessLiveObserverEvents(AGENT, [ + lifecycleEvent(1, "waking"), + lifecycleEvent(2, "ready"), + ]); + await Promise.resolve(); + + resetAgentObserverStore(); + releaseWaking(); + await processing; + + assert.deepEqual(started, ["waking"]); + assert.deepEqual(getAgentObserverSnapshot(AGENT).events, []); +}); + +test("newest-first replay does not regress ready to stale waking", async () => { + const started = []; + const tauriInternals = { + invoke: (command, args) => { + assert.equal(command, "put_managed_agent_runtime_lifecycle"); + started.push(args.payload.lifecycle); + return Promise.resolve({}); + }, + }; + globalThis.__TAURI_INTERNALS__ = tauriInternals; + globalThis.window = { __TAURI_INTERNALS__: tauriInternals }; + + await _testProcessLiveObserverEvents(AGENT, [lifecycleEvent(2, "ready")]); + await _testProcessLiveObserverEvents(AGENT, [lifecycleEvent(1, "waking")]); + + assert.deepEqual(started, ["ready"]); +}); + +test("a new startNonce restarts the lifecycle sequence domain", async () => { + const started = []; + const tauriInternals = { + invoke: (_command, args) => { + started.push(`${args.payload.startNonce}:${args.payload.lifecycle}`); + return Promise.resolve({}); + }, + }; + globalThis.__TAURI_INTERNALS__ = tauriInternals; + globalThis.window = { __TAURI_INTERNALS__: tauriInternals }; + + await _testProcessLiveObserverEvents(AGENT, [ + lifecycleEvent(2, "ready", "gen-1"), + ]); + await _testProcessLiveObserverEvents(AGENT, [ + lifecycleEvent(1, "waking", "gen-2"), + ]); + + assert.deepEqual(started, ["gen-1:ready", "gen-2:waking"]); +}); diff --git a/desktop/src/features/agents/observerRelayStore.ts b/desktop/src/features/agents/observerRelayStore.ts index 68fa290ad25..a3b5fdb9e43 100644 --- a/desktop/src/features/agents/observerRelayStore.ts +++ b/desktop/src/features/agents/observerRelayStore.ts @@ -25,6 +25,20 @@ import { createEmptyTranscriptState, processTranscriptEvent, } from "./ui/agentSessionTranscript"; +import { + compareObserverEvents, + isObserverEventAfter, +} from "./observerEventOrder"; +import { + recordAppliedLifecycleFrame, + resetAppliedLifecycleSeq, + shouldApplyLifecycleFrame, +} from "./observerLifecycleSeq"; + +export { + compareObserverEvents, + isObserverEventAfter, +} from "./observerEventOrder"; const MAX_OBSERVER_EVENTS = 3000; // Length the per-agent journal is evicted down to when it overflows @@ -110,7 +124,6 @@ function liveSessionKey(agentPubkey: string, channelId: string | null): string { return `${normalizePubkey(agentPubkey)}:${channelId ?? ""}`; } -/** Read the latest-live-session-id for a (agent, channel) pair. */ export function getLatestLiveSessionId( agentPubkey: string | null | undefined, channelId: string | null | undefined, @@ -400,42 +413,6 @@ export function getArchivedChannelEvents( ); } -export function compareObserverEvents( - left: ObserverEvent, - right: ObserverEvent, -) { - const leftTime = Date.parse(left.timestamp); - const rightTime = Date.parse(right.timestamp); - if (Number.isFinite(leftTime) && Number.isFinite(rightTime)) { - const timeDiff = leftTime - rightTime; - if (timeDiff !== 0) { - return timeDiff; - } - } - - return left.seq - right.seq; -} - -/** - * Returns true if `candidate` sorts strictly after `stored` using the same - * two-key ordering as `compareObserverEvents`: later timestamp wins; equal - * timestamp falls back to higher seq. Extracted so latest-live advancement - * cannot drift from transcript ordering. - */ -export function isObserverEventAfter( - candidate: { timestamp: string; seq: number }, - stored: { timestamp: string; seq: number }, -): boolean { - const candidateTime = Date.parse(candidate.timestamp); - const storedTime = Date.parse(stored.timestamp); - if (Number.isFinite(candidateTime) && Number.isFinite(storedTime)) { - if (candidateTime !== storedTime) { - return candidateTime > storedTime; - } - } - return candidate.seq > stored.seq; -} - // Observer event kind for a batch envelope wrapping multiple events. The ACP // harness publishes one frame per second; everything that accumulated between // ticks arrives as `{ kind: "batch", payload: { events: [...] } }` with every @@ -460,10 +437,13 @@ function unwrapObserverBatch(parsed: ObserverEvent): ObserverEvent[] { // Per-event processing shared by every event a live frame carries (one for a // plain frame, many for a batch envelope). -function processLiveObserverEvents( +async function processLiveObserverEvents( agentPubkey: string, events: readonly ObserverEvent[], + activeGeneration: number, ) { + if (activeGeneration !== generation) return; + // Commit the full envelope before dispatching synchronous specialized // callbacks. Those callbacks historically observed their triggering frame // in the raw/transcript stores; batching must preserve that visibility while @@ -514,17 +494,21 @@ function processLiveObserverEvents( // count a terminal switch result once per distinct channel. dispatchControlResult(agentPubkey, parsed.payload, parsed.channelId); } else if (parsed.kind === "managed_agent_runtime_lifecycle") { - void putManagedAgentRuntimeLifecycle(agentPubkey, parsed.payload).catch( - (error) => { - console.debug("Late/untracked lifecycle frame dropped:", error); - }, - ); + if (!shouldApplyLifecycleFrame(agentPubkey, parsed)) continue; + try { + await putManagedAgentRuntimeLifecycle(agentPubkey, parsed.payload); + if (activeGeneration !== generation) return; + recordAppliedLifecycleFrame(agentPubkey, parsed); + } catch (error) { + if (activeGeneration !== generation) return; + console.debug("Late/untracked lifecycle frame dropped:", error); + } } } // Preserve the harness's envelope backpressure: retained state was committed // before specialized callbacks, but external-store subscribers publish once. - if (accepted) { + if (accepted && activeGeneration === generation) { notifyListeners({ agentPubkey, events: accepted }); } } @@ -563,7 +547,11 @@ async function handleRelayObserverEvent( if (activeGeneration !== generation) { return; } - processLiveObserverEvents(agentPubkey, unwrapObserverBatch(parsed)); + await processLiveObserverEvents( + agentPubkey, + unwrapObserverBatch(parsed), + activeGeneration, + ); } catch (error) { if (activeGeneration !== generation) { return; @@ -918,6 +906,7 @@ export function resetAgentObserverStore() { knownAgentsBySubscription.clear(); pendingUnknownAgentFrames.length = 0; latestLiveSessionByAgentChannel.clear(); + resetAppliedLifecycleSeq(); agentManagementListeners.clear(); onSessionConfigCaptured = null; connectionState = "idle"; @@ -942,8 +931,8 @@ export function _testRegisterKnownAgents( export function _testProcessLiveObserverEvents( agentPubkey: string, events: readonly ObserverEvent[], -): void { - processLiveObserverEvents(agentPubkey, events); +): Promise { + return processLiveObserverEvents(agentPubkey, events, generation); } /**