From 58efd45c3308e12dfdc1c37ca41a88033ee2c1d1 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Fri, 4 Sep 2026 02:24:03 -0400 Subject: [PATCH 1/2] fix(federation): negotiate structured read by method, not advertisement - close an exited remote worker's terminal before labelling the release closed_exited_terminal; a host-certified exit keeps the verdict when the kill stops nothing - replace the structured-read and fleet-snapshot capability probes with the optimistic call plus method_not_found, so hosts that serve federationReadOutput without advertising it stop downgrading to a scrape - drop forceProbe so an unchanged peer at an unchanged epoch probes once - distinct fleet reasons for home-side budget exhaustion and peer_changed; keep a host-supplied unverifiable reason and lastObservedAt for exited - delete three compile-time-true self-capability checks, the advertised but never-read federation-release capability, and decode the pull page with zod --- .../db/federation/federation-relay-ack.ts | 6 +- .../federation-ack-checkpoints.test.ts | 61 ++++++ .../federation-sync-capability.ts | 1 - .../federation-sync-test-harness.ts | 109 ++++++++++ .../orchestration/federation-sync.test.ts | 196 +++--------------- .../runtime/orchestration/federation-sync.ts | 40 +++- .../orchestration-peer-capability-cache.ts | 28 ++- ...estration-federated-fleet-snapshot.test.ts | 178 +++++++--------- .../orchestration-federated-fleet-snapshot.ts | 95 ++++++--- ...estration-federated-release-safety.test.ts | 36 ++++ .../orchestration-federated-worker-read.ts | 34 ++- ...estration-federated-worker-release-host.ts | 66 +++--- .../orchestration-federated-worker-release.ts | 2 + ...ration-federation-liveness-verdict.test.ts | 3 +- .../orchestration-federation-output.test.ts | 30 ++- .../methods/orchestration-federation-relay.ts | 18 +- .../rpc/methods/orchestration-send-remote.ts | 11 +- ...chestration-worker-list-pagination.test.ts | 35 +--- src/shared/orchestration-fleet-projection.ts | 6 + src/shared/protocol-version.ts | 3 - 20 files changed, 518 insertions(+), 440 deletions(-) create mode 100644 src/main/runtime/orchestration/federation-ack-checkpoints.test.ts create mode 100644 src/main/runtime/orchestration/federation-sync-test-harness.ts diff --git a/src/main/runtime/orchestration/db/federation/federation-relay-ack.ts b/src/main/runtime/orchestration/db/federation/federation-relay-ack.ts index 8c7f3dcbb10..37ebd8b2b6f 100644 --- a/src/main/runtime/orchestration/db/federation/federation-relay-ack.ts +++ b/src/main/runtime/orchestration/db/federation/federation-relay-ack.ts @@ -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( diff --git a/src/main/runtime/orchestration/federation-ack-checkpoints.test.ts b/src/main/runtime/orchestration/federation-ack-checkpoints.test.ts new file mode 100644 index 00000000000..984f3d2c075 --- /dev/null +++ b/src/main/runtime/orchestration/federation-ack-checkpoints.test.ts @@ -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) + }) +}) diff --git a/src/main/runtime/orchestration/federation-sync-capability.ts b/src/main/runtime/orchestration/federation-sync-capability.ts index 9eda17900e2..57dd583dab2 100644 --- a/src/main/runtime/orchestration/federation-sync-capability.ts +++ b/src/main/runtime/orchestration/federation-sync-capability.ts @@ -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, diff --git a/src/main/runtime/orchestration/federation-sync-test-harness.ts b/src/main/runtime/orchestration/federation-sync-test-harness.ts new file mode 100644 index 00000000000..f41c4e0b83e --- /dev/null +++ b/src/main/runtime/orchestration/federation-sync-test-harness.ts @@ -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 } | null = null + let blockedPull: { reached: () => void; released: Promise } | 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((resolve) => (noteReached = resolve)) + const released = new Promise((resolve) => (release = resolve)) + blockedAck = { reached: noteReached, released } + return { reached, release } + }, + blockPull: () => { + let noteReached!: () => void + let release!: () => void + const reached = new Promise((resolve) => (noteReached = resolve)) + const released = new Promise((resolve) => (release = resolve)) + blockedPull = { reached: noteReached, released } + return { reached, release } + } + } +} diff --git a/src/main/runtime/orchestration/federation-sync.test.ts b/src/main/runtime/orchestration/federation-sync.test.ts index 26223ed1a64..847bc7cdb67 100644 --- a/src/main/runtime/orchestration/federation-sync.test.ts +++ b/src/main/runtime/orchestration/federation-sync.test.ts @@ -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 } | null = null - let blockedPull: { reached: () => void; released: Promise } | 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((resolve) => (noteReached = resolve)) - const released = new Promise((resolve) => (release = resolve)) - blockedAck = { reached: noteReached, released } - return { reached, release } - }, - blockPull: () => { - let noteReached!: () => void - let release!: () => void - const reached = new Promise((resolve) => (noteReached = resolve)) - const released = new Promise((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) - }) }) diff --git a/src/main/runtime/orchestration/federation-sync.ts b/src/main/runtime/orchestration/federation-sync.ts index a2a0d77d6ab..41e275a2ff1 100644 --- a/src/main/runtime/orchestration/federation-sync.ts +++ b/src/main/runtime/orchestration/federation-sync.ts @@ -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 } } diff --git a/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts b/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts index 391e259441a..736792add91 100644 --- a/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts +++ b/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts @@ -51,8 +51,6 @@ export class OrchestrationPeerCapabilityCache { expectedRuntimeEpoch: string | null capability: RuntimeCapability probe: () => Promise - /** Require a status probe even when the expected epoch has a cached answer. */ - forceProbe?: boolean }): Promise { return this.resolveAttempt(args, 1) } @@ -63,16 +61,14 @@ export class OrchestrationPeerCapabilityCache { expectedRuntimeEpoch: string | null capability: RuntimeCapability probe: () => Promise - forceProbe?: boolean }, staleRetriesRemaining: number ): Promise { 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, diff --git a/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.test.ts b/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.test.ts index c228628a4eb..e81c0e5c804 100644 --- a/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.test.ts @@ -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).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).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' ) }) }) diff --git a/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.ts b/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.ts index 75ab763f922..451e5586edb 100644 --- a/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.ts +++ b/src/main/runtime/rpc/methods/orchestration-federated-fleet-snapshot.ts @@ -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 - } - }) - 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() @@ -203,7 +197,13 @@ export function applyFederatedFleetObservations( federated: Awaited>, 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 diff --git a/src/main/runtime/rpc/methods/orchestration-federated-release-safety.test.ts b/src/main/runtime/rpc/methods/orchestration-federated-release-safety.test.ts index 0a5842f8f6e..fe1aa5d69e0 100644 --- a/src/main/runtime/rpc/methods/orchestration-federated-release-safety.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-federated-release-safety.test.ts @@ -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') diff --git a/src/main/runtime/rpc/methods/orchestration-federated-worker-read.ts b/src/main/runtime/rpc/methods/orchestration-federated-worker-read.ts index 7789dbd5f9f..71e063c10e0 100644 --- a/src/main/runtime/rpc/methods/orchestration-federated-worker-read.ts +++ b/src/main/runtime/rpc/methods/orchestration-federated-worker-read.ts @@ -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 - }) - 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 } } diff --git a/src/main/runtime/rpc/methods/orchestration-federated-worker-release-host.ts b/src/main/runtime/rpc/methods/orchestration-federated-worker-release-host.ts index 2fdb7777738..b1b365f012d 100644 --- a/src/main/runtime/rpc/methods/orchestration-federated-worker-release-host.ts +++ b/src/main/runtime/rpc/methods/orchestration-federated-worker-release-host.ts @@ -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({ diff --git a/src/main/runtime/rpc/methods/orchestration-federated-worker-release.ts b/src/main/runtime/rpc/methods/orchestration-federated-worker-release.ts index 55c52cadff8..0398be06093 100644 --- a/src/main/runtime/rpc/methods/orchestration-federated-worker-release.ts +++ b/src/main/runtime/rpc/methods/orchestration-federated-worker-release.ts @@ -46,6 +46,8 @@ export async function releaseFederatedWorker(args: { requestId: string }): Promise { 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, diff --git a/src/main/runtime/rpc/methods/orchestration-federation-liveness-verdict.test.ts b/src/main/runtime/rpc/methods/orchestration-federation-liveness-verdict.test.ts index c41082e38a1..6084613988f 100644 --- a/src/main/runtime/rpc/methods/orchestration-federation-liveness-verdict.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-federation-liveness-verdict.test.ts @@ -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() } diff --git a/src/main/runtime/rpc/methods/orchestration-federation-output.test.ts b/src/main/runtime/rpc/methods/orchestration-federation-output.test.ts index 2e101b6fcb1..75f27f1f52a 100644 --- a/src/main/runtime/rpc/methods/orchestration-federation-output.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-federation-output.test.ts @@ -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') diff --git a/src/main/runtime/rpc/methods/orchestration-federation-relay.ts b/src/main/runtime/rpc/methods/orchestration-federation-relay.ts index ea8b09432e8..235286fcbc8 100644 --- a/src/main/runtime/rpc/methods/orchestration-federation-relay.ts +++ b/src/main/runtime/rpc/methods/orchestration-federation-relay.ts @@ -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 ? {} : { diff --git a/src/main/runtime/rpc/methods/orchestration-send-remote.ts b/src/main/runtime/rpc/methods/orchestration-send-remote.ts index 31396bfc0ee..9243b977d2e 100644 --- a/src/main/runtime/rpc/methods/orchestration-send-remote.ts +++ b/src/main/runtime/rpc/methods/orchestration-send-remote.ts @@ -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', diff --git a/src/main/runtime/rpc/methods/orchestration-worker-list-pagination.test.ts b/src/main/runtime/rpc/methods/orchestration-worker-list-pagination.test.ts index 780a648dd61..3dbe11f30cc 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-list-pagination.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-list-pagination.test.ts @@ -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) => void - const status = new Promise>((resolve) => { - resolveStatus = resolve + let resolveSnapshot!: () => void + const snapshotGate = new Promise((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({ @@ -420,15 +417,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 - } -} diff --git a/src/shared/orchestration-fleet-projection.ts b/src/shared/orchestration-fleet-projection.ts index ad624131f12..c7740c77d72 100644 --- a/src/shared/orchestration-fleet-projection.ts +++ b/src/shared/orchestration-fleet-projection.ts @@ -57,6 +57,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' } diff --git a/src/shared/protocol-version.ts b/src/shared/protocol-version.ts index 16b1a8bea83..95e4bfc37b7 100644 --- a/src/shared/protocol-version.ts +++ b/src/shared/protocol-version.ts @@ -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, From 58adae375e2c3fb48a986794fbf9125d8d76d14b Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Fri, 4 Sep 2026 02:30:46 -0400 Subject: [PATCH 2/2] fix(transcript): wire the attested WSL distro into exact worker session selection A WSL pane's PTY is local (connectionId null), so every WSL hook status was filtered out and worker-read always fell back to screen scraping. Pass the same distro expression the headless terminal state already uses. Also deletes the dead v1 transcript_pin archive path (nothing has written version 1, and its reader read a possibly-remote path from the local filesystem) along with the endOffset thread it was the only caller of, and adds Archived:/Liveness:/Worker: lines so a released archive read no longer prints identically to a live one. --- .../orchestration-module-boundaries.test.ts | 1 + .../orchestration/worker-output.test.ts | 29 ++++++ .../handlers/orchestration/worker-output.ts | 8 ++ ...time-exact-worker-provider-session.test.ts | 57 +++++++++++ ...a-runtime-get-terminal-interactive-wait.ts | 7 ++ .../orchestration/worker-output-archive.ts | 10 -- .../worker-transcript-local-read.ts | 18 +--- .../worker-transcript-read.test.ts | 43 --------- .../orchestration/worker-transcript-read.ts | 4 +- .../orchestration-worker-archive-read.ts | 95 ++----------------- ...chestration-worker-release-archive.test.ts | 77 +-------------- 11 files changed, 116 insertions(+), 233 deletions(-) create mode 100644 src/main/runtime/orca-runtime-exact-worker-provider-session.test.ts diff --git a/src/cli/handlers/orchestration-module-boundaries.test.ts b/src/cli/handlers/orchestration-module-boundaries.test.ts index cea8a934120..b9a58f71ed7 100644 --- a/src/cli/handlers/orchestration-module-boundaries.test.ts +++ b/src/cli/handlers/orchestration-module-boundaries.test.ts @@ -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' ) diff --git a/src/cli/handlers/orchestration/worker-output.test.ts b/src/cli/handlers/orchestration/worker-output.test.ts index cc58cf188a3..bcffb1815ff 100644 --- a/src/cli/handlers/orchestration/worker-output.test.ts +++ b/src/cli/handlers/orchestration/worker-output.test.ts @@ -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( diff --git a/src/cli/handlers/orchestration/worker-output.ts b/src/cli/handlers/orchestration/worker-output.ts index 04e630aae4a..cf2541ccb84 100644 --- a/src/cli/handlers/orchestration/worker-output.ts +++ b/src/cli/handlers/orchestration/worker-output.ts @@ -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}`) } diff --git a/src/main/runtime/orca-runtime-exact-worker-provider-session.test.ts b/src/main/runtime/orca-runtime-exact-worker-provider-session.test.ts new file mode 100644 index 00000000000..e0c599983c0 --- /dev/null +++ b/src/main/runtime/orca-runtime-exact-worker-provider-session.test.ts @@ -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() + }) +}) diff --git a/src/main/runtime/orca-runtime-get-terminal-interactive-wait.ts b/src/main/runtime/orca-runtime-get-terminal-interactive-wait.ts index f428404de90..059f4ffd7b3 100644 --- a/src/main/runtime/orca-runtime-get-terminal-interactive-wait.ts +++ b/src/main/runtime/orca-runtime-get-terminal-interactive-wait.ts @@ -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?.() ?? [] }) diff --git a/src/main/runtime/orchestration/worker-output-archive.ts b/src/main/runtime/orchestration/worker-output-archive.ts index 0f7b9f9146b..20aa3423c6d 100644 --- a/src/main/runtime/orchestration/worker-output-archive.ts +++ b/src/main/runtime/orchestration/worker-output-archive.ts @@ -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 diff --git a/src/main/runtime/orchestration/worker-transcript-local-read.ts b/src/main/runtime/orchestration/worker-transcript-local-read.ts index 75f4f037a4d..50bdc1a924f 100644 --- a/src/main/runtime/orchestration/worker-transcript-local-read.ts +++ b/src/main/runtime/orchestration/worker-transcript-local-read.ts @@ -40,17 +40,13 @@ type LocalTranscriptReadResult = export async function readInitialLocalWorkerTranscriptPage( filePath: string, limit: number, - decode: NativeChatLineDecoder, - endOffset?: number + decode: NativeChatLineDecoder ): Promise { 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 { 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 { diff --git a/src/main/runtime/orchestration/worker-transcript-read.test.ts b/src/main/runtime/orchestration/worker-transcript-read.test.ts index 997f88d874e..06068ab2423 100644 --- a/src/main/runtime/orchestration/worker-transcript-read.test.ts +++ b/src/main/runtime/orchestration/worker-transcript-read.test.ts @@ -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`) diff --git a/src/main/runtime/orchestration/worker-transcript-read.ts b/src/main/runtime/orchestration/worker-transcript-read.ts index bfee9fb24e2..cb6f0467c4a 100644 --- a/src/main/runtime/orchestration/worker-transcript-read.ts +++ b/src/main/runtime/orchestration/worker-transcript-read.ts @@ -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) { diff --git a/src/main/runtime/rpc/methods/orchestration-worker-archive-read.ts b/src/main/runtime/rpc/methods/orchestration-worker-archive-read.ts index 806d6ce9897..8709e675ccc 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-archive-read.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-archive-read.ts @@ -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[0], - pin: WorkerTranscriptPinArchive -): Promise { - 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[0], archive: WorkerTerminalArchiveRow diff --git a/src/main/runtime/rpc/methods/orchestration-worker-release-archive.test.ts b/src/main/runtime/rpc/methods/orchestration-worker-release-archive.test.ts index a92b14d14d7..7eb795beba3 100644 --- a/src/main/runtime/rpc/methods/orchestration-worker-release-archive.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-worker-release-archive.test.ts @@ -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-harness' 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-'))