From 033416546c5f116e540e410a89922a33cfdca2c2 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Tue, 15 Sep 2026 11:32:12 -0700 Subject: [PATCH] feat(agent-status): compose host replicas across SSH --- .../hook-status-session-tabs-republish.ts | 24 +- src/main/agent-hooks/server/server-cleanup.ts | 5 +- src/main/ipc/agent-hooks.ts | 9 +- src/main/orcad/orcad-entry.ts | 22 +- .../agent-status-host-replica-store.test.ts | 201 +++++++++++++++ .../agent-status-host-replica-store.ts | 244 ++++++++++++++++++ .../orca-runtime-preserved-branch-cleanup.ts | 2 + src/main/runtime/orca-runtime-state-fields.ts | 11 + .../rpc/methods/agent-status-store.test.ts | 91 +++++++ .../runtime/rpc/methods/agent-status-store.ts | 34 ++- .../runtime/rpc/rpc-params-type-parity.ts | 7 +- .../runtime-mobile-agent-status-projection.ts | 12 +- .../runtime-rpc-mobile-method-allowlist.ts | 2 + .../ssh-relay-agent-status-test-support.ts | 44 ++++ src/main/ssh/ssh-relay-deploy.ts | 18 +- ...ay-session-agent-hooks.integration.test.ts | 150 ++++++++++- src/main/ssh/ssh-relay-session.ts | 198 +++++++++++++- .../startup/main-process-runtime-service.ts | 19 +- src/main/startup/main-window-agent-status.ts | 39 +++ src/relay/agent-hook-cache-actions.ts | 55 ++++ src/relay/agent-hook-http-handler.ts | 85 ++++++ src/relay/agent-hook-server.test.ts | 38 +++ src/relay/agent-hook-server.ts | 141 ++++------ src/relay/agent-hook-status-store-source.ts | 99 +++++++ src/relay/relay-agent-hook-runtime.ts | 84 ++++-- src/relay/relay-agent-status-store.ts | 82 ------ src/relay/relay-daemon.ts | 3 +- src/relay/relay-launch-options.ts | 10 +- .../agent-status-event-applicator.ts | 3 + ...web-session-tabs-sync-agent-status.test.ts | 76 +++++- .../agent-status-patch.ts | 13 +- .../apply-final-patch.ts | 10 +- .../runtime/web-session-tabs-sync/state.ts | 1 + .../src/store/slices/agent-status-contract.ts | 2 + .../slices/agent-status-live-entry-builder.ts | 3 + src/shared/agent-status-ipc-payload.ts | 2 + src/shared/agent-status-store-publisher.ts | 7 +- .../agent-status-store-replication.test.ts | 43 +++ src/shared/agent-status-store-replication.ts | 81 +++++- src/shared/protocol-version.ts | 3 +- .../rpc-params-catalog.generated.ts | 2 + 41 files changed, 1706 insertions(+), 269 deletions(-) create mode 100644 src/main/runtime/agent-status-host-replica-store.test.ts create mode 100644 src/main/runtime/agent-status-host-replica-store.ts create mode 100644 src/main/runtime/rpc/methods/agent-status-store.test.ts create mode 100644 src/main/ssh/ssh-relay-agent-status-test-support.ts create mode 100644 src/relay/agent-hook-cache-actions.ts create mode 100644 src/relay/agent-hook-http-handler.ts create mode 100644 src/relay/agent-hook-status-store-source.ts delete mode 100644 src/relay/relay-agent-status-store.ts diff --git a/src/main/agent-hooks/hook-status-session-tabs-republish.ts b/src/main/agent-hooks/hook-status-session-tabs-republish.ts index b53e4e90501..c25b9d17dbb 100644 --- a/src/main/agent-hooks/hook-status-session-tabs-republish.ts +++ b/src/main/agent-hooks/hook-status-session-tabs-republish.ts @@ -7,7 +7,8 @@ type SessionTabsRepublisher = { touchMobileSessionTabsForWorktree(worktreeId: string): void } -type StatusStore = Pick +type StatusStore = Pick & + Partial> /** * Republish `session.tabs` whenever a pane's status row changes. @@ -50,16 +51,17 @@ export function installHookStatusSessionTabsRepublish( runtime.touchMobileSessionTabsForWorktree(worktreeId) } }) - const unsubscribeFreshness = statusStore.subscribeStatusFreshness((status) => { - const runtime = getRuntime() - if (!runtime) { - return - } - const worktreeId = resolveWorktreeId(status, runtime) - if (worktreeId) { - runtime.scheduleMobileSessionTabsAgentStatusHeartbeatForWorktree(worktreeId) - } - }) + const unsubscribeFreshness = + statusStore.subscribeStatusFreshness?.((status) => { + const runtime = getRuntime() + if (!runtime) { + return + } + const worktreeId = resolveWorktreeId(status, runtime) + if (worktreeId) { + runtime.scheduleMobileSessionTabsAgentStatusHeartbeatForWorktree(worktreeId) + } + }) ?? (() => {}) return () => { unsubscribeMutations() unsubscribeFreshness() diff --git a/src/main/agent-hooks/server/server-cleanup.ts b/src/main/agent-hooks/server/server-cleanup.ts index 9f36d676143..3d33729d542 100644 --- a/src/main/agent-hooks/server/server-cleanup.ts +++ b/src/main/agent-hooks/server/server-cleanup.ts @@ -145,10 +145,10 @@ export abstract class AgentHookServerCleanup extends AgentHookServerAuthorityFen } /** Clear statuses proven to belong to one lost SSH transport. */ - clearStatusEntriesForConnection(connectionId: string): void { + clearStatusEntriesForConnection(connectionId: string): number | null { const normalizedConnectionId = connectionId.trim() if (normalizedConnectionId.length === 0) { - return + return null } const clearedAt = Math.max( Date.now(), @@ -195,6 +195,7 @@ export abstract class AgentHookServerCleanup extends AgentHookServerAuthorityFen connectionId: normalizedConnectionId, clearedAt }) + return clearedAt } protected deleteStatusEntry( diff --git a/src/main/ipc/agent-hooks.ts b/src/main/ipc/agent-hooks.ts index f460be06b5f..bbd5c4d087a 100644 --- a/src/main/ipc/agent-hooks.ts +++ b/src/main/ipc/agent-hooks.ts @@ -19,12 +19,16 @@ type AgentHookHandlerDependencies = { getPtyIdForPaneKey?: (paneKey: string) => string | undefined } +type AgentStatusRuntime = AgentStatusRuntimeEnrichment & { + getAgentStatusSnapshot?: () => AgentStatusIpcPayload[] +} + // Why: install/remove are intentionally not exposed to the renderer. Orca // auto-installs managed hooks at app startup (see src/main/index.ts), so a // renderer-triggered remove would be silently reverted on the next launch // and mislead the user. export function registerAgentHookHandlers( - runtime?: AgentStatusRuntimeEnrichment, + runtime?: AgentStatusRuntime, dependencies: AgentHookHandlerDependencies = {} ): void { // Why: matches the defensive pattern in src/main/ipc/pty.ts so re-registration @@ -50,8 +54,7 @@ export function registerAgentHookHandlers( // lose replayed statuses while its local store is still empty. Match the // live push enrichment in main/index.ts so parent/child rows survive replay. return ( - agentHookServer - .getStatusSnapshot() + (runtime?.getAgentStatusSnapshot?.() ?? agentHookServer.getStatusSnapshot()) // Same rule as the live push: the renderer's feed bridge owns structured rows for now. .filter((entry) => entry.structuredHost === undefined) .map((entry) => enrichAgentStatusIpcPayload(entry, runtime)) diff --git a/src/main/orcad/orcad-entry.ts b/src/main/orcad/orcad-entry.ts index c5c328204fe..aa5f8bac9e9 100644 --- a/src/main/orcad/orcad-entry.ts +++ b/src/main/orcad/orcad-entry.ts @@ -156,6 +156,7 @@ async function startOrcadRuntime( await import('../agent-hooks/hook-status-session-tabs-republish') const { AgentStatusObservedPaneIdentities, AgentStatusObservedPaneIdentityCapture } = await import('../runtime/agent-status-observed-pane-identity') + const { AgentStatusHostReplicaStore } = await import('../runtime/agent-status-host-replica-store') let rpc: InstanceType | null = null let uninstallHookStatusRepublish = (): void => {} @@ -186,6 +187,11 @@ async function startOrcadRuntime( const agentStatusStorePublisher = agentHookServer.createStatusStorePublisher({ executionHostId: 'local' }) + const agentStatusHostReplicaStore = new AgentStatusHostReplicaStore() + const getAgentStatusSnapshot = () => [ + ...agentHookServer.getStatusSnapshot(), + ...agentStatusHostReplicaStore.getStatusSnapshot() + ] // Why a real Store: without one every persistence-backed RPC throws `runtime_unavailable` // and the read paths that use `this.store?.x ?? []` quietly answer "empty" instead — // a server that pairs and lists nothing looks healthy and is not. @@ -231,11 +237,12 @@ async function startOrcadRuntime( // Why here too and not only on the desktop: orcad serves `worktree.ps` and `agentSession.*`, // so without these a headless host publishes its structured chats nowhere and lists no agents. getAgentStatusSnapshot: () => - agentHookServer.getStatusSnapshot().filter((entry) => entry.providerSessionOnly !== true), + getAgentStatusSnapshot().filter((entry) => entry.providerSessionOnly !== true), agentStatusStorePublisher, - getAgentProviderSessionSnapshot: () => agentHookServer.getStatusSnapshot(), + agentStatusHostReplicaStore, + getAgentProviderSessionSnapshot: getAgentStatusSnapshot, getAgentProviderSessionRowsForPane: (paneKey) => - agentHookServer.getStatusSnapshotForPane(paneKey), + getAgentStatusSnapshot().filter((entry) => entry.paneKey === paneKey), // Why captured rather than resolved at read: the fleet snapshot remints cached rows on every // read, so a row observed under one process otherwise acquires whatever process owns the pane now. readObservedAgentStatusPaneIdentity: (paneKey) => observedPaneIdentities.read(paneKey), @@ -262,6 +269,15 @@ async function startOrcadRuntime( agentHookServer, () => runtime ) + const uninstallReplicaStatusRepublish = installHookStatusSessionTabsRepublish( + agentStatusHostReplicaStore, + () => runtime + ) + const uninstallStatusRepublish = uninstallHookStatusRepublish + uninstallHookStatusRepublish = () => { + uninstallStatusRepublish() + uninstallReplicaStatusRepublish() + } // Why the headless entry point rather than registerPtyHandlers directly: this is the // same call `--serve` makes, and it threads the store through. Without the store the diff --git a/src/main/runtime/agent-status-host-replica-store.test.ts b/src/main/runtime/agent-status-host-replica-store.test.ts new file mode 100644 index 00000000000..16fadec35ba --- /dev/null +++ b/src/main/runtime/agent-status-host-replica-store.test.ts @@ -0,0 +1,201 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' + +import type { AgentStatusIpcPayload } from '../../shared/agent-status-types' +import type { + AgentStatusStoreDelta, + AgentStatusStoreSnapshot +} from '../../shared/agent-status-store-replication' +import { AgentStatusHostReplicaStore } from './agent-status-host-replica-store' + +function row( + paneKey: string, + state: AgentStatusIpcPayload['state'] = 'working', + revision = 1, + receivedAt = 1 +): AgentStatusIpcPayload { + return { + paneKey, + connectionId: null, + receivedAt, + stateStartedAt: receivedAt, + state, + prompt: '', + observation: { + origin: 'hook', + authorityId: 'remote-authority', + incarnation: 1, + revision, + observedAt: receivedAt + } + } +} + +function snapshot( + executionHostId: AgentStatusStoreSnapshot['executionHostId'], + overrides: Partial = {} +): AgentStatusStoreSnapshot { + return { + type: 'snapshot', + executionHostId, + ownerEpoch: 'epoch-a', + cursor: 0, + complete: true, + rows: [], + ...overrides + } +} + +function delta( + executionHostId: AgentStatusStoreDelta['executionHostId'], + overrides: Partial = {} +): AgentStatusStoreDelta { + return { + type: 'delta', + executionHostId, + ownerEpoch: 'epoch-a', + previousCursor: 0, + cursor: 1, + changes: [], + ...overrides + } +} + +describe('AgentStatusHostReplicaStore', () => { + afterEach(() => { + vi.useRealTimers() + }) + + it('replaces complete membership for only the scoped host', () => { + const store = new AgentStatusHostReplicaStore() + store.apply(snapshot('ssh:first', { rows: [row('first-a'), row('first-b')] }), { + executionHostId: 'ssh:first', + connectionId: 'first' + }) + store.apply(snapshot('ssh:second', { rows: [row('second-a')] }), { + executionHostId: 'ssh:second', + connectionId: 'second' + }) + + store.apply(snapshot('ssh:first', { cursor: 1, rows: [] }), { + executionHostId: 'ssh:first', + connectionId: 'first' + }) + + expect(store.getHostSnapshot('ssh:first').rows).toEqual([]) + expect(store.getHostSnapshot('ssh:second').rows).toEqual([ + expect.objectContaining({ paneKey: 'second-a', connectionId: 'second' }) + ]) + }) + + it('retains omitted rows through an incomplete census', () => { + const store = new AgentStatusHostReplicaStore() + const routing = { executionHostId: 'ssh:first' as const, connectionId: 'first' } + store.apply(snapshot('ssh:first', { rows: [row('a'), row('b')] }), routing) + + store.apply( + snapshot('ssh:first', { + cursor: 1, + complete: false, + rows: [row('a', 'waiting', 2)] + }), + routing + ) + + expect(store.getHostSnapshot('ssh:first')).toMatchObject({ + membershipConfirmed: false, + rows: [ + expect.objectContaining({ paneKey: 'a', state: 'waiting' }), + expect.objectContaining({ paneKey: 'b', state: 'working' }) + ] + }) + }) + + it('retains rows and unconfirms membership across gaps and owner restarts', () => { + const store = new AgentStatusHostReplicaStore() + const routing = { executionHostId: 'ssh:first' as const, connectionId: 'first' } + store.apply(snapshot('ssh:first', { rows: [row('a')] }), routing) + + expect(store.apply(delta('ssh:first', { previousCursor: 1, cursor: 2 }), routing)).toBe( + 'resnapshot-required' + ) + expect( + store.apply( + delta('ssh:first', { ownerEpoch: 'epoch-b', previousCursor: 2, cursor: 3 }), + routing + ) + ).toBe('resnapshot-required') + expect(store.getHostSnapshot('ssh:first')).toMatchObject({ + ownerEpoch: 'epoch-b', + cursor: null, + membershipConfirmed: false, + rows: [expect.objectContaining({ paneKey: 'a' })] + }) + }) + + it('rejects frames from a different execution host', () => { + const store = new AgentStatusHostReplicaStore() + const mutations = vi.fn() + store.subscribeStatusRowMutations(mutations) + + expect( + store.apply(snapshot('ssh:other', { rows: [row('wrong')] }), { + executionHostId: 'ssh:expected', + connectionId: 'expected' + }) + ).toBe('wrong-host') + expect(store.getStatusSnapshot()).toEqual([]) + expect(mutations).not.toHaveBeenCalled() + }) + + it('keeps a waiting row on contact loss and removes it only on explicit abandon', () => { + const store = new AgentStatusHostReplicaStore() + store.apply(snapshot('ssh:first', { rows: [row('question', 'waiting')] }), { + executionHostId: 'ssh:first', + connectionId: 'first' + }) + + store.setContact('ssh:first', 'unverifiable') + expect(store.getHostSnapshot('ssh:first')).toMatchObject({ + contact: 'unverifiable', + rows: [expect.objectContaining({ paneKey: 'question', state: 'waiting' })] + }) + + store.abandonHost('ssh:first') + expect(store.getHostSnapshot('ssh:first')).toMatchObject({ + contact: 'unverifiable', + membershipConfirmed: false, + rows: [] + }) + }) + + it('uses a stable replica clock for duplicate evidence and advances it for new evidence', () => { + vi.useFakeTimers() + vi.setSystemTime(1_000) + const store = new AgentStatusHostReplicaStore() + const routing = { executionHostId: 'ssh:first' as const, connectionId: 'first' } + store.apply(snapshot('ssh:first', { rows: [row('a', 'working', 1, 99_000)] }), routing) + const first = store.getHostSnapshot('ssh:first').rows[0] + + vi.setSystemTime(5_000) + store.apply( + snapshot('ssh:first', { cursor: 1, rows: [row('a', 'working', 1, 99_000)] }), + routing + ) + const duplicate = store.getHostSnapshot('ssh:first').rows[0] + + vi.setSystemTime(7_000) + store.apply( + snapshot('ssh:first', { cursor: 2, rows: [row('a', 'waiting', 2, -99_000)] }), + routing + ) + const changed = store.getHostSnapshot('ssh:first').rows[0] + + expect(first).toMatchObject({ receivedAt: 1_000, replicaEvidenceReceivedAt: 1_000 }) + expect(duplicate).toMatchObject({ receivedAt: 1_000, replicaEvidenceReceivedAt: 1_000 }) + expect(changed).toMatchObject({ + state: 'waiting', + receivedAt: 7_000, + replicaEvidenceReceivedAt: 7_000 + }) + }) +}) diff --git a/src/main/runtime/agent-status-host-replica-store.ts b/src/main/runtime/agent-status-host-replica-store.ts new file mode 100644 index 00000000000..bea2eb086d7 --- /dev/null +++ b/src/main/runtime/agent-status-host-replica-store.ts @@ -0,0 +1,244 @@ +import { isDeepStrictEqual } from 'node:util' + +import type { AgentStatusIpcPayload } from '../../shared/agent-status-types' +import type { ExecutionHostId } from '../../shared/execution-host' +import { + AgentStatusStoreReplica, + agentStatusStoreRowIdentity, + agentStatusStoreRowKey, + type AgentStatusStoreFrame, + type AgentStatusStoreReplicaApplyResult, + type AgentStatusStoreReplicaContact, + type AgentStatusStoreReplicaHostSnapshot +} from '../../shared/agent-status-store-replication' + +export type AgentStatusHostReplicaRouting = { + executionHostId: ExecutionHostId + connectionId: string +} + +export type AgentStatusHostReplicaRowIdentity = { + paneKey: string + worktreeId?: string + terminalHandle?: string +} + +export type AgentStatusHostReplicaRowMutation = { + executionHostId: ExecutionHostId + before: AgentStatusHostReplicaRowIdentity | null + after: AgentStatusHostReplicaRowIdentity | null +} + +export type AgentStatusHostReplicaContactMutation = { + executionHostId: ExecutionHostId + contact: AgentStatusStoreReplicaContact +} + +export type AgentStatusHostReplicaApplyResult = AgentStatusStoreReplicaApplyResult | 'wrong-host' + +type RowReceipt = { + evidenceKey: string + evidenceReceivedAt: number + deliveredAt: number +} + +function scopedRowKey(executionHostId: ExecutionHostId, row: AgentStatusIpcPayload): string { + return `${executionHostId}\0${agentStatusStoreRowKey(agentStatusStoreRowIdentity(row))}` +} + +function evidenceKey(ownerEpoch: string | null, row: AgentStatusIpcPayload): string { + const observation = row.observation + if (observation) { + return `${ownerEpoch ?? ''}\0${observation.authorityId}\0${observation.incarnation}\0${observation.revision}` + } + return `${ownerEpoch ?? ''}\0${row.receivedAt}\0${row.evidenceObservedAt ?? ''}\0${row.stateStartedAt}\0${row.state}` +} + +function mutationIdentity(row: AgentStatusIpcPayload): AgentStatusHostReplicaRowIdentity { + return { + paneKey: row.paneKey, + ...(row.worktreeId ? { worktreeId: row.worktreeId } : {}), + ...(row.terminalHandle ? { terminalHandle: row.terminalHandle } : {}) + } +} + +/** Runtime-owned replicas for execution hosts reached through this host. Row semantics stay remote-owned. */ +export class AgentStatusHostReplicaStore { + private readonly replica = new AgentStatusStoreReplica() + private readonly routingByHost = new Map() + private readonly receiptsByRow = new Map() + private readonly rowMutationListeners = new Set< + (mutation: AgentStatusHostReplicaRowMutation) => void + >() + private readonly contactMutationListeners = new Set< + (mutation: AgentStatusHostReplicaContactMutation) => void + >() + private deliveredAt = 0 + + apply( + frame: AgentStatusStoreFrame, + routing: AgentStatusHostReplicaRouting, + options: { deliveryFloor?: number } = {} + ): AgentStatusHostReplicaApplyResult { + if (frame.executionHostId !== routing.executionHostId) { + return 'wrong-host' + } + const previousContact = this.replica.getHostSnapshot(routing.executionHostId).contact + const before = this.rowsByKey(routing.executionHostId) + this.routingByHost.set(routing.executionHostId, routing) + const result = this.replica.apply(frame) + if (previousContact !== 'live') { + this.emitContactMutation(routing.executionHostId, 'live') + } + if (result !== 'applied') { + return result + } + const snapshot = this.replica.getHostSnapshot(routing.executionHostId) + const rawAfter = new Map( + snapshot.rows.map((row) => [agentStatusStoreRowKey(agentStatusStoreRowIdentity(row)), row]) + ) + const deliveryFloor = options.deliveryFloor ?? -1 + for (const row of rawAfter.values()) { + const receiptKey = scopedRowKey(routing.executionHostId, row) + const nextEvidenceKey = evidenceKey(snapshot.ownerEpoch, row) + const previous = this.receiptsByRow.get(receiptKey) + const unchanged = previous?.evidenceKey === nextEvidenceKey + if (unchanged) { + continue + } + this.deliveredAt = Math.max(Date.now(), deliveryFloor + 1, this.deliveredAt + 1) + this.receiptsByRow.set(receiptKey, { + evidenceKey: nextEvidenceKey, + evidenceReceivedAt: this.deliveredAt, + deliveredAt: this.deliveredAt + }) + } + for (const [rowKey, row] of before) { + if (!rawAfter.has(rowKey)) { + this.receiptsByRow.delete(scopedRowKey(routing.executionHostId, row)) + } + } + const after = this.rowsByKey(routing.executionHostId) + this.emitRowDiff(routing.executionHostId, before, after) + return result + } + + setContact(executionHostId: ExecutionHostId, contact: AgentStatusStoreReplicaContact): void { + const previous = this.replica.getHostSnapshot(executionHostId).contact + this.replica.setContact(executionHostId, contact) + if (previous === contact) { + return + } + this.emitContactMutation(executionHostId, contact) + } + + private emitContactMutation( + executionHostId: ExecutionHostId, + contact: AgentStatusStoreReplicaContact + ): void { + for (const listener of this.contactMutationListeners) { + try { + listener({ executionHostId, contact }) + } catch (error) { + console.warn('[agent-status-replica] contact listener failed', error) + } + } + } + + abandonHost(executionHostId: ExecutionHostId): void { + const before = this.rowsByKey(executionHostId) + this.replica.abandonHost(executionHostId) + this.routingByHost.delete(executionHostId) + for (const row of before.values()) { + this.receiptsByRow.delete(scopedRowKey(executionHostId, row)) + } + this.emitRowDiff(executionHostId, before, new Map()) + } + + getStatusSnapshot(): AgentStatusIpcPayload[] { + const executionHostIds = new Set( + this.replica.getRows().map(({ executionHostId }) => executionHostId) + ) + return [...executionHostIds].flatMap( + (executionHostId) => this.getHostSnapshot(executionHostId).rows + ) + } + + getStatusSnapshotForPane(paneKey: string): AgentStatusIpcPayload[] { + return this.getStatusSnapshot().filter((row) => row.paneKey === paneKey) + } + + getHostSnapshot(executionHostId: ExecutionHostId): AgentStatusStoreReplicaHostSnapshot { + const snapshot = this.replica.getHostSnapshot(executionHostId) + const routing = this.routingByHost.get(executionHostId) + return { + ...snapshot, + rows: snapshot.rows.map((row) => this.toReaderRow(executionHostId, row, routing)) + } + } + + subscribeStatusRowMutations( + listener: (mutation: AgentStatusHostReplicaRowMutation) => void + ): () => void { + this.rowMutationListeners.add(listener) + return () => this.rowMutationListeners.delete(listener) + } + + subscribeContactMutations( + listener: (mutation: AgentStatusHostReplicaContactMutation) => void + ): () => void { + this.contactMutationListeners.add(listener) + return () => this.contactMutationListeners.delete(listener) + } + + private rowsByKey(executionHostId: ExecutionHostId): Map { + const snapshot = this.getHostSnapshot(executionHostId) + return new Map( + snapshot.rows.map((row) => [agentStatusStoreRowKey(agentStatusStoreRowIdentity(row)), row]) + ) + } + + private toReaderRow( + executionHostId: ExecutionHostId, + row: AgentStatusIpcPayload, + routing: AgentStatusHostReplicaRouting | undefined + ): AgentStatusIpcPayload { + const receipt = this.receiptsByRow.get(scopedRowKey(executionHostId, row)) + return { + ...row, + connectionId: routing?.connectionId ?? row.connectionId, + ...(receipt + ? { + receivedAt: receipt.deliveredAt, + replicaEvidenceReceivedAt: receipt.evidenceReceivedAt + } + : {}) + } + } + + private emitRowDiff( + executionHostId: ExecutionHostId, + before: ReadonlyMap, + after: ReadonlyMap + ): void { + for (const rowKey of new Set([...before.keys(), ...after.keys()])) { + const previous = before.get(rowKey) ?? null + const next = after.get(rowKey) ?? null + if (isDeepStrictEqual(previous, next)) { + continue + } + const mutation: AgentStatusHostReplicaRowMutation = { + executionHostId, + before: previous ? mutationIdentity(previous) : null, + after: next ? mutationIdentity(next) : null + } + for (const listener of this.rowMutationListeners) { + try { + listener(mutation) + } catch (error) { + console.warn('[agent-status-replica] row listener failed', error) + } + } + } + } +} diff --git a/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts b/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts index 04c8a87db1a..a6178e6c056 100644 --- a/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts +++ b/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts @@ -11,6 +11,7 @@ import type { import type { TerminalSideEffectBatch } from '../../shared/terminal-side-effect-facts' import type { AgentStatusIpcPayload } from '../../shared/agent-status-types' import type { AgentStatusStorePublisher } from '../../shared/agent-status-store-publisher' +import type { AgentStatusHostReplicaStore } from './agent-status-host-replica-store' import type { StructuredAgentSessionStatusSink } from '../native-chat/agent-session-wire/structured-agent-session-status-feed' import type { ObservedAgentStatusPaneIdentity } from '../ipc/agent-status-ipc-boundary' import type { AgentHookAuthorityAttestation } from '../agent-hooks/server' @@ -68,6 +69,7 @@ export class OrcaRuntimeWithPreservedBranchCleanup extends OrcaRuntimeWithTermin protected readonly getAgentStatusSnapshotFn: (() => AgentStatusIpcPayload[]) | null protected readonly agentStatusStorePublisherFn: AgentStatusStorePublisher | null + protected readonly agentStatusHostReplicaStoreFn: AgentStatusHostReplicaStore | null protected readonly structuredAgentStatusSinkFn: StructuredAgentSessionStatusSink | null diff --git a/src/main/runtime/orca-runtime-state-fields.ts b/src/main/runtime/orca-runtime-state-fields.ts index a8d0f411e11..db83b7171f3 100644 --- a/src/main/runtime/orca-runtime-state-fields.ts +++ b/src/main/runtime/orca-runtime-state-fields.ts @@ -7,6 +7,7 @@ import type { RuntimeTerminalAgentStatusEvent } from './runtime-terminal-contrac import type { TerminalSideEffectBatch } from '../../shared/terminal-side-effect-facts' import type { AgentStatusIpcPayload } from '../../shared/agent-status-types' import type { AgentStatusStorePublisher } from '../../shared/agent-status-store-publisher' +import type { AgentStatusHostReplicaStore } from './agent-status-host-replica-store' import type { StructuredAgentSessionStatusSink } from '../native-chat/agent-session-wire/structured-agent-session-status-feed' import type { ObservedAgentStatusPaneIdentity } from '../ipc/agent-status-ipc-boundary' import type { AgentHookAuthorityAttestation } from '../agent-hooks/server' @@ -42,6 +43,14 @@ export class OrcaRuntimeWithStateFields extends OrcaRuntimeWithLinearCommands { return this.agentStatusStorePublisherFn } + getAgentStatusHostReplicaStore(): AgentStatusHostReplicaStore | null { + return this.agentStatusHostReplicaStoreFn + } + + getAgentStatusSnapshot(): AgentStatusIpcPayload[] { + return this.getAgentProviderSessionSnapshotFn?.() ?? this.getAgentStatusSnapshotFn?.() ?? [] + } + constructor( store: RuntimeStore | null = null, stats?: StatsCollector, @@ -56,6 +65,7 @@ export class OrcaRuntimeWithStateFields extends OrcaRuntimeWithLinearCommands { // same inline agent rows the desktop sidebar does — same source, 1:1. getAgentStatusSnapshot?: () => AgentStatusIpcPayload[] agentStatusStorePublisher?: AgentStatusStorePublisher + agentStatusHostReplicaStore?: AgentStatusHostReplicaStore /** Where structured (native chat) sessions publish into that same store, so the snapshot * above lists them like every other agent. */ structuredAgentStatusSink?: StructuredAgentSessionStatusSink @@ -199,6 +209,7 @@ export class OrcaRuntimeWithStateFields extends OrcaRuntimeWithLinearCommands { } this.getAgentStatusSnapshotFn = deps?.getAgentStatusSnapshot ?? null this.agentStatusStorePublisherFn = deps?.agentStatusStorePublisher ?? null + this.agentStatusHostReplicaStoreFn = deps?.agentStatusHostReplicaStore ?? null this.structuredAgentStatusSinkFn = deps?.structuredAgentStatusSink ?? null this.readObservedAgentStatusPaneIdentityFn = deps?.readObservedAgentStatusPaneIdentity ?? (() => ({ kind: 'unobserved' })) diff --git a/src/main/runtime/rpc/methods/agent-status-store.test.ts b/src/main/runtime/rpc/methods/agent-status-store.test.ts new file mode 100644 index 00000000000..bfa398cf1da --- /dev/null +++ b/src/main/runtime/rpc/methods/agent-status-store.test.ts @@ -0,0 +1,91 @@ +import { describe, expect, it, vi } from 'vitest' + +import type { AgentStatusStoreFrame } from '../../../../shared/agent-status-store-replication' +import { AGENT_STATUS_STORE_REPLICA_CAPABILITY } from '../../../../shared/protocol-version' +import type { OrcaRuntimeService } from '../../orca-runtime' +import { AGENT_STATUS_STORE_METHODS } from './agent-status-store' + +const snapshot: AgentStatusStoreFrame = { + type: 'snapshot', + executionHostId: 'local', + ownerEpoch: 'epoch-a', + cursor: 0, + complete: true, + rows: [] +} + +function testRuntime( + options: { + subscribe?: (emit: (frame: AgentStatusStoreFrame) => void) => () => void + } = {} +): OrcaRuntimeService { + const publisher = { + snapshot: () => snapshot, + subscribe: options.subscribe ?? (() => () => {}) + } + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The methods under test read only this explicitly modeled runtime seam. + return { getAgentStatusStorePublisher: () => publisher } as OrcaRuntimeService +} + +describe('agent status store RPC methods', () => { + it('rejects snapshot reads from clients that did not negotiate the contract', () => { + const method = AGENT_STATUS_STORE_METHODS[0] + + expect(() => method.handler(undefined, { runtime: testRuntime() })).toThrow( + 'agent_status_store_capability_required' + ) + }) + + it('requires a lifecycle signal for subscriptions', async () => { + const method = AGENT_STATUS_STORE_METHODS[1] + + await expect( + method.handler( + undefined, + { + runtime: testRuntime(), + clientCapabilities: [AGENT_STATUS_STORE_REPLICA_CAPABILITY] + }, + vi.fn() + ) + ).rejects.toThrow('agent_status_store_subscription_signal_required') + }) + + it('does not subscribe an already-aborted request', async () => { + const subscribe = vi.fn(() => vi.fn()) + const controller = new AbortController() + controller.abort() + + await AGENT_STATUS_STORE_METHODS[1].handler( + undefined, + { + runtime: testRuntime({ subscribe }), + clientCapabilities: [AGENT_STATUS_STORE_REPLICA_CAPABILITY], + signal: controller.signal + }, + vi.fn() + ) + + expect(subscribe).not.toHaveBeenCalled() + }) + + it('releases the exact subscription when its request aborts', async () => { + const unsubscribe = vi.fn() + const subscribe = vi.fn(() => unsubscribe) + const controller = new AbortController() + const running = AGENT_STATUS_STORE_METHODS[1].handler( + undefined, + { + runtime: testRuntime({ subscribe }), + clientCapabilities: [AGENT_STATUS_STORE_REPLICA_CAPABILITY], + signal: controller.signal + }, + vi.fn() + ) + + expect(subscribe).toHaveBeenCalledOnce() + controller.abort() + await running + expect(unsubscribe).toHaveBeenCalledOnce() + }) +}) diff --git a/src/main/runtime/rpc/methods/agent-status-store.ts b/src/main/runtime/rpc/methods/agent-status-store.ts index f4a44abb229..cd035d19a17 100644 --- a/src/main/runtime/rpc/methods/agent-status-store.ts +++ b/src/main/runtime/rpc/methods/agent-status-store.ts @@ -33,22 +33,28 @@ export const AGENT_STATUS_STORE_METHODS = [ if (!clientCapabilities?.includes(AGENT_STATUS_STORE_REPLICA_CAPABILITY)) { throw new Error('agent_status_store_capability_required') } - const publisher = requirePublisher(runtime) - let unsubscribe = (): void => {} - const abort = (): void => unsubscribe() - unsubscribe = publisher.subscribe(emit) - signal?.addEventListener('abort', abort, { once: true }) - if (signal?.aborted) { - abort() + if (!signal) { + throw new Error('agent_status_store_subscription_signal_required') } - await new Promise((resolve) => { - signal?.addEventListener('abort', () => resolve(), { once: true }) - if (!signal) { - resolve() - } + if (signal.aborted) { + return + } + const publisher = requirePublisher(runtime) + let resolveAbort = (): void => {} + const aborted = new Promise((resolve) => { + resolveAbort = resolve }) - signal?.removeEventListener('abort', abort) - unsubscribe() + signal.addEventListener('abort', resolveAbort, { once: true }) + const unsubscribe = publisher.subscribe(emit) + try { + if (signal.aborted) { + resolveAbort() + } + await aborted + } finally { + signal.removeEventListener('abort', resolveAbort) + unsubscribe() + } } }) ] as const diff --git a/src/main/runtime/rpc/rpc-params-type-parity.ts b/src/main/runtime/rpc/rpc-params-type-parity.ts index 421aedbb42d..263a7fe38cd 100644 --- a/src/main/runtime/rpc/rpc-params-type-parity.ts +++ b/src/main/runtime/rpc/rpc-params-type-parity.ts @@ -8,12 +8,7 @@ import type { ALL_RPC_METHODS } from './methods' type RegisteredMethod = (typeof ALL_RPC_METHODS)[number] // These schemas reach into src/main and have no shared catalog entry. -type UncataloguedMethod = - | 'emulator.install' - | 'orchestration.send' - | 'orchestration.taskUpdate' - | 'agentStatus.getStoreSnapshot' - | 'agentStatus.subscribeStore' +type UncataloguedMethod = 'emulator.install' | 'orchestration.send' | 'orchestration.taskUpdate' type IsAny = 0 extends 1 & T ? true : false diff --git a/src/main/runtime/runtime-mobile-agent-status-projection.ts b/src/main/runtime/runtime-mobile-agent-status-projection.ts index b21fd8bb3ff..1160ae21b43 100644 --- a/src/main/runtime/runtime-mobile-agent-status-projection.ts +++ b/src/main/runtime/runtime-mobile-agent-status-projection.ts @@ -122,6 +122,8 @@ export function selectRuntimeHookAgentRowForPane( let live: AgentStatusIpcPayload | null = null const freshAfter = Date.now() - AGENT_STATUS_STALE_AFTER_MS for (const entry of rows) { + const evidenceAt = + entry.replicaEvidenceReceivedAt ?? entry.evidenceObservedAt ?? entry.receivedAt if (entry.providerSession && (!session || entry.receivedAt > session.receivedAt)) { session = entry } @@ -129,7 +131,7 @@ export function selectRuntimeHookAgentRowForPane( entry.agentType && (entry.providerSessionOnly !== true || (entry.agentType === 'pi' && entry.providerSession != null)) && - (entry.evidenceObservedAt ?? entry.receivedAt) >= freshAfter && + evidenceAt >= freshAfter && (!agent || entry.receivedAt > agent.receivedAt) ) { agent = entry @@ -138,7 +140,7 @@ export function selectRuntimeHookAgentRowForPane( entry.providerSessionOnly !== true && // Restored rows cannot prove liveness because the turn may have ended while offline (#12346). entry.restoredUnconfirmed !== true && - (entry.evidenceObservedAt ?? entry.receivedAt) >= freshAfter && + evidenceAt >= freshAfter && (!live || entry.receivedAt > live.receivedAt) ) { live = entry @@ -154,8 +156,10 @@ export function selectRuntimeHookAgentRowForPane( ? { payload: pickParsedAgentStatusPayload(live), updatedAt: live.receivedAt, - ...(live.evidenceObservedAt !== undefined - ? { evidenceObservedAt: live.evidenceObservedAt } + ...((live.replicaEvidenceReceivedAt ?? live.evidenceObservedAt) !== undefined + ? { + evidenceObservedAt: live.replicaEvidenceReceivedAt ?? live.evidenceObservedAt + } : {}), stateStartedAt: live.stateStartedAt ?? live.receivedAt, ...(live.worktreeId ? { worktreeId: live.worktreeId } : {}) diff --git a/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts b/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts index e536f76260c..7c33e0795f7 100644 --- a/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts +++ b/src/main/runtime/runtime-rpc/runtime-rpc-mobile-method-allowlist.ts @@ -11,6 +11,8 @@ export const MOBILE_RPC_METHOD_ALLOWLIST = new Set([ 'aiVault.searchStatus', 'aiVault.resolveSessionTitles', 'aiVault.prepareSessionResume', + 'agentStatus.getStoreSnapshot', + 'agentStatus.subscribeStore', 'browser.back', 'browser.dialogAccept', 'browser.dialogDismiss', diff --git a/src/main/ssh/ssh-relay-agent-status-test-support.ts b/src/main/ssh/ssh-relay-agent-status-test-support.ts new file mode 100644 index 00000000000..8b547333ee6 --- /dev/null +++ b/src/main/ssh/ssh-relay-agent-status-test-support.ts @@ -0,0 +1,44 @@ +import type { RelayDispatcher } from '../../relay/dispatcher' +import type { + AgentStatusStoreFrame, + AgentStatusStoreSnapshot +} from '../../shared/agent-status-store-replication' +import { + AGENT_STATUS_STORE_FRAME_NOTIFICATION, + AGENT_STATUS_STORE_REPLICA_CAPABILITY, + AGENT_STATUS_STORE_SNAPSHOT_METHOD, + AGENT_STATUS_STORE_SUBSCRIBE_METHOD +} from '../../shared/agent-status-store-replication' + +export function registerFakeAgentStatusRelayHandlers( + dispatcher: RelayDispatcher, + options: { statusSnapshot?: AgentStatusStoreSnapshot }, + onSubscribed: () => void +): void { + dispatcher.onRequest(AGENT_STATUS_STORE_SNAPSHOT_METHOD, async (params) => { + if (params.capability !== AGENT_STATUS_STORE_REPLICA_CAPABILITY) { + throw new Error('agent_status_store_capability_required') + } + if (!options.statusSnapshot) { + throw Object.assign(new Error('Method not found'), { code: -32601 }) + } + return options.statusSnapshot + }) + dispatcher.onRequest(AGENT_STATUS_STORE_SUBSCRIBE_METHOD, async (params) => { + if (params.capability !== AGENT_STATUS_STORE_REPLICA_CAPABILITY) { + throw new Error('agent_status_store_capability_required') + } + onSubscribed() + return { subscribed: true } + }) +} + +export function notifyFakeAgentStatusFrame( + dispatcher: RelayDispatcher, + subscribed: boolean, + frame: AgentStatusStoreFrame +): void { + if (subscribed) { + dispatcher.notify(AGENT_STATUS_STORE_FRAME_NOTIFICATION, frame) + } +} diff --git a/src/main/ssh/ssh-relay-deploy.ts b/src/main/ssh/ssh-relay-deploy.ts index 903d318e872..691a4fd89d8 100644 --- a/src/main/ssh/ssh-relay-deploy.ts +++ b/src/main/ssh/ssh-relay-deploy.ts @@ -114,6 +114,7 @@ import { MAX_SSH_RELAY_GRACE_PERIOD_SECONDS, MIN_SSH_RELAY_GRACE_PERIOD_SECONDS } from '../../shared/ssh-types' +import { toSshExecutionHostId } from '../../shared/execution-host' export type RelayDeployResult = { transport: MultiplexerTransport @@ -1700,6 +1701,7 @@ async function launchRelay( const escapedNode = shellEscape(nodePath) // Why: remoteRelayDir is shared across Orca targets for one account; hashing the target ID into the socket name stops cross-target attach. const sockName = relaySocketNameForInstanceId(relayInstanceId) + const executionHostId = relayInstanceId ? toSshExecutionHostId(relayInstanceId) : undefined const defaultSockFile = relayEndpointForHost(hostPlatform, remoteDir, sockName) const endpointDir = relayHookEndpointDirForHost(hostPlatform, remoteDir, defaultSockFile) const credentialFile = joinRemotePath(hostPlatform, remoteDir, `${sockName}.credential`) @@ -1733,7 +1735,8 @@ async function launchRelay( graceTime, activePipeMarkerPath, reconnectFallback: fallbackEndpoint, - credentialFile + credentialFile, + executionHostId }, signal ) @@ -1799,7 +1802,10 @@ async function launchRelay( // Why: the relay derives its hook endpoint dir from the socket path; pin it back under the relay dir when the socket moved to /tmp. const endpointDirArg = sockFile === defaultSockFile ? '' : ` --endpoint-dir ${shellEscape(endpointDir)}` - const launchCmd = `cd ${escapedDir} && nohup ${escapedNode} relay.js --detached --grace-time ${graceTime} --sock-path ${shellEscape(sockFile)}${endpointDirArg} --credential-file ${shellEscape(credentialFile)} --log-file ${shellEscape(logFile)} > ${shellEscape(logFile)} 2>&1 ${shellEscape(logFile)} 2>&1 {}) launchChannel.on('error', () => {}) @@ -1991,6 +1997,7 @@ type WindowsRelayLaunchOptions = { graceTime: number activePipeMarkerPath: string credentialFile: string + executionHostId?: string } & WindowsRelayEndpoint & { reconnectFallback?: WindowsRelayEndpoint } @@ -2073,7 +2080,8 @@ async function launchWindowsRelay( launchOpts.graceTime, logFile, errFile, - launchOpts.credentialFile + launchOpts.credentialFile, + launchOpts.executionHostId ), { signal } ) @@ -2165,7 +2173,8 @@ function windowsRelayLaunchCommand( graceTime: number, logFile: string, errFile: string, - credentialFile: string + credentialFile: string, + executionHostId?: string ): string { const relayScript = joinRemotePath(hostPlatform, remoteDir, 'relay.js') // Why: Windows sshd kills the exec channel's process tree on close; WMI re-parents the detached relay to survive. @@ -2180,6 +2189,7 @@ function windowsRelayLaunchCommand( quoted(sockPath), '--credential-file', quoted(credentialFile), + ...(executionHostId ? ['--execution-host-id', quoted(executionHostId)] : []), '--endpoint-dir', quoted(endpointDir), // Why: --log-file owns rotation; shell redirects still capture pre-JS boot/crash output. diff --git a/src/main/ssh/ssh-relay-session-agent-hooks.integration.test.ts b/src/main/ssh/ssh-relay-session-agent-hooks.integration.test.ts index 87d5c77f205..830308b886b 100644 --- a/src/main/ssh/ssh-relay-session-agent-hooks.integration.test.ts +++ b/src/main/ssh/ssh-relay-session-agent-hooks.integration.test.ts @@ -5,6 +5,11 @@ import type { SshPortForwardManager } from './ssh-port-forward' import type { SshConnection } from './ssh-connection' import type { MultiplexerTransport } from './ssh-channel-multiplexer' import type { AgentHookRelayEnvelope } from '../../shared/agent-hook-relay' +import type { + AgentStatusStoreFrame, + AgentStatusStoreSnapshot +} from '../../shared/agent-status-store-replication' +import { toSshExecutionHostId } from '../../shared/execution-host' import { RelayDispatcher } from '../../relay/dispatcher' import { AGENT_HOOK_NOTIFICATION_METHOD, @@ -16,6 +21,12 @@ import { agentHookServer, _internals as agentHookInternals } from '../agent-hook import { getSshPtyProvider } from '../ipc/pty' import { toAppSshPtyId } from '../providers/ssh-pty-id' import { DEFAULT_PTY_SOURCE_WINDOW_SU } from '../../shared/pty-source-credit-contract' +import { AgentStatusHostReplicaStore } from '../runtime/agent-status-host-replica-store' +import type { OrcaRuntimeService } from '../runtime/orca-runtime' +import { + notifyFakeAgentStatusFrame, + registerFakeAgentStatusRelayHandlers +} from './ssh-relay-agent-status-test-support' const { getCohortAtEmitMock, trackMock } = vi.hoisted(() => ({ getCohortAtEmitMock: vi.fn(), @@ -64,16 +75,16 @@ type FakeRelay = { replayEnvelopes: AgentHookRelayEnvelope[] notifyAgentHook: (envelope: AgentHookRelayEnvelope | Record) => void dispose: () => void + notifyStatusFrame: (frame: AgentStatusStoreFrame) => void } -// Why: mock below SSH at the relay transport boundary so CI covers session, -// mux, provider, and hook-ingest wiring without relying on a local sshd. -function createFakeRelay(): FakeRelay { +function createFakeRelay(options: { statusSnapshot?: AgentStatusStoreSnapshot } = {}): FakeRelay { let relayFeed: ((data: Buffer) => void) | null = null const clientDataCallbacks: ((data: Buffer) => void)[] = [] const clientCloseCallbacks: (() => void)[] = [] const ptySpawnRequests: Record[] = [] const replayEnvelopes: AgentHookRelayEnvelope[] = [] + let statusSubscribed = false const transport: MultiplexerTransport = { write: (data) => { @@ -137,6 +148,9 @@ function createFakeRelay(): FakeRelay { } return { replayed: replayEnvelopes.length } }) + registerFakeAgentStatusRelayHandlers(dispatcher, options, () => { + statusSubscribed = true + }) return { transport, @@ -146,11 +160,15 @@ function createFakeRelay(): FakeRelay { notifyAgentHook: (envelope) => { dispatcher.notify(AGENT_HOOK_NOTIFICATION_METHOD, envelope as Record) }, - dispose: () => dispatcher.dispose() + dispose: () => dispatcher.dispose(), + notifyStatusFrame: (frame) => notifyFakeAgentStatusFrame(dispatcher, statusSubscribed, frame) } } -function createSession(targetId: string): InstanceType { +function createSession( + targetId: string, + runtime?: OrcaRuntimeService +): InstanceType { const store = { getRepos: vi.fn().mockReturnValue([]), getSshPtyConsumerRecovery: vi.fn().mockReturnValue(null), @@ -175,7 +193,7 @@ function createSession(targetId: string): InstanceType { isDestroyed: () => false, webContents: { send: vi.fn() } }) - return new SshRelaySession(targetId, getMainWindow, store, portForwardManager) + return new SshRelaySession(targetId, getMainWindow, store, portForwardManager, runtime) } async function waitForStatusCount(events: CapturedStatus[], count: number): Promise { @@ -306,6 +324,126 @@ describe('SshRelaySession agent hooks over a fake relay transport', () => { }) }) + it('uses the host-owned replica stream for a capable relay instead of legacy ingest', async () => { + const targetId = 'conn-capable' + const executionHostId = toSshExecutionHostId(targetId) + const statusSnapshot: AgentStatusStoreSnapshot = { + type: 'snapshot', + executionHostId, + ownerEpoch: 'epoch-capable', + cursor: 0, + complete: true, + rows: [ + { + paneKey: `tab-capable:${SSH_LEAF_ID}`, + connectionId: null, + receivedAt: 1, + stateStartedAt: 1, + state: 'working', + prompt: 'host-owned' + } + ] + } + relay = createFakeRelay({ statusSnapshot }) + vi.mocked(deployAndLaunchRelay).mockResolvedValue({ + transport: relay.transport, + serverBuildId: 'test-relay-build', + platform: 'linux-x64' + }) + const replicaStore = new AgentStatusHostReplicaStore() + const ingestSpy = vi.spyOn(agentHookServer, 'ingestRemote') + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The integration seam only reads the replica accessor during status negotiation. + session = createSession(targetId, { + getAgentStatusHostReplicaStore: () => replicaStore + } as unknown as OrcaRuntimeService) + + await session.establish({} as SshConnection) + expect(ingestSpy).not.toHaveBeenCalled() + expect(replicaStore.getHostSnapshot(executionHostId)).toMatchObject({ + membershipConfirmed: true, + contact: 'live', + rows: [ + expect.objectContaining({ paneKey: statusSnapshot.rows[0].paneKey, connectionId: targetId }) + ] + }) + + relay.notifyStatusFrame({ + type: 'delta', + executionHostId, + ownerEpoch: 'epoch-capable', + previousCursor: 0, + cursor: 1, + changes: [ + { + type: 'set', + row: { + ...statusSnapshot.rows[0], + state: 'waiting', + receivedAt: 2, + stateStartedAt: 2 + } + } + ] + }) + await vi.waitFor(() => + expect(replicaStore.getHostSnapshot(executionHostId).rows[0]).toMatchObject({ + state: 'waiting' + }) + ) + + expect(relay).not.toBeNull() + const activeRelay = relay! + activeRelay.transport.close?.() + activeRelay.dispose() + expect(replicaStore.getHostSnapshot(executionHostId)).toMatchObject({ + contact: 'unverifiable', + rows: [expect.objectContaining({ state: 'waiting' })] + }) + ingestSpy.mockRestore() + }) + + it('abandons stale replica membership when a reconnect falls back to a legacy relay', async () => { + const targetId = 'conn-legacy-fallback' + const executionHostId = toSshExecutionHostId(targetId) + relay = createFakeRelay() + vi.mocked(deployAndLaunchRelay).mockResolvedValue({ + transport: relay.transport, + serverBuildId: 'test-relay-build', + platform: 'linux-x64' + }) + const replicaStore = new AgentStatusHostReplicaStore() + replicaStore.apply( + { + type: 'snapshot', + executionHostId, + ownerEpoch: 'old-epoch', + cursor: 0, + complete: true, + rows: [ + { + paneKey: `tab-legacy:${SSH_LEAF_ID}`, + connectionId: null, + receivedAt: 1, + stateStartedAt: 1, + state: 'working', + prompt: 'stale replica' + } + ] + }, + { executionHostId, connectionId: targetId } + ) + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The integration seam only reads the replica accessor during status negotiation. + session = createSession(targetId, { + getAgentStatusHostReplicaStore: () => replicaStore + } as unknown as OrcaRuntimeService) + + await session.establish({} as SshConnection) + expect(replicaStore.getHostSnapshot(executionHostId)).toMatchObject({ + membershipConfirmed: false, + rows: [] + }) + }) + it('preserves Claude monitoring mode across the SSH relay boundary', async () => { relay = createFakeRelay() vi.mocked(deployAndLaunchRelay).mockResolvedValue({ diff --git a/src/main/ssh/ssh-relay-session.ts b/src/main/ssh/ssh-relay-session.ts index ef33f2ecadd..06379f771bf 100644 --- a/src/main/ssh/ssh-relay-session.ts +++ b/src/main/ssh/ssh-relay-session.ts @@ -139,6 +139,16 @@ import { } from './ssh-pty-consumer-recovery' import { classifySshPtyFrameRejection, SshPtyFrameRejectionLog } from './ssh-pty-frame-rejection' import { SshPtyTargetedReattachQueue } from './ssh-pty-targeted-reattach-queue' +import { + AGENT_STATUS_STORE_FRAME_NOTIFICATION, + AGENT_STATUS_STORE_REPLICA_BUFFER_MAX, + AGENT_STATUS_STORE_REPLICA_CAPABILITY, + AGENT_STATUS_STORE_SNAPSHOT_METHOD, + AGENT_STATUS_STORE_SUBSCRIBE_METHOD, + isAgentStatusStoreFrame, + type AgentStatusStoreFrame, + type AgentStatusStoreSnapshot +} from '../../shared/agent-status-store-replication' export type RelaySessionState = 'idle' | 'deploying' | 'ready' | 'reconnecting' | 'disposed' @@ -288,6 +298,15 @@ function normalizeRelayGracePeriodSeconds(graceTimeSeconds: number | undefined): ) } +function isAgentStatusStoreSnapshot(value: unknown): value is AgentStatusStoreSnapshot { + return isAgentStatusStoreFrame(value) && value.type === 'snapshot' +} + +function isMissingAgentStatusStoreMethod(error: unknown): boolean { + const code = typeof error === 'object' && error !== null ? Reflect.get(error, 'code') : undefined + return code === -32601 || code === 'CONNECTION_LOST' || code === 'DISPOSED' +} + // Why: teardown barriers are independent, so one failing store write must not hide the others — // settle them all and aggregate, rather than rethrowing only whichever rejected first. async function settleSshSessionTeardown( @@ -370,6 +389,9 @@ export class SshRelaySession { private readonly ptyConsumerClientInstanceId: string private ptyConsumerSessionState: SshPtyConsumerSessionState | null = null private activeCompatibilityAttachmentIds = new Set() + private agentStatusReplicaNotificationCleanup: (() => void) | null = null + private agentStatusReplicaDisposeCleanup: (() => void) | null = null + private agentStatusReplicaRefresh: Promise | null = null constructor( readonly targetId: string, @@ -1206,12 +1228,180 @@ export class SshRelaySession { registerSshGitProvider(this.targetId, gitProvider) this.wireUpPtyEvents(ptyProvider, mux, providerGeneration) - this.wireUpAgentHookEvents(mux) + const replicatedStatus = await this.wireUpAgentStatusStore(mux, shouldContinue) + if (!replicatedStatus) { + this.wireUpAgentHookEvents(mux) + } this.wireUpRemoteWorkspaceEvents(mux) void this.installManagedHooksOnRemote(mux, shouldContinue) return true } + /** Negotiate the host-owned status projection before accepting legacy hook envelopes. */ + private async wireUpAgentStatusStore( + mux: SshChannelMultiplexer, + shouldContinue?: () => boolean + ): Promise { + this.teardownAgentStatusReplica() + const replicaStore = this.runtime?.getAgentStatusHostReplicaStore() + const expectedHostId = toSshExecutionHostId(this.targetId) + if (!replicaStore) { + return false + } + if (shouldContinue && !shouldContinue()) { + return false + } + let snapshot: AgentStatusStoreSnapshot + try { + const result = await mux.request(AGENT_STATUS_STORE_SNAPSHOT_METHOD, { + capability: AGENT_STATUS_STORE_REPLICA_CAPABILITY + }) + if (!isAgentStatusStoreSnapshot(result) || result.executionHostId !== expectedHostId) { + replicaStore.abandonHost(expectedHostId) + return false + } + snapshot = result + } catch (error) { + if ((!shouldContinue || shouldContinue()) && !mux.isDisposed()) { + replicaStore.abandonHost(expectedHostId) + } + if (!isMissingAgentStatusStoreMethod(error) && !mux.isDisposed()) { + console.warn( + `[ssh-relay-session] status store negotiation failed for ${this.targetId}: ${ + error instanceof Error ? error.message : String(error) + }` + ) + } + return false + } + if (shouldContinue && !shouldContinue()) { + return false + } + const pendingFrames: AgentStatusStoreFrame[] = [] + let pendingOverflow = false + let publicationActive = false + const applyFrame = (frame: AgentStatusStoreFrame, deliveryFloor?: number): void => { + const result = replicaStore.apply( + frame, + { executionHostId: expectedHostId, connectionId: this.targetId }, + { deliveryFloor } + ) + if (result === 'resnapshot-required') { + this.requestAgentStatusReplicaSnapshot(mux, expectedHostId) + } + } + this.agentStatusReplicaNotificationCleanup = mux.onNotificationByMethod( + AGENT_STATUS_STORE_FRAME_NOTIFICATION, + (params) => { + if (!isAgentStatusStoreFrame(params) || params.executionHostId !== expectedHostId) { + return + } + if (!publicationActive) { + if (pendingFrames.length >= AGENT_STATUS_STORE_REPLICA_BUFFER_MAX) { + pendingFrames.length = 0 + pendingOverflow = true + } else if (!pendingOverflow) { + pendingFrames.push(params) + } + return + } + applyFrame(params) + } + ) + try { + await mux.request(AGENT_STATUS_STORE_SUBSCRIBE_METHOD, { + capability: AGENT_STATUS_STORE_REPLICA_CAPABILITY + }) + } catch (error) { + this.teardownAgentStatusReplica() + replicaStore.abandonHost(expectedHostId) + if (!isMissingAgentStatusStoreMethod(error) && !mux.isDisposed()) { + console.warn( + `[ssh-relay-session] status store subscription failed for ${this.targetId}: ${ + error instanceof Error ? error.message : String(error) + }` + ) + } + return false + } + if (shouldContinue && !shouldContinue()) { + this.teardownAgentStatusReplica() + return false + } + const priorHost = replicaStore.getHostSnapshot(expectedHostId) + const legacyClearFloor = + priorHost.ownerEpoch === null + ? (agentHookServer.clearStatusEntriesForConnection(this.targetId) ?? undefined) + : undefined + publicationActive = true + applyFrame(snapshot, legacyClearFloor) + if (pendingOverflow) { + this.requestAgentStatusReplicaSnapshot(mux, expectedHostId) + } else { + for (const frame of pendingFrames) { + applyFrame(frame) + } + } + this.agentStatusReplicaDisposeCleanup = mux.onDispose(() => { + replicaStore.setContact(expectedHostId, 'unverifiable') + }) + return true + } + + private requestAgentStatusReplicaSnapshot( + mux: SshChannelMultiplexer, + expectedHostId: ExecutionHostId + ): void { + if (this.agentStatusReplicaRefresh) { + return + } + const pending = this.refreshAgentStatusReplicaSnapshot(mux, expectedHostId).finally(() => { + if (this.agentStatusReplicaRefresh === pending) { + this.agentStatusReplicaRefresh = null + } + }) + this.agentStatusReplicaRefresh = pending + } + + private async refreshAgentStatusReplicaSnapshot( + mux: SshChannelMultiplexer, + expectedHostId: ExecutionHostId + ): Promise { + const replicaStore = this.runtime?.getAgentStatusHostReplicaStore() + if (!replicaStore) { + return + } + try { + const result = await mux.request(AGENT_STATUS_STORE_SNAPSHOT_METHOD, { + capability: AGENT_STATUS_STORE_REPLICA_CAPABILITY + }) + if ( + this.mux !== mux || + mux.isDisposed() || + !isAgentStatusStoreSnapshot(result) || + result.executionHostId !== expectedHostId + ) { + return + } + replicaStore.apply(result, { + executionHostId: expectedHostId, + connectionId: this.targetId + }) + } catch { + if (this.mux === mux && !mux.isDisposed()) { + replicaStore.setContact(expectedHostId, 'unverifiable') + } + } + } + + private teardownAgentStatusReplica(): void { + this.agentStatusReplicaNotificationCleanup?.() + this.agentStatusReplicaNotificationCleanup = null + this.agentStatusReplicaDisposeCleanup?.() + this.agentStatusReplicaDisposeCleanup = null + this.agentStatusReplicaRefresh = null + } + private activePtyConsumerOwner(): SshPtyConsumerOwnerState | null { const state = this.ptyConsumerSessionState return state && state.mode !== 'legacy-fallback' ? state : null @@ -1638,6 +1828,12 @@ export class SshRelaySession { outputGenerationReason: string = reason ): void { this.releaseRelayLossWatcher() + if (reason === 'connection_lost') { + this.runtime + ?.getAgentStatusHostReplicaStore() + ?.setContact(toSshExecutionHostId(this.targetId), 'unverifiable') + } + this.teardownAgentStatusReplica() this.muxNotificationCleanup?.() this.muxNotificationCleanup = null for (const cleanup of this.ptyRecoveryNotificationCleanups) { diff --git a/src/main/startup/main-process-runtime-service.ts b/src/main/startup/main-process-runtime-service.ts index 2432e190544..a7bc4c6f37e 100644 --- a/src/main/startup/main-process-runtime-service.ts +++ b/src/main/startup/main-process-runtime-service.ts @@ -27,6 +27,8 @@ import { AgentStatusObservedPaneIdentities, recordObservedAgentStatusPaneIdentity } from '../runtime/agent-status-observed-pane-identity' +import { AgentStatusHostReplicaStore } from '../runtime/agent-status-host-replica-store' +import { installHookStatusSessionTabsRepublish } from '../agent-hooks/hook-status-session-tabs-republish' export function getDesktopWindowStatus(): RuntimeDesktopWindowStatus { const activation = state.desktopActivationGate @@ -71,6 +73,11 @@ export function initializeMainProcessRuntime(): OrcaRuntimeService { const agentStatusStorePublisher = agentHookServer.createStatusStorePublisher({ executionHostId: 'local' }) + const agentStatusHostReplicaStore = new AgentStatusHostReplicaStore() + const getAgentStatusSnapshot = () => [ + ...agentHookServer.getStatusSnapshot(), + ...agentStatusHostReplicaStore.getStatusSnapshot() + ] const runtime = new OrcaRuntimeService(store, stats, { agentSessionClaimSigner: loadAgentSessionClaimSigner( getProfileUserDataPath(), @@ -91,8 +98,9 @@ export function initializeMainProcessRuntime(): OrcaRuntimeService { getDesktopWindowStatus, // Why: worktree.ps pulls hook-reported agent status (same source as the desktop sidebar) at query time so mobile shows the same agents. getAgentStatusSnapshot: () => - agentHookServer.getStatusSnapshot().filter((entry) => entry.providerSessionOnly !== true), + getAgentStatusSnapshot().filter((entry) => entry.providerSessionOnly !== true), agentStatusStorePublisher, + agentStatusHostReplicaStore, // Why: structured chats have no hooks, so the host writes their projections here itself; the // snapshot above then lists them for the CLI and mobile without a second store. structuredAgentStatusSink: { @@ -105,9 +113,9 @@ export function initializeMainProcessRuntime(): OrcaRuntimeService { // Why: the filter above hides resume-identity rows from the live-agent views, but // those rows carry the provider session mobile native chat addresses transcripts // by — Pi publishes identity that way and would otherwise be unreachable. - getAgentProviderSessionSnapshot: () => agentHookServer.getStatusSnapshot(), + getAgentProviderSessionSnapshot: getAgentStatusSnapshot, getAgentProviderSessionRowsForPane: (paneKey) => - agentHookServer.getStatusSnapshotForPane(paneKey), + getAgentStatusSnapshot().filter((entry) => entry.paneKey === paneKey), attestAgentHookCompatibilityAuthority: (candidate) => agentHookServer.attestCompatibilityAuthority(candidate), retireAgentHookCompatibilityAuthority: (paneKey) => @@ -144,6 +152,11 @@ export function initializeMainProcessRuntime(): OrcaRuntimeService { }) app.once('will-quit', () => sessionSearch?.dispose()) state.runtime = runtime + const uninstallReplicaStatusRepublish = installHookStatusSessionTabsRepublish( + agentStatusHostReplicaStore, + () => runtime + ) + app.once('will-quit', uninstallReplicaStatusRepublish) agentHookServer.subscribeEnrichedStatus((enriched) => recordObservedAgentStatusPaneIdentity(observedPaneIdentities, enriched.paneKey, runtime) ) diff --git a/src/main/startup/main-window-agent-status.ts b/src/main/startup/main-window-agent-status.ts index 0f10a373bf9..cb9a2dd3d32 100644 --- a/src/main/startup/main-window-agent-status.ts +++ b/src/main/startup/main-window-agent-status.ts @@ -13,6 +13,9 @@ import { stopAllSyntheticTitleSpinners } from './synthetic-title-runtime' import { mainProcessState as state } from './main-process-state' +import { enrichAgentStatusIpcPayload } from '../ipc/agent-status-ipc-boundary' + +let uninstallReplicaStatusListener: (() => void) | null = null export type MainWindowAgentStatusOptions = { window: BrowserWindow @@ -27,6 +30,8 @@ export type MainWindowAgentStatusOptions = { } export function installMainWindowAgentStatusListeners(options: MainWindowAgentStatusOptions): void { + uninstallReplicaStatusListener?.() + uninstallReplicaStatusListener = null agentHookServer.setListener( ({ paneKey, @@ -139,6 +144,38 @@ export function installMainWindowAgentStatusListeners(options: MainWindowAgentSt }) } }) + const replicaStore = state.runtime?.getAgentStatusHostReplicaStore() + if (replicaStore) { + uninstallReplicaStatusListener = replicaStore.subscribeStatusRowMutations((mutation) => { + const window = state.mainWindow + const runtime = state.runtime + if (!window || window.isDestroyed() || !runtime) { + return + } + const paneKeys = new Set( + [mutation.before?.paneKey, mutation.after?.paneKey].filter( + (paneKey): paneKey is string => paneKey !== undefined + ) + ) + for (const paneKey of paneKeys) { + const rows = runtime + .getAgentStatusSnapshot() + .filter((row) => row.paneKey === paneKey && row.structuredHost === undefined) + if (rows.length === 0) { + window.webContents.send('agentStatus:clear', { paneKey }) + getDashboardPopoutWindow()?.webContents.send('agentStatus:clear', { paneKey }) + continue + } + for (const row of rows) { + const statusEvent = enrichAgentStatusIpcPayload(row, runtime) + window.webContents.send('agentStatus:set', statusEvent) + if (row.providerSessionOnly !== true) { + getDashboardPopoutWindow()?.webContents.send('agentStatus:set', statusEvent) + } + } + } + }) + } } export function clearMainWindowAgentStatusListeners(): void { @@ -146,6 +183,8 @@ export function clearMainWindowAgentStatusListeners(): void { agentHookServer.setListener(null) agentHookServer.setPaneStatusClearListener(null) setMigrationUnsupportedPtyListener(null) + uninstallReplicaStatusListener?.() + uninstallReplicaStatusListener = null // Why: stop the spinner timer here — it would fire into destroyed webContents, and per-pane teardown may never run for restored-but-untorn panes. stopAllSyntheticTitleSpinners() } diff --git a/src/relay/agent-hook-cache-actions.ts b/src/relay/agent-hook-cache-actions.ts new file mode 100644 index 00000000000..7a324f1f0b4 --- /dev/null +++ b/src/relay/agent-hook-cache-actions.ts @@ -0,0 +1,55 @@ +import type { AgentHookEventPayload } from '../shared/agent-hook-listener/listener-event' +import type { HookListenerState } from '../shared/agent-hook-listener/listener-state' +import { normalizeHookPayload } from '../shared/agent-hook-listener' +import { + isAgentHookSource, + type AgentHookRelayEnvelope, + type AgentHookSource +} from '../shared/agent-hook-relay' +import { buildSpoolHookBody, type SpoolRecord } from '../shared/agent-hook-spool' +import { buildRelayHookEnvelope, hookBodyEnv, hookBodyVersion } from './agent-hook-envelope-build' +import { selectReplayableCachedPanes } from './agent-hook-cached-pane-status' + +export function replayCachedRelayPayloads(options: { + cachedByPaneKey: ReadonlyMap + metaByPaneKey: ReadonlyMap + isPaneSurfaceRetired: (paneKey: string) => boolean + dropPane: (paneKey: string) => void + forward: (envelope: AgentHookRelayEnvelope) => void +}): number { + const replayable = selectReplayableCachedPanes(options) + for (const { event, meta } of replayable) { + options.forward( + buildRelayHookEnvelope(event, meta.source, meta.env, meta.version, { isReplay: true }) + ) + } + return replayable.length +} + +export function ingestRelaySpoolRecord(options: { + record: SpoolRecord + state: HookListenerState + env: string + applyEvent: ( + event: AgentHookEventPayload, + source: AgentHookSource, + env?: string, + version?: string, + options?: { isReplay?: boolean } + ) => void +}): void { + const { record } = options + if (!isAgentHookSource(record.source)) { + return + } + const body = buildSpoolHookBody(record) + const event = normalizeHookPayload(options.state, record.source, body, options.env, { + deferCompactOwnershipToClient: true + }) + if (!event) { + return + } + options.applyEvent(event, record.source, hookBodyEnv(body), hookBodyVersion(body), { + isReplay: true + }) +} diff --git a/src/relay/agent-hook-http-handler.ts b/src/relay/agent-hook-http-handler.ts new file mode 100644 index 00000000000..df247d2536c --- /dev/null +++ b/src/relay/agent-hook-http-handler.ts @@ -0,0 +1,85 @@ +import type { IncomingMessage, ServerResponse } from 'node:http' + +import { normalizeHookPayload } from '../shared/agent-hook-listener' +import { HOOK_REQUEST_SLOWLORIS_MS } from '../shared/agent-hook-listener/listener-limits' +import { mergeAgentHookRequestHeaders } from '../shared/agent-hook-listener/hook-envelope' +import type { AgentHookEventPayload } from '../shared/agent-hook-listener/listener-event' +import type { HookListenerState } from '../shared/agent-hook-listener/listener-state' +import { readRequestBody } from '../shared/agent-hook-listener/request-body' +import { resolveHookSource } from '../shared/agent-hook-listener/source-routing' +import type { AgentHookSource } from '../shared/agent-hook-relay' +import { + isHookRequestTruncatedError, + type HookTransportInterferenceTracker +} from '../shared/agent-hook-transport-interference' +import { hookBodyEnv, hookBodyVersion } from './agent-hook-envelope-build' +import type { AgentHookResultRetryScheduler } from './agent-hook-result-retry-scheduler' + +export type RelayHookHttpHandlerOptions = { + token: string + env: string + state: HookListenerState + retryScheduler: AgentHookResultRetryScheduler + transportInterference: HookTransportInterferenceTracker + applyEvent: ( + event: AgentHookEventPayload, + source: AgentHookSource, + env?: string, + version?: string + ) => void +} + +export async function handleRelayHookHttpRequest( + req: IncomingMessage, + res: ServerResponse, + options: RelayHookHttpHandlerOptions +): Promise { + if (req.method !== 'POST') { + res.writeHead(404) + res.end() + return + } + if (req.headers['x-orca-agent-hook-token'] !== options.token) { + res.writeHead(403) + res.end() + return + } + let destroyedBySlowlorisCap = false + req.setTimeout(HOOK_REQUEST_SLOWLORIS_MS, () => { + destroyedBySlowlorisCap = true + req.destroy() + }) + try { + const pathname = new URL(req.url ?? '/', 'http://127.0.0.1').pathname + const source = resolveHookSource(pathname) + if (!source) { + res.writeHead(404) + res.end() + return + } + const body = await readRequestBody(req) + const hookBody = mergeAgentHookRequestHeaders(body, req.headers) + const event = normalizeHookPayload(options.state, source, hookBody, options.env, { + deferCompactOwnershipToClient: true + }) + if (event) { + // TODO: source env/version from normalizeHookPayload once its result carries validated metadata. + const env = hookBodyEnv(hookBody) + const version = hookBodyVersion(hookBody) + options.applyEvent(event, source, env, version) + options.retryScheduler.scheduleAssistantMessageRetry(source, hookBody, event, env, version) + options.retryScheduler.scheduleCodexSubagentPoll(source, hookBody, event, env, version) + } + res.writeHead(204) + res.end() + } catch (error) { + if (isHookRequestTruncatedError(error) && !destroyedBySlowlorisCap) { + options.transportInterference.record({ source: null, error }) + } + process.stderr.write( + `[relay-hook-server] hook request failed: ${error instanceof Error ? error.message : String(error)}\n` + ) + res.writeHead(204) + res.end() + } +} diff --git a/src/relay/agent-hook-server.test.ts b/src/relay/agent-hook-server.test.ts index c3a99b8a20c..96bf9d93dab 100644 --- a/src/relay/agent-hook-server.test.ts +++ b/src/relay/agent-hook-server.test.ts @@ -80,6 +80,44 @@ describe('RelayAgentHookServer', () => { } }) + it('preserves state start across same-state delivery and resets it on transition', async () => { + const server = new RelayAgentHookServer({ endpointDir: dir, forward: vi.fn() }) + await server.start() + try { + const { port, token } = server.getCoordinates() + const post = (hookEventName: 'UserPromptSubmit' | 'Stop') => + fetch(`http://127.0.0.1:${port}/hook/claude`, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'X-Orca-Agent-Hook-Token': token + }, + body: JSON.stringify({ + paneKey: PANE_KEY, + payload: { + hook_event_name: hookEventName, + prompt: 'status timing' + } + }) + }) + + expect((await post('UserPromptSubmit')).status).toBe(204) + const first = server.getStatusSnapshot()[0] + expect((await post('UserPromptSubmit')).status).toBe(204) + const repeated = server.getStatusSnapshot()[0] + expect((await post('Stop')).status).toBe(204) + const transitioned = server.getStatusSnapshot()[0] + + expect(repeated.state).toBe('working') + expect(repeated.stateStartedAt).toBe(first.stateStartedAt) + expect(repeated.receivedAt).toBeGreaterThan(first.receivedAt) + expect(transitioned.state).toBe('done') + expect(transitioned.stateStartedAt).toBeGreaterThan(repeated.stateStartedAt) + } finally { + server.stop() + } + }) + it('normalizes and forwards raw spooled hooks on startup', async () => { const spoolDir = join(dir, 'spool') const spoolFile = join(spoolDir, 'pane-codex.jsonl') diff --git a/src/relay/agent-hook-server.ts b/src/relay/agent-hook-server.ts index 29a1f723195..9f41e2ef203 100644 --- a/src/relay/agent-hook-server.ts +++ b/src/relay/agent-hook-server.ts @@ -16,35 +16,24 @@ import { getEndpointFileName, writeEndpointFile } from '../shared/agent-hook-listener/endpoint-publication' -import { HOOK_REQUEST_SLOWLORIS_MS } from '../shared/agent-hook-listener/listener-limits' -import { normalizeHookPayload } from '../shared/agent-hook-listener' -import { mergeAgentHookRequestHeaders } from '../shared/agent-hook-listener/hook-envelope' -import { readRequestBody } from '../shared/agent-hook-listener/request-body' -import { resolveHookSource } from '../shared/agent-hook-listener/source-routing' import type { AgentHookEventPayload } from '../shared/agent-hook-listener/listener-event' import { createHookTransportInterferenceTracker, - describeHookTransportInterference, - isHookRequestTruncatedError + describeHookTransportInterference } from '../shared/agent-hook-transport-interference' import { - isAgentHookSource, REMOTE_AGENT_HOOK_ENV, type AgentHookRelayEnvelope, type AgentHookSource } from '../shared/agent-hook-relay' -import { - buildSpoolHookBody, - drainAgentHookSpool, - type SpoolRecord -} from '../shared/agent-hook-spool' +import { drainAgentHookSpool } from '../shared/agent-hook-spool' import { buildRelayHookPtyEnv, defaultEndpointDir } from './agent-hook-endpoint-coordinates' -import { buildRelayHookEnvelope, hookBodyEnv, hookBodyVersion } from './agent-hook-envelope-build' +import { buildRelayHookEnvelope } from './agent-hook-envelope-build' import { AgentHookResultRetryScheduler } from './agent-hook-result-retry-scheduler' -import { - evictCachedPanesOverCap, - selectReplayableCachedPanes -} from './agent-hook-cached-pane-status' +import { evictCachedPanesOverCap } from './agent-hook-cached-pane-status' +import { RelayAgentStatusStoreSource } from './agent-hook-status-store-source' +import { handleRelayHookHttpRequest } from './agent-hook-http-handler' +import { ingestRelaySpoolRecord, replayCachedRelayPayloads } from './agent-hook-cache-actions' export type RelayHookForward = (envelope: AgentHookRelayEnvelope) => void @@ -95,6 +84,7 @@ export class RelayAgentHookServer { private preferredPort: number private portFallbackApplied = false private retryScheduler: AgentHookResultRetryScheduler + private readonly statusStoreSource: RelayAgentStatusStoreSource constructor(options: RelayHookServerOptions) { this.env = options.env ?? REMOTE_AGENT_HOOK_ENV @@ -104,6 +94,7 @@ export class RelayAgentHookServer { this.preferredPort = options.preferredPort ?? 0 this.forward = options.forward this.isPaneSurfaceRetired = options.isPaneSurfaceRetired ?? (() => false) + this.statusStoreSource = new RelayAgentStatusStoreSource(this.state) this.retryScheduler = new AgentHookResultRetryScheduler({ state: this.state, env: this.env, @@ -125,7 +116,14 @@ export class RelayAgentHookServer { drainAgentHookSpool({ endpointDir: this.endpointDir, getPersistedLaunchTokenHash: () => undefined, - ingest: (record) => this.ingestSpoolRecord(record) + ingest: (record) => + ingestRelaySpoolRecord({ + record, + state: this.state, + env: this.env, + applyEvent: (event, source, env, version, eventOptions) => + this.applyEvent(event, source, env, version, eventOptions) + }) }) } catch (err) { // Why: a downstream relay failure must not prevent the loopback listener from starting; @@ -204,31 +202,38 @@ export class RelayAgentHookServer { this.retryScheduler.clearAll() clearAllListenerCaches(this.state) this.lastEnvelopeMetaByPaneKey.clear() + this.statusStoreSource.reset() } - /** Request-driven replay: re-forwards each cached paneKey payload as a fresh notification. Forwards are - * issued before the request handler returns, so the response trails all replayed notifications. */ + createStatusStorePublisher( + options: Parameters[0] + ) { + return this.statusStoreSource.createPublisher(options) + } + getStatusSnapshot() { + return this.statusStoreSource.getSnapshot() + } + + getStatusSnapshotForPane(paneKey: string) { + return this.statusStoreSource.getRowsForPane(paneKey) + } replayCachedPayloadsForPanes(): number { - const replayable = selectReplayableCachedPanes({ + return replayCachedRelayPayloads({ cachedByPaneKey: this.state.lastStatusByPaneKey, metaByPaneKey: this.lastEnvelopeMetaByPaneKey, isPaneSurfaceRetired: this.isPaneSurfaceRetired, - dropPane: (paneKey) => this.clearPaneState(paneKey) + dropPane: (paneKey) => this.clearPaneState(paneKey), + forward: this.forward }) - for (const { event, meta } of replayable) { - this.forward( - buildRelayHookEnvelope(event, meta.source, meta.env, meta.version, { isReplay: true }) - ) - } - return replayable.length } - /** Drop a paneKey's cached entries on PTY exit so a terminated pane can't resurface as a ghost event on reconnect. */ clearPaneState(paneKey: string): void { + const previous = this.state.lastStatusByPaneKey.get(paneKey) this.retryScheduler.clearAssistantMessageRetry(paneKey) this.retryScheduler.clearCodexSubagentPoll(paneKey) clearPaneCacheState(this.state, paneKey) this.lastEnvelopeMetaByPaneKey.delete(paneKey) + this.statusStoreSource.clearPane(paneKey, previous) } /** Env vars to inject into relay-spawned PTYs so the hook script/plugin POSTs back to this loopback server. */ @@ -250,58 +255,14 @@ export class RelayAgentHookServer { // ─── Private ────────────────────────────────────────────────────── private async handleRequest(req: IncomingMessage, res: ServerResponse): Promise { - if (req.method !== 'POST') { - res.writeHead(404) - res.end() - return - } - if (req.headers['x-orca-agent-hook-token'] !== this.token) { - res.writeHead(403) - res.end() - return - } - // Why: track our own destroy so the slowloris cap can't be misread as outside interference. - let destroyedBySlowlorisCap = false - req.setTimeout(HOOK_REQUEST_SLOWLORIS_MS, () => { - destroyedBySlowlorisCap = true - req.destroy() + return handleRelayHookHttpRequest(req, res, { + token: this.token, + env: this.env, + state: this.state, + retryScheduler: this.retryScheduler, + transportInterference: this.transportInterference, + applyEvent: (event, source, env, version) => this.applyEvent(event, source, env, version) }) - try { - const pathname = new URL(req.url ?? '/', 'http://127.0.0.1').pathname - const source = resolveHookSource(pathname) - if (!source) { - res.writeHead(404) - res.end() - return - } - const body = await readRequestBody(req) - const hookBody = mergeAgentHookRequestHeaders(body, req.headers) - const event = normalizeHookPayload(this.state, source, hookBody, this.env, { - deferCompactOwnershipToClient: true - }) - if (event) { - // TODO: once normalizeHookPayload returns validated env/version, drop bodyEnv/bodyVersion and source them from the listener result. - const env = hookBodyEnv(hookBody) - const version = hookBodyVersion(hookBody) - this.applyEvent(event, source, env, version) - this.retryScheduler.scheduleAssistantMessageRetry(source, hookBody, event, env, version) - this.retryScheduler.scheduleCodexSubagentPoll(source, hookBody, event, env, version) - } - res.writeHead(204) - res.end() - } catch (err) { - // Why (#11217): a remote host can run the same IDS; count truncations here so a blocked SSH - // relay reports the cause instead of an anonymous "hook request failed". - if (isHookRequestTruncatedError(err) && !destroyedBySlowlorisCap) { - this.transportInterference.record({ source: null, error: err }) - } - // Why: hooks fail open (204 on any error) so a buggy agent never blocks the run; still log so the 204 doesn't mask bugs. - process.stderr.write( - `[relay-hook-server] hook request failed: ${err instanceof Error ? err.message : String(err)}\n` - ) - res.writeHead(204) - res.end() - } } private applyEvent( @@ -326,28 +287,14 @@ export class RelayAgentHookServer { // it reconnects. Stripping it would let a cold relay replay a completion as an ordinary `done` // row and resurrect a pane that the client had already retired. const cachedEvent = event + const previous = this.state.lastStatusByPaneKey.get(event.paneKey) // Why: delete-then-set makes Map insertion order = recency, so the cap below evicts the longest-idle pane. this.state.lastStatusByPaneKey.delete(event.paneKey) this.state.lastStatusByPaneKey.set(event.paneKey, cachedEvent) this.lastEnvelopeMetaByPaneKey.delete(event.paneKey) this.lastEnvelopeMetaByPaneKey.set(event.paneKey, { source, env, version }) + this.statusStoreSource.recordEvent(event, previous) evictCachedPanesOverCap(this.state.lastStatusByPaneKey, (key) => this.clearPaneState(key)) this.forward(buildRelayHookEnvelope(event, source, env, version, options)) } - - private ingestSpoolRecord(record: SpoolRecord): void { - if (!isAgentHookSource(record.source)) { - return - } - const body = buildSpoolHookBody(record) - const event = normalizeHookPayload(this.state, record.source, body, this.env, { - deferCompactOwnershipToClient: true - }) - if (!event) { - return - } - this.applyEvent(event, record.source, hookBodyEnv(body), hookBodyVersion(body), { - isReplay: true - }) - } } diff --git a/src/relay/agent-hook-status-store-source.ts b/src/relay/agent-hook-status-store-source.ts new file mode 100644 index 00000000000..48ed750adc6 --- /dev/null +++ b/src/relay/agent-hook-status-store-source.ts @@ -0,0 +1,99 @@ +import type { AgentHookEventPayload } from '../shared/agent-hook-listener/listener-event' +import type { HookListenerState } from '../shared/agent-hook-listener/listener-state' +import type { AgentStatusIpcPayload } from '../shared/agent-status-types' +import { + AgentStatusStorePublisher, + type AgentStatusStorePublisherOptions, + type AgentStatusStoreSourceMutation +} from '../shared/agent-status-store-publisher' + +/** Delivery source for the relay's hook cache; it does not own status semantics. */ +export class RelayAgentStatusStoreSource { + private readonly receivedAtByPaneKey = new Map() + private readonly stateStartedAtByPaneKey = new Map() + private readonly listeners = new Set<(mutation: AgentStatusStoreSourceMutation) => void>() + + constructor(private readonly state: Pick) {} + + createPublisher( + options: Omit + ): AgentStatusStorePublisher { + return new AgentStatusStorePublisher({ + ...options, + source: { + getSnapshot: () => this.getSnapshot(), + getRowsForPane: (paneKey) => this.getRowsForPane(paneKey), + subscribeMutations: (listener) => { + this.listeners.add(listener) + return () => this.listeners.delete(listener) + } + } + }) + } + + getSnapshot(): AgentStatusIpcPayload[] { + return Array.from(this.state.lastStatusByPaneKey.values(), (event) => this.toStatusRow(event)) + } + + getRowsForPane(paneKey: string): AgentStatusIpcPayload[] { + const event = this.state.lastStatusByPaneKey.get(paneKey) + return event ? [this.toStatusRow(event)] : [] + } + + reset(): void { + this.receivedAtByPaneKey.clear() + this.stateStartedAtByPaneKey.clear() + this.listeners.clear() + } + + recordEvent(event: AgentHookEventPayload, previous: AgentHookEventPayload | undefined): void { + const receivedAt = Math.max(Date.now(), (this.receivedAtByPaneKey.get(event.paneKey) ?? -1) + 1) + this.receivedAtByPaneKey.set(event.paneKey, receivedAt) + const stateStartedAt = + previous && previous.payload.state === event.payload.state + ? (this.stateStartedAtByPaneKey.get(event.paneKey) ?? receivedAt) + : receivedAt + this.stateStartedAtByPaneKey.set(event.paneKey, stateStartedAt) + this.emit({ + before: previous ? { paneKey: event.paneKey } : null, + after: { paneKey: event.paneKey } + }) + } + + clearPane(paneKey: string, previous: AgentHookEventPayload | undefined): void { + this.receivedAtByPaneKey.delete(paneKey) + this.stateStartedAtByPaneKey.delete(paneKey) + if (previous) { + this.emit({ before: { paneKey }, after: null }) + } + } + + private toStatusRow(event: AgentHookEventPayload): AgentStatusIpcPayload { + const receivedAt = this.receivedAtByPaneKey.get(event.paneKey) ?? 0 + const stateStartedAt = this.stateStartedAtByPaneKey.get(event.paneKey) ?? receivedAt + return { + ...event.payload, + paneKey: event.paneKey, + ...(event.launchToken ? { launchToken: event.launchToken } : {}), + ...(event.tabId ? { tabId: event.tabId } : {}), + ...(event.worktreeId ? { worktreeId: event.worktreeId } : {}), + connectionId: null, + receivedAt, + stateStartedAt, + ...(event.providerSession ? { providerSession: event.providerSession } : {}), + ...(event.providerSessionOnly ? { providerSessionOnly: true } : {}) + } + } + + private emit(mutation: AgentStatusStoreSourceMutation): void { + for (const listener of this.listeners) { + try { + listener(mutation) + } catch (error) { + process.stderr.write( + `[relay-hook-server] status-row listener failed: ${error instanceof Error ? error.message : String(error)}\n` + ) + } + } + } +} diff --git a/src/relay/relay-agent-hook-runtime.ts b/src/relay/relay-agent-hook-runtime.ts index 88ae6bf2933..4ecf4f612e8 100644 --- a/src/relay/relay-agent-hook-runtime.ts +++ b/src/relay/relay-agent-hook-runtime.ts @@ -17,35 +17,40 @@ import { import { resolveSetupAgentSequenceLaunchCommand } from '../shared/setup-agent-sequencing' import { relayLogLine } from './relay-diagnostic-log' import { registerManagedHookInstaller } from './managed-hook-installer' -import { RelayAgentStatusStore } from './relay-agent-status-store' -import { AGENT_STATUS_STORE_FRAME_NOTIFICATION } from '../shared/agent-status-store-replication' +import { + AGENT_STATUS_STORE_FRAME_NOTIFICATION, + AGENT_STATUS_STORE_SNAPSHOT_METHOD, + AGENT_STATUS_STORE_SUBSCRIBE_METHOD +} from '../shared/agent-status-store-replication' import { AGENT_STATUS_STORE_REPLICA_CAPABILITY } from '../shared/protocol-version' +import type { ExecutionHostId } from '../shared/execution-host' export class RelayAgentHookRuntime { private readonly hookServer: RelayAgentHookServer - private readonly statusStore = new RelayAgentStatusStore() private readonly statusStorePublisher - private stopStatusPublication: (() => void) | null = null private readonly pluginOverlay = new PluginOverlayManager() + private readonly statusStoreSubscriptions = new Map void>() + private disposeStatusStoreDetachListener: (() => void) | null = null + private disposeStatusStoreDispatcherListener: (() => void) | null = null constructor( private readonly dispatcher: RelayDispatcher, private readonly ptyHandler: PtyHandler, sockPath: string, - endpointDir?: string + endpointDir?: string, + executionHostId: ExecutionHostId = 'local' ) { this.hookServer = new RelayAgentHookServer({ endpointDir: endpointDir ?? endpointDirForRelaySocket(sockPath), forward: (envelope) => { - this.statusStore.apply(envelope) publishAgentHookEnvelope(dispatcher, envelope) }, // Why: the PTY handler is the only component that knows which panes still have a client // surface, so it — not the client — decides whether a hook post describes a live pane. isPaneSurfaceRetired: (paneKey) => ptyHandler.isPaneSurfaceRetired(paneKey) }) - this.statusStorePublisher = this.statusStore.createPublisher({ - executionHostId: 'local' + this.statusStorePublisher = this.hookServer.createStatusStorePublisher({ + executionHostId }) } @@ -59,14 +64,6 @@ export class RelayAgentHookRuntime { } this.registerPtyEnvironment() this.registerHandlers() - this.stopStatusPublication = this.statusStorePublisher.subscribe((frame) => { - if (typeof this.dispatcher.notify === 'function') { - this.dispatcher.notify( - AGENT_STATUS_STORE_FRAME_NOTIFICATION, - Object.fromEntries(Object.entries(frame)) - ) - } - }) } publishEndpointFile(): void { @@ -74,8 +71,14 @@ export class RelayAgentHookRuntime { } stop(): void { - this.stopStatusPublication?.() - this.stopStatusPublication = null + for (const unsubscribe of this.statusStoreSubscriptions.values()) { + unsubscribe() + } + this.statusStoreSubscriptions.clear() + this.disposeStatusStoreDetachListener?.() + this.disposeStatusStoreDetachListener = null + this.disposeStatusStoreDispatcherListener?.() + this.disposeStatusStoreDispatcherListener = null this.statusStorePublisher.dispose() this.hookServer.stop() } @@ -86,7 +89,6 @@ export class RelayAgentHookRuntime { this.ptyHandler.setExitListener(({ paneKey, id }) => { if (paneKey) { this.hookServer.clearPaneState(paneKey) - this.statusStore.drop(paneKey) } this.pluginOverlay.clearOverlay(paneKey ?? id) }) @@ -95,7 +97,6 @@ export class RelayAgentHookRuntime { // reconnecting client cannot be handed a replay of an agent nobody owns. this.ptyHandler.setSurfaceRetiredListener(({ paneKey }) => { this.hookServer.clearPaneState(paneKey) - this.statusStore.drop(paneKey) }) } @@ -162,13 +163,54 @@ export class RelayAgentHookRuntime { } private registerHandlers(): void { - this.dispatcher.onRequest('agentStatus.getStoreSnapshot', async (params) => { + this.disposeStatusStoreDetachListener?.() + this.disposeStatusStoreDetachListener = + this.dispatcher.onClientDetached?.((clientId) => { + this.statusStoreSubscriptions.get(clientId)?.() + this.statusStoreSubscriptions.delete(clientId) + }) ?? null + this.disposeStatusStoreDispatcherListener?.() + this.disposeStatusStoreDispatcherListener = + this.dispatcher.onDisposed?.(() => { + for (const unsubscribe of this.statusStoreSubscriptions.values()) { + unsubscribe() + } + this.statusStoreSubscriptions.clear() + }) ?? null + this.dispatcher.onRequest(AGENT_STATUS_STORE_SNAPSHOT_METHOD, async (params) => { if (params.capability !== AGENT_STATUS_STORE_REPLICA_CAPABILITY) { throw new Error('agent_status_store_capability_required') } const snapshot = this.statusStorePublisher.snapshot() return snapshot }) + this.dispatcher.onRequest(AGENT_STATUS_STORE_SUBSCRIBE_METHOD, async (params, context) => { + if (params.capability !== AGENT_STATUS_STORE_REPLICA_CAPABILITY) { + throw new Error('agent_status_store_capability_required') + } + this.statusStoreSubscriptions.get(context.clientId)?.() + const unsubscribe = this.statusStorePublisher.subscribe((frame) => { + if (!context.isStale()) { + this.dispatcher.notifyClient( + context.clientId, + AGENT_STATUS_STORE_FRAME_NOTIFICATION, + frame + ) + } + }) + this.statusStoreSubscriptions.set(context.clientId, unsubscribe) + context.signal?.addEventListener( + 'abort', + () => { + if (this.statusStoreSubscriptions.get(context.clientId) === unsubscribe) { + this.statusStoreSubscriptions.delete(context.clientId) + unsubscribe() + } + }, + { once: true } + ) + return { subscribed: true } + }) this.dispatcher.onRequest(AGENT_HOOK_REQUEST_REPLAY_METHOD, async () => ({ replayed: this.hookServer.replayCachedPayloadsForPanes() })) diff --git a/src/relay/relay-agent-status-store.ts b/src/relay/relay-agent-status-store.ts deleted file mode 100644 index 9fe37e288de..00000000000 --- a/src/relay/relay-agent-status-store.ts +++ /dev/null @@ -1,82 +0,0 @@ -import type { AgentHookRelayEnvelope } from '../shared/agent-hook-relay' -import type { AgentStatusIpcPayload } from '../shared/agent-status-types' -import { AgentStatusStorePublisher } from '../shared/agent-status-store-publisher' -import type { - AgentStatusStorePublisherOptions, - AgentStatusStoreSourceMutation -} from '../shared/agent-status-store-publisher' - -function toStatusRow(envelope: AgentHookRelayEnvelope, receivedAt: number): AgentStatusIpcPayload { - return { - ...envelope.payload, - paneKey: envelope.paneKey, - ...(envelope.launchToken ? { launchToken: envelope.launchToken } : {}), - ...(envelope.tabId ? { tabId: envelope.tabId } : {}), - ...(envelope.worktreeId ? { worktreeId: envelope.worktreeId } : {}), - connectionId: null, - receivedAt, - stateStartedAt: receivedAt, - ...(envelope.providerSession ? { providerSession: envelope.providerSession } : {}), - ...(envelope.providerSessionOnly ? { providerSessionOnly: true } : {}) - } -} - -/** Relay-local projection fed by normalized hook envelopes; it owns delivery ordering only. */ -export class RelayAgentStatusStore { - private readonly rows = new Map() - private readonly mutations = new Set<(mutation: AgentStatusStoreSourceMutation) => void>() - private readonly receivedAtByPaneKey = new Map() - - apply(envelope: AgentHookRelayEnvelope): void { - const previous = this.rows.get(envelope.paneKey) - const receivedAt = Math.max( - Date.now(), - (this.receivedAtByPaneKey.get(envelope.paneKey) ?? -1) + 1 - ) - this.receivedAtByPaneKey.set(envelope.paneKey, receivedAt) - this.rows.set(envelope.paneKey, toStatusRow(envelope, receivedAt)) - this.emit({ - before: previous ? { paneKey: previous.paneKey } : null, - after: { paneKey: envelope.paneKey } - }) - } - - drop(paneKey: string): void { - if (!this.rows.delete(paneKey)) { - return - } - this.receivedAtByPaneKey.delete(paneKey) - this.emit({ before: { paneKey }, after: null }) - } - - createPublisher( - options: Omit - ): AgentStatusStorePublisher { - return new AgentStatusStorePublisher({ - ...options, - source: { - getSnapshot: () => [...this.rows.values()], - getRowsForPane: (paneKey) => { - const row = this.rows.get(paneKey) - return row ? [row] : [] - }, - subscribeMutations: (listener) => { - this.mutations.add(listener) - return () => this.mutations.delete(listener) - } - } - }) - } - - private emit(mutation: AgentStatusStoreSourceMutation): void { - for (const listener of this.mutations) { - try { - listener(mutation) - } catch (error) { - process.stderr.write( - `[relay-agent-status] mutation listener failed: ${error instanceof Error ? error.message : String(error)}\n` - ) - } - } - } -} diff --git a/src/relay/relay-daemon.ts b/src/relay/relay-daemon.ts index 27983176abf..ccc98817f45 100644 --- a/src/relay/relay-daemon.ts +++ b/src/relay/relay-daemon.ts @@ -53,7 +53,8 @@ export async function runRelayDaemon(options: RelayLaunchOptions): Promise primaryChannel.dispatcher, runtime.ptyHandler, options.sockPath, - options.endpointDir + options.endpointDir, + options.executionHostId ) const lifecycle = new RelayGraceLifecycle({ dispatcher: primaryChannel.dispatcher, diff --git a/src/relay/relay-launch-options.ts b/src/relay/relay-launch-options.ts index 95fafae818a..de14b1033f6 100644 --- a/src/relay/relay-launch-options.ts +++ b/src/relay/relay-launch-options.ts @@ -1,6 +1,7 @@ import { chmodSync, readFileSync } from 'node:fs' import { join } from 'node:path' import { DEFAULT_SSH_RELAY_GRACE_PERIOD_SECONDS } from '../shared/ssh-types' +import { normalizeExecutionHostId, type ExecutionHostId } from '../shared/execution-host' const DEFAULT_GRACE_MS = DEFAULT_SSH_RELAY_GRACE_PERIOD_SECONDS * 1000 const DEFAULT_SOCKET_NAME = 'relay.sock' @@ -22,6 +23,8 @@ export type RelayLaunchOptions = { endpointDir?: string logFile?: string credentialFile?: string + /** Explicit host scope stamped into replicated status frames. Legacy launches omit it and stay local. */ + executionHostId?: ExecutionHostId } export function parseRelayLaunchOptions(argv: string[]): RelayLaunchOptions { @@ -33,6 +36,7 @@ export function parseRelayLaunchOptions(argv: string[]): RelayLaunchOptions { let endpointDir: string | undefined let logFile: string | undefined let credentialFile: string | undefined + let executionHostId: ExecutionHostId | undefined for (let i = 2; i < argv.length; i++) { if (argv[i] === '--grace-time' && argv[i + 1]) { const parsed = Number.parseInt(argv[i + 1], 10) @@ -59,6 +63,9 @@ export function parseRelayLaunchOptions(argv: string[]): RelayLaunchOptions { } else if (argv[i] === '--credential-file' && argv[i + 1]) { credentialFile = argv[i + 1] i++ + } else if (argv[i] === '--execution-host-id' && argv[i + 1]) { + executionHostId = normalizeExecutionHostId(argv[i + 1]) ?? undefined + i++ } } if (!sockPath) { @@ -72,7 +79,8 @@ export function parseRelayLaunchOptions(argv: string[]): RelayLaunchOptions { sockPath, endpointDir, logFile, - credentialFile + credentialFile, + ...(executionHostId ? { executionHostId } : {}) } } diff --git a/src/renderer/src/hooks/ipc-events/agent-status-event-applicator.ts b/src/renderer/src/hooks/ipc-events/agent-status-event-applicator.ts index a4872822252..c2aace936cc 100644 --- a/src/renderer/src/hooks/ipc-events/agent-status-event-applicator.ts +++ b/src/renderer/src/hooks/ipc-events/agent-status-event-applicator.ts @@ -231,6 +231,9 @@ export function createAgentStatusEventApplicator(args: { ...(data.evidenceObservedAt !== undefined ? { evidenceObservedAt: data.evidenceObservedAt } : {}), + ...(data.replicaEvidenceReceivedAt !== undefined + ? { mirroredEvidenceReceivedAt: data.replicaEvidenceReceivedAt } + : {}), stateStartedAt: data.stateStartedAt }, routing: { diff --git a/src/renderer/src/runtime/web-session-tabs-sync-agent-status.test.ts b/src/renderer/src/runtime/web-session-tabs-sync-agent-status.test.ts index 64a59a80aa1..9ae98796a21 100644 --- a/src/renderer/src/runtime/web-session-tabs-sync-agent-status.test.ts +++ b/src/renderer/src/runtime/web-session-tabs-sync-agent-status.test.ts @@ -1,5 +1,6 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' import { makePaneKey } from '../../../shared/stable-pane-id' +import { AGENT_STATUS_STORE_REPLICA_CAPABILITY } from '../../../shared/protocol-version' import { applyWebSessionTabsSnapshot, type WebSessionTabsSyncState } from './web-session-tabs-sync' import { ENV, @@ -52,7 +53,7 @@ describe('applyWebSessionTabsSnapshot', () => { ]), ENV, NOW - ) as Partial + ) const mirroredId = patch.tabsByWorktree?.[WT]?.[0]?.id const mirroredPaneKey = makePaneKey(mirroredId!, LEAF_ID) @@ -71,6 +72,79 @@ describe('applyWebSessionTabsSnapshot', () => { expect(patch.sortEpoch).toBe(1) }) + it('lets a capable host replace a newer client row and reconcile an omitted row', () => { + const hostPaneKey = makePaneKey('host-tab-1', LEAF_ID) + const hostTerminal = { + type: 'terminal' as const, + id: HOST_SURFACE_ID, + title: 'codex [working]', + parentTabId: 'host-tab-1', + leafId: LEAF_ID, + isActive: true, + status: 'ready' as const, + terminal: 'terminal-1', + agentStatus: { + state: 'working' as const, + prompt: 'host evidence', + updatedAt: NOW, + stateStartedAt: NOW, + agentType: 'codex', + paneKey: hostPaneKey, + worktreeId: WT, + stateHistory: [] + } + } + const capableState = makeState({ + runtimeStatusByEnvironmentId: new Map([ + [ + ENV, + { + checkedAt: NOW, + status: { + runtimeId: ENV, + rendererGraphEpoch: 1, + graphStatus: 'ready', + authoritativeWindowId: null, + liveTabCount: 0, + liveLeafCount: 0, + capabilities: [AGENT_STATUS_STORE_REPLICA_CAPABILITY] + } + } + ] + ]) + }) + const initial = applyWebSessionTabsSnapshot( + capableState, + makeSnapshot([hostTerminal]), + ENV, + NOW + ) + const mirroredPaneKey = Object.keys(initial.agentStatusByPaneKey ?? {})[0]! + const newerClientRow = { + ...initial.agentStatusByPaneKey![mirroredPaneKey]!, + state: 'waiting' as const, + updatedAt: NOW + 10_000, + stateStartedAt: NOW + 10_000 + } + const hostWins = applyWebSessionTabsSnapshot( + { ...capableState, ...initial, agentStatusByPaneKey: { [mirroredPaneKey]: newerClientRow } }, + makeSnapshot([{ ...hostTerminal, agentStatus: { ...hostTerminal.agentStatus } }], { + snapshotVersion: 2 + }), + ENV, + NOW + 1 + ) + expect(hostWins.agentStatusByPaneKey?.[mirroredPaneKey]?.state).toBe('working') + + const omitted = applyWebSessionTabsSnapshot( + { ...capableState, ...initial, agentStatusByPaneKey: { [mirroredPaneKey]: newerClientRow } }, + makeSnapshot([{ ...hostTerminal, agentStatus: undefined }], { snapshotVersion: 3 }), + ENV, + NOW + 2 + ) + expect(omitted.agentStatusByPaneKey?.[mirroredPaneKey]).toBeUndefined() + }) + it('clears stale tool-output provenance when a newer host preview is assistant prose', () => { const hostPaneKey = makePaneKey('host-tab-1', LEAF_ID) const initialSnapshot = makeSnapshot([ diff --git a/src/renderer/src/runtime/web-session-tabs-sync/agent-status-patch.ts b/src/renderer/src/runtime/web-session-tabs-sync/agent-status-patch.ts index 4b29c800ced..e4d53329c0c 100644 --- a/src/renderer/src/runtime/web-session-tabs-sync/agent-status-patch.ts +++ b/src/renderer/src/runtime/web-session-tabs-sync/agent-status-patch.ts @@ -61,7 +61,8 @@ export function buildMirroredAgentStatusPatch( terminalSurfaceTabs: readonly TerminalSurface[], mirroredTerminalTabs: readonly MirroredTerminalTab[], now: number, - batchContext?: WebSessionTabsBatchContext + batchContext?: WebSessionTabsBatchContext, + options: { hostOwnsAgentStatus?: boolean } = {} ): Pick | null { const mirroredTabIds = new Set() for (const tab of currentTerminalTabs) { @@ -112,10 +113,13 @@ export function buildMirroredAgentStatusPatch( // state (still adopting the host's identity fields below) unless the host // carries a state class the client's bytes can never see. const clientOwnsEntry = + options.hostOwnsAgentStatus !== true && isFencedClientAgentStatus(entry.paneKey, existing, now) && !hostAgentStatusPiercesClientAuthority(entry) const nextEntry = - existing && (clientOwnsEntry || existing.updatedAt > entry.updatedAt) + options.hostOwnsAgentStatus !== true && + existing && + (clientOwnsEntry || existing.updatedAt > entry.updatedAt) ? { ...normalizeCompatibleAgentStatusEntryForOwner(existing, entry.agentType), ...(clientOwnsEntry && existing.state === 'working' && entry.state === 'working' @@ -161,7 +165,10 @@ export function buildMirroredAgentStatusPatch( // there is nothing to arbitrate, and a client asleep past the stale // boundary would otherwise erase every pane it owns on the first snapshot // after wake (STA-3107) instead of decaying it like a local pane. - if (isClientOwnedAgentStatus(paneKey, state.agentStatusByPaneKey[paneKey])) { + if ( + options.hostOwnsAgentStatus !== true && + isClientOwnedAgentStatus(paneKey, state.agentStatusByPaneKey[paneKey]) + ) { continue } if (nextAgentStatusByPaneKey === state.agentStatusByPaneKey) { diff --git a/src/renderer/src/runtime/web-session-tabs-sync/apply-final-patch.ts b/src/renderer/src/runtime/web-session-tabs-sync/apply-final-patch.ts index 4ef51cd52f7..4027eb36cd8 100644 --- a/src/renderer/src/runtime/web-session-tabs-sync/apply-final-patch.ts +++ b/src/renderer/src/runtime/web-session-tabs-sync/apply-final-patch.ts @@ -6,6 +6,7 @@ import { buildRetractedMirroredTabSweepPatch } from './agent-status-primitives' import { isWebSessionTabsWorktreeRemovalFrame } from './session-tabs-inventory-absence' +import { AGENT_STATUS_STORE_REPLICA_CAPABILITY } from '../../../../shared/protocol-version' type FinalPatchContext = ReturnType @@ -15,6 +16,7 @@ export function buildWebSessionTabsFinalPatch( ): WebSessionTabsSyncState | Partial { const { state, + environmentId, snapshot, worktreeId, now, @@ -58,7 +60,13 @@ export function buildWebSessionTabsFinalPatch( terminalSurfaceTabs, mirroredTerminalTabs, now, - batchContext + batchContext, + { + hostOwnsAgentStatus: + state.runtimeStatusByEnvironmentId + ?.get(environmentId) + ?.status?.capabilities?.includes(AGENT_STATUS_STORE_REPLICA_CAPABILITY) === true + } ) // A tombstone clears all environments' view of a worktree; it is not a terminal retraction. const retractedTabSweepPatch = isWebSessionTabsWorktreeRemovalFrame(snapshot) diff --git a/src/renderer/src/runtime/web-session-tabs-sync/state.ts b/src/renderer/src/runtime/web-session-tabs-sync/state.ts index f44d4a3d18c..03006149e00 100644 --- a/src/renderer/src/runtime/web-session-tabs-sync/state.ts +++ b/src/renderer/src/runtime/web-session-tabs-sync/state.ts @@ -210,6 +210,7 @@ export type WebSessionTabsSyncState = Pick< | 'recentlyRetiredAgentStatusPaneKeys' | 'retainedAgentsByPaneKey' | 'retentionSuppressedPaneKeys' + | 'runtimeStatusByEnvironmentId' > > diff --git a/src/renderer/src/store/slices/agent-status-contract.ts b/src/renderer/src/store/slices/agent-status-contract.ts index eaab5bb78ae..85a24f1fa26 100644 --- a/src/renderer/src/store/slices/agent-status-contract.ts +++ b/src/renderer/src/store/slices/agent-status-contract.ts @@ -98,6 +98,8 @@ export type AgentStatusTiming = { updatedAt?: number /** Observation clock for staleness; see `AgentStatusEntry.evidenceObservedAt`. */ evidenceObservedAt?: number + /** Local receipt clock for evidence replicated from another execution host. */ + mirroredEvidenceReceivedAt?: number stateStartedAt?: number } diff --git a/src/renderer/src/store/slices/agent-status-live-entry-builder.ts b/src/renderer/src/store/slices/agent-status-live-entry-builder.ts index b161e51bc04..5aa314f19a4 100644 --- a/src/renderer/src/store/slices/agent-status-live-entry-builder.ts +++ b/src/renderer/src/store/slices/agent-status-live-entry-builder.ts @@ -227,6 +227,9 @@ export function buildAgentStatusLiveEntry( ...(timing?.evidenceObservedAt !== undefined ? { evidenceObservedAt: timing.evidenceObservedAt } : {}), + ...(timing?.mirroredEvidenceReceivedAt !== undefined + ? { mirroredEvidenceReceivedAt: timing.mirroredEvidenceReceivedAt } + : {}), ...(metadata?.structuredHostOwned === true ? { structuredHostOwned: true as const } : {}), stateStartedAt, agentType: identity.agentType, diff --git a/src/shared/agent-status-ipc-payload.ts b/src/shared/agent-status-ipc-payload.ts index 3bff35272c2..5f8de14085a 100644 --- a/src/shared/agent-status-ipc-payload.ts +++ b/src/shared/agent-status-ipc-payload.ts @@ -53,6 +53,8 @@ export type AgentStatusIpcPayload = ParsedAgentStatusPayload & { evidenceObservedAt?: number /** Timestamp (ms) when the current state first appeared for this pane. */ stateStartedAt: number + /** Replica-local receipt of the current evidence. Never published by the execution host. */ + replicaEvidenceReceivedAt?: number orchestration?: AgentStatusOrchestrationContext providerSession?: AgentProviderSessionMetadata /** Resume identity update only; the status-shaped fields are transport placeholders. */ diff --git a/src/shared/agent-status-store-publisher.ts b/src/shared/agent-status-store-publisher.ts index b182ee24aa0..3f33c632b95 100644 --- a/src/shared/agent-status-store-publisher.ts +++ b/src/shared/agent-status-store-publisher.ts @@ -90,8 +90,11 @@ export class AgentStatusStorePublisher { cursor: this.cursor, reason: 'overflow' }) - this.subscribers.delete(subscriber) - return () => {} + subscriber.buffered = [] + subscriber.buffering = false + return () => { + this.subscribers.delete(subscriber) + } } subscriber.emit(snapshot) for (const delta of subscriber.buffered) { diff --git a/src/shared/agent-status-store-replication.test.ts b/src/shared/agent-status-store-replication.test.ts index a768e0e13b2..2c08cfe78c6 100644 --- a/src/shared/agent-status-store-replication.test.ts +++ b/src/shared/agent-status-store-replication.test.ts @@ -112,6 +112,37 @@ describe('AgentStatusStoreReplica', () => { }) }) + it('resets the cursor when a new-owner resnapshot marker arrives', () => { + const replica = new AgentStatusStoreReplica() + replica.apply(snapshot({ rows: [row('a')], cursor: 7 })) + + expect( + replica.apply({ + type: 'resnapshot-required', + executionHostId: 'local', + ownerEpoch: 'epoch-b', + cursor: 2, + reason: 'owner-restart' + }) + ).toBe('resnapshot-required') + expect(replica.getHostSnapshot('local')).toMatchObject({ + ownerEpoch: 'epoch-b', + cursor: null, + membershipConfirmed: false, + rows: [row('a')] + }) + expect( + replica.apply( + delta({ + ownerEpoch: 'epoch-b', + previousCursor: 2, + cursor: 3, + changes: [{ type: 'drop', identity: { paneKey: 'a' } }] + }) + ) + ).toBe('resnapshot-required') + }) + it('does not turn contact loss into row deletion or execution evidence', () => { const replica = new AgentStatusStoreReplica() replica.apply(snapshot({ rows: [row('question', 'waiting')] })) @@ -209,5 +240,17 @@ describe('AgentStatusStorePublisher', () => { reason: 'overflow' } ]) + + for (const listener of listeners) { + listener({ before: null, after: { paneKey: 'c' } }) + } + expect(frames[1]).toEqual({ + type: 'delta', + executionHostId: 'local', + ownerEpoch: 'epoch-a', + previousCursor: 2, + cursor: 3, + changes: [{ type: 'set', row: row('c') }] + }) }) }) diff --git a/src/shared/agent-status-store-replication.ts b/src/shared/agent-status-store-replication.ts index 46e5c16868a..6ddb9c5318d 100644 --- a/src/shared/agent-status-store-replication.ts +++ b/src/shared/agent-status-store-replication.ts @@ -1,9 +1,11 @@ -import type { AgentStatusIpcPayload } from './agent-status-types' -import type { ExecutionHostId } from './execution-host' +import { normalizeAgentStatusPayload, type AgentStatusIpcPayload } from './agent-status-types' +import { normalizeExecutionHostId, type ExecutionHostId } from './execution-host' export const AGENT_STATUS_STORE_REPLICA_CAPABILITY = 'agent-status.store-replica.v1' as const export const AGENT_STATUS_STORE_REPLICA_BUFFER_MAX = 256 export const AGENT_STATUS_STORE_FRAME_NOTIFICATION = 'agentStatus.storeFrame' as const +export const AGENT_STATUS_STORE_SUBSCRIBE_METHOD = 'agentStatus.subscribeStore' as const +export const AGENT_STATUS_STORE_SNAPSHOT_METHOD = 'agentStatus.getStoreSnapshot' as const export type AgentStatusStoreRowIdentity = { paneKey: string @@ -65,6 +67,72 @@ export type AgentStatusStoreReplicaHostSnapshot = { export type AgentStatusStoreReplicaApplyResult = 'applied' | 'ignored-stale' | 'resnapshot-required' +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +function isCursor(value: unknown): value is number { + return Number.isSafeInteger(value) && Number(value) >= 0 +} + +function isRowIdentity(value: unknown): value is AgentStatusStoreRowIdentity { + return ( + isRecord(value) && + typeof value.paneKey === 'string' && + value.paneKey.length > 0 && + (value.runId === undefined || typeof value.runId === 'string') + ) +} + +function isStatusRow(value: unknown): value is AgentStatusIpcPayload { + return ( + isRecord(value) && + typeof value.paneKey === 'string' && + value.paneKey.length > 0 && + (value.connectionId === null || typeof value.connectionId === 'string') && + typeof value.receivedAt === 'number' && + Number.isFinite(value.receivedAt) && + typeof value.stateStartedAt === 'number' && + Number.isFinite(value.stateStartedAt) && + normalizeAgentStatusPayload(value) !== null + ) +} + +/** Validates replication frames at transport boundaries before they can mutate a replica. */ +export function isAgentStatusStoreFrame(value: unknown): value is AgentStatusStoreFrame { + if ( + !isRecord(value) || + normalizeExecutionHostId( + typeof value.executionHostId === 'string' ? value.executionHostId : undefined + ) !== value.executionHostId || + typeof value.ownerEpoch !== 'string' || + value.ownerEpoch.length === 0 || + value.ownerEpoch.length > 128 || + !isCursor(value.cursor) + ) { + return false + } + if (value.type === 'snapshot') { + return ( + typeof value.complete === 'boolean' && + Array.isArray(value.rows) && + value.rows.every(isStatusRow) + ) + } + if (value.type === 'resnapshot-required') { + return value.reason === 'gap' || value.reason === 'overflow' || value.reason === 'owner-restart' + } + if (value.type !== 'delta' || !isCursor(value.previousCursor) || !Array.isArray(value.changes)) { + return false + } + return value.changes.every( + (change) => + isRecord(change) && + ((change.type === 'set' && isStatusRow(change.row)) || + (change.type === 'drop' && isRowIdentity(change.identity))) + ) +} + export function agentStatusStoreRowKey(identity: AgentStatusStoreRowIdentity): string { return identity.runId ? `run\0${identity.runId}` : `pane\0${identity.paneKey}` } @@ -106,6 +174,10 @@ export class AgentStatusStoreReplica { return this.applySnapshot(state, frame) } if (frame.type === 'resnapshot-required') { + // The marker is authoritative about the publisher epoch. Reset the cursor as + // well as membership so a stale delta from the prior owner cannot be accepted. + state.ownerEpoch = frame.ownerEpoch + state.cursor = null state.membershipConfirmed = false return 'resnapshot-required' } @@ -139,6 +211,11 @@ export class AgentStatusStoreReplica { state.contact = contact } + /** Explicitly abandons one mirrored host without implying that any represented process exited. */ + abandonHost(executionHostId: ExecutionHostId): void { + this.hosts.delete(executionHostId) + } + getHostSnapshot(executionHostId: ExecutionHostId): AgentStatusStoreReplicaHostSnapshot { const state = this.hosts.get(executionHostId) ?? newReplicaHostState() return { diff --git a/src/shared/protocol-version.ts b/src/shared/protocol-version.ts index 0d7ae39a85b..43f5b76f001 100644 --- a/src/shared/protocol-version.ts +++ b/src/shared/protocol-version.ts @@ -1,4 +1,5 @@ import { REMOTE_SERVER_UPDATE_CAPABILITY } from './remote-server-update' +import { AGENT_STATUS_STORE_REPLICA_CAPABILITY } from './agent-status-store-replication' import { SKILL_BUNDLE_INSTALL_CAPABILITY, SKILL_DELETE_CAPABILITY, @@ -132,7 +133,7 @@ export const AGENT_SESSION_BOUNDARY_RUNTIME_CAPABILITY = export { REMOTE_SERVER_UPDATE_CAPABILITY } from './remote-server-update' export const AGENT_SESSION_HOST_AUTHORITY_RUNTIME_CAPABILITY = 'agent-session.host-authority.v1' as const -export const AGENT_STATUS_STORE_REPLICA_CAPABILITY = 'agent-status.store-replica.v1' as const +export { AGENT_STATUS_STORE_REPLICA_CAPABILITY } from './agent-status-store-replication' export const AGENT_SESSION_OMP_RESUME_PATH_RUNTIME_CAPABILITY = 'agent-session.omp-resume-path.v1' as const // Why: structured sessions are journal-backed, not PTY-backed, so an incapable client must not diff --git a/src/shared/rpc-contract/rpc-params-catalog.generated.ts b/src/shared/rpc-contract/rpc-params-catalog.generated.ts index deae9bc0eaa..064c1243218 100644 --- a/src/shared/rpc-contract/rpc-params-catalog.generated.ts +++ b/src/shared/rpc-contract/rpc-params-catalog.generated.ts @@ -577,6 +577,8 @@ export const RPC_PARAMS_BY_METHOD = { 'agentSession.subscribe': SubscribeParams, 'agentSession.subscribeStatus': null, 'agentSession.unsubscribe': UnsubscribeParams, + 'agentStatus.getStoreSnapshot': null, + 'agentStatus.subscribeStore': null, 'agentTeams.prepareLaunch': AgentTeamsPrepareLaunch, 'agentTeams.tmuxCompat': AgentTeamsTmuxCompat, 'aiVault.listSessions': AiVaultListSessionsParams,