perf(orchestration): project explicit columns so the existing cache covers the hot path

Replaces the branch's second statement cache. The six graph-publish reads were
uncacheable only because they were spelled `SELECT *` / `SELECT t.*`, which
SyncDatabase refuses to cache (node:sqlite can build the first row after a schema
change from stale column names). Spelling the projection out from type-checked
column tuples makes them cacheable by the SyncDatabase LRU that is already merged,
already bounded, and already clears on DDL — so the WeakMap and its documented
cross-connection ALTER hazard both go away.

Drift is caught at build time: `satisfies readonly (keyof Row)[]` plus an
`Exclude<keyof Row, Cols[number]> extends never` assertion pins list vs type at tsc,
and a PRAGMA table_info test against a freshly migrated OrchestrationDb pins list
vs schema.

Same win, verified: 6 compilations per publish -> 2 total then 0, identical to the
WeakMap branch; 92/96/91 us CPU per 2-pane publish before, 11-12 us after on both.
This commit is contained in:
Neil
2026-09-03 19:57:49 -07:00
parent 2579469201
commit 037ef0502f
7 changed files with 182 additions and 87 deletions
@@ -6,24 +6,23 @@ import {
paneKeyMatchSuffix
} from '../pane-key-match'
import type { OrchestrationDb } from '../orchestration-db'
import { prepareCachedOrchestrationRead } from '../prepared-statement-cache'
import { DISPATCH_CONTEXT_COLUMN_LIST } from '../row-column-lists'
// Why: hoisted so the graph-publish fan-out hits one stable cache key per lookup shape.
// Why: hoisted and wildcard-free so the graph-publish fan-out hits the SyncDatabase statement cache.
const ACTIVE_DISPATCH_BY_HANDLE_SQL =
// Why: newest-first like the pane lookups below — an unordered LIMIT 1 could pin a stale row if a handle ever has two active dispatches.
`SELECT * FROM dispatch_contexts
`SELECT ${DISPATCH_CONTEXT_COLUMN_LIST} FROM dispatch_contexts
WHERE assignee_handle = ? AND status IN ('pending', 'dispatched')
ORDER BY rowid DESC LIMIT 1`
const ACTIVE_DISPATCH_BY_PANE_KEY_SQL = `SELECT * FROM dispatch_contexts
const ACTIVE_DISPATCH_BY_PANE_KEY_SQL = `SELECT ${DISPATCH_CONTEXT_COLUMN_LIST} FROM dispatch_contexts
WHERE assignee_pane_key = ? AND status IN ('pending', 'dispatched')
ORDER BY rowid DESC LIMIT 1`
const ACTIVE_DISPATCH_BY_PANE_SUFFIX_SQL = `SELECT * FROM dispatch_contexts
const ACTIVE_DISPATCH_BY_PANE_SUFFIX_SQL = `SELECT ${DISPATCH_CONTEXT_COLUMN_LIST} FROM dispatch_contexts
WHERE assignee_pane_key IS NOT NULL
AND status IN ('pending', 'dispatched') AND instr(assignee_pane_key, ':') > 1
AND ${DISPATCH_PANE_KEY_MATCH_SUFFIX_SQL} = ?
ORDER BY rowid DESC LIMIT 1`
const LATEST_DISPATCH_BY_HANDLE_SQL =
'SELECT * FROM dispatch_contexts WHERE assignee_handle = ? ORDER BY rowid DESC LIMIT 1'
const LATEST_DISPATCH_BY_HANDLE_SQL = `SELECT ${DISPATCH_CONTEXT_COLUMN_LIST} FROM dispatch_contexts WHERE assignee_handle = ? ORDER BY rowid DESC LIMIT 1`
export function getActiveDispatchForTerminal(
this: OrchestrationDb,
@@ -134,9 +133,9 @@ export function findActiveDispatchForAssignee(
assigneeHandle: string,
assigneePaneKey?: string
): DispatchContextRow | undefined {
const byHandle = prepareCachedOrchestrationRead(this.db, ACTIVE_DISPATCH_BY_HANDLE_SQL).get(
assigneeHandle
) as DispatchContextRow | undefined
const byHandle = this.db.prepare(ACTIVE_DISPATCH_BY_HANDLE_SQL).get(assigneeHandle) as
| DispatchContextRow
| undefined
if (byHandle) {
return byHandle
}
@@ -145,25 +144,25 @@ export function findActiveDispatchForAssignee(
return undefined
}
const exactPane = prepareCachedOrchestrationRead(this.db, ACTIVE_DISPATCH_BY_PANE_KEY_SQL).get(
assigneePaneKey
) as DispatchContextRow | undefined
const exactPane = this.db.prepare(ACTIVE_DISPATCH_BY_PANE_KEY_SQL).get(assigneePaneKey) as
| DispatchContextRow
| undefined
if (exactPane) {
return exactPane
}
if (!parsePaneKey(assigneePaneKey)) {
return undefined
}
return prepareCachedOrchestrationRead(this.db, ACTIVE_DISPATCH_BY_PANE_SUFFIX_SQL).get(
paneKeyMatchSuffix(assigneePaneKey)
) as DispatchContextRow | undefined
return this.db
.prepare(ACTIVE_DISPATCH_BY_PANE_SUFFIX_SQL)
.get(paneKeyMatchSuffix(assigneePaneKey)) as DispatchContextRow | undefined
}
export function getLatestDispatchForTerminal(
this: OrchestrationDb,
handle: string
): DispatchContextRow | undefined {
return prepareCachedOrchestrationRead(this.db, LATEST_DISPATCH_BY_HANDLE_SQL).get(handle) as
return this.db.prepare(LATEST_DISPATCH_BY_HANDLE_SQL).get(handle) as
| DispatchContextRow
| undefined
}
@@ -7,7 +7,6 @@ import type { OrchestrationCompatibilityTerminalAuthority } from '../../runtime-
import type { RuntimeLeafRecord } from '../../runtime-terminal-state-records'
import { OrchestrationDb } from '../db'
import { createRootDispatch } from './root-dispatch-test-fixture'
import { prepareCachedOrchestrationRead } from './prepared-statement-cache'
const COORDINATOR_HANDLE = 'term_coordinator'
const COORDINATOR_PANE = 'tab_c:leaf_c'
@@ -15,7 +14,9 @@ const WORKER_HANDLE = 'term_worker'
const WORKER_PANE = 'tab_w:leaf_w'
const IDLE_HANDLE = 'term_idle'
const IDLE_PANE = 'tab_i:leaf_i'
const RUN_BY_ID_SQL = 'SELECT * FROM runs WHERE id = ?'
// Why: mirrors SyncDatabase's `isStatementCacheable` — aggregate `(*)` is fine, any other `*` is not.
const WILDCARD_PROJECTION = /(?<!\(\s*)\*/
const openDatabases: OrchestrationDb[] = []
const temporaryDirectories: string[] = []
@@ -40,7 +41,7 @@ function openDatabase(path: string): OrchestrationDb {
}
function temporaryDatabasePath(): string {
const directory = mkdtempSync(join(tmpdir(), 'orca-orchestration-cache-'))
const directory = mkdtempSync(join(tmpdir(), 'orca-orchestration-hot-path-'))
temporaryDirectories.push(directory)
return join(directory, 'orchestration.db')
}
@@ -104,7 +105,7 @@ function buildProjection(db: OrchestrationDb): RuntimeAgentOrchestrationProjecti
})
}
describe('orchestration prepared-statement cache', () => {
describe('orchestration hot-path statement compilation', () => {
it('compiles each hot-path SQL exactly once across repeated graph publishes', () => {
const db = openDatabase(':memory:')
seedDispatchedWorker(db)
@@ -130,26 +131,18 @@ describe('orchestration prepared-statement cache', () => {
}
})
it('reuses one statement object per connection', () => {
const db = openDatabase(':memory:')
expect(prepareCachedOrchestrationRead(db.db, RUN_BY_ID_SQL)).toBe(
prepareCachedOrchestrationRead(db.db, RUN_BY_ID_SQL)
)
})
it('does not serve a statement from the previous connection after reopen', () => {
// Why: `SELECT *` is what made these statements uncacheable, and a retained wildcard is the only
// way node:sqlite could build a row from stale column names after another connection's ALTER.
// Seeds on one connection and publishes on a second so every compilation here is hot-path SQL.
it('publishes without compiling a single wildcard projection', () => {
const path = temporaryDatabasePath()
const first = openDatabase(path)
const beforeClose = prepareCachedOrchestrationRead(first.db, RUN_BY_ID_SQL)
first.close()
seedDispatchedWorker(openDatabase(path))
const second = openDatabase(path)
const afterReopen = prepareCachedOrchestrationRead(second.db, RUN_BY_ID_SQL)
const reader = openDatabase(path)
const compiled = trackCompiledSql(reader)
buildProjection(reader).buildByPaneKey()
// Why Object.is: the closed connection's statement throws on any property read, so vitest
// cannot format it as a matcher operand.
expect(Object.is(afterReopen, beforeClose)).toBe(false)
expect(() => beforeClose.get('run_missing')).toThrow()
expect(() => afterReopen.get('run_missing')).not.toThrow()
expect(compiled.length).toBeGreaterThan(0)
expect(compiled.filter((sql) => WILDCARD_PROJECTION.test(sql))).toEqual([])
})
})
@@ -1,33 +0,0 @@
import type Database from '../../../sqlite/sync-database'
// Why: SyncDatabase deliberately refuses to cache any `SELECT *` (node:sqlite can build the first
// row after a schema change from stale column names), so orchestration's wildcard reads recompile
// their SQL on every call — and the graph publish runs that fan-out once per pane. Orchestration
// DDL only runs in the OrchestrationDb constructor (createTables/migrate/trigger) and resets are
// DELETE-only, so a connection's column shape is frozen before the first read.
const statementsByDatabase = new WeakMap<Database.Database, Map<string, Database.Statement>>()
/**
* Compile `sql` once per open connection. Keyed by Database instance, so a reopened database
* starts empty and never serves a statement bound to the previous connection.
*
* Only for read statements with fixed SQL text. Never for DDL, dynamic `IN (?,?,…)` arities, or
* statements passed to `.iterate()` (better-sqlite3/node:sqlite reject re-entrant iteration).
*/
export function prepareCachedOrchestrationRead(
db: Database.Database,
sql: string
): Database.Statement {
let statements = statementsByDatabase.get(db)
if (!statements) {
statements = new Map()
statementsByDatabase.set(db, statements)
}
const cached = statements.get(sql)
if (cached) {
return cached
}
const statement = db.prepare(sql)
statements.set(sql, statement)
return statement
}
@@ -0,0 +1,57 @@
import { afterEach, describe, expect, it } from 'vitest'
import { OrchestrationDb } from './orchestration-db'
import {
DISPATCH_CONTEXT_COLUMNS,
RUN_COLUMNS,
selectColumns,
TASK_COLUMNS
} from './row-column-lists'
let db: OrchestrationDb | undefined
afterEach(() => {
db?.close()
db = undefined
})
function tableColumns(table: string): string[] {
const rows = (db as OrchestrationDb).db.pragma(`table_info(${table})`) as { name: string }[]
return rows.map((row) => row.name).sort()
}
describe('row column lists', () => {
// Why: these lists replaced `SELECT *`, so a column added to the schema without being listed here
// would silently stop being read. tsc pins list↔type; this pins list↔schema.
it.each([
['runs', RUN_COLUMNS],
['tasks', TASK_COLUMNS],
['dispatch_contexts', DISPATCH_CONTEXT_COLUMNS]
])('projects every %s column the migrated schema declares', (table, columns) => {
db = new OrchestrationDb(':memory:')
expect([...columns].sort()).toEqual(tableColumns(table))
})
it('qualifies each column when the statement joins under an alias', () => {
expect(selectColumns(['id', 'run_id'])).toBe('id, run_id')
expect(selectColumns(['id', 'run_id'], 't')).toBe('t.id, t.run_id')
})
// Why: an alias-qualified projection must key the returned row by the bare column name, exactly as
// the `t.*` it replaced did — otherwise every lineage consumer reads undefined.
it('returns bare column names for an alias-qualified projection', () => {
db = new OrchestrationDb(':memory:')
const run = db.createRun({
objective: 'demo',
coordinatorHandle: 'term_c',
coordinatorPaneKey: 'tab_c:leaf_c'
})
const task = db.createTask({ spec: 'work', runId: run.id })
const row = db.db
.prepare(`SELECT ${selectColumns(TASK_COLUMNS, 't')} FROM tasks t WHERE t.id = ?`)
.get(task.id) as Record<string, unknown>
expect(Object.keys(row).sort()).toEqual([...TASK_COLUMNS].sort())
})
})
@@ -0,0 +1,82 @@
import type { DispatchContextRow, RunRow, TaskRow } from '../types'
// Why: `SyncDatabase` refuses to cache any `SELECT *` (node:sqlite can build the first row after a
// schema change from stale column names), so a wildcard read recompiles its SQL on every call.
// Spelling the projection out makes the hot-path statements cacheable by that existing LRU.
// Drift is caught twice: `satisfies` + the exhaustiveness assertions below pin list↔type at tsc,
// and `row-column-lists.test.ts` pins list↔schema against a freshly migrated database.
export const RUN_COLUMNS = [
'id',
'objective',
'home_database',
'coordinator_handle',
'coordinator_pane_key',
'consumer_generation',
'legacy',
'created_at',
'updated_at'
] as const satisfies readonly (keyof RunRow)[]
export const TASK_COLUMNS = [
'id',
'run_id',
'parent_id',
'created_by_terminal_handle',
'created_by_pane_key',
'created_by_process_incarnation',
'created_by_run_generation',
'task_title',
'display_name',
'spec',
'status',
'deps',
'result',
'created_at',
'completed_at'
] as const satisfies readonly (keyof TaskRow)[]
export const DISPATCH_CONTEXT_COLUMNS = [
'id',
'run_id',
'task_id',
'contract_version',
'launch_token_hash',
'assignee_handle',
'assignee_pane_key',
'capability_hash',
'process_incarnation',
'capability_revoked_at',
'status',
'failure_count',
'last_failure',
'termination_reason',
'depth',
'dispatched_at',
'completed_at',
'created_at',
'last_heartbeat_at'
] as const satisfies readonly (keyof DispatchContextRow)[]
// Compile check: a row field added without its column here would silently vanish from the
// projection that used to be `SELECT *`, so the missing key must fail the build.
type UnprojectedRunColumn = Exclude<keyof RunRow, (typeof RUN_COLUMNS)[number]>
type UnprojectedTaskColumn = Exclude<keyof TaskRow, (typeof TASK_COLUMNS)[number]>
type UnprojectedDispatchContextColumn = Exclude<
keyof DispatchContextRow,
(typeof DISPATCH_CONTEXT_COLUMNS)[number]
>
const assertEveryRowColumnProjected: [
UnprojectedRunColumn extends never ? true : never,
UnprojectedTaskColumn extends never ? true : never,
UnprojectedDispatchContextColumn extends never ? true : never
] = [true, true, true]
void assertEveryRowColumnProjected
/** Projection list for a `SELECT`; `alias` qualifies each name for a joined table (`t.id, …`). */
export function selectColumns(columns: readonly string[], alias?: string): string {
return columns.map((column) => (alias ? `${alias}.${column}` : column)).join(', ')
}
export const RUN_COLUMN_LIST = selectColumns(RUN_COLUMNS)
export const DISPATCH_CONTEXT_COLUMN_LIST = selectColumns(DISPATCH_CONTEXT_COLUMNS)
@@ -9,16 +9,16 @@ import { exposeRunTimestamps } from '../utc-timestamp'
import { encodeRunListCursor, decodeRunListCursor } from '../run-list-cursor'
import type { RunListPage } from '../run-list-page'
import type { OrchestrationDb } from '../orchestration-db'
import { prepareCachedOrchestrationRead } from '../prepared-statement-cache'
import { RUN_COLUMN_LIST } from '../row-column-lists'
export type LegacyAdoptedMailboxOwner = {
runId: string
terminalHandle: string
}
// Why: hoisted so the per-publish run lookups hit one stable cache key each.
const RUN_BY_ID_SQL = 'SELECT * FROM runs WHERE id = ?'
const RUNS_BOUND_TO_PANE_SQL = `SELECT * FROM runs
// Why: hoisted and wildcard-free so the per-publish run lookups hit the SyncDatabase statement cache.
const RUN_BY_ID_SQL = `SELECT ${RUN_COLUMN_LIST} FROM runs WHERE id = ?`
const RUNS_BOUND_TO_PANE_SQL = `SELECT ${RUN_COLUMN_LIST} FROM runs
WHERE coordinator_pane_key IS NOT NULL AND legacy = 0
AND ${RUN_PANE_KEY_MATCH_SUFFIX_SQL} = ?
ORDER BY rowid`
@@ -111,9 +111,7 @@ export function getCurrentRunForPane(this: OrchestrationDb, paneKey: string): Ru
// reminted tab halves keep matching and unparseable keys keep requiring an exact match.
export function runsBoundToPane(this: OrchestrationDb, paneKey: string): RunRow[] {
return (
prepareCachedOrchestrationRead(this.db, RUNS_BOUND_TO_PANE_SQL).all(
paneKeyMatchSuffix(paneKey)
) as RunRow[]
this.db.prepare(RUNS_BOUND_TO_PANE_SQL).all(paneKeyMatchSuffix(paneKey)) as RunRow[]
).filter(
(run) =>
run.coordinator_pane_key !== null && isEquivalentPaneKey(run.coordinator_pane_key, paneKey)
@@ -121,7 +119,7 @@ export function runsBoundToPane(this: OrchestrationDb, paneKey: string): RunRow[
}
export function getRunRaw(this: OrchestrationDb, id: string): RunRow | undefined {
return prepareCachedOrchestrationRead(this.db, RUN_BY_ID_SQL).get(id) as RunRow | undefined
return this.db.prepare(RUN_BY_ID_SQL).get(id) as RunRow | undefined
}
export function unbindOtherRunsForPane(
@@ -5,7 +5,7 @@ import { LEGACY_RUN_ID } from '../contract-constants'
import { generateId } from '../generated-id'
import type { TaskRuntimeLineageRow } from '../run-list-page'
import type { OrchestrationDb } from '../orchestration-db'
import { prepareCachedOrchestrationRead } from '../prepared-statement-cache'
import { selectColumns, TASK_COLUMNS } from '../row-column-lists'
// ── Tasks ──
@@ -82,8 +82,8 @@ export function createTask(
return this.db.prepare('SELECT * FROM tasks WHERE id = ?').get(id) as TaskRow
}
// Why: hoisted so the per-publish lineage lookup hits one stable cache key.
const TASK_RUNTIME_LINEAGE_SQL = `SELECT t.*,
// Why: hoisted and wildcard-free so the per-publish lineage lookup hits the SyncDatabase statement cache.
const TASK_RUNTIME_LINEAGE_SQL = `SELECT ${selectColumns(TASK_COLUMNS, 't')},
creator.id AS creator_dispatch_id,
creator.run_id AS creator_dispatch_run_id,
creator.assignee_pane_key AS creator_dispatch_pane_key,
@@ -115,10 +115,9 @@ export function getTask(
if (dispatchRunId === undefined) {
return this.db.prepare('SELECT * FROM tasks WHERE id = ?').get(id) as TaskRow | undefined
}
return prepareCachedOrchestrationRead(this.db, TASK_RUNTIME_LINEAGE_SQL).get(
dispatchRunId,
id
) as TaskRuntimeLineageRow | undefined
return this.db.prepare(TASK_RUNTIME_LINEAGE_SQL).get(dispatchRunId, id) as
| TaskRuntimeLineageRow
| undefined
}
export function listTasks(