diff --git a/.changeset/remove-legacy-sqlite-state.md b/.changeset/remove-legacy-sqlite-state.md new file mode 100644 index 000000000..3b1ba392f --- /dev/null +++ b/.changeset/remove-legacy-sqlite-state.md @@ -0,0 +1,5 @@ +--- +"@agent-bundle/runtime": minor +--- + +Stop adopting pre-#201 durable SQLite state stores in `createSqliteStateDriver`: a root-mode store still named `-<12 hex of utf8(id)>.sqlite` is left on disk unread and a fresh `-.sqlite` store opens beside it, and a journal table without a `result_state` column is no longer upgraded in place but rejected on open with a typed `corrupt` error. Delete or recreate old stores before upgrading. (#837) diff --git a/packages/rsc-runtime/src/state/sqlite.ts b/packages/rsc-runtime/src/state/sqlite.ts index c9b4c2292..ea9c5c437 100644 --- a/packages/rsc-runtime/src/state/sqlite.ts +++ b/packages/rsc-runtime/src/state/sqlite.ts @@ -1,9 +1,5 @@ import { createHash } from 'node:crypto'; -import { - existsSync, - mkdirSync, - renameSync, -} from 'node:fs'; +import { mkdirSync } from 'node:fs'; import { dirname, join, resolve } from 'node:path'; // node:sqlite emits an ExperimentalWarning on load (documented in the README): // the module is Node's built-in SQLite binding, stable enough for Node >= 22.13 @@ -103,6 +99,19 @@ const KERNEL_FORMAT = 1; const COMPACTED_KERNEL_FORMAT = 2; const READABLE_KERNEL_FORMATS: readonly number[] = Object.freeze([KERNEL_FORMAT, COMPACTED_KERNEL_FORMAT]); +/** + * Column order of every table `initialize` creates. A pre-existing table + * whose columns differ is not this kernel's schema and fails closed as + * `corrupt` instead of surfacing a raw SQLite error on the first statement + * that names a missing column. + */ +const TABLE_COLUMNS: Readonly> = Object.freeze({ + agent_state_head: ['id', 'revision', 'state'], + agent_state_journal: ['revision', 'kind', 'name', 'payload', 'state', 'result_state', 'to_version', 'idempotency_key', 'committed_at'], + agent_state_meta: ['id', 'definition_id', 'schema_version', 'kernel_format'], + agent_state_pruned_keys: ['idempotency_key', 'revision', 'canonical_input'], +}); + export interface SqliteStateDriverOptions { /** * SQLite lock wait budget per operation in milliseconds (default 5000). @@ -237,7 +246,7 @@ interface JournalRow { readonly kind: string; readonly name: string | null; readonly payload: string | null; - readonly result_state: string | null; + readonly result_state: string; readonly revision: number; readonly state: string | null; readonly to_version: number | null; @@ -282,9 +291,6 @@ const recordFromRow = (definitionId: string, row: JournalRow): AgentStateJournal const sanitizedFileName = (definitionId: string): string => `${definitionId.replace(/[^a-zA-Z0-9._-]+/gu, '-')}-${createHash('sha256').update(definitionId, 'utf8').digest('hex').slice(0, 16)}.sqlite`; -const legacySanitizedFileName = (definitionId: string): string => - `${definitionId.replace(/[^a-zA-Z0-9._-]+/gu, '-')}-${Buffer.from(definitionId, 'utf8').toString('hex').slice(0, 12)}.sqlite`; - class SqliteConnection extends Context.Service()( '@agent-bundle/runtime/state/SqliteConnection', ) {} @@ -413,14 +419,14 @@ class SqliteStore implements Age db: DatabaseSync, key: string, ): - | { readonly kind: 'committed'; readonly record: AgentStateJournalRecord; readonly resultStateText: string | null } + | { readonly kind: 'committed'; readonly record: AgentStateJournalRecord; readonly resultStateText: string } | { readonly canonicalInput: string; readonly kind: 'pruned'; readonly revision: number } | undefined { const row = this.#prepare(db, 'SELECT * FROM agent_state_journal WHERE idempotency_key = ?').get(key) as | JournalRow | undefined; if (row !== undefined) { - return { kind: 'committed', record: recordFromRow(this.#definition.id, row), resultStateText: row.result_state ?? row.state }; + return { kind: 'committed', record: recordFromRow(this.#definition.id, row), resultStateText: row.result_state }; } const pruned = this .#prepare(db, 'SELECT revision, canonical_input FROM agent_state_pruned_keys WHERE idempotency_key = ?') @@ -433,17 +439,12 @@ class SqliteStore implements Age /** * Recovers the state a committed record produced. Every record stores its * post-commit state (migrated forward on schema migrations), so replay - * does not depend on exact-revision history; rows written before post- - * commit states were stored fall back to journal replay. + * does not depend on exact-revision history. */ #committedState( - db: DatabaseSync, - committed: { readonly kind: 'committed'; readonly record: AgentStateJournalRecord; readonly resultStateText: string | null }, + committed: { readonly kind: 'committed'; readonly record: AgentStateJournalRecord; readonly resultStateText: string }, ): TState { - const raw = - committed.resultStateText !== null - ? parseStoredJson(this.#definition.id, 'result state', committed.record.revision, committed.resultStateText) - : this.#replayTo(db, committed.record.revision); + const raw = parseStoredJson(this.#definition.id, 'result state', committed.record.revision, committed.resultStateText); const parsed = this.#definition.schema.safeParse(raw); if (!parsed.success) { throw new AgentStateError( @@ -572,7 +573,7 @@ class SqliteStore implements Age return Object.freeze({ replayed: true, revision: committed.record.revision, - state: this.#committedState(db, committed), + state: this.#committedState(committed), }); } const head = this.#headState(db, 'commit'); @@ -829,7 +830,7 @@ class SqliteStore implements Age name TEXT, payload TEXT, state TEXT, - result_state TEXT, + result_state TEXT NOT NULL, to_version INTEGER, idempotency_key TEXT NOT NULL UNIQUE, committed_at TEXT NOT NULL @@ -845,13 +846,17 @@ class SqliteStore implements Age canonical_input TEXT NOT NULL ); `); - const journalColumns = transactionDb.prepare('PRAGMA table_info(agent_state_journal)').all() as unknown as { - readonly name: string; - }[]; - if (!journalColumns.some((column) => column.name === 'result_state')) { - transactionDb.exec('ALTER TABLE agent_state_journal ADD COLUMN result_state TEXT'); - } const definition = this.#definition; + for (const [table, columns] of Object.entries(TABLE_COLUMNS)) { + const actual = (transactionDb.prepare(`PRAGMA table_info(${table})`).all() as unknown as { readonly name: string }[]) + .map((column) => column.name); + if (actual.join(',') !== columns.join(',')) { + throw new AgentStateError( + 'corrupt', + `State '${definition.id}' table ${table} at '${this.location}' does not match the current schema (has ${actual.join(', ')}; expected ${columns.join(', ')})`, + ); + } + } const meta = transactionDb .prepare('SELECT definition_id, schema_version, kernel_format FROM agent_state_meta WHERE id = 1') .get() as { definition_id: string; kernel_format: number; schema_version: number } | undefined; @@ -906,10 +911,10 @@ class SqliteStore implements Age } // A pending migration cannot replay records written under the older // definition; verify the head against the last stored post-commit - // state instead (rows predating stored event states leave it null). - const lastStateText = rows.length === 0 ? null : (rows[rows.length - 1] as JournalRow).state; - if (lastStateText !== null) { - const lastState = parseStoredJson(definition.id, 'state', journalHead, lastStateText); + // state instead. + const lastRow = rows[rows.length - 1]; + if (lastRow !== undefined) { + const lastState = parseStoredJson(definition.id, 'result state', journalHead, lastRow.result_state); if (canonicalJson(lastState) !== canonicalJson(rawHead)) { throw new AgentStateError( 'corrupt', @@ -922,27 +927,13 @@ class SqliteStore implements Age expectMigrationWithinStateBudget(definition, migratedStateText); // Journal records retain the original commit input for dedupe. Their // committed results migrate separately, matching the memory driver's - // `{ record, state }` split. A legacy journal-head result can be - // recovered from the authoritative materialized head. Earlier missing - // results cannot be reconstructed with the current-version reducer. + // `{ record, state }` split. const updateResult = transactionDb.prepare( 'UPDATE agent_state_journal SET result_state = ? WHERE revision = ?', ); for (const row of rows) { - const storedResultText = row.result_state ?? row.state; - let migratedResult: TState; - if (storedResultText !== null) { - const storedResult = parseStoredJson(definition.id, 'result state', row.revision, storedResultText); - migratedResult = runStateMigrations(definition, meta.schema_version, storedResult); - } else if (row.revision === journalHead) { - migratedResult = migrated; - } else { - throw new AgentStateError( - 'migration-failure', - `State '${definition.id}' legacy journal row at revision ${String(row.revision)} has no recoverable committed result; restore a compatible backup or materialize the result with the version ${String(meta.schema_version)} definition before migrating to version ${String(definition.version)}`, - ); - } - updateResult.run(canonicalJson(migratedResult), row.revision); + const storedResult = parseStoredJson(definition.id, 'result state', row.revision, row.result_state); + updateResult.run(canonicalJson(runStateMigrations(definition, meta.schema_version, storedResult)), row.revision); } const record: AgentStateJournalRecord = { committedAt: this.#now().toISOString(), @@ -1043,32 +1034,7 @@ export const createSqliteStateDriver = (options: SqliteStateDriverOptions): Agen ); } if (options.file !== undefined) return resolve(options.file); - const root = options.root as string; - const currentFile = resolve(join(root, sanitizedFileName(definition.id))); - const legacyFile = resolve(join(root, legacySanitizedFileName(definition.id))); - mkdirSync(dirname(currentFile), { recursive: true }); - if (!existsSync(currentFile) && existsSync(legacyFile)) { - for (const suffix of ['-wal', '-shm']) { - const legacySidecar = `${legacyFile}${suffix}`; - if (!existsSync(legacySidecar)) continue; - try { - renameSync(legacySidecar, `${currentFile}${suffix}`); - } catch (error) { - // A concurrent adopter may have moved this sidecar after - // the existence check. Other failures must remain visible. - if ((error as SqliteErrorShape).code !== 'ENOENT') throw error; - } - } - try { - renameSync(legacyFile, currentFile); - } catch (error) { - // Another opener may have atomically adopted the same - // legacy file after both observed it. The winner's current - // path is authoritative; otherwise preserve the failure. - if (!existsSync(currentFile)) throw error; - } - } - return currentFile; + return resolve(join(options.root as string, sanitizedFileName(definition.id))); }, true), ); const connection = Effect.acquireRelease( diff --git a/packages/rsc-runtime/tests/state-sqlite.test.ts b/packages/rsc-runtime/tests/state-sqlite.test.ts index 4d5e002fb..09531b1a1 100644 --- a/packages/rsc-runtime/tests/state-sqlite.test.ts +++ b/packages/rsc-runtime/tests/state-sqlite.test.ts @@ -1,5 +1,4 @@ -import { createHash } from 'node:crypto'; -import { access, mkdtemp, writeFile } from 'node:fs/promises'; +import { mkdtemp, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { DatabaseSync } from 'node:sqlite'; @@ -94,116 +93,6 @@ const migratingCounterDefinition = ( version: 2, }); -const legacyFileName = (definitionId: string): string => - `${definitionId.replace(/[^a-zA-Z0-9._-]+/gu, '-')}-${Buffer.from(definitionId, 'utf8').toString('hex').slice(0, 12)}.sqlite`; - -const currentFileName = (definitionId: string): string => - `${definitionId.replace(/[^a-zA-Z0-9._-]+/gu, '-')}-${createHash('sha256').update(definitionId, 'utf8').digest('hex').slice(0, 16)}.sqlite`; - -const createLegacyMigrationDatabase = (file: string, definitionId: string): void => { - const db = new DatabaseSync(file); - try { - db.exec(` - CREATE TABLE agent_state_meta ( - id INTEGER PRIMARY KEY CHECK (id = 1), - definition_id TEXT NOT NULL, - schema_version INTEGER NOT NULL, - kernel_format INTEGER NOT NULL - ); - CREATE TABLE agent_state_journal ( - revision INTEGER PRIMARY KEY, - kind TEXT NOT NULL CHECK (kind IN ('event', 'reset', 'migrate')), - name TEXT, - payload TEXT, - state TEXT, - to_version INTEGER, - idempotency_key TEXT NOT NULL UNIQUE, - committed_at TEXT NOT NULL - ); - CREATE TABLE agent_state_head ( - id INTEGER PRIMARY KEY CHECK (id = 1), - revision INTEGER NOT NULL, - state TEXT NOT NULL - ); - `); - db.prepare( - 'INSERT INTO agent_state_meta (id, definition_id, schema_version, kernel_format) VALUES (1, ?, 1, 1)', - ).run(definitionId); - const insert = db.prepare( - 'INSERT INTO agent_state_journal (revision, kind, name, payload, state, to_version, idempotency_key, committed_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)', - ); - insert.run(1, 'event', 'bumped', '{"by":2}', '{"count":2}', null, 'legacy:event', '2026-01-01T00:00:00.000Z'); - insert.run(2, 'reset', null, null, '{"count":5}', null, 'legacy:reset', '2026-01-01T00:00:01.000Z'); - insert.run(3, 'event', 'bumped', '{"by":1}', '{"count":6}', null, 'legacy:event-2', '2026-01-01T00:00:02.000Z'); - db.prepare('INSERT INTO agent_state_head (id, revision, state) VALUES (1, 3, ?)').run('{"count":6}'); - } finally { - db.close(); - } -}; - -const clearLegacyEventResults = (file: string): void => { - const db = new DatabaseSync(file); - try { - db.exec("UPDATE agent_state_journal SET state = NULL WHERE kind = 'event'"); - } finally { - db.close(); - } -}; - -interface ValueState { - readonly value: number; -} - -const valueCounterDefinition = ( - id = 'state-sqlite-test/value-counter', -): AgentStateDefinition => - defineState({ - events: counterEvents, - id, - initial: { value: 0 }, - lifetime: 'workspace-durable', - migrations: { - 2: (persisted) => ({ value: (persisted as CounterState).count * 10 }), - }, - reduce: (state, event) => ({ value: state.value + event.payload.by }), - schema: z.object({ value: z.number().int() }).strict(), - version: 2, - }); - -const createLegacyHeadOnlyDatabase = (file: string, definitionId: string): void => { - createLegacyMigrationDatabase(file, definitionId); - const db = new DatabaseSync(file); - try { - db.exec(` - DELETE FROM agent_state_journal WHERE revision > 1; - UPDATE agent_state_journal SET state = NULL WHERE revision = 1; - UPDATE agent_state_head SET revision = 1, state = '{"count":2}' WHERE id = 1; - `); - } finally { - db.close(); - } -}; - -const holdUncheckpointedLegacyEvent = (file: string): DatabaseSync => { - const keeper = new DatabaseSync(file); - keeper.exec('PRAGMA journal_mode = WAL; PRAGMA wal_autocheckpoint = 0; BEGIN DEFERRED'); - keeper.prepare('SELECT revision FROM agent_state_head WHERE id = 1').get(); - const writer = new DatabaseSync(file); - try { - writer.exec('PRAGMA wal_autocheckpoint = 0; BEGIN IMMEDIATE'); - writer - .prepare( - 'INSERT INTO agent_state_journal (revision, kind, name, payload, state, to_version, idempotency_key, committed_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)', - ) - .run(4, 'event', 'bumped', '{"by":1}', null, null, 'legacy:event-3', '2026-01-01T00:00:03.000Z'); - writer.prepare('UPDATE agent_state_head SET revision = 4, state = ? WHERE id = 1').run('{"count":7}'); - writer.exec('COMMIT'); - } finally { - writer.close(); - } - return keeper; -}; - const otherDefinition = (): AgentStateDefinition => defineState({ events: counterEvents, @@ -256,75 +145,26 @@ describe('sqlite driver storage behavior', () => { await expect(first.read()).rejects.toMatchObject({ code: 'store-closed' }); })); - it('adopts the legacy root database and live WAL sidecars without data loss', () => + it('keeps the original commit input while migrating each committed result', () => withRoot(async (root) => { - const definition = migratingCounterDefinition(); - const legacyFile = join(root, legacyFileName(definition.id)); - const currentFile = join(root, currentFileName(definition.id)); - createLegacyMigrationDatabase(legacyFile, definition.id); - const keeper = holdUncheckpointedLegacyEvent(legacyFile); - try { - await access(`${legacyFile}-wal`); - await access(`${legacyFile}-shm`); - - const driver = createSqliteStateDriver({ root }); - const store = await driver.open(definition); - - expect(store.location).toBe(currentFile); - expect(await store.read()).toEqual({ revision: 5, state: { count: 70 } }); - await expect(access(legacyFile)).rejects.toMatchObject({ code: 'ENOENT' }); - await expect(access(`${legacyFile}-wal`)).rejects.toMatchObject({ code: 'ENOENT' }); - await expect(access(`${legacyFile}-shm`)).rejects.toMatchObject({ code: 'ENOENT' }); - await access(`${currentFile}-wal`); - await access(`${currentFile}-shm`); - await driver.close(); - } finally { - keeper.exec('ROLLBACK'); - keeper.close(); - } - })); - - it('fails closed when a non-head legacy event has no recoverable committed result', () => - withRoot(async (root) => { - const definition = migratingCounterDefinition(); - const file = join(root, 'legacy-null-event.sqlite'); - createLegacyMigrationDatabase(file, definition.id); - clearLegacyEventResults(file); - - await expect(createSqliteStateDriver({ file }).open(definition)).rejects.toMatchObject({ - code: 'migration-failure', - message: expect.stringContaining('has no recoverable committed result'), - name: 'AgentStateError', - }); - })); - - it('migrates a legacy journal-head result from the materialized head without using the current reducer', () => - withRoot(async (root) => { - const definition = valueCounterDefinition(); - const file = join(root, 'legacy-head-event.sqlite'); - createLegacyHeadOnlyDatabase(file, definition.id); - - const store = await createSqliteStateDriver({ file }).open(definition); + const file = join(root, 'migrating-reset.sqlite'); + const v1 = await createSqliteStateDriver({ file }).open(counterDefinition('state-sqlite-test/migrating-counter')); + await v1.dispatch('bumped', { by: 2 }, { idempotencyKey: 'v1:event' }); + await v1.reset({ idempotencyKey: 'v1:reset', seed: { count: 5 } }); + await v1.close(); + + const store = await createSqliteStateDriver({ file }).open(migratingCounterDefinition()); + await expect(store.read()).resolves.toEqual({ revision: 3, state: { count: 50 } }); await expect( - store.dispatch('bumped', { by: 2 }, { idempotencyKey: 'legacy:event' }), - ).resolves.toEqual({ replayed: true, revision: 1, state: { value: 20 } }); - await expect(store.read()).resolves.toEqual({ revision: 2, state: { value: 20 } }); - await store.close(); - })); - - it('preserves legacy reset input while migrating its idempotent result', () => - withRoot(async (root) => { - const definition = migratingCounterDefinition(); - const file = join(root, 'legacy-reset.sqlite'); - createLegacyMigrationDatabase(file, definition.id); - - const store = await createSqliteStateDriver({ file }).open(definition); + store.dispatch('bumped', { by: 2 }, { idempotencyKey: 'v1:event' }), + ).resolves.toEqual({ replayed: true, revision: 1, state: { count: 20 } }); await expect( - store.reset({ idempotencyKey: 'legacy:reset', seed: { count: 5 } }), + store.reset({ idempotencyKey: 'v1:reset', seed: { count: 5 } }), ).resolves.toEqual({ replayed: true, revision: 2, state: { count: 50 } }); const db = new DatabaseSync(file); try { - expect(db.prepare('SELECT state FROM agent_state_journal WHERE revision = 2').get()).toEqual({ + expect(db.prepare('SELECT state, result_state FROM agent_state_journal WHERE revision = 2').get()).toEqual({ + result_state: '{"count":50}', state: '{"count":5}', }); } finally { @@ -333,6 +173,22 @@ describe('sqlite driver storage behavior', () => { await store.close(); })); + it('fails closed with a typed corrupt error when a table does not match the current schema', () => + withRoot(async (root) => { + const file = join(root, 'state.sqlite'); + const store = await createSqliteStateDriver({ file }).open(counterDefinition()); + await store.dispatch('bumped', { by: 1 }, { idempotencyKey: 'k1' }); + await store.close(); + const db = new DatabaseSync(file); + db.exec('ALTER TABLE agent_state_journal DROP COLUMN result_state'); + db.close(); + await expect(createSqliteStateDriver({ file }).open(counterDefinition())).rejects.toMatchObject({ + code: 'corrupt', + message: expect.stringContaining('table agent_state_journal') as string, + name: 'AgentStateError', + }); + })); + it('rejects a pending open when the driver closes before initialization resumes', () => withRoot(async (root) => { const driver = createSqliteStateDriver({ root });