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

This commit is contained in:
Jinwoo-H
2026-09-04 02:23:26 -04:00
27 changed files with 778 additions and 276 deletions
+19
View File
@@ -4,6 +4,7 @@ import { describeQuoteStrippedJsonFlag } from './quote-stripped-json-flag'
export function getRequiredStringFlag(flags: Map<string, string | boolean>, name: string): string {
const value = flags.get(name)
rejectValuelessFlag(value, name)
if (typeof value === 'string' && value.length > 0) {
return value
}
@@ -15,6 +16,7 @@ export function getRequiredStringFlagAllowingEmpty(
name: string
): string {
const value = flags.get(name)
rejectValuelessFlag(value, name)
if (typeof value === 'string') {
return value
}
@@ -26,9 +28,24 @@ export function getOptionalStringFlag(
name: string
): string | undefined {
const value = flags.get(name)
rejectValuelessFlag(value, name)
return typeof value === 'string' && value.length > 0 ? value : undefined
}
/**
* A valued flag whose value the shell (or a missing variable) ate parses as `true`. Dropping it
* silently mints a fresh mutation identity and can deliver a prompt twice (#15180), so every
* valued-flag accessor refuses the damaged shape by name.
*/
export function rejectValuelessFlag(value: string | boolean | undefined, name: string): void {
if (value === true) {
throw new RuntimeClientError(
'invalid_argument',
`--${name} requires a value; it was passed with none.`
)
}
}
/**
* A JSON-valued flag, rejected up front when a native argv boundary stripped its quotes so the
* error names the shell instead of the user's value (#16706). The value itself is still parsed
@@ -64,6 +81,7 @@ export function getOptionalNumberFlag(
name: string
): number | undefined {
const value = flags.get(name)
rejectValuelessFlag(value, name)
if (typeof value !== 'string' || value.length === 0) {
return undefined
}
@@ -131,6 +149,7 @@ export function getOptionalNullableNumberFlag(
name: string
): number | null | undefined {
const value = flags.get(name)
rejectValuelessFlag(value, name)
if (value === 'null') {
return null
}
@@ -75,7 +75,7 @@ describe('orchestration worker-start CLI contract', () => {
['timeout-ms', '90000'],
['run', 'run_1'],
['from', 'term_coord'],
['retry-request', 'request_1']
['retry-request', '44444444-4444-4444-8444-444444444444']
])
)
@@ -99,7 +99,7 @@ describe('orchestration worker-start CLI contract', () => {
from: 'term_coord',
devMode: false
},
{ orchestrationRequestId: 'request_1' }
{ orchestrationRequestId: '44444444-4444-4444-8444-444444444444' }
)
expect(process.exitCode).toBeUndefined()
})
+5 -2
View File
@@ -122,14 +122,17 @@ describe('orchestration send structured payload flags', () => {
['type', 'heartbeat'],
['dispatch-id', 'ctx_1'],
['dispatch-capability', 'dcap_secret'],
['retry-request', 'mutation_1']
['retry-request', '33333333-3333-4333-8333-333333333333']
])
)
expect(callMock).toHaveBeenCalledWith(
'orchestration.send',
expect.not.objectContaining({ dispatchCapability: expect.anything() }),
{ orchestrationCapability: 'dcap_secret', orchestrationRequestId: 'mutation_1' }
{
orchestrationCapability: 'dcap_secret',
orchestrationRequestId: '33333333-3333-4333-8333-333333333333'
}
)
})
@@ -1,5 +1,5 @@
import type { RuntimeClient } from '../../runtime-client'
import { getOptionalStringFlag } from '../../flags'
import { readRetryRequestFlag } from '../../retry-request-flag'
import { orchestrationMutationRecoveryError } from '../../orchestration-mutation-recovery'
export function callOrchestrationMutation<TResult>(
@@ -9,7 +9,7 @@ export function callOrchestrationMutation<TResult>(
params: unknown,
options?: { timeoutMs?: number; orchestrationCapability?: string }
) {
const requestId = getOptionalStringFlag(flags, 'retry-request')
const requestId = readRetryRequestFlag(flags)
const result = requestId
? client.call<TResult>(method, params, { ...options, orchestrationRequestId: requestId })
: options
+2 -1
View File
@@ -3,6 +3,7 @@ import { TERMINAL_PROMPT_DELIVERY_RUNTIME_CAPABILITY } from '../../shared/protoc
import type { CommandHandler } from '../dispatch'
import { formatTerminalSend, printResult } from '../format'
import { getOptionalPositiveIntegerFlag, getOptionalStringFlag } from '../flags'
import { readRetryRequestFlag } from '../retry-request-flag'
import { RuntimeClientError } from '../runtime-client'
import { attachUnverifiedTerminalPromptRecovery } from '../runtime/terminal-prompt-mutation-recovery'
import { getTerminalHandle } from '../selectors'
@@ -12,7 +13,7 @@ export const terminalSendHandler: CommandHandler = async ({ flags, client, cwd,
const enter = flags.get('enter') === true
const interrupt = flags.get('interrupt') === true
const promptCandidate = !!text && enter && !interrupt
const retryRequest = getOptionalStringFlag(flags, 'retry-request')
const retryRequest = readRetryRequestFlag(flags)
const waitSubmitSeconds = getOptionalPositiveIntegerFlag(flags, 'wait-submit')
if ((retryRequest || waitSubmitSeconds) && !promptCandidate) {
throw new RuntimeClientError(
+11 -11
View File
@@ -278,7 +278,7 @@ describe('terminal send CLI', () => {
accepted: true,
bytesWritten: 7,
prompt: {
requestId: 'prompt-1',
requestId: '11111111-1111-4111-8111-111111111111',
stages: ['input_accepted'],
provider: 'codex',
observation: 'supported',
@@ -405,7 +405,7 @@ describe('terminal send CLI', () => {
accepted: true,
bytesWritten: 8,
prompt: {
requestId: 'prompt-1',
requestId: '11111111-1111-4111-8111-111111111111',
stages: ['input_accepted'],
provider: 'codex',
observation: 'supported',
@@ -423,7 +423,7 @@ describe('terminal send CLI', () => {
['terminal', 'term-1'],
['text', 'continue'],
['enter', true],
['retry-request', 'prompt-1'],
['retry-request', '11111111-1111-4111-8111-111111111111'],
['wait-submit', '3']
]),
client: promptClient(call, true),
@@ -436,7 +436,7 @@ describe('terminal send CLI', () => {
expect.objectContaining({ agentPrompt: true, waitSubmitMs: 3_000 }),
{
terminalPromptPreflight: { runtimeId: 'runtime-current' },
orchestrationRequestId: 'prompt-1',
orchestrationRequestId: '11111111-1111-4111-8111-111111111111',
timeoutMs: 13_000
}
)
@@ -467,7 +467,7 @@ describe('terminal send CLI', () => {
['terminal', 'term-1'],
['text', 'continue'],
['enter', true],
['retry-request', 'prompt-1'],
['retry-request', '11111111-1111-4111-8111-111111111111'],
['wait-submit', '3']
]),
client,
@@ -482,7 +482,7 @@ describe('terminal send CLI', () => {
expect.objectContaining({ agentPrompt: true, waitSubmitMs: 3_000 }),
{
terminalPromptPreflight: { runtimeId: 'new-runtime-before-restart' },
orchestrationRequestId: 'prompt-1',
orchestrationRequestId: '11111111-1111-4111-8111-111111111111',
timeoutMs: 13_000
}
)
@@ -572,7 +572,7 @@ describe('terminal send CLI', () => {
['terminal', 'term-1'],
['text', 'review'],
['enter', true],
['retry-request', 'prompt-1']
['retry-request', '11111111-1111-4111-8111-111111111111']
]),
client: {
call,
@@ -605,7 +605,7 @@ describe('terminal send CLI', () => {
accepted: true,
bytesWritten: 13,
prompt: {
requestId: 'prompt-retry',
requestId: '22222222-2222-4222-8222-222222222222',
stages: ['input_accepted'],
provider: 'codex',
observation: 'supported',
@@ -621,7 +621,7 @@ describe('terminal send CLI', () => {
['terminal', 'term-1'],
['text', 'retry safely'],
['enter', true],
['retry-request', 'prompt-retry']
['retry-request', '22222222-2222-4222-8222-222222222222']
])
vi.spyOn(console, 'log').mockImplementation(() => {})
@@ -639,11 +639,11 @@ describe('terminal send CLI', () => {
expect(call.mock.calls.map((args) => args[2])).toEqual([
{
terminalPromptPreflight: { runtimeId: 'runtime-current' },
orchestrationRequestId: 'prompt-retry'
orchestrationRequestId: '22222222-2222-4222-8222-222222222222'
},
{
terminalPromptPreflight: { runtimeId: 'runtime-current' },
orchestrationRequestId: 'prompt-retry'
orchestrationRequestId: '22222222-2222-4222-8222-222222222222'
}
])
})
+148
View File
@@ -0,0 +1,148 @@
import { describe, expect, it, vi } from 'vitest'
import { parseArgs } from './args'
import { COMMAND_SPECS } from './specs'
import { TERMINAL_PROMPT_DELIVERY_RUNTIME_CAPABILITY } from '../shared/protocol-version'
import type { RuntimeClient } from './runtime-client'
import { TERMINAL_HANDLERS } from './handlers/terminal'
import { ORCHESTRATION_HANDLERS } from './handlers/orchestration'
import { readRetryRequestFlag } from './retry-request-flag'
const PATHS = COMMAND_SPECS.map((spec) => spec.path)
function promptClient() {
const call = vi.fn().mockResolvedValue({
result: {
send: {
handle: 'term-1',
accepted: true,
bytesWritten: 2,
prompt: {
requestId: '11111111-1111-4111-8111-111111111111',
stages: ['input_accepted', 'turn_started'],
provider: 'claude',
observation: 'supported',
processIncarnation: 'inc-1',
generation: 1,
baselineWorkingSequence: 3
}
}
},
_meta: { runtimeId: 'runtime-1' }
})
const client = {
call,
getCliStatus: vi.fn().mockResolvedValue({
result: {
runtime: {
reachable: true,
runtimeId: 'runtime-1',
capabilities: [TERMINAL_PROMPT_DELIVERY_RUNTIME_CAPABILITY]
}
}
})
} as unknown as RuntimeClient
return { client, call }
}
async function sendWith(
argv: string[]
): Promise<{ error: unknown; call: ReturnType<typeof vi.fn> }> {
const { client, call } = promptClient()
vi.spyOn(console, 'log').mockImplementation(() => {})
const error = await TERMINAL_HANDLERS['terminal send']({
flags: parseArgs(argv, PATHS).flags,
client,
cwd: '/tmp/worktree',
json: true
})
.then(() => undefined)
.catch((caught: unknown) => caught)
return { error, call }
}
describe('--retry-request and --wait-submit value damage', () => {
it('parses a value-less flag as boolean true', () => {
const parsed = parseArgs(
['terminal', 'send', '--terminal', 'term-1', '--text', 'hi', '--enter', '--retry-request'],
PATHS
)
expect(parsed.flags.get('retry-request')).toBe(true)
})
it('rejects a value-less --retry-request instead of minting a fresh identity', async () => {
const { error, call } = await sendWith([
'terminal',
'send',
'--terminal',
'term-1',
'--text',
'hi',
'--enter',
'--retry-request'
])
expect(error).toMatchObject({ code: 'invalid_argument' })
expect((error as Error).message).toContain('--retry-request requires a value')
expect(call).not.toHaveBeenCalled()
})
it('rejects an empty --retry-request= value', async () => {
const { error, call } = await sendWith([
'terminal',
'send',
'--terminal',
'term-1',
'--text',
'hi',
'--enter',
'--retry-request='
])
expect(error).toMatchObject({ code: 'invalid_argument' })
expect((error as Error).message).toContain('--retry-request must be the UUID')
expect(call).not.toHaveBeenCalled()
})
it('rejects a non-UUID --retry-request value', () => {
expect(() => readRetryRequestFlag(new Map([['retry-request', 'prompt-1']]))).toThrow(
'--retry-request must be the UUID'
)
expect(
readRetryRequestFlag(new Map([['retry-request', '11111111-1111-4111-8111-111111111111']]))
).toBe('11111111-1111-4111-8111-111111111111')
})
it('rejects a value-less --wait-submit instead of silently not waiting', async () => {
const { error, call } = await sendWith([
'terminal',
'send',
'--terminal',
'term-1',
'--text',
'hi',
'--enter',
'--wait-submit'
])
expect(error).toMatchObject({ code: 'invalid_argument' })
expect((error as Error).message).toContain('--wait-submit requires a value')
expect(call).not.toHaveBeenCalled()
})
it('rejects a damaged --retry-request on an orchestration verb', async () => {
const call = vi.fn()
const client = { call } as unknown as RuntimeClient
for (const value of [true as const, 'worker-stop-1']) {
const error = await ORCHESTRATION_HANDLERS['orchestration worker-stop']({
flags: new Map<string, string | boolean>([
['dispatch', 'ctx_1'],
['retry-request', value]
]),
client,
cwd: '/tmp/worktree',
json: true
})
.then(() => undefined)
.catch((caught: unknown) => caught)
expect(error).toMatchObject({ code: 'invalid_argument' })
}
expect(call).not.toHaveBeenCalled()
})
})
+24
View File
@@ -0,0 +1,24 @@
import { rejectValuelessFlag } from './flags'
import { RuntimeClientError } from './runtime/types'
const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i
/**
* `--retry-request` carries the mutation identity that makes a replay idempotent. A damaged value
* must never fall through to `undefined`, because the client would then mint a fresh identity and
* re-apply a mutation that may already have taken effect (#15180).
*/
export function readRetryRequestFlag(flags: Map<string, string | boolean>): string | undefined {
const value = flags.get('retry-request')
rejectValuelessFlag(value, 'retry-request')
if (value === undefined) {
return undefined
}
if (typeof value !== 'string' || !UUID_PATTERN.test(value)) {
throw new RuntimeClientError(
'invalid_argument',
'--retry-request must be the UUID Orca reported for the original request; pass it exactly as printed, or omit the flag to start a new request.'
)
}
return value
}
+45
View File
@@ -199,6 +199,50 @@ describe('RuntimeClient orchestration recovery identity', () => {
expect((error as Error).message).toContain('--retry-request prompt-current')
})
it('keeps the prompt retry ID when the attested runtime times out in transport', async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-rt-timeout-'))
const endpoint = join(userDataPath, 'runtime.sock')
let receivedRequest: Record<string, unknown> | undefined
const server = createServer((socket) => {
socket.once('data', (data) => {
receivedRequest = JSON.parse(String(data).trim()) as Record<string, unknown>
})
})
servers.add(server)
await new Promise<void>((resolve) => server.listen(endpoint, resolve))
writeRuntimeConnection(userDataPath, endpoint, 'runtime-current')
const client = new RuntimeClient(userDataPath, 200, null, null, 'orca')
const error = await client
.call(
'terminal.send',
{
terminal: 'term-current',
text: 'review',
enter: true,
interrupt: false,
agentPrompt: true,
client: { id: 'orca-cli', type: 'desktop' }
},
{
terminalPromptPreflight: { runtimeId: 'runtime-current' },
orchestrationRequestId: 'prompt-transport-timeout'
}
)
.then(() => undefined)
.catch((caught: unknown) => caught)
expect(receivedRequest?.orchestrationRequestId).toBe('prompt-transport-timeout')
expect(error).toBeInstanceOf(RuntimeClientError)
expect(error).not.toBeInstanceOf(RuntimeRpcFailureError)
expect((error as RuntimeClientError).code).toBe('runtime_timeout')
expect(error).toMatchObject({ data: { orchestrationRequestId: 'prompt-transport-timeout' } })
expect((error as Error).message).toContain(
'--retry-request prompt-transport-timeout --wait-submit <seconds>'
)
expect((error as RuntimeClientError).data).not.toHaveProperty('retrySafe')
})
it('blocks retry when a downgraded runtime rejects after capability preflight', async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-runtime-downgraded-prompt-'))
const endpoint = join(userDataPath, 'runtime.sock')
@@ -285,6 +329,7 @@ describe('RuntimeClient orchestration recovery identity', () => {
expect(error).toBeInstanceOf(RuntimeClientError)
expect(error).not.toBeInstanceOf(RuntimeRpcFailureError)
expectPromptRetryBlockedJson(error, 'prompt-downgraded-lost-reply')
expect(JSON.stringify((error as RuntimeClientError).data)).not.toContain('Update Orca')
})
it('reports an unknown legacy prompt outcome without advertising an unsafe retry', async () => {
+11 -7
View File
@@ -18,7 +18,7 @@ import {
attachDurableMutationRecovery,
attachLegacyTerminalPromptRecovery,
attachUnverifiedTerminalPromptRecovery,
isFailureFromPreflightRuntime
didAnotherRuntimeHandleTerminalPrompt
} from './terminal-prompt-mutation-recovery'
import { markEnvironmentUsed } from './environments'
import { resolveRemotePairing } from './runtime-remote-pairing'
@@ -104,14 +104,18 @@ export class RuntimeClient {
const originalCommand = durableMutation
? buildOrchestrationRecoveryCommand(method, params, this.cliExecutable, this.originalArgs)
: undefined
const recover = (error: unknown) => {
const recover = (error: unknown, targetRuntimeId: string | null) => {
if (legacyTerminalPrompt) {
return attachLegacyTerminalPromptRecovery(error)
}
if (
terminalPromptMutation &&
options?.terminalPromptPreflight &&
!isFailureFromPreflightRuntime(error, options.terminalPromptPreflight.runtimeId)
didAnotherRuntimeHandleTerminalPrompt(
error,
options.terminalPromptPreflight.runtimeId,
targetRuntimeId
)
) {
return attachUnverifiedTerminalPromptRecovery(error)
}
@@ -145,10 +149,10 @@ export class RuntimeClient {
envelope
})
} catch (error) {
throw recover(error)
throw recover(error, null)
}
if (response.ok === false) {
throw recover(new RuntimeRpcFailureError(response))
throw recover(new RuntimeRpcFailureError(response), null)
}
if (this.environmentSelector) {
markEnvironmentUsed(this.userDataPath, this.environmentSelector, {
@@ -162,10 +166,10 @@ export class RuntimeClient {
try {
response = await sendRequest<TResult>(metadata, method, params, effectiveTimeoutMs, envelope)
} catch (error) {
throw recover(error)
throw recover(error, metadata.runtimeId ?? null)
}
if (response.ok === false) {
throw recover(new RuntimeRpcFailureError(response))
throw recover(new RuntimeRpcFailureError(response), metadata.runtimeId ?? null)
}
return response
}
@@ -1,6 +1,8 @@
import { attachMutationRecovery } from './client-error-recovery'
import { RuntimeClientError, RuntimeRpcFailureError } from './types'
const INSPECT_STEP = 'Inspect the terminal output and agent state without sending input.'
export function attachDurableMutationRecovery(
error: unknown,
requestId: string | undefined,
@@ -10,7 +12,7 @@ export function attachDurableMutationRecovery(
if (method !== 'terminal.send' || !requestId || !(error instanceof RuntimeClientError)) {
return attachMutationRecovery(error, requestId, originalCommand)
}
const message = `${error.message} Terminal prompt request ID: ${requestId}. Re-issue the exact command with --retry-request ${requestId}; do not retry it without that ID.`
const message = `${error.message} Terminal prompt request ID: ${requestId}. Re-issue the exact command with --retry-request ${requestId} --wait-submit <seconds>; do not retry it without that ID.`
const data = {
...(error.data && typeof error.data === 'object' ? error.data : {}),
orchestrationRequestId: requestId,
@@ -31,7 +33,11 @@ export function attachLegacyTerminalPromptRecovery(error: unknown): unknown {
}
return attachUnknownTerminalPromptRecovery(
error,
'The legacy host cannot prove whether the prompt was delivered'
'The legacy host cannot prove whether the prompt was delivered',
[
INSPECT_STEP,
'Update Orca on the execution host before future prompt sends that need durable retry.'
]
)
}
@@ -45,35 +51,49 @@ export function attachUnverifiedTerminalPromptRecovery(error: unknown): RuntimeC
)
return attachUnknownTerminalPromptRecovery(
normalized,
'Orca cannot prove whether the prompt was delivered by the prompt-delivery-capable runtime from the preflight'
'Orca cannot prove whether the prompt was delivered by the prompt-delivery-capable runtime from the preflight',
[
INSPECT_STEP,
'A different Orca runtime answered than the one whose prompt-delivery support was verified; confirm which runtime serves this host before sending again.'
]
)
}
export function isFailureFromPreflightRuntime(
/**
* The prompt request ID survives a failed send unless a runtime other than the preflight's
* prompt-delivery host handled it; a transport failure alone means nobody else answered, and the
* attested host still holds the durable pending receipt that makes `--retry-request` idempotent.
*/
export function didAnotherRuntimeHandleTerminalPrompt(
error: unknown,
preflightRuntimeId: string | null
preflightRuntimeId: string | null,
targetRuntimeId: string | null
): boolean {
const handledBy =
error instanceof RuntimeRpcFailureError
? (error.response._meta?.runtimeId ?? null)
: targetRuntimeId
if (handledBy === null) {
return false
}
return (
typeof preflightRuntimeId === 'string' &&
preflightRuntimeId.length > 0 &&
error instanceof RuntimeRpcFailureError &&
error.response._meta?.runtimeId === preflightRuntimeId
typeof preflightRuntimeId !== 'string' ||
preflightRuntimeId.length === 0 ||
handledBy !== preflightRuntimeId
)
}
function attachUnknownTerminalPromptRecovery(
error: RuntimeClientError,
reason: string
reason: string,
nextSteps: string[]
): RuntimeClientError {
const message = `${error.message} ${reason}; inspect the terminal before deciding what to do, and do not resend automatically.`
const data: Record<string, unknown> = {
...(error.data && typeof error.data === 'object' ? error.data : {}),
deliveryOutcome: 'unknown',
retrySafe: false,
nextSteps: [
'Inspect the terminal output and agent state without sending input.',
'Update Orca on the execution host before future prompt sends that need durable retry.'
]
nextSteps
}
delete data.orchestrationRequestId
delete data.originalCommand
+1 -1
View File
@@ -303,7 +303,7 @@ describe('orca skills CLI', () => {
await main(['skills', 'install', '--skill'], '/tmp/repo')
expect(process.exitCode).toBe(1)
expect(errorSpy).toHaveBeenCalledWith('Missing required --skill')
expect(errorSpy).toHaveBeenCalledWith('--skill requires a value; it was passed with none.')
expect(spawnMock).not.toHaveBeenCalled()
})
+24 -2
View File
@@ -70,7 +70,7 @@ describe('formatTerminalSend', () => {
bytesWritten: 8,
prompt: {
requestId: 'prompt-healthy',
stages: ['input_accepted', 'submission_observed', 'turn_started'],
stages: ['input_accepted', 'turn_started'],
provider: 'codex',
observation: 'supported',
processIncarnation: 'inc-1',
@@ -81,7 +81,7 @@ describe('formatTerminalSend', () => {
})
).toBe(
[
'Prompt prompt-healthy on term_worker: input_accepted -> submission_observed -> turn_started.',
'Prompt prompt-healthy on term_worker: input_accepted -> turn_started.',
'provider: codex',
'delivery observation: supported'
].join('\n')
@@ -123,4 +123,26 @@ describe('formatTerminalSend', () => {
expect(output).toContain(warning)
expect(output).toContain(nextStep)
})
it('names the next command when a supported send never reached turn_started', () => {
const output = formatTerminalSend({
send: {
handle: 'term_worker',
accepted: true,
bytesWritten: 8,
prompt: {
requestId: 'prompt-swallowed',
stages: ['input_accepted'],
provider: 'claude',
observation: 'supported',
processIncarnation: 'inc-1',
generation: 1,
baselineWorkingSequence: 0
}
}
})
expect(output).toContain('no turn start was observed')
expect(output).toContain('--retry-request prompt-swallowed --wait-submit <seconds>')
})
})
+3
View File
@@ -209,6 +209,9 @@ function formatTerminalPromptObservationWarning(
? 'this host predates durable prompt receipts. Update Orca on the execution host, and inspect the terminal before retrying an ambiguous send.'
: 'input was accepted, but this provider cannot report delivery. Inspect the terminal before retrying.'
}
if (!prompt.stages.includes('turn_started')) {
return `input was accepted but no turn start was observed, so the Enter may have been swallowed. Confirm delivery by reissuing the exact command with --retry-request ${prompt.requestId} --wait-submit <seconds>; the same request ID replays the receipt instead of sending the prompt again.`
}
return null
}
@@ -55,17 +55,14 @@ describe('agent prompt receipt correlation', () => {
Date.now()
)
const writesAfterSubmission = writes.length
await expect(
runtime.observeTerminalAgentPrompt(handle, second.prompt!, 0)
).resolves.toMatchObject({
stages: ['input_accepted', 'submission_observed', 'turn_started']
})
).resolves.toMatchObject({ stages: ['input_accepted', 'turn_started'] })
await expect(
runtime.observeTerminalAgentPrompt(handle, first.prompt!, 0)
).resolves.toMatchObject({
stages: ['input_accepted', 'submission_observed', 'turn_started']
})
const writesAfterSubmission = writes.length
).resolves.toMatchObject({ stages: ['input_accepted', 'turn_started'] })
// Observing a queued receipt must never write to the PTY again.
expect(writes).toHaveLength(writesAfterSubmission)
})
})
@@ -0,0 +1,85 @@
import { describe, expect, it } from 'vitest'
import { AgentPromptRequestCorrelation } from './agent-prompt-request-correlation'
const PTY = 'pty-1'
const GENERATION = 1
function lifecycle(workingSequence: number) {
return { kind: 'lifecycle' as const, workingSequence }
}
function register(
correlation: AgentPromptRequestCorrelation,
requestId: string,
baselineWorkingSequence: number
): void {
correlation.register(PTY, {
generation: GENERATION,
requestId,
baselineWorkingSequence,
baselineExplicitWorkingStartedAt: null
})
}
describe('agent prompt request correlation', () => {
it('gives one lifecycle transition to exactly one queued request', () => {
const correlation = new AgentPromptRequestCorrelation()
register(correlation, 'first', 0)
register(correlation, 'second', 0)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'first', 0, null, lifecycle(1))).toBe(true)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'second', 0, null, lifecycle(1))).toBe(
false
)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'second', 0, null, lifecycle(2))).toBe(true)
})
it('still allocates a later request when an earlier one has no free sequence', () => {
const correlation = new AgentPromptRequestCorrelation()
register(correlation, 'owner-of-6', 5)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'owner-of-6', 5, null, lifecycle(6))).toBe(
true
)
// `late` can only take sequence 6, which is taken; `early` can still take 3.
register(correlation, 'late', 5)
register(correlation, 'early', 2)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'early', 2, null, lifecycle(6))).toBe(true)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'late', 5, null, lifecycle(6))).toBe(false)
})
it('reserves a hook turn start for the oldest eligible request', () => {
const correlation = new AgentPromptRequestCorrelation()
register(correlation, 'oldest', 0)
register(correlation, 'newest', 0)
const hook = { kind: 'hook' as const, workingStartedAt: 500 }
expect(correlation.acceptTurnStart(PTY, GENERATION, 'newest', 0, null, hook)).toBe(false)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'oldest', 0, null, hook)).toBe(true)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'newest', 0, null, hook)).toBe(false)
})
it('refuses a request the PTY no longer holds', () => {
const correlation = new AgentPromptRequestCorrelation()
register(correlation, 'cleared', 0)
correlation.clearForPty(PTY)
expect(correlation.acceptTurnStart(PTY, GENERATION, 'cleared', 0, null, lifecycle(1))).toBe(
false
)
})
it('scopes claims to the generation that recorded them', () => {
const correlation = new AgentPromptRequestCorrelation()
register(correlation, 'gen-1', 0)
correlation.register(PTY, {
generation: 2,
requestId: 'gen-2',
baselineWorkingSequence: 0,
baselineExplicitWorkingStartedAt: null
})
expect(correlation.acceptTurnStart(PTY, GENERATION, 'gen-1', 0, null, lifecycle(1))).toBe(true)
expect(correlation.acceptTurnStart(PTY, 2, 'gen-2', 0, null, lifecycle(1))).toBe(true)
})
})
@@ -0,0 +1,219 @@
import type { AgentPromptTurnStartEvidence } from './agent-prompt-submission-verification'
/**
* Per-PTY ledger that decides which queued prompt owns an observed turn start.
*
* Turn evidence is PTY-wide, so without an owner one observed turn would settle every queued
* prompt that shares its baseline. Registrations stay in arrival order per PTY: the oldest
* eligible request claims the next turn, and a claimed turn can never change hands.
*/
// A stalled request is only dropped when its PTY or generation goes away, so cap the backlog.
const REQUESTS_PER_PTY_LIMIT = 1_024
export type AgentPromptRequestBaseline = {
generation: number
requestId: string
baselineWorkingSequence: number
baselineExplicitWorkingStartedAt: number | null
}
type TurnStartClaim = {
generation: number
kind: 'hook' | 'lifecycle'
/** Hook turn-start timestamp, or the lifecycle working sequence the turn was attributed to. */
value: number
requestId: string
}
export class AgentPromptRequestCorrelation {
private readonly requestsByPty = new Map<string, AgentPromptRequestBaseline[]>()
private readonly claimsByPty = new Map<string, TurnStartClaim[]>()
register(ptyId: string, request: AgentPromptRequestBaseline): void {
const requests = this.requestsByPty.get(ptyId) ?? []
const existing = requests.findIndex(
(candidate) =>
candidate.generation === request.generation && candidate.requestId === request.requestId
)
if (existing !== -1) {
requests.splice(existing, 1)
}
requests.push(request)
if (requests.length > REQUESTS_PER_PTY_LIMIT) {
requests.splice(0, requests.length - REQUESTS_PER_PTY_LIMIT)
}
this.requestsByPty.set(ptyId, requests)
}
forget(ptyId: string, generation: number, requestId: string): void {
const requests = this.requestsByPty.get(ptyId)
const index = requests?.findIndex(
(candidate) => candidate.generation === generation && candidate.requestId === requestId
)
if (requests && index !== undefined && index !== -1) {
requests.splice(index, 1)
}
}
clearForPty(ptyId: string): void {
this.requestsByPty.delete(ptyId)
this.claimsByPty.delete(ptyId)
}
acceptTurnStart(
ptyId: string,
generation: number,
requestId: string,
baselineWorkingSequence: number,
baselineExplicitWorkingStartedAt: number | null,
evidence: AgentPromptTurnStartEvidence
): boolean {
if (
!isTurnStartAfterBaseline(evidence, {
baselineWorkingSequence,
baselineExplicitWorkingStartedAt
})
) {
return false
}
const requests = this.requestsByPty.get(ptyId) ?? []
const request = requests.find(
(candidate) => candidate.generation === generation && candidate.requestId === requestId
)
// A receipt restored after a runtime restart has no in-memory registration;
// leave it queued rather than attributing an unrelated turn to it.
if (
!request ||
request.baselineWorkingSequence !== baselineWorkingSequence ||
request.baselineExplicitWorkingStartedAt !== baselineExplicitWorkingStartedAt
) {
return false
}
let claim: TurnStartClaim | null
if (evidence.kind === 'lifecycle') {
this.allocateLifecycleClaims(ptyId, generation, evidence)
claim = this.findClaim(ptyId, generation, requestId)
} else {
const first = requests.find(
(candidate) =>
candidate.generation === generation && isTurnStartAfterBaseline(evidence, candidate)
)
if (first && first.requestId !== requestId) {
return false
}
claim = this.nextFreeClaim(ptyId, generation, baselineWorkingSequence, evidence, requestId)
}
if (!claim) {
return false
}
const owner = this.claimOwner(ptyId, claim)
if (owner && owner !== requestId) {
return false
}
this.recordClaim(ptyId, claim)
this.forget(ptyId, generation, requestId)
return true
}
private allocateLifecycleClaims(
ptyId: string,
generation: number,
evidence: Extract<AgentPromptTurnStartEvidence, { kind: 'lifecycle' }>
): void {
for (const candidate of this.requestsByPty.get(ptyId) ?? []) {
if (
candidate.generation !== generation ||
!isTurnStartAfterBaseline(evidence, candidate) ||
this.findClaim(ptyId, generation, candidate.requestId)
) {
continue
}
// A candidate with a later baseline can run out of free sequences while an
// earlier-baselined one still has room, so keep scanning the queue.
const claim = this.nextFreeClaim(
ptyId,
generation,
candidate.baselineWorkingSequence,
evidence,
candidate.requestId
)
if (claim) {
this.recordClaim(ptyId, claim)
}
}
}
private findClaim(ptyId: string, generation: number, requestId: string): TurnStartClaim | null {
return (
this.claimsByPty
.get(ptyId)
?.find((claim) => claim.generation === generation && claim.requestId === requestId) ?? null
)
}
private claimOwner(ptyId: string, claim: TurnStartClaim): string | null {
return (
this.claimsByPty
.get(ptyId)
?.find(
(existing) =>
existing.generation === claim.generation &&
existing.kind === claim.kind &&
existing.value === claim.value
)?.requestId ?? null
)
}
private recordClaim(ptyId: string, claim: TurnStartClaim): void {
const claims = this.claimsByPty.get(ptyId) ?? []
const existing = claims.findIndex(
(candidate) =>
candidate.generation === claim.generation &&
candidate.kind === claim.kind &&
candidate.value === claim.value
)
if (existing === -1) {
claims.push(claim)
if (claims.length > REQUESTS_PER_PTY_LIMIT) {
claims.splice(0, claims.length - REQUESTS_PER_PTY_LIMIT)
}
} else {
claims[existing] = claim
}
this.claimsByPty.set(ptyId, claims)
}
private nextFreeClaim(
ptyId: string,
generation: number,
baselineWorkingSequence: number,
evidence: AgentPromptTurnStartEvidence,
requestId: string
): TurnStartClaim | null {
if (evidence.kind === 'hook') {
return { generation, kind: 'hook', value: evidence.workingStartedAt, requestId }
}
const claimed = new Set(
(this.claimsByPty.get(ptyId) ?? [])
.filter((claim) => claim.generation === generation && claim.kind === 'lifecycle')
.map((claim) => claim.value)
)
let sequence = baselineWorkingSequence + 1
while (claimed.has(sequence)) {
sequence += 1
}
return sequence <= evidence.workingSequence
? { generation, kind: 'lifecycle', value: sequence, requestId }
: null
}
}
function isTurnStartAfterBaseline(
evidence: AgentPromptTurnStartEvidence,
baseline: { baselineWorkingSequence: number; baselineExplicitWorkingStartedAt: number | null }
): boolean {
return evidence.kind === 'lifecycle'
? evidence.workingSequence > baseline.baselineWorkingSequence
: evidence.workingStartedAt > (baseline.baselineExplicitWorkingStartedAt ?? 0)
}
@@ -464,7 +464,7 @@ describe('agent prompt submission runtime', () => {
await vi.runAllTimersAsync()
await expect(submission).resolves.toMatchObject({
prompt: { stages: ['input_accepted', 'queued_pending_turn'] }
prompt: { stages: ['input_accepted'] }
})
})
@@ -573,7 +573,7 @@ describe('agent prompt submission runtime', () => {
})
await vi.runAllTimersAsync()
const first = await firstPromise
expect(first.prompt?.stages).toEqual(['input_accepted', 'queued_pending_turn'])
expect(first.prompt?.stages).toEqual(['input_accepted'])
const firstObserved = runtime.observeTerminalAgentPrompt(handle, first.prompt!, 20_000)
runtime.setPtyController({
@@ -596,11 +596,11 @@ describe('agent prompt submission runtime', () => {
await vi.runAllTimersAsync()
await expect(firstObserved).resolves.toMatchObject({
stages: ['input_accepted', 'submission_observed', 'turn_started']
stages: ['input_accepted', 'turn_started']
})
const second = await secondPromise
expect(second).toMatchObject({
prompt: { stages: ['input_accepted', 'queued_pending_turn'] }
prompt: { stages: ['input_accepted'] }
})
const secondObserved = runtime.observeTerminalAgentPrompt(handle, second.prompt!, 1_000)
@@ -608,7 +608,7 @@ describe('agent prompt submission runtime', () => {
await vi.advanceTimersByTimeAsync(50)
await expect(secondObserved).resolves.toMatchObject({
stages: ['input_accepted', 'submission_observed', 'turn_started']
stages: ['input_accepted', 'turn_started']
})
})
@@ -790,10 +790,10 @@ describe('agent prompt submission runtime', () => {
await vi.runAllTimersAsync()
await expect(firstObserved).resolves.toMatchObject({
stages: ['input_accepted', 'submission_observed', 'turn_started']
stages: ['input_accepted', 'turn_started']
})
await expect(secondObserved).resolves.toMatchObject({
stages: ['input_accepted', 'queued_pending_turn']
stages: ['input_accepted']
})
})
@@ -7,22 +7,10 @@ import type {
AgentPromptWaitTextCache
} from './agent-prompt-submission-verification'
import { verifyAgentPromptSubmission } from './agent-prompt-submission-verification'
const AGENT_PROMPT_CORRELATION_LIMIT_PER_PTY = 1_024
type AgentPromptRequestBaseline = {
ptyId: string
generation: number
requestId: string
baselineWorkingSequence: number
baselineExplicitWorkingStartedAt: number | null
}
import { AgentPromptRequestCorrelation } from './agent-prompt-request-correlation'
export class OrcaRuntimeWithAgentPromptRequestCorrelation extends OrcaRuntimeWithSerializeAgentPromptSubmission {
// Turn evidence is PTY-wide; keep a request owner so one observed turn
// cannot settle every queued prompt that shares the same baseline.
private agentPromptRequestBaselines = new Map<string, AgentPromptRequestBaseline>()
private agentPromptTurnStartClaims = new Map<string, string>()
private readonly agentPromptCorrelation = new AgentPromptRequestCorrelation()
// Declared, not defined: both live further up the mixin chain, so this link cannot see them.
declare protected getLivePtyForHandle: (
handle: string
@@ -95,11 +83,7 @@ export class OrcaRuntimeWithAgentPromptRequestCorrelation extends OrcaRuntimeWit
timeoutMs
})
this.forgetAgentPromptRequest(binding.ptyId, binding.generation, prompt.requestId)
return {
...prompt,
stages: ['input_accepted', 'submission_observed', 'turn_started'],
observation: 'supported'
}
return { ...prompt, stages: ['input_accepted', 'turn_started'], observation: 'supported' }
} catch (error) {
if (error instanceof Error && error.message === 'agent_prompt_stalled') {
return prompt
@@ -119,28 +103,16 @@ export class OrcaRuntimeWithAgentPromptRequestCorrelation extends OrcaRuntimeWit
baselineWorkingSequence: number,
baselineExplicitWorkingStartedAt: number | null
): void {
const requestKey = this.agentPromptRequestKey(ptyId, generation, requestId)
this.agentPromptRequestBaselines.delete(requestKey)
this.agentPromptRequestBaselines.set(requestKey, {
ptyId,
this.agentPromptCorrelation.register(ptyId, {
generation,
requestId,
baselineWorkingSequence,
baselineExplicitWorkingStartedAt
})
const prefix = `${ptyId}\u0000${generation}\u0000`
const matchingKeys = [...this.agentPromptRequestBaselines.keys()].filter((key) =>
key.startsWith(prefix)
)
for (const staleKey of matchingKeys.slice(0, -AGENT_PROMPT_CORRELATION_LIMIT_PER_PTY)) {
this.agentPromptRequestBaselines.delete(staleKey)
}
}
protected forgetAgentPromptRequest(ptyId: string, generation: number, requestId: string): void {
this.agentPromptRequestBaselines.delete(
this.agentPromptRequestKey(ptyId, generation, requestId)
)
this.agentPromptCorrelation.forget(ptyId, generation, requestId)
}
protected acceptAgentPromptTurnStart(
@@ -151,167 +123,17 @@ export class OrcaRuntimeWithAgentPromptRequestCorrelation extends OrcaRuntimeWit
baselineExplicitWorkingStartedAt: number | null,
evidence: AgentPromptTurnStartEvidence
): boolean {
if (
!this.isAgentPromptTurnStartAfterBaseline(evidence, {
baselineWorkingSequence,
baselineExplicitWorkingStartedAt
})
) {
return false
}
// A receipt restored after a runtime restart has no in-memory registration;
// leave it queued rather than attributing an unrelated turn to it.
const requestKey = this.agentPromptRequestKey(ptyId, generation, requestId)
const request = this.agentPromptRequestBaselines.get(requestKey)
if (
!request ||
request.ptyId !== ptyId ||
request.generation !== generation ||
request.baselineWorkingSequence !== baselineWorkingSequence ||
request.baselineExplicitWorkingStartedAt !== baselineExplicitWorkingStartedAt
) {
return false
}
let claimKey: string | null
if (evidence.kind === 'lifecycle') {
this.allocateAgentPromptLifecycleClaims(ptyId, generation, evidence)
claimKey = this.findAgentPromptTurnClaimKey(ptyId, generation, requestId)
} else {
for (const [candidateKey, candidate] of this.agentPromptRequestBaselines) {
if (
candidate.ptyId === ptyId &&
candidate.generation === generation &&
this.isAgentPromptTurnStartAfterBaseline(evidence, candidate)
) {
if (candidateKey !== requestKey) {
return false
}
break
}
}
claimKey = this.getAgentPromptTurnClaimKey(
ptyId,
generation,
baselineWorkingSequence,
evidence
)
}
if (!claimKey) {
return false
}
const owner = this.agentPromptTurnStartClaims.get(claimKey)
if (owner && owner !== requestId) {
return false
}
this.agentPromptTurnStartClaims.set(claimKey, requestId)
this.agentPromptRequestBaselines.delete(requestKey)
const claimPrefix = `${ptyId}\u0000${generation}\u0000`
const matchingClaimKeys = [...this.agentPromptTurnStartClaims.keys()].filter((key) =>
key.startsWith(claimPrefix)
return this.agentPromptCorrelation.acceptTurnStart(
ptyId,
generation,
requestId,
baselineWorkingSequence,
baselineExplicitWorkingStartedAt,
evidence
)
for (const staleKey of matchingClaimKeys.slice(0, -AGENT_PROMPT_CORRELATION_LIMIT_PER_PTY)) {
this.agentPromptTurnStartClaims.delete(staleKey)
}
return true
}
private allocateAgentPromptLifecycleClaims(
ptyId: string,
generation: number,
evidence: Extract<AgentPromptTurnStartEvidence, { kind: 'lifecycle' }>
): void {
for (const candidate of this.agentPromptRequestBaselines.values()) {
if (
candidate.ptyId !== ptyId ||
candidate.generation !== generation ||
!this.isAgentPromptTurnStartAfterBaseline(evidence, candidate)
) {
continue
}
if (this.findAgentPromptTurnClaimKey(ptyId, generation, candidate.requestId)) {
continue
}
const claimKey = this.getAgentPromptTurnClaimKey(
ptyId,
generation,
candidate.baselineWorkingSequence,
evidence
)
if (!claimKey) {
return
}
this.agentPromptTurnStartClaims.set(claimKey, candidate.requestId)
}
}
private findAgentPromptTurnClaimKey(
ptyId: string,
generation: number,
requestId: string
): string | null {
const prefix = `${ptyId}\u0000${generation}\u0000`
for (const [claimKey, owner] of this.agentPromptTurnStartClaims) {
if (claimKey.startsWith(prefix) && owner === requestId) {
return claimKey
}
}
return null
}
private isAgentPromptTurnStartAfterBaseline(
evidence: AgentPromptTurnStartEvidence,
baseline: {
baselineWorkingSequence: number
baselineExplicitWorkingStartedAt: number | null
}
): boolean {
return evidence.kind === 'lifecycle'
? evidence.workingSequence > baseline.baselineWorkingSequence
: evidence.workingStartedAt > (baseline.baselineExplicitWorkingStartedAt ?? 0)
}
private getAgentPromptTurnClaimKey(
ptyId: string,
generation: number,
baselineWorkingSequence: number,
evidence: AgentPromptTurnStartEvidence
): string | null {
const prefix = `${ptyId}\u0000${generation}\u0000`
if (evidence.kind === 'hook') {
return `${prefix}hook:${evidence.workingStartedAt}`
}
const lifecyclePrefix = `${prefix}lifecycle:`
const claimedSequences = new Set<number>()
for (const key of this.agentPromptTurnStartClaims.keys()) {
if (!key.startsWith(lifecyclePrefix)) {
continue
}
const sequence = Number(key.slice(lifecyclePrefix.length))
if (Number.isFinite(sequence)) {
claimedSequences.add(sequence)
}
}
let sequence = baselineWorkingSequence + 1
while (claimedSequences.has(sequence)) {
sequence += 1
}
return sequence <= evidence.workingSequence ? `${lifecyclePrefix}${sequence}` : null
}
private agentPromptRequestKey(ptyId: string, generation: number, requestId: string): string {
return `${ptyId}\u0000${generation}\u0000${requestId}`
}
protected clearAgentPromptCorrelationForPty(ptyId: string): void {
for (const key of this.agentPromptRequestBaselines.keys()) {
if (key.startsWith(`${ptyId}\u0000`)) {
this.agentPromptRequestBaselines.delete(key)
}
}
for (const key of this.agentPromptTurnStartClaims.keys()) {
if (key.startsWith(`${ptyId}\u0000`)) {
this.agentPromptTurnStartClaims.delete(key)
}
}
this.agentPromptCorrelation.clearForPty(ptyId)
}
}
@@ -111,10 +111,7 @@ export class OrcaRuntimeWithWriteTerminalAgentPrompt extends OrcaRuntimeWithReso
: null
const inputAccepted: RuntimeTerminalPromptDelivery = {
requestId: options.requestId,
stages:
baseline.status === 'working'
? ['input_accepted', 'queued_pending_turn']
: ['input_accepted'],
stages: ['input_accepted'],
provider: settlementAgent ?? 'unsupported',
observation: settlementAgent ? 'supported' : 'unsupported',
processIncarnation: binding.processIncarnation,
@@ -165,7 +162,7 @@ export class OrcaRuntimeWithWriteTerminalAgentPrompt extends OrcaRuntimeWithReso
submits: 1,
prompt: {
...inputAccepted,
stages: ['input_accepted', 'submission_observed', 'turn_started']
stages: ['input_accepted', 'turn_started']
}
}
} catch (error) {
@@ -296,7 +296,7 @@ describe('orchestration worker-start prompt contract', () => {
expect(send).toMatchObject({
prompt: {
requestId: 'busy-swallowed',
stages: ['input_accepted', 'queued_pending_turn']
stages: ['input_accepted']
}
})
const observed = runtime.observeTerminalAgentPrompt(handle, send.prompt!, 1_000)
@@ -306,7 +306,7 @@ describe('orchestration worker-start prompt contract', () => {
await vi.runAllTimersAsync()
await expect(observed).resolves.toMatchObject({
stages: ['input_accepted', 'queued_pending_turn']
stages: ['input_accepted']
})
})
@@ -52,7 +52,7 @@ export function ensureUnsupportedTerminalPromptReceipt(
if (send.prompt) {
return send
}
const binding = runtime.getTerminalPromptRequestBinding(handle)!
const binding = runtime.getTerminalPromptRequestBinding(handle)
return {
...send,
prompt: {
@@ -37,12 +37,20 @@ function createHarness() {
const db = new OrchestrationDb(':memory:')
const runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(db)
vi.spyOn(runtime, 'getTerminalPromptRequestBinding').mockReturnValue({
const binding = vi.spyOn(runtime, 'getTerminalPromptRequestBinding').mockReturnValue({
ptyId: 'pty-prompt',
processIncarnation: 'incarnation-1',
generation: 1
})
return { db, executor: new OrchestrationMutationExecutor(runtime) }
// Every handle for this PTY resolves to one pane, so a re-minted handle is the same terminal.
vi.spyOn(runtime, 'getTerminalPaneKey').mockReturnValue('window-1:leaf-prompt')
return {
db,
executor: new OrchestrationMutationExecutor(runtime),
bindTerminal: (next: { generation: number; processIncarnation: string }) => {
binding.mockReturnValue({ ptyId: 'pty-prompt', ...next })
}
}
}
describe('terminal prompt mutation receipt retry boundary', () => {
@@ -109,21 +117,62 @@ describe('terminal prompt mutation receipt retry boundary', () => {
const invoke = vi
.fn()
.mockResolvedValueOnce({
send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } }
send: { prompt: { stages: ['input_accepted'] } }
})
.mockRejectedValueOnce(new Error('terminal was parked'))
await expect(harness.executor.run(request, params, invoke)).resolves.toMatchObject({
send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } },
send: { prompt: { stages: ['input_accepted'] } },
mutation: { replayed: false }
})
await expect(harness.executor.run(request, params, invoke)).resolves.toMatchObject({
send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } },
send: { prompt: { stages: ['input_accepted'] } },
mutation: { replayed: true }
})
expect(invoke).toHaveBeenCalledTimes(2)
})
it('reports a replay as incarnation_replaced once the PTY generation advances', async () => {
const harness = createHarness()
databases.push(harness.db)
const requestId = 'stale-binding-replay'
const invoke = vi.fn().mockResolvedValue({
send: { prompt: { stages: ['input_accepted', 'turn_started'], observation: 'supported' } }
})
await expect(
harness.executor.run(promptRequest(requestId), promptParams, invoke)
).resolves.toMatchObject({ send: { prompt: { observation: 'supported' } } })
harness.bindTerminal({ generation: 2, processIncarnation: 'incarnation-2' })
await expect(
harness.executor.run(promptRequest(requestId), promptParams, invoke)
).resolves.toMatchObject({
send: { prompt: { observation: 'incarnation_replaced' } },
mutation: { replayed: true }
})
expect(invoke).toHaveBeenCalledOnce()
})
it('replays a byte-identical prompt after the handle is re-minted', async () => {
const harness = createHarness()
databases.push(harness.db)
const requestId = 'rebound-handle-replay'
const invoke = vi.fn().mockResolvedValue({
send: { prompt: { stages: ['input_accepted', 'turn_started'], observation: 'supported' } }
})
await harness.executor.run(promptRequest(requestId), promptParams, invoke)
const reminted = { ...promptParams, terminal: 'term_00000000-0000-4000-8000-000000000000' }
const request = { ...promptRequest(requestId), params: reminted }
await expect(harness.executor.run(request, reminted, invoke)).resolves.toMatchObject({
send: { prompt: { observation: 'supported' } },
mutation: { replayed: true }
})
expect(invoke).toHaveBeenCalledOnce()
})
it('keeps an uncheckpointed pending worker_done fenced after restart', async () => {
const harness = createHarness()
databases.push(harness.db)
@@ -12,7 +12,9 @@ import {
getPendingWorkerStartRecovery,
hashCanonical,
isResumablePendingWorkerDone,
markReplayedPromptIncarnationReplaced,
readPromptBasePayloadHash,
readPromptBindingPayloadHash,
replayStableCallerParams,
shouldObserveCompletedMutation
} from './orchestration-mutation-receipt'
@@ -77,6 +79,15 @@ export class OrchestrationMutationExecutor {
`Mutation request ${requestId} was already used with different input.`
)
}
const recordedPromptBindingHash = existingPromptReceipt
? readPromptBindingPayloadHash(existingPromptReceipt.payload_hash)
: null
// The recorded observation is only true while the prompt's terminal incarnation survives, so
// every replay re-checks the binding rather than only the --wait-submit ones.
const promptBindingChanged =
recordedPromptBindingHash !== null &&
recordedPromptBindingHash !==
this.readTerminalPromptBindingHash((params as { terminal: string }).terminal)
const payloadHash = existingPromptReceipt
? existingPromptReceipt.payload_hash
: isPromptMutation
@@ -132,6 +143,13 @@ export class OrchestrationMutationExecutor {
return attachMutationReceipt(await active.promise, requestId, true)
}
const receipt = JSON.parse(begun.row.receipt ?? 'null')
if (promptBindingChanged) {
return attachMutationReceipt(
markReplayedPromptIncarnationReplaced(receipt),
requestId,
true
)
}
if (!shouldObserveCompletedMutation(request.method, params, receipt)) {
return attachMutationReceipt(receipt, requestId, true)
}
@@ -236,6 +254,15 @@ export class OrchestrationMutationExecutor {
}
}
// A replayed prompt may name a terminal that is gone; an unreadable binding is a changed one.
private readTerminalPromptBindingHash(handle: string): string | null {
try {
return hashCanonical(this.runtime.getTerminalPromptRequestBinding(handle))
} catch {
return null
}
}
getLocalAuthenticatedCallerFingerprint(): string {
return this.runtime.getOrchestrationDb().getOrCreateLocalMutationCallerFingerprint()
}
@@ -20,7 +20,7 @@ export function replayStableCallerParams(runtime: OrcaRuntimeService, params: un
const source = params as Record<string, unknown>
const result = { ...source }
delete result.waitSubmitMs
for (const property of ['from', 'callerTerminalHandle'] as const) {
for (const property of ['from', 'callerTerminalHandle', 'terminal'] as const) {
const handle = source[property]
if (typeof handle !== 'string') {
continue
@@ -47,6 +47,27 @@ export function readPromptBasePayloadHash(payloadHash: string): string {
return payloadHash.split(':', 1)[0] ?? payloadHash
}
/** Absent on receipts recorded before the binding was hashed into the payload. */
export function readPromptBindingPayloadHash(payloadHash: string): string | null {
const separator = payloadHash.indexOf(':')
return separator === -1 ? null : payloadHash.slice(separator + 1)
}
/** A stored `observation` only describes the incarnation the prompt was written to. */
export function markReplayedPromptIncarnationReplaced(receipt: unknown): unknown {
if (!receipt || typeof receipt !== 'object' || Array.isArray(receipt)) {
return receipt
}
const send = (receipt as { send?: { prompt?: { observation?: string } } }).send
if (!send?.prompt) {
return receipt
}
return {
...(receipt as Record<string, unknown>),
send: { ...send, prompt: { ...send.prompt, observation: 'incarnation_replaced' } }
}
}
export function shouldObserveCompletedMutation(
method: string,
params: unknown,
@@ -112,7 +112,7 @@ describe('durable terminal prompt delivery receipts', () => {
send: {
prompt: {
provider: agent,
stages: ['input_accepted', 'submission_observed', 'turn_started']
stages: ['input_accepted', 'turn_started']
}
}
}
@@ -138,7 +138,7 @@ describe('durable terminal prompt delivery receipts', () => {
expect(response).toMatchObject({
ok: true,
result: {
send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } }
send: { prompt: { stages: ['input_accepted'] } }
}
})
}
@@ -245,7 +245,7 @@ describe('durable terminal prompt delivery receipts', () => {
result: {
send: {
prompt: {
stages: ['input_accepted', 'submission_observed', 'turn_started']
stages: ['input_accepted', 'turn_started']
}
},
mutation: { replayed: true }
@@ -270,11 +270,11 @@ describe('durable terminal prompt delivery receipts', () => {
const second = await secondPromise
expect(first).toMatchObject({
ok: true,
result: { send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } } }
result: { send: { prompt: { stages: ['input_accepted'] } } }
})
expect(second).toMatchObject({
ok: true,
result: { send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } } }
result: { send: { prompt: { stages: ['input_accepted'] } } }
})
harness.runtime.onPtyData(
@@ -296,13 +296,13 @@ describe('durable terminal prompt delivery receipts', () => {
await expect(firstObserved).resolves.toMatchObject({
ok: true,
result: {
send: { prompt: { stages: ['input_accepted', 'submission_observed', 'turn_started'] } }
send: { prompt: { stages: ['input_accepted', 'turn_started'] } }
}
})
await expect(secondObserved).resolves.toMatchObject({
ok: true,
result: {
send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } }
send: { prompt: { stages: ['input_accepted'] } }
}
})
harness.db.close()
@@ -340,12 +340,12 @@ describe('durable terminal prompt delivery receipts', () => {
await expect(secondObserved).resolves.toMatchObject({
ok: true,
result: { send: { prompt: { stages: ['input_accepted', 'queued_pending_turn'] } } }
result: { send: { prompt: { stages: ['input_accepted'] } } }
})
await expect(firstObserved).resolves.toMatchObject({
ok: true,
result: {
send: { prompt: { stages: ['input_accepted', 'submission_observed', 'turn_started'] } }
send: { prompt: { stages: ['input_accepted', 'turn_started'] } }
}
})
harness.db.close()
@@ -378,7 +378,7 @@ describe('durable terminal prompt delivery receipts', () => {
result: {
send: {
prompt: {
stages: ['input_accepted', 'queued_pending_turn'],
stages: ['input_accepted'],
observation: 'incarnation_replaced'
}
},
+1 -5
View File
@@ -218,11 +218,7 @@ export type RuntimeTerminalSend = {
prompt?: RuntimeTerminalPromptDelivery
}
export type RuntimeTerminalPromptStage =
| 'input_accepted'
| 'queued_pending_turn'
| 'submission_observed'
| 'turn_started'
export type RuntimeTerminalPromptStage = 'input_accepted' | 'turn_started'
export type RuntimeTerminalPromptDelivery = {
requestId: string