diff --git a/packages/session/session-persistence-sqlite/src/store.ts b/packages/session/session-persistence-sqlite/src/store.ts index e1e3f8b651..02223ee780 100644 --- a/packages/session/session-persistence-sqlite/src/store.ts +++ b/packages/session/session-persistence-sqlite/src/store.ts @@ -15,6 +15,7 @@ import { type SessionHeader, } from '@deepseek-ai/dsh-session' import { + createStoredEventRead, decodeStoredSessionHeader, SessionPersistenceRevision, SessionPersistenceRevisionConflictError, @@ -22,7 +23,6 @@ import { type SessionPersistenceRevision as PersistenceRevision, type SessionPersistenceSnapshot, type StoredEventRead, - type StoredEventReadCompletion, type StoredEventReadOptions, type StoredSessionSource, } from '@deepseek-ai/dsh-session-persistence' @@ -71,36 +71,6 @@ interface SqliteStoredSuffix { readonly revision: PersistenceRevision } -/** - * Build the standard lazy event stream and EOF metadata around one backend - * read, mirroring the service's protected helper for standalone stores. - * @param load - revision-checked batch loader owned by the backend. - * @param include - whether one loaded event belongs in this physical read. - * @param signal - optional cancellation checked between yielded events. - * @returns an independently consumable event read. - */ -function createStoredEventRead( - load: () => Promise<{ readonly events: readonly unknown[]; readonly tornMarker?: TornMarker }>, - include: (event: unknown) => boolean, - signal?: AbortSignal, -): StoredEventRead { - const completed = Promise.withResolvers>() - const events = (async function* (): AsyncIterable { - try { - const batch = await load() - for (const event of batch.events) { - signal?.throwIfAborted() - if (include(event)) yield event - } - completed.resolve(batch.tornMarker === undefined ? {} : { tornMarker: batch.tornMarker }) - } catch (error: unknown) { - completed.reject(error) - throw error - } - })() - return { events, completed: completed.promise } -} - /** SQLite implementation of the coordinator's physical backend hooks. */ export class SqliteStore implements PersistenceBackend { readonly name = 'session-persistence-sqlite' @@ -180,6 +150,12 @@ export class SqliteStore implements PersistenceBackend { } } + /** + * Load one row's complete validated prefix at a single snapshot. + * @param id - persisted session id to resolve. + * @param signal - optional cancellation for backend read work. + * @returns the stored prefix, or `undefined` when the session has no stored row. + */ async loadStored(id: SessionId, signal?: AbortSignal): Promise { await this.observe(signal) const snapshot = this.readTransaction(() => { @@ -206,6 +182,13 @@ export class SqliteStore implements PersistenceBackend { return row === undefined ? undefined : sqliteRevision(this.storeIdentity, row) } + /** + * Load one row's physical suffix at or past a sequence at a single snapshot. + * @param id - persisted session id to resolve. + * @param fromSeq - first physical event sequence to include. + * @param signal - optional cancellation for backend read work. + * @returns the stored suffix, or `undefined` when the session has no stored row. + */ async loadStoredFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise { await this.observe(signal) const snapshot = this.readTransaction(() => { @@ -373,18 +356,7 @@ export class SqliteStore implements PersistenceBackend { this.db.prepare(sql('delete-events-from')).run(meta.id, 0) const insert = this.insertStatement() for (const record of records) this.insertRecord(insert, meta.id, bindRecord(record)) - this.db.prepare(sql('upsert-session')).run( - meta.id, - meta.version, - meta.createdAt, - meta.cwd ?? null, - meta.parentSession ?? null, - meta.seedLength ?? null, - meta.origin ?? null, - meta.delegationDepth ?? null, - meta.agentPreset ?? null, - randomUUID(), - ) + this.writeRow(meta) this.incrementRevision(meta.id) this.db.exec(sql('commit')) } catch (error: unknown) { diff --git a/packages/session/session-persistence-sqlite/tests/resources/sql/count-session-events.sql b/packages/session/session-persistence-sqlite/tests/resources/sql/count-session-events.sql new file mode 100644 index 0000000000..da02575e16 --- /dev/null +++ b/packages/session/session-persistence-sqlite/tests/resources/sql/count-session-events.sql @@ -0,0 +1,3 @@ +SELECT COUNT(*) AS n +FROM events +WHERE session_id = ?; diff --git a/packages/session/session-persistence-sqlite/tests/resources/sql/create-temp-replace-trigger.sql b/packages/session/session-persistence-sqlite/tests/resources/sql/create-temp-replace-trigger.sql new file mode 100644 index 0000000000..1fbf6ae0c7 --- /dev/null +++ b/packages/session/session-persistence-sqlite/tests/resources/sql/create-temp-replace-trigger.sql @@ -0,0 +1,5 @@ +CREATE TEMP TRIGGER fail_format_replace +BEFORE UPDATE ON sessions +BEGIN + SELECT RAISE(ABORT, 'simulated format replacement failure'); +END diff --git a/packages/session/session-persistence-sqlite/tests/resources/sql/delete-session-by-id.sql b/packages/session/session-persistence-sqlite/tests/resources/sql/delete-session-by-id.sql new file mode 100644 index 0000000000..afe9d6c0ec --- /dev/null +++ b/packages/session/session-persistence-sqlite/tests/resources/sql/delete-session-by-id.sql @@ -0,0 +1,2 @@ +DELETE FROM sessions +WHERE id = ?; diff --git a/packages/session/session-persistence-sqlite/tests/resources/sql/drop-temp-replace-trigger.sql b/packages/session/session-persistence-sqlite/tests/resources/sql/drop-temp-replace-trigger.sql new file mode 100644 index 0000000000..b41f647451 --- /dev/null +++ b/packages/session/session-persistence-sqlite/tests/resources/sql/drop-temp-replace-trigger.sql @@ -0,0 +1 @@ +DROP TRIGGER fail_format_replace; diff --git a/packages/session/session-persistence-sqlite/tests/resources/sql/update-session-cwd.sql b/packages/session/session-persistence-sqlite/tests/resources/sql/update-session-cwd.sql new file mode 100644 index 0000000000..7586325b16 --- /dev/null +++ b/packages/session/session-persistence-sqlite/tests/resources/sql/update-session-cwd.sql @@ -0,0 +1,3 @@ +UPDATE sessions +SET cwd = ? +WHERE id = ?; diff --git a/packages/session/session-persistence-sqlite/tests/resources/sql/update-session-revision.sql b/packages/session/session-persistence-sqlite/tests/resources/sql/update-session-revision.sql new file mode 100644 index 0000000000..2cfbcb2b82 --- /dev/null +++ b/packages/session/session-persistence-sqlite/tests/resources/sql/update-session-revision.sql @@ -0,0 +1,3 @@ +UPDATE sessions +SET revision = revision + 1 +WHERE id = ?; diff --git a/packages/session/session-persistence-sqlite/tests/sqlite.spec.ts b/packages/session/session-persistence-sqlite/tests/sqlite.spec.ts index 448bff85ee..55d6e26862 100644 --- a/packages/session/session-persistence-sqlite/tests/sqlite.spec.ts +++ b/packages/session/session-persistence-sqlite/tests/sqlite.spec.ts @@ -871,7 +871,7 @@ describe('SessionPersistenceSqlite stored-source and replacement primitives', () const removed = await store.openStored(m.id) if (removed === undefined) throw new Error('test session must remain materialized') const db = (store as unknown as { db: DatabaseSync }).db - db.prepare('DELETE FROM sessions WHERE id = ?').run(m.id) + db.prepare(testSql('delete-session-by-id')).run(m.id) const removedRead = removed.readEvents({ fromSeq: 1 }) const removedCompletion = removedRead.completed.catch((error: unknown) => error) await expect((async () => { for await (const _event of removedRead.events) { /* consume */ } })()) @@ -896,7 +896,7 @@ describe('SessionPersistenceSqlite stored-source and replacement primitives', () }) await expect(store.loadStoredFrom(m.id, 1)).rejects.toThrow('simulated suffix SELECT failure') spy.mockRestore() - expect((db.prepare('SELECT COUNT(*) AS n FROM events WHERE session_id = ?').get(m.id) as { n: number }).n) + expect((db.prepare(testSql('count-session-events')).get(m.id) as { n: number }).n) .toBe(oneTurnLog().length) await store.close() }) @@ -944,11 +944,11 @@ describe('SessionPersistenceSqlite stored-source and replacement primitives', () const db = (store as unknown as { db: DatabaseSync }).db const changesDuringStaging = (async function* (): AsyncIterable { yield* oneTurnLog() - db.prepare('UPDATE sessions SET cwd = ? WHERE id = ?').run('/raced', m.id) + db.prepare(testSql('update-session-cwd')).run('/raced', m.id) })() await expect(store.replaceStored(source.revision, m, changesDuringStaging)) .rejects.toThrow(/changes its stored identity/) - expect((db.prepare('SELECT COUNT(*) AS n FROM events WHERE session_id = ?').get(m.id) as { n: number }).n) + expect((db.prepare(testSql('count-session-events')).get(m.id) as { n: number }).n) .toBe(oneTurnLog().length) await store.close() }) @@ -963,12 +963,12 @@ describe('SessionPersistenceSqlite stored-source and replacement primitives', () const db = (store as unknown as { db: DatabaseSync }).db const changesDuringStaging = (async function* (): AsyncIterable { yield* oneTurnLog() - db.prepare('UPDATE sessions SET revision = revision + 1 WHERE id = ?').run(m.id) + db.prepare(testSql('update-session-revision')).run(m.id) })() await expect(store.replaceStored(source.revision, m, changesDuringStaging)) .rejects.toBeInstanceOf(SessionPersistenceRevisionConflictError) - expect((db.prepare('SELECT COUNT(*) AS n FROM events WHERE session_id = ?').get(m.id) as { n: number }).n) + expect((db.prepare(testSql('count-session-events')).get(m.id) as { n: number }).n) .toBe(oneTurnLog().length) await store.close() }) @@ -981,18 +981,12 @@ describe('SessionPersistenceSqlite stored-source and replacement primitives', () const source = await store.openStored(m.id) if (source === undefined) throw new Error('test session must be materialized') const db = (store as unknown as { db: DatabaseSync }).db - db.exec(` - CREATE TEMP TRIGGER fail_format_replace - BEFORE UPDATE ON sessions - BEGIN - SELECT RAISE(ABORT, 'simulated format replacement failure'); - END - `) + db.exec(testSql('create-temp-replace-trigger')) await expect( store.replaceStored(source.revision, m, replacementEvents([])), ).rejects.toThrow(/simulated format replacement failure/) - db.exec('DROP TRIGGER fail_format_replace') + db.exec(testSql('drop-temp-replace-trigger')) expect((await store.loadStored(m.id))?.events).toEqual(oneTurnLog()) await store.close() diff --git a/packages/session/session-persistence-sqlite/tests/test-sql.ts b/packages/session/session-persistence-sqlite/tests/test-sql.ts index 77b53a404e..beb11d6f65 100644 --- a/packages/session/session-persistence-sqlite/tests/test-sql.ts +++ b/packages/session/session-persistence-sqlite/tests/test-sql.ts @@ -8,10 +8,14 @@ export type TestSqlName = | 'count-ignorable-events' | 'count-packed-events' | 'count-physical-types' + | 'count-session-events' | 'create-loose-schema' + | 'create-temp-replace-trigger' | 'create-unrelated-table' | 'delete-persistence-state' + | 'delete-session-by-id' | 'delete-session-events' + | 'drop-temp-replace-trigger' | 'empty-store-id' | 'insert-corrupt-event' | 'measure-write-traffic' @@ -25,6 +29,8 @@ export type TestSqlName = | 'set-user-version-16' | 'set-user-version-17' | 'update-invalid-session-metadata' + | 'update-session-cwd' + | 'update-session-revision' /** Load one fixed test SQL resource. */ export function testSql(name: TestSqlName): string { diff --git a/packages/session/session-persistence/src/format-decoder.ts b/packages/session/session-persistence/src/format-decoder.ts index 3a136441d0..b5e38c4e98 100644 --- a/packages/session/session-persistence/src/format-decoder.ts +++ b/packages/session/session-persistence/src/format-decoder.ts @@ -89,6 +89,36 @@ export interface StoredSessionSource { readEvents(options?: StoredEventReadOptions): StoredEventRead } +/** + * Build the standard lazy event stream and EOF metadata around one backend + * read, shared by every first-party backend. + * @param load - revision-checked batch loader owned by the backend. + * @param include - whether one loaded event belongs in this physical read. + * @param signal - optional cancellation checked between yielded events. + * @returns an independently consumable event read. + */ +export function createStoredEventRead( + load: () => Promise<{ readonly events: readonly unknown[]; readonly tornMarker?: TornMarker }>, + include: (event: unknown) => boolean, + signal?: AbortSignal, +): StoredEventRead { + const completed = Promise.withResolvers>() + const events = (async function* (): AsyncIterable { + try { + const batch = await load() + for (const event of batch.events) { + signal?.throwIfAborted() + if (include(event)) yield event + } + completed.resolve(batch.tornMarker === undefined ? {} : { tornMarker: batch.tornMarker }) + } catch (error: unknown) { + completed.reject(error) + throw error + } + })() + return { events, completed: completed.promise } +} + /** One decoded current-format read bound to an exact stored revision. */ export interface DecodedSession { /** Validated current-format header. */ diff --git a/packages/session/session-persistence/src/index.ts b/packages/session/session-persistence/src/index.ts index 297a768b75..1fd2c3a73d 100644 --- a/packages/session/session-persistence/src/index.ts +++ b/packages/session/session-persistence/src/index.ts @@ -8,7 +8,7 @@ import { Context, Service } from '@deepseek-ai/cordis' import { SessionPreparation, type SessionEvent, type SessionId, type SessionHeader } from '@deepseek-ai/dsh-session' import type { SessionPersistenceRevision } from './revision.ts' -import type { StoredEventRead, StoredEventReadCompletion } from './format-decoder.ts' +import { createStoredEventRead, type StoredEventRead } from './format-decoder.ts' // Re-export the metadata vocabulary so Consumers import it from the Service Definition. export type { SessionHeader } from '@deepseek-ai/dsh-session' @@ -53,6 +53,7 @@ export type { PersistenceCoordinatorOptions, } from './coordinator.ts' export { + createStoredEventRead, decodeStoredSessionHeader, SessionFormatUnsupportedError, sessionFormatVersionRefusal, @@ -98,21 +99,7 @@ export abstract class SessionPersistence extends Service { include: (event: unknown) => boolean, signal?: AbortSignal, ): StoredEventRead { - const completed = Promise.withResolvers>() - const events = (async function* (): AsyncIterable { - try { - const batch = await load() - for (const event of batch.events) { - signal?.throwIfAborted() - if (include(event)) yield event - } - completed.resolve(batch.tornMarker === undefined ? {} : { tornMarker: batch.tornMarker }) - } catch (error: unknown) { - completed.reject(error) - throw error - } - })() - return { events, completed: completed.promise } + return createStoredEventRead(load, include, signal) } /**