Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion src/main/runtime/orchestration/db/contract-constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
49 changes: 49 additions & 0 deletions src/main/runtime/orchestration/db/dispatch-depth.test.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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)
})
})
2 changes: 2 additions & 0 deletions src/main/runtime/orchestration/db/orchestration-db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down
4 changes: 4 additions & 0 deletions src/main/runtime/orchestration/db/row-column-lists.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down Expand Up @@ -45,13 +47,15 @@ export const DISPATCH_CONTEXT_COLUMNS = [
'launch_token_hash',
'assignee_handle',
'assignee_pane_key',
'assignee_orca_session_id',
'capability_hash',
'process_incarnation',
'capability_revoked_at',
'retry_of_dispatch_id',
'creator_dispatch_id',
'creator_handle',
'creator_pane_key',
'creator_orca_session_id',
'host_scope',
'status',
'failure_count',
Expand Down
4 changes: 3 additions & 1 deletion src/main/runtime/orchestration/db/runs/run-binding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 = ?`
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { OrchestrationDb } from '../orchestration-db'
import { currentRunCoordinatorSessionAddressSql } from './run-coordinator-orca-session'

export function rememberRunCoordinatorHandle(
this: OrchestrationDb,
Expand All @@ -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;
`)
}

Expand Down
Original file line number Diff line number Diff line change
@@ -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:<id>` 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)})`
}
3 changes: 2 additions & 1 deletion src/main/runtime/orchestration/db/runs/run-lookup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 = ?`
Expand Down
10 changes: 10 additions & 0 deletions src/main/runtime/orchestration/db/schema/create-core-tables-sql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')),
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:<id> address.
assignee_orca_session_id TEXT,
capability_hash TEXT,
process_incarnation TEXT,
capability_revoked_at TEXT,
Expand All @@ -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')),
Expand Down
72 changes: 72 additions & 0 deletions src/main/runtime/orchestration/db/schema/migrate-v42.ts
Original file line number Diff line number Diff line change
@@ -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:<id>`, 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;
`)
}
2 changes: 2 additions & 0 deletions src/main/runtime/orchestration/db/schema/migrate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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)
Expand Down
Loading
Loading