11import { createHash } from 'node:crypto' ;
2- import {
3- existsSync ,
4- mkdirSync ,
5- renameSync ,
6- } from 'node:fs' ;
2+ import { mkdirSync } from 'node:fs' ;
73import { dirname , join , resolve } from 'node:path' ;
84// node:sqlite emits an ExperimentalWarning on load (documented in the README):
95// the module is Node's built-in SQLite binding, stable enough for Node >= 22.13
@@ -103,6 +99,19 @@ const KERNEL_FORMAT = 1;
10399const COMPACTED_KERNEL_FORMAT = 2 ;
104100const READABLE_KERNEL_FORMATS : readonly number [ ] = Object . freeze ( [ KERNEL_FORMAT , COMPACTED_KERNEL_FORMAT ] ) ;
105101
102+ /**
103+ * Column order of every table `initialize` creates. A pre-existing table
104+ * whose columns differ is not this kernel's schema and fails closed as
105+ * `corrupt` instead of surfacing a raw SQLite error on the first statement
106+ * that names a missing column.
107+ */
108+ const TABLE_COLUMNS : Readonly < Record < string , readonly string [ ] > > = Object . freeze ( {
109+ agent_state_head : [ 'id' , 'revision' , 'state' ] ,
110+ agent_state_journal : [ 'revision' , 'kind' , 'name' , 'payload' , 'state' , 'result_state' , 'to_version' , 'idempotency_key' , 'committed_at' ] ,
111+ agent_state_meta : [ 'id' , 'definition_id' , 'schema_version' , 'kernel_format' ] ,
112+ agent_state_pruned_keys : [ 'idempotency_key' , 'revision' , 'canonical_input' ] ,
113+ } ) ;
114+
106115export interface SqliteStateDriverOptions {
107116 /**
108117 * SQLite lock wait budget per operation in milliseconds (default 5000).
@@ -237,7 +246,7 @@ interface JournalRow {
237246 readonly kind : string ;
238247 readonly name : string | null ;
239248 readonly payload : string | null ;
240- readonly result_state : string | null ;
249+ readonly result_state : string ;
241250 readonly revision : number ;
242251 readonly state : string | null ;
243252 readonly to_version : number | null ;
@@ -282,9 +291,6 @@ const recordFromRow = (definitionId: string, row: JournalRow): AgentStateJournal
282291const sanitizedFileName = ( definitionId : string ) : string =>
283292 `${ definitionId . replace ( / [ ^ a - z A - Z 0 - 9 . _ - ] + / gu, '-' ) } -${ createHash ( 'sha256' ) . update ( definitionId , 'utf8' ) . digest ( 'hex' ) . slice ( 0 , 16 ) } .sqlite` ;
284293
285- const legacySanitizedFileName = ( definitionId : string ) : string =>
286- `${ definitionId . replace ( / [ ^ a - z A - Z 0 - 9 . _ - ] + / gu, '-' ) } -${ Buffer . from ( definitionId , 'utf8' ) . toString ( 'hex' ) . slice ( 0 , 12 ) } .sqlite` ;
287-
288294class SqliteConnection extends Context . Service < SqliteConnection , DatabaseSync > ( ) (
289295 '@agent-bundle/runtime/state/SqliteConnection' ,
290296) { }
@@ -413,14 +419,14 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
413419 db : DatabaseSync ,
414420 key : string ,
415421 ) :
416- | { readonly kind : 'committed' ; readonly record : AgentStateJournalRecord ; readonly resultStateText : string | null }
422+ | { readonly kind : 'committed' ; readonly record : AgentStateJournalRecord ; readonly resultStateText : string }
417423 | { readonly canonicalInput : string ; readonly kind : 'pruned' ; readonly revision : number }
418424 | undefined {
419425 const row = this . #prepare( db , 'SELECT * FROM agent_state_journal WHERE idempotency_key = ?' ) . get ( key ) as
420426 | JournalRow
421427 | undefined ;
422428 if ( row !== undefined ) {
423- return { kind : 'committed' , record : recordFromRow ( this . #definition. id , row ) , resultStateText : row . result_state ?? row . state } ;
429+ return { kind : 'committed' , record : recordFromRow ( this . #definition. id , row ) , resultStateText : row . result_state } ;
424430 }
425431 const pruned = this
426432 . #prepare( db , 'SELECT revision, canonical_input FROM agent_state_pruned_keys WHERE idempotency_key = ?' )
@@ -433,17 +439,12 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
433439 /**
434440 * Recovers the state a committed record produced. Every record stores its
435441 * post-commit state (migrated forward on schema migrations), so replay
436- * does not depend on exact-revision history; rows written before post-
437- * commit states were stored fall back to journal replay.
442+ * does not depend on exact-revision history.
438443 */
439444 #committedState(
440- db : DatabaseSync ,
441- committed : { readonly kind : 'committed' ; readonly record : AgentStateJournalRecord ; readonly resultStateText : string | null } ,
445+ committed : { readonly kind : 'committed' ; readonly record : AgentStateJournalRecord ; readonly resultStateText : string } ,
442446 ) : TState {
443- const raw =
444- committed . resultStateText !== null
445- ? parseStoredJson ( this . #definition. id , 'result state' , committed . record . revision , committed . resultStateText )
446- : this . #replayTo( db , committed . record . revision ) ;
447+ const raw = parseStoredJson ( this . #definition. id , 'result state' , committed . record . revision , committed . resultStateText ) ;
447448 const parsed = this . #definition. schema . safeParse ( raw ) ;
448449 if ( ! parsed . success ) {
449450 throw new AgentStateError (
@@ -572,7 +573,7 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
572573 return Object . freeze ( {
573574 replayed : true ,
574575 revision : committed . record . revision ,
575- state : this . #committedState( db , committed ) ,
576+ state : this . #committedState( committed ) ,
576577 } ) ;
577578 }
578579 const head = this . #headState( db , 'commit' ) ;
@@ -829,7 +830,7 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
829830 name TEXT,
830831 payload TEXT,
831832 state TEXT,
832- result_state TEXT,
833+ result_state TEXT NOT NULL ,
833834 to_version INTEGER,
834835 idempotency_key TEXT NOT NULL UNIQUE,
835836 committed_at TEXT NOT NULL
@@ -845,13 +846,17 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
845846 canonical_input TEXT NOT NULL
846847 );
847848 ` ) ;
848- const journalColumns = transactionDb . prepare ( 'PRAGMA table_info(agent_state_journal)' ) . all ( ) as unknown as {
849- readonly name : string ;
850- } [ ] ;
851- if ( ! journalColumns . some ( ( column ) => column . name === 'result_state' ) ) {
852- transactionDb . exec ( 'ALTER TABLE agent_state_journal ADD COLUMN result_state TEXT' ) ;
853- }
854849 const definition = this . #definition;
850+ for ( const [ table , columns ] of Object . entries ( TABLE_COLUMNS ) ) {
851+ const actual = ( transactionDb . prepare ( `PRAGMA table_info(${ table } )` ) . all ( ) as unknown as { readonly name : string } [ ] )
852+ . map ( ( column ) => column . name ) ;
853+ if ( actual . join ( ',' ) !== columns . join ( ',' ) ) {
854+ throw new AgentStateError (
855+ 'corrupt' ,
856+ `State '${ definition . id } ' table ${ table } at '${ this . location } ' does not match the current schema (has ${ actual . join ( ', ' ) } ; expected ${ columns . join ( ', ' ) } )` ,
857+ ) ;
858+ }
859+ }
855860 const meta = transactionDb
856861 . prepare ( 'SELECT definition_id, schema_version, kernel_format FROM agent_state_meta WHERE id = 1' )
857862 . get ( ) as { definition_id : string ; kernel_format : number ; schema_version : number } | undefined ;
@@ -906,10 +911,10 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
906911 }
907912 // A pending migration cannot replay records written under the older
908913 // definition; verify the head against the last stored post-commit
909- // state instead (rows predating stored event states leave it null) .
910- const lastStateText = rows . length === 0 ? null : ( rows [ rows . length - 1 ] as JournalRow ) . state ;
911- if ( lastStateText !== null ) {
912- const lastState = parseStoredJson ( definition . id , 'state' , journalHead , lastStateText ) ;
914+ // state instead.
915+ const lastRow = rows [ rows . length - 1 ] ;
916+ if ( lastRow !== undefined ) {
917+ const lastState = parseStoredJson ( definition . id , 'result state' , journalHead , lastRow . result_state ) ;
913918 if ( canonicalJson ( lastState ) !== canonicalJson ( rawHead ) ) {
914919 throw new AgentStateError (
915920 'corrupt' ,
@@ -922,27 +927,13 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
922927 expectMigrationWithinStateBudget ( definition , migratedStateText ) ;
923928 // Journal records retain the original commit input for dedupe. Their
924929 // committed results migrate separately, matching the memory driver's
925- // `{ record, state }` split. A legacy journal-head result can be
926- // recovered from the authoritative materialized head. Earlier missing
927- // results cannot be reconstructed with the current-version reducer.
930+ // `{ record, state }` split.
928931 const updateResult = transactionDb . prepare (
929932 'UPDATE agent_state_journal SET result_state = ? WHERE revision = ?' ,
930933 ) ;
931934 for ( const row of rows ) {
932- const storedResultText = row . result_state ?? row . state ;
933- let migratedResult : TState ;
934- if ( storedResultText !== null ) {
935- const storedResult = parseStoredJson ( definition . id , 'result state' , row . revision , storedResultText ) ;
936- migratedResult = runStateMigrations ( definition , meta . schema_version , storedResult ) ;
937- } else if ( row . revision === journalHead ) {
938- migratedResult = migrated ;
939- } else {
940- throw new AgentStateError (
941- 'migration-failure' ,
942- `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 ) } ` ,
943- ) ;
944- }
945- updateResult . run ( canonicalJson ( migratedResult ) , row . revision ) ;
935+ const storedResult = parseStoredJson ( definition . id , 'result state' , row . revision , row . result_state ) ;
936+ updateResult . run ( canonicalJson ( runStateMigrations ( definition , meta . schema_version , storedResult ) ) , row . revision ) ;
946937 }
947938 const record : AgentStateJournalRecord = {
948939 committedAt : this . #now( ) . toISOString ( ) ,
@@ -1043,32 +1034,7 @@ export const createSqliteStateDriver = (options: SqliteStateDriverOptions): Agen
10431034 ) ;
10441035 }
10451036 if ( options . file !== undefined ) return resolve ( options . file ) ;
1046- const root = options . root as string ;
1047- const currentFile = resolve ( join ( root , sanitizedFileName ( definition . id ) ) ) ;
1048- const legacyFile = resolve ( join ( root , legacySanitizedFileName ( definition . id ) ) ) ;
1049- mkdirSync ( dirname ( currentFile ) , { recursive : true } ) ;
1050- if ( ! existsSync ( currentFile ) && existsSync ( legacyFile ) ) {
1051- for ( const suffix of [ '-wal' , '-shm' ] ) {
1052- const legacySidecar = `${ legacyFile } ${ suffix } ` ;
1053- if ( ! existsSync ( legacySidecar ) ) continue ;
1054- try {
1055- renameSync ( legacySidecar , `${ currentFile } ${ suffix } ` ) ;
1056- } catch ( error ) {
1057- // A concurrent adopter may have moved this sidecar after
1058- // the existence check. Other failures must remain visible.
1059- if ( ( error as SqliteErrorShape ) . code !== 'ENOENT' ) throw error ;
1060- }
1061- }
1062- try {
1063- renameSync ( legacyFile , currentFile ) ;
1064- } catch ( error ) {
1065- // Another opener may have atomically adopted the same
1066- // legacy file after both observed it. The winner's current
1067- // path is authoritative; otherwise preserve the failure.
1068- if ( ! existsSync ( currentFile ) ) throw error ;
1069- }
1070- }
1071- return currentFile ;
1037+ return resolve ( join ( options . root as string , sanitizedFileName ( definition . id ) ) ) ;
10721038 } , true ) ,
10731039 ) ;
10741040 const connection = Effect . acquireRelease (
0 commit comments