mirror of
https://github.com/deepseek-ai/deepseek-harness.git
synced 2026-08-29 04:26:38 +00:00
fix(session-persistence): repair persistence CI gates
- document the exported SqliteStore prefix/suffix loaders (verify-export-jsdoc) - share createStoredEventRead from session-persistence so the standalone SQLite store stops duplicating the service helper (duplication gate) - route replaceStored header upserts through writeRow (duplication gate) - move replacement/conflict test SQL into closed test resources so the SQLite SQL resource boundary test passes
This commit is contained in:
@@ -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<TornMarker>(
|
||||
load: () => Promise<{ readonly events: readonly unknown[]; readonly tornMarker?: TornMarker }>,
|
||||
include: (event: unknown) => boolean,
|
||||
signal?: AbortSignal,
|
||||
): StoredEventRead<TornMarker> {
|
||||
const completed = Promise.withResolvers<StoredEventReadCompletion<TornMarker>>()
|
||||
const events = (async function* (): AsyncIterable<unknown> {
|
||||
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<number> {
|
||||
readonly name = 'session-persistence-sqlite'
|
||||
@@ -180,6 +150,12 @@ export class SqliteStore implements PersistenceBackend<number> {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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<SqliteStoredPrefix | undefined> {
|
||||
await this.observe(signal)
|
||||
const snapshot = this.readTransaction(() => {
|
||||
@@ -206,6 +182,13 @@ export class SqliteStore implements PersistenceBackend<number> {
|
||||
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<SqliteStoredSuffix | undefined> {
|
||||
await this.observe(signal)
|
||||
const snapshot = this.readTransaction(() => {
|
||||
@@ -373,18 +356,7 @@ export class SqliteStore implements PersistenceBackend<number> {
|
||||
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) {
|
||||
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
SELECT COUNT(*) AS n
|
||||
FROM events
|
||||
WHERE session_id = ?;
|
||||
+5
@@ -0,0 +1,5 @@
|
||||
CREATE TEMP TRIGGER fail_format_replace
|
||||
BEFORE UPDATE ON sessions
|
||||
BEGIN
|
||||
SELECT RAISE(ABORT, 'simulated format replacement failure');
|
||||
END
|
||||
+2
@@ -0,0 +1,2 @@
|
||||
DELETE FROM sessions
|
||||
WHERE id = ?;
|
||||
+1
@@ -0,0 +1 @@
|
||||
DROP TRIGGER fail_format_replace;
|
||||
@@ -0,0 +1,3 @@
|
||||
UPDATE sessions
|
||||
SET cwd = ?
|
||||
WHERE id = ?;
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
UPDATE sessions
|
||||
SET revision = revision + 1
|
||||
WHERE id = ?;
|
||||
@@ -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<SessionEvent> {
|
||||
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<SessionEvent> {
|
||||
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()
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -89,6 +89,36 @@ export interface StoredSessionSource<TornMarker> {
|
||||
readEvents(options?: StoredEventReadOptions): StoredEventRead<TornMarker>
|
||||
}
|
||||
|
||||
/**
|
||||
* 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<TornMarker>(
|
||||
load: () => Promise<{ readonly events: readonly unknown[]; readonly tornMarker?: TornMarker }>,
|
||||
include: (event: unknown) => boolean,
|
||||
signal?: AbortSignal,
|
||||
): StoredEventRead<TornMarker> {
|
||||
const completed = Promise.withResolvers<StoredEventReadCompletion<TornMarker>>()
|
||||
const events = (async function* (): AsyncIterable<unknown> {
|
||||
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<TornMarker> {
|
||||
/** Validated current-format header. */
|
||||
|
||||
@@ -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<TornMarker> {
|
||||
const completed = Promise.withResolvers<StoredEventReadCompletion<TornMarker>>()
|
||||
const events = (async function* (): AsyncIterable<unknown> {
|
||||
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)
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user