From 586e6eaf8fec87e52c1942fd3bbea954e7095e38 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Fri, 4 Sep 2026 19:38:32 -0400 Subject: [PATCH] fix(orchestration): advance capability epochs after peer restart --- ...rchestration-peer-capability-cache.test.ts | 13 ++++ .../orchestration-peer-capability-cache.ts | 12 ++-- .../federation/federated-fleet-host-groups.ts | 47 +++++++++++++++ .../federation/federated-fleet-snapshot.ts | 59 ++----------------- .../federation/federated-worker-read.ts | 6 +- 5 files changed, 77 insertions(+), 60 deletions(-) create mode 100644 src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-host-groups.ts diff --git a/src/main/runtime/orchestration/orchestration-peer-capability-cache.test.ts b/src/main/runtime/orchestration/orchestration-peer-capability-cache.test.ts index 207a3435b9b..77c0df93846 100644 --- a/src/main/runtime/orchestration/orchestration-peer-capability-cache.test.ts +++ b/src/main/runtime/orchestration/orchestration-peer-capability-cache.test.ts @@ -345,6 +345,19 @@ describe('OrchestrationPeerCapabilityCache', () => { cached: true }) }) + + it('accepts a response that advances the epoch it was sent against', () => { + const cache = new OrchestrationPeerCapabilityCache() + cache.remember('peer-a', 'epoch-a', capability, true) + + cache.remember('peer-a', 'epoch-b', capability, false, 'epoch-a') + + expect(cache.knownSupport('peer-a', null, capability)).toEqual({ + runtimeEpoch: 'epoch-b', + supported: false, + cached: true + }) + }) }) function runtimeStatus(runtimeId: string, supported: boolean) { diff --git a/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts b/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts index 0a2e0132f50..f9bca1dca4b 100644 --- a/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts +++ b/src/main/runtime/orchestration/orchestration-peer-capability-cache.ts @@ -135,12 +135,16 @@ export class OrchestrationPeerCapabilityCache { peerFingerprint: string, runtimeEpoch: string, capability: RuntimeCapability, - supported: boolean + supported: boolean, + expectedRuntimeEpoch?: string | null ): void { const latestEpoch = this.latestEpochs.get(peerFingerprint) - // A late answer from a retired epoch used to mint the highest sequence, so it always cleared - // the stale guard and evicted the live epoch's maps on its way back in. - if (latestEpoch !== undefined && latestEpoch !== runtimeEpoch) { + // Advance only from the epoch this call targeted; late answers cannot replace a newer epoch. + if ( + latestEpoch !== undefined && + latestEpoch !== runtimeEpoch && + latestEpoch !== expectedRuntimeEpoch + ) { return } const generation = this.touchPeer(peerFingerprint) diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-host-groups.ts b/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-host-groups.ts new file mode 100644 index 00000000000..3c17a3fdf37 --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-host-groups.ts @@ -0,0 +1,47 @@ +import { ORCHESTRATION_FLEET_PAGE_MAX } from '../../../../../../shared/orchestration-fleet-projection' +import type { OrchestrationDb } from '../../../../orchestration/db' +import type { FederatedDispatchRow } from '../../../../orchestration/types' +import type { OrcaRuntimeService } from '../../../../orca-runtime' + +type HostGroup = { + environmentId: string + name: string + dispatches: FederatedDispatchRow[] +} + +export function groupFederatedDispatches(args: { + runtime: OrcaRuntimeService + db: OrchestrationDb + dispatchIds: readonly string[] +}): HostGroup[] { + const groups = new Map() + const federatedByDispatchId = new Map( + args.db + .listFederatedDispatchesByIds(args.dispatchIds) + .map((dispatch) => [dispatch.dispatch_id, dispatch]) + ) + for (const dispatchId of args.dispatchIds) { + const dispatch = federatedByDispatchId.get(dispatchId) + if (!dispatch) { + continue + } + const groupKey = `${dispatch.environment_id}\u0000${dispatch.peer_fingerprint}` + const group = groups.get(groupKey) ?? { + environmentId: dispatch.environment_id, + name: dispatch.environment_name, + dispatches: [] + } + group.dispatches.push(dispatch) + groups.set(groupKey, group) + } + return [...groups.values()].flatMap((group) => { + const batches: HostGroup[] = [] + for (let offset = 0; offset < group.dispatches.length; offset += ORCHESTRATION_FLEET_PAGE_MAX) { + batches.push({ + ...group, + dispatches: group.dispatches.slice(offset, offset + ORCHESTRATION_FLEET_PAGE_MAX) + }) + } + return batches + }) +} diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-snapshot.ts b/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-snapshot.ts index df54243ec96..be91c6b968e 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-snapshot.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federated-fleet-snapshot.ts @@ -1,7 +1,7 @@ +import { groupFederatedDispatches } from './federated-fleet-host-groups' import { mapWithConcurrency } from '../../../../../../shared/map-with-concurrency' import { ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY } from '../../../../../../shared/protocol-version' import { - ORCHESTRATION_FLEET_PAGE_MAX, refreshOrchestrationFleetLivenessAttention, type FleetDurableWorker, type OrchestrationFleetPage @@ -10,7 +10,6 @@ import { projectFleetNextAction } from '../../../../../../shared/orchestration-f import { getOrchestrationPeerCapabilityCache } from '../../../../orchestration/orchestration-peer-capability-cache' import type { OrchestrationDb } from '../../../../orchestration/db' import { OrchestrationError } from '../../../../orchestration/orchestration-error' -import type { FederatedDispatchRow } from '../../../../orchestration/types' import type { OrcaRuntimeService } from '../../../../orca-runtime' import { resolvePinnedFederatedServer } from '../worker/worker-observation' @@ -31,12 +30,6 @@ export type FederatedFleetHostError = { dispatchIds: string[] } -type HostGroup = { - environmentId: string - name: string - dispatches: FederatedDispatchRow[] -} - export async function readFederatedFleetSnapshots(args: { runtime: OrcaRuntimeService db: OrchestrationDb @@ -59,7 +52,6 @@ export async function readFederatedFleetSnapshots(args: { }) const remaining = deadline - Date.now() if (remaining <= 0) { - // 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) @@ -68,8 +60,7 @@ export async function readFederatedFleetSnapshots(args: { let observedCapabilityEpoch: string | null = null try { const server = resolvePinnedFederatedServer(args.runtime, first) - // `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. + // Shipped hosts serve this method without advertising it. const known = cache.knownSupport( first.peer_fingerprint, first.remote_runtime_epoch, @@ -101,7 +92,8 @@ export async function readFederatedFleetSnapshots(args: { first.peer_fingerprint, snapshot.runtimeEpoch, ORCHESTRATION_FEDERATION_FLEET_SNAPSHOT_RUNTIME_CAPABILITY, - true + true, + observedCapabilityEpoch ) const projectedDispatches = projectFleetRuntimeEpochs( args.db, @@ -119,7 +111,6 @@ export async function readFederatedFleetSnapshots(args: { ? item : { ...item, - // A non-exact identity can never prove either liveness or exit. observation: { ...item.observation, status: 'unverifiable' as const } } ), @@ -138,7 +129,6 @@ export async function readFederatedFleetSnapshots(args: { } return { observations: [], error: error('capability_unsupported') } } - // The environment was repointed at another Orca server; that is an identity fact, not silence. return { observations: [], error: error( @@ -231,14 +221,12 @@ export function applyFederatedFleetObservations( ? { verdict: 'exited', source: 'execution_host' } : { verdict: 'unverifiable', reason: hostReportedReason(observation.reason) } worker.evidence.liveStatus = observation.status === 'live' ? 'fresh' : 'unavailable' - // The host answered for both `live` and `exited`, so both carry a real observation time. worker.evidence.lastObservedAt = observation.status === 'unverifiable' ? null : observedAt refreshFleetWorkerVerdict(worker, durable) } } -/** The host verdict replaces the local one, so everything derived from liveness has to - * follow it: a stale `inspect` outranked the `recover` a proven remote exit owes. */ +// Recompute every projection derived from the host's verdict. function refreshFleetWorkerVerdict( worker: OrchestrationFleetPage['workers'][number], durable: ReadonlyMap @@ -276,40 +264,3 @@ function hostReportedReason( ? (reason as 'missing_status' | 'stale_status' | 'future_status' | 'restored_unconfirmed') : 'host_indeterminate' } - -function groupFederatedDispatches(args: { - runtime: OrcaRuntimeService - db: OrchestrationDb - dispatchIds: readonly string[] -}): HostGroup[] { - const groups = new Map() - const federatedByDispatchId = new Map( - args.db - .listFederatedDispatchesByIds(args.dispatchIds) - .map((dispatch) => [dispatch.dispatch_id, dispatch]) - ) - for (const dispatchId of args.dispatchIds) { - const dispatch = federatedByDispatchId.get(dispatchId) - if (!dispatch) { - continue - } - const groupKey = `${dispatch.environment_id}\u0000${dispatch.peer_fingerprint}` - const group = groups.get(groupKey) ?? { - environmentId: dispatch.environment_id, - name: dispatch.environment_name, - dispatches: [] - } - group.dispatches.push(dispatch) - groups.set(groupKey, group) - } - return [...groups.values()].flatMap((group) => { - const batches: HostGroup[] = [] - for (let offset = 0; offset < group.dispatches.length; offset += ORCHESTRATION_FLEET_PAGE_MAX) { - batches.push({ - ...group, - dispatches: group.dispatches.slice(offset, offset + ORCHESTRATION_FLEET_PAGE_MAX) - }) - } - return batches - }) -} diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-read.ts b/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-read.ts index d16250f4407..563f7e81bb8 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-read.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federated-worker-read.ts @@ -63,7 +63,8 @@ export async function readFederatedWorkerOutput(args: { args.federated.peer_fingerprint, remote.runtimeEpoch, ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY, - true + true, + expectedRuntimeEpoch ) projectRemoteRuntimeEpoch(args.db, observationFence, remote.runtimeEpoch) return { @@ -80,7 +81,8 @@ export async function readFederatedWorkerOutput(args: { args.federated.peer_fingerprint, legacy.remoteRuntimeEpoch, ORCHESTRATION_FEDERATION_STRUCTURED_READ_RUNTIME_CAPABILITY, - false + false, + expectedRuntimeEpoch ) projectRemoteRuntimeEpoch(args.db, observationFence, legacy.remoteRuntimeEpoch) return legacy