diff --git a/config/reliability-gates.jsonc b/config/reliability-gates.jsonc index 0a7a6feadf8..c56dd9941e8 100644 --- a/config/reliability-gates.jsonc +++ b/config/reliability-gates.jsonc @@ -13289,7 +13289,7 @@ "https://github.com/stablyai/orca/issues/13821", "https://github.com/stablyai/orca/issues/14347" ], - "invariant": "Injected orchestration task prompts for recognized agent CLIs must send the prompt body inside one bracketed-paste frame, sanitize embedded ESC bytes, preserve chunk boundaries without losing the frame, and submit exactly once only after the agent can accept Enter. A successful orchestration.workerStart must durably record exactly one accepted and started turn; a swallowed Enter must fail with agent_prompt_stalled and never trigger a blind rescue Enter. Claude and Codex must emit a post-paste composer marker and then settle, or reach the bounded fallback first; every other agent retains the platform delay.", + "invariant": "Injected orchestration task prompts for recognized agent CLIs must send the prompt body inside one bracketed-paste frame, sanitize embedded ESC bytes, preserve chunk boundaries without losing the frame, and submit exactly once only after the agent can accept Enter. Local worker-start with supported observation must preserve an unobserved turn as start_unknown without revoking authority, closing questions, or triggering a rescue Enter; a worker report during observation must settle normally. Claude and Codex must emit a post-paste composer marker and then settle, or reach the bounded fallback first; every other agent retains the platform delay.", "oracle": "Runtime tests assert the exact PTY write sequence, failure cleanup, Claude/Codex marker-gated multi-frame renders, and the legacy platform delay for every other configured agent. The candidate resets settlement on later frames, gives a late marker a fresh bounded window, and still submits once at the hard deadline if output never settles. The worker-start contract drives the production RPC through a delayed fake Codex composer and independently checks exact turn/Enter counts plus reopened SQLite Task, Dispatch, worker receipt, and mutation receipt state for accepted and swallowed outcomes. Other orchestration tests assert dispatch/coordinator use the agent prompt path; the live CLI harness covers long Codex-like framing.", "commands": [ "pnpm exec vitest run --config config/vitest.config.ts src/shared/agent-prompt-injection.test.ts src/main/runtime/orca-runtime.test.ts src/main/runtime/rpc/methods/orchestration/runs/tasks-dispatch.test.ts src/main/runtime/orchestration/coordinator.test.ts", @@ -13339,7 +13339,8 @@ "file": "src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts", "assertions": [ "delayed composer readiness produces exactly one submitted and started turn with no premature Enter and durable ready receipts", - "a swallowed Enter records agent_prompt_stalled across Task, Dispatch, worker, and mutation receipts without a rescue Enter" + "a swallowed Enter durably records start_unknown without a rescue Enter or capability revocation", + "early worker reports settle during observation, and outstanding questions survive observation uncertainty" ] }, { diff --git a/config/scripts/locale-collator-sort-benchmark.mjs b/config/scripts/locale-collator-sort-benchmark.mjs index 3d34466532a..1a0b08e1643 100644 --- a/config/scripts/locale-collator-sort-benchmark.mjs +++ b/config/scripts/locale-collator-sort-benchmark.mjs @@ -129,6 +129,7 @@ for (const count of [36, 50, 250]) { const issues = makeJiraIssues(count) const before = () => [...issues] + // oxlint-disable-next-line sort-comparator-performance/no-repeated-collator -- Baseline measures per-comparison setup against a reused collator. .sort((a, b) => a.key.localeCompare(b.key, undefined, { numeric: true })) .map((issue) => issue.key) const after = () => sortJiraIssues(issues, 'key', 'asc').map((issue) => issue.key) @@ -142,6 +143,7 @@ for (const count of [36, 50, 250]) { for (const count of [10, 50, 250]) { const values = makeBaseSensitivityValues(count) const before = () => + // oxlint-disable-next-line sort-comparator-performance/no-repeated-collator -- Baseline measures per-comparison setup against a reused collator. [...values].sort((a, b) => a.localeCompare(b, undefined, { sensitivity: 'base' })) const after = () => [...values].sort(compareBaseSensitivityLocaleText) assertSameOrder(before, after, `base ${count}`) diff --git a/src/cli/handlers/orchestration-worker-cli.test.ts b/src/cli/handlers/orchestration-worker-cli.test.ts index cddd36a4cf7..931c4477855 100644 --- a/src/cli/handlers/orchestration-worker-cli.test.ts +++ b/src/cli/handlers/orchestration-worker-cli.test.ts @@ -104,6 +104,34 @@ describe('orchestration worker-start CLI contract', () => { expect(process.exitCode).toBeUndefined() }) + it.each(['succeeded', 'failed'])( + 'accepts a successful start whose task already %s', + async (workerOutcome) => { + const receipt = { + taskId: 'task_1', + dispatchId: 'ctx_1', + state: 'ready', + stage: 'settled', + workerOutcome, + effects: [], + residualResources: [] + } + callMock.mockResolvedValue({ result: receipt }) + await invokeWorkerStart( + new Map([ + ['task', 'task_1'], + ['from', 'term_coord'] + ]) + ) + expect(process.exitCode).toBeUndefined() + expect(printResult).toHaveBeenCalledWith( + expect.objectContaining({ result: receipt }), + true, + expect.any(Function) + ) + } + ) + it('capability-gates and forwards per-invocation launch preferences', async () => { callMock .mockResolvedValueOnce({ diff --git a/src/main/claude/claude-structured-dispatch-content.test.ts b/src/main/claude/claude-structured-dispatch-content.test.ts new file mode 100644 index 00000000000..02ee6fa721d --- /dev/null +++ b/src/main/claude/claude-structured-dispatch-content.test.ts @@ -0,0 +1,165 @@ +import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { describe, expect, it } from 'vitest' +import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types' +import { + claudeDispatchInvokesSlashCommand, + claudeDispatchMessageContent +} from './claude-structured-dispatch-content' + +const PNG = Buffer.from( + 'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mP8z8BQDwAEhQGAhKmMIQAAAABJRU5ErkJggg==', + 'base64' +) + +function userMessage(blocks: AgentJournalMessageItem['blocks']): AgentJournalMessageItem { + return { kind: 'message', role: 'user', blocks } +} + +const REMOTE_IMAGE = { type: 'image-ref' as const, url: 'https://example.test/a.png' } + +describe('claudeDispatchMessageContent', () => { + it('puts the text block last so a slash command still expands with an attachment', async () => { + const content = await claudeDispatchMessageContent( + // The composer builds text-then-images; Claude only treats a leading `/` as a + // command when the LAST block is text. + userMessage([{ type: 'text', text: '/goal ship the parser' }, REMOTE_IMAGE]) + ) + + expect(content).toEqual([ + { type: 'image', source: { type: 'url', url: 'https://example.test/a.png' } }, + { type: 'text', text: '/goal ship the parser' } + ]) + }) + + it('keeps every image ahead of the text and preserves each side’s order', async () => { + const second = { type: 'image-ref' as const, url: 'https://example.test/b.png' } + + const content = await claudeDispatchMessageContent( + userMessage([{ type: 'text', text: 'look' }, REMOTE_IMAGE, second]) + ) + + expect(content.map((part) => (part as { type: string }).type)).toEqual([ + 'image', + 'image', + 'text' + ]) + expect(content[0]).toEqual({ + type: 'image', + source: { type: 'url', url: 'https://example.test/a.png' } + }) + expect(content[1]).toEqual({ + type: 'image', + source: { type: 'url', url: 'https://example.test/b.png' } + }) + }) + + it('sends text alone unchanged', async () => { + const content = await claudeDispatchMessageContent(userMessage([{ type: 'text', text: 'hi' }])) + + expect(content).toEqual([{ type: 'text', text: 'hi' }]) + }) + + it('sends an image with no text', async () => { + const content = await claudeDispatchMessageContent(userMessage([REMOTE_IMAGE])) + + expect(content).toEqual([ + { type: 'image', source: { type: 'url', url: 'https://example.test/a.png' } } + ]) + }) + + it('rejects a message with no renderable block', async () => { + await expect( + claudeDispatchMessageContent(userMessage([{ type: 'text', text: '' }])) + ).rejects.toThrow('Claude dispatch requires text or an image') + }) + + it('rejects a non-user message', async () => { + await expect( + claudeDispatchMessageContent({ + ...userMessage([{ type: 'text', text: 'hi' }]), + role: 'assistant' + }) + ).rejects.toThrow('Claude dispatch accepts only user messages') + }) + + it('joins several text blocks so a command is not stranded ahead of trailing prose', async () => { + // Appending each block would leave `thanks` trailing, and Claude reads only that block. + const content = await claudeDispatchMessageContent( + userMessage([ + { type: 'text', text: '/goal ship' }, + REMOTE_IMAGE, + { type: 'text', text: 'thanks' } + ]) + ) + + expect(content).toEqual([ + { type: 'image', source: { type: 'url', url: 'https://example.test/a.png' } }, + { type: 'text', text: '/goal ship\nthanks' } + ]) + expect(claudeDispatchInvokesSlashCommand(content)).toBe(true) + }) + + it('puts a locally attached image ahead of the text, the shape the composer sends', async () => { + const dir = await mkdtemp(join(tmpdir(), 'claude-dispatch-content-')) + const path = join(dir, 'shot.png') + await writeFile(path, PNG) + + try { + const content = await claudeDispatchMessageContent( + userMessage([ + { type: 'text', text: '/goal ship' }, + { type: 'image-ref', path } + ]) + ) + + expect(content).toEqual([ + { + type: 'image', + source: { type: 'base64', media_type: 'image/png', data: PNG.toString('base64') } + }, + { type: 'text', text: '/goal ship' } + ]) + } finally { + await rm(dir, { recursive: true, force: true }) + } + }) +}) + +describe('claudeDispatchInvokesSlashCommand', () => { + it('reads the trailing prompt Claude recovers, not any text block', () => { + expect( + claudeDispatchInvokesSlashCommand([ + { type: 'image', source: { type: 'url', url: 'https://example.test/a.png' } }, + { type: 'text', text: '/goal ship' } + ]) + ).toBe(true) + // The pre-fix order: Claude recovers no prompt at all, so no command runs. + expect( + claudeDispatchInvokesSlashCommand([ + { type: 'text', text: '/goal ship' }, + { type: 'image', source: { type: 'url', url: 'https://example.test/a.png' } } + ]) + ).toBe(false) + }) + + it('reads the joined prompt, so a command behind leading prose is not one', async () => { + // Keeping the blocks separate would leave `/goal ship` trailing and falsely claim a command. + const content = await claudeDispatchMessageContent( + userMessage([ + { type: 'text', text: 'take a look' }, + { type: 'text', text: '/goal ship' } + ]) + ) + + expect(content).toEqual([{ type: 'text', text: 'take a look\n/goal ship' }]) + expect(claudeDispatchInvokesSlashCommand(content)).toBe(false) + }) + + it('matches untrimmed, as Claude does, and ignores a promptless turn', () => { + expect(claudeDispatchInvokesSlashCommand([{ type: 'text', text: ' /goal ship' }])).toBe(false) + expect(claudeDispatchInvokesSlashCommand([{ type: 'text', text: 'ship it' }])).toBe(false) + expect(claudeDispatchInvokesSlashCommand([])).toBe(false) + }) +}) diff --git a/src/main/claude/claude-structured-dispatch-content.ts b/src/main/claude/claude-structured-dispatch-content.ts index 71f180bc3ac..7df84091b7c 100644 --- a/src/main/claude/claude-structured-dispatch-content.ts +++ b/src/main/claude/claude-structured-dispatch-content.ts @@ -3,6 +3,7 @@ import { open } from 'node:fs/promises' import { extname } from 'node:path' import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types' import type { NativeChatBlock } from '../../shared/native-chat-types' +import { claudeRecord } from './claude-structured-item-translation' const MAX_IMAGE_BYTES = 5 * 1024 * 1024 const MAX_IMAGE_COUNT = 20 @@ -88,27 +89,49 @@ async function imageContent( } } +/** + * Claude encodes a user turn as attachment blocks followed by the typed text, and recovers the + * typed prompt by reading only the trailing text block. Verified against the real CLI over + * stream-json: a body ending in an image has no recoverable prompt, so its `/command` reaches + * the model as prose instead of being expanded. + */ export async function claudeDispatchMessageContent( body: AgentJournalMessageItem ): Promise { if (body.role !== 'user') { throw new Error('Claude dispatch accepts only user messages') } - const content: unknown[] = [] + const images: unknown[] = [] + const texts: string[] = [] const imageBudget: ImageBudget = { count: 0, localBytes: 0 } for (const block of body.blocks as NativeChatBlock[]) { if (block.type === 'text' && block.text.length > 0) { - content.push({ type: 'text', text: block.text }) + texts.push(block.text) } else if (block.type === 'image-ref') { - content.push(await imageContent(block, imageBudget)) + images.push(await imageContent(block, imageBudget)) } } + // Join rather than append each block: only the trailing text is read as the prompt, so several + // text blocks would silently discard every one but the last. + const content = texts.length > 0 ? [...images, { type: 'text', text: texts.join('\n') }] : images if (content.length === 0) { throw new Error('Claude dispatch requires text or an image') } return content } +/** The prompt Claude recovers from a dispatch, or null when the turn carries no prompt. */ +function claudeDispatchPrompt(content: readonly unknown[]): string | null { + const last = claudeRecord(content.at(-1)) + return last?.type === 'text' && typeof last.text === 'string' ? last.text : null +} + +/** Mirrors how Claude decides a turn is a command. Untrimmed on purpose: Claude does not trim + * here either, so leading whitespace really does mean no command runs. */ +export function claudeDispatchInvokesSlashCommand(content: readonly unknown[]): boolean { + return claudeDispatchPrompt(content)?.startsWith('/') === true +} + /** * Keep waiter metadata bounded even when a dispatch contains large base64 images. * The digest is only diagnostic: replay acknowledgement must use provider identity. @@ -117,10 +140,7 @@ export function claudeDispatchContentKey(content: readonly unknown[]): string { const digest = createHash('sha256') const summary = content .map((part) => { - const record = - typeof part === 'object' && part !== null && !Array.isArray(part) - ? (part as Record) - : null + const record = claudeRecord(part) const type = typeof record?.type === 'string' ? record.type : 'unknown' if (type === 'text') { return `text:${typeof record?.text === 'string' ? record.text.length : 0}` @@ -136,10 +156,7 @@ export function claudeDispatchContentKey(content: readonly unknown[]): string { }) .join(',') for (const [index, part] of content.entries()) { - const record = - typeof part === 'object' && part !== null && !Array.isArray(part) - ? (part as Record) - : null + const record = claudeRecord(part) const type = typeof record?.type === 'string' ? record.type : 'unknown' digest.update(`${index}:${type}:`) if (type === 'text' && typeof record?.text === 'string') { diff --git a/src/main/claude/claude-structured-dispatch.test.ts b/src/main/claude/claude-structured-dispatch.test.ts index ad09357df58..933bd77774b 100644 --- a/src/main/claude/claude-structured-dispatch.test.ts +++ b/src/main/claude/claude-structured-dispatch.test.ts @@ -423,6 +423,73 @@ describe('Claude structured dispatch image limits', () => { }) }) + it('accepts a slash command sent with an attachment from its result receipt', async () => { + const session = sessionFor() + const dispatched = dispatchClaudeTurn( + session, + { + clientMessageId: 'client-1', + body: userMessage([ + { type: 'text', text: '/permissions' }, + { type: 'image-ref', url: 'https://example.test/a.png' } + ]) + }, + 100 + ) + await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1)) + // The mapper moves the image ahead of the prompt, so Claude runs the command and replies + // with a result receipt instead of a user replay. + expect( + resolveClaudeReplayWaiter(session, { + type: 'result', + subtype: 'success', + session_id: 'provider-session', + uuid: 'command-result-uuid' + }) + ).toBe(false) + + await expect(dispatched).resolves.toMatchObject({ + state: 'accepted', + providerIdentity: { uuid: 'command-result-uuid' } + }) + // The sent order is the fix: the waiter's verdict alone was already what it is today. + expect(session.connection.send).toHaveBeenCalledWith( + expect.objectContaining({ + message: { + role: 'user', + content: [ + { type: 'image', source: { type: 'url', url: 'https://example.test/a.png' } }, + { type: 'text', text: '/permissions' } + ] + } + }) + ) + }) + + it('does not take a result receipt for leading whitespace Claude never reads as a command', async () => { + const session = sessionFor() + const dispatched = dispatchClaudeTurn( + session, + { + clientMessageId: 'client-1', + body: userMessage([{ type: 'text', text: ' /permissions' }]) + }, + 100 + ) + await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1)) + + expect( + resolveClaudeReplayWaiter(session, { + type: 'result', + subtype: 'success', + session_id: 'provider-session', + uuid: 'unrelated-result-uuid' + }) + ).toBe(false) + + await expect(dispatched).resolves.toMatchObject({ state: 'unknown' }) + }) + it('correlates a later slash-command result by user_message_uuid despite a timed-out slash waiter', async () => { const session = sessionFor() const first = dispatchClaudeTurn( diff --git a/src/main/claude/claude-structured-dispatch.ts b/src/main/claude/claude-structured-dispatch.ts index b7619a1e94e..3080b5479e4 100644 --- a/src/main/claude/claude-structured-dispatch.ts +++ b/src/main/claude/claude-structured-dispatch.ts @@ -12,6 +12,7 @@ import type { ClaudeDispatchWaiter, ClaudeSession } from './claude-structured-se import { readClaudeFrameString } from './claude-structured-init-proof' import { claudeDispatchContentKey, + claudeDispatchInvokesSlashCommand, claudeDispatchMessageContent } from './claude-structured-dispatch-content' @@ -231,9 +232,9 @@ export async function dispatchClaudeTurn( return { state: 'rejected', reason: (error as Error).message } } const dispatchSequence = ++session.dispatchSequence - const acceptsResult = input.body.blocks.some( - (block) => block.type === 'text' && block.text.trimStart().startsWith('/') - ) + // Read the sent content, not the journal blocks: only the mapped trailing prompt decides + // whether Claude runs a command, so the two cannot disagree about which frame settles this. + const acceptsResult = claudeDispatchInvokesSlashCommand(content) const sentUuid = randomUUID() const replay = waitForReplay( session, diff --git a/src/main/codex/codex-structured-item-translation.test.ts b/src/main/codex/codex-structured-item-translation.test.ts index 45afb0c9fa5..201df58bc90 100644 --- a/src/main/codex/codex-structured-item-translation.test.ts +++ b/src/main/codex/codex-structured-item-translation.test.ts @@ -584,6 +584,16 @@ describe('codex item bodies', () => { }) }) + it('preserves plan prose documents byte-for-byte as status text', () => { + const text = + ' # Implementation plan\r\n\r\n- [ ] Preserve prose\r\n- [x] Keep café → 日本語\r\n\r\n```ts\r\nconst task = "pending"\r\n```\r\n ' + + expect(codexJournalItem({ type: 'plan', id: 'plan-document', text })).toEqual({ + body: { kind: 'status', text, presentation: 'plan-document' }, + handled: true + }) + }) + it('renders reasoning as status and exposes an unknown item as a provider frame', () => { expect(codexItemBody({ type: 'reasoning', id: 'r', text: 'thinking' })).toEqual({ kind: 'status', diff --git a/src/main/runtime/orchestration/db/contract-constants.ts b/src/main/runtime/orchestration/db/contract-constants.ts index 47e0cd6165a..56c2e2d542f 100644 --- a/src/main/runtime/orchestration/db/contract-constants.ts +++ b/src/main/runtime/orchestration/db/contract-constants.ts @@ -1,7 +1,17 @@ -import { ORCHESTRATION_LEGACY_RUN_ID } from '../../../../shared/orchestration-rpc-contract' +import { + ORCHESTRATION_LEGACY_RUN_ID, + ORCHESTRATION_UNBOUND_RUN_ID +} from '../../../../shared/orchestration-rpc-contract' import { ORCHESTRATION_CONTRACT_VERSION } from '../../../../shared/protocol-version' export const LEGACY_RUN_ID = ORCHESTRATION_LEGACY_RUN_ID +export const UNBOUND_RUN_ID = ORCHESTRATION_UNBOUND_RUN_ID + +// Why: a v1.4.198 coordinator sends no Run id, so its remote workers file mail under a per-attachment stub Run. +export const FEDERATED_STUB_HOME_RUN_ID_PREFIX = 'run_federated_' +export function federatedStubHomeRunId(dispatchId: string): string { + return `${FEDERATED_STUB_HOME_RUN_ID_PREFIX}${dispatchId}` +} export const LEGACY_CONTRACT_VERSION = 0 export const CURRENT_CONTRACT_VERSION = ORCHESTRATION_CONTRACT_VERSION diff --git a/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts b/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts index 3584a2e59a6..ddb0e383143 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/worker-report-settlement.ts @@ -109,9 +109,13 @@ export function settleWorkerReportInTransaction( (dispatch.status === 'pending' || dispatch.status === 'dispatched') && task.status === 'blocked' && reportingWorker?.state === 'start_unknown' + const reportingStart = + dispatch.status === 'pending' && + task.status === 'dispatched' && + reportingWorker?.state === 'starting' const previousDispatchStatus = settledByUnobservedPrompt ? 'failed' - : reconnectingStart + : reconnectingStart || reportingStart ? dispatch.status : 'dispatched' const previousTaskStatus = settledByUnobservedPrompt @@ -198,7 +202,7 @@ export function settleWorkerReportInTransaction( const dispatchTransition = transitionLifecycleWithDb(this.db, { entity: 'dispatch', id: params.dispatchId, - from: reconnectingStart ? ['pending', 'dispatched'] : 'dispatched', + from: reconnectingStart || reportingStart ? ['pending', 'dispatched'] : 'dispatched', to: expectedDispatchStatus, projection: { completed_at: new Date().toISOString(), @@ -234,11 +238,11 @@ export function settleWorkerReportInTransaction( projection: { stage: 'settled', updated_at: new Date().toISOString() }, correction: 'unobserved_prompt_report' }) - } else if (reconnectingStart && params.outcome === 'succeeded') { + } else if ((reconnectingStart || reportingStart) && params.outcome === 'succeeded') { transitionLifecycleWithDb(this.db, { entity: 'worker', id: params.dispatchId, - from: 'start_unknown', + from: reportingStart ? 'starting' : 'start_unknown', to: 'ready' }) transitionLifecycleWithDb(this.db, { @@ -253,7 +257,7 @@ export function settleWorkerReportInTransaction( entity: 'worker', id: params.dispatchId, // A start_unknown success report reconnects through 'ready' above; only failure settles here. - from: params.outcome === 'succeeded' ? 'ready' : ['ready', 'start_unknown'], + from: params.outcome === 'succeeded' ? 'ready' : ['ready', 'start_unknown', 'starting'], to: params.outcome === 'succeeded' ? 'succeeded' : 'failed', projection: { stage: 'settled', updated_at: new Date().toISOString() } }) diff --git a/src/main/runtime/orchestration/db/federation/federated-stub-home-run-backfill.ts b/src/main/runtime/orchestration/db/federation/federated-stub-home-run-backfill.ts new file mode 100644 index 00000000000..47c596e260d --- /dev/null +++ b/src/main/runtime/orchestration/db/federation/federated-stub-home-run-backfill.ts @@ -0,0 +1,16 @@ +import type Database from '../../../../sqlite/sync-database' +import { FEDERATED_STUB_HOME_RUN_ID_PREFIX } from '../contract-constants' + +// Why: a rolled-back v1.4.198 host inserts attachments with home_run_id='' after user_version is +// already 40, so this idempotent repair runs on every open, not only inside the v40 migration. +export function backfillFederatedStubHomeRuns(db: Database.Database): void { + db.exec(` + INSERT OR IGNORE INTO runs (id, objective, home_database, consumer_generation, legacy) + SELECT '${FEDERATED_STUB_HOME_RUN_ID_PREFIX}' || dispatch_id, + 'Coordinated from ' || home_peer_fingerprint, 'remote', 0, 0 + FROM remote_dispatch_attachments WHERE home_run_id = ''; + UPDATE remote_dispatch_attachments + SET home_run_id = '${FEDERATED_STUB_HOME_RUN_ID_PREFIX}' || dispatch_id + WHERE home_run_id = ''; + `) +} diff --git a/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-create.ts b/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-create.ts index d927cf96838..5a6e89b65ca 100644 --- a/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-create.ts +++ b/src/main/runtime/orchestration/db/federation/remote-dispatch-attachment-create.ts @@ -2,13 +2,15 @@ import type { WorkerDispatchState, RemoteDispatchAttachmentRow } from '../../typ import { OrchestrationError } from '../../orchestration-error' import { ensureMutationReceiptCapacity } from '../../mutation-receipt-capacity' import type { OrchestrationDb } from '../orchestration-db' +import { federatedStubHomeRunId } from '../contract-constants' import { insertRemoteDispatchAttachmentRow } from '../dispatch-row-writer' export function createRemoteDispatchAttachment( this: OrchestrationDb, params: { dispatchId: string - runId: string + /** Absent from a v1.4.198 coordinator; replaced by a per-attachment stub Run. */ + runId?: string taskId: string homePeerFingerprint: string protocolVersion: number @@ -44,7 +46,8 @@ export function createRemoteDispatchAttachment( `Remote attachment request ${params.mutationReceipt.requestId} already exists.` ) } - if (!params.runId?.trim()) { + const runId = params.runId ?? federatedStubHomeRunId(params.dispatchId) + if (!runId.trim()) { throw new OrchestrationError('invalid_argument', 'Missing Run ID') } this.db @@ -52,8 +55,8 @@ export function createRemoteDispatchAttachment( `INSERT OR IGNORE INTO runs (id, objective, home_database, consumer_generation, legacy) VALUES (?, ?, 'remote', 0, 0)` ) - .run(params.runId, `Coordinated from ${params.homePeerFingerprint}`) - this.requireRun(params.runId) + .run(runId, `Coordinated from ${params.homePeerFingerprint}`) + this.requireRun(runId) ensureMutationReceiptCapacity(this.db) this.db .prepare( @@ -70,7 +73,7 @@ export function createRemoteDispatchAttachment( ) insertRemoteDispatchAttachmentRow(this.db, { dispatchId: params.dispatchId, - runId: params.runId, + runId, taskId: params.taskId, homePeerFingerprint: params.homePeerFingerprint, protocolVersion: params.protocolVersion, diff --git a/src/main/runtime/orchestration/db/messages/message-insert.ts b/src/main/runtime/orchestration/db/messages/message-insert.ts index 82573305479..51ce7aaaa50 100644 --- a/src/main/runtime/orchestration/db/messages/message-insert.ts +++ b/src/main/runtime/orchestration/db/messages/message-insert.ts @@ -3,6 +3,7 @@ import { generateId } from '../generated-id' import { exposeMessageTimestamps } from '../utc-timestamp' import type { OrchestrationDb } from '../orchestration-db' import { runLifecycleWriteTransaction } from '../lifecycle-write-transaction-runner' +import { UNBOUND_RUN_ID } from '../contract-constants' // ── Messages ── @@ -25,9 +26,17 @@ export type MessageInsert = { } export function insertMessage(this: OrchestrationDb, msg: MessageInsert): MessageRow { - const runId = msg.runId - if (!runId) { - throw new Error('Run is required') + // A sender in no Run (two plain terminals, `send --to `) still gets durable mail. It is + // filed under the unbound Run, never the legacy one, which the schema-skew probe reads as pre-Runs. + // Created on first use so `run list` shows it only to a user who has such mail. + const runId = msg.runId ?? UNBOUND_RUN_ID + if (msg.runId == null) { + this.db + .prepare( + `INSERT OR IGNORE INTO runs (id, objective, home_database, consumer_generation, legacy) + VALUES (?, 'Mail from terminals in no Run', 'this_database', 0, 0)` + ) + .run(UNBOUND_RUN_ID) } const deliveryContract = msg.deliveryContract ?? 'current_delivery' this.requireRun(runId) diff --git a/src/main/runtime/orchestration/db/orchestration-db.ts b/src/main/runtime/orchestration/db/orchestration-db.ts index 2a970841a38..1ce52e96c4a 100644 --- a/src/main/runtime/orchestration/db/orchestration-db.ts +++ b/src/main/runtime/orchestration/db/orchestration-db.ts @@ -1,6 +1,7 @@ import Database from '../../../sqlite/sync-database' import { attachOrchestrationDbMethods } from './attach-orchestration-db-methods' import { hardenOrchestrationDatabaseFiles } from './database-file-permissions' +import { backfillFederatedStubHomeRuns } from './federation/federated-stub-home-run-backfill' import type { OrchestrationDbMethods } from './orchestration-db-methods' import { createCoordinatorMailRoutingTrigger, @@ -28,6 +29,7 @@ class OrchestrationDbCore { this.db.pragma('busy_timeout = 5000') createTables.call(this as unknown as OrchestrationDb) migrate.call(this as unknown as OrchestrationDb) + backfillFederatedStubHomeRuns(this.db) createCoordinatorMailRoutingTrigger.call(this as unknown as OrchestrationDb) rememberCurrentRunCoordinatorHandles.call(this as unknown as OrchestrationDb) hardenOrchestrationDatabaseFiles(dbPath) diff --git a/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts b/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts index 0897cb852de..d674a5298e7 100644 --- a/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts +++ b/src/main/runtime/orchestration/db/schema/create-graph-tables-sql.ts @@ -38,7 +38,8 @@ CREATE TABLE IF NOT EXISTS federated_dispatches ( ); CREATE TABLE IF NOT EXISTS remote_dispatch_attachments ( - home_run_id TEXT NOT NULL, + -- DEFAULT: a rolled-back v1.4.198 host still inserts here without a home Run. + home_run_id TEXT NOT NULL DEFAULT '', dispatch_id TEXT PRIMARY KEY, task_id TEXT NOT NULL, home_peer_fingerprint TEXT NOT NULL, diff --git a/src/main/runtime/orchestration/db/schema/federated-home-run-migration.test.ts b/src/main/runtime/orchestration/db/schema/federated-home-run-migration.test.ts index b970435223a..4fcd78634f3 100644 --- a/src/main/runtime/orchestration/db/schema/federated-home-run-migration.test.ts +++ b/src/main/runtime/orchestration/db/schema/federated-home-run-migration.test.ts @@ -1,26 +1,105 @@ -import { afterEach, describe, expect, it } from 'vitest' +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { ORCHESTRATION_CONTRACT_VERSION } from '../../../../../shared/protocol-version' import { OrchestrationDb } from '../orchestration-db' +import { SCHEMA_VERSION, federatedStubHomeRunId } from '../contract-constants' import { migrateV40 } from './migrate-v40' import { importFederatedControlMessage } from '../../federation-control-message' describe('federated home Run migration', () => { - const db = new OrchestrationDb(':memory:') + let db: OrchestrationDb + beforeEach(() => { + db = new OrchestrationDb(':memory:') + }) afterEach(() => db.close()) - it('adds the home Run column and refuses mail for a development placeholder', () => { + function importInstruction(target: OrchestrationDb, dispatchId: string, messageId: string): void { + expect( + importFederatedControlMessage(target, { + dispatchId, + messageId, + payload: JSON.stringify({ from: 'home', subject: 'Instruction', body: '', type: 'status' }) + }) + ).toEqual({ imported: true, type: 'status' }) + } + + function attachWithoutRunId(target: OrchestrationDb, dispatchId: string, runId?: string): void { + target.createRemoteDispatchAttachment({ + dispatchId, + runId, + taskId: `task_${dispatchId}`, + homePeerFingerprint: 'home', + protocolVersion: ORCHESTRATION_CONTRACT_VERSION, + runtimeEpoch: 'epoch', + mutationReceipt: { + callerFingerprint: 'home', + requestId: `request_${dispatchId}`, + method: 'orchestration.federationAttachStart', + payloadHash: `payload_${dispatchId}` + } + }) + } + + it('backfills a pre-upgrade attachment with a stub home Run that keeps its mailbox', () => { db.db.exec('ALTER TABLE remote_dispatch_attachments DROP COLUMN home_run_id') db.db.exec(`INSERT INTO remote_dispatch_attachments (dispatch_id, task_id, home_peer_fingerprint, runtime_epoch) VALUES ('ctx_old', 'task_old', 'home', 'epoch')`) migrateV40.call(db, 39) - expect(db.getRemoteDispatchAttachment('ctx_old')?.home_run_id).toBe('') - expect(() => - importFederatedControlMessage(db, { - dispatchId: 'ctx_old', - messageId: 'message_old', - payload: JSON.stringify({ from: 'home', subject: 'Instruction', body: '', type: 'message' }) - }) - ).toThrow('Run not found:') - expect(db.getMessageById('message_old')).toBeUndefined() + const stubRunId = federatedStubHomeRunId('ctx_old') + expect(db.getRemoteDispatchAttachment('ctx_old')?.home_run_id).toBe(stubRunId) + expect(db.getRunRaw(stubRunId)).toMatchObject({ home_database: 'remote', legacy: 0 }) + importInstruction(db, 'ctx_old', 'message_old') + expect(db.getMessageById('message_old')?.run_id).toBe(stubRunId) + }) + + it('mints a stub home Run when a v1.4.198 coordinator attaches without a Run id', () => { + attachWithoutRunId(db, 'ctx_legacy_home') + const stubRunId = federatedStubHomeRunId('ctx_legacy_home') + expect(db.getRemoteDispatchAttachment('ctx_legacy_home')?.home_run_id).toBe(stubRunId) + importInstruction(db, 'ctx_legacy_home', 'message_legacy_home') + expect(db.getMessageById('message_legacy_home')?.run_id).toBe(stubRunId) + }) + + it('rejects a whitespace-only Run id instead of minting a stub', () => { + expect(() => attachWithoutRunId(db, 'ctx_blank', ' ')).toThrow('Missing Run ID') + expect(db.getRemoteDispatchAttachment('ctx_blank')).toBeUndefined() + }) + + it('repairs rows a rolled-back v1.4.198 host inserted after user_version reached 40', () => { + const dir = mkdtempSync(join(tmpdir(), 'orca-federated-home-run-')) + const dbPath = join(dir, 'orchestration.db') + try { + const upgraded = new OrchestrationDb(dbPath) + expect(upgraded.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + // v1.4.198's insert shape: no home_run_id column, so the DEFAULT '' lands. + upgraded.db.exec(`INSERT INTO remote_dispatch_attachments + (dispatch_id, task_id, home_peer_fingerprint, runtime_epoch) + VALUES ('ctx_rolled_back', 'task_rolled_back', 'home', 'epoch')`) + expect(upgraded.getRemoteDispatchAttachment('ctx_rolled_back')?.home_run_id).toBe('') + expect(() => + importFederatedControlMessage(upgraded, { + dispatchId: 'ctx_rolled_back', + messageId: 'message_refused', + payload: JSON.stringify({ from: 'home', subject: 'x', body: '', type: 'status' }) + }) + ).toThrow(/Run not found/) + upgraded.close() + + const reopened = new OrchestrationDb(dbPath) + try { + const stubRunId = federatedStubHomeRunId('ctx_rolled_back') + expect(reopened.getRemoteDispatchAttachment('ctx_rolled_back')?.home_run_id).toBe(stubRunId) + expect(reopened.getRunRaw(stubRunId)).toMatchObject({ home_database: 'remote', legacy: 0 }) + importInstruction(reopened, 'ctx_rolled_back', 'message_rolled_back') + expect(reopened.getMessageById('message_rolled_back')?.run_id).toBe(stubRunId) + } finally { + reopened.close() + } + } finally { + rmSync(dir, { recursive: true, force: true }) + } }) }) diff --git a/src/main/runtime/orchestration/db/schema/migrate-v40.ts b/src/main/runtime/orchestration/db/schema/migrate-v40.ts index 50ef46f82cc..e3cd46ca30d 100644 --- a/src/main/runtime/orchestration/db/schema/migrate-v40.ts +++ b/src/main/runtime/orchestration/db/schema/migrate-v40.ts @@ -1,11 +1,15 @@ import type { OrchestrationDb } from '../orchestration-db' +import { backfillFederatedStubHomeRuns } from '../federation/federated-stub-home-run-backfill' export function migrateV40(this: OrchestrationDb, current: number): void { - if (current >= 40 || this.hasColumn('remote_dispatch_attachments', 'home_run_id')) { + if (current >= 40) { return } - // Federation is unreleased; any development-only rows fail Run validation until reattached. - this.db.exec( - "ALTER TABLE remote_dispatch_attachments ADD COLUMN home_run_id TEXT NOT NULL DEFAULT ''" - ) + if (!this.hasColumn('remote_dispatch_attachments', 'home_run_id')) { + this.db.exec( + "ALTER TABLE remote_dispatch_attachments ADD COLUMN home_run_id TEXT NOT NULL DEFAULT ''" + ) + } + // Why: workers attached by v1.4.198 keep a mailbox; without a Run their control mail is refused. + backfillFederatedStubHomeRuns(this.db) } diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts index 5e0f5f5b142..57ec354a8cc 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-authority.ts @@ -15,9 +15,9 @@ export function prepareStartingWorkerAuthority( effects: unknown[] setupState: string hostScope?: string | null - // 'created': this worker-start operation created the agent terminal (including agent-first - // worktree creation, whose effects receipt says 'reused_agent_terminal'). 'external': an - // explicit --terminal reuse; ownership transfers only from an exact owned settled resource. + // 'created': this worker-start operation created the agent terminal (agent-first worktree + // creation included; its pre-rename effects rows said 'reused_agent_terminal'). 'external': + // an explicit --terminal reuse; ownership transfers only from an exact owned settled resource. terminalOwnership?: 'created' | 'external' } ): string { diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts index 6544f519497..99873949f92 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-outcome.ts @@ -110,7 +110,8 @@ export function markWorkerStartUnknown( this: OrchestrationDb, dispatchId: string, stage: string, - reason: string + reason: string, + effects?: unknown[] ): WorkerDispatchRow { this.db.exec('BEGIN IMMEDIATE') try { @@ -124,7 +125,12 @@ export function markWorkerStartUnknown( id: dispatchId, from: 'starting', to: 'start_unknown', - projection: { stage, last_error: reason, updated_at: new Date().toISOString() } + projection: { + stage, + last_error: reason, + updated_at: new Date().toISOString(), + ...(effects ? { effects: JSON.stringify(effects) } : {}) + } }) transitionLifecycleWithDb(this.db, { entity: 'dispatch', @@ -138,7 +144,7 @@ export function markWorkerStartUnknown( from: 'dispatched', to: 'blocked' }) - this.closeQuestionsForDispatch(dispatchId) + // Authority survives uncertainty, so its outstanding questions must remain answerable. this.db.exec('COMMIT') return this.getWorkerDispatch(dispatchId) as WorkerDispatchRow } catch (error) { diff --git a/src/main/runtime/orchestration/db/writer-run-required.test.ts b/src/main/runtime/orchestration/db/writer-run-required.test.ts index 21fcbeb28d6..682fabb1b60 100644 --- a/src/main/runtime/orchestration/db/writer-run-required.test.ts +++ b/src/main/runtime/orchestration/db/writer-run-required.test.ts @@ -1,5 +1,6 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { OrchestrationDb } from './orchestration-db' +import { UNBOUND_RUN_ID } from './contract-constants' describe('writers require a Run', () => { let db: OrchestrationDb @@ -8,11 +9,11 @@ describe('writers require a Run', () => { }) afterEach(() => db.close()) - it('rejects a message without a Run instead of using the legacy Run', () => { - expect(() => db.insertMessage({ from: 'sender', to: 'worker', subject: 'mail' })).toThrow( - 'Run is required' - ) - expect(db.db.prepare('SELECT id FROM messages').all()).toEqual([]) + it('files a message without a Run under the unbound Run, never the legacy one', () => { + const message = db.insertMessage({ from: 'sender', to: 'worker', subject: 'mail' }) + expect(message.run_id).toBe(UNBOUND_RUN_ID) + expect(db.getRun(UNBOUND_RUN_ID)).toMatchObject({ legacy: 0 }) + expect(db.getUnreadMessages('worker').map((row) => row.id)).toEqual([message.id]) }) it('rejects a Task without a Run instead of using the legacy Run', () => { diff --git a/src/main/runtime/orchestration/orchestration-federated-legacy-probe.test.ts b/src/main/runtime/orchestration/orchestration-federated-legacy-probe.test.ts index ad08a7383f0..5eff7569d27 100644 --- a/src/main/runtime/orchestration/orchestration-federated-legacy-probe.test.ts +++ b/src/main/runtime/orchestration/orchestration-federated-legacy-probe.test.ts @@ -3,7 +3,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import { afterEach, describe, expect, it } from 'vitest' import { LEGACY_RUN_ID, OrchestrationDb } from './db' -import { SCHEMA_VERSION } from './db/contract-constants' +import { federatedStubHomeRunId, SCHEMA_VERSION, UNBOUND_RUN_ID } from './db/contract-constants' import { resolveOrchestrationMigrationStartVersion } from './orchestration-schema-version-skew' describe('federated mailbox legacy-adoption probe', () => { @@ -17,15 +17,22 @@ describe('federated mailbox legacy-adoption probe', () => { } }) - function seedMailbox(handle: string, kind: 'message' | 'delivery'): string { + function seedMailbox( + handle: string, + kind: 'message' | 'delivery', + homeRunId = 'run_home', + mailRunId = LEGACY_RUN_ID + ): string { directory = mkdtempSync(join(tmpdir(), 'orca-federated-legacy-probe-')) const path = join(directory, 'orchestration.db') db = new OrchestrationDb(path) - db.db.exec(` - INSERT INTO remote_dispatch_attachments ( - dispatch_id, task_id, home_peer_fingerprint, home_run_id, runtime_epoch, state - ) VALUES ('ctx_remote', 'task_remote', 'peer_home', 'run_home', 'epoch', 'ready'); - `) + db.db + .prepare( + `INSERT INTO remote_dispatch_attachments ( + dispatch_id, task_id, home_peer_fingerprint, home_run_id, runtime_epoch, state + ) VALUES ('ctx_remote', 'task_remote', 'peer_home', ?, 'epoch', 'ready')` + ) + .run(homeRunId) if (kind === 'message') { db.db .prepare( @@ -33,14 +40,14 @@ describe('federated mailbox legacy-adoption probe', () => { id, run_id, delivery_contract, from_handle, to_handle, subject, type ) VALUES ('msg_probe', ?, 'current_delivery', 'term_home', ?, 'continue', 'dispatch')` ) - .run(LEGACY_RUN_ID, handle) + .run(mailRunId, handle) } else { db.db .prepare( `INSERT INTO deliveries (id, run_id, mailbox_handle, consumer_generation, message_ids) VALUES ('delivery_probe', ?, ?, 0, '[]')` ) - .run(LEGACY_RUN_ID, handle) + .run(mailRunId, handle) } return path } @@ -70,6 +77,27 @@ describe('federated mailbox legacy-adoption probe', () => { } ) + it.each(['message', 'delivery'] as const)( + 'does not treat a stub-home-Run attachment %s as pre-Runs evidence', + (kind) => { + const stubRunId = federatedStubHomeRunId('ctx_remote') + const path = seedMailbox('dispatch:ctx_remote', kind, stubRunId, stubRunId) + db!.db + .prepare( + `INSERT INTO runs (id, objective, home_database, consumer_generation, legacy) + VALUES (?, 'Coordinated from peer_home', 'remote', 0, 0)` + ) + .run(stubRunId) + expect( + resolveOrchestrationMigrationStartVersion(db!.db, SCHEMA_VERSION, SCHEMA_VERSION) + ).toBe(SCHEMA_VERSION) + db!.close() + db = new OrchestrationDb(path) + expect(db.getLegacyAdoption()).toBeUndefined() + expect(db.getRemoteDispatchAttachment('ctx_remote')?.home_run_id).toBe(stubRunId) + } + ) + it.each(['message', 'delivery'] as const)( 'still replays adoption for a genuine legacy %s', (kind) => { @@ -94,4 +122,19 @@ describe('federated mailbox legacy-adoption probe', () => { } } ) + + it('keeps mail from a terminal in no Run across a reopen without replaying adoption', () => { + directory = mkdtempSync(join(tmpdir(), 'orca-unbound-mail-probe-')) + const path = join(directory, 'orchestration.db') + db = new OrchestrationDb(path) + const sent = db.insertMessage({ from: 'term_a', to: 'term_b', subject: 'hi' }) + expect(sent.run_id).toBe(UNBOUND_RUN_ID) + expect(resolveOrchestrationMigrationStartVersion(db.db, SCHEMA_VERSION, SCHEMA_VERSION)).toBe( + SCHEMA_VERSION + ) + db.close() + db = new OrchestrationDb(path) + expect(db.getLegacyAdoption()).toBeUndefined() + expect(db.getUnreadMessages('term_b').map((row) => row.id)).toEqual([sent.id]) + }) }) diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federation-effects.ts b/src/main/runtime/rpc/methods/orchestration/federation/federation-effects.ts index da0bcf92894..36d25552d80 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federation-effects.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federation-effects.ts @@ -29,7 +29,9 @@ export function appendFederationTerminalEffects( : terminal.handle === setupHandle ? 'setup' : 'configured_tab', - action: terminal.handle === agentHandle ? 'reused_agent_terminal' : 'created', + // Remote agent-first worktree creation made every listed terminal, the agent one included; + // the verb is a lifecycle fact, not a role marker. + action: 'created', id: terminal.handle, tabId: terminal.tabId, leafId: terminal.leafId diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.test.ts b/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.test.ts new file mode 100644 index 00000000000..f00e5102d6d --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.test.ts @@ -0,0 +1,41 @@ +import { describe, expect, it } from 'vitest' +import { FederationAttachStartParams } from './federation-start-schema' + +// The request shape a v1.4.198 coordinator sends: no runId field at all. +const legacyRequest = { + dispatchId: 'ctx_legacy', + taskId: 'task_legacy', + taskSpec: 'Do the thing', + protocolVersion: 3, + worktree: 'feature-branch' +} + +describe('FederationAttachStartParams', () => { + it('parses a v1.4.198 request that carries no runId', () => { + const result = FederationAttachStartParams.safeParse(legacyRequest) + expect(result.success, result.success ? undefined : JSON.stringify(result.error.issues)).toBe( + true + ) + expect(result.success && result.data.runId).toBeUndefined() + }) + + it('keeps a v1.4.199 runId verbatim', () => { + const result = FederationAttachStartParams.parse({ ...legacyRequest, runId: 'run_home' }) + expect(result.runId).toBe('run_home') + }) + + // Pins OptionalString: '' and non-strings drop to undefined (a stub Run is minted downstream); + // whitespace-only passes the schema and is refused by createRemoteDispatchAttachment. + it('maps an empty or non-string runId to undefined but passes whitespace through', () => { + expect(FederationAttachStartParams.parse({ ...legacyRequest, runId: '' }).runId).toBeUndefined() + expect(FederationAttachStartParams.parse({ ...legacyRequest, runId: 7 }).runId).toBeUndefined() + expect(FederationAttachStartParams.parse({ ...legacyRequest, runId: ' ' }).runId).toBe(' ') + }) + + it('still requires the dispatch, task, spec, and worktree fields', () => { + for (const field of ['dispatchId', 'taskId', 'taskSpec', 'worktree'] as const) { + const { [field]: _dropped, ...rest } = legacyRequest + expect(FederationAttachStartParams.safeParse(rest).success, field).toBe(false) + } + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.ts b/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.ts index 1e7257df27d..84ed57d58cc 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federation-start-schema.ts @@ -3,7 +3,8 @@ import { OptionalFiniteNumber, OptionalString, requiredString } from '../../../s import { OptionalWorkerLaunchPreference } from '../worker/worker-start-schema' export const FederationAttachStartParams = z.object({ - runId: requiredString('Missing Run ID'), + /** Omitted by v1.4.198 coordinators; the worker host then mints a stub home Run. */ + runId: OptionalString, dispatchId: requiredString('Missing Dispatch ID'), taskId: requiredString('Missing Task ID'), taskSpec: requiredString('Missing Task spec'), diff --git a/src/main/runtime/rpc/methods/orchestration/messaging/send-unbound-terminals.test.ts b/src/main/runtime/rpc/methods/orchestration/messaging/send-unbound-terminals.test.ts new file mode 100644 index 00000000000..84cb250a19e --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/messaging/send-unbound-terminals.test.ts @@ -0,0 +1,31 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { createOrchestrationRpcHarness } from '../rpc-test-harness' + +describe('orchestration.send between terminals in no Run', () => { + const h = createOrchestrationRpcHarness() + afterEach(() => h.cleanup()) + + it('delivers terminal-to-terminal mail when neither terminal is in a Run', async () => { + // Two plain panes and `send --to `: the first command the guide teaches. #19542 + // refused this with a bare "Run is required"; it files under the unbound Run instead. + const { db, runtime, ctx } = h.setup(false) + vi.spyOn(runtime, 'getTerminalPaneKey').mockImplementation((handle) => + handle === 'term_a' ? 'tab_a:leaf_a' : handle === 'term_b' ? 'tab_b:leaf_b' : null + ) + vi.spyOn(runtime, 'deliverPendingMessagesForHandle').mockImplementation(() => {}) + + const result = (await h.call( + 'orchestration.send', + { from: 'term_a', to: 'term_b', subject: 'hello from no Run' }, + ctx + )) as { message: { id: string; run_id: string; to_handle: string } } + + expect(result.message).toMatchObject({ run_id: 'run_unbound', to_handle: 'term_b' }) + expect(db.getRun('run_unbound')).toMatchObject({ legacy: 0 }) + expect(db.getUnreadMessages('term_b').map((row) => row.id)).toEqual([result.message.id]) + const checked = (await h.call('orchestration.check', { terminal: 'term_b' }, ctx)) as { + messages: { id: string }[] + } + expect(checked.messages.map((row) => row.id)).toEqual([result.message.id]) + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/composed-workers.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/composed-workers.test.ts index 32863c09371..aee45e25259 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/composed-workers.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/composed-workers.test.ts @@ -572,7 +572,7 @@ describe('orchestration RPC methods', () => { ) expect(result.effects).toEqual( expect.arrayContaining([ - expect.objectContaining({ role: 'agent', action: 'reused_agent_terminal' }), + expect.objectContaining({ role: 'agent', action: 'created' }), expect.objectContaining({ role: 'setup', action: 'created' }), expect.objectContaining({ role: 'configured_tab', action: 'created' }) ]) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/local-worker-start.ts b/src/main/runtime/rpc/methods/orchestration/worker/local-worker-start.ts index 71370a4b3e9..ec188695f5c 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/local-worker-start.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/local-worker-start.ts @@ -19,11 +19,11 @@ import { import { failWorkerStartWithReceipt } from './worker-start-receipt' import { parseTaskDeps } from './task-deps-argument' import { assertExplicitWorkerTerminalUsable } from './explicit-worker-terminal-validation' -import { deliverWorkerDispatchPreamble } from './deliver-worker-dispatch-preamble' import { recordCreatedWorkerTerminalCustody } from './created-worker-terminal-custody' import { tearDownFailedWorkerStart } from './failed-worker-start-teardown' -import { monitorWorkerSetup, requireWorkerAuthority, type WorkerEffect } from './worker-topology' +import { requireWorkerAuthority, type WorkerEffect } from './worker-topology' import { prepareLocalWorkerStart } from './worker-start-validation' +import { deliverAndSettleWorkerStartReadiness } from './worker-start-readiness-settlement' type WorkerStartMutation = { callerFingerprint: string @@ -204,50 +204,30 @@ export async function startLocalWorker(args: { terminalOwnership: params.terminal ? 'external' : 'created' }) - failedStage = 'dispatch_input' - const promptDelivery = await deliverWorkerDispatchPreamble({ + return await deliverAndSettleWorkerStartReadiness({ runtime, - structuredSession, - terminalHandle, + db, + run, + task, dispatchId: started.dispatch.id, dispatchDepth: started.dispatch.depth, - taskId: task.id, - taskSpec: task.spec, + structuredSession, + terminalHandle, coordinatorHandle: params.from, dispatchCapability: capability, devMode: params.devMode, - requestId: orchestrationMutation?.requestId ?? started.dispatch.id - }) - effects.push({ - kind: 'dispatch_input', - role: 'agent', - id: terminalHandle, - state: 'accepted' - }) - const worker = db.markWorkerDispatchReady(started.dispatch.id, effects) - monitorWorkerSetup({ - runtime, - db, - runId: run.id, - dispatchId: started.dispatch.id, + requestId: orchestrationMutation?.requestId ?? started.dispatch.id, + agent: agent ?? null, setupReceipt, - effects - }) - return { - runId: run.id, - taskId: task.id, - dispatchId: started.dispatch.id, - state: worker.state, - stage: worker.stage, - setup: setupReceipt, - launch: launch.receipt, + launchReceipt: launch.receipt, mode, timeoutMs: params.timeoutMs ?? 60_000, effects, - ...(promptDelivery ? { prompt: promptDelivery } : {}), - residualResources: [], - ...(placed.warning ? { warning: placed.warning } : {}) - } + terminalRevealWarning: placed.warning, + onStage: (stage) => { + failedStage = stage + } + }) } catch (error) { await tearDownFailedWorkerStart({ runtime, diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-setup-gate.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-setup-gate.ts index 98407cdb1b4..52a70433867 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-setup-gate.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-setup-gate.ts @@ -6,6 +6,8 @@ import { } from './worker-topology' function residualWorkerEffects(effects: WorkerEffect[]): WorkerEffect[] { + // 'reused_agent_terminal' is the retired verb agent-first creation used for its own agent + // terminal; rows persisted before the rename still carry it. return effects.filter( (effect) => effect.action?.startsWith('created') || effect.action === 'reused_agent_terminal' ) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts index dc645c97c88..593d097580f 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-prompt-contract.test.ts @@ -41,6 +41,8 @@ const openDatabases: OrchestrationDb[] = [] const temporaryRoots: string[] = [] type PromptContractHarness = { + runtime: Awaited>['runtime'] + handle: string db: OrchestrationDb dbPath: string dispatcher: RpcDispatcher @@ -131,6 +133,8 @@ async function createPromptContractHarness( vi.spyOn(runtime, 'getTerminalOrchestrationCliCommand').mockReturnValue('orca') return { + runtime, + handle, db, dbPath, dispatcher: new RpcDispatcher({ runtime, methods: ORCHESTRATION_METHODS }), @@ -227,22 +231,100 @@ describe('orchestration worker-start prompt contract', () => { }) }) - it('keeps a swallowed Enter queued without revoking the worker or retrying input', async () => { + it.each([ + ['succeeded', false], + ['succeeded', true], + ['failed', false], + ['failed', true] + ] as const)('preserves an early %s report with turn evidence=%s', async (outcome, observed) => { + vi.useFakeTimers() + const harness = await createPromptContractHarness('swallowed') + vi.spyOn(harness.runtime, 'observeTerminalAgentPrompt').mockImplementation( + async (_handle, prompt) => { + const dispatch = harness.db.findActiveDispatchForAssignee(harness.handle) + expect(dispatch).toBeDefined() + expect( + harness.db.settleWorkerReport({ + taskId: harness.taskId, + dispatchId: dispatch!.id, + outcome, + result: 'Finished before the hook arrived' + }) + ).toMatchObject({ action: 'settled', outcome }) + return observed ? { ...prompt, stages: ['input_accepted', 'turn_started'] } : prompt + } + ) + const pending = harness.dispatcher.dispatch(harness.request) + await vi.runAllTimersAsync() + expect(await pending).toMatchObject({ + ok: true, + result: { state: 'ready', stage: 'settled', workerOutcome: outcome } + }) + expect(harness.db.getTask(harness.taskId)?.status).toBe( + outcome === 'succeeded' ? 'completed' : 'failed' + ) + }) + + it('retains accepted authority when the observation binding becomes stale', async () => { + vi.useFakeTimers() + const harness = await createPromptContractHarness('swallowed') + vi.spyOn(harness.runtime, 'observeTerminalAgentPrompt').mockRejectedValue( + new Error('terminal_handle_stale') + ) + const pending = harness.dispatcher.dispatch(harness.request) + await vi.runAllTimersAsync() + expect(await pending).toMatchObject({ ok: true, result: { state: 'outcome_unknown' } }) + expect(harness.db.findActiveDispatchForAssignee(harness.handle)).toMatchObject({ + status: 'pending', + capability_hash: expect.any(String), + capability_revoked_at: null + }) + expect(harness.submittedTurns()).toBe(1) + expect(vi.getTimerCount()).toBe(0) + }) + + it('keeps a worker question answerable after turn observation times out', async () => { + vi.useFakeTimers() + const harness = await createPromptContractHarness('swallowed') + let questionId = '' + vi.spyOn(harness.runtime, 'observeTerminalAgentPrompt').mockImplementation( + async (_handle, prompt) => { + const dispatch = harness.db.findActiveDispatchForAssignee(harness.handle)! + questionId = harness.db.createQuestion({ + runId: dispatch.run_id!, + dispatchId: dispatch.id, + askerHandle: harness.handle, + question: 'Which target should I use?' + }).question.message_id + return prompt + } + ) + const pending = harness.dispatcher.dispatch(harness.request) + await vi.runAllTimersAsync() + expect(await pending).toMatchObject({ ok: true, result: { state: 'outcome_unknown' } }) + expect(harness.db.getQuestion(questionId)?.status).toBe('pending') + }) + + it('reports a swallowed Enter as start_unknown while keeping the worker and its capability', async () => { vi.useFakeTimers() const harness = await createPromptContractHarness('swallowed') const pending = harness.dispatcher.dispatch(harness.request) await vi.runAllTimersAsync() const response = await pending + // Codex supports turn-start observation and no turn started, so ready would be a lie: the + // paste can sit unsent in the composer while the receipt looks like a healthy dispatch. expect(response).toMatchObject({ ok: true, result: { - state: 'ready', - stage: 'input_accepted', + state: 'outcome_unknown', + stage: 'turn_start_unobserved', + turnStart: 'unobserved', prompt: { requestId: harness.requestId, stages: ['input_accepted'] }, + nextCommands: expect.arrayContaining([expect.stringContaining('worker-show')]), mutation: { requestId: harness.requestId, replayed: false } } }) @@ -251,29 +333,42 @@ describe('orchestration worker-start prompt contract', () => { } const dispatchId = (response.result as { dispatchId: string }).dispatchId await vi.advanceTimersByTimeAsync(20_000) + // Unverifiable is not failure: exactly one submit, no blind retry, nothing torn down. expect(harness.submittedTurns()).toBe(1) expect(harness.startedTurns()).toBe(0) expect(harness.prematureSubmits()).toBe(0) expect(harness.writes.filter((data) => data === '\r')).toHaveLength(1) const persisted = reopenPromptContractDb(harness) - expect(persisted.getTask(harness.taskId)?.status).toBe('dispatched') + expect(persisted.getTask(harness.taskId)?.status).toBe('blocked') expect(persisted.getDispatchContextById(dispatchId)).toMatchObject({ - status: 'dispatched', + status: 'pending', last_failure: null, + // The capability survives so a worker that recovers can still report; worker-report + // settlement reconnects a start_unknown worker through 'ready'. + capability_hash: expect.any(String), capability_revoked_at: null }) expect(persisted.getWorkerDispatch(dispatchId)).toMatchObject({ - state: 'ready', - stage: 'input_accepted', - last_error: null + state: 'start_unknown', + stage: 'turn_start_unobserved', + last_error: expect.stringContaining('turn start could not be verified') }) + const persistedEffects = JSON.parse( + persisted.getWorkerDispatch(dispatchId)?.effects ?? '[]' + ) as { kind?: string; state?: string }[] + expect(persistedEffects).toEqual( + expect.arrayContaining([ + expect.objectContaining({ kind: 'dispatch_input', state: 'accepted' }), + expect.objectContaining({ kind: 'dispatch_input', state: 'turn_unobserved' }) + ]) + ) const callerFingerprint = persisted.getOrCreateLocalMutationCallerFingerprint() const receipt = persisted.getMutationReceipt(callerFingerprint, harness.requestId) expect(receipt).toMatchObject({ state: 'completed' }) expect(JSON.parse(receipt?.receipt ?? 'null')).toMatchObject({ dispatchId, - state: 'ready', - stage: 'input_accepted', + state: 'outcome_unknown', + stage: 'turn_start_unobserved', prompt: { requestId: harness.requestId, stages: ['input_accepted'] diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts new file mode 100644 index 00000000000..a3c57a357d2 --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-readiness-settlement.ts @@ -0,0 +1,156 @@ +import type { OrcaRuntimeService } from '../../../../orca-runtime' +import type { OrchestrationDb } from '../../../../orchestration/db' +import type { RunRow, TaskRow } from '../../../../orchestration/types' +import type { WorkerStartModeReceipt } from '../../orchestration-worker-start-mode' +import { deliverWorkerDispatchPreamble } from './deliver-worker-dispatch-preamble' +import type { OrchestrationWorkerLaunchReceipt } from './worker-launch-preferences' +import { + describeUnobservedWorkerTurnStart, + observeWorkerTurnStart, + type WorkerTurnStartObservation +} from './worker-start-turn-observation' +import { + monitorWorkerSetup, + type createStructuredWorkerSessionForWorktree, + type WorkerEffect, + type WorkerSetupReceipt +} from './worker-topology' + +/** + * Delivers the dispatch preamble and settles the worker's start state on the strongest + * evidence available: `ready` only with a positive turn-start (or a provider that cannot + * prove one), `start_unknown` when observation is supported and nothing started. + */ +export async function deliverAndSettleWorkerStartReadiness(args: { + runtime: OrcaRuntimeService + db: OrchestrationDb + run: RunRow + task: TaskRow + dispatchId: string + dispatchDepth: number + structuredSession: Awaited> | null + terminalHandle: string + coordinatorHandle: string + dispatchCapability: string + devMode: boolean | undefined + requestId: string + agent: string | null + setupReceipt: WorkerSetupReceipt + launchReceipt: OrchestrationWorkerLaunchReceipt + mode: WorkerStartModeReceipt + timeoutMs: number + effects: WorkerEffect[] + terminalRevealWarning: string | undefined + /** Keeps the caller's failure receipt naming the stage that actually failed. */ + onStage: (stage: 'dispatch_input' | 'turn_observation') => void +}): Promise { + const { runtime, db, run, task, structuredSession, terminalHandle, effects } = args + + args.onStage('dispatch_input') + const promptDelivery = await deliverWorkerDispatchPreamble({ + runtime, + structuredSession, + terminalHandle, + dispatchId: args.dispatchId, + dispatchDepth: args.dispatchDepth, + taskId: task.id, + taskSpec: task.spec, + coordinatorHandle: args.coordinatorHandle, + dispatchCapability: args.dispatchCapability, + devMode: args.devMode, + requestId: args.requestId + }) + effects.push({ + kind: 'dispatch_input', + role: 'agent', + id: terminalHandle, + state: 'accepted' + }) + + args.onStage('turn_observation') + // The write above was accepted without waiting on provider hooks; now demand the positive + // evidence the receipt claims is observable. A worker whose turn never starts must not be + // reported ready — a wedged agent and a working one looked identical before this gate. + // A structured preamble send is acknowledged by the provider or throws, so it is already + // positive evidence. + const turnStart: WorkerTurnStartObservation = structuredSession + ? { verdict: 'observed' } + : await observeWorkerTurnStart({ runtime, terminalHandle, prompt: promptDelivery }) + const deliveredPrompt = turnStart.prompt ?? promptDelivery + monitorWorkerSetup({ + runtime, + db, + runId: run.id, + dispatchId: args.dispatchId, + setupReceipt: args.setupReceipt, + effects + }) + // A worker report can settle the dispatch while turn observation is outstanding. + const currentWorker = db.getWorkerDispatch(args.dispatchId) + const alreadySettled = currentWorker && currentWorker.state !== 'starting' + if (turnStart.verdict === 'unobserved' && !alreadySettled) { + // Honest `unverifiable`: keep the dispatch capability and the terminal — the worker may + // still recover and report (worker-report settlement reconnects a start_unknown worker) — + // but never claim ready for a turn nobody observed. + effects.push({ + kind: 'dispatch_input', + role: 'agent', + id: terminalHandle, + state: 'turn_unobserved' + }) + const reason = describeUnobservedWorkerTurnStart(args.agent) + const worker = db.markWorkerStartUnknown( + args.dispatchId, + 'turn_start_unobserved', + reason, + effects + ) + return { + runId: run.id, + taskId: task.id, + dispatchId: args.dispatchId, + state: 'outcome_unknown', + stage: worker.stage, + turnStart: turnStart.verdict, + lastError: reason, + setup: args.setupReceipt, + launch: args.launchReceipt, + mode: args.mode, + timeoutMs: args.timeoutMs, + effects, + ...(deliveredPrompt ? { prompt: deliveredPrompt } : {}), + residualResources: JSON.parse(worker.residual_resources) as unknown[], + nextCommands: [ + `orca orchestration worker-show --dispatch ${args.dispatchId} --json`, + `orca terminal read --terminal ${terminalHandle} --screen`, + `orca orchestration worker-abandon --dispatch ${args.dispatchId} --json` + ], + ...(args.terminalRevealWarning ? { warning: args.terminalRevealWarning } : {}) + } + } + const worker = alreadySettled + ? currentWorker + : db.markWorkerDispatchReady(args.dispatchId, effects) + // A completed task proves start succeeded; older callers use only 'ready' as start success. + const reportedOutcome = + worker.stage === 'settled' && (worker.state === 'succeeded' || worker.state === 'failed') + ? worker.state + : undefined + return { + runId: run.id, + taskId: task.id, + dispatchId: args.dispatchId, + state: reportedOutcome ? 'ready' : worker.state, + stage: worker.stage, + ...(reportedOutcome ? { workerOutcome: reportedOutcome } : {}), + turnStart: turnStart.verdict, + setup: args.setupReceipt, + launch: args.launchReceipt, + mode: args.mode, + timeoutMs: args.timeoutMs, + effects, + ...(deliveredPrompt ? { prompt: deliveredPrompt } : {}), + residualResources: [], + ...(args.terminalRevealWarning ? { warning: args.terminalRevealWarning } : {}) + } +} diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts new file mode 100644 index 00000000000..f1992eb9b9c --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.test.ts @@ -0,0 +1,120 @@ +import { describe, expect, it, vi } from 'vitest' +import type { RuntimeTerminalPromptDelivery } from '../../../../../../shared/runtime-terminal-contracts' +import type { OrcaRuntimeService } from '../../../../orca-runtime' +import { observeWorkerTurnStart } from './worker-start-turn-observation' + +function delivery( + overrides: Partial = {} +): RuntimeTerminalPromptDelivery { + return { + requestId: 'req-1', + stages: ['input_accepted'], + provider: 'codex', + observation: 'supported', + processIncarnation: 'inc-1', + generation: 1, + baselineWorkingSequence: 0, + ...overrides + } +} + +function runtimeObserving(result: RuntimeTerminalPromptDelivery): { + runtime: OrcaRuntimeService + observe: ReturnType +} { + const observe = vi.fn().mockResolvedValue(result) + return { + runtime: { observeTerminalAgentPrompt: observe } as unknown as OrcaRuntimeService, + observe + } +} + +describe('observeWorkerTurnStart', () => { + it('treats a missing receipt as unsupported observation, never as failure', async () => { + const { runtime, observe } = runtimeObserving(delivery()) + await expect( + observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt: undefined }) + ).resolves.toEqual({ verdict: 'unsupported' }) + expect(observe).not.toHaveBeenCalled() + }) + + it('accepts a first-stage turn_started without a second observation pass', async () => { + const prompt = delivery({ stages: ['input_accepted', 'turn_started'] }) + const { runtime, observe } = runtimeObserving(prompt) + await expect( + observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt }) + ).resolves.toEqual({ verdict: 'observed', prompt }) + expect(observe).not.toHaveBeenCalled() + }) + + it('reports observed when the second-stage observer sees the turn start', async () => { + const observed = delivery({ stages: ['input_accepted', 'turn_started'] }) + const { runtime, observe } = runtimeObserving(observed) + await expect( + observeWorkerTurnStart({ + runtime, + terminalHandle: 'term_w', + prompt: delivery(), + timeoutMs: 5 + }) + ).resolves.toEqual({ verdict: 'observed', prompt: observed }) + expect(observe).toHaveBeenCalledWith('term_w', delivery(), 5) + }) + + it('reports unobserved — not dead — when a supported observation stalls', async () => { + const { runtime } = runtimeObserving(delivery()) + await expect( + observeWorkerTurnStart({ + runtime, + terminalHandle: 'term_w', + prompt: delivery(), + timeoutMs: 5 + }) + ).resolves.toMatchObject({ verdict: 'unobserved' }) + }) + + it('reports permission as positive liveness', async () => { + const observed = delivery({ observation: 'permission' }) + const { runtime } = runtimeObserving(observed) + await expect( + observeWorkerTurnStart({ + runtime, + terminalHandle: 'term_w', + prompt: delivery(), + timeoutMs: 5 + }) + ).resolves.toEqual({ verdict: 'permission', prompt: observed }) + }) + + it('treats a replaced incarnation as unobserved rather than unsupported', async () => { + const observed = delivery({ observation: 'incarnation_replaced' }) + const { runtime } = runtimeObserving(observed) + await expect( + observeWorkerTurnStart({ + runtime, + terminalHandle: 'term_w', + prompt: delivery(), + timeoutMs: 5 + }) + ).resolves.toEqual({ verdict: 'unobserved', prompt: observed }) + }) + + it('leaves an unsupported provider on the accepted receipt', async () => { + const prompt = delivery({ provider: 'unsupported', observation: 'unsupported' }) + const { runtime, observe } = runtimeObserving(prompt) + await expect( + observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt }) + ).resolves.toEqual({ verdict: 'unsupported', prompt }) + expect(observe).not.toHaveBeenCalled() + }) + + it('preserves uncertainty when observation loses the terminal binding', async () => { + const prompt = delivery() + const { runtime, observe } = runtimeObserving(prompt) + observe.mockRejectedValue(new Error('terminal_handle_stale')) + await expect( + observeWorkerTurnStart({ runtime, terminalHandle: 'term_w', prompt }) + ).resolves.toEqual({ verdict: 'unobserved', prompt }) + expect(observe).toHaveBeenCalledTimes(1) + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts new file mode 100644 index 00000000000..d1c9e59adfa --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-start-turn-observation.ts @@ -0,0 +1,86 @@ +import { AGENT_PROMPT_EFFECT_TIMEOUT_MS } from '../../../../../../shared/orchestration-timing-budgets' +import type { RuntimeTerminalPromptDelivery } from '../../../../../../shared/runtime-terminal-contracts' +import type { OrcaRuntimeService } from '../../../../orca-runtime' + +/** + * Turn-start verdict for a dispatched worker prompt, in the execution-boundary vocabulary: + * + * - 'observed': the provider proved a turn started for this request. Positive liveness. + * - 'permission': the agent rendered an approval prompt after the write. Positive liveness, + * but the turn is blocked on a human. + * - 'unsupported': this provider exposes no turn-start signal; the accepted write is the + * strongest receipt that can exist. Never treated as failure. + * - 'unobserved': observation IS supported and no turn started within the window. This is + * `unverifiable`, never evidence of death — the bytes were written, but the agent may be + * wedged at startup or holding the spec unsent in its composer. + */ +export type WorkerTurnStartVerdict = 'observed' | 'permission' | 'unsupported' | 'unobserved' + +export type WorkerTurnStartObservation = { + verdict: WorkerTurnStartVerdict + prompt?: RuntimeTerminalPromptDelivery +} + +function classifyPromptDelivery(prompt: RuntimeTerminalPromptDelivery): WorkerTurnStartVerdict { + if (prompt.stages.includes('turn_started')) { + return 'observed' + } + if (prompt.observation === 'permission') { + return 'permission' + } + if (prompt.observation === 'supported') { + return 'unobserved' + } + // 'unsupported' (and an old host's missing observation) leaves acceptance as the best receipt. + return 'unsupported' +} + +/** + * Second-stage turn-start observation for a worker prompt that was accepted without waiting on + * provider hooks. Reuses the same observer that terminal.send receipts replay through, so the + * evidence rules (lifecycle edge or hook turn-start, never output bytes) stay in one place. + * + * The observation window is `AGENT_PROMPT_EFFECT_TIMEOUT_MS`, which worker-start's client RPC + * grace already budgets for (see orchestration-worker-start-prompt-budget.ts). + */ +export async function observeWorkerTurnStart(args: { + runtime: OrcaRuntimeService + terminalHandle: string + prompt: RuntimeTerminalPromptDelivery | undefined + timeoutMs?: number +}): Promise { + if (!args.prompt) { + return { verdict: 'unsupported' } + } + const verdict = classifyPromptDelivery(args.prompt) + if (verdict !== 'unobserved') { + return { verdict, prompt: args.prompt } + } + let observed: RuntimeTerminalPromptDelivery + try { + observed = await args.runtime.observeTerminalAgentPrompt( + args.terminalHandle, + args.prompt, + args.timeoutMs ?? AGENT_PROMPT_EFFECT_TIMEOUT_MS + ) + } catch { + // Observation failure cannot revoke authority for input that was already accepted. + return { verdict: 'unobserved', prompt: args.prompt } + } + if (observed.observation === 'incarnation_replaced') { + // The PTY under this handle changed mid-observation; the accepted write is unproven. + return { verdict: 'unobserved', prompt: observed } + } + return { verdict: classifyPromptDelivery(observed), prompt: observed } +} + +export function describeUnobservedWorkerTurnStart(agent: string | null): string { + const name = agent ?? 'the agent' + return ( + `Dispatch input was written and submitted, but ${name}'s turn start could not be verified ` + + `during observation (up to ${Math.round(AGENT_PROMPT_EFFECT_TIMEOUT_MS / 1000)}s). This is unverifiable, not proof the ` + + 'worker is dead: the agent may still be starting, may be wedged (for example waiting on ' + + 'network), or may be holding the task unsent in its composer. If the worker recovers and ' + + 'reports, this Dispatch settles normally.' + ) +} diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-worktree-creation.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-worktree-creation.ts index 0f9435bb018..fe0c11e374d 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-worktree-creation.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-worktree-creation.ts @@ -108,7 +108,10 @@ export async function createWorkerWorktree(args: { : terminal.handle === setupTerminalHandle ? 'setup' : 'configured_tab', - action: terminal.handle === terminalHandle ? 'reused_agent_terminal' : 'created', + // Every terminal listed here — the agent terminal included — was created by this call's + // agent-first worktree creation. The old 'reused_agent_terminal' verb on the agent row + // conflated the role test with a lifecycle claim and misdiagnosed at least one incident. + action: 'created', id: terminal.handle, tabId: terminal.tabId, leafId: terminal.leafId diff --git a/src/main/runtime/rpc/methods/orchestration/worker/workers-new-worktree.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/workers-new-worktree.test.ts index a0d74031b72..5164ac0b22f 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/workers-new-worktree.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/workers-new-worktree.test.ts @@ -157,7 +157,7 @@ describe('orchestration new-worktree workers', () => { expect.objectContaining({ kind: 'terminal', role: 'agent', - action: 'reused_agent_terminal', + action: 'created', id: 'term_worker' }) ]) diff --git a/src/renderer/src/components/native-chat/NativeChatBackgroundTasksStatus.tsx b/src/renderer/src/components/native-chat/NativeChatBackgroundTasksStatus.tsx index 9f8f08e1efd..ce3718a967b 100644 --- a/src/renderer/src/components/native-chat/NativeChatBackgroundTasksStatus.tsx +++ b/src/renderer/src/components/native-chat/NativeChatBackgroundTasksStatus.tsx @@ -49,7 +49,7 @@ export function NativeChatBackgroundTasksStatus(props: { - + {translate( 'components.native-chat.backgroundTasks.monitoring', 'Monitoring background tasks' diff --git a/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx b/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx new file mode 100644 index 00000000000..e34c1afd131 --- /dev/null +++ b/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx @@ -0,0 +1,225 @@ +// @vitest-environment happy-dom + +import '@testing-library/jest-dom/vitest' + +import { cleanup, fireEvent, render, screen, within } from '@testing-library/react' +import { afterEach, describe, expect, it, vi } from 'vitest' +import type { AgentJournalStatusItem } from '../../../../shared/agent-session-journal-types' +import { projectStructuredItemToNativeChat } from '../../../../shared/structured-agent-session-projection' +import type { NativeChatMessage } from '../../../../shared/native-chat-types' +import { NativeChatMessageList } from './NativeChatMessageList' +import { projectNativeChatTaskListFrames } from './native-chat-task-list-frames' + +afterEach(cleanup) + +function frame(id: number, status: string, overrides: { kind?: string; truncated?: boolean } = {}) { + const kind = overrides.kind ?? 'notification:turn/plan/updated' + const head = JSON.stringify({ + threadId: 'thread', + turnId: 'turn', + explanation: 'Keep verification visible', + plan: [{ step: 'Verify', status }] + }) + const body: AgentJournalStatusItem = { + kind: 'status', + text: `codex · ${kind}`, + providerFrame: { + provider: 'codex', + kind, + payload: { + head, + byteLength: new TextEncoder().encode(head).byteLength, + digest: 'fixture-digest', + truncated: overrides.truncated ?? false + } + } + } + const message = projectStructuredItemToNativeChat({ + itemId: `frame-${id}`, + revision: 1, + sequence: id, + observedAt: id, + body + }) + if (!message) { + throw new Error('Expected a projected message') + } + return message +} + +function transcript(messages: NativeChatMessage[], sessionId = 'live-codex') { + return ( + + ) +} + +describe('live Codex checklist frames', () => { + it('updates one pinned checklist from journal notifications without rewinding on pagination', () => { + const first = frame(1, 'pending') + const active = frame(2, 'inProgress') + const last = frame(3, 'completed') + const { rerender, container } = render(transcript([first])) + const toggle = screen.getByRole('button', { name: 'Tasks 0 of 1 tasks completed' }) + const viewport = container.querySelector('.overflow-y-auto')! + expect(viewport.contains(toggle)).toBe(false) + fireEvent.click(toggle) + rerender(transcript([first, active])) + expect(within(toggle.parentElement!).getByText('Verify').closest('li')).toHaveClass( + 'text-foreground' + ) + expect(within(viewport as HTMLElement).getByText('Started Verify')).toBeInTheDocument() + rerender(transcript([last])) + expect( + within( + screen.getByRole('button', { name: 'Tasks 1 of 1 tasks completed' }).parentElement! + ).getByText('Verify') + ).toHaveClass('line-through') + expect(screen.getAllByText('Keep verification visible')).toHaveLength(2) + rerender(transcript([first, active, last])) + expect(screen.getAllByText('Verify')).toHaveLength(2) + expect(screen.getByRole('button', { name: 'Tasks 1 of 1 tasks completed' })).toBe(toggle) + expect(screen.getByText('Started Verify')).toBeInTheDocument() + expect(screen.queryByText('notification:turn/plan/updated')).toBeNull() + expect(projectNativeChatTaskListFrames([last])[0]).toBe( + projectNativeChatTaskListFrames([last])[0] + ) + }) + + it('uses the latest complete snapshot across Codex tool calls and notifications', () => { + const tool: NativeChatMessage = { + id: 'tool', + role: 'assistant', + timestamp: 1, + source: 'transcript', + blocks: [ + { + type: 'tool-call', + name: 'update_plan', + input: { plan: [{ step: 'Verify', status: 'pending' }] } + } + ] + } + render(transcript([tool, frame(3, 'completed')])) + fireEvent.click(screen.getByRole('button', { name: 'Tasks 1 of 1 tasks completed' })) + expect(screen.getAllByText('Verify')).toHaveLength(2) + expect( + within( + screen.getByRole('button', { name: 'Tasks 1 of 1 tasks completed' }).parentElement! + ).getByText('Verify') + ).toHaveClass('line-through') + }) + + it('keeps malformed, truncated, other-provider, and plan-document frames unchanged', () => { + const truncated = frame(1, 'pending', { truncated: true }) + const document = frame(2, 'pending', { kind: 'item:plan' }) + const malformed = frame(3, 'pending') + const otherProvider = frame(4, 'pending') + const malformedBlock = malformed.blocks[0] + const otherBlock = otherProvider.blocks[0] + if (malformedBlock.type === 'text' && malformedBlock.providerFrame) { + malformedBlock.providerFrame.payload.head = '{"plan":null}' + } + if (otherBlock.type === 'text' && otherBlock.providerFrame) { + otherBlock.providerFrame.provider = 'claude' + } + const messages = [truncated, document, malformed, otherProvider] + const projected = projectNativeChatTaskListFrames(messages) + projected.forEach((message, index) => expect(message).toBe(messages[index])) + render(transcript([truncated])) + expect(screen.getByText('notification:turn/plan/updated')).toBeInTheDocument() + expect(screen.queryByText('Tasks')).toBeNull() + }) + + it('does not consume a neighboring tool failure as a notification result', () => { + const command: NativeChatMessage = { + id: 'command', + role: 'assistant', + timestamp: 2, + source: 'transcript', + blocks: [ + { type: 'tool-call', name: 'shell', input: { command: 'verify' }, state: 'failed' }, + { type: 'tool-result', output: 'Verification failed', isError: true } + ] + } + render(transcript([frame(1, 'pending'), command])) + expect(screen.getByRole('button', { name: 'Tasks 0 of 1 tasks completed' })).toBeInTheDocument() + expect(screen.getByText('Verification failed', { selector: 'pre' })).toHaveClass( + 'text-destructive' + ) + }) +}) + +describe('NativeChatMessageList task list history', () => { + it('keeps the latest Claude state after pagination and resets disclosure between sessions', () => { + const first = { + id: 'first-list', + role: 'assistant' as const, + timestamp: 1, + source: 'transcript' as const, + blocks: [ + { + type: 'tool-call' as const, + name: 'TodoWrite', + input: { + todos: [ + { content: 'Read', status: 'pending' }, + { content: 'Test', status: 'pending' } + ] + } + } + ] + } + const last = { + ...first, + id: 'last-list', + timestamp: 3, + blocks: [ + { type: 'text' as const, text: 'Ready for verification' }, + { + type: 'tool-call' as const, + name: 'TodoWrite', + input: { + todos: [ + { content: 'Read', status: 'completed' }, + { content: 'Test', status: 'pending' } + ] + } + } + ] + } + const { rerender } = render(transcript([last])) + fireEvent.click(screen.getByRole('button', { name: 'Tasks 1 of 2 tasks completed' })) + rerender(transcript([first, last])) + expect(screen.getAllByText('Read')).toHaveLength(2) + expect(screen.getByText('Completed Read')).toBeInTheDocument() + expect( + within( + screen.getByRole('button', { name: 'Tasks 1 of 2 tasks completed' }).parentElement! + ).getByText('Read') + ).toHaveClass('line-through') + expect(screen.getByText('Ready for verification')).toBeInTheDocument() + rerender(transcript([first], 'two')) + expect(screen.getByRole('button', { name: 'Tasks 0 of 2 tasks completed' })).toHaveAttribute( + 'aria-expanded', + 'false' + ) + expect(screen.getAllByText('Read')).toHaveLength(1) + rerender(transcript([], 'three')) + expect(screen.queryByText('Tasks')).toBeNull() + }) +}) diff --git a/src/renderer/src/components/native-chat/NativeChatMessageList.tsx b/src/renderer/src/components/native-chat/NativeChatMessageList.tsx index d8cee7f7e7c..c6da9030a2c 100644 --- a/src/renderer/src/components/native-chat/NativeChatMessageList.tsx +++ b/src/renderer/src/components/native-chat/NativeChatMessageList.tsx @@ -5,6 +5,10 @@ import { translate } from '@/i18n/i18n' import type { NativeChatLiveSession } from './use-native-chat-live-session' import { createNativeChatMessageListProjection } from './native-chat-message-list-projection' import { isNearBottom, shouldShowJumpToLatest, type ScrollGeometry } from './native-chat-autoscroll' +import { nativeChatTaskListState } from './native-chat-task-list-state' +import { nativeChatTaskListPredecessors } from './native-chat-task-list-history' +import { NativeChatTaskList } from './NativeChatTaskList' +import { projectNativeChatTaskListFrames } from './native-chat-task-list-frames' import { MessageRow } from './NativeChatMessageRow' import { shouldShowNativeChatTypingIndicator } from './native-chat-typing-indicator' import { NativeChatWorkingStatus } from './NativeChatWorkingStatus' @@ -112,9 +116,11 @@ export function NativeChatMessageList({ [session.agent, session.sessionId] ) const messages = useMemo( - () => projectMessages(session.messages), + () => projectNativeChatTaskListFrames(projectMessages(session.messages)), [projectMessages, session.messages] ) + const taskListPredecessors = useMemo(() => nativeChatTaskListPredecessors(messages), [messages]) + const taskListState = useMemo(() => nativeChatTaskListState(messages), [messages]) const showTypingIndicator = showTurnStatus ? isWorking : shouldShowNativeChatTypingIndicator({ messages, isWorking }) @@ -224,123 +230,140 @@ export function NativeChatMessageList({ }, [handleScroll, scrollToBottom]) return ( -
-
+
+
- {hasMore ? ( -
- -
- ) : null} - {messages.map((message, index) => { - const turnKey = turnKeys[index] - const isCurrentTurn = currentTurnKey - ? turnKey === currentTurnKey - : turnKey === undefined - const status = - index === latestUserIndex - ? turnStatuses.active - : message.role === 'user' && turnKey - ? turnStatuses.completedByTurn[turnKey] - : undefined - const receipt = receipts.get(message.id) - const turnDiff = - turnKey && turnKeys[index + 1] !== turnKey ? turnDiffs.get(turnKey) : undefined - return ( - - {receipt ? ( - - ) : ( - - )} - {showTurnStatus && - status && - (index !== latestUserIndex || showTypingIndicator || !isWorking) ? ( - toggleExpandedTurn(turnKey) - : undefined - } - /> - ) : null} - {turnDiff ? ( - - ) : null} - - ) - })} - {showTurnStatus && - latestUserIndex === -1 && - turnStatuses.active && - showTypingIndicator ? ( - - ) : null} - {showTurnStatus && isWorking ? ( - - ) : null} - {!showTurnStatus && showTypingIndicator ? : null} +
+ {hasMore ? ( +
+ +
+ ) : null} + {messages.map((message, index) => { + const turnKey = turnKeys[index] + const isCurrentTurn = currentTurnKey + ? turnKey === currentTurnKey + : turnKey === undefined + const status = + index === latestUserIndex + ? turnStatuses.active + : message.role === 'user' && turnKey + ? turnStatuses.completedByTurn[turnKey] + : undefined + const receipt = receipts.get(message.id) + const turnDiff = + turnKey && turnKeys[index + 1] !== turnKey ? turnDiffs.get(turnKey) : undefined + return ( + + {receipt ? ( + + ) : ( + + )} + {showTurnStatus && + status && + (index !== latestUserIndex || showTypingIndicator || !isWorking) ? ( + toggleExpandedTurn(turnKey) + : undefined + } + /> + ) : null} + {turnDiff ? ( + + ) : null} + + ) + })} + {showTurnStatus && + latestUserIndex === -1 && + turnStatuses.active && + showTypingIndicator ? ( + + ) : null} + {showTurnStatus && isWorking ? ( + + ) : null} + {!showTurnStatus && showTypingIndicator ? : null} +
+ {showJump ? ( + + ) : null}
- {showJump ? ( - + {taskListState.list && taskListState.list.tasks.length > 0 ? ( +
+
+ +
+
) : null}
) diff --git a/src/renderer/src/components/native-chat/NativeChatMessageRow.tsx b/src/renderer/src/components/native-chat/NativeChatMessageRow.tsx index e3a4e78efb4..773e6f2c2ff 100644 --- a/src/renderer/src/components/native-chat/NativeChatMessageRow.tsx +++ b/src/renderer/src/components/native-chat/NativeChatMessageRow.tsx @@ -8,7 +8,11 @@ import { isSubagentGroupFallbackText, subagentGroupBlocks } from '../../../../shared/native-chat-subagent-summary' -import { isSubagentGroupBlock, type NativeChatMessage } from '../../../../shared/native-chat-types' +import { + isSubagentGroupBlock, + type NativeChatMessage, + type NativeChatToolCallBlock +} from '../../../../shared/native-chat-types' import { splitNativeChatBlocks } from './native-chat-tool-fold' import { NativeChatToolRun } from './NativeChatToolRun' import { NativeChatNoticeRow } from './NativeChatNoticeRow' @@ -29,6 +33,8 @@ import type { RuntimeFileOperationArgs } from '@/runtime/runtime-file-client' * keep their block identity, so only the changed row re-renders. */ export const MessageRow = memo(function MessageRow({ message, + previousTodoWrite, + previousUpdatePlan, revealedDiff, expandSignal, activeTurnIsWorking, @@ -41,6 +47,8 @@ export const MessageRow = memo(function MessageRow({ runtimeContext }: { message: NativeChatMessage + previousTodoWrite?: NativeChatToolCallBlock + previousUpdatePlan?: NativeChatToolCallBlock revealedDiff?: NativeChatDiffReveal expandSignal: boolean activeTurnIsWorking?: boolean @@ -202,6 +210,8 @@ export const MessageRow = memo(function MessageRow({ {tools.length > 0 || subagentGroups.length > 0 ? ( - - {headline} - + {headline} {verdict} {alert === null ? null : ` +${alert}`} diff --git a/src/renderer/src/components/native-chat/NativeChatTaskList.test.tsx b/src/renderer/src/components/native-chat/NativeChatTaskList.test.tsx new file mode 100644 index 00000000000..588f392f713 --- /dev/null +++ b/src/renderer/src/components/native-chat/NativeChatTaskList.test.tsx @@ -0,0 +1,75 @@ +// @vitest-environment happy-dom +import '@testing-library/jest-dom/vitest' +import { cleanup, fireEvent, render, screen } from '@testing-library/react' +import { afterEach, describe, expect, it } from 'vitest' +import { NativeChatTaskList } from './NativeChatTaskList' +import type { NativeChatTaskList as TaskList } from '../../../../shared/native-chat-task-list' + +afterEach(cleanup) +const previous: TaskList = { + tasks: [ + { content: 'Read', status: 'in_progress', activeForm: 'Reading' }, + { content: 'Write', status: 'pending', activeForm: 'Writing' }, + { content: 'Test', status: 'pending' } + ] +} +const current: TaskList = { + tasks: [ + { content: 'Read', status: 'completed', activeForm: 'Reading' }, + { content: 'Write', status: 'in_progress', activeForm: 'Writing' }, + { content: 'Test', status: 'pending' } + ] +} + +describe('NativeChatTaskList', () => { + it('shows tri-state glyphs, progress, and activeForm in the first checklist', () => { + const { container } = render() + expect(screen.getByText('Read')).toHaveClass('line-through') + expect(screen.getByText('Writing').closest('li')).toHaveClass('text-foreground') + expect(screen.getByText('Test')).toBeInTheDocument() + expect(screen.getByLabelText('1 of 3 tasks completed')).toHaveTextContent('1/3') + for (const glyph of ['circle', 'circle-dot', 'circle-check']) { + expect(container.querySelector(`.lucide-${glyph}`)).not.toBeNull() + } + expect(screen.getByText('In progress:')).toHaveClass('sr-only') + }) + + it('leads with the diff and expands the complete checklist on demand', () => { + render() + expect(screen.getByText('Completed Read')).toBeInTheDocument() + expect(screen.getByText('Started Write')).toBeInTheDocument() + expect(screen.queryByText('Test')).toBeNull() + const disclosure = screen.getByRole('button', { name: 'Full task list' }) + expect(disclosure).toHaveAttribute('aria-expanded', 'false') + fireEvent.click(disclosure) + expect(disclosure).toHaveAttribute('aria-expanded', 'true') + expect(screen.getByText('Writing')).toBeInTheDocument() + expect(screen.getByText('Test')).toBeInTheDocument() + }) + + it('shows unchanged feedback and the current explanation', () => { + render( + + ) + expect(screen.getByText('Tasks unchanged')).toBeInTheDocument() + expect(screen.getByText('Continuing verification')).toBeInTheDocument() + expect(screen.queryByText('Test')).toBeNull() + }) + + it('renders empty lists without claiming any task completed', () => { + render() + expect(screen.getByText('No tasks')).toBeInTheDocument() + expect(screen.getByLabelText('0 of 0 tasks completed')).toHaveTextContent('0/0') + }) + + it('switches from full list to diff when earlier history supplies a predecessor', () => { + const { rerender } = render() + expect(screen.getByText('Test')).toBeInTheDocument() + rerender() + expect(screen.queryByText('Test')).toBeNull() + expect(screen.getByText('Started Write')).toBeInTheDocument() + }) +}) diff --git a/src/renderer/src/components/native-chat/NativeChatTaskList.tsx b/src/renderer/src/components/native-chat/NativeChatTaskList.tsx new file mode 100644 index 00000000000..3ddaa8959e8 --- /dev/null +++ b/src/renderer/src/components/native-chat/NativeChatTaskList.tsx @@ -0,0 +1,183 @@ +import { Circle, CircleCheck, CircleDot, ChevronRight, ListChecks } from 'lucide-react' +import { Collapsible, CollapsibleContent, CollapsibleTrigger } from '@/components/ui/collapsible' +import { cn } from '@/lib/utils' +import { translate } from '@/i18n/i18n' +import { + diffNativeChatTaskLists, + nativeChatTaskLabel, + type NativeChatTask, + type NativeChatTaskChange, + type NativeChatTaskList as TaskList +} from '../../../../shared/native-chat-task-list' + +function statusLabel(task: NativeChatTask): string { + if (task.status === 'completed') { + return translate('components.native-chat.taskList.completed', 'Completed') + } + if (task.status === 'in_progress') { + return translate('components.native-chat.taskList.inProgress', 'In progress') + } + return translate('components.native-chat.taskList.pending', 'Pending') +} + +function changeLabel(change: NativeChatTaskChange): string { + const values = { task: change.task.content } + switch (change.kind) { + case 'added': + return translate('components.native-chat.taskList.added', 'Added {{task}}', values) + case 'removed': + return translate('components.native-chat.taskList.removed', 'Removed {{task}}', values) + case 'started': + return translate('components.native-chat.taskList.started', 'Started {{task}}', values) + case 'completed': + return translate('components.native-chat.taskList.finished', 'Completed {{task}}', values) + case 'pending': + return translate('components.native-chat.taskList.reset', 'Marked pending: {{task}}', values) + case 'updated': + return translate('components.native-chat.taskList.updated', 'Updated {{task}}', { + task: nativeChatTaskLabel(change.task) + }) + } +} + +function TaskRow({ task, label }: { task: NativeChatTask; label?: string }): React.JSX.Element { + const Icon = + task.status === 'completed' ? CircleCheck : task.status === 'in_progress' ? CircleDot : Circle + return ( +
  • + + {statusLabel(task)}: + + {label ?? nativeChatTaskLabel(task)} + +
  • + ) +} + +function Checklist({ list }: { list: TaskList }): React.JSX.Element { + return list.tasks.length === 0 ? ( +

    + {translate('components.native-chat.taskList.empty', 'No tasks')} +

    + ) : ( +
      + {list.tasks.map((task, index) => ( + + ))} +
    + ) +} + +export function NativeChatTaskList({ + list, + previous, + presentation = 'inline' +}: { + list: TaskList + previous?: TaskList + presentation?: 'inline' | 'composer' +}): React.JSX.Element { + const completed = list.tasks.filter((task) => task.status === 'completed').length + if (presentation === 'composer') { + return ( + + + + + {translate('components.native-chat.taskList.title', 'Tasks')} + + + {completed}/{list.tasks.length} + + + + +
    + + {list.explanation ? ( +

    + {list.explanation} +

    + ) : null} +
    +
    +
    + ) + } + const changes = previous ? diffNativeChatTaskLists(previous, list) : null + return ( +
    +
    + + + {translate('components.native-chat.taskList.title', 'Tasks')} + + + {completed}/{list.tasks.length} + +
    + {changes ? ( + <> + {changes.length > 0 ? ( +
      + {changes.map((change, index) => ( + + ))} +
    + ) : ( +

    + {translate('components.native-chat.taskList.unchanged', 'Tasks unchanged')} +

    + )} + + + + {translate('components.native-chat.taskList.showAll', 'Full task list')} + + + + + + + ) : ( + + )} + {list.explanation ? ( +

    + {list.explanation} +

    + ) : null} +
    + ) +} diff --git a/src/renderer/src/components/native-chat/NativeChatToolAnnotations.tsx b/src/renderer/src/components/native-chat/NativeChatToolAnnotations.tsx index 7da3345fae4..25c0cd14d89 100644 --- a/src/renderer/src/components/native-chat/NativeChatToolAnnotations.tsx +++ b/src/renderer/src/components/native-chat/NativeChatToolAnnotations.tsx @@ -41,7 +41,7 @@ export function NativeChatCommandMetadata({ return null } return ( - + {exitCode !== undefined ? ( {translate('components.native-chat.tool.exitCode', 'exit {{value0}}', { diff --git a/src/renderer/src/components/native-chat/NativeChatToolRun.test.tsx b/src/renderer/src/components/native-chat/NativeChatToolRun.test.tsx index cd819c5ced1..f38f63df4c6 100644 --- a/src/renderer/src/components/native-chat/NativeChatToolRun.test.tsx +++ b/src/renderer/src/components/native-chat/NativeChatToolRun.test.tsx @@ -730,3 +730,58 @@ describe('NativeChatToolRun', () => { expect(screen.getByTitle('ls')).toHaveTextContent('ls') }) }) + +describe('NativeChatToolRun task lists', () => { + it('renders task updates instead of JSON and consumes successful results', () => { + const blocks: NativeChatBlock[] = [ + { + type: 'tool-call', + name: 'update_plan', + input: { + plan: [ + { step: 'Read', status: 'in_progress' }, + { step: 'Test', status: 'pending' } + ] + } + }, + { type: 'tool-result', output: 'Plan updated' }, + { + type: 'tool-call', + name: 'update_plan', + input: { + plan: [ + { step: 'Read', status: 'completed' }, + { step: 'Test', status: 'in_progress' } + ] + } + } + ] + const { container } = render() + expect(screen.getByText('Completed Read')).toBeInTheDocument() + expect(screen.getByText('Started Test')).toBeInTheDocument() + expect(screen.getByText('1/2')).toBeInTheDocument() + expect(screen.queryByText('Plan updated')).toBeNull() + expect(container.querySelector('pre')).toBeNull() + }) + + it('keeps malformed calls and failed results visible in the generic view', () => { + render( + + ) + expect(screen.getByText('Invalid arguments', { selector: 'pre' })).toBeInTheDocument() + expect(screen.getByText('Update rejected', { selector: 'pre' })).toBeInTheDocument() + expect(screen.queryByText('1/1')).toBeNull() + }) +}) diff --git a/src/renderer/src/components/native-chat/NativeChatToolRun.tsx b/src/renderer/src/components/native-chat/NativeChatToolRun.tsx index 5f579a2778b..6feb44bd703 100644 --- a/src/renderer/src/components/native-chat/NativeChatToolRun.tsx +++ b/src/renderer/src/components/native-chat/NativeChatToolRun.tsx @@ -12,7 +12,8 @@ import { isToolCallBlock, isToolResultBlock, type NativeChatBlock, - type NativeChatSubagentGroupBlock + type NativeChatSubagentGroupBlock, + type NativeChatToolCallBlock } from '../../../../shared/native-chat-types' import { isRenderableSubagentGroup } from '../../../../shared/native-chat-subagent-summary' import { diffFromText, diffFromToolCall, type DiffLine } from './native-chat-diff' @@ -31,6 +32,8 @@ import { selectActiveToolCall } from '../../../../shared/native-chat-tool-activity' import { nativeChatToolRunIconName } from '../../../../shared/native-chat-tool-icon' +import { NativeChatTaskList } from './NativeChatTaskList' +import { buildNativeChatTaskListRows } from './native-chat-task-list-history' import { NativeChatDiffView } from './NativeChatDiffView' import { NativeChatSubagentRun } from './NativeChatSubagentRun' import { NativeChatToolIcon, NativeChatToolRunIcon } from './NativeChatToolIcon' @@ -158,6 +161,8 @@ function ToolLine({ * toolbar toggle drive every run at once while still allowing per-run override. */ export function NativeChatToolRun({ blocks, + previousTodoWrite, + previousUpdatePlan, revealedDiff, onRevealDiff, subagentGroups = NO_SUBAGENT_GROUPS, @@ -168,6 +173,8 @@ export function NativeChatToolRun({ onLinkClick }: { blocks: NativeChatBlock[] + previousTodoWrite?: NativeChatToolCallBlock + previousUpdatePlan?: NativeChatToolCallBlock revealedDiff?: NativeChatDiffReveal onRevealDiff?: (element: HTMLElement) => void /** Spawn-group rosters that belong with this run's activity, one row each. */ @@ -231,6 +238,18 @@ export function NativeChatToolRun({ // The turn caret opens the activity group, while each child tool remains // collapsed. The global expand toolbar still opens child details together. const expandToolLines = expandOverride === undefined ? open : false + // Diffing every edit is the run's most expensive work, so a collapsed run — + // which renders none of it — never pays for it. + const taskLists = useMemo( + () => + open + ? buildNativeChatTaskListRows(blocks, { + todowrite: previousTodoWrite, + update_plan: previousUpdatePlan + }) + : null, + [open, blocks, previousTodoWrite, previousUpdatePlan] + ) // Rollups cache counts only; detailed diff rows are built when the run opens. const { editCards, consumedResults } = useMemo( () => (open ? buildEditCards(blocks) : NO_EDIT_CARDS), @@ -297,7 +316,7 @@ export function NativeChatToolRun({ rowWord={latestActiveCall.name} className="text-muted-foreground" /> - + {nativeChatToolActivityLabel(latestActiveCall)} {open ? : null} @@ -381,7 +400,14 @@ export function NativeChatToolRun({
    {(() => { const seen = new Map() - return blocks.map((block) => { + return blocks.map((block, blockIndex) => { + const taskList = taskLists?.rows.get(block) + if (taskList) { + return + } + if (taskLists?.consumedResults.has(block)) { + return null + } const edit = editCards.get(block) if (edit) { return ( diff --git a/src/renderer/src/components/native-chat/NativeChatTurnDiffRollup.tsx b/src/renderer/src/components/native-chat/NativeChatTurnDiffRollup.tsx index 8050710a177..6a7359ce4a6 100644 --- a/src/renderer/src/components/native-chat/NativeChatTurnDiffRollup.tsx +++ b/src/renderer/src/components/native-chat/NativeChatTurnDiffRollup.tsx @@ -29,7 +29,7 @@ export function NativeChatTurnDiffRollup({ ) : null} diff --git a/src/renderer/src/components/native-chat/native-chat-task-list-frames.ts b/src/renderer/src/components/native-chat/native-chat-task-list-frames.ts new file mode 100644 index 00000000000..f58c4a2bbdd --- /dev/null +++ b/src/renderer/src/components/native-chat/native-chat-task-list-frames.ts @@ -0,0 +1,36 @@ +import { normalizeNativeChatTaskList } from '../../../../shared/native-chat-task-list' +import type { NativeChatMessage } from '../../../../shared/native-chat-types' + +const projectedFrames = new WeakMap() + +/** Project after tool folding so a notification never takes another call's result. */ +export function projectNativeChatTaskListFrames( + messages: readonly NativeChatMessage[] +): NativeChatMessage[] { + return messages.map((message) => { + const cached = projectedFrames.get(message) + if (cached) { + return cached + } + const block = message.blocks.length === 1 ? message.blocks[0] : undefined + const frame = block?.type === 'text' ? block.providerFrame : undefined + if ( + message.role !== 'system' || + frame?.provider !== 'codex' || + frame.kind !== 'notification:turn/plan/updated' || + frame.payload.truncated || + !normalizeNativeChatTaskList('update_plan', frame.payload.head) + ) { + return message + } + const projected: NativeChatMessage = { + ...message, + role: 'assistant', + blocks: [ + { type: 'tool-call', name: 'update_plan', input: frame.payload.head, state: 'completed' } + ] + } + projectedFrames.set(message, projected) + return projected + }) +} diff --git a/src/renderer/src/components/native-chat/native-chat-task-list-history.test.ts b/src/renderer/src/components/native-chat/native-chat-task-list-history.test.ts new file mode 100644 index 00000000000..ef8eae291ad --- /dev/null +++ b/src/renderer/src/components/native-chat/native-chat-task-list-history.test.ts @@ -0,0 +1,114 @@ +import { describe, expect, it } from 'vitest' +import type { + NativeChatBlock, + NativeChatMessage, + NativeChatToolCallBlock +} from '../../../../shared/native-chat-types' +import { + buildNativeChatTaskListRows, + nativeChatTaskListPredecessors +} from './native-chat-task-list-history' + +function call(name = 'TodoWrite', status = 'pending'): NativeChatToolCallBlock { + return { + type: 'tool-call', + name, + input: + name === 'TodoWrite' + ? { todos: [{ content: 'Test', status }] } + : { plan: [{ step: 'Test', status }] } + } +} +function message( + id: string, + blocks: NativeChatBlock[], + role: NativeChatMessage['role'] = 'assistant' +): NativeChatMessage { + return { id, blocks, role, timestamp: 1, source: 'transcript' } +} + +describe('native chat task list history', () => { + it('carries predecessors across prose, ordinary tools, and user turns', () => { + const first = call() + const next = call('TodoWrite', 'completed') + const history = nativeChatTaskListPredecessors([ + message('a', [first]), + message('b', [{ type: 'text', text: 'Continue' }], 'user'), + message('c', [{ type: 'tool-call', name: 'Read', input: {} }]), + message('d', [next]) + ]) + expect(history.get('d')?.todowrite).toBe(first) + expect( + buildNativeChatTaskListRows([next], history.get('d')).rows.get(next)?.previous?.tasks[0] + .status + ).toBe('pending') + }) + + it('keeps interleaved tool families separate and ignores MCP lookalikes', () => { + const claude = call() + const codex = call('update_plan') + const next = call('TodoWrite', 'completed') + const model = buildNativeChatTaskListRows([claude, codex, call('mcp__x__TodoWrite'), next]) + expect(model.rows.get(codex)?.previous).toBeUndefined() + expect(model.rows.get(next)?.previous).toEqual(model.rows.get(claude)?.list) + const history = nativeChatTaskListPredecessors([ + message('a', [claude]), + message('b', [codex]), + message('c', [next]) + ]) + expect(history.get('c')).toEqual({ todowrite: claude, update_plan: codex }) + }) + + it('skips failed and malformed calls and keeps errors unconsumed', () => { + const first = call() + const failed = { ...call(), state: 'failed' as const } + const rejected = call('TodoWrite', 'completed') + const error: NativeChatBlock = { type: 'tool-result', output: 'Rejected', isError: true } + const next = call('TodoWrite', 'in_progress') + const blocks: NativeChatBlock[] = [ + first, + { type: 'tool-result', output: 'ok' }, + failed, + { type: 'tool-result', output: 'failed' }, + rejected, + error, + { ...call(), input: '{' }, + next + ] + const model = buildNativeChatTaskListRows(blocks) + expect(model.rows.has(failed)).toBe(false) + expect(model.rows.has(rejected)).toBe(false) + expect(model.consumedResults.has(error)).toBe(false) + expect(model.rows.get(next)?.previous).toEqual(model.rows.get(first)?.list) + const history = nativeChatTaskListPredecessors([ + message('a', blocks.slice(0, -1)), + message('b', [next]) + ]) + expect(history.get('b')?.todowrite).toBe(first) + }) + + it('updates predecessor identity after pagination and remains stable on rerender', () => { + const first = call() + const second = call('TodoWrite', 'in_progress') + const tail = message('b', [second]) + expect(nativeChatTaskListPredecessors([tail]).get('b')?.todowrite).toBeUndefined() + const history = nativeChatTaskListPredecessors([message('a', [first]), tail]) + expect(history.get('b')?.todowrite).toBe(first) + expect(nativeChatTaskListPredecessors([message('a', [first]), tail]).get('b')?.todowrite).toBe( + history.get('b')?.todowrite + ) + expect(nativeChatTaskListPredecessors([tail]).get('b')?.todowrite).toBeUndefined() + }) + + it('diffs a running call before its result arrives and consumes a successful result', () => { + const first = call() + const running = { ...call('TodoWrite', 'in_progress'), state: 'running' as const } + const result: NativeChatBlock = { type: 'tool-result', output: 'ok' } + const model = buildNativeChatTaskListRows([running, result], { + todowrite: first, + update_plan: undefined + }) + expect(model.rows.get(running)?.previous).toBeDefined() + expect(model.consumedResults.has(result)).toBe(true) + }) +}) diff --git a/src/renderer/src/components/native-chat/native-chat-task-list-history.ts b/src/renderer/src/components/native-chat/native-chat-task-list-history.ts new file mode 100644 index 00000000000..62582e2a7cf --- /dev/null +++ b/src/renderer/src/components/native-chat/native-chat-task-list-history.ts @@ -0,0 +1,83 @@ +import { + nativeChatTaskListTool, + normalizeNativeChatTaskList, + type NativeChatTaskList, + type NativeChatTaskListTool +} from '../../../../shared/native-chat-task-list' +import type { + NativeChatBlock, + NativeChatMessage, + NativeChatToolCallBlock +} from '../../../../shared/native-chat-types' +import { pairToolBlocks } from './native-chat-tool-fold' + +export type NativeChatTaskListPredecessors = Partial< + Record +> +export type NativeChatTaskListRow = { list: NativeChatTaskList; previous?: NativeChatTaskList } + +function taskListFromCall(call: NativeChatToolCallBlock): NativeChatTaskList | null { + return call.state === 'failed' ? null : normalizeNativeChatTaskList(call.name, call.input) +} + +/** Store call identities so unchanged rows stay memoized, while prepends replace their context. */ +export function nativeChatTaskListPredecessors( + messages: readonly NativeChatMessage[] +): Map { + const history = new Map() + const previous: NativeChatTaskListPredecessors = {} + for (const message of messages) { + history.set(message.id, { ...previous }) + if (message.role === 'user') { + continue + } + for (const { call, result } of pairToolBlocks(message.blocks)) { + if (!call || result?.isError) { + continue + } + const tool = nativeChatTaskListTool(call.name) + if (tool && taskListFromCall(call)) { + previous[tool] = call + } + } + } + return history +} + +export function buildNativeChatTaskListRows( + blocks: readonly NativeChatBlock[], + predecessors: NativeChatTaskListPredecessors = {} +): { + rows: Map + consumedResults: Set +} { + const rows = new Map() + const consumedResults = new Set() + const previous = new Map() + for (const call of Object.values(predecessors)) { + if (!call) { + continue + } + const tool = nativeChatTaskListTool(call.name) + const list = taskListFromCall(call) + if (tool && list) { + previous.set(tool, list) + } + } + for (const { call, result } of pairToolBlocks(blocks)) { + if (!call || result?.isError) { + continue + } + const tool = nativeChatTaskListTool(call.name) + const list = taskListFromCall(call) + if (!tool || !list) { + continue + } + rows.set(call, { list, previous: previous.get(tool) }) + previous.set(tool, list) + if (result) { + consumedResults.add(result) + } + } + return { rows, consumedResults } +} diff --git a/src/renderer/src/components/native-chat/native-chat-task-list-state.test.ts b/src/renderer/src/components/native-chat/native-chat-task-list-state.test.ts new file mode 100644 index 00000000000..c282c4d4532 --- /dev/null +++ b/src/renderer/src/components/native-chat/native-chat-task-list-state.test.ts @@ -0,0 +1,73 @@ +import { describe, expect, it } from 'vitest' +import type { NativeChatBlock, NativeChatMessage } from '../../../../shared/native-chat-types' +import { nativeChatTaskListState } from './native-chat-task-list-state' + +function message(id: string, blocks: NativeChatBlock[]): NativeChatMessage { + return { id, role: 'assistant', timestamp: 1, source: 'transcript', blocks } +} +function call(content: string, status = 'pending'): NativeChatBlock { + return { type: 'tool-call', name: 'TodoWrite', input: { todos: [{ content, status }] } } +} + +describe('nativeChatTaskListState', () => { + it('projects one latest snapshot, preserves prose and leaves source messages unchanged', () => { + const first = message('first', [call('Read')]) + const last = message('last', [ + { type: 'text', text: 'Here is the result' }, + call('Read', 'completed'), + { type: 'tool-result', output: 'Updated todos' } + ]) + const result = nativeChatTaskListState([first, last]) + expect(result.list?.tasks).toEqual([{ content: 'Read', status: 'completed' }]) + expect(result.messages[0]).toBe(first) + expect(result.messages[1]).toBe(last) + expect(first.blocks).toHaveLength(1) + expect(last.blocks).toHaveLength(3) + expect(nativeChatTaskListState([first, last]).messages[1]).toBe(result.messages[1]) + }) + + it('preserves latest state across user follow-ups and clears it on an explicit empty list', () => { + const first = message('first', [call('Read')]) + const user = { ...message('user', [{ type: 'text', text: 'Continue' }]), role: 'user' as const } + const empty = message('empty', [{ type: 'tool-call', name: 'TodoWrite', input: { todos: [] } }]) + expect(nativeChatTaskListState([first, user]).list?.tasks).toHaveLength(1) + expect(nativeChatTaskListState([first, user, empty]).list?.tasks).toEqual([]) + expect(nativeChatTaskListState([]).list).toBeNull() + }) + + it('does not replace valid state with malformed or failed calls, and retains their diagnostics', () => { + const first = message('first', [call('Read')]) + const malformed = message('malformed', [ + { type: 'tool-call', name: 'TodoWrite', input: '{' }, + { type: 'tool-result', output: 'Invalid arguments', isError: true } + ]) + const failed = message('failed', [ + call('Wrong', 'completed'), + { type: 'tool-result', output: 'Update rejected', isError: true } + ]) + const failedCall = message('failed-call', [ + { type: 'tool-call', name: 'TodoWrite', state: 'failed', input: { todos: [] } } + ]) + const result = nativeChatTaskListState([first, malformed, failed, failedCall]) + expect(result.list?.tasks[0].content).toBe('Read') + expect(result.messages.slice(1)).toEqual([malformed, failed, failedCall]) + }) + + it('retains task history and unrelated errors while selecting the paired snapshot', () => { + const tasks: NativeChatBlock = call('Read') + const shell: NativeChatBlock = { + type: 'tool-call', + name: 'shell', + input: {} + } + const error: NativeChatBlock = { + type: 'tool-result', + output: 'Failed', + isError: true + } + const success: NativeChatBlock = { type: 'tool-result', output: 'Updated' } + const result = nativeChatTaskListState([message('mixed', [tasks, success, shell, error])]) + expect(result.list?.tasks[0].content).toBe('Read') + expect(result.messages[0].blocks).toEqual([tasks, success, shell, error]) + }) +}) diff --git a/src/renderer/src/components/native-chat/native-chat-task-list-state.ts b/src/renderer/src/components/native-chat/native-chat-task-list-state.ts new file mode 100644 index 00000000000..f45d6531600 --- /dev/null +++ b/src/renderer/src/components/native-chat/native-chat-task-list-state.ts @@ -0,0 +1,43 @@ +import { + normalizeNativeChatTaskList, + type NativeChatTaskList +} from '../../../../shared/native-chat-task-list' +import type { NativeChatMessage } from '../../../../shared/native-chat-types' +import { pairToolBlocks } from './native-chat-tool-fold' + +const snapshots = new WeakMap() + +function latestSnapshot(message: NativeChatMessage): NativeChatTaskList | null { + if (snapshots.has(message)) { + return snapshots.get(message) ?? null + } + let list: NativeChatTaskList | null = null + if (message.role === 'assistant') { + for (const { call, result } of pairToolBlocks(message.blocks)) { + if (!call || call.state === 'failed' || result?.isError) { + continue + } + const snapshot = normalizeNativeChatTaskList(call.name, call.input) + if (snapshot) { + list = snapshot + } + } + } + snapshots.set(message, list) + return list +} + +/** Select composer progress without consuming historical transcript updates. */ +export function nativeChatTaskListState(messages: readonly NativeChatMessage[]): { + messages: readonly NativeChatMessage[] + list: NativeChatTaskList | null +} { + let list: NativeChatTaskList | null = null + for (const message of messages) { + const snapshot = latestSnapshot(message) + if (snapshot) { + list = snapshot + } + } + return { messages, list } +} diff --git a/src/renderer/src/i18n/locales/en.json b/src/renderer/src/i18n/locales/en.json index 2f33a2d9745..4f283dc4f47 100644 --- a/src/renderer/src/i18n/locales/en.json +++ b/src/renderer/src/i18n/locales/en.json @@ -16994,6 +16994,22 @@ "empty": "No users found" }, "native-chat": { + "taskList": { + "title": "Tasks", + "completed": "Completed", + "inProgress": "In progress", + "pending": "Pending", + "empty": "No tasks", + "progress": "{{completed}} of {{total}} tasks completed", + "added": "Added {{task}}", + "removed": "Removed {{task}}", + "started": "Started {{task}}", + "finished": "Completed {{task}}", + "reset": "Marked pending: {{task}}", + "updated": "Updated {{task}}", + "unchanged": "Tasks unchanged", + "showAll": "Full task list" + }, "turnDiff": { "one": "1 changed file", "many": "{{count}} changed files", diff --git a/src/shared/native-chat-task-list.test.ts b/src/shared/native-chat-task-list.test.ts new file mode 100644 index 00000000000..0bbb7077f3c --- /dev/null +++ b/src/shared/native-chat-task-list.test.ts @@ -0,0 +1,138 @@ +import { describe, expect, it } from 'vitest' +import { + diffNativeChatTaskLists, + nativeChatTaskLabel, + normalizeNativeChatTaskList, + type NativeChatTask, + type NativeChatTaskList +} from './native-chat-task-list' + +const task = (content: string, status: NativeChatTask['status'] = 'pending'): NativeChatTask => ({ + content, + status +}) +const list = (...tasks: NativeChatTask[]): NativeChatTaskList => ({ tasks }) + +describe('normalizeNativeChatTaskList', () => { + it('normalizes Claude tasks and uses activeForm only while in progress', () => { + const result = normalizeNativeChatTaskList('TodoWrite', { + todos: [ + { content: 'Read', status: 'completed', activeForm: 'Reading' }, + { content: 'Write', status: 'in_progress', activeForm: 'Writing' }, + { content: 'Test', status: 'pending', activeForm: 'Testing' } + ] + })! + expect(result.tasks.map(nativeChatTaskLabel)).toEqual(['Read', 'Writing', 'Test']) + expect(result.tasks.map((entry) => entry.status)).toEqual([ + 'completed', + 'in_progress', + 'pending' + ]) + }) + + it('normalizes Codex JSON-string arguments and explanation', () => { + expect( + normalizeNativeChatTaskList( + 'update_plan', + JSON.stringify({ + explanation: 'Proceed with verification', + plan: [{ step: 'Test', status: 'in_progress' }] + }) + ) + ).toEqual({ explanation: 'Proceed with verification', tasks: [task('Test', 'in_progress')] }) + }) + + it('defaults unknown/missing statuses and ignores invalid entries', () => { + expect( + normalizeNativeChatTaskList(' TodoWrite ', { + todos: [ + null, + [], + 4, + {}, + { content: ' ' }, + { content: 7 }, + { content: ' One ', status: 'unknown', activeForm: 4 }, + { content: 'Two' } + ] + }) + ).toEqual(list(task('One'), task('Two'))) + }) + + it.each([undefined, null, 42, [], '{', '{}', { todos: null }, { todos: [{}] }])( + 'returns null for malformed input %j', + (input) => { + expect(normalizeNativeChatTaskList('TodoWrite', input)).toBeNull() + } + ) + + it('keeps empty lists valid and recognizes only exact tool families', () => { + expect(normalizeNativeChatTaskList('update_plan', { plan: [] })).toEqual(list()) + expect(normalizeNativeChatTaskList('TodoWrite', { todos: [] })).toEqual(list()) + expect(normalizeNativeChatTaskList('mcp__server__TodoWrite', { todos: [] })).toBeNull() + expect(normalizeNativeChatTaskList('ExitPlanMode', { plan: [] })).toBeNull() + expect(normalizeNativeChatTaskList('update_plan', { todos: [] })).toBeNull() + }) +}) + +describe('diffNativeChatTaskLists', () => { + it('reports completions and starts, omitting unchanged tasks', () => { + expect( + diffNativeChatTaskLists( + list(task('Read', 'in_progress'), task('Write'), task('Test')), + list(task('Read', 'completed'), task('Write', 'in_progress'), task('Test')) + ) + ).toEqual([ + { kind: 'completed', task: task('Read', 'completed') }, + { kind: 'started', task: task('Write', 'in_progress') } + ]) + }) + + it('ignores reorder-only updates and explanation changes', () => { + expect( + diffNativeChatTaskLists(list(task('A'), task('B')), { + tasks: [task('B'), task('A')], + explanation: 'Reordered' + }) + ).toEqual([]) + }) + + it('matches duplicate contents by occurrence', () => { + expect( + diffNativeChatTaskLists( + list(task('A'), task('A', 'in_progress')), + list(task('A', 'completed'), task('A', 'in_progress')) + ) + ).toEqual([{ kind: 'completed', task: task('A', 'completed') }]) + }) + + it('reports renamed content as an addition and removal', () => { + expect(diffNativeChatTaskLists(list(task('Old')), list(task('New')))).toEqual([ + { kind: 'added', task: task('New') }, + { kind: 'removed', task: task('Old') } + ]) + }) + + it('reports resets, reopening, and activeForm-only edits', () => { + const changed = { ...task('C', 'in_progress'), activeForm: 'Checking C' } + expect( + diffNativeChatTaskLists( + list(task('A', 'completed'), task('B', 'completed'), task('C', 'in_progress')), + list(task('A'), task('B', 'in_progress'), changed) + ) + ).toEqual([ + { kind: 'pending', task: task('A') }, + { kind: 'started', task: task('B', 'in_progress') }, + { kind: 'updated', task: changed } + ]) + }) + + it('reports clearing a list and removing a duplicate', () => { + expect(diffNativeChatTaskLists(list(task('A')), list())).toEqual([ + { kind: 'removed', task: task('A') } + ]) + expect(diffNativeChatTaskLists(list(task('A'), task('A')), list(task('A')))).toEqual([ + { kind: 'removed', task: task('A') } + ]) + }) +}) diff --git a/src/shared/native-chat-task-list.ts b/src/shared/native-chat-task-list.ts new file mode 100644 index 00000000000..4d4d470b246 --- /dev/null +++ b/src/shared/native-chat-task-list.ts @@ -0,0 +1,122 @@ +export type NativeChatTaskStatus = 'pending' | 'in_progress' | 'completed' +export type NativeChatTask = { + content: string + status: NativeChatTaskStatus + activeForm?: string +} +export type NativeChatTaskList = { tasks: NativeChatTask[]; explanation?: string } +export type NativeChatTaskChange = { + kind: 'added' | 'removed' | 'started' | 'completed' | 'pending' | 'updated' + task: NativeChatTask +} +export type NativeChatTaskListTool = 'todowrite' | 'update_plan' + +export function nativeChatTaskListTool(name: string): NativeChatTaskListTool | null { + const normalized = name.trim().toLowerCase() + return normalized === 'todowrite' || normalized === 'update_plan' ? normalized : null +} + +function record(value: unknown): Record | null { + return value !== null && typeof value === 'object' && !Array.isArray(value) + ? (value as Record) + : null +} + +function nonemptyString(value: unknown): string | undefined { + return typeof value === 'string' && value.trim() ? value.trim() : undefined +} + +export function normalizeNativeChatTaskList( + name: string, + input: unknown +): NativeChatTaskList | null { + const tool = nativeChatTaskListTool(name) + if (!tool) { + return null + } + if (typeof input === 'string') { + try { + input = JSON.parse(input) + } catch { + return null + } + } + const value = record(input) + const entries = tool === 'todowrite' ? value?.todos : value?.plan + if (!Array.isArray(entries)) { + return null + } + const tasks: NativeChatTask[] = [] + for (const entry of entries) { + const item = record(entry) + const content = nonemptyString(tool === 'todowrite' ? item?.content : item?.step) + if (!item || !content) { + continue + } + const status = + item.status === 'in_progress' || (tool === 'update_plan' && item.status === 'inProgress') + ? 'in_progress' + : item.status === 'completed' + ? 'completed' + : 'pending' + const activeForm = tool === 'todowrite' ? nonemptyString(item.activeForm) : undefined + tasks.push({ content, status, ...(activeForm ? { activeForm } : {}) }) + } + if (entries.length > 0 && tasks.length === 0) { + return null + } + const explanation = tool === 'update_plan' ? nonemptyString(value?.explanation) : undefined + return { tasks, ...(explanation ? { explanation } : {}) } +} + +export function nativeChatTaskLabel(task: NativeChatTask): string { + return task.status === 'in_progress' && task.activeForm ? task.activeForm : task.content +} + +/** Content plus occurrence is the only identity the providers give these entries. */ +export function diffNativeChatTaskLists( + previous: NativeChatTaskList, + current: NativeChatTaskList +): NativeChatTaskChange[] { + const byContent = new Map() + for (const task of previous.tasks) { + const matches = byContent.get(task.content) + if (matches) { + matches.push(task) + } else { + byContent.set(task.content, [task]) + } + } + const occurrences = new Map() + const consumed = new Set() + const changes: NativeChatTaskChange[] = [] + for (const task of current.tasks) { + const occurrence = occurrences.get(task.content) ?? 0 + occurrences.set(task.content, occurrence + 1) + const before = byContent.get(task.content)?.[occurrence] + if (!before) { + changes.push({ kind: 'added', task }) + continue + } + consumed.add(before) + if (before.status !== task.status) { + changes.push({ + kind: + task.status === 'completed' + ? 'completed' + : task.status === 'in_progress' + ? 'started' + : 'pending', + task + }) + } else if (before.activeForm !== task.activeForm) { + changes.push({ kind: 'updated', task }) + } + } + for (const task of previous.tasks) { + if (!consumed.has(task)) { + changes.push({ kind: 'removed', task }) + } + } + return changes +} diff --git a/src/shared/native-chat-tool-icon.test.ts b/src/shared/native-chat-tool-icon.test.ts index 695400c94ca..f8cabed950c 100644 --- a/src/shared/native-chat-tool-icon.test.ts +++ b/src/shared/native-chat-tool-icon.test.ts @@ -50,6 +50,9 @@ describe('native chat tool icons', () => { expect(nativeChatToolCategory('list')).toBe('listFiles') expect(nativeChatToolCategory('shell')).toBe('unknown') expect(nativeChatToolCategory('apply_patch')).toBe('fileChange') + expect(nativeChatToolCategory('update_plan')).toBe('todoList') + expect(nativeChatToolIconName('update_plan')).toBe('list-checks') + expect(nativeChatToolRunIconName([{ name: 'update_plan' }])).toBe('list-checks') expect(nativeChatToolCategory('web search')).toBe('webSearch') }) diff --git a/src/shared/native-chat-tool-icon.ts b/src/shared/native-chat-tool-icon.ts index d546a76e865..d52df9e52e9 100644 --- a/src/shared/native-chat-tool-icon.ts +++ b/src/shared/native-chat-tool-icon.ts @@ -81,6 +81,7 @@ const CATEGORY_BY_ROW_WORD = new Map([ ['task', 'subAgentActivity'], ['webfetch', 'webSearch'], ['todowrite', 'todoList'], + ['update_plan', 'todoList'], ['web search', 'webSearch'], ['websearch', 'webSearch'], ['web_search', 'webSearch'] diff --git a/src/shared/orchestration-rpc-contract.ts b/src/shared/orchestration-rpc-contract.ts index 3067d5ad6b0..fec4d4ed4e7 100644 --- a/src/shared/orchestration-rpc-contract.ts +++ b/src/shared/orchestration-rpc-contract.ts @@ -14,6 +14,9 @@ export const ORCHESTRATION_SKILL_COMMAND_ARGS = [ ] as const export const ORCHESTRATION_LEGACY_RUN_ID = 'run_legacy_local' +/** Files mail sent from a terminal in no Run. Distinct from the legacy Run on purpose: rows under + * the legacy id read as a pre-Runs database and trigger adoption on the next open. */ +export const ORCHESTRATION_UNBOUND_RUN_ID = 'run_unbound' const ORCHESTRATION_MUTATION_METHODS = new Set([ 'orchestration.runCreate', diff --git a/tests/e2e/orchestration-idle-mail-delivery.spec.ts b/tests/e2e/orchestration-idle-mail-delivery.spec.ts index f679bdba608..cffa6415208 100644 --- a/tests/e2e/orchestration-idle-mail-delivery.spec.ts +++ b/tests/e2e/orchestration-idle-mail-delivery.spec.ts @@ -350,6 +350,9 @@ test.describe('orchestration push-on-idle mail delivery', () => { await expectSubmitted(pane) }) + // #19542 deleted the legacy-Run write fallback, so a sender in no Run has + // nowhere to file mail to a bare handle: the send is refused outright, which + // is what keeps an unsafe pointer out of the pane on the next idle frame. test('keeps unbound direct mail durable without pointing to an unsafe check', async ({ orcaPage, electronApp @@ -359,6 +362,8 @@ test.describe('orchestration push-on-idle mail delivery', () => { const pane = await openAgentPane() await driveToLiveIdle(client, pane) + // Two plain terminals, neither in a Run: `send --to ` must still land durably. It + // files under the unbound Run, so a reopen never reads it as pre-Runs state (#19542 regression). const stdinBeforeScan = pane.agent.readStdin() const messageId = await sendMail(client, pane.handle, { subject: 'Unbound direct mail' }) pane.agent.setTitle(CODEX_WORKING_TITLE) @@ -366,8 +371,10 @@ test.describe('orchestration push-on-idle mail delivery', () => { pane.agent.setTitle(CODEX_IDLE_TITLE) await waitForObservedTitle(client, pane.handle, CODEX_IDLE_TITLE) + await orcaPage.waitForTimeout(NO_DELIVERY_SETTLE_MS) expect(readMailRow(userDataDir, messageId)).toMatchObject({ to_handle: pane.handle, + run_id: 'run_unbound', read: 0, delivered_at: null })