Skip to content

Commit 941aa08

Browse files
fix(state): preserve legacy results and harden driver close (#201 follow-ups) (#213)
1 parent 7dbfacf commit 941aa08

3 files changed

Lines changed: 145 additions & 34 deletions

File tree

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
---
2+
"@agent-bundle/runtime": patch
3+
---
4+
5+
Fail closed instead of replaying unrecoverable legacy state with a newer
6+
reducer, recover journal-head results from the materialized sqlite head, and
7+
preserve lifecycle errors while closing every open sqlite store.

‎packages/rsc-runtime/src/state/sqlite.ts‎

Lines changed: 23 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,6 @@ import type {
3131
AgentStateDefinition,
3232
AgentStateDispatchOptions,
3333
AgentStateDriver,
34-
AgentStateEvent,
3534
AgentStateEventSchemas,
3635
AgentStateJournalRecord,
3736
AgentStateReadOptions,
@@ -699,38 +698,24 @@ class SqliteStore<TState, TEvents extends AgentStateEventSchemas> implements Age
699698
const migrated = runStateMigrations(definition, meta.schema_version, rawHead);
700699
// Journal records retain the original commit input for dedupe. Their
701700
// committed results migrate separately, matching the memory driver's
702-
// `{ record, state }` split. Legacy event rows without a result are
703-
// replayed before the new migration baseline makes old revisions
704-
// unavailable.
701+
// `{ record, state }` split. A legacy journal-head result can be
702+
// recovered from the authoritative materialized head. Earlier missing
703+
// results cannot be reconstructed with the current-version reducer.
705704
const updateResult = transactionDb.prepare(
706705
'UPDATE agent_state_journal SET result_state = ? WHERE revision = ?',
707706
);
708-
let replayState: unknown = definition.initial;
709-
for (const [index, row] of rows.entries()) {
710-
const record = records[index] as AgentStateJournalRecord;
707+
for (const row of rows) {
711708
const storedResultText = row.result_state ?? row.state;
712709
let migratedResult: TState;
713710
if (storedResultText !== null) {
714-
replayState = parseStoredJson(definition.id, 'result state', row.revision, storedResultText);
715-
migratedResult = runStateMigrations(definition, meta.schema_version, replayState);
716-
} else if (record.kind === 'event') {
717-
try {
718-
replayState = definition.reduce(
719-
replayState as TState,
720-
{ name: record.name, payload: record.payload } as AgentStateEvent<TEvents>,
721-
);
722-
} catch (error) {
723-
throw new AgentStateError(
724-
'migration-failure',
725-
`State '${definition.id}' could not recover legacy result at revision ${String(record.revision)}`,
726-
{ cause: error },
727-
);
728-
}
729-
migratedResult = runStateMigrations(definition, meta.schema_version, replayState);
711+
const storedResult = parseStoredJson(definition.id, 'result state', row.revision, storedResultText);
712+
migratedResult = runStateMigrations(definition, meta.schema_version, storedResult);
713+
} else if (row.revision === journalHead) {
714+
migratedResult = migrated;
730715
} else {
731716
throw new AgentStateError(
732-
'corrupt',
733-
`State '${definition.id}' journal row at revision ${String(record.revision)} has no committed result`,
717+
'migration-failure',
718+
`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)}`,
734719
);
735720
}
736721
updateResult.run(canonicalJson(migratedResult), row.revision);
@@ -787,10 +772,16 @@ export const createSqliteStateDriver = (options: SqliteStateDriverOptions): Agen
787772
closed = true;
788773
closing = (async () => {
789774
await pendingOpens.settle();
775+
const closeErrors: unknown[] = [];
790776
for (const store of [...openStores]) {
791-
await store.close();
777+
try {
778+
await store.close();
779+
} catch (error) {
780+
closeErrors.push(error);
781+
}
792782
}
793783
openStores.clear();
784+
if (closeErrors.length > 0) throw closeErrors[0];
794785
})();
795786
return closing;
796787
},
@@ -869,7 +860,12 @@ export const createSqliteStateDriver = (options: SqliteStateDriverOptions): Agen
869860
);
870861
await store.initialize(busyTimeoutMs);
871862
} catch (error) {
872-
await runtime.close();
863+
try {
864+
await runtime.close();
865+
} catch {
866+
// Initialization is the caller-visible failure. The scoped
867+
// finalizer still attempts close, but must not replace it.
868+
}
873869
throw error;
874870
}
875871
if (closed) {

‎packages/rsc-runtime/tests/state-sqlite.test.ts‎

Lines changed: 115 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -131,15 +131,58 @@ const createLegacyMigrationDatabase = (file: string, definitionId: string): void
131131
const insert = db.prepare(
132132
'INSERT INTO agent_state_journal (revision, kind, name, payload, state, to_version, idempotency_key, committed_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
133133
);
134-
insert.run(1, 'event', 'bumped', '{"by":2}', null, null, 'legacy:event', '2026-01-01T00:00:00.000Z');
134+
insert.run(1, 'event', 'bumped', '{"by":2}', '{"count":2}', null, 'legacy:event', '2026-01-01T00:00:00.000Z');
135135
insert.run(2, 'reset', null, null, '{"count":5}', null, 'legacy:reset', '2026-01-01T00:00:01.000Z');
136-
insert.run(3, 'event', 'bumped', '{"by":1}', null, null, 'legacy:event-2', '2026-01-01T00:00:02.000Z');
136+
insert.run(3, 'event', 'bumped', '{"by":1}', '{"count":6}', null, 'legacy:event-2', '2026-01-01T00:00:02.000Z');
137137
db.prepare('INSERT INTO agent_state_head (id, revision, state) VALUES (1, 3, ?)').run('{"count":6}');
138138
} finally {
139139
db.close();
140140
}
141141
};
142142

143+
const clearLegacyEventResults = (file: string): void => {
144+
const db = new DatabaseSync(file);
145+
try {
146+
db.exec("UPDATE agent_state_journal SET state = NULL WHERE kind = 'event'");
147+
} finally {
148+
db.close();
149+
}
150+
};
151+
152+
interface ValueState {
153+
readonly value: number;
154+
}
155+
156+
const valueCounterDefinition = (
157+
id = 'state-sqlite-test/value-counter',
158+
): AgentStateDefinition<ValueState, typeof counterEvents> =>
159+
defineState({
160+
events: counterEvents,
161+
id,
162+
initial: { value: 0 },
163+
lifetime: 'workspace-durable',
164+
migrations: {
165+
2: (persisted) => ({ value: (persisted as CounterState).count * 10 }),
166+
},
167+
reduce: (state, event) => ({ value: state.value + event.payload.by }),
168+
schema: z.object({ value: z.number().int() }).strict(),
169+
version: 2,
170+
});
171+
172+
const createLegacyHeadOnlyDatabase = (file: string, definitionId: string): void => {
173+
createLegacyMigrationDatabase(file, definitionId);
174+
const db = new DatabaseSync(file);
175+
try {
176+
db.exec(`
177+
DELETE FROM agent_state_journal WHERE revision > 1;
178+
UPDATE agent_state_journal SET state = NULL WHERE revision = 1;
179+
UPDATE agent_state_head SET revision = 1, state = '{"count":2}' WHERE id = 1;
180+
`);
181+
} finally {
182+
db.close();
183+
}
184+
};
185+
143186
const holdUncheckpointedLegacyEvent = (file: string): DatabaseSync => {
144187
const keeper = new DatabaseSync(file);
145188
keeper.exec('PRAGMA journal_mode = WAL; PRAGMA wal_autocheckpoint = 0; BEGIN DEFERRED');
@@ -240,19 +283,31 @@ describe('sqlite driver storage behavior', () => {
240283
}
241284
}));
242285

243-
it('backfills legacy NULL event results before migration for idempotent replay', () =>
286+
it('fails closed when a non-head legacy event has no recoverable committed result', () =>
244287
withRoot(async (root) => {
245288
const definition = migratingCounterDefinition();
246289
const file = join(root, 'legacy-null-event.sqlite');
247290
createLegacyMigrationDatabase(file, definition.id);
291+
clearLegacyEventResults(file);
292+
293+
await expect(createSqliteStateDriver({ file }).open(definition)).rejects.toMatchObject({
294+
code: 'migration-failure',
295+
message: expect.stringContaining('has no recoverable committed result'),
296+
name: 'AgentStateError',
297+
});
298+
}));
299+
300+
it('migrates a legacy journal-head result from the materialized head without using the current reducer', () =>
301+
withRoot(async (root) => {
302+
const definition = valueCounterDefinition();
303+
const file = join(root, 'legacy-head-event.sqlite');
304+
createLegacyHeadOnlyDatabase(file, definition.id);
248305

249306
const store = await createSqliteStateDriver({ file }).open(definition);
250307
await expect(
251308
store.dispatch('bumped', { by: 2 }, { idempotencyKey: 'legacy:event' }),
252-
).resolves.toEqual({ replayed: true, revision: 1, state: { count: 20 } });
253-
await expect(
254-
store.dispatch('bumped', { by: 1 }, { idempotencyKey: 'legacy:event-2' }),
255-
).resolves.toEqual({ replayed: true, revision: 3, state: { count: 60 } });
309+
).resolves.toEqual({ replayed: true, revision: 1, state: { value: 20 } });
310+
await expect(store.read()).resolves.toEqual({ revision: 2, state: { value: 20 } });
256311
await store.close();
257312
}));
258313

@@ -372,6 +427,59 @@ describe('sqlite driver storage behavior', () => {
372427
}
373428
}));
374429

430+
it('preserves the initialization error when database close also fails', () =>
431+
withRoot(async (root) => {
432+
const file = join(root, 'state.sqlite');
433+
const seed = await createSqliteStateDriver({ file }).open(counterDefinition());
434+
await seed.close();
435+
const db = new DatabaseSync(file);
436+
db.exec('UPDATE agent_state_meta SET kernel_format = 99');
437+
db.close();
438+
439+
const closeFailure = new Error('database close failed');
440+
const originalClose = DatabaseSync.prototype.close;
441+
DatabaseSync.prototype.close = function close(this: DatabaseSync): void {
442+
originalClose.call(this);
443+
throw closeFailure;
444+
};
445+
try {
446+
await expect(createSqliteStateDriver({ file }).open(counterDefinition())).rejects.toMatchObject({
447+
code: 'corrupt',
448+
message: expect.stringContaining('kernel format 99'),
449+
name: 'AgentStateError',
450+
});
451+
} finally {
452+
DatabaseSync.prototype.close = originalClose;
453+
}
454+
}));
455+
456+
it('attempts every store close before propagating the first close failure', () =>
457+
withRoot(async (root) => {
458+
const driver = createSqliteStateDriver({ root });
459+
const first = await driver.open(counterDefinition());
460+
const second = await driver.open(otherDefinition());
461+
const closeFailure = new Error('first database close failed');
462+
const originalClose = DatabaseSync.prototype.close;
463+
let closeAttempts = 0;
464+
DatabaseSync.prototype.close = function close(this: DatabaseSync): void {
465+
originalClose.call(this);
466+
closeAttempts += 1;
467+
if (closeAttempts === 1) throw closeFailure;
468+
};
469+
try {
470+
await expect(driver.close()).rejects.toBe(closeFailure);
471+
expect(closeAttempts).toBe(2);
472+
await expect(driver.close()).rejects.toBe(closeFailure);
473+
expect(closeAttempts).toBe(2);
474+
await expect(first.read()).rejects.toMatchObject({ code: 'store-closed' });
475+
await expect(second.read()).rejects.toMatchObject({ code: 'store-closed' });
476+
} finally {
477+
DatabaseSync.prototype.close = originalClose;
478+
await first.close().catch(() => undefined);
479+
await second.close().catch(() => undefined);
480+
}
481+
}));
482+
375483
it('fails closed with a typed corrupt error when the file is not a database', () =>
376484
withRoot(async (root) => {
377485
const file = join(root, 'state.sqlite');

0 commit comments

Comments
 (0)