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
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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 },
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down
46 changes: 43 additions & 3 deletions packages/storage/src/sqlite-session-metadata-schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"}';

Expand Down Expand Up @@ -1156,15 +1156,55 @@ const MIGRATIONS: ReadonlyMap<number, string> = 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');
Expand Down
12 changes: 9 additions & 3 deletions packages/storage/src/sqlite-session-metadata-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down