diff --git a/src/main/runtime/orchestration/db/contract-constants.ts b/src/main/runtime/orchestration/db/contract-constants.ts index a8ddf9994d0c..f2a8fc12fd1b 100644 --- a/src/main/runtime/orchestration/db/contract-constants.ts +++ b/src/main/runtime/orchestration/db/contract-constants.ts @@ -18,4 +18,5 @@ 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, v39 structured session journal archives. // v41: derive outstanding deliveries from unread messages. -export const SCHEMA_VERSION = 41 +// v42: structured-session Orca session id columns. +export const SCHEMA_VERSION = 42 diff --git a/src/main/runtime/orchestration/db/dispatch-depth.test.ts b/src/main/runtime/orchestration/db/dispatch-depth.test.ts index 013cedffe91b..28d94e9556b1 100644 --- a/src/main/runtime/orchestration/db/dispatch-depth.test.ts +++ b/src/main/runtime/orchestration/db/dispatch-depth.test.ts @@ -1,6 +1,12 @@ import { afterEach, describe, expect, it } from 'vitest' +import { + mintStructuredWorkerHandle, + mintStructuredWorkerPaneKey, + structuredWorkerProcessIncarnation +} from '../../structured-worker-identity' import { OrchestrationDb } from '../db' import { AmbiguousDispatchParentError } from './dispatch-depth' +import { backfillStructuredWorkerOrcaSessionIds } from './schema/structured-worker-orca-session-backfill' /** * These pin the fence Orca documented but never enforced: before this feature a @@ -283,4 +289,47 @@ describe('nested worker depth', () => { expect(row.process_incarnation).toBeNull() expect(db.resolveCreatorDepth({ kind: 'terminal', handle: 'term_ctx' })).toBe(1) }) + + // Pinned for the reader that switches self-dispatch detection to Orca session id equality: equal + // creator and assignee ids must keep meaning bookkeeping, and different ones delegation. + it('records equal Orca session ids exactly when a structured session dispatches to itself', () => { + db = new OrchestrationDb(':memory:') + const sessionId = '5c7e9a1d-3f6b-4c8e-8d2a-4b6c8e0a2d36' + const self = { + kind: 'terminal', + handle: mintStructuredWorkerHandle(), + paneKey: mintStructuredWorkerPaneKey(sessionId) + } as const + const own = db.createDispatchContext({ + taskId: db.createTask({ runId: 'run_legacy_local', spec: 'own bookkeeping' }).id, + assigneeHandle: self.handle, + assigneePaneKey: self.paneKey, + processIncarnation: structuredWorkerProcessIncarnation(sessionId), + creator: self, + maxDepth: UNCAPPED + }) + const delegated = db.createDispatchContext({ + taskId: db.createTask({ runId: 'run_legacy_local', spec: 'delegated' }).id, + assigneeHandle: 'term_delegate', + assigneePaneKey: 'tab_delegate:22222222-2222-4222-8222-222222222222', + creator: self, + maxDepth: UNCAPPED + }) + backfillStructuredWorkerOrcaSessionIds(db.db) + + const ownRow = db.getDispatchContextById(own.id) + expect(ownRow?.creator_orca_session_id).toBe(sessionId) + expect(ownRow?.assignee_orca_session_id).toBe(ownRow?.creator_orca_session_id) + expect(db.resolveCreatorDepth(self)).toBe(0) + const delegatedRow = db.getDispatchContextById(delegated.id) + expect(delegatedRow?.creator_orca_session_id).toBe(sessionId) + expect(delegatedRow?.assignee_orca_session_id).toBeNull() + expect( + db.resolveCreatorDepth({ + kind: 'terminal', + handle: 'term_delegate', + paneKey: 'tab_delegate:22222222-2222-4222-8222-222222222222' + }) + ).toBe(1) + }) }) diff --git a/src/main/runtime/orchestration/db/orchestration-db.ts b/src/main/runtime/orchestration/db/orchestration-db.ts index 1ce52e96c4a5..7bd0d65517f4 100644 --- a/src/main/runtime/orchestration/db/orchestration-db.ts +++ b/src/main/runtime/orchestration/db/orchestration-db.ts @@ -9,6 +9,7 @@ import { } from './runs/run-coordinator-mail-routing' import { createTables } from './schema/create-tables' import { migrate } from './schema/migrate' +import { backfillStructuredWorkerOrcaSessionIds } from './schema/structured-worker-orca-session-backfill' class OrchestrationDbCore { db: Database.Database @@ -30,6 +31,7 @@ class OrchestrationDbCore { createTables.call(this as unknown as OrchestrationDb) migrate.call(this as unknown as OrchestrationDb) backfillFederatedStubHomeRuns(this.db) + backfillStructuredWorkerOrcaSessionIds(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/row-column-lists.ts b/src/main/runtime/orchestration/db/row-column-lists.ts index 368c97ede3c2..4ee3740c3834 100644 --- a/src/main/runtime/orchestration/db/row-column-lists.ts +++ b/src/main/runtime/orchestration/db/row-column-lists.ts @@ -13,6 +13,8 @@ export const RUN_COLUMNS = [ 'home_database', 'coordinator_handle', 'coordinator_pane_key', + 'coordinator_orca_session_id', + 'coordinator_orca_session_id_generation', 'consumer_generation', 'legacy', 'created_at', @@ -45,6 +47,7 @@ export const DISPATCH_CONTEXT_COLUMNS = [ 'launch_token_hash', 'assignee_handle', 'assignee_pane_key', + 'assignee_orca_session_id', 'capability_hash', 'process_incarnation', 'capability_revoked_at', @@ -52,6 +55,7 @@ export const DISPATCH_CONTEXT_COLUMNS = [ 'creator_dispatch_id', 'creator_handle', 'creator_pane_key', + 'creator_orca_session_id', 'host_scope', 'status', 'failure_count', diff --git a/src/main/runtime/orchestration/db/runs/run-binding.ts b/src/main/runtime/orchestration/db/runs/run-binding.ts index e2ff77ec8182..c5e6f34abfad 100644 --- a/src/main/runtime/orchestration/db/runs/run-binding.ts +++ b/src/main/runtime/orchestration/db/runs/run-binding.ts @@ -133,10 +133,12 @@ export function bindRun( this.setLegacyCompatibilityPrincipalStatus(coordinatorPrincipal.id, 'revoked') } } + // The Orca session id belongs to the coordinator being replaced; nothing here resolves the new one's. this.db .prepare( `UPDATE runs - SET coordinator_handle = ?, coordinator_pane_key = ?, + SET coordinator_handle = ?, coordinator_pane_key = ?, coordinator_orca_session_id = NULL, + coordinator_orca_session_id_generation = NULL, consumer_generation = consumer_generation + 1, updated_at = datetime('now') WHERE id = ?` diff --git a/src/main/runtime/orchestration/db/runs/run-coordinator-mail-routing.ts b/src/main/runtime/orchestration/db/runs/run-coordinator-mail-routing.ts index 121a051cba3c..07a6044d5aaa 100644 --- a/src/main/runtime/orchestration/db/runs/run-coordinator-mail-routing.ts +++ b/src/main/runtime/orchestration/db/runs/run-coordinator-mail-routing.ts @@ -1,4 +1,5 @@ import type { OrchestrationDb } from '../orchestration-db' +import { currentRunCoordinatorSessionAddressSql } from './run-coordinator-orca-session' export function rememberRunCoordinatorHandle( this: OrchestrationDb, @@ -12,11 +13,17 @@ export function rememberRunCoordinatorHandle( .run(runId, terminalHandle) } +const CURRENT_COORDINATOR_SESSION_ADDRESS_SQL = currentRunCoordinatorSessionAddressSql('runs') + +// Every address the coordinator has, its handle and its session address, as the migrate-v42 triggers. export function rememberCurrentRunCoordinatorHandles(this: OrchestrationDb): void { this.db.exec(` INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) SELECT id, coordinator_handle FROM runs - WHERE legacy = 0 AND coordinator_handle IS NOT NULL + WHERE legacy = 0 AND coordinator_handle IS NOT NULL; + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + SELECT id, ${CURRENT_COORDINATOR_SESSION_ADDRESS_SQL} FROM runs + WHERE legacy = 0 AND ${CURRENT_COORDINATOR_SESSION_ADDRESS_SQL} IS NOT NULL; `) } diff --git a/src/main/runtime/orchestration/db/runs/run-coordinator-orca-session.ts b/src/main/runtime/orchestration/db/runs/run-coordinator-orca-session.ts new file mode 100644 index 000000000000..2ed9d9809a29 --- /dev/null +++ b/src/main/runtime/orchestration/db/runs/run-coordinator-orca-session.ts @@ -0,0 +1,31 @@ +import { ORCA_SESSION_ADDRESS_PREFIX } from '../../../../../shared/orca-session-address' +import type { RunRow } from '../../types' + +type RunCoordinatorOrcaSessionFields = Pick< + RunRow, + 'coordinator_orca_session_id' | 'coordinator_orca_session_id_generation' | 'consumer_generation' +> + +/** + * A Run coordinator's Orca session id counts only at the `consumer_generation` it was written at. + * Every write that rebinds or unbinds a Run bumps that generation, including one from a binary that + * predates the column, so an id such a write leaves behind stops counting with nothing to clear it. + */ +export function currentRunCoordinatorOrcaSessionId( + run: RunCoordinatorOrcaSessionFields +): string | null { + return run.coordinator_orca_session_id_generation === run.consumer_generation + ? run.coordinator_orca_session_id + : null +} + +/** The same rule in SQL, for a `runs` row named `row` (a table name, alias, or `NEW`). */ +export function currentRunCoordinatorOrcaSessionIdSql(row: string): string { + return `(CASE WHEN ${row}.coordinator_orca_session_id_generation = ${row}.consumer_generation + THEN ${row}.coordinator_orca_session_id END)` +} + +/** The coordinator's `session:` address in SQL; NULL when it has no current Orca session id. */ +export function currentRunCoordinatorSessionAddressSql(row: string): string { + return `('${ORCA_SESSION_ADDRESS_PREFIX}' || ${currentRunCoordinatorOrcaSessionIdSql(row)})` +} diff --git a/src/main/runtime/orchestration/db/runs/run-lookup.ts b/src/main/runtime/orchestration/db/runs/run-lookup.ts index 194cceffcf0c..0a60b4e7f9d3 100644 --- a/src/main/runtime/orchestration/db/runs/run-lookup.ts +++ b/src/main/runtime/orchestration/db/runs/run-lookup.ts @@ -135,7 +135,8 @@ export function unbindOtherRunsForPane( this.db .prepare( `UPDATE runs - SET coordinator_handle = NULL, coordinator_pane_key = NULL, + SET coordinator_handle = NULL, coordinator_pane_key = NULL, coordinator_orca_session_id = NULL, + coordinator_orca_session_id_generation = NULL, consumer_generation = consumer_generation + 1, updated_at = datetime('now') WHERE id = ?` diff --git a/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts b/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts index f02e7e8c15ac..ecc99241e35b 100644 --- a/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts +++ b/src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts @@ -8,6 +8,12 @@ CREATE TABLE IF NOT EXISTS runs ( home_database TEXT NOT NULL DEFAULT 'this_database', coordinator_handle TEXT, coordinator_pane_key TEXT, + -- Bare Orca session id the coordinator is addressed by, when it has one (today only structured + -- sessions); for a /clear'd chat, its lineage root's. + coordinator_orca_session_id TEXT, + -- The consumer_generation coordinator_orca_session_id was written at; the id counts only while they + -- are equal (run-coordinator-orca-session). So bump consumer_generation for a rebind or unbind only. + coordinator_orca_session_id_generation INTEGER, consumer_generation INTEGER NOT NULL DEFAULT 0, legacy INTEGER NOT NULL DEFAULT 0, created_at TEXT NOT NULL DEFAULT (datetime('now')), @@ -55,6 +61,10 @@ CREATE TABLE IF NOT EXISTS run_coordinator_handles ( CREATE INDEX IF NOT EXISTS idx_run_coordinator_handles_handle ON run_coordinator_handles(terminal_handle, run_id); +-- Handle-only on purpose; migrate-v42 replaces both triggers with a form that also remembers the +-- coordinator's session address. This SQL runs before migrate on every open, so it +-- must compile against a pre-v42 runs table: a trigger naming coordinator_orca_session_id there makes +-- the next INSERT INTO runs fail to prepare mid-migration. CREATE TRIGGER IF NOT EXISTS trg_runs_remember_coordinator_insert AFTER INSERT ON runs WHEN NEW.legacy = 0 AND NEW.coordinator_handle IS NOT NULL 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 d674a5298e72..e0a7a0bef2d4 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 @@ -145,6 +145,9 @@ CREATE TABLE IF NOT EXISTS dispatch_contexts ( launch_token_hash TEXT, assignee_handle TEXT, assignee_pane_key TEXT, + -- Bare Orca session id the agent is addressed by, when it has one (today only structured + -- sessions); for a /clear'd chat, its lineage root's. Not its session: address. + assignee_orca_session_id TEXT, capability_hash TEXT, process_incarnation TEXT, capability_revoked_at TEXT, @@ -155,6 +158,8 @@ CREATE TABLE IF NOT EXISTS dispatch_contexts ( -- so it must not count as a nesting parent. Null on rows written before v37 and for Orca's loop. creator_handle TEXT, creator_pane_key TEXT, + -- Same form as assignee_orca_session_id: the id the creator is addressed by (a lineage root's). + creator_orca_session_id TEXT, host_scope TEXT, status TEXT NOT NULL DEFAULT 'pending' CHECK(status IN ('pending', 'dispatched', 'completed', 'failed', 'circuit_broken')), diff --git a/src/main/runtime/orchestration/db/schema/migrate-v42.ts b/src/main/runtime/orchestration/db/schema/migrate-v42.ts new file mode 100644 index 000000000000..e8d55037abc8 --- /dev/null +++ b/src/main/runtime/orchestration/db/schema/migrate-v42.ts @@ -0,0 +1,72 @@ +import type { OrchestrationDb } from '../orchestration-db' +import { currentRunCoordinatorSessionAddressSql } from '../runs/run-coordinator-orca-session' + +const ORCA_SESSION_ID_COLUMNS = [ + ['runs', 'coordinator_orca_session_id', 'TEXT'], + ['runs', 'coordinator_orca_session_id_generation', 'INTEGER'], + ['dispatch_contexts', 'assignee_orca_session_id', 'TEXT'], + ['dispatch_contexts', 'creator_orca_session_id', 'TEXT'] +] as const + +const NEW_COORDINATOR_SESSION_ADDRESS_SQL = currentRunCoordinatorSessionAddressSql('NEW') + +/** + * Orca session id columns (bare ids, see orca-session-address) on a Run's coordinator and a + * Dispatch's assignee and creator: the Orca session id the agent is addressed by, when it has one + * (today only structured sessions); for a `/clear`ed chat, its lineage root's, not the live one. + * Existing structured-worker rows get their id from `backfillStructuredWorkerOrcaSessionIds`, which + * runs after migrate on every open. A coordinator's id carries the consumer generation it was + * written at and counts only at that generation. + * + * Dev databases stamped v42 by earlier builds hold `*_principal` or `*_actor` columns instead. They + * are unsupported: the version-skew probe finds a column missing and replays the chain, which adds + * these columns and leaves the stale ones unread. + */ +export function migrateV42(this: OrchestrationDb, current: number): void { + if (current >= 42) { + return + } + // Guarded because createTables runs first on every open and already gives a fresh database these. + for (const [table, column, type] of ORCA_SESSION_ID_COLUMNS) { + if (!this.hasColumn(table, column)) { + this.db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${type}`) + } + } + // A Run lookup that ORs a pane-leaf match with an Orca session id match scans every Run without this index. + this.db.exec(` + CREATE INDEX IF NOT EXISTS idx_runs_coordinator_orca_session_id + ON runs(coordinator_orca_session_id) WHERE coordinator_orca_session_id IS NOT NULL; + CREATE INDEX IF NOT EXISTS idx_dispatch_assignee_orca_session_id + ON dispatch_contexts(assignee_orca_session_id) WHERE assignee_orca_session_id IS NOT NULL; + `) + // Every address the coordinator has is remembered, as bindRun remembers them: its handle and its + // session address, `session:`, so every reader of this cache matches either by string + // equality, unchanged, and neither address takes precedence. This step owns the trigger form: the + // static createTables SQL must stay handle-only (see create-core-tables-sql), and CREATE TRIGGER IF + // NOT EXISTS never replaces an existing database's triggers, so they are dropped and recreated by name. + this.db.exec(` + DROP TRIGGER IF EXISTS trg_runs_remember_coordinator_insert; + DROP TRIGGER IF EXISTS trg_runs_remember_coordinator_update; + CREATE TRIGGER trg_runs_remember_coordinator_insert + AFTER INSERT ON runs + WHEN NEW.legacy = 0 + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + SELECT NEW.id, NEW.coordinator_handle WHERE NEW.coordinator_handle IS NOT NULL; + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + SELECT NEW.id, ${NEW_COORDINATOR_SESSION_ADDRESS_SQL} + WHERE ${NEW_COORDINATOR_SESSION_ADDRESS_SQL} IS NOT NULL; + END; + CREATE TRIGGER trg_runs_remember_coordinator_update + AFTER UPDATE OF coordinator_handle, coordinator_orca_session_id, + coordinator_orca_session_id_generation ON runs + WHEN NEW.legacy = 0 + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + SELECT NEW.id, NEW.coordinator_handle WHERE NEW.coordinator_handle IS NOT NULL; + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + SELECT NEW.id, ${NEW_COORDINATOR_SESSION_ADDRESS_SQL} + WHERE ${NEW_COORDINATOR_SESSION_ADDRESS_SQL} IS NOT NULL; + END; + `) +} diff --git a/src/main/runtime/orchestration/db/schema/migrate.ts b/src/main/runtime/orchestration/db/schema/migrate.ts index bcf1b830ebc2..c02eeaf7cb38 100644 --- a/src/main/runtime/orchestration/db/schema/migrate.ts +++ b/src/main/runtime/orchestration/db/schema/migrate.ts @@ -12,6 +12,7 @@ import { migrateV38 } from './migrate-v38' import { migrateV39 } from './migrate-v39' import { migrateV40 } from './migrate-v40' import { DERIVED_DELIVERY_SCHEMA_SQL, migrateV41 } from './migrate-v41' +import { migrateV42 } from './migrate-v42' // 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 { @@ -38,6 +39,7 @@ export function migrate(this: OrchestrationDb): void { migrateV40.call(this, current) // Why: older steps recreate the unique index; v41 must run after them. migrateV41.call(this, current) + migrateV42.call(this, current) this.createMailboxDeliveryIndexesIfPossible() // Why: rebuild steps above RENAME the table, which SQLite refuses while a view names it. this.db.exec(DERIVED_DELIVERY_SCHEMA_SQL) diff --git a/src/main/runtime/orchestration/db/schema/structured-worker-orca-session-backfill.test.ts b/src/main/runtime/orchestration/db/schema/structured-worker-orca-session-backfill.test.ts new file mode 100644 index 000000000000..57bf53e7026a --- /dev/null +++ b/src/main/runtime/orchestration/db/schema/structured-worker-orca-session-backfill.test.ts @@ -0,0 +1,155 @@ +import { afterEach, describe, expect, it } from 'vitest' +import { + mintStructuredWorkerHandle, + mintStructuredWorkerPaneKey, + structuredWorkerProcessIncarnation +} from '../../../structured-worker-identity' +import { OrchestrationDb } from '../orchestration-db' +import { backfillStructuredWorkerOrcaSessionIds } from './structured-worker-orca-session-backfill' + +const SESSION_A = '0d2f4b6a-8c1e-4a3b-9d5f-7e0a2c4b6d81' +const SESSION_B = '1e3a5c7b-9d2f-4b4c-8e6a-0f1b3d5c7e92' +const SESSION_C = '2f4b6d8c-0e3a-4c5d-9f7b-1a2c4e6d8fa3' +const UNCAPPED = Number.MAX_SAFE_INTEGER +const SYSTEM = { kind: 'system' } as const + +describe('structured worker Orca session id backfill', () => { + let db: OrchestrationDb + + afterEach(() => db?.close()) + + function dispatch(params: { + handle: string + paneKey: string + incarnation?: string + creator?: { kind: 'terminal'; handle: string; paneKey: string } + }): string { + const task = db.createTask({ runId: 'run_legacy_local', spec: `work for ${params.handle}` }) + return db.createDispatchContext({ + taskId: task.id, + assigneeHandle: params.handle, + assigneePaneKey: params.paneKey, + processIncarnation: params.incarnation, + creator: params.creator ?? SYSTEM, + maxDepth: UNCAPPED + }).id + } + + function orcaSessionIds(dispatchId: string): { assignee: string | null; creator: string | null } { + const row = db.getDispatchContextById(dispatchId) + return { + assignee: row?.assignee_orca_session_id ?? null, + creator: row?.creator_orca_session_id ?? null + } + } + + it('proves a handle through the session this host recorded against it', () => { + db = new OrchestrationDb(':memory:') + const handle = mintStructuredWorkerHandle() + const pane = mintStructuredWorkerPaneKey(SESSION_B) + // The worker's own row carries no incarnation; its terminal resource row does. + const handleOnly = dispatch({ handle, paneKey: pane }) + db.createWorkerTerminalResourceStatement({ + dispatchId: handleOnly, + worktreeId: 'wt_1', + terminalHandle: handle, + paneKey: pane, + processIncarnation: structuredWorkerProcessIncarnation(SESSION_B), + hostScope: JSON.stringify({ kind: 'local', hostId: 'local' }), + ownership: 'owned' + }) + + backfillStructuredWorkerOrcaSessionIds(db.db) + + expect(orcaSessionIds(handleOnly)).toEqual({ assignee: SESSION_B, creator: null }) + }) + + it('leaves every row it cannot tie to exactly one valid session NULL', () => { + db = new OrchestrationDb(':memory:') + const unrecorded = mintStructuredWorkerHandle() + const noRecord = dispatch({ + handle: unrecorded, + paneKey: mintStructuredWorkerPaneKey(SESSION_A) + }) + const invalidId = dispatch({ + handle: mintStructuredWorkerHandle(), + paneKey: 'tab_x:44444444-4444-4444-8444-444444444444', + incarnation: 'structured:not a session id' + }) + // A terminal in a structured session's tab: the pane key names a session, the process does not. + const terminalInSessionTab = dispatch({ + handle: 'term_tui', + paneKey: mintStructuredWorkerPaneKey(SESSION_C), + incarnation: 'pty_proc_1a2b:777' + }) + const shared = mintStructuredWorkerHandle() + const conflicting = dispatch({ + handle: shared, + paneKey: mintStructuredWorkerPaneKey(SESSION_A), + incarnation: structuredWorkerProcessIncarnation(SESSION_A) + }) + db.createWorkerTerminalResourceStatement({ + dispatchId: conflicting, + worktreeId: 'wt_1', + terminalHandle: shared, + paneKey: null, + processIncarnation: structuredWorkerProcessIncarnation(SESSION_B), + ownership: 'owned' + }) + + backfillStructuredWorkerOrcaSessionIds(db.db) + + for (const id of [noRecord, invalidId, terminalInSessionTab, conflicting]) { + expect(orcaSessionIds(id), id).toEqual({ assignee: null, creator: null }) + } + }) + + it('fills only NULLs and never rewrites an Orca session id a writer recorded', () => { + db = new OrchestrationDb(':memory:') + const handle = mintStructuredWorkerHandle() + const id = dispatch({ + handle, + paneKey: mintStructuredWorkerPaneKey(SESSION_A), + incarnation: structuredWorkerProcessIncarnation(SESSION_A) + }) + db.db + .prepare('UPDATE dispatch_contexts SET assignee_orca_session_id = ? WHERE id = ?') + .run(SESSION_C, id) + + backfillStructuredWorkerOrcaSessionIds(db.db) + + expect(orcaSessionIds(id).assignee).toBe(SESSION_C) + }) + + it("fills a worker-coordinated Run over an Orca session id an older binding's generation left behind", () => { + db = new OrchestrationDb(':memory:') + const runFor = (sessionId: string): string => { + const handle = mintStructuredWorkerHandle() + const paneKey = mintStructuredWorkerPaneKey(sessionId) + dispatch({ handle, paneKey, incarnation: structuredWorkerProcessIncarnation(sessionId) }) + return db.createRun({ + objective: sessionId, + coordinatorHandle: handle, + coordinatorPaneKey: paneKey + }).id + } + const stale = runFor(SESSION_A) + const recorded = runFor(SESSION_B) + const setOrcaSessionId = db.db.prepare( + `UPDATE runs SET coordinator_orca_session_id = ?, + coordinator_orca_session_id_generation = consumer_generation - ? + WHERE id = ?` + ) + // An older binary rebound this Run to the worker over session C's id, which it cannot see. + setOrcaSessionId.run(SESSION_C, 1, stale) + // A writer recorded this one at the current generation. + setOrcaSessionId.run(SESSION_C, 0, recorded) + + backfillStructuredWorkerOrcaSessionIds(db.db) + + const filled = db.getRunRaw(stale) + expect(filled?.coordinator_orca_session_id).toBe(SESSION_A) + expect(filled?.coordinator_orca_session_id_generation).toBe(filled?.consumer_generation) + expect(db.getRunRaw(recorded)?.coordinator_orca_session_id).toBe(SESSION_C) + }) +}) diff --git a/src/main/runtime/orchestration/db/schema/structured-worker-orca-session-backfill.ts b/src/main/runtime/orchestration/db/schema/structured-worker-orca-session-backfill.ts new file mode 100644 index 000000000000..59a0bdb1c7f7 --- /dev/null +++ b/src/main/runtime/orchestration/db/schema/structured-worker-orca-session-backfill.ts @@ -0,0 +1,135 @@ +import type Database from '../../../../sqlite/sync-database' +import { isOrcaSessionId } from '../../../../../shared/orca-session-address' +import { + STRUCTURED_WORKER_HANDLE_PREFIX, + STRUCTURED_WORKER_INCARNATION_PREFIX, + isStructuredWorkerHandle, + sessionIdFromStructuredWorkerIncarnation +} from '../../../structured-worker-identity' +import { currentRunCoordinatorOrcaSessionIdSql } from '../runs/run-coordinator-orca-session' + +const CURRENT_COORDINATOR_ORCA_SESSION_ID_SQL = currentRunCoordinatorOrcaSessionIdSql('runs') + +// GLOB is a case-sensitive prefix filter; the canonical predicates still decide every row. +const HANDLE_GLOB = `${STRUCTURED_WORKER_HANDLE_PREFIX}*` +const INCARNATION_GLOB = `${STRUCTURED_WORKER_INCARNATION_PREFIX}*` + +const RECORDED_WORKER_SESSIONS_SQL = ` + SELECT assignee_handle AS handle, process_incarnation AS incarnation FROM dispatch_contexts + WHERE assignee_handle GLOB ? AND process_incarnation GLOB ? + UNION + SELECT terminal_handle, process_incarnation FROM worker_terminal_resources + WHERE terminal_handle GLOB ? AND process_incarnation GLOB ?` + +/** + * Fills the Orca session id on rows that provably belong to a structured worker session and have + * none: a `structured:` process incarnation, or a `structworker_` handle this host + * recorded against such an incarnation. Every other row stays NULL, PTY rows included, and evidence + * naming more than one session proves none. Pane keys are never read: a pane outlives the agent in it. + * + * Both markers are minted only for a local, non-WSL session (`structuredWorkerHostScope`), so the + * rows carrying them were written by this host. + * + * Runs after migrate on every open, not only once at v42: a binary rolled back past v42 keeps + * writing structured-worker rows without an Orca session id after user_version is already 42. It + * fills only rows with no id that counts, so an id a writer recorded is never rewritten. + */ +export function backfillStructuredWorkerOrcaSessionIds(db: Database.Database): void { + let recordedSessions: Map> | undefined + const orcaSessionIdFor = (handle: unknown, incarnation: unknown): string | null => { + if (incarnation != null && typeof incarnation !== 'string') { + return null + } + const sessions = new Set() + if (incarnation != null) { + const sessionId = sessionIdFromStructuredWorkerIncarnation(incarnation) + if (!sessionId) { + // A live non-structured incarnation says this row is some other process. + return null + } + sessions.add(sessionId) + } + if (typeof handle === 'string' && isStructuredWorkerHandle(handle)) { + recordedSessions ??= recordedWorkerSessionsByHandle(db) + for (const sessionId of recordedSessions.get(handle) ?? []) { + sessions.add(sessionId) + } + } + const [sessionId, ...others] = sessions + return sessionId && others.length === 0 && isOrcaSessionId(sessionId) ? sessionId : null + } + + const assignees = db + .prepare( + `SELECT id, assignee_handle, process_incarnation FROM dispatch_contexts + WHERE assignee_orca_session_id IS NULL + AND (process_incarnation GLOB ? OR assignee_handle GLOB ?)` + ) + .all(INCARNATION_GLOB, HANDLE_GLOB) + const setAssignee = db.prepare( + `UPDATE dispatch_contexts SET assignee_orca_session_id = ? + WHERE id = ? AND assignee_orca_session_id IS NULL` + ) + for (const row of assignees) { + const orcaSessionId = orcaSessionIdFor(row.assignee_handle, row.process_incarnation) + if (orcaSessionId && typeof row.id === 'string') { + setAssignee.run(orcaSessionId, row.id) + } + } + + const creators = db + .prepare( + `SELECT id, creator_handle FROM dispatch_contexts + WHERE creator_orca_session_id IS NULL AND creator_handle GLOB ?` + ) + .all(HANDLE_GLOB) + const setCreator = db.prepare( + `UPDATE dispatch_contexts SET creator_orca_session_id = ? + WHERE id = ? AND creator_orca_session_id IS NULL` + ) + for (const row of creators) { + const orcaSessionId = orcaSessionIdFor(row.creator_handle, null) + if (orcaSessionId && typeof row.id === 'string') { + setCreator.run(orcaSessionId, row.id) + } + } + + // A coordinator id left at an older generation counts as none, so the handle's session fills it. + const coordinators = db + .prepare( + `SELECT id, coordinator_handle FROM runs + WHERE ${CURRENT_COORDINATOR_ORCA_SESSION_ID_SQL} IS NULL AND coordinator_handle GLOB ?` + ) + .all(HANDLE_GLOB) + const setCoordinator = db.prepare( + `UPDATE runs SET coordinator_orca_session_id = ?, + coordinator_orca_session_id_generation = consumer_generation + WHERE id = ? AND ${CURRENT_COORDINATOR_ORCA_SESSION_ID_SQL} IS NULL` + ) + for (const row of coordinators) { + const orcaSessionId = orcaSessionIdFor(row.coordinator_handle, null) + if (orcaSessionId && typeof row.id === 'string') { + setCoordinator.run(orcaSessionId, row.id) + } + } +} + +function recordedWorkerSessionsByHandle(db: Database.Database): Map> { + const sessionsByHandle = new Map>() + const rows = db + .prepare(RECORDED_WORKER_SESSIONS_SQL) + .all(HANDLE_GLOB, INCARNATION_GLOB, HANDLE_GLOB, INCARNATION_GLOB) + for (const row of rows) { + const sessionId = + typeof row.incarnation === 'string' + ? sessionIdFromStructuredWorkerIncarnation(row.incarnation) + : null + if (typeof row.handle !== 'string' || !sessionId) { + continue + } + const sessions = sessionsByHandle.get(row.handle) ?? new Set() + sessions.add(sessionId) + sessionsByHandle.set(row.handle, sessions) + } + return sessionsByHandle +} diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-assignee-orca-session.test.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-assignee-orca-session.test.ts new file mode 100644 index 000000000000..b43c996742f8 --- /dev/null +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-dispatch-assignee-orca-session.test.ts @@ -0,0 +1,76 @@ +import { afterEach, describe, expect, it } from 'vitest' +import { isOrcaSessionId } from '../../../../../shared/orca-session-address' +import { mintStructuredWorkerHandle } from '../../../structured-worker-identity' +import { OrchestrationDb } from '../orchestration-db' + +const EARLIER_ORCA_SESSION_ID = '7d9f1b3e-5a2c-4e6b-8f0a-1c3e5a7b9d42' +const WORKER_PANE = 'tab_worker:88888888-8888-4888-8888-888888888888' + +describe('assignee identity writers', () => { + let db: OrchestrationDb | undefined + + afterEach(() => { + db?.close() + db = undefined + }) + + /** A starting Dispatch whose row already names an Orca session id, standing in for any earlier writer. */ + function startingDispatchWithOrcaSessionId(target: OrchestrationDb): string { + const task = target.createTask({ runId: 'run_legacy_local', spec: 'worker' }) + const started = target.createStartingWorkerDispatch({ + creator: { kind: 'system' }, + maxDepth: Number.MAX_SAFE_INTEGER, + taskId: task.id, + startOptions: {} + }) + target.db + .prepare('UPDATE dispatch_contexts SET assignee_orca_session_id = ? WHERE id = ?') + .run(EARLIER_ORCA_SESSION_ID, started.dispatch.id) + return started.dispatch.id + } + + it('clears the Orca session id when worker authority names the assignee', () => { + db = new OrchestrationDb(':memory:') + const dispatchId = startingDispatchWithOrcaSessionId(db) + + db.prepareStartingWorkerAuthority({ + dispatchId, + handle: 'term_worker', + paneKey: WORKER_PANE, + processIncarnation: 'pty_proc_9c1d:31', + worktreeId: 'wt_1', + setupState: 'not_applicable', + effects: [] + }) + + expect(db.getDispatchContextById(dispatchId)).toMatchObject({ + assignee_handle: 'term_worker', + assignee_orca_session_id: null + }) + }) + + it('clears the Orca session id when a failed start records the terminal it owned', () => { + db = new OrchestrationDb(':memory:') + const dispatchId = startingDispatchWithOrcaSessionId(db) + db.recordCreatedWorkerTerminalCustody({ + dispatchId, + handle: 'term_worker', + paneKey: WORKER_PANE, + processIncarnation: 'pty_proc_9c1d:31', + worktreeId: 'wt_1' + }) + db.recordWorkerStage({ dispatchId, stage: 'agent_readiness', terminalHandle: 'term_worker' }) + + db.failWorkerStart(dispatchId, 'agent_readiness', 'agent never became ready') + + expect(db.getDispatchContextById(dispatchId)).toMatchObject({ + assignee_handle: 'term_worker', + assignee_orca_session_id: null + }) + }) + + it('refuses a minted structured-worker handle as an Orca session id', () => { + // Ties the handle-prefix refusal to the handle this runtime actually mints. + expect(isOrcaSessionId(mintStructuredWorkerHandle())).toBe(false) + }) +}) 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 b33a8798dc3e..166ccd7193e6 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 @@ -54,7 +54,7 @@ export function prepareStartingWorkerAuthority( .prepare( `UPDATE dispatch_contexts SET assignee_handle = ?, assignee_pane_key = ?, process_incarnation = ?, - host_scope = ?, + assignee_orca_session_id = NULL, host_scope = ?, capability_hash = ?, launch_token_hash = COALESCE(launch_token_hash, ?), capability_revoked_at = NULL, consumer_generation = consumer_generation + 1 diff --git a/src/main/runtime/orchestration/db/worker-terminal/failed-start-dispatch-identity.ts b/src/main/runtime/orchestration/db/worker-terminal/failed-start-dispatch-identity.ts index 671b12c39a71..c9b9b5d0afdf 100644 --- a/src/main/runtime/orchestration/db/worker-terminal/failed-start-dispatch-identity.ts +++ b/src/main/runtime/orchestration/db/worker-terminal/failed-start-dispatch-identity.ts @@ -21,7 +21,8 @@ export function recordFailedStartDispatchIdentity( db.db .prepare( `UPDATE dispatch_contexts - SET assignee_handle = ?, assignee_pane_key = ?, process_incarnation = ?, host_scope = ? + SET assignee_handle = ?, assignee_pane_key = ?, process_incarnation = ?, host_scope = ?, + assignee_orca_session_id = NULL WHERE id = ? AND status = 'failed' AND capability_hash IS NULL` ) .run( diff --git a/src/main/runtime/orchestration/orchestration-orca-session-column-migration.test.ts b/src/main/runtime/orchestration/orchestration-orca-session-column-migration.test.ts new file mode 100644 index 000000000000..69186334828b --- /dev/null +++ b/src/main/runtime/orchestration/orchestration-orca-session-column-migration.test.ts @@ -0,0 +1,674 @@ +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 { + mintStructuredWorkerHandle, + mintStructuredWorkerPaneKey, + structuredWorkerProcessIncarnation +} from '../structured-worker-identity' +import { OrchestrationDb } from './db' +import { SCHEMA_VERSION } from './db/contract-constants' +import { formatOrcaSessionAddress } from '../../../shared/orca-session-address' +import { RUN_PANE_KEY_MATCH_SUFFIX_SQL } from './db/pane-key-match' +import { + currentRunCoordinatorOrcaSessionId, + currentRunCoordinatorOrcaSessionIdSql +} from './db/runs/run-coordinator-orca-session' +import { resolveOrchestrationMigrationStartVersion } from './orchestration-schema-version-skew' + +const SESSION_ID = '5f0c1d9e-2b7a-4c3e-8f61-0a9d2e7b4c11' +const CHAT_SESSION_ID = '9a4e7c1b-3d2f-4b6a-8e5c-7f1d0b2a6c93' +const CHAT_SESSION_ADDRESS = formatOrcaSessionAddress(CHAT_SESSION_ID) +const ORCA_SESSION_ID_COLUMNS = [ + 'assignee_orca_session_id', + 'coordinator_orca_session_id', + 'coordinator_orca_session_id_generation', + 'creator_orca_session_id' +] +const COORDINATOR_PANE = 'tab_coord:11111111-1111-4111-8111-111111111111' +const PTY_WORKER_PANE = 'tab_pty:22222222-2222-4222-8222-222222222222' +const NESTED_PANE = 'tab_nested:33333333-3333-4333-8333-333333333333' +const UNCAPPED = Number.MAX_SAFE_INTEGER + +// The coordinator-address triggers exactly as main stamped them at v41 and before. +const HANDLE_ONLY_COORDINATOR_TRIGGERS_SQL = ` + CREATE TRIGGER trg_runs_remember_coordinator_insert + AFTER INSERT ON runs + WHEN NEW.legacy = 0 AND NEW.coordinator_handle IS NOT NULL + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + VALUES (NEW.id, NEW.coordinator_handle); + END; + CREATE TRIGGER trg_runs_remember_coordinator_update + AFTER UPDATE OF coordinator_handle ON runs + WHEN NEW.legacy = 0 AND NEW.coordinator_handle IS NOT NULL + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + VALUES (NEW.id, NEW.coordinator_handle); + END;` + +const V41_RUN_COLUMNS = + 'id, objective, home_database, coordinator_handle, coordinator_pane_key, consumer_generation, legacy, created_at, updated_at' + +type SeededRows = { + ptyRunId: string + structuredRunId: string + structuredDispatchId: string + ptyDispatchId: string + nestedDispatchId: string + workerHandle: string +} + +/** One PTY coordinator with a structured worker and a PTY worker; the structured worker runs a nested Run. */ +function seedStructuredAndPtyRows(db: OrchestrationDb): SeededRows { + const workerHandle = mintStructuredWorkerHandle() + const workerPane = mintStructuredWorkerPaneKey(SESSION_ID) + const coordinator = { kind: 'terminal', handle: 'term_coord', paneKey: COORDINATOR_PANE } as const + const ptyRun = db.createRun({ + objective: 'pty coordinator', + coordinatorHandle: 'term_coord', + coordinatorPaneKey: COORDINATOR_PANE + }) + const structuredDispatch = db.createDispatchContext({ + taskId: db.createTask({ runId: ptyRun.id, spec: 'structured worker' }).id, + assigneeHandle: workerHandle, + assigneePaneKey: workerPane, + processIncarnation: structuredWorkerProcessIncarnation(SESSION_ID), + creator: coordinator, + maxDepth: UNCAPPED + }) + const ptyDispatch = db.createDispatchContext({ + taskId: db.createTask({ runId: ptyRun.id, spec: 'pty worker' }).id, + assigneeHandle: 'term_pty_worker', + assigneePaneKey: PTY_WORKER_PANE, + processIncarnation: 'pty_proc_7f3a:4242', + creator: coordinator, + maxDepth: UNCAPPED + }) + const structuredRun = db.createRun({ + objective: 'structured worker coordinates a nested run', + coordinatorHandle: workerHandle, + coordinatorPaneKey: workerPane + }) + const nestedDispatch = db.createDispatchContext({ + taskId: db.createTask({ runId: structuredRun.id, spec: 'nested' }).id, + assigneeHandle: 'term_nested', + assigneePaneKey: NESTED_PANE, + creator: { kind: 'terminal', handle: workerHandle, paneKey: workerPane }, + maxDepth: UNCAPPED + }) + return { + ptyRunId: ptyRun.id, + structuredRunId: structuredRun.id, + structuredDispatchId: structuredDispatch.id, + ptyDispatchId: ptyDispatch.id, + nestedDispatchId: nestedDispatch.id, + workerHandle + } +} + +/** Strips a current database back to the shape main stamps at v41 and stamps `version`. */ +function stripOrcaSessionSchema(path: string, version: number): void { + const raw = new Database(path) + raw.exec(` + DROP INDEX idx_runs_coordinator_orca_session_id; + DROP INDEX idx_dispatch_assignee_orca_session_id; + DROP TRIGGER trg_runs_remember_coordinator_insert; + DROP TRIGGER trg_runs_remember_coordinator_update; + ALTER TABLE runs DROP COLUMN coordinator_orca_session_id; + ALTER TABLE runs DROP COLUMN coordinator_orca_session_id_generation; + ALTER TABLE dispatch_contexts DROP COLUMN assignee_orca_session_id; + ALTER TABLE dispatch_contexts DROP COLUMN creator_orca_session_id; + ${HANDLE_ONLY_COORDINATOR_TRIGGERS_SQL} + `) + raw.pragma(`user_version = ${version}`) + raw.close() +} + +/** The plan for finding a caller's Runs by pane leaf or by coordinator Orca session id in one statement. */ +function coordinatorLookupPlan(db: Database.Database): string { + return db + .prepare( + `EXPLAIN QUERY PLAN SELECT id FROM runs + WHERE legacy = 0 AND ( + (coordinator_pane_key IS NOT NULL AND ${RUN_PANE_KEY_MATCH_SUFFIX_SQL} = ?) + OR coordinator_orca_session_id = ? + )` + ) + .all('leaf', CHAT_SESSION_ID) + .map((row) => String(row.detail)) + .join(' | ') +} + +function coordinatorTriggerSql(db: Database.Database): string[] { + return db + .prepare( + `SELECT sql FROM sqlite_master WHERE type = 'trigger' + AND name IN ('trg_runs_remember_coordinator_insert', 'trg_runs_remember_coordinator_update') + ORDER BY name` + ) + .all() + .map((row) => String(row.sql)) +} + +function orcaSessionColumns(db: Database.Database): string[] { + return db + .prepare( + `SELECT name FROM pragma_table_info('runs') + UNION ALL SELECT name FROM pragma_table_info('dispatch_contexts')` + ) + .all() + .map((column) => String(column.name)) + .filter((name) => name.includes('_orca_session_id')) + .sort() +} + +function coordinatorAddresses(db: Database.Database, runIds: string[]): string[] { + return db + .prepare( + `SELECT run_id, terminal_handle FROM run_coordinator_handles + WHERE run_id IN (${runIds.map(() => '?').join(', ')})` + ) + .all(...runIds) + .map((row) => `${String(row.run_id)} ${String(row.terminal_handle)}`) + .sort() +} + +describe('orchestration Orca session id column migration', () => { + const tempRoots: string[] = [] + + afterEach(() => { + for (const root of tempRoots.splice(0)) { + rmSync(root, { recursive: true, force: true }) + } + }) + + function tempDbPath(): string { + const root = mkdtempSync(join(tmpdir(), 'orca-session-column-migration-')) + tempRoots.push(root) + return join(root, 'orchestration.db') + } + + it('starts a v41 database at v41 and gives exactly its structured-worker rows an Orca session id', () => { + const path = tempDbPath() + const seed = new OrchestrationDb(path) + const rows = seedStructuredAndPtyRows(seed) + seed.close() + stripOrcaSessionSchema(path, 41) + + const probe = new Database(path) + try { + // Why: a v42 skew entry registered under v41 makes this 6 and replays the whole chain. + expect(resolveOrchestrationMigrationStartVersion(probe, 41, SCHEMA_VERSION)).toBe(41) + } finally { + probe.close() + } + + const db = new OrchestrationDb(path) + try { + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + expect(db.getDispatchContextById(rows.structuredDispatchId)).toMatchObject({ + assignee_orca_session_id: SESSION_ID, + creator_orca_session_id: null + }) + expect(db.getDispatchContextById(rows.ptyDispatchId)).toMatchObject({ + assignee_orca_session_id: null, + creator_orca_session_id: null + }) + expect(db.getDispatchContextById(rows.nestedDispatchId)).toMatchObject({ + assignee_orca_session_id: null, + creator_orca_session_id: SESSION_ID + }) + expect( + db.db + .prepare('SELECT id FROM dispatch_contexts WHERE assignee_orca_session_id IS NOT NULL') + .all() + ).toEqual([{ id: rows.structuredDispatchId }]) + expect(db.getRunRaw(rows.ptyRunId)?.coordinator_orca_session_id).toBeNull() + expect(db.getRunRaw(rows.structuredRunId)?.coordinator_orca_session_id).toBe(SESSION_ID) + // Every address a coordinator has: the PTY one its handle, the structured worker its handle + // and the session address its backfilled id gives it. + expect(coordinatorAddresses(db.db, [rows.ptyRunId, rows.structuredRunId])).toEqual( + [ + `${rows.ptyRunId} term_coord`, + `${rows.structuredRunId} ${rows.workerHandle}`, + `${rows.structuredRunId} ${formatOrcaSessionAddress(SESSION_ID)}` + ].sort() + ) + // CREATE TRIGGER IF NOT EXISTS alone would have kept the handle-only form here. + for (const sql of coordinatorTriggerSql(db.db)) { + expect(sql).toContain( + 'NEW.coordinator_orca_session_id_generation = NEW.consumer_generation' + ) + } + } finally { + db.close() + } + }) + + it('starts a v40 database at v40 and runs v41 before v42', () => { + const path = tempDbPath() + const seed = new OrchestrationDb(path) + const rows = seedStructuredAndPtyRows(seed) + seed.close() + stripOrcaSessionSchema(path, 40) + const raw = new Database(path) + // The v40 shape of the one object v41 changed: a unique outstanding-delivery index. + raw.exec(` + DROP TRIGGER trg_deliveries_one_outstanding; + DROP VIEW outstanding_deliveries; + DROP INDEX idx_deliveries_one_outstanding; + CREATE UNIQUE INDEX idx_deliveries_one_outstanding + ON deliveries(mailbox_handle) WHERE status = 'outstanding' AND mailbox_handle != ''; + `) + try { + expect(resolveOrchestrationMigrationStartVersion(raw, 40, SCHEMA_VERSION)).toBe(40) + } finally { + raw.close() + } + + const db = new OrchestrationDb(path) + try { + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + expect(orcaSessionColumns(db.db)).toEqual(ORCA_SESSION_ID_COLUMNS) + const index = db.db + .prepare("SELECT sql FROM sqlite_master WHERE name = 'idx_deliveries_one_outstanding'") + .get() + expect(String(index?.sql)).not.toContain('UNIQUE') + expect(db.getDispatchContextById(rows.structuredDispatchId)?.assignee_orca_session_id).toBe( + SESSION_ID + ) + expect(db.getDispatchContextById(rows.ptyDispatchId)?.assignee_orca_session_id).toBeNull() + } finally { + db.close() + } + }) + + it('upgrades a database from before the coordinator cache through the static triggers', () => { + const path = tempDbPath() + const seed = new OrchestrationDb(path) + const rows = seedStructuredAndPtyRows(seed) + seed.close() + stripOrcaSessionSchema(path, 27) + const raw = new Database(path) + raw.exec(` + DROP TRIGGER trg_runs_remember_coordinator_insert; + DROP TRIGGER trg_runs_remember_coordinator_update; + DROP TABLE run_coordinator_handles; + `) + raw.close() + + // createTables installs its static triggers before this chain's v40 step inserts into runs, so + // a static form naming coordinator_orca_session_id would fail to prepare here. + const db = new OrchestrationDb(path) + try { + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + for (const sql of coordinatorTriggerSql(db.db)) { + expect(sql).toContain( + 'NEW.coordinator_orca_session_id_generation = NEW.consumer_generation' + ) + } + expect(db.getRunRaw(rows.structuredRunId)?.coordinator_orca_session_id).toBe(SESSION_ID) + expect(coordinatorAddresses(db.db, [rows.ptyRunId])).toEqual([`${rows.ptyRunId} term_coord`]) + } finally { + db.close() + } + }) + + it('lets a v41 binary read and write a v42 database with Orca session ids in it', () => { + const path = tempDbPath() + const seed = new OrchestrationDb(path) + const rows = seedStructuredAndPtyRows(seed) + seed.close() + const upgraded = new OrchestrationDb(path) + expect(upgraded.getRunRaw(rows.structuredRunId)?.coordinator_orca_session_id).toBe(SESSION_ID) + upgraded.db + .prepare( + `INSERT INTO runs ( + id, objective, coordinator_orca_session_id, coordinator_orca_session_id_generation, consumer_generation, legacy + ) VALUES ('run_session', 'session coordinator', ?, 1, 1, 0)` + ) + .run(CHAT_SESSION_ID) + upgraded.close() + + // A raw connection stands in for the v41 binary; each statement below is v41's own SQL. + const v41 = new Database(path) + try { + // v41's migrate returns early on a newer stamp, so nothing rewrites the v42 objects. + expect(resolveOrchestrationMigrationStartVersion(v41, SCHEMA_VERSION, 41)).toBe( + SCHEMA_VERSION + ) + expect( + v41.prepare(`SELECT ${V41_RUN_COLUMNS} FROM runs WHERE id = ?`).get('run_session') + ).toMatchObject({ id: 'run_session', coordinator_handle: null, coordinator_pane_key: null }) + expect( + v41.prepare(`SELECT ${V41_RUN_COLUMNS} FROM runs WHERE id = ?`).get(rows.ptyRunId) + ).toMatchObject({ coordinator_handle: 'term_coord', coordinator_pane_key: COORDINATOR_PANE }) + // v41's listRuns reads `SELECT *`; the extra column rides along and every v41 column is intact. + const listed = v41.prepare('SELECT * FROM runs ORDER BY created_at DESC, id DESC').all() + expect(listed.map((run) => run.id).sort()).toEqual( + ['run_legacy_local', 'run_session', rows.ptyRunId, rows.structuredRunId].sort() + ) + + // v41's createTables runs these on every open. + v41.exec( + HANDLE_ONLY_COORDINATOR_TRIGGERS_SQL.replaceAll( + 'CREATE TRIGGER', + 'CREATE TRIGGER IF NOT EXISTS' + ) + ) + v41 + .prepare( + `INSERT INTO runs (id, objective, coordinator_handle, coordinator_pane_key, consumer_generation, legacy) + VALUES ('run_v41', 'written by v41', 'term_v41', 'tab_v41:44444444-4444-4444-8444-444444444444', 1, 0)` + ) + .run() + // An older binary rebinding a structured-coordinated Run cannot clear an id it cannot see. + v41 + .prepare( + `UPDATE runs SET coordinator_handle = ?, coordinator_pane_key = ?, + consumer_generation = consumer_generation + 1, updated_at = datetime('now') + WHERE id = ?` + ) + .run('term_taker', PTY_WORKER_PANE, rows.structuredRunId) + v41.exec(`INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + SELECT id, coordinator_handle FROM runs WHERE legacy = 0 AND coordinator_handle IS NOT NULL`) + + const readOrcaSessionId = v41.prepare( + 'SELECT coordinator_orca_session_id FROM runs WHERE id = ?' + ) + expect(readOrcaSessionId.get('run_v41')).toEqual({ coordinator_orca_session_id: null }) + expect(readOrcaSessionId.get(rows.structuredRunId)).toEqual({ + coordinator_orca_session_id: SESSION_ID + }) + // The session address was remembered by the v42 open, before v41's rebind made the id stale. + expect(coordinatorAddresses(v41, ['run_v41', rows.structuredRunId])).toEqual( + [ + `${rows.structuredRunId} ${rows.workerHandle}`, + `${rows.structuredRunId} ${formatOrcaSessionAddress(SESSION_ID)}`, + `${rows.structuredRunId} term_taker`, + 'run_v41 term_v41' + ].sort() + ) + // The v42 trigger form survives v41's IF NOT EXISTS create. + for (const sql of coordinatorTriggerSql(v41)) { + expect(sql).toContain('coordinator_orca_session_id') + } + } finally { + v41.close() + } + + const rolledForward = new OrchestrationDb(path) + try { + expect(rolledForward.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + expect(rolledForward.getRun('run_v41')?.coordinator_handle).toBe('term_v41') + expect(rolledForward.getRunMailboxOwnerIdsForHandle(CHAT_SESSION_ADDRESS)).toEqual([ + 'run_session' + ]) + // v41's rebind bumped the generation, so the id it could not clear no longer counts. + const rebound = rolledForward.getRunRaw(rows.structuredRunId) + expect(rebound?.coordinator_orca_session_id).toBe(SESSION_ID) + expect(rebound && currentRunCoordinatorOrcaSessionId(rebound)).toBeNull() + } finally { + rolledForward.close() + } + }) + + it('fills structured-worker rows written after the stamp reached v42 on the next open', () => { + const path = tempDbPath() + const first = new OrchestrationDb(path) + // No writer records an Orca session id yet, which is also the shape a binary rolled back past v42 writes. + const rows = seedStructuredAndPtyRows(first) + expect( + first.getDispatchContextById(rows.structuredDispatchId)?.assignee_orca_session_id + ).toBeNull() + first.close() + + const reopened = new OrchestrationDb(path) + try { + expect( + reopened.getDispatchContextById(rows.structuredDispatchId)?.assignee_orca_session_id + ).toBe(SESSION_ID) + expect(reopened.getDispatchContextById(rows.nestedDispatchId)?.creator_orca_session_id).toBe( + SESSION_ID + ) + expect(reopened.getRunRaw(rows.structuredRunId)?.coordinator_orca_session_id).toBe(SESSION_ID) + expect( + reopened.getDispatchContextById(rows.ptyDispatchId)?.assignee_orca_session_id + ).toBeNull() + } finally { + reopened.close() + } + }) + + it('drives a dev database stamped v42 with prototype principal columns to add the Orca session ids', () => { + const path = tempDbPath() + const seed = new OrchestrationDb(path) + const rows = seedStructuredAndPtyRows(seed) + seed.close() + stripOrcaSessionSchema(path, 41) + const raw = new Database(path) + // An unmerged prototype stamped v42 with differently named columns and triggers over them. + raw.exec(` + ALTER TABLE runs ADD COLUMN coordinator_principal TEXT; + ALTER TABLE dispatch_contexts ADD COLUMN assignee_principal TEXT; + ALTER TABLE dispatch_contexts ADD COLUMN creator_principal TEXT; + ALTER TABLE worker_terminal_resources ADD COLUMN principal TEXT; + DROP TRIGGER trg_runs_remember_coordinator_insert; + DROP TRIGGER trg_runs_remember_coordinator_update; + CREATE TRIGGER trg_runs_remember_coordinator_insert + AFTER INSERT ON runs + WHEN NEW.legacy = 0 AND (NEW.coordinator_handle IS NOT NULL OR NEW.coordinator_principal IS NOT NULL) + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + VALUES (NEW.id, COALESCE(NEW.coordinator_handle, NEW.coordinator_principal)); + END; + CREATE TRIGGER trg_runs_remember_coordinator_update + AFTER UPDATE OF coordinator_handle, coordinator_principal ON runs + WHEN NEW.legacy = 0 AND (NEW.coordinator_handle IS NOT NULL OR NEW.coordinator_principal IS NOT NULL) + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + VALUES (NEW.id, COALESCE(NEW.coordinator_handle, NEW.coordinator_principal)); + END; + `) + raw.pragma('user_version = 42') + try { + expect(resolveOrchestrationMigrationStartVersion(raw, 42, SCHEMA_VERSION)).toBe(6) + } finally { + raw.close() + } + + const db = new OrchestrationDb(path) + try { + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + expect(orcaSessionColumns(db.db)).toEqual(ORCA_SESSION_ID_COLUMNS) + expect(coordinatorTriggerSql(db.db).join('\n')).not.toContain('principal') + expect(db.getDispatchContextById(rows.structuredDispatchId)?.assignee_orca_session_id).toBe( + SESSION_ID + ) + expect(() => + db.createRun({ + objective: 'after the replay', + coordinatorHandle: 'term_after', + coordinatorPaneKey: 'tab_after:55555555-5555-4555-8555-555555555555' + }) + ).not.toThrow() + } finally { + db.close() + } + }) + + it("stops counting a chat coordinator's Orca session id once a v41 binary rebinds and then unbinds the Run", () => { + const path = tempDbPath() + const seeded = new OrchestrationDb(path) + seeded.db + .prepare( + `INSERT INTO runs ( + id, objective, coordinator_orca_session_id, coordinator_orca_session_id_generation, consumer_generation, legacy + ) VALUES ('run_chat', 'chat coordinated', ?, 1, 1, 0)` + ) + .run(CHAT_SESSION_ID) + // The cache row this binding wrote; a later open must not be able to write it back. + seeded.db.prepare('DELETE FROM run_coordinator_handles WHERE run_id = ?').run('run_chat') + seeded.close() + + // Each statement below is v41's own SQL. + const v41 = new Database(path) + // A terminal takes the Run over (bindRun)... + v41 + .prepare( + `UPDATE runs SET coordinator_handle = ?, coordinator_pane_key = ?, + consumer_generation = consumer_generation + 1, updated_at = datetime('now') + WHERE id = ?` + ) + .run('term_taker', PTY_WORKER_PANE, 'run_chat') + // ...then claims another Run from the same pane, which unbinds this one (unbindOtherRunsForPane). + v41 + .prepare( + `UPDATE runs SET coordinator_handle = NULL, coordinator_pane_key = NULL, + consumer_generation = consumer_generation + 1, updated_at = datetime('now') + WHERE id = ?` + ) + .run('run_chat') + // Without the generation this row is byte-identical to a live chat binding. + expect( + v41 + .prepare( + `SELECT coordinator_handle, coordinator_pane_key, coordinator_orca_session_id, + consumer_generation + FROM runs WHERE id = ?` + ) + .get('run_chat') + ).toEqual({ + coordinator_handle: null, + coordinator_pane_key: null, + coordinator_orca_session_id: CHAT_SESSION_ID, + consumer_generation: 3 + }) + v41.close() + + const reopened = new OrchestrationDb(path) + try { + const run = reopened.getRunRaw('run_chat') + expect(run && currentRunCoordinatorOrcaSessionId(run)).toBeNull() + expect( + reopened.db + .prepare(`SELECT id FROM runs WHERE ${currentRunCoordinatorOrcaSessionIdSql('runs')} = ?`) + .all(CHAT_SESSION_ID) + ).toEqual([]) + expect(coordinatorAddresses(reopened.db, ['run_chat'])).toEqual(['run_chat term_taker']) + } finally { + reopened.close() + } + }) + + it('finds Runs by coordinator Orca session id through an index on fresh and upgraded databases', () => { + const expectIndexedLookup = (path: string): void => { + const db = new OrchestrationDb(path) + try { + const plan = coordinatorLookupPlan(db.db) + expect(plan).toContain('USING INDEX idx_runs_coordinator_orca_session_id') + expect(plan).not.toContain('SCAN runs') + } finally { + db.close() + } + } + expectIndexedLookup(tempDbPath()) + + const upgradedPath = tempDbPath() + const seed = new OrchestrationDb(upgradedPath) + seedStructuredAndPtyRows(seed) + seed.close() + stripOrcaSessionSchema(upgradedPath, 41) + expectIndexedLookup(upgradedPath) + }) + + it('replays a dev database stamped v42 with the earlier *_actor columns to add the Orca session ids', () => { + const path = tempDbPath() + const seed = new OrchestrationDb(path) + const rows = seedStructuredAndPtyRows(seed) + seed.close() + stripOrcaSessionSchema(path, 41) + const raw = new Database(path) + // An earlier build of this step stamped v42 with `session:` values in differently named + // columns, their indexes, and address triggers over them. + const earlierAddress = `COALESCE(NEW.coordinator_handle, (CASE WHEN + NEW.coordinator_actor_generation = NEW.consumer_generation THEN NEW.coordinator_actor END))` + raw.exec(` + ALTER TABLE runs ADD COLUMN coordinator_actor TEXT; + ALTER TABLE runs ADD COLUMN coordinator_actor_generation INTEGER; + ALTER TABLE dispatch_contexts ADD COLUMN assignee_actor TEXT; + ALTER TABLE dispatch_contexts ADD COLUMN creator_actor TEXT; + CREATE INDEX idx_runs_coordinator_actor + ON runs(coordinator_actor) WHERE coordinator_actor IS NOT NULL; + CREATE INDEX idx_dispatch_assignee_actor + ON dispatch_contexts(assignee_actor) WHERE assignee_actor IS NOT NULL; + DROP TRIGGER trg_runs_remember_coordinator_insert; + DROP TRIGGER trg_runs_remember_coordinator_update; + CREATE TRIGGER trg_runs_remember_coordinator_insert + AFTER INSERT ON runs + WHEN NEW.legacy = 0 AND ${earlierAddress} IS NOT NULL + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + VALUES (NEW.id, ${earlierAddress}); + END; + CREATE TRIGGER trg_runs_remember_coordinator_update + AFTER UPDATE OF coordinator_handle, coordinator_actor, coordinator_actor_generation ON runs + WHEN NEW.legacy = 0 AND ${earlierAddress} IS NOT NULL + BEGIN + INSERT OR IGNORE INTO run_coordinator_handles (run_id, terminal_handle) + VALUES (NEW.id, ${earlierAddress}); + END; + `) + raw + .prepare( + `INSERT INTO runs ( + id, objective, coordinator_actor, coordinator_actor_generation, consumer_generation, legacy + ) VALUES ('run_earlier_chat', 'chat coordinated', ?, 1, 1, 0)` + ) + .run(CHAT_SESSION_ADDRESS) + // The cache row the earlier trigger wrote; only the stale column could write it back. + raw.prepare('DELETE FROM run_coordinator_handles WHERE run_id = ?').run('run_earlier_chat') + raw.pragma('user_version = 42') + try { + expect(resolveOrchestrationMigrationStartVersion(raw, 42, SCHEMA_VERSION)).toBe(6) + } finally { + raw.close() + } + + const db = new OrchestrationDb(path) + try { + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) + expect(orcaSessionColumns(db.db)).toEqual(ORCA_SESSION_ID_COLUMNS) + for (const sql of coordinatorTriggerSql(db.db)) { + expect(sql).toContain( + 'NEW.coordinator_orca_session_id_generation = NEW.consumer_generation' + ) + expect(sql).not.toContain('coordinator_actor') + } + expect(db.getDispatchContextById(rows.structuredDispatchId)?.assignee_orca_session_id).toBe( + SESSION_ID + ) + const run = db.getRunRaw(rows.structuredRunId) + expect(run && currentRunCoordinatorOrcaSessionId(run)).toBe(SESSION_ID) + // The stale column stays where it was and nothing reads it. + expect( + db.db + .prepare('SELECT coordinator_actor, coordinator_orca_session_id FROM runs WHERE id = ?') + .get('run_earlier_chat') + ).toEqual({ coordinator_actor: CHAT_SESSION_ADDRESS, coordinator_orca_session_id: null }) + expect(coordinatorAddresses(db.db, ['run_earlier_chat'])).toEqual([]) + expect(() => + db.createRun({ + objective: 'after the replay', + coordinatorHandle: 'term_after', + coordinatorPaneKey: 'tab_after:55555555-5555-4555-8555-555555555555' + }) + ).not.toThrow() + } finally { + db.close() + } + }) +}) diff --git a/src/main/runtime/orchestration/orchestration-schema-version-skew.ts b/src/main/runtime/orchestration/orchestration-schema-version-skew.ts index 85e0a78b3ae4..5d44066ca8f9 100644 --- a/src/main/runtime/orchestration/orchestration-schema-version-skew.ts +++ b/src/main/runtime/orchestration/orchestration-schema-version-skew.ts @@ -43,7 +43,11 @@ const VERSIONED_POST_V6_COLUMNS = [ { version: 36, table: 'remote_dispatch_attachments', column: 'consumer_generation' }, { version: 37, table: 'dispatch_contexts', column: 'creator_handle' }, { version: 37, table: 'dispatch_contexts', column: 'creator_pane_key' }, - { version: 40, table: 'remote_dispatch_attachments', column: 'home_run_id' } + { version: 40, table: 'remote_dispatch_attachments', column: 'home_run_id' }, + { version: 42, table: 'runs', column: 'coordinator_orca_session_id' }, + { version: 42, table: 'runs', column: 'coordinator_orca_session_id_generation' }, + { version: 42, table: 'dispatch_contexts', column: 'assignee_orca_session_id' }, + { version: 42, table: 'dispatch_contexts', column: 'creator_orca_session_id' } ] as const // Why: v34 shipped without these two, so a v34 stamp proves nothing about them; v35 repairs both diff --git a/src/main/runtime/orchestration/run-coordinator-orca-session-address.test.ts b/src/main/runtime/orchestration/run-coordinator-orca-session-address.test.ts new file mode 100644 index 000000000000..f93e1e4a4b5f --- /dev/null +++ b/src/main/runtime/orchestration/run-coordinator-orca-session-address.test.ts @@ -0,0 +1,262 @@ +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import { + formatOrcaSessionAddress, + parseOrcaSessionAddress +} from '../../../shared/orca-session-address' +import { + mintStructuredWorkerHandle, + mintStructuredWorkerPaneKey, + structuredWorkerProcessIncarnation +} from '../structured-worker-identity' +import { OrchestrationDb } from './db' +import { backfillStructuredWorkerOrcaSessionIds } from './db/schema/structured-worker-orca-session-backfill' + +const CHAT_SESSION_ID = '3a5c7e9b-1d4f-4a6c-8b0e-2f4a6c8e0b14' +const CHAT_ADDRESS = formatOrcaSessionAddress(CHAT_SESSION_ID) +const WORKER_SESSION_ID = '4b6d8f0c-2e5a-4b7d-9c1f-3a5b7d9f1c25' +const WORKER_ADDRESS = formatOrcaSessionAddress(WORKER_SESSION_ID) +const PTY_PANE = 'tab_pty:66666666-6666-4666-8666-666666666666' + +function addressesFor(db: OrchestrationDb, runId: string): string[] { + return db.db + .prepare('SELECT terminal_handle FROM run_coordinator_handles WHERE run_id = ?') + .all(runId) + .map((row) => String(row.terminal_handle)) + .sort() +} + +function tempDbPath(tempRoots: string[]): string { + const root = mkdtempSync(join(tmpdir(), 'orca-run-coordinator-orca-session-')) + tempRoots.push(root) + return join(root, 'orchestration.db') +} + +/** Drops the Run's remembered addresses and reopens, so only the on-open refill can write them back. */ +function refillAfterReopen(db: OrchestrationDb, path: string, runId: string): OrchestrationDb { + db.db.prepare('DELETE FROM run_coordinator_handles WHERE run_id = ?').run(runId) + db.close() + return new OrchestrationDb(path) +} + +/** A handle-less coordinator row; no writer records one until the caller resolver lands. */ +function insertSessionCoordinatedRun(db: OrchestrationDb, runId: string): void { + db.db + .prepare( + `INSERT INTO runs ( + id, objective, coordinator_orca_session_id, coordinator_orca_session_id_generation, + consumer_generation, legacy + ) VALUES (?, 'coordinated by a structured session', ?, 1, 1, 0)` + ) + .run(runId, CHAT_SESSION_ID) +} + +describe('Run coordinator Orca session address', () => { + let db: OrchestrationDb | undefined + const tempRoots: string[] = [] + + afterEach(() => { + db?.close() + db = undefined + for (const root of tempRoots.splice(0)) { + rmSync(root, { recursive: true, force: true }) + } + }) + + it('remembers a handle-less session coordinator by the address derived from its bare id', () => { + db = new OrchestrationDb(':memory:') + insertSessionCoordinatedRun(db, 'run_session') + + // The column holds the bare id; only the remembered address carries the session: prefix. + expect(db.getRunRaw('run_session')?.coordinator_orca_session_id).toBe(CHAT_SESSION_ID) + expect(addressesFor(db, 'run_session')).toEqual([CHAT_ADDRESS]) + expect(parseOrcaSessionAddress(addressesFor(db, 'run_session')[0])).toBe(CHAT_SESSION_ID) + expect(db.getRunMailboxOwnerIdsForHandle(CHAT_ADDRESS)).toEqual(['run_session']) + expect(db.getRunMailboxOwnerIdsForHandle(CHAT_SESSION_ID)).toEqual([]) + // The existing routing trigger matches the address by string equality, unchanged. + const reply = db.insertMessage({ + runId: 'run_session', + from: 'term_worker', + to: CHAT_ADDRESS, + subject: 'done', + type: 'worker_done' + }) + expect(db.getMessageById(reply.id)?.to_handle).toBe('run:run_session') + }) + + it('remembers an Orca session id bound by update, and again on reopen when the cache row is gone', () => { + const path = tempDbPath(tempRoots) + db = new OrchestrationDb(path) + db.db + .prepare( + `INSERT INTO runs (id, objective, consumer_generation, legacy) + VALUES ('run_unbound', 'bound later', 1, 0)` + ) + .run() + db.db + .prepare( + `UPDATE runs SET coordinator_orca_session_id = ?, + coordinator_orca_session_id_generation = consumer_generation + WHERE id = ?` + ) + .run(CHAT_SESSION_ID, 'run_unbound') + expect(addressesFor(db, 'run_unbound')).toEqual([CHAT_ADDRESS]) + + db = refillAfterReopen(db, path, 'run_unbound') + expect(addressesFor(db, 'run_unbound')).toEqual([CHAT_ADDRESS]) + }) + + it('remembers a PTY coordinator written by insert, update and refill by exactly its handle', () => { + const path = tempDbPath(tempRoots) + db = new OrchestrationDb(path) + db.db + .prepare( + `INSERT INTO runs (id, objective, coordinator_handle, consumer_generation, legacy) + VALUES ('run_pty', 'pty', 'term_first', 1, 0)` + ) + .run() + expect(addressesFor(db, 'run_pty')).toEqual(['term_first']) + db.db + .prepare( + `UPDATE runs SET coordinator_handle = 'term_second', + consumer_generation = consumer_generation + 1 WHERE id = 'run_pty'` + ) + .run() + expect(addressesFor(db, 'run_pty')).toEqual(['term_first', 'term_second']) + + db = refillAfterReopen(db, path, 'run_pty') + expect(addressesFor(db, 'run_pty')).toEqual(['term_second']) + }) + + it('remembers a structured-worker coordinator by its handle and its session address', () => { + const path = tempDbPath(tempRoots) + db = new OrchestrationDb(path) + const handle = mintStructuredWorkerHandle() + db.db + .prepare( + `INSERT INTO runs ( + id, objective, coordinator_handle, coordinator_orca_session_id, + coordinator_orca_session_id_generation, consumer_generation, legacy + ) VALUES ('run_inserted', 'structured worker', ?, ?, 1, 1, 0)` + ) + .run(handle, WORKER_SESSION_ID) + expect(addressesFor(db, 'run_inserted')).toEqual([WORKER_ADDRESS, handle].sort()) + + db.db + .prepare( + `INSERT INTO runs (id, objective, coordinator_handle, consumer_generation, legacy) + VALUES ('run_updated', 'id recorded later', ?, 1, 0)` + ) + .run(handle) + db.db + .prepare( + `UPDATE runs SET coordinator_orca_session_id = ?, + coordinator_orca_session_id_generation = consumer_generation + WHERE id = 'run_updated'` + ) + .run(WORKER_SESSION_ID) + expect(addressesFor(db, 'run_updated')).toEqual([WORKER_ADDRESS, handle].sort()) + + db = refillAfterReopen(db, path, 'run_updated') + expect(addressesFor(db, 'run_updated')).toEqual([WORKER_ADDRESS, handle].sort()) + expect(db.getRunMailboxOwnerIdsForHandle(WORKER_ADDRESS)).toEqual( + db.getRunMailboxOwnerIdsForHandle(handle) + ) + }) + + it('adds no session address for an Orca session id at a stale generation', () => { + const path = tempDbPath(tempRoots) + db = new OrchestrationDb(path) + db.db + .prepare( + `INSERT INTO runs ( + id, objective, coordinator_handle, coordinator_orca_session_id, + coordinator_orca_session_id_generation, consumer_generation, legacy + ) VALUES ('run_stale_insert', 'stale id', 'term_stale', ?, 1, 2, 0)` + ) + .run(WORKER_SESSION_ID) + expect(addressesFor(db, 'run_stale_insert')).toEqual(['term_stale']) + + db.db + .prepare( + `INSERT INTO runs (id, objective, consumer_generation, legacy) + VALUES ('run_stale_update', 'written stale', 2, 0)` + ) + .run() + db.db + .prepare( + `UPDATE runs SET coordinator_orca_session_id = ?, coordinator_orca_session_id_generation = 1 + WHERE id = 'run_stale_update'` + ) + .run(WORKER_SESSION_ID) + expect(addressesFor(db, 'run_stale_update')).toEqual([]) + + db = refillAfterReopen(db, path, 'run_stale_insert') + expect(addressesFor(db, 'run_stale_insert')).toEqual(['term_stale']) + expect(db.getRunMailboxOwnerIdsForHandle(WORKER_ADDRESS)).toEqual([]) + }) + + it('keeps PTY coordinators remembered by handle alone', () => { + db = new OrchestrationDb(':memory:') + const run = db.createRun({ + objective: 'pty', + coordinatorHandle: 'term_first', + coordinatorPaneKey: PTY_PANE + }) + db.bindRun({ + runId: run.id, + coordinatorHandle: 'term_second', + coordinatorPaneKey: 'tab_second:77777777-7777-4777-8777-777777777777' + }) + + expect(db.getRunRaw(run.id)?.coordinator_orca_session_id).toBeNull() + expect(addressesFor(db, run.id)).toEqual(['term_first', 'term_second']) + }) + + it("never leaves a replaced structured coordinator's Orca session id on the Run", () => { + db = new OrchestrationDb(':memory:') + const handle = mintStructuredWorkerHandle() + const pane = mintStructuredWorkerPaneKey(WORKER_SESSION_ID) + const ownTask = db.createTask({ runId: 'run_legacy_local', spec: 'structured worker' }) + db.createDispatchContext({ + taskId: ownTask.id, + assigneeHandle: handle, + assigneePaneKey: pane, + processIncarnation: structuredWorkerProcessIncarnation(WORKER_SESSION_ID), + creator: { kind: 'system' }, + maxDepth: Number.MAX_SAFE_INTEGER + }) + const first = db.createRun({ + objective: 'first', + coordinatorHandle: handle, + coordinatorPaneKey: pane + }) + backfillStructuredWorkerOrcaSessionIds(db.db) + expect(db.getRunRaw(first.id)?.coordinator_orca_session_id).toBe(WORKER_SESSION_ID) + + // A second Run from the same pane unbinds the first. + const second = db.createRun({ + objective: 'second', + coordinatorHandle: handle, + coordinatorPaneKey: pane + }) + expect(db.getRunRaw(first.id)).toMatchObject({ + coordinator_handle: null, + coordinator_orca_session_id: null + }) + + backfillStructuredWorkerOrcaSessionIds(db.db) + expect(db.getRunRaw(second.id)?.coordinator_orca_session_id).toBe(WORKER_SESSION_ID) + db.bindRun({ runId: second.id, coordinatorHandle: 'term_taker', coordinatorPaneKey: PTY_PANE }) + expect(db.getRunRaw(second.id)).toMatchObject({ + coordinator_handle: 'term_taker', + coordinator_orca_session_id: null + }) + // A remembered address is never forgotten, so the worker's session address reaches exactly the + // Runs its handle does. + expect(db.getRunMailboxOwnerIdsForHandle(WORKER_ADDRESS)).toEqual([first.id, second.id].sort()) + expect(db.getRunMailboxOwnerIdsForHandle(handle)).toEqual([first.id, second.id].sort()) + }) +}) diff --git a/src/main/runtime/orchestration/types.ts b/src/main/runtime/orchestration/types.ts index 85d5dcfc1596..d00ea480c7dd 100644 --- a/src/main/runtime/orchestration/types.ts +++ b/src/main/runtime/orchestration/types.ts @@ -46,6 +46,10 @@ export type RunRow = { home_database: string coordinator_handle: string | null coordinator_pane_key: string | null + /** Bare Orca session id the coordinator is addressed by, when it has one (today only structured sessions); a `/clear`ed chat's lineage root. */ + coordinator_orca_session_id: string | null + /** The consumer_generation the id was written at; see currentRunCoordinatorOrcaSessionId. */ + coordinator_orca_session_id_generation: number | null consumer_generation: number legacy: number created_at: string @@ -278,6 +282,8 @@ export type DispatchContextRow = { launch_token_hash: string | null assignee_handle: string | null assignee_pane_key: string | null + /** Bare Orca session id the assignee is addressed by, when it has one (today only structured sessions); a `/clear`ed chat's lineage root. */ + assignee_orca_session_id: string | null capability_hash: string | null process_incarnation: string | null capability_revoked_at: string | null @@ -287,6 +293,8 @@ export type DispatchContextRow = { /** Creator identity; equal to the assignee means a self-dispatch, which adds no nesting depth. */ creator_handle: string | null creator_pane_key: string | null + /** Bare Orca session id the creator is addressed by, when it has one (today only structured sessions); a `/clear`ed chat's lineage root. */ + creator_orca_session_id: string | null host_scope: string | null status: DispatchStatus failure_count: number diff --git a/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.test.ts b/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.test.ts index 83434c631194..49afdfe119a7 100644 --- a/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.test.ts @@ -9,6 +9,8 @@ const RUN_ROW: RunRow = { home_database: '/tmp/orca/orchestration.db', coordinator_handle: 'term_coord', coordinator_pane_key: 'tab_coord:11111111-1111-4111-8111-111111111111', + coordinator_orca_session_id: '22222222-2222-4222-8222-222222222222', + coordinator_orca_session_id_generation: 3, consumer_generation: 3, legacy: 0, created_at: '2026-09-04T18:53:07Z', @@ -30,6 +32,8 @@ describe('exposeRun', () => { ]) expect(exposed).not.toHaveProperty('home_database') expect(exposed).not.toHaveProperty('coordinator_pane_key') + expect(exposed).not.toHaveProperty('coordinator_orca_session_id') + expect(exposed).not.toHaveProperty('coordinator_orca_session_id_generation') }) it('preserves every published column by value', () => { @@ -54,8 +58,15 @@ describe('exposeRun', () => { }) it('strips the columns even when they are null', () => { - const exposed = exposeRun({ ...RUN_ROW, coordinator_pane_key: null }) + const exposed = exposeRun({ + ...RUN_ROW, + coordinator_pane_key: null, + coordinator_orca_session_id: null, + coordinator_orca_session_id_generation: null + }) expect(exposed).not.toHaveProperty('coordinator_pane_key') + expect(exposed).not.toHaveProperty('coordinator_orca_session_id') + expect(exposed).not.toHaveProperty('coordinator_orca_session_id_generation') }) }) diff --git a/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.ts b/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.ts index 30a22e2fc18c..b376d205874b 100644 --- a/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.ts +++ b/src/main/runtime/rpc/methods/orchestration/runs/run-receipt.ts @@ -1,7 +1,13 @@ import type { RunRow } from '../../../../orchestration/types' // Why: home_database and coordinator_pane_key are runtime routing state; no caller reads them. -const INTERNAL_RUN_COLUMNS = ['home_database', 'coordinator_pane_key'] as const +// The coordinator's Orca session id stays off the wire until a reader needs it; publishing it is a wire change. +const INTERNAL_RUN_COLUMNS = [ + 'home_database', + 'coordinator_pane_key', + 'coordinator_orca_session_id', + 'coordinator_orca_session_id_generation' +] as const export type RunReceipt = Omit diff --git a/src/main/runtime/structured-worker-identity.ts b/src/main/runtime/structured-worker-identity.ts index b29d68c297a7..ef2fec6a8fa4 100644 --- a/src/main/runtime/structured-worker-identity.ts +++ b/src/main/runtime/structured-worker-identity.ts @@ -32,8 +32,8 @@ import { // Deliberately not `term_`: `issueHandle` revalidates the renderer graph epoch against the // renderer-driven leaves map, so a main-minted `term_` leaf evaporates on the next window reload. -const STRUCTURED_WORKER_HANDLE_PREFIX = 'structworker_' -const STRUCTURED_WORKER_INCARNATION_PREFIX = 'structured:' +export const STRUCTURED_WORKER_HANDLE_PREFIX = 'structworker_' +export const STRUCTURED_WORKER_INCARNATION_PREFIX = 'structured:' export type StructuredWorkerIdentity = { handle: string diff --git a/src/shared/orca-session-address.test.ts b/src/shared/orca-session-address.test.ts new file mode 100644 index 000000000000..96b8a2268e5e --- /dev/null +++ b/src/shared/orca-session-address.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from 'vitest' +import { + ORCA_SESSION_ADDRESS_PREFIX, + formatOrcaSessionAddress, + isOrcaSessionId, + parseOrcaSessionAddress +} from './orca-session-address' + +const SESSION_ID = '0b7e4c2a-5f1d-4e8a-9c3b-2d6f8a1e4b70' +const ADDRESS = `session:${SESSION_ID}` + +describe('Orca session address', () => { + it('addresses an Orca session id as session: and parses the bare id back', () => { + expect(ORCA_SESSION_ADDRESS_PREFIX).toBe('session:') + expect(formatOrcaSessionAddress(SESSION_ID)).toBe(ADDRESS) + expect(parseOrcaSessionAddress(ADDRESS)).toBe(SESSION_ID) + expect(formatOrcaSessionAddress(parseOrcaSessionAddress(ADDRESS) ?? '')).toBe(ADDRESS) + }) + + it('reads only the addressed spelling when parsing an address', () => { + // A bare id is what the columns store, not an address. + expect(parseOrcaSessionAddress(SESSION_ID)).toBeNull() + expect(parseOrcaSessionAddress(null)).toBeNull() + expect(parseOrcaSessionAddress(undefined)).toBeNull() + expect(parseOrcaSessionAddress('')).toBeNull() + }) + + it.each([ + ['an unknown prefix', `pane:${SESSION_ID}`], + ['the Run mailbox namespace', 'run:run_123'], + ['the Dispatch mailbox namespace', 'dispatch:ctx_123'], + ['an empty prefix', `:${SESSION_ID}`], + ['an empty id', 'session:'], + ['an id with a separator', `session:${SESSION_ID}:extra`], + ['an id the session predicate rejects', 'session:short'], + ['a terminal handle', 'term_4f2c9a'] + ])('refuses %s', (_label, value) => { + expect(parseOrcaSessionAddress(value)).toBeNull() + }) + + it.each([ + ['a PTY terminal handle', 'term_4f2c9a1b-7d3e-4a5f-8b6c-9d0e1f2a3b4c'], + ['a short PTY terminal handle', 'term_4f2c9a'], + ['a structured-worker handle', 'structworker_4f2c9a1b-7d3e-4a5f-8b6c-9d0e1f2a3b4c'] + ])('never treats %s as an Orca session id', (_label, handle) => { + // Handles share the session-id charset, so the session-record predicate alone would accept them. + expect(isOrcaSessionId(handle)).toBe(false) + expect(parseOrcaSessionAddress(`session:${handle}`)).toBeNull() + }) + + it('validates an Orca session id with the session-record predicate', () => { + expect(isOrcaSessionId(SESSION_ID)).toBe(true) + expect(isOrcaSessionId('has space in it')).toBe(false) + expect(isOrcaSessionId('x'.repeat(129))).toBe(false) + }) +}) diff --git a/src/shared/orca-session-address.ts b/src/shared/orca-session-address.ts new file mode 100644 index 000000000000..f2ccd4fe00cc --- /dev/null +++ b/src/shared/orca-session-address.ts @@ -0,0 +1,36 @@ +import { isAgentSessionId } from './agent-session-record' + +/** + * The Orca session id is the id Orca minted for a structured session (its session record id), never + * the provider's own session id. Orchestration stores, bare, the one the agent is addressed by: for + * a `/clear`ed chat, its lineage root's, not the live session's. Mail addresses the session as + * `session:`, beside `run:` and `dispatch:`, and derives that spelling here rather than + * storing it. + * + * Where the session runs is not part of the id; it is read from the session record when needed. PTY + * agents have none today, and never a pane-keyed one: a pane outlives the agent in it, so such an id + * would be inherited by the pane's next occupant. + */ +export const ORCA_SESSION_ADDRESS_PREFIX = 'session:' + +// Terminal handles (`term_` from the PTY runtime, `structworker_` from structured-worker-identity) +// share the session-id charset. A handle is never a session, so one handed over by mistake must not +// become a durable Orca session id. +const TERMINAL_HANDLE_PREFIXES = ['term_', 'structworker_'] as const + +export function isOrcaSessionId(id: string): boolean { + return isAgentSessionId(id) && !TERMINAL_HANDLE_PREFIXES.some((prefix) => id.startsWith(prefix)) +} + +export function formatOrcaSessionAddress(orcaSessionId: string): string { + return `${ORCA_SESSION_ADDRESS_PREFIX}${orcaSessionId}` +} + +/** The bare Orca session id of a `session:` address; anything else reads as null. */ +export function parseOrcaSessionAddress(address: string | null | undefined): string | null { + if (!address?.startsWith(ORCA_SESSION_ADDRESS_PREFIX)) { + return null + } + const id = address.slice(ORCA_SESSION_ADDRESS_PREFIX.length) + return isOrcaSessionId(id) ? id : null +} diff --git a/tests/e2e/cross-version-wire/orchestration-delivery-downgrade.unit.test.ts b/tests/e2e/cross-version-wire/orchestration-delivery-downgrade.unit.test.ts index e53b3bc9e6cf..8126d33d6010 100644 --- a/tests/e2e/cross-version-wire/orchestration-delivery-downgrade.unit.test.ts +++ b/tests/e2e/cross-version-wire/orchestration-delivery-downgrade.unit.test.ts @@ -3,12 +3,13 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import { expect, test } from 'vitest' import { OrchestrationDb } from '../../../src/main/runtime/orchestration/db' +import { SCHEMA_VERSION } from '../../../src/main/runtime/orchestration/db/contract-constants' import { importReleaseCheckoutModule, materializeReleaseCheckout } from './release-checkout' // Pin the last pre-v41 implementation: this contract specifically exercises status-only readers. const PRE_V41 = 'aac38d698ff75ac4c8658addab48ef5a83617619' -test('pre-v41 code opens, acknowledges and writes a v41 database, then current code reopens it', async () => { +test('pre-v41 code opens, acknowledges and writes a current-schema database, then current code reopens it', async () => { const checkout = await materializeReleaseCheckout(PRE_V41) const baseline = await importReleaseCheckoutModule( checkout, @@ -47,7 +48,8 @@ test('pre-v41 code opens, acknowledges and writes a v41 database, then current c db.close() db = new OldDb(path) - expect(db.db.pragma('user_version', { simple: true })).toBe(41) + // Old code leaves a newer stamp alone, so the reopen still reads the current schema version. + expect(db.db.pragma('user_version', { simple: true })).toBe(SCHEMA_VERSION) // Old readers retain their original replay semantics, but can acknowledge either stored batch. db.acknowledgeRunDelivery({ ...params, deliveryId: oldBatch.delivery.id }) expect(db.getOrCreateRunDelivery(params)?.delivery.id).toBe(currentBatch.delivery.id)