mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 16:02:03 +00:00
fix(orchestration): read a fresh paired-host pane as live and close questions on settlement
Two findings from the final live cross-host smoke on b082443e1f.
A worker just started on a paired server projected unverifiable /
host_indeterminate with requiresAction for ~3 minutes, including after
its own worker_done. The host's verdict register only fills on the first
inventory sweep or exit frame, so the federation observation returned
missing_liveness_verdict for every freshly spawned PTY. The host owns a
connected local pane, so its own connected flag is host evidence of
life, exactly as worker-show already reads it; a disconnected pane or
one the host reaches over SSH still stays unverifiable, never exited.
Six pre-v3 completed dispatches still carried an input attention
category: settling through the task-status path or failDispatch never
closed the Dispatch's pending question threads, and nothing can answer a
question on a settled Dispatch. Both paths close them now, and schema
v38 closes the threads already left pending on settled rows.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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')
|
||||
)`
|
||||
)
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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')
|
||||
})
|
||||
})
|
||||
+19
@@ -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,
|
||||
|
||||
+51
@@ -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.
|
||||
|
||||
@@ -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')
|
||||
|
||||
Reference in New Issue
Block a user