Merge fix-send-federation-transcript into integrate-fixes

This commit is contained in:
Jinwoo-H
2026-09-04 02:35:09 -04:00
31 changed files with 634 additions and 673 deletions
@@ -86,6 +86,7 @@ describe('extracted orchestration worker formatting', () => {
} as never)
).toBe(
'Source: transcript (provider=codex)\n' +
'Archived: false\n' +
'Continuation cursor (opaque; pass unchanged to --cursor): owr1_next\n\n' +
'[assistant] working\n[tool inspect] [unserializable input]\n[tool result error] failed\n[image] https://example.test/proof.png'
)
@@ -94,6 +94,8 @@ describe('worker-read plain formatting', () => {
)
).toBe(
'Source: transcript (provider=codex)\n' +
'Worker: ready\n' +
'Archived: false\n' +
'Source exact: true\n' +
'Content complete: false\n' +
'Clipping: message_limit_or_scan_window\n' +
@@ -126,6 +128,8 @@ describe('worker-read plain formatting', () => {
)
).toBe(
'Source: terminal\n' +
'Worker: ready\n' +
'Archived: false\n' +
'Source exact: false\n' +
'Fallback reason: session_not_reported\n' +
'Content complete: false\n' +
@@ -159,12 +163,37 @@ describe('worker-read plain formatting', () => {
)
).toBe(
'Source: transcript (provider=codex)\n' +
'Worker: ready\n' +
'Archived: false\n' +
'Source exact: true\n' +
'Content complete: true\n' +
'Continuation cursor (opaque; pass unchanged to --cursor): owr1_empty\n\n' +
'No transcript messages returned. This exact transcript read did not request terminal evidence.'
)
})
it('distinguishes a released archive read from a live one', () => {
const output = formatWorkerRead({
dispatchId: 'dispatch_1',
status: { worker: 'succeeded', terminal: 'released', liveness: 'unverifiable' },
source: 'terminal',
sourceIdentity: 'private-source-identity',
terminal: {
handle: 'term_worker',
status: 'exited',
tail: ['archived tail'],
truncated: false,
nextCursor: null
},
cursor: null,
fallbackReason: null,
warnings: [],
archived: true
} as unknown as OrchestrationWorkerReadResult)
expect(output).toContain('Archived: true')
expect(output).toContain('Liveness: unverifiable')
expect(output).toContain('Worker: succeeded')
})
})
function workerReadResult(
@@ -68,6 +68,14 @@ function formatWorkerReadDetails(value: OrchestrationWorkerReadResult): string {
? `Source: transcript (provider=${value.provider})`
: 'Source: terminal'
const lines = [source]
// A released archive read otherwise prints identically to a live one.
if (value.status?.worker) {
lines.push(`Worker: ${value.status.worker}`)
}
lines.push(`Archived: ${value.archived === true}`)
if (value.status?.liveness) {
lines.push(`Liveness: ${value.status.liveness}`)
}
if (value.sourceExact !== undefined) {
lines.push(`Source exact: ${value.sourceExact}`)
}
@@ -0,0 +1,57 @@
import { describe, expect, it } from 'vitest'
import { wslHookRelayConnectionId } from '../../shared/wsl-hook-relay-contract'
import { OrcaRuntimeWithGetTerminalInteractiveWait } from './orca-runtime-get-terminal-interactive-wait'
const PANE_KEY = 'tab:worker'
const PTY_ID = 'pty-wsl'
type ExactWorkerProviderSessionHost = {
getExactWorkerProviderSession: (handle: string, observedAfter: number) => unknown
}
/** Drives the shipping method, not the selector helper: the wiring is what regressed. */
function selectThroughRuntime(statusConnectionId: string | null): unknown {
const runtime = {
getTerminalPaneKey: () => PANE_KEY,
getTerminalProcessIncarnation: () => 'pty-wsl:inc-1',
getTerminalAgentStatusPtyId: () => PTY_ID,
ptysById: new Map([
[PTY_ID, { connectionId: null, launchToken: 'launch-1', wslDistro: 'Ubuntu' }]
]),
wslDistroByPtyId: new Map([[PTY_ID, 'Ubuntu']]),
getAgentStatusSnapshotFn: () => [
{
paneKey: PANE_KEY,
connectionId: statusConnectionId,
launchToken: 'launch-1',
agentType: 'codex',
receivedAt: 500,
providerSession: { key: 'session_id', id: 's1', transcriptPath: '/t.jsonl' }
}
]
}
return (
OrcaRuntimeWithGetTerminalInteractiveWait.prototype as unknown as ExactWorkerProviderSessionHost
).getExactWorkerProviderSession.call(runtime as never, 'term_wsl', 0)
}
describe('exact worker provider session wiring', () => {
it('selects a local hook status for a local pane', () => {
expect(selectThroughRuntime(null)).toMatchObject({
agent: 'codex',
providerSession: { id: 's1' }
})
})
it('selects the WSL-relayed hook status for the same local pane', () => {
expect(selectThroughRuntime(wslHookRelayConnectionId('Ubuntu'))).toMatchObject({
agent: 'codex',
wslDistro: 'Ubuntu',
providerSession: { id: 's1' }
})
})
it('rejects a relay from a different distro', () => {
expect(selectThroughRuntime(wslHookRelayConnectionId('Debian'))).toBeNull()
})
})
@@ -143,21 +143,28 @@ export class OrcaRuntimeWithGetTerminalInteractiveWait extends OrcaRuntimeWithAd
}
let connectionId: string | null | undefined
let launchToken: string | null | undefined
let wslDistro: string | undefined
try {
const ptyId = this.getTerminalAgentStatusPtyId(handle)
const pty = this.ptysById.get(ptyId)
connectionId = pty?.connectionId ?? null
launchToken = pty?.launchToken ?? null
// A WSL pane's PTY is local, so its hook events only match once the distro is supplied.
wslDistro = pty?.connectionId
? undefined
: (this.wslDistroByPtyId.get(ptyId) ?? pty?.wslDistro ?? undefined)
} catch {
// Exact worker validation rejects this in production; test/legacy providers may not expose PTY metadata.
connectionId = undefined
launchToken = undefined
wslDistro = undefined
}
return selectExactWorkerProviderSession({
paneKey,
processIncarnation,
connectionId,
launchToken,
wslDistro,
observedAfter,
statuses: this.getAgentStatusSnapshotFn?.() ?? []
})
@@ -53,8 +53,6 @@ export function acknowledgeFederationRelay(
direction: FederationRelayDirection
throughSequence: number
settleRemoteReports?: { sequence: number; outcome?: WorkerReportOutcome }[]
/** Current runtime capability; omitted callers retain persisted-protocol behavior. */
supportsLifecycleSettlement?: boolean
}
): void {
this.db.exec('BEGIN IMMEDIATE')
@@ -83,9 +81,7 @@ export function acknowledgeFederationRelay(
if (
params.direction === 'to_home' &&
attachment !== undefined &&
(params.supportsLifecycleSettlement ??
attachment.protocol_version >=
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION)
attachment.protocol_version >= ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION
) {
const acknowledgedReports = this.db
.prepare(
@@ -0,0 +1,61 @@
import { describe, expect, it } from 'vitest'
import type { OrcaRuntimeService } from '../orca-runtime'
import {
acquireFederationAckLease,
clearFederationAckCheckpoints,
getFederationAckedThrough,
recordFederationAckCheckpoint,
type FederationAckIdentity
} from './federation-ack-checkpoints'
describe('federation acknowledgment checkpoints', () => {
it('matches checkpoints only to their exact remote identity and never moves backward', () => {
const runtime = {} as OrcaRuntimeService
const identity: FederationAckIdentity = {
environmentId: 'environment_windows',
peerFingerprint: 'windows_peer_fingerprint',
remoteRuntimeEpoch: 'remote_epoch_1'
}
const lease = acquireFederationAckLease(runtime, 'dispatch_remote')
recordFederationAckCheckpoint(runtime, lease, {
...identity,
throughSequence: 2
})
recordFederationAckCheckpoint(runtime, lease, {
...identity,
throughSequence: 3
})
recordFederationAckCheckpoint(runtime, lease, {
...identity,
throughSequence: 2
})
expect(getFederationAckedThrough(lease, identity)).toBe(3)
expect(
getFederationAckedThrough(lease, { ...identity, remoteRuntimeEpoch: 'remote_epoch_2' })
).toBe(0)
expect(
getFederationAckedThrough(lease, { ...identity, peerFingerprint: 'replacement_peer' })
).toBe(0)
expect(getFederationAckedThrough(lease, { ...identity, environmentId: 'replacement' })).toBe(0)
})
it('fences delayed writes after runtime reset', () => {
const runtime = {} as OrcaRuntimeService
const identity: FederationAckIdentity = {
environmentId: 'environment_windows',
peerFingerprint: 'windows_peer_fingerprint',
remoteRuntimeEpoch: 'remote_epoch_1'
}
const staleRuntimeLease = acquireFederationAckLease(runtime, 'dispatch_remote')
clearFederationAckCheckpoints(runtime)
recordFederationAckCheckpoint(runtime, staleRuntimeLease, {
...identity,
throughSequence: 2
})
expect(
getFederationAckedThrough(acquireFederationAckLease(runtime, 'dispatch_remote'), identity)
).toBe(0)
})
})
@@ -19,7 +19,6 @@ export async function resolveFederatedLifecycleSettlementCapability(
peerFingerprint: federated.peer_fingerprint,
expectedRuntimeEpoch: federated.remote_runtime_epoch,
capability: ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY,
forceProbe: true,
probe: () =>
runtime.callOrchestrationWorkerServer(
federated.environment_id,
@@ -0,0 +1,109 @@
import { vi } from 'vitest'
import { ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY } from '../../../shared/protocol-version'
import { OrcaRuntimeService } from '../orca-runtime'
export function createIdleSyncHarness(initialSequence = 2, protocolVersion?: 1 | 2 | 3) {
let remoteRuntimeEpoch = 'remote_epoch_1'
let remoteCapabilities: string[] = [
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY
]
let blockedAck: { reached: () => void; released: Promise<void> } | null = null
let blockedPull: { reached: () => void; released: Promise<void> } | null = null
let relayEligible = true
const federated = {
environment_id: 'environment_windows',
environment_name: 'windows',
peer_fingerprint: 'windows_peer_fingerprint',
remote_runtime_epoch: remoteRuntimeEpoch,
...(protocolVersion ? { protocol_version: protocolVersion } : {}),
to_home_imported_sequence: initialSequence,
to_home_acknowledged_sequence: 0
}
const createDb = () =>
({
getFederatedDispatch: () => federated,
getDispatchContextById: () => ({ run_id: 'run_home', task_id: 'task_home' }),
getWorkerDispatch: () => ({ state: 'ready' }),
listPendingFederationRelay: () => [],
isFederatedDispatchRelayEligible: () => relayEligible,
recordFederatedHomeAcknowledgment: (params: {
remoteRuntimeEpoch: string
sequence: number
}) => {
federated.remote_runtime_epoch = params.remoteRuntimeEpoch
federated.to_home_acknowledged_sequence = params.sequence
},
updateFederatedDispatchRuntimeEpoch: (_dispatchId: string, runtimeEpoch: string) => {
federated.remote_runtime_epoch = runtimeEpoch
}
}) as never
const runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(createDb())
vi.spyOn(runtime, 'resolveOrchestrationWorkerServer').mockReturnValue({
peerFingerprint: federated.peer_fingerprint
} as never)
const remoteCall = vi
.spyOn(runtime, 'callOrchestrationWorkerServer')
.mockImplementation(async (_environmentId, method) => {
if (method === 'orchestration.federationPull') {
const gate = blockedPull
if (gate) {
gate.reached()
await gate.released
if (blockedPull === gate) {
blockedPull = null
}
}
return { runtimeEpoch: remoteRuntimeEpoch, items: [] }
}
if (method === 'status.get') {
return { runtimeId: remoteRuntimeEpoch, capabilities: remoteCapabilities }
}
if (method === 'orchestration.federationAck') {
const gate = blockedAck
if (gate) {
gate.reached()
await gate.released
if (blockedAck === gate) {
blockedAck = null
}
}
return { acknowledgedThrough: federated.to_home_imported_sequence }
}
throw new Error(`Unexpected method ${method}`)
})
return {
runtime,
remoteCall,
advanceCursor: () => {
federated.to_home_imported_sequence += 1
},
restartRemote: () => {
remoteRuntimeEpoch = 'remote_epoch_2'
},
getPersistedRemoteRuntimeEpoch: () => federated.remote_runtime_epoch,
settleDispatch: () => {
relayEligible = false
},
setRemoteCapabilities: (capabilities: string[]) => {
remoteCapabilities = capabilities
},
replaceDb: () => runtime.setOrchestrationDb(createDb()),
blockAck: () => {
let noteReached!: () => void
let release!: () => void
const reached = new Promise<void>((resolve) => (noteReached = resolve))
const released = new Promise<void>((resolve) => (release = resolve))
blockedAck = { reached: noteReached, released }
return { reached, release }
},
blockPull: () => {
let noteReached!: () => void
let release!: () => void
const reached = new Promise<void>((resolve) => (noteReached = resolve))
const released = new Promise<void>((resolve) => (release = resolve))
blockedPull = { reached: noteReached, released }
return { reached, release }
}
}
}
@@ -1,126 +1,19 @@
import { describe, expect, it, vi } from 'vitest'
import {
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY
} from '../../../shared/protocol-version'
import { OrcaRuntimeService } from '../orca-runtime'
import { OrchestrationDb } from './db'
import {
acquireFederationAckLease,
clearFederationAckCheckpoints,
getFederationAckedThrough,
recordFederationAckCheckpoint,
type FederationAckIdentity
} from './federation-ack-checkpoints'
import { createIdleSyncHarness } from './federation-sync-test-harness'
import { parseRelayedMessage, syncFederatedDispatch } from './federation-sync'
import { getOrchestrationPeerCapabilityCache } from './orchestration-peer-capability-cache'
function createIdleSyncHarness(initialSequence = 2, protocolVersion?: 1 | 2 | 3) {
let remoteRuntimeEpoch = 'remote_epoch_1'
let remoteCapabilities: string[] = [
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY
]
let blockedAck: { reached: () => void; released: Promise<void> } | null = null
let blockedPull: { reached: () => void; released: Promise<void> } | null = null
let relayEligible = true
const federated = {
environment_id: 'environment_windows',
environment_name: 'windows',
peer_fingerprint: 'windows_peer_fingerprint',
remote_runtime_epoch: remoteRuntimeEpoch,
...(protocolVersion ? { protocol_version: protocolVersion } : {}),
to_home_imported_sequence: initialSequence,
to_home_acknowledged_sequence: 0
}
const createDb = () =>
({
getFederatedDispatch: () => federated,
getDispatchContextById: () => ({ run_id: 'run_home', task_id: 'task_home' }),
getWorkerDispatch: () => ({ state: 'ready' }),
listPendingFederationRelay: () => [],
isFederatedDispatchRelayEligible: () => relayEligible,
recordFederatedHomeAcknowledgment: (params: {
remoteRuntimeEpoch: string
sequence: number
}) => {
federated.remote_runtime_epoch = params.remoteRuntimeEpoch
federated.to_home_acknowledged_sequence = params.sequence
},
updateFederatedDispatchRuntimeEpoch: (_dispatchId: string, runtimeEpoch: string) => {
federated.remote_runtime_epoch = runtimeEpoch
}
}) as never
const runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(createDb())
vi.spyOn(runtime, 'resolveOrchestrationWorkerServer').mockReturnValue({
peerFingerprint: federated.peer_fingerprint
} as never)
const remoteCall = vi
.spyOn(runtime, 'callOrchestrationWorkerServer')
.mockImplementation(async (_environmentId, method) => {
if (method === 'orchestration.federationPull') {
const gate = blockedPull
if (gate) {
gate.reached()
await gate.released
if (blockedPull === gate) {
blockedPull = null
}
}
return { runtimeEpoch: remoteRuntimeEpoch, items: [] }
}
if (method === 'status.get') {
return { runtimeId: remoteRuntimeEpoch, capabilities: remoteCapabilities }
}
if (method === 'orchestration.federationAck') {
const gate = blockedAck
if (gate) {
gate.reached()
await gate.released
if (blockedAck === gate) {
blockedAck = null
}
}
return { acknowledgedThrough: federated.to_home_imported_sequence }
}
throw new Error(`Unexpected method ${method}`)
})
return {
runtime,
remoteCall,
advanceCursor: () => {
federated.to_home_imported_sequence += 1
},
restartRemote: () => {
remoteRuntimeEpoch = 'remote_epoch_2'
},
getPersistedRemoteRuntimeEpoch: () => federated.remote_runtime_epoch,
settleDispatch: () => {
relayEligible = false
},
setRemoteCapabilities: (capabilities: string[]) => {
remoteCapabilities = capabilities
},
replaceDb: () => runtime.setOrchestrationDb(createDb()),
blockAck: () => {
let noteReached!: () => void
let release!: () => void
const reached = new Promise<void>((resolve) => (noteReached = resolve))
const released = new Promise<void>((resolve) => (release = resolve))
blockedAck = { reached: noteReached, released }
return { reached, release }
},
blockPull: () => {
let noteReached!: () => void
let release!: () => void
const reached = new Promise<void>((resolve) => (noteReached = resolve))
const released = new Promise<void>((resolve) => (release = resolve))
blockedPull = { reached: noteReached, released }
return { reached, release }
}
}
}
describe('federation relay parsing', () => {
it('accepts a supported message type', () => {
expect(
@@ -388,30 +281,52 @@ describe('federation relay acknowledgments', () => {
await cache.resolve({
peerFingerprint: 'windows_peer_fingerprint',
expectedRuntimeEpoch: 'remote_epoch_1',
capability: ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY,
capability: ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY,
probe: vi.fn().mockResolvedValue({ runtimeId: 'remote_epoch_1', capabilities: [] })
})
restartRemote()
setRemoteCapabilities([
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY
])
await runtime.syncOrchestrationFederatedDispatch('dispatch_remote')
expect(getPersistedRemoteRuntimeEpoch()).toBe('remote_epoch_2')
await expect(
// The restart dropped the old epoch's answers, so the next resolve re-probes once and
// then serves the new epoch from cache.
const probe = vi.fn().mockResolvedValue({
runtimeId: 'remote_epoch_2',
capabilities: [ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY]
})
const resolveRelease = () =>
cache.resolve({
peerFingerprint: 'windows_peer_fingerprint',
expectedRuntimeEpoch: 'remote_epoch_1',
capability: ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY,
probe: vi.fn()
capability: ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY,
probe
})
).resolves.toMatchObject({
await expect(resolveRelease()).resolves.toMatchObject({
runtimeEpoch: 'remote_epoch_2',
supported: true,
cached: false
})
await expect(resolveRelease()).resolves.toMatchObject({
runtimeEpoch: 'remote_epoch_2',
supported: true,
cached: true
})
expect(probe).toHaveBeenCalledOnce()
})
it('probes an unchanged peer once across repeated syncs', async () => {
const { runtime, remoteCall } = createIdleSyncHarness(0)
await runtime.syncOrchestrationFederatedDispatch('dispatch_remote')
await runtime.syncOrchestrationFederatedDispatch('dispatch_remote')
await runtime.syncOrchestrationFederatedDispatch('dispatch_remote')
expect(remoteCall.mock.calls.filter(([, method]) => method === 'status.get')).toHaveLength(1)
})
it('does not wake a waiter for an acknowledged duplicate replay', async () => {
@@ -625,7 +540,6 @@ describe('federation relay acknowledgments', () => {
'orchestration.federationPull',
'orchestration.federationAck',
'orchestration.federationImport',
'status.get',
'orchestration.federationPull',
'orchestration.federationAck'
])
@@ -779,54 +693,4 @@ describe('federation relay acknowledgments', () => {
expect(ackCalls()).toHaveLength(1)
})
it('matches checkpoints only to their exact remote identity and never moves backward', () => {
const runtime = {} as OrcaRuntimeService
const identity: FederationAckIdentity = {
environmentId: 'environment_windows',
peerFingerprint: 'windows_peer_fingerprint',
remoteRuntimeEpoch: 'remote_epoch_1'
}
const lease = acquireFederationAckLease(runtime, 'dispatch_remote')
recordFederationAckCheckpoint(runtime, lease, {
...identity,
throughSequence: 2
})
recordFederationAckCheckpoint(runtime, lease, {
...identity,
throughSequence: 3
})
recordFederationAckCheckpoint(runtime, lease, {
...identity,
throughSequence: 2
})
expect(getFederationAckedThrough(lease, identity)).toBe(3)
expect(
getFederationAckedThrough(lease, { ...identity, remoteRuntimeEpoch: 'remote_epoch_2' })
).toBe(0)
expect(
getFederationAckedThrough(lease, { ...identity, peerFingerprint: 'replacement_peer' })
).toBe(0)
expect(getFederationAckedThrough(lease, { ...identity, environmentId: 'replacement' })).toBe(0)
})
it('fences delayed writes after runtime reset', () => {
const runtime = {} as OrcaRuntimeService
const identity: FederationAckIdentity = {
environmentId: 'environment_windows',
peerFingerprint: 'windows_peer_fingerprint',
remoteRuntimeEpoch: 'remote_epoch_1'
}
const staleRuntimeLease = acquireFederationAckLease(runtime, 'dispatch_remote')
clearFederationAckCheckpoints(runtime)
recordFederationAckCheckpoint(runtime, staleRuntimeLease, {
...identity,
throughSequence: 2
})
expect(
getFederationAckedThrough(acquireFederationAckLease(runtime, 'dispatch_remote'), identity)
).toBe(0)
})
})
@@ -1,3 +1,4 @@
import { z } from 'zod'
import type { OrcaRuntimeService } from '../orca-runtime'
import type { FederatedLifecycleSettlement } from './federation-lifecycle-settlement'
import { ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION } from '../../../shared/protocol-version'
@@ -16,14 +17,25 @@ export { parseRelayedMessage } from './federation-sync-message'
const FEDERATION_PULL_PAGE_SIZE = 50
const MAX_FEDERATION_PULL_PAGES_PER_SYNC = 6
type PulledRelayItem = {
dispatch_id: string
direction: 'to_home'
sequence: number
message_id: string
kind: string
payload: string
}
// Peer payloads are untrusted input: decode them so a malformed page fails as an
// orchestration error instead of a TypeError deep inside the import loop.
const PulledRelayPage = z
.object({
runtimeEpoch: z.string().min(1),
items: z.array(
z
.object({
dispatch_id: z.string(),
direction: z.literal('to_home'),
sequence: z.number(),
message_id: z.string(),
kind: z.string(),
payload: z.string()
})
.passthrough()
)
})
.passthrough()
export async function syncFederatedDispatch(
runtime: OrcaRuntimeService,
@@ -77,7 +89,7 @@ async function syncFederatedDispatchPages(
(federated.protocol_version >= ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION &&
(federated.to_home_acknowledged_sequence ?? 0) < federated.to_home_imported_sequence)
const pulled = (await runtime.callOrchestrationWorkerServer(
const pulledResponse = await runtime.callOrchestrationWorkerServer(
federated.environment_id,
'orchestration.federationPull',
{
@@ -89,7 +101,15 @@ async function syncFederatedDispatchPages(
15_000,
undefined,
{ expectedEnvironmentPairingRevision: currentServer.pairingRevision }
)) as { runtimeEpoch: string; items: PulledRelayItem[] }
)
const parsedPull = PulledRelayPage.safeParse(pulledResponse)
if (!parsedPull.success) {
throw new OrchestrationError(
'invalid_runtime_response',
`The execution host returned an invalid federation relay page for ${dispatchId}.`
)
}
const pulled = parsedPull.data
if (!isCurrent()) {
return { imported: 0, acknowledgedThrough: federated.to_home_imported_sequence }
}
@@ -51,8 +51,6 @@ export class OrchestrationPeerCapabilityCache {
expectedRuntimeEpoch: string | null
capability: RuntimeCapability
probe: () => Promise<RuntimeStatus>
/** Require a status probe even when the expected epoch has a cached answer. */
forceProbe?: boolean
}): Promise<PeerCapabilityDecision> {
return this.resolveAttempt(args, 1)
}
@@ -63,16 +61,14 @@ export class OrchestrationPeerCapabilityCache {
expectedRuntimeEpoch: string | null
capability: RuntimeCapability
probe: () => Promise<RuntimeStatus>
forceProbe?: boolean
},
staleRetriesRemaining: number
): Promise<PeerCapabilityDecision> {
const generation = this.touchPeer(args.peerFingerprint)
const knownEpoch = this.latestEpochs.get(args.peerFingerprint) ?? args.expectedRuntimeEpoch
const cached =
!args.forceProbe && knownEpoch
? this.cached(args.peerFingerprint, knownEpoch, args.capability)
: null
const cached = knownEpoch
? this.cached(args.peerFingerprint, knownEpoch, args.capability)
: null
if (cached) {
return cached
}
@@ -117,6 +113,24 @@ export class OrchestrationPeerCapabilityCache {
return { runtimeEpoch: status.runtimeId, supported, cached: false }
}
/**
* What the peer's own answers proved, or null when nothing has. Deliberately ignores the
* advertised capability list: shipped hosts serve federation methods they never advertise, so
* only a real `method_not_found` may downgrade one.
*/
knownSupport(
peerFingerprint: string,
expectedRuntimeEpoch: string | null,
capability: RuntimeCapability
): PeerCapabilityDecision | null {
const epoch = this.latestEpochs.get(peerFingerprint) ?? expectedRuntimeEpoch
const state = epoch ? this.states.get(this.key(peerFingerprint, epoch))?.get(capability) : null
if (!state || (!state.supported && (state.negativeExpiresAt ?? 0) <= this.now())) {
return null
}
return { runtimeEpoch: state.runtimeEpoch, supported: state.supported, cached: true }
}
remember(
peerFingerprint: string,
runtimeEpoch: string,
@@ -17,16 +17,6 @@ import { isWslHookRelayConnectionId } from '../../../shared/wsl-hook-relay-contr
// Bound the durable copy of raw terminal output; the tail end is the evidence that matters.
const TERMINAL_ARCHIVE_MAX_CHARS = 262_144
export type WorkerTranscriptPinArchive = {
agent: AgentType
providerSessionKey: string
providerSessionId: string
transcriptPath: string | null
processIncarnation: string
observedAfter: number
endOffset?: number
}
export type WorkerTranscriptSnapshotArchive = {
version: 2
agent: AgentType
@@ -40,17 +40,13 @@ type LocalTranscriptReadResult =
export async function readInitialLocalWorkerTranscriptPage(
filePath: string,
limit: number,
decode: NativeChatLineDecoder,
endOffset?: number
decode: NativeChatLineDecoder
): Promise<LocalTranscriptReadResult> {
const before = await readLocalTranscriptSourceIdentity(filePath)
if (!before) {
return { ok: false, reason: 'transcript_unreadable', warnings: [] }
}
if (endOffset !== undefined && before.size < endOffset) {
return sourceChanged()
}
const page = await readNativeChatTranscriptTailFile(filePath, limit, decode, false, endOffset)
const page = await readNativeChatTranscriptTailFile(filePath, limit, decode, false)
const after = await readLocalTranscriptSourceIdentity(filePath)
if (workerTranscriptSourceChanged(before, after, page.consumedTo)) {
return sourceChanged()
@@ -81,17 +77,13 @@ export async function readForwardLocalWorkerTranscriptPage(
startOffset: number,
limit: number,
decode: NativeChatLineDecoder,
endOffset?: number,
expectedBoundaryCheckpoint?: string
): Promise<LocalTranscriptReadResult> {
const sourceIdentity = await readLocalTranscriptSourceIdentity(filePath)
if (!sourceIdentity) {
return { ok: false, reason: 'transcript_unreadable', warnings: [] }
}
if (endOffset !== undefined && sourceIdentity.size < endOffset) {
return sourceChanged()
}
const fileSize = Math.min(sourceIdentity.size, endOffset ?? Number.MAX_SAFE_INTEGER)
const fileSize = sourceIdentity.size
if (startOffset > fileSize) {
return sourceChanged()
}
@@ -120,7 +112,6 @@ export async function readForwardLocalWorkerTranscriptPage(
startOffset,
scanEnd,
fileSize,
endOffset,
limit,
decode
})
@@ -131,7 +122,7 @@ export async function readForwardLocalWorkerTranscriptPage(
)
const handleAfter = localWorkerTranscriptSourceIdentity(await handle.stat({ bigint: true }))
const pathAfter = await readLocalTranscriptSourceIdentity(filePath)
const minimumSize = Math.max(page.nextOffset, endOffset ?? 0)
const minimumSize = page.nextOffset
return !afterCheckpoint ||
(expectedBoundaryCheckpoint !== undefined &&
afterCheckpoint !== expectedBoundaryCheckpoint) ||
@@ -152,7 +143,6 @@ async function scanForwardPage(args: {
startOffset: number
scanEnd: number
fileSize: number
endOffset?: number
limit: number
decode: NativeChatLineDecoder
}): Promise<LocalTranscriptPage> {
@@ -119,49 +119,6 @@ describe('worker transcript reads', () => {
).resolves.toEqual({ ok: false, reason: 'source_changed', warnings: [] })
})
it('pins archived reads to the transcript offset observed before release', async () => {
await writeFile(
transcriptPath,
`${codexMessage('one', 'before release')}\n${codexMessage('two', 'release boundary')}\n`
)
const snapshot = await readWorkerTranscript({
agent: 'codex',
sessionId: 'session-exact',
transcriptPath,
limit: 1
})
if (!snapshot.ok) {
throw new Error('Expected the release transcript probe')
}
await appendFile(transcriptPath, `${codexMessage('three', 'after release')}\n`)
await expect(
readWorkerTranscript({
agent: 'codex',
sessionId: 'session-exact',
transcriptPath,
endOffset: snapshot.nextOffset,
limit: 10
})
).resolves.toMatchObject({
ok: true,
messages: [
{ id: 'one', blocks: [{ type: 'text', text: 'before release' }] },
{ id: 'two', blocks: [{ type: 'text', text: 'release boundary' }] }
]
})
await expect(
readWorkerTranscript({
agent: 'codex',
sessionId: 'session-exact',
transcriptPath,
offset: snapshot.nextOffset,
endOffset: snapshot.nextOffset,
limit: 10
})
).resolves.toMatchObject({ ok: true, messages: [], nextOffset: snapshot.nextOffset })
})
it('reports source changes and unsupported providers without guessing', async () => {
await writeFile(transcriptPath, `${codexMessage('one', 'first')}\n`)
@@ -41,7 +41,6 @@ export async function readWorkerTranscript(args: {
/** Attested local WSL distro. Keeps host path translation on the selected guest. */
wslDistro?: string
offset?: number
endOffset?: number
limit?: number
/** Prior file identity from the cursor owner, when it retains that evidence. */
expectedSourceFingerprint?: string
@@ -92,13 +91,12 @@ export async function readWorkerTranscript(args: {
try {
const page =
args.offset === undefined
? await readInitialLocalWorkerTranscriptPage(filePath, limit, decode, args.endOffset)
? await readInitialLocalWorkerTranscriptPage(filePath, limit, decode)
: await readForwardLocalWorkerTranscriptPage(
filePath,
args.offset,
limit,
decode,
args.endOffset,
args.expectedBoundaryCheckpoint
)
if (!page.ok) {
@@ -4,7 +4,6 @@ import type { OrcaRuntimeService } from '../../orca-runtime'
import type { OrchestrationDb } from '../../orchestration/db'
import { OrchestrationError } from '../../orchestration/orchestration-error'
import type { FederatedDispatchRow } from '../../orchestration/types'
import { getOrchestrationPeerCapabilityCache } from '../../orchestration/orchestration-peer-capability-cache'
import { projectOrchestrationFleet } from '../../../../shared/orchestration-fleet-projection'
import {
applyFederatedFleetObservations,
@@ -61,96 +60,23 @@ describe('federated fleet snapshots', () => {
expect(result.observations).toHaveLength(101)
})
it('keeps a capability retry inside the fleet total deadline', async () => {
const dispatch = federatedDispatch('dispatch-deadline', 'peer-deadline', 'epoch-a')
it('asks the snapshot method directly instead of probing status.get', async () => {
const dispatch = federatedDispatch('dispatch-optimistic', 'peer-optimistic', 'epoch-a')
const db = {
getFederatedDispatch: () => dispatch,
updateFederatedDispatchRuntimeEpoch: vi.fn(),
...observationFenceMethods()
} as unknown as OrchestrationDb
const runtime = {
resolveOrchestrationWorkerServer: () => ({
environmentId: dispatch.environment_id,
name: dispatch.environment_name,
peerFingerprint: dispatch.peer_fingerprint,
pairingRevision: 1
}),
callOrchestrationWorkerServer: vi.fn()
} as unknown as OrcaRuntimeService
const cache = getOrchestrationPeerCapabilityCache(runtime)
let now = 1_000
const dateNow = vi.spyOn(Date, 'now').mockImplementation(() => now)
const timeouts: number[] = []
let statusCalls = 0
;(runtime.callOrchestrationWorkerServer as ReturnType<typeof vi.fn>).mockImplementation(
async (_environmentId: string, method: string, _params: unknown, timeoutMs: number) => {
timeouts.push(timeoutMs)
if (method === 'status.get') {
statusCalls += 1
if (statusCalls === 1) {
// A concurrent observer notices a restart while this probe is in flight.
cache.observeEpoch(dispatch.peer_fingerprint, 'epoch-b')
now += 4_900
return runtimeStatus('epoch-a')
}
return runtimeStatus('epoch-b')
}
return {
runtimeEpoch: 'epoch-b',
items: [
{
dispatchId: dispatch.dispatch_id,
observation: { status: 'live' as const, exactWorker: true }
}
]
}
}
)
try {
const result = await readFederatedFleetSnapshots({
runtime,
db,
dispatchIds: [dispatch.dispatch_id]
})
expect(result.observations.get(dispatch.dispatch_id)).toEqual({
status: 'live',
exactWorker: true
})
expect(statusCalls).toBe(2)
expect(timeouts).toEqual([3_000, 100, 100])
} finally {
dateNow.mockRestore()
}
})
it('does not grant a snapshot call budget after the fleet deadline expires', async () => {
const dispatch = federatedDispatch('dispatch-expired', 'peer-expired', 'epoch-a')
const db = {
getFederatedDispatch: () => dispatch,
updateFederatedDispatchRuntimeEpoch: vi.fn(),
...observationFenceMethods()
} as unknown as OrchestrationDb
const runtime = {
resolveOrchestrationWorkerServer: () => ({
environmentId: dispatch.environment_id,
name: dispatch.environment_name,
peerFingerprint: dispatch.peer_fingerprint,
pairingRevision: 1
}),
callOrchestrationWorkerServer: vi.fn()
} as unknown as OrcaRuntimeService
let now = 1_000
const dateNow = vi.spyOn(Date, 'now').mockImplementation(() => now)
const methods: string[] = []
;(runtime.callOrchestrationWorkerServer as ReturnType<typeof vi.fn>).mockImplementation(
async (_environmentId: string, method: string) => {
const runtime = {
resolveOrchestrationWorkerServer: () => ({
environmentId: dispatch.environment_id,
name: dispatch.environment_name,
peerFingerprint: dispatch.peer_fingerprint,
pairingRevision: 1
}),
callOrchestrationWorkerServer: vi.fn(async (_environmentId: string, method: string) => {
methods.push(method)
if (method === 'status.get') {
now += 5_001
return runtimeStatus('epoch-a')
}
return {
runtimeEpoch: 'epoch-a',
items: [
@@ -160,21 +86,72 @@ describe('federated fleet snapshots', () => {
}
]
}
}
})
} as unknown as OrcaRuntimeService
const result = await readFederatedFleetSnapshots({
runtime,
db,
dispatchIds: [dispatch.dispatch_id]
})
expect(methods).toEqual(['orchestration.federationFleetSnapshot'])
expect(result.observations.get(dispatch.dispatch_id)).toEqual({
status: 'live',
exactWorker: true
})
})
it('does not grant a snapshot call budget after the fleet deadline expires', async () => {
// Five distinct peers exceed the host concurrency, so the last one only starts after the
// first wave has already spent the whole fleet budget.
const dispatchIds = Array.from({ length: 5 }, (_, index) => `dispatch-expired-${index}`)
const dispatches = new Map(
dispatchIds.map((dispatchId) => [
dispatchId,
{
...federatedDispatch(dispatchId, `peer-${dispatchId}`, 'epoch-a'),
environment_id: dispatchId
}
])
)
const db = {
getFederatedDispatch: (dispatchId: string) => dispatches.get(dispatchId),
updateFederatedDispatchRuntimeEpoch: vi.fn(),
...observationFenceMethods()
} as unknown as OrchestrationDb
let now = 1_000
const dateNow = vi.spyOn(Date, 'now').mockImplementation(() => now)
const runtime = {
resolveOrchestrationWorkerServer: (environmentId: string) => ({
environmentId,
name: 'repointed',
peerFingerprint: `peer-${environmentId}`,
pairingRevision: 1
}),
callOrchestrationWorkerServer: vi.fn(
async (_environmentId: string, _method: string, params: unknown) => {
now += 5_001
return {
runtimeEpoch: 'epoch-a',
items: (params as { dispatchIds: string[] }).dispatchIds.map((dispatchId) => ({
dispatchId,
observation: { status: 'live' as const, exactWorker: true }
}))
}
}
)
} as unknown as OrcaRuntimeService
try {
const result = await readFederatedFleetSnapshots({
runtime,
db,
dispatchIds: [dispatch.dispatch_id]
})
const result = await readFederatedFleetSnapshots({ runtime, db, dispatchIds })
expect(result.observations).toEqual(new Map())
expect(result.errors).toEqual([
expect.objectContaining({ code: 'host_unavailable', dispatchIds: [dispatch.dispatch_id] })
])
expect(methods).toEqual(['status.get'])
// Orca never contacted these hosts, so calling them unavailable would fabricate a verdict.
expect(result.errors.length).toBeGreaterThan(0)
expect(result.errors.map((error) => error.code)).toEqual(
result.errors.map(() => 'home_budget_exhausted')
)
expect(result.errors.flatMap((error) => error.dispatchIds)).toContain('dispatch-expired-4')
} finally {
dateNow.mockRestore()
}
@@ -226,7 +203,7 @@ describe('federated fleet snapshots', () => {
expect(result.errors).toEqual([
expect.objectContaining({
environmentId: 'environment-repointed',
code: 'host_unavailable',
code: 'peer_changed',
dispatchIds: ['dispatch-a']
})
])
@@ -335,7 +312,7 @@ describe('federated fleet snapshots', () => {
expect(updateFederatedDispatchRuntimeEpoch).not.toHaveBeenCalled()
})
it('records a method-not-found result at the probed runtime epoch', async () => {
it('records a method-not-found result at the pinned runtime epoch', async () => {
const dispatch = federatedDispatch('dispatch-unsupported', 'peer-a', 'epoch-old')
const updateFederatedDispatchRuntimeEpoch = vi.fn()
const db = {
@@ -350,10 +327,7 @@ describe('federated fleet snapshots', () => {
peerFingerprint: dispatch.peer_fingerprint,
pairingRevision: 1
}),
callOrchestrationWorkerServer: vi.fn(async (_environmentId: string, method: string) => {
if (method === 'status.get') {
return runtimeStatus('epoch-new')
}
callOrchestrationWorkerServer: vi.fn(async () => {
throw new OrchestrationError('method_not_found', 'fleet snapshot unavailable')
})
} as unknown as OrcaRuntimeService
@@ -367,7 +341,7 @@ describe('federated fleet snapshots', () => {
expect(result.errors).toEqual([expect.objectContaining({ code: 'capability_unsupported' })])
expect(updateFederatedDispatchRuntimeEpoch).toHaveBeenCalledWith(
dispatch.dispatch_id,
'epoch-new'
'epoch-old'
)
})
})
@@ -1,6 +1,5 @@
import { mapWithConcurrency } from '../../../../shared/map-with-concurrency'
import { ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY } from '../../../../shared/protocol-version'
import type { RuntimeStatus } from '../../../../shared/runtime-types'
import {
ORCHESTRATION_FLEET_PAGE_MAX,
refreshOrchestrationFleetLivenessAttention,
@@ -26,7 +25,7 @@ export type FederatedFleetObservation = {
export type FederatedFleetHostError = {
environmentId: string
name: string
code: 'capability_unsupported' | 'host_unavailable'
code: 'capability_unsupported' | 'host_unavailable' | 'home_budget_exhausted' | 'peer_changed'
dispatchIds: string[]
}
@@ -63,7 +62,8 @@ export async function readFederatedFleetSnapshots(args: {
})
const remaining = deadline - Date.now()
if (remaining <= 0) {
return { observations: [], error: error('host_unavailable') }
// Orca never contacted this host; that is a home-side budget fact, not host silence.
return { observations: [], error: error('home_budget_exhausted') }
}
const timeoutMs = Math.min(FLEET_HOST_TIMEOUT_MS, remaining)
const first = group.dispatches[0]
@@ -71,29 +71,15 @@ export async function readFederatedFleetSnapshots(args: {
let observedCapabilityEpoch: string | null = null
try {
const server = resolvePinnedFederatedServer(args.runtime, first)
const capability = await cache.resolve({
peerFingerprint: first.peer_fingerprint,
expectedRuntimeEpoch: first.remote_runtime_epoch,
capability: ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY,
// Capability negotiation may retry after a peer restart; each probe
// must consume the same fleet deadline instead of restarting its budget.
probe: async () => {
const probeTimeoutMs = Math.min(FLEET_HOST_TIMEOUT_MS, deadline - Date.now())
if (probeTimeoutMs <= 0) {
throw new Error('Federated fleet deadline exceeded during capability negotiation')
}
return args.runtime.callOrchestrationWorkerServer(
server.environmentId,
'status.get',
undefined,
probeTimeoutMs,
undefined,
{ expectedEnvironmentPairingRevision: server.pairingRevision }
) as Promise<RuntimeStatus>
}
})
observedCapabilityEpoch = capability.runtimeEpoch
if (!capability.supported) {
// `method_not_found` is the single downgrade signal; a status probe would spend the
// fleet budget on a round trip that still misses hosts serving the unadvertised method.
const known = cache.knownSupport(
first.peer_fingerprint,
first.remote_runtime_epoch,
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY
)
observedCapabilityEpoch = known?.runtimeEpoch ?? first.remote_runtime_epoch
if (known?.supported === false) {
if (observedCapabilityEpoch) {
projectFleetRuntimeEpochs(args.db, observationFences, observedCapabilityEpoch)
}
@@ -101,7 +87,7 @@ export async function readFederatedFleetSnapshots(args: {
}
const snapshotRemainingMs = deadline - Date.now()
if (snapshotRemainingMs <= 0) {
return { observations: [], error: error('host_unavailable') }
return { observations: [], error: error('home_budget_exhausted') }
}
const snapshot = (await args.runtime.callOrchestrationWorkerServer(
server.environmentId,
@@ -155,7 +141,15 @@ export async function readFederatedFleetSnapshots(args: {
}
return { observations: [], error: error('capability_unsupported') }
}
return { observations: [], error: error('host_unavailable') }
// The environment was repointed at another Orca server; that is an identity fact, not silence.
return {
observations: [],
error: error(
caught instanceof OrchestrationError && caught.code === 'peer_changed'
? 'peer_changed'
: 'host_unavailable'
)
}
}
})
const observations = new Map<string, FederatedFleetObservation>()
@@ -203,7 +197,13 @@ export function applyFederatedFleetObservations(
federated: Awaited<ReturnType<typeof readFederatedFleetSnapshots>>,
observedAt = Date.now()
): void {
const unavailableDispatches = new Set(federated.errors.flatMap((error) => error.dispatchIds))
const unavailableDispatches = new Map(
federated.errors.flatMap((error) =>
error.dispatchIds.map(
(dispatchId) => [dispatchId, unavailableLivenessReason(error.code)] as const
)
)
)
for (const worker of fleet.workers) {
const hostId = federated.hosts.get(worker.dispatchId)
if (hostId) {
@@ -211,11 +211,12 @@ export function applyFederatedFleetObservations(
}
const observation = federated.observations.get(worker.dispatchId)
if (!observation) {
if (unavailableDispatches.has(worker.dispatchId)) {
const unavailableReason = unavailableDispatches.get(worker.dispatchId)
if (unavailableReason) {
if (worker.liveness.verdict === 'exited') {
continue
}
worker.liveness = { verdict: 'unverifiable', reason: 'host_unavailable' }
worker.liveness = { verdict: 'unverifiable', reason: unavailableReason }
worker.evidence.liveStatus = 'unavailable'
worker.evidence.lastObservedAt = null
refreshOrchestrationFleetLivenessAttention(worker)
@@ -230,13 +231,41 @@ export function applyFederatedFleetObservations(
? { verdict: 'live', observedAt, source: 'execution_host' }
: observation.status === 'exited'
? { verdict: 'exited', source: 'execution_host' }
: { verdict: 'unverifiable', reason: 'host_unavailable' }
: { verdict: 'unverifiable', reason: hostReportedReason(observation.reason) }
worker.evidence.liveStatus = observation.status === 'live' ? 'fresh' : 'unavailable'
worker.evidence.lastObservedAt = observation.status === 'live' ? observedAt : null
// The host answered for both `live` and `exited`, so both carry a real observation time.
worker.evidence.lastObservedAt = observation.status === 'unverifiable' ? null : observedAt
refreshOrchestrationFleetLivenessAttention(worker)
}
}
function unavailableLivenessReason(
code: FederatedFleetHostError['code']
): 'home_budget_exhausted' | 'peer_changed' | 'host_unavailable' {
return code === 'home_budget_exhausted' || code === 'peer_changed' ? code : 'host_unavailable'
}
const HOST_REPORTED_REASONS = new Set([
'missing_status',
'stale_status',
'future_status',
'restored_unconfirmed'
])
/** The host answered; contact was never lost, so never relabel its verdict as host_unavailable. */
function hostReportedReason(
reason: string | undefined
):
| 'host_indeterminate'
| 'missing_status'
| 'stale_status'
| 'future_status'
| 'restored_unconfirmed' {
return reason && HOST_REPORTED_REASONS.has(reason)
? (reason as 'missing_status' | 'stale_status' | 'future_status' | 'restored_unconfirmed')
: 'host_indeterminate'
}
function groupFederatedDispatches(args: {
runtime: OrcaRuntimeService
db: OrchestrationDb
@@ -81,6 +81,42 @@ describe('federated worker release ownership', () => {
expect(runtime.closeTerminal).not.toHaveBeenCalled()
})
it('closes an exited remote terminal before reporting closed_exited_terminal', async () => {
vi.mocked(runtime.showTerminal).mockResolvedValue({
handle: TERMINAL_HANDLE,
worktreeId: 'repo::remote',
connected: false,
status: 'exited'
} as never)
// Host-owned evidence: the execution host certifies this PTY exited.
vi.mocked(runtime.getTerminalLivenessVerdict).mockReturnValue({
status: 'exited',
ptyIds: [TERMINAL_HANDLE]
} as never)
vi.mocked(runtime.closeTerminal).mockResolvedValue({
handle: TERMINAL_HANDLE,
tabId: 'tab-remote',
ptyKilled: true
} as never)
vi.spyOn(runtime, 'readTerminal').mockResolvedValue({
handle: TERMINAL_HANDLE,
status: 'exited',
tail: ['worker output'],
truncated: false,
entries: [{ cursor: 1, text: 'worker output' }],
nextCursor: '1',
limited: false
} as never)
createAttachment('ctx_exited', 'created')
settleAttachment('ctx_exited')
await expect(call('orchestration.federationRelease', 'ctx_exited')).resolves.toMatchObject({
state: 'released',
processAction: 'closed_exited_terminal'
})
expect(runtime.closeTerminal).toHaveBeenCalledWith(TERMINAL_HANDLE)
})
it('fails closed for a settled legacy attachment without an ownership lease', async () => {
createAttachment('ctx_legacy')
settleAttachment('ctx_legacy')
@@ -3,7 +3,6 @@ import type {
OrchestrationWorkerReadResult
} from '../../../../shared/orchestration-worker-output'
import { ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY } from '../../../../shared/protocol-version'
import type { RuntimeStatus } from '../../../../shared/runtime-types'
import type { OrchestrationDb } from '../../orchestration/db'
import { OrchestrationError } from '../../orchestration/orchestration-error'
import { getOrchestrationPeerCapabilityCache } from '../../orchestration/orchestration-peer-capability-cache'
@@ -30,24 +29,18 @@ export async function readFederatedWorkerOutput(args: {
)
}
const capabilities = getOrchestrationPeerCapabilityCache(args.runtime)
const capability = await capabilities.resolve({
peerFingerprint: args.federated.peer_fingerprint,
expectedRuntimeEpoch: args.federated.remote_runtime_epoch,
capability: ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY,
probe: () =>
args.runtime.callOrchestrationWorkerServer(
args.server.environmentId,
'status.get',
undefined,
15_000,
undefined,
{ expectedEnvironmentPairingRevision: args.server.pairingRevision }
) as Promise<RuntimeStatus>
})
if (!capability.supported) {
// Hosts that serve `orchestration.federationReadOutput` shipped before the capability string
// did, so ask the method itself and let `method_not_found` be the only downgrade signal.
const known = capabilities.knownSupport(
args.federated.peer_fingerprint,
args.federated.remote_runtime_epoch,
ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY
)
const expectedRuntimeEpoch = known?.runtimeEpoch ?? args.federated.remote_runtime_epoch
if (known?.supported === false) {
const legacy = await readLegacy(args)
projectRemoteRuntimeEpoch(args.db, observationFence, legacy.remoteRuntimeEpoch)
if (legacy.remoteRuntimeEpoch !== capability.runtimeEpoch) {
if (legacy.remoteRuntimeEpoch !== expectedRuntimeEpoch) {
capabilities.observeEpoch(args.federated.peer_fingerprint, legacy.remoteRuntimeEpoch)
}
return legacy
@@ -82,17 +75,14 @@ export async function readFederatedWorkerOutput(args: {
if (!(error instanceof OrchestrationError) || error.code !== 'method_not_found') {
throw error
}
const legacy = await readLegacy(args)
capabilities.remember(
args.federated.peer_fingerprint,
capability.runtimeEpoch,
legacy.remoteRuntimeEpoch,
ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY,
false
)
const legacy = await readLegacy(args)
projectRemoteRuntimeEpoch(args.db, observationFence, legacy.remoteRuntimeEpoch)
if (legacy.remoteRuntimeEpoch !== capability.runtimeEpoch) {
capabilities.observeEpoch(args.federated.peer_fingerprint, legacy.remoteRuntimeEpoch)
}
return legacy
}
}
@@ -207,42 +207,44 @@ export async function releaseRemoteAttachment(args: {
output
}
}
if (observation.status !== 'exited') {
try {
const close = await runtime.closeTerminal(observation.terminal.handle)
if (!close.ptyKilled) {
const reason = describeUnconfirmedAgentStop(close)
return {
dispatchId: attachment.dispatch_id,
state: 'release_unknown',
processAction: 'closed_agent_terminal',
lastError: reason,
recovery: releaseUnknownRecovery(attachment.dispatch_id),
archive: archiveSummary(db.markWorkerTerminalReleaseUnknown(resource.id, reason)),
output: projectArchivedOutputLiveness(
output,
close.ptyStopVerdict === 'live' ? 'live' : 'unverifiable'
)
}
}
} catch (error) {
const closeError = classifyWorkerTerminalCloseError(error)
// An exited worker still owns a terminal record and tab on the host; close it before
// reporting `closed_exited_terminal`, exactly as the local release path does.
try {
const close = await runtime.closeTerminal(observation.terminal.handle)
// A host-certified exit already proved the process is gone, so a kill that stops nothing
// is not new doubt; anything else that survives the close still is.
if (!close.ptyKilled && observation.status !== 'exited') {
const reason = describeUnconfirmedAgentStop(close)
return {
dispatchId: attachment.dispatch_id,
state: closeError.transient ? 'release_pending' : 'release_unknown',
processAction: 'none',
lastError: closeError.reason,
recovery: closeError.transient
? TRANSIENT_WORKER_RELEASE_RECOVERY
: releaseUnknownRecovery(attachment.dispatch_id),
archive: archiveSummary(
closeError.transient
? releasing
: db.markWorkerTerminalReleaseUnknown(resource.id, closeError.reason)
),
output: projectArchivedOutputLiveness(output, 'unverifiable')
state: 'release_unknown',
processAction: 'closed_agent_terminal',
lastError: reason,
recovery: releaseUnknownRecovery(attachment.dispatch_id),
archive: archiveSummary(db.markWorkerTerminalReleaseUnknown(resource.id, reason)),
output: projectArchivedOutputLiveness(
output,
close.ptyStopVerdict === 'live' ? 'live' : 'unverifiable'
)
}
}
} catch (error) {
const closeError = classifyWorkerTerminalCloseError(error)
return {
dispatchId: attachment.dispatch_id,
state: closeError.transient ? 'release_pending' : 'release_unknown',
processAction: 'none',
lastError: closeError.reason,
recovery: closeError.transient
? TRANSIENT_WORKER_RELEASE_RECOVERY
: releaseUnknownRecovery(attachment.dispatch_id),
archive: archiveSummary(
closeError.transient
? releasing
: db.markWorkerTerminalReleaseUnknown(resource.id, closeError.reason)
),
output: projectArchivedOutputLiveness(output, 'unverifiable')
}
}
const released = db.settleWorkerTerminalRelease(resource.id)
db.recordRemoteAttachmentStage({
@@ -46,6 +46,8 @@ export async function releaseFederatedWorker(args: {
requestId: string
}): Promise<WorkerReleaseReceipt & { remoteOutput?: unknown }> {
const cache = getOrchestrationPeerCapabilityCache(args.runtime)
// This capability states that the host writes a durable archive before it closes anything;
// `method_not_found` cannot express that, so release still asks the advertisement.
const capability = await cache.resolve({
peerFingerprint: args.federated.peer_fingerprint,
expectedRuntimeEpoch: args.federated.remote_runtime_epoch,
@@ -248,7 +248,8 @@ describe('federation host liveness verdicts', () => {
archive: { source: 'terminal', status: 'empty' }
})
expect(host.hostDb.getWorkerTerminalArchive(DISPATCH_ID)).toBeDefined()
expect(closeTerminal).not.toHaveBeenCalled()
// The exited worker still owns a terminal record and tab; release must close it.
expect(closeTerminal).toHaveBeenCalledOnce()
} finally {
host.hostDb.close()
}
@@ -8,7 +8,6 @@ import {
ORCHESTRATION_CONTRACT_VERSION,
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_ARCHIVE_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY
} from '../../../../shared/protocol-version'
import { OrcaRuntimeService } from '../../orca-runtime'
@@ -78,7 +77,7 @@ describe('orchestration federated worker output', () => {
[
ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY
].includes(capability as never)) ||
((!workerAdvertisesNewCapabilities || !workerAdvertisesDurableRelease) &&
capability === ORCHESTRATION_FEDERATION_RELEASE_ARCHIVE_RUNTIME_CAPABILITY)
@@ -95,6 +94,17 @@ describe('orchestration federated worker output', () => {
error: { code: 'method_not_found', message: `Unknown method: ${method}` }
}
}
if (
method === 'orchestration.federationFleetSnapshot' &&
!workerAdvertisesNewCapabilities
) {
// A host old enough to lack the capability lacks the method too.
return {
id: `remote_${method}`,
ok: false,
error: { code: 'method_not_found', message: `Unknown method: ${method}` }
}
}
if (method === 'orchestration.federationFleetSnapshot' && workerFleetUnavailable) {
return {
id: `remote_${method}`,
@@ -687,6 +697,9 @@ describe('orchestration federated worker output', () => {
it('keeps reads, fleet snapshots, and release on legacy fallbacks for an old peer', async () => {
workerAdvertisesNewCapabilities = false
// A shipped host that does not advertise structured read still has to be asked; only its
// own method_not_found may downgrade the read to a terminal scrape.
workerSupportsStructuredRead = false
const dispatchId = await startRemoteWorker()
remoteCalls = []
const read = await homeDispatcher.dispatch({
@@ -722,8 +735,8 @@ describe('orchestration federated worker output', () => {
ok: true,
result: { state: 'retained', reason: 'federation_unsupported' }
})
expect(remoteCalls).not.toContain('orchestration.federationReadOutput')
expect(remoteCalls).not.toContain('orchestration.federationFleetSnapshot')
expect(remoteCalls).toContain('orchestration.federationReadOutput')
expect(remoteCalls).toContain('orchestration.federationFleetSnapshot')
expect(remoteCalls).not.toContain('orchestration.federationRelease')
})
@@ -761,7 +774,8 @@ describe('orchestration federated worker output', () => {
method: 'orchestration.workerRead',
params: { dispatch: dispatchId }
})
expect(remoteCalls).not.toContain('orchestration.federationReadOutput')
// The unadvertised capability never blocks the call; only method_not_found would.
expect(remoteCalls).toContain('orchestration.federationReadOutput')
const oldEpoch = homeDb.getFederatedDispatch(dispatchId)?.remote_runtime_epoch
expect(oldEpoch).toBe(workerRuntime.getRuntimeId())
@@ -789,6 +803,8 @@ describe('orchestration federated worker output', () => {
method: 'orchestration.workerList',
params: { includeRemote: true }
})
// Read and fleet negotiate through the methods themselves, so neither spends a probe.
expect(remoteCalls.filter((method) => method === 'status.get')).toHaveLength(0)
const release = await homeDispatcher.dispatch({
id: 'rpc_restarted_peer_release',
authToken: 'coordinator-token',
@@ -801,9 +817,7 @@ describe('orchestration federated worker output', () => {
expect(read).toMatchObject({ ok: true, result: { source: 'terminal' } })
expect(fleet).toMatchObject({ ok: true })
expect(release).toMatchObject({ ok: true, result: { state: 'released' } })
// One forced probe after the restart seeds all capability decisions for the
// new epoch; read, fleet, and release reuse that result.
expect(remoteCalls.filter((method) => method === 'status.get')).toHaveLength(1)
// Only release still probes: its capability asserts a durable archive, not method existence.
expect(remoteCalls).toContain('orchestration.federationReadOutput')
expect(remoteCalls).toContain('orchestration.federationFleetSnapshot')
expect(remoteCalls).toContain('orchestration.federationRelease')
@@ -1,8 +1,5 @@
import { z } from 'zod'
import {
ORCHESTRATION_FEDERATION_CONTROL_MAIL_PROTOCOL_VERSION,
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY
} from '../../../../shared/protocol-version'
import { ORCHESTRATION_FEDERATION_CONTROL_MAIL_PROTOCOL_VERSION } from '../../../../shared/protocol-version'
import { importFederatedControlMessage } from '../../orchestration/federation-control-message'
import { OrchestrationError } from '../../orchestration/orchestration-error'
import {
@@ -84,11 +81,7 @@ export const ORCHESTRATION_FEDERATION_RELAY_METHODS: RpcMethod[] = [
name: 'orchestration.federationAck',
params: FederationAckParams,
handler: (params, { runtime, authenticatedCallerFingerprint }) => {
const attachment = requireHomeAttachment(
runtime,
params.dispatchId,
authenticatedCallerFingerprint
)
requireHomeAttachment(runtime, params.dispatchId, authenticatedCallerFingerprint)
const receivedSettlements = (params.settlements ?? []).filter(
(settlement) => settlement.sequence <= params.throughSequence
)
@@ -124,13 +117,6 @@ export const ORCHESTRATION_FEDERATION_RELAY_METHODS: RpcMethod[] = [
dispatchId: params.dispatchId,
direction: 'to_home',
throughSequence: params.throughSequence,
supportsLifecycleSettlement:
attachment.protocol_version >= 3 &&
runtime
.getStatus()
.capabilities?.includes(
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY
) === true,
...(settlements.length === 0
? {}
: {
@@ -3,10 +3,7 @@ import type { OrcaRuntimeService } from '../../orca-runtime'
import { OrchestrationError } from '../../orchestration/orchestration-error'
import { waitForFederatedLifecycleSettlement } from '../../orchestration/federation-lifecycle-settlement'
import { bindCoordinatorMutationPayload } from '../../orchestration/dispatch-message-binding'
import {
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION,
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY
} from '../../../../shared/protocol-version'
import { ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION } from '../../../../shared/protocol-version'
import type { z } from 'zod'
import { parseRemoteWorkerPayload } from './orchestration-schemas'
import type { SendParams } from './orchestration-schemas'
@@ -70,11 +67,7 @@ export async function sendRemoteMessage(args: {
const supportsLifecycleSettlement =
remoteAttachment.protocol_version >=
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION &&
runtime
.getStatus()
.capabilities?.includes(ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_RUNTIME_CAPABILITY) ===
true
ORCHESTRATION_FEDERATION_LIFECYCLE_SETTLEMENT_PROTOCOL_VERSION
const relay = db.enqueueFederationRelay({
dispatchId: remoteAttachment.dispatch_id,
direction: 'to_home',
@@ -11,7 +11,6 @@ import type {
} from '../../orchestration/worker-terminal-ownership'
import type {
WorkerTerminalTailArchive,
WorkerTranscriptPinArchive,
WorkerTranscriptSnapshotArchive
} from '../../orchestration/worker-output-archive'
import { clampWorkerTranscriptLimit } from '../../orchestration/worker-transcript-payload'
@@ -20,13 +19,12 @@ import {
decodeWorkerOutputCursor,
encodeWorkerOutputCursor
} from '../../orchestration/worker-output-cursor'
import { readWorkerTranscript } from '../../orchestration/worker-transcript-read'
const ARCHIVED_TERMINAL_PAGE_LINES = 2_000
// Serves the frozen output source after the live PTY is gone. Transcript pins read the exact
// provider transcript directly; terminal archives page the stored redacted tail. Cursors stay
// Dispatch-scoped and source-pinned exactly like live reads.
// Serves the frozen output source after the live PTY is gone: a decoded transcript snapshot or
// the stored redacted terminal tail. Cursors stay Dispatch-scoped and source-pinned exactly like
// live reads.
export async function readArchivedWorkerOutput(args: {
db: OrchestrationDb
dispatchId: string
@@ -52,12 +50,11 @@ export async function readArchivedWorkerOutput(args: {
`Dispatch ${args.dispatchId} preserved structured transcript output only; terminal output was released.`
)
}
const content = JSON.parse(archive.content) as
| WorkerTranscriptPinArchive
| WorkerTranscriptSnapshotArchive
return isTranscriptSnapshot(content)
? readFrozenTranscript(args, archive, content)
: readLegacyPinnedTranscript(args, content)
return readFrozenTranscript(
args,
archive,
JSON.parse(archive.content) as WorkerTranscriptSnapshotArchive
)
}
if (args.source === 'transcript') {
throw new OrchestrationError(
@@ -119,82 +116,6 @@ function readFrozenTranscript(
}
}
async function readLegacyPinnedTranscript(
args: Parameters<typeof readArchivedWorkerOutput>[0],
pin: WorkerTranscriptPinArchive
): Promise<OrchestrationWorkerReadResult> {
const cursor = decodeWorkerOutputCursor(args.cursor, args.dispatchId)
if (cursor && (cursor.source !== 'transcript' || !cursor.boundaryCheckpoint)) {
throw sourceChanged()
}
const transcript = await readWorkerTranscript({
agent: pin.agent,
sessionId: pin.providerSessionId,
transcriptPath: pin.transcriptPath ?? undefined,
offset: cursor?.position,
endOffset: pin.endOffset,
expectedBoundaryCheckpoint: cursor?.boundaryCheckpoint ?? undefined,
limit: args.limit
})
if (!transcript.ok) {
if (transcript.reason === 'source_changed') {
throw sourceChanged()
}
throw new OrchestrationError(
'transcript_required',
`The pinned transcript for released Dispatch ${args.dispatchId} is unavailable: ${transcript.reason}.`,
{ reason: transcript.reason }
)
}
const sourceIdentity = createWorkerOutputSourceIdentity([
'released-transcript',
pin.processIncarnation,
pin.agent,
pin.providerSessionKey,
pin.providerSessionId,
pin.transcriptPath ?? '',
String(pin.endOffset),
transcript.sourceFingerprint
])
if (cursor && cursor.sourceIdentity !== sourceIdentity) {
throw sourceChanged()
}
const nextCursor = encodeWorkerOutputCursor(
args.dispatchId,
'transcript',
sourceIdentity,
transcript.nextOffset,
transcript.boundaryCheckpoint
)
const status = archivedStatus(args)
return {
dispatchId: args.dispatchId,
source: 'transcript',
sourceIdentity,
provider: pin.agent,
transcript: {
messages: transcript.messages,
nextCursor,
limited: transcript.limited,
returnedMessageCount: transcript.messages.length
},
cursor: nextCursor,
status,
fallbackReason: null,
sourceExact: true,
contentComplete: !transcript.limited,
...(transcript.clipping.length > 0 ? { clipping: transcript.clipping } : {}),
warnings: transcript.warnings,
archived: true
}
}
function isTranscriptSnapshot(
content: WorkerTranscriptPinArchive | WorkerTranscriptSnapshotArchive
): content is WorkerTranscriptSnapshotArchive {
return 'version' in content && content.version === 2
}
function readArchivedTerminalTail(
args: Parameters<typeof readArchivedWorkerOutput>[0],
archive: WorkerTerminalArchiveRow
@@ -1,5 +1,4 @@
import { afterEach, describe, expect, it, vi } from 'vitest'
import { ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY } from '../../../../shared/protocol-version'
import type Database from '../../../sqlite/sync-database'
import { OrchestrationDb } from '../../orchestration/db'
import type { FederatedDispatchRow } from '../../orchestration/types'
@@ -180,17 +179,15 @@ describe('orchestration worker-list pagination', () => {
peerFingerprint: 'peer-remote',
pairingRevision: 1
})
let resolveStatus!: (status: ReturnType<typeof fleetRuntimeStatus>) => void
const status = new Promise<ReturnType<typeof fleetRuntimeStatus>>((resolve) => {
resolveStatus = resolve
let resolveSnapshot!: () => void
const snapshotGate = new Promise<void>((resolve) => {
resolveSnapshot = resolve
})
const remoteCall = vi
.spyOn(runtime, 'callOrchestrationWorkerServer')
.mockImplementation((_environmentId, method) => {
if (method === 'status.get') {
return status
}
return Promise.resolve({
.mockImplementation(async () => {
await snapshotGate
return {
runtimeEpoch: 'epoch-remote',
items: [
{
@@ -198,7 +195,7 @@ describe('orchestration worker-list pagination', () => {
observation: { status: 'live', exactWorker: true }
}
]
})
}
})
const pending = callWorkerList(runtime, {
@@ -210,8 +207,8 @@ describe('orchestration worker-list pagination', () => {
await vi.waitFor(() =>
expect(remoteCall).toHaveBeenCalledWith(
'environment-remote',
'status.get',
undefined,
'orchestration.federationFleetSnapshot',
{ dispatchIds: ['dispatch-a'] },
expect.any(Number),
undefined,
{ expectedEnvironmentPairingRevision: 1 }
@@ -220,7 +217,7 @@ describe('orchestration worker-list pagination', () => {
for (let call = 0; call < 32; call += 1) {
await callWorkerList(runtime, { run: run.id, terminalState: 'retained', limit: 1 })
}
resolveStatus(fleetRuntimeStatus())
resolveSnapshot()
const first = await pending
expect(first).toMatchObject({
@@ -520,15 +517,3 @@ function federatedDispatch(dispatchId: string): FederatedDispatchRow {
updated_at: '2026-08-27 00:00:00'
}
}
function fleetRuntimeStatus() {
return {
runtimeId: 'epoch-remote',
capabilities: [ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY],
rendererGraphEpoch: 0,
graphStatus: 'ready' as const,
authoritativeWindowId: null,
liveTabCount: 0,
liveLeafCount: 0
}
}
@@ -1,8 +1,7 @@
import { appendFile, mkdtemp, rm, stat, writeFile } from 'node:fs/promises'
import { mkdtemp, rm, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { encodeWorkerOutputCursor } from '../../orchestration/worker-output-cursor'
import { createOrchestrationWorkerReleaseHarness } from './orchestration-worker-release.test-support'
function codexMessage(id: string, text: string): string {
@@ -157,80 +156,6 @@ describe('orchestration worker release archive', () => {
}
})
it('pins legacy transcript cursors to their content boundary', async () => {
h.setup()
const directory = await mkdtemp(join(tmpdir(), 'orca-worker-release-legacy-pin-'))
const transcriptPath = join(directory, 'rollout.jsonl')
const original = `${codexMessage('one', 'first')}\n${codexMessage('two', 'second')}\n`
try {
await writeFile(transcriptPath, original)
const { dispatchId } = await h.startSettledWorker()
await h.call('orchestration.workerRelease', { dispatch: dispatchId })
const resource = h.db.getWorkerTerminalResourceByOwner(dispatchId)
if (!resource) {
throw new Error('Expected a released worker terminal resource')
}
h.db.storeWorkerTerminalArchive({
dispatchId,
resourceId: resource.id,
kind: 'transcript_pin',
content: JSON.stringify({
agent: 'codex',
providerSessionKey: 'codex:legacy-pin',
providerSessionId: 'legacy-pin',
transcriptPath,
processIncarnation: resource.process_incarnation,
observedAfter: 0,
endOffset: Buffer.byteLength(original)
})
})
const page = (await h.call('orchestration.workerRead', {
dispatch: dispatchId,
limit: 1
})) as {
sourceIdentity: string
transcript: { messages: { id: string }[] }
cursor: string
}
expect(page.transcript.messages).toMatchObject([{ id: 'two' }])
const checkpointFreeCursor = encodeWorkerOutputCursor(
dispatchId,
'transcript',
page.sourceIdentity,
Buffer.byteLength(original)
)
await expect(
h.call('orchestration.workerRead', { dispatch: dispatchId, cursor: checkpointFreeCursor })
).rejects.toMatchObject({
code: 'source_changed',
message:
'The worker output source changed. Start a fresh worker-read without the old cursor.'
})
await appendFile(transcriptPath, `${codexMessage('three', 'after release')}\n`)
await expect(
h.call('orchestration.workerRead', { dispatch: dispatchId, cursor: page.cursor })
).resolves.toMatchObject({ transcript: { messages: [] }, contentComplete: true })
const before = await stat(transcriptPath, { bigint: true })
const replacement = `${codexMessage('red', 'other')}\n${codexMessage('new', 'REWRIT')}\n`
expect(Buffer.byteLength(replacement)).toBe(Buffer.byteLength(original))
await writeFile(transcriptPath, replacement)
const after = await stat(transcriptPath, { bigint: true })
expect(after.dev).toBe(before.dev)
expect(after.ino).toBe(before.ino)
expect(Number(after.size)).toBeGreaterThanOrEqual(Buffer.byteLength(original))
await expect(
h.call('orchestration.workerRead', { dispatch: dispatchId, cursor: page.cursor })
).rejects.toMatchObject({ code: 'source_changed' })
} finally {
await rm(directory, { recursive: true, force: true })
}
})
it('preserves payload clipping metadata in the released transcript snapshot', async () => {
h.setup()
const directory = await mkdtemp(join(tmpdir(), 'orca-worker-release-clipped-snapshot-'))
@@ -58,6 +58,12 @@ export type FleetLiveness =
| 'future_status'
| 'restored_unconfirmed'
| 'host_unavailable'
/** Orca's own fleet budget ran out before it asked the host anything. */
| 'home_budget_exhausted'
/** The host answered and could not tell; contact was never lost. */
| 'host_indeterminate'
/** The saved environment now identifies a different Orca server. */
| 'peer_changed'
observedAt?: number
}
| { verdict: 'exited'; source: 'resource_release' | 'worker_stop' | 'execution_host' }
-3
View File
@@ -56,8 +56,6 @@ export const ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY =
'orchestration.federation-structured-read.v1' as const
export const ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY =
'orchestration.federation-fleet-snapshot.v1' as const
export const ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY =
'orchestration.federation-release.v1' as const
export const ORCHESTRATION_FEDERATION_RELEASE_ARCHIVE_RUNTIME_CAPABILITY =
'orchestration.federation-release-archive.v1' as const
export const ORCHESTRATION_FEDERATION_CONTROL_MAIL_PROTOCOL_VERSION = 2 as const
@@ -202,7 +200,6 @@ export const RUNTIME_CAPABILITIES = [
ORCHESTRATION_WORKER_LAUNCH_PREFERENCES_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_RUNTIME_CAPABILITY,
ORCHESTRATION_FEDERATION_RELEASE_ARCHIVE_RUNTIME_CAPABILITY,
ORCHESTRATION_CONTRACT_RUNTIME_CAPABILITY,
BROWSER_SCREENCAST_RUNTIME_CAPABILITY,