fix(orchestration): close CI durability regressions

This commit is contained in:
Jinwoo-H
2026-08-31 15:50:58 -04:00
parent f42dcb3094
commit df91d88bd1
9 changed files with 55 additions and 77 deletions
+1 -1
View File
@@ -203,7 +203,7 @@ Terminal rules:
- A text-plus-Enter agent prompt returns a durable request ID and additive stages: `input_accepted`, optionally `queued_pending_turn`, then `submission_observed` and `turn_started` when proven. Raw text-only, bare Enter, interrupt, and terminal query replies keep their existing direct-input behavior.
- `--wait-submit <seconds>` only observes the same accepted prompt. A timeout returns queued/input-accepted truth without resending; after an ambiguous transport failure, repeat the exact command with the reported `--retry-request <id>`.
- An older host reports a legacy `old-host` fallback for an ordinary send and refuses `--wait-submit` or `--retry-request` before input, because it cannot provide durable replay.
- For structured coordination, invoke the `orchestration` skill; it uses `orca orchestration ...` commands for messages, handoffs, task DAGs, dispatches, inbox/reply flows, and coordinator loops. A receiving agent can run `orca orchestration check --unread --inject` to render its unread mail in agent-readable form; this checks the caller's inbox and does not remotely deliver input to another terminal.
- For structured coordination, invoke the `orchestration` skill; it uses `orca orchestration ...` commands for messages, handoffs, task DAGs, dispatches, inbox/reply flows, and coordinator loops. A receiving agent can run `orca orchestration check --peek --format --json` to render its unread mail in agent-readable form; this checks the caller's inbox and does not remotely deliver input to another terminal.
- Use `terminal create --worktree active --command "<agent>"` for a fresh agent in the current worktree. Use `worktree create --agent <agent>` only for a separate checkout (agent in the first terminal — do not also `terminal create` the same agent).
- Use `terminal wait --for tui-idle` for agent CLIs such as Claude Code, Gemini, Codex, OMP, Pi, and Grok; always pass `--timeout-ms`.
- Terminal handles are runtime-scoped. Use `startupTerminal.handle` as the sole agent handle when `worktree create --agent` returns it; if Orca restarts, omits the handle, or returns `terminal_handle_stale`, reacquire with `terminal list` and continue with the replacement only.
File diff suppressed because one or more lines are too long
@@ -122,7 +122,8 @@ describe('Codex WSL scan gate', () => {
await expect(
resolveSessionFilePath('codex', 'session-id', {
transcriptPath: `${WSL_SESSIONS_DIR}\\2026\\rollout-1-session-id.jsonl`,
codexSessionsDirs: [DEBIAN_SESSIONS_DIR]
codexSessionsDirs: [DEBIAN_SESSIONS_DIR],
wslDistro: 'Ubuntu'
})
).rejects.toBe(refusal)
expect(mocks.walk).not.toHaveBeenCalled()
@@ -11,6 +11,11 @@ const DEBIAN_ROLLOUT_UNC = ROLLOUT_UNC.replace('Ubuntu', 'Debian')
vi.mock('../wsl', () => ({
listWslDistrosAsync: vi.fn(async () => ['Ubuntu', 'Debian']),
listRunningWslDistrosAsync: vi.fn(async () => ['Ubuntu', 'Debian']),
listRunningWslHomeDirsAsync: vi.fn(async () => [
UBUNTU_HOME,
UBUNTU_HOME.replace('Ubuntu', 'Debian')
]),
getWslHomeAsync: vi.fn(async (distro: string) => UBUNTU_HOME.replace('Ubuntu', distro))
}))
+6 -1
View File
@@ -6351,7 +6351,12 @@ export class OrcaRuntimeService {
if (current?.db === db) {
return current.promise
}
const sync = syncFederatedDispatch(this, dispatchId)
const generation = this.orchestrationFederationRelayGeneration
const sync = syncFederatedDispatch(
this,
dispatchId,
() => generation === this.orchestrationFederationRelayGeneration
)
.then(() => {
if (this.orchestrationFederationSyncs.get(dispatchId)?.promise === sync) {
this.orchestrationFederationWarnings.delete(dispatchId)
@@ -27,16 +27,26 @@ type PulledRelayItem = {
export async function syncFederatedDispatch(
runtime: OrcaRuntimeService,
dispatchId: string
dispatchId: string,
isCurrent: () => boolean = () => true
): Promise<{ imported: number; acknowledgedThrough: number }> {
return syncFederatedDispatchPages(runtime, dispatchId, MAX_FEDERATION_PULL_PAGES_PER_SYNC)
return syncFederatedDispatchPages(
runtime,
dispatchId,
MAX_FEDERATION_PULL_PAGES_PER_SYNC,
isCurrent
)
}
async function syncFederatedDispatchPages(
runtime: OrcaRuntimeService,
dispatchId: string,
remainingPages: number
remainingPages: number,
isCurrent: () => boolean
): Promise<{ imported: number; acknowledgedThrough: number }> {
if (!isCurrent()) {
return { imported: 0, acknowledgedThrough: 0 }
}
const db = runtime.getOrchestrationDb()
const federated = db.getFederatedDispatch(dispatchId)
const dispatch = db.getDispatchContextById(dispatchId)
@@ -80,6 +90,9 @@ async function syncFederatedDispatchPages(
undefined,
{ expectedEnvironmentPairingRevision: currentServer.pairingRevision }
)) as { runtimeEpoch: string; items: PulledRelayItem[] }
if (!isCurrent()) {
return { imported: 0, acknowledgedThrough: federated.to_home_imported_sequence }
}
let cursor =
shouldReplayUnacknowledged && pulled.items.length > 0
? pulled.items[0].sequence - 1
@@ -157,8 +170,9 @@ async function syncFederatedDispatchPages(
? lifecycleAcknowledgmentBarrier - 1
: cursor
if (
isCurrent() &&
acknowledgmentCursor >
Math.max(getFederationAckedThrough(ackLease, ackIdentity), durableAcknowledgedThrough)
Math.max(getFederationAckedThrough(ackLease, ackIdentity), durableAcknowledgedThrough)
) {
const delivered = (await runtime.callOrchestrationWorkerServer(
federated.environment_id,
@@ -214,11 +228,17 @@ async function syncFederatedDispatchPages(
})
}
if (
isCurrent() &&
pulled.items.length === FEDERATION_PULL_PAGE_SIZE &&
remainingPages > 1 &&
lifecycleAcknowledgmentBarrier === undefined
) {
const next = await syncFederatedDispatchPages(runtime, dispatchId, remainingPages - 1)
const next = await syncFederatedDispatchPages(
runtime,
dispatchId,
remainingPages - 1,
isCurrent
)
return {
imported: imported + next.imported,
acknowledgedThrough: next.acknowledgedThrough
+12 -67
View File
@@ -14,18 +14,15 @@ import { emulatorProbe, emulatorProbeError } from '../../emulator/emulator-probe
import type { OrcaRuntimeService } from '../orca-runtime'
import {
getOrchestrationMutationExecutor,
type OrchestrationMutationExecutor,
type DurableMutationInvocation
type OrchestrationMutationExecutor
} from './orchestration-mutation-executor'
import { orchestrationMigrationFence } from './orchestration-contract-fence'
import { recordRuntimeFeatureInteraction } from './runtime-feature-interaction'
import { OrchestrationLegacyCompatibility } from './orchestration-legacy-compatibility'
import type { RpcDispatchStreamingOptions } from './dispatcher-stream-options'
import { mapDispatcherError } from './dispatcher-error-response'
import { parseRpcRequestParams } from './dispatcher-request-parsing'
import { routeDispatcherClientHostedBrowserRpc } from './dispatcher-client-browser-routing'
import { needsLocalCallerFingerprint } from './dispatcher-caller-fingerprint'
import { RpcStreamingDispatcher } from './rpc-streaming-dispatcher'
import { invokeDispatcherUnaryMethod } from './dispatcher-unary-method-invocation'
export type DispatcherOptions = { runtime: OrcaRuntimeService; methods?: readonly RpcAnyMethod[] }
@@ -88,42 +85,12 @@ export class RpcDispatcher {
emulatorProbe(`rpc ${request.method}`, request.params)
}
try {
const clientHostedBrowser = await routeDispatcherClientHostedBrowserRpc(
this.runtime,
request.method,
parsedParams.value
)
if (clientHostedBrowser.handled) {
recordRuntimeFeatureInteraction(
this.runtime,
request.method,
clientHostedBrowser.result,
undefined,
request.params
)
return successResponse(request.id, meta, clientHostedBrowser.result)
}
const compatibility = await this.legacyOrchestration.tryHandle(
const result = await invokeDispatcherUnaryMethod({
runtime: this.runtime,
request,
parsedParams.value,
options?.signal
)
if (compatibility.handled) {
return successResponse(request.id, meta, compatibility.result)
}
const effectiveParams = compatibility.params ?? parsedParams.value
const legacyCoordinator = this.legacyOrchestration.createCoordinatorInvocation(
request,
compatibility.legacyCoordinatorAuthority
)
const authenticatedCallerFingerprint =
options?.authenticatedCallerFingerprint ??
(needsLocalCallerFingerprint(request, effectiveParams)
? this.orchestrationMutations.getLocalAuthenticatedCallerFingerprint()
: undefined)
const invoke = (mutation?: DurableMutationInvocation) => {
const legacyCoordinatorRunId = legacyCoordinator?.revalidate()
return method.handler(effectiveParams, {
method,
params: parsedParams.value,
context: {
runtime: this.runtime,
signal: options?.signal,
connectionId: options?.connectionId,
@@ -132,33 +99,11 @@ export class RpcDispatcher {
clientKind: options?.clientKind,
clientCapabilities: options?.clientCapabilities,
orchestrationCapability: request.orchestrationCapability,
authenticatedCallerFingerprint:
mutation?.identity.callerFingerprint ??
legacyCoordinator?.mutationCallerFingerprint ??
authenticatedCallerFingerprint,
recordMutationReceipt: mutation?.recordReceipt,
orchestrationMutation: mutation?.identity,
legacyCoordinatorRunId,
legacyCoordinatorAuthority: legacyCoordinator?.authority,
revalidateLegacyCoordinator: legacyCoordinator?.revalidate,
orchestrationCompatibilityCallerAuthority:
compatibility.orchestrationCompatibilityCallerAuthority,
orchestrationCompatibilityEvidence: request.orchestrationCompatibilityEvidence
})
}
const result = await this.orchestrationMutations.run(
request,
effectiveParams,
invoke,
legacyCoordinator?.mutationCallerFingerprint ?? authenticatedCallerFingerprint
)
recordRuntimeFeatureInteraction(
this.runtime,
request.method,
result,
undefined,
request.params
)
authenticatedCallerFingerprint: options?.authenticatedCallerFingerprint
},
orchestrationMutations: this.orchestrationMutations,
legacyOrchestration: this.legacyOrchestration
})
return successResponse(request.id, meta, result)
} catch (error) {
if (request.method.startsWith('emulator.')) {
@@ -26,7 +26,7 @@ describe('orchestration RPC methods', () => {
it('registers all expected methods', () => {
const registry = buildRegistry(ORCHESTRATION_METHODS)
expect(registry.size).toBe(40)
expect(registry.size).toBe(41)
expect(registry.has('orchestration.workerRelease')).toBe(true)
expect(registry.has('orchestration.workerRetain')).toBe(true)
expect(registry.has('orchestration.workerList')).toBe(true)
@@ -178,6 +178,8 @@ export class OrchestrationMutationExecutor {
...identity,
receipt: JSON.stringify(attachMutationReceipt(result, requestId, resumedPendingMutation))
})
// Keep completed receipts when post-commit notification fails; retries replay the durable effect.
effectPossible = true
}
let effectPossible = false
const active = Promise.resolve().then(() =>