diff --git a/src/main/runtime/orchestration/db/contract-constants.ts b/src/main/runtime/orchestration/db/contract-constants.ts index 71875cc9167..4f9975529e4 100644 --- a/src/main/runtime/orchestration/db/contract-constants.ts +++ b/src/main/runtime/orchestration/db/contract-constants.ts @@ -7,4 +7,4 @@ export const LEGACY_CONTRACT_VERSION = 0 export const CURRENT_CONTRACT_VERSION = ORCHESTRATION_CONTRACT_VERSION // Schema versions: v2 'heartbeat'+last_heartbeat_at, v3 delivered_at, v4 task-creator terminal, v5 task_title/display_name, v6 pane identity, v7 lightweight Runs, v8 crash-safe Run deliveries, v9 durable question threads, v10 Dispatch capabilities, v11 durable mutation receipts, v12 composed worker state, v18 post-v6 version-skew repair, v19 adopted legacy Runs and compatibility receipts, v20 legacy question backfill, v21 legacy scheduler-loss provenance, v22 dispatch assignee lookup, v23 worker terminal resource ownership, v24 creator-incarnation authority, v25 active Dispatch handle lookup, v26 indexed mutation receipt capacity, v27 durable federation acknowledgments, v28 durable local mutation caller identity, v31 dispatch/resource identity links, v32 bounded worker-terminal recovery metadata, v33 durable mailbox pointer Enter state, v34 role-addressed mailbox deliveries, v35 mailbox delivery default and index-predicate repair, v36 dispatch mailbox consumer generation, v37 recorded dispatch creator identity. -export const SCHEMA_VERSION = 37 +export const SCHEMA_VERSION = 38 diff --git a/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts b/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts index 2e281ca00f9..be1cf99d2b4 100644 --- a/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts +++ b/src/main/runtime/orchestration/db/dispatch-context/dispatch-completion.ts @@ -29,6 +29,9 @@ export function completeDispatch(this: OrchestrationDb, ctxId: string): void { capability_revoked_at: dispatch.capability_revoked_at ?? new Date().toISOString() } }) + // Why: a settled Dispatch can never be answered, and a pending thread on it kept the fleet row + // demanding input after the work was done. + this.closeQuestionsForDispatch(ctxId) this.db.exec('RELEASE complete_dispatch_transition') } catch (error) { this.db.exec('ROLLBACK TO complete_dispatch_transition') @@ -63,6 +66,7 @@ export function settleActiveDispatchesForTask( capability_revoked_at: row.capability_revoked_at ?? new Date().toISOString() } }) + db.closeQuestionsForDispatch(row.id) } } @@ -202,6 +206,7 @@ export function failDispatch( projection: { completed_at: taskStatus === 'failed' ? new Date().toISOString() : null } }) } + this.closeQuestionsForDispatch(ctxId) const updated = this.db.prepare('SELECT * FROM dispatch_contexts WHERE id = ?').get(ctxId) as | DispatchContextRow | undefined diff --git a/src/main/runtime/orchestration/db/schema/migrate-v38.ts b/src/main/runtime/orchestration/db/schema/migrate-v38.ts new file mode 100644 index 00000000000..0fa1ddcb12c --- /dev/null +++ b/src/main/runtime/orchestration/db/schema/migrate-v38.ts @@ -0,0 +1,21 @@ +import type { OrchestrationDb } from '../orchestration-db' + +/** + * Settling a Dispatch through the task-status path never closed its pending question threads, so + * a completed pre-v3 row kept an `input` attention category forever. Nothing can answer a question + * on a settled Dispatch (`answerQuestion` refuses closed threads and the Dispatch is inactive), so + * closing them is the only reading that matches the row. + */ +export function migrateV38(this: OrchestrationDb, current: number): void { + if (current >= 38) { + return + } + this.db.exec( + `UPDATE question_threads + SET status = 'closed', closed_at = datetime('now') + WHERE status = 'pending' + AND dispatch_id IN ( + SELECT id FROM dispatch_contexts WHERE status NOT IN ('pending', 'dispatched') + )` + ) +} diff --git a/src/main/runtime/orchestration/db/schema/migrate.ts b/src/main/runtime/orchestration/db/schema/migrate.ts index 16c4a520ecc..fade2bf15e4 100644 --- a/src/main/runtime/orchestration/db/schema/migrate.ts +++ b/src/main/runtime/orchestration/db/schema/migrate.ts @@ -8,6 +8,7 @@ import { migrateRoleMailboxDeliveryV34 } from './migrate-role-mailbox-delivery-v import { migrateV35 } from './migrate-v35' import { migrateV36 } from './migrate-v36' import { migrateV37 } from './migrate-v37' +import { migrateV38 } from './migrate-v38' // Why: CREATE TABLE IF NOT EXISTS won't alter existing DBs; migrate in a txn that bumps user_version only on success (atomic all-or-nothing). export function migrate(this: OrchestrationDb): void { @@ -26,6 +27,7 @@ export function migrate(this: OrchestrationDb): void { migrateV35.call(this, current) migrateV36.call(this, current) migrateV37.call(this, current) + migrateV38.call(this, current) this.db.pragma(`user_version = ${SCHEMA_VERSION}`) this.db.exec('COMMIT') } catch (err) { diff --git a/src/main/runtime/orchestration/settled-question-threads-migration.test.ts b/src/main/runtime/orchestration/settled-question-threads-migration.test.ts new file mode 100644 index 00000000000..6932e7dc229 --- /dev/null +++ b/src/main/runtime/orchestration/settled-question-threads-migration.test.ts @@ -0,0 +1,67 @@ +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import Database from '../../sqlite/sync-database' +import { OrchestrationDb } from './db' +import { SCHEMA_VERSION } from './db/contract-constants' +import { createRootDispatch } from './db/root-dispatch-test-fixture' + +/** v38 closes question threads left pending on Dispatches that settled through the task path. */ +describe('OrchestrationDb v37 to v38 migration', () => { + let db: OrchestrationDb | undefined + let tempDir: string | undefined + + afterEach(() => { + db?.close() + db = undefined + if (tempDir) { + rmSync(tempDir, { recursive: true, force: true }) + tempDir = undefined + } + }) + + /** A v37 database with one pending question on a settled Dispatch and one on an active one. */ + function createV37Database(): { path: string; settled: string; active: string } { + tempDir = mkdtempSync(join(tmpdir(), 'orca-db-v38-')) + const dbPath = join(tempDir, 'orchestration.db') + const seed = new OrchestrationDb(dbPath) + const run = seed.createRun({ + objective: 'pre-v38 run', + coordinatorHandle: 'term_coord', + coordinatorPaneKey: 'tab_coord:cccccccc-cccc-4ccc-8ccc-cccccccccccc' + }) + const ask = (dispatchId: string) => + seed.createQuestion({ + runId: run.id, + dispatchId, + askerHandle: 'term_worker', + question: 'still pending?' + }).question.message_id + const settledTask = seed.createTask({ spec: 'settled before v38', runId: run.id }) + const settledDispatch = createRootDispatch(seed, settledTask.id, 'term_worker') + const settled = ask(settledDispatch.id) + const activeTask = seed.createTask({ spec: 'still running', runId: run.id }) + const active = ask(createRootDispatch(seed, activeTask.id, 'term_worker_2').id) + seed.close() + + // Why: pre-v38 settlement left the thread pending; recreate that on-disk shape directly. + const raw = new Database(dbPath) + raw + .prepare("UPDATE dispatch_contexts SET status = 'completed' WHERE id = ?") + .run(settledDispatch.id) + raw.prepare("UPDATE question_threads SET status = 'pending', closed_at = NULL").run() + raw.pragma('user_version = 37') + raw.close() + return { path: dbPath, settled, active } + } + + it('closes pending questions on settled dispatches and keeps active ones pending', () => { + const v37 = createV37Database() + db = new OrchestrationDb(v37.path) + + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + expect(db.getQuestion(v37.settled)?.status).toBe('closed') + expect(db.getQuestion(v37.active)?.status).toBe('pending') + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federation-attachment-observation.ts b/src/main/runtime/rpc/methods/orchestration/federation/federation-attachment-observation.ts index 0fe506c6d7c..df9059609ec 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federation-attachment-observation.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federation-attachment-observation.ts @@ -1,6 +1,7 @@ import type { RuntimeTerminalInteractiveWait } from '../../../../../../shared/runtime-types' import type { OrcaRuntimeService } from '../../../../orca-runtime' import { OrchestrationError } from '../../../../orchestration/orchestration-error' +import { parseWorkerTerminalHostScope } from '../../../../orchestration/worker-terminal-process-liveness' import type { RemoteDispatchAttachmentRow } from '../../../../orchestration/types' export function requireHomeAttachment( @@ -54,6 +55,24 @@ export async function inspectRemoteAttachment( return { terminal, exact, status: 'unverifiable', reason: verdict.reason, agentWait } } if (!verdict) { + // Why: the verdict register only fills on the first inventory sweep or exit frame, so a PTY + // this host just spawned has none for minutes and every fleet row read host_indeterminate. + // The host owns a connected local pane, so its own connected flag is host evidence of life, + // exactly as worker-show reads it. Nothing weaker earns a claim: a disconnected pane or an + // SSH-scoped one (contact, not the process) stays unverifiable, never `exited`. + const currentHostScope = runtime.getOrchestrationDispatchAuthority?.( + attachment.terminal_handle + )?.hostScope + const persistedHostScope = parseWorkerTerminalHostScope( + db.getWorkerTerminalResourceByOwner(dispatchId)?.host_scope ?? null + ) + const provenLocal = + currentHostScope !== undefined && + currentHostScope.kind !== 'ssh' && + persistedHostScope?.kind !== 'ssh' + if (provenLocal && terminal.connected !== false) { + return { terminal, exact, status: 'live', agentWait } + } return { terminal, exact, diff --git a/src/main/runtime/rpc/methods/orchestration/federation/federation-liveness-verdict.test.ts b/src/main/runtime/rpc/methods/orchestration/federation/federation-liveness-verdict.test.ts index 58be92d997c..20ae3135ec2 100644 --- a/src/main/runtime/rpc/methods/orchestration/federation/federation-liveness-verdict.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/federation/federation-liveness-verdict.test.ts @@ -270,6 +270,57 @@ describe('federation host liveness verdicts', () => { } }) + // Why: the verdict register only fills on the first inventory sweep, so a PTY this host just + // spawned has none for minutes; the fleet row read host_indeterminate the whole time. + it('reads a freshly spawned local pane from its own connected flag before any verdict', async () => { + vi.spyOn(runtime, 'showTerminal').mockResolvedValue({ + handle: HANDLE, + worktreeId: 'repo::remote-worktree', + connected: true, + status: 'running' + } as never) + vi.spyOn(runtime, 'getTerminalLivenessVerdict').mockReturnValue(null) + vi.spyOn(runtime, 'getOrchestrationDispatchAuthority').mockReturnValue({ + hostScope: { kind: 'local', hostId: 'local' } + } as never) + + await expect( + call('orchestration.federationFleetSnapshot', { dispatchIds: [DISPATCH_ID] }) + ).resolves.toMatchObject({ + items: [{ dispatchId: DISPATCH_ID, observation: { status: 'live', exactWorker: true } }] + }) + }) + + it('keeps a disconnected verdict-less pane unverifiable rather than exited', async () => { + vi.spyOn(runtime, 'getTerminalLivenessVerdict').mockReturnValue(null) + vi.spyOn(runtime, 'getOrchestrationDispatchAuthority').mockReturnValue(null) + + await expect( + call('orchestration.federationShow', { dispatchId: DISPATCH_ID }) + ).resolves.toMatchObject({ + observation: { status: 'unverifiable', exactWorker: true, reason: 'missing_liveness_verdict' } + }) + }) + + it('keeps a verdict-less pane the host reaches over SSH unverifiable', async () => { + vi.spyOn(runtime, 'showTerminal').mockResolvedValue({ + handle: HANDLE, + worktreeId: 'repo::remote-worktree', + connected: true, + status: 'running' + } as never) + vi.spyOn(runtime, 'getTerminalLivenessVerdict').mockReturnValue(null) + vi.spyOn(runtime, 'getOrchestrationDispatchAuthority').mockReturnValue({ + hostScope: { kind: 'ssh', targetId: 'ssh-hop' } + } as never) + + await expect( + call('orchestration.federationShow', { dispatchId: DISPATCH_ID }) + ).resolves.toMatchObject({ + observation: { status: 'unverifiable', exactWorker: true, reason: 'missing_liveness_verdict' } + }) + }) + it('keeps an old peer without a liveness verdict unverifiable', async () => { // Legacy hosts can return an exited-looking terminal summary but have no // verdict API; relay/contact state is not proof that the process exited. diff --git a/src/main/runtime/rpc/methods/orchestration/worker/legacy-dispatch-projection.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/legacy-dispatch-projection.test.ts index 845bd0d7744..4df73994beb 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/legacy-dispatch-projection.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/legacy-dispatch-projection.test.ts @@ -63,6 +63,35 @@ describe('pre-v3 dispatch rows in worker-list', () => { expect(worker.projection.nextAction.kind).toBe('none') }) + it.each(['completed', 'failed'] as const)( + 'closes a pending question when a legacy dispatch settles as %s', + async (status) => { + h.setup() + const task = h.db.createTask({ spec: `legacy ${status} with question`, runId: h.activeRunId }) + const dispatch = createRootDispatch(h.db, task.id, `term_legacy_q_${status}`) + const asked = h.db.createQuestion({ + runId: h.activeRunId, + dispatchId: dispatch.id, + askerHandle: `term_legacy_q_${status}`, + question: 'Which branch?' + }) + // Both settlement paths a pre-v3 dispatch can take: the task-status path and failDispatch. + if (status === 'completed') { + h.db.updateTaskStatus(task.id, 'completed', 'done') + } else { + h.db.failDispatch(dispatch.id, 'legacy failure') + } + + const worker = (await listWorkers()).get(dispatch.id)! + + expect(h.db.getQuestion(asked.question.message_id)?.status).toBe('closed') + expect(worker.dispatchStatus).toBe(status) + expect(worker.projection.attention.categories).not.toContain('input') + // Nothing can answer a question on a settled Dispatch, so `input` must not outlive it. + expect(worker.projection.attention.requiresAction).toBe(status === 'failed') + } + ) + it('keeps a legacy failed dispatch actionable on the failure, not on absence', async () => { h.setup() const failed = createLegacyDispatch('failed')