feat(agent-status): compose host replicas across SSH

This commit is contained in:
Brennan Benson
2026-09-15 11:33:23 -07:00
parent 94bda37fc8
commit 033416546c
41 changed files with 1706 additions and 269 deletions
@@ -7,7 +7,8 @@ type SessionTabsRepublisher = {
touchMobileSessionTabsForWorktree(worktreeId: string): void
}
type StatusStore = Pick<AgentHookServer, 'subscribeStatusFreshness' | 'subscribeStatusRowMutations'>
type StatusStore = Pick<AgentHookServer, 'subscribeStatusRowMutations'> &
Partial<Pick<AgentHookServer, 'subscribeStatusFreshness'>>
/**
* 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()
@@ -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(
+6 -3
View File
@@ -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))
+19 -3
View File
@@ -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<typeof OrcaRuntimeRpcServer> | 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
@@ -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> = {}
): AgentStatusStoreSnapshot {
return {
type: 'snapshot',
executionHostId,
ownerEpoch: 'epoch-a',
cursor: 0,
complete: true,
rows: [],
...overrides
}
}
function delta(
executionHostId: AgentStatusStoreDelta['executionHostId'],
overrides: Partial<AgentStatusStoreDelta> = {}
): 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
})
})
})
@@ -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<ExecutionHostId, AgentStatusHostReplicaRouting>()
private readonly receiptsByRow = new Map<string, RowReceipt>()
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<string, AgentStatusIpcPayload> {
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<string, AgentStatusIpcPayload>,
after: ReadonlyMap<string, AgentStatusIpcPayload>
): 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)
}
}
}
}
}
@@ -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
@@ -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' }))
@@ -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()
})
})
@@ -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<void>((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<void>((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
@@ -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<T> = 0 extends 1 & T ? true : false
@@ -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 } : {})
@@ -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',
@@ -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)
}
}
+14 -4
View File
@@ -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 </dev/null &`
const executionHostIdArg = executionHostId
? ` --execution-host-id ${shellEscape(executionHostId)}`
: ''
const launchCmd = `cd ${escapedDir} && nohup ${escapedNode} relay.js --detached --grace-time ${graceTime} --sock-path ${shellEscape(sockFile)}${endpointDirArg} --credential-file ${shellEscape(credentialFile)}${executionHostIdArg} --log-file ${shellEscape(logFile)} > ${shellEscape(logFile)} 2>&1 </dev/null &`
const launchChannel = await conn.exec(launchCmd, { signal })
launchChannel.on('data', () => {})
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.
@@ -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<string, unknown>) => 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<string, unknown>[] = []
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<string, unknown>)
},
dispose: () => dispatcher.dispose()
dispose: () => dispatcher.dispose(),
notifyStatusFrame: (frame) => notifyFakeAgentStatusFrame(dispatcher, statusSubscribed, frame)
}
}
function createSession(targetId: string): InstanceType<typeof SshRelaySession> {
function createSession(
targetId: string,
runtime?: OrcaRuntimeService
): InstanceType<typeof SshRelaySession> {
const store = {
getRepos: vi.fn().mockReturnValue([]),
getSshPtyConsumerRecovery: vi.fn().mockReturnValue(null),
@@ -175,7 +193,7 @@ function createSession(targetId: string): InstanceType<typeof SshRelaySession> {
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<void> {
@@ -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({
+197 -1
View File
@@ -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<string>()
private agentStatusReplicaNotificationCleanup: (() => void) | null = null
private agentStatusReplicaDisposeCleanup: (() => void) | null = null
private agentStatusReplicaRefresh: Promise<void> | 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<boolean> {
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<void> {
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) {
@@ -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)
)
@@ -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()
}
+55
View File
@@ -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<string, AgentHookEventPayload>
metaByPaneKey: ReadonlyMap<string, { source: AgentHookSource; env?: string; version?: string }>
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
})
}
+85
View File
@@ -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<void> {
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()
}
}
+38
View File
@@ -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')
+44 -97
View File
@@ -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<RelayAgentStatusStoreSource['createPublisher']>[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<void> {
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
})
}
}
@@ -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<string, number>()
private readonly stateStartedAtByPaneKey = new Map<string, number>()
private readonly listeners = new Set<(mutation: AgentStatusStoreSourceMutation) => void>()
constructor(private readonly state: Pick<HookListenerState, 'lastStatusByPaneKey'>) {}
createPublisher(
options: Omit<AgentStatusStorePublisherOptions, 'source'>
): 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`
)
}
}
}
}
+63 -21
View File
@@ -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<number, () => 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()
}))
-82
View File
@@ -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<string, AgentStatusIpcPayload>()
private readonly mutations = new Set<(mutation: AgentStatusStoreSourceMutation) => void>()
private readonly receivedAtByPaneKey = new Map<string, number>()
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<AgentStatusStorePublisherOptions, 'source'>
): 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`
)
}
}
}
}
+2 -1
View File
@@ -53,7 +53,8 @@ export async function runRelayDaemon(options: RelayLaunchOptions): Promise<void>
primaryChannel.dispatcher,
runtime.ptyHandler,
options.sockPath,
options.endpointDir
options.endpointDir,
options.executionHostId
)
const lifecycle = new RelayGraceLifecycle({
dispatcher: primaryChannel.dispatcher,
+9 -1
View File
@@ -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 } : {})
}
}
@@ -231,6 +231,9 @@ export function createAgentStatusEventApplicator(args: {
...(data.evidenceObservedAt !== undefined
? { evidenceObservedAt: data.evidenceObservedAt }
: {}),
...(data.replicaEvidenceReceivedAt !== undefined
? { mirroredEvidenceReceivedAt: data.replicaEvidenceReceivedAt }
: {}),
stateStartedAt: data.stateStartedAt
},
routing: {
@@ -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<WebSessionTabsSyncState>
)
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([
@@ -61,7 +61,8 @@ export function buildMirroredAgentStatusPatch(
terminalSurfaceTabs: readonly TerminalSurface[],
mirroredTerminalTabs: readonly MirroredTerminalTab[],
now: number,
batchContext?: WebSessionTabsBatchContext
batchContext?: WebSessionTabsBatchContext,
options: { hostOwnsAgentStatus?: boolean } = {}
): Pick<WebSessionTabsSyncState, 'agentStatusByPaneKey' | 'agentStatusEpoch' | 'sortEpoch'> | null {
const mirroredTabIds = new Set<string>()
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) {
@@ -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<typeof applyActiveStateUpdates>
@@ -15,6 +16,7 @@ export function buildWebSessionTabsFinalPatch(
): WebSessionTabsSyncState | Partial<WebSessionTabsSyncState> {
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)
@@ -210,6 +210,7 @@ export type WebSessionTabsSyncState = Pick<
| 'recentlyRetiredAgentStatusPaneKeys'
| 'retainedAgentsByPaneKey'
| 'retentionSuppressedPaneKeys'
| 'runtimeStatusByEnvironmentId'
>
>
@@ -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
}
@@ -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,
+2
View File
@@ -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. */
+5 -2
View File
@@ -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) {
@@ -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') }]
})
})
})
+79 -2
View File
@@ -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<string, unknown> {
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 {
+2 -1
View File
@@ -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
+2
View File
@@ -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,