diff --git a/packages/runtime-host/src/__tests__/execution-host-queue.test.ts b/packages/runtime-host/src/__tests__/execution-host-queue.test.ts index d918317c93..40f8c22a18 100644 --- a/packages/runtime-host/src/__tests__/execution-host-queue.test.ts +++ b/packages/runtime-host/src/__tests__/execution-host-queue.test.ts @@ -245,8 +245,11 @@ test('production UDS admission commits one transcript before the root handoff', }); assert.equal(started.disposition, 'turn_started'); if (started.disposition !== 'turn_started') return; - const active = await client.queryTurn({ sessionId: fixture.sessionId, turnId: started.turnId }); - await client.stopTurn({ + const active = await client.request('turn.query', { + sessionId: fixture.sessionId, + turnId: started.turnId, + }); + await client.request('turn.stop', { sessionId: fixture.sessionId, turnId: started.turnId, runId: active.runId, @@ -268,7 +271,7 @@ test('a Host crash after queue admission recovers the durable successor once', a const firstHost = await fixture.startHost(); const first = await connectClient(fixture.root); const started = requireStartedTurn( - await first.startTurn({ + await first.request('turn.start', { sessionId: fixture.sessionId, turnId: randomUUID(), content: { text: FAKE_ASK_USER_QUESTION_PROMPT }, diff --git a/packages/storage/src/__tests__/sqlite-session-metadata-store.test.ts b/packages/storage/src/__tests__/sqlite-session-metadata-store.test.ts index 87d57957ba..f939ab6099 100644 --- a/packages/storage/src/__tests__/sqlite-session-metadata-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-session-metadata-store.test.ts @@ -54,6 +54,73 @@ import { import { SQLITE_AGENT_GRAPH_CONTROL_TABLES } from '../sqlite-session-metadata-schema.js'; describe('SqliteSessionMetadataStore', () => { + for (const version30Shape of ['admissions-only', 'coordination-only', 'complete'] as const) { + test(`converges the ${version30Shape} version-30 schema after the merge`, async () => { + const root = await mkdtemp(join(tmpdir(), `maka-session-v30-${version30Shape}-`)); + const path = join(root, 'state.sqlite'); + try { + const setup = createSqliteSessionMetadataStore(path); + setup.close(); + + const version30 = new DatabaseSync(path); + try { + if (version30Shape === 'admissions-only') { + version30.exec('DROP INDEX session_metadata_one_workhub_coordination_session'); + } else if (version30Shape === 'coordination-only') { + version30.exec(` + DROP TABLE cancelled_message_admissions; + DROP TABLE message_admissions; + `); + } + version30 + .prepare( + `UPDATE session_metadata_schema SET version = 30 WHERE scope = 'session_metadata'`, + ) + .run(); + } finally { + version30.close(); + } + + const converged = createSqliteSessionMetadataStore(path); + try { + assert.equal(converged.schemaVersion(), SQLITE_SESSION_METADATA_SCHEMA_VERSION); + } finally { + converged.close(); + } + + const schema = new DatabaseSync(path, { readOnly: true }); + try { + const objects = schema + .prepare( + ` + SELECT name + FROM sqlite_schema + WHERE name IN ( + 'message_admissions', + 'message_admissions_by_session_order', + 'cancelled_message_admissions', + 'session_metadata_one_workhub_coordination_session' + ) + ORDER BY name + `, + ) + .all() + .map((row) => (row as { name: string }).name); + assert.deepEqual(objects, [ + 'cancelled_message_admissions', + 'message_admissions', + 'message_admissions_by_session_order', + 'session_metadata_one_workhub_coordination_session', + ]); + } finally { + schema.close(); + } + } finally { + await rm(root, { recursive: true, force: true }); + } + }); + } + test('migrates v27 metadata to the current schema without backfilling external origin', async () => { const root = await mkdtemp(join(tmpdir(), 'maka-session-metadata-v27-')); const path = join(root, 'state.sqlite'); diff --git a/packages/storage/src/sqlite-session-metadata-schema.ts b/packages/storage/src/sqlite-session-metadata-schema.ts index 5e90ec82be..2cc54783af 100644 --- a/packages/storage/src/sqlite-session-metadata-schema.ts +++ b/packages/storage/src/sqlite-session-metadata-schema.ts @@ -19,7 +19,7 @@ import type { DatabaseSync } from 'node:sqlite'; -export const SQLITE_SESSION_METADATA_SCHEMA_VERSION = 30; +export const SQLITE_SESSION_METADATA_SCHEMA_VERSION = 31; export const SQLITE_SESSION_MESSAGE_CHUNK_BYTES = 64 * 1024; export const SQLITE_SESSION_MESSAGE_CHUNK_MARKER = '{"$maka":"session-message-chunks-v1"}'; @@ -1156,15 +1156,55 @@ const MIGRATIONS: ReadonlyMap = new Map([ `, ], [ - 30, + 31, ` - CREATE UNIQUE INDEX session_metadata_one_workhub_coordination_session + CREATE TABLE IF NOT EXISTS message_admissions ( + sequence INTEGER PRIMARY KEY AUTOINCREMENT, + session_id TEXT NOT NULL, + turn_id TEXT NOT NULL, + run_id TEXT NOT NULL, + message_id TEXT NOT NULL, + content_json TEXT NOT NULL, + submitted_content_digest TEXT NOT NULL, + submitted_placement TEXT NOT NULL + CHECK (submitted_placement IN ('current_turn', 'next_turn')), + placement TEXT NOT NULL CHECK (placement IN ('current_turn', 'next_turn')), + disposition TEXT NOT NULL CHECK (disposition IN ('steering', 'followup')), + queue_order INTEGER NOT NULL CHECK (queue_order >= 0), + admitted_at INTEGER NOT NULL CHECK (admitted_at >= 0), + UNIQUE (session_id, message_id), + FOREIGN KEY(session_id) REFERENCES session_metadata(session_id) ON DELETE CASCADE + ); + + CREATE INDEX IF NOT EXISTS message_admissions_by_session_order + ON message_admissions(session_id, queue_order, sequence); + + CREATE TABLE IF NOT EXISTS cancelled_message_admissions ( + session_id TEXT NOT NULL, + message_id TEXT NOT NULL, + submitted_content_digest TEXT NOT NULL, + submitted_placement TEXT NOT NULL + CHECK (submitted_placement IN ('current_turn', 'next_turn')), + PRIMARY KEY (session_id, message_id), + FOREIGN KEY(session_id) REFERENCES session_metadata(session_id) ON DELETE CASCADE + ); + + CREATE UNIQUE INDEX IF NOT EXISTS session_metadata_one_workhub_coordination_session ON session_metadata(json_extract(payload_json, '$.role')) WHERE json_extract(payload_json, '$.role') = 'workhub_coordination'; `, ], ]); +if (MIGRATIONS.size !== SQLITE_SESSION_METADATA_SCHEMA_VERSION) { + throw new Error('SQLite session metadata migrations contain a duplicate or missing version'); +} +for (let version = 1; version <= SQLITE_SESSION_METADATA_SCHEMA_VERSION; version += 1) { + if (!MIGRATIONS.has(version)) { + throw new Error(`Missing SQLite session metadata migration ${version}`); + } +} + export function configureSqliteSessionMetadataDatabase(db: DatabaseSync): void { db.exec('PRAGMA busy_timeout = 5000'); db.exec('PRAGMA journal_mode = WAL'); diff --git a/packages/storage/src/sqlite-session-metadata-store.ts b/packages/storage/src/sqlite-session-metadata-store.ts index f48aab82c8..8943dea71f 100644 --- a/packages/storage/src/sqlite-session-metadata-store.ts +++ b/packages/storage/src/sqlite-session-metadata-store.ts @@ -429,9 +429,15 @@ export class SqliteSessionMetadataStore { return; } const { DatabaseSync } = loadSqliteModule(); - this.db = new DatabaseSync(path); - configureSqliteSessionMetadataDatabase(this.db); - migrateSqliteSessionMetadataDatabase(this.db); + const database = new DatabaseSync(path); + try { + configureSqliteSessionMetadataDatabase(database); + migrateSqliteSessionMetadataDatabase(database); + } catch (error) { + database.close(); + throw error; + } + this.db = database; this.now = options.now ?? Date.now; }