feat(session): refuse session logs a build cannot faithfully read
Old runtimes meeting a newer session format now fail loud instead of misreading: version refusal names the direction (newer: upgrade the harness; older: no upgrade path) and points at the raw JSONL log, and an event type outside the generated known vocabulary refuses resume unless its envelope carries the new ignorable: true marker (default: required, so a forgotten marker over-refuses instead of silently resuming a gutted session). gen-persistence-catalog now also emits KNOWN_SESSION_EVENT_TYPES; SQLite stores the marker in a dedicated column (SCHEMA_VERSION 15). The versioning design (monotonic integer, n->n+1 upgrader chain, migrate-on-continue) is recorded in the session-log-version-mechanism Agent Note.
This commit is contained in:
@@ -186,6 +186,23 @@ describe('SessionPersistenceJsonl: format helpers', () => {
|
||||
})
|
||||
await fiber.dispose()
|
||||
})
|
||||
|
||||
it('points a format refusal at the raw log path', async () => {
|
||||
const absoluteRoot = await freshRoot()
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const fiber = await ctx.plugin(SessionPersistenceJsonl, { root: absoluteRoot, compression: 'none' })
|
||||
const m = { ...meta('newer-format', '/work'), version: 7 }
|
||||
await ctx.sessionPersistence.create(m)
|
||||
await ctx.sessionPersistence.append(m.id, [
|
||||
{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
|
||||
{ type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } },
|
||||
])
|
||||
const failure = await ctx.sessionPersistence.load(m.id).then(() => undefined, (error: unknown) => error as Error)
|
||||
expect(failure?.name).toBe('SessionFormatUnsupportedError')
|
||||
expect(failure?.message).toContain(`(raw log: ${rawLogPath(resolve(absoluteRoot), '/work', m.id)})`)
|
||||
await fiber.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
describe('SessionPersistenceJsonl: durability and crash semantics', () => {
|
||||
|
||||
@@ -28,15 +28,18 @@ import {
|
||||
export { SCHEMA_VERSION } from './schema.ts'
|
||||
|
||||
/**
|
||||
* Serialize an event's surface-metadata fields for SQL binding. Both fields are
|
||||
* nullable TEXT columns — null when the event has no surface metadata (non-surface
|
||||
* events, events written before surface support).
|
||||
* Serialize an event's optional envelope fields for SQL binding. The surface
|
||||
* fields are nullable TEXT columns — null when the event has no surface
|
||||
* metadata (non-surface events, events written before surface support); the
|
||||
* ignorable marker is a nullable INTEGER column — `1` iff the envelope carries
|
||||
* `ignorable: true`.
|
||||
*/
|
||||
function surfaceBindings(event: SessionEvent): [string | null, string | null] {
|
||||
function envelopeBindings(event: SessionEvent): [string | null, string | null, number | null] {
|
||||
const se = event as SessionEvent<SurfaceEventType>
|
||||
return [
|
||||
se.sourceEventSeqs ? JSON.stringify(se.sourceEventSeqs) : null,
|
||||
se.surfaceOp !== undefined ? JSON.stringify(se.surfaceOp) : null,
|
||||
event.ignorable === true ? 1 : null,
|
||||
]
|
||||
}
|
||||
|
||||
@@ -225,7 +228,7 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
if (row === undefined) return undefined
|
||||
const meta = rowToMeta(row)
|
||||
const eventRows = this.db
|
||||
.prepare('SELECT seq, type, time, data, source_event_seqs, surface_op FROM events WHERE session_id = ? AND seq >= ? ORDER BY seq')
|
||||
.prepare('SELECT seq, type, time, data, source_event_seqs, surface_op, ignorable FROM events WHERE session_id = ? AND seq >= ? ORDER BY seq')
|
||||
.all(id, fromSeq) as unknown as EventRow[]
|
||||
signal?.throwIfAborted()
|
||||
const { preserved } = scanRows(eventRows, fromSeq)
|
||||
@@ -247,7 +250,7 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
const row = this.rowFor(id)
|
||||
if (row !== undefined) {
|
||||
const eventRows = this.db
|
||||
.prepare('SELECT seq, type, time, data, source_event_seqs, surface_op FROM events WHERE session_id = ? ORDER BY seq')
|
||||
.prepare('SELECT seq, type, time, data, source_event_seqs, surface_op, ignorable FROM events WHERE session_id = ? ORDER BY seq')
|
||||
.all(id) as unknown as EventRow[]
|
||||
snapshot = { row, eventRows }
|
||||
}
|
||||
@@ -279,14 +282,14 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
async appendBatch(meta: SessionHeader, events: readonly SessionEvent[], isMaterialized: boolean): Promise<void> {
|
||||
await this.ready
|
||||
const insertEvent = this.db.prepare(
|
||||
'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)',
|
||||
'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op, ignorable) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
|
||||
)
|
||||
this.db.exec('BEGIN')
|
||||
try {
|
||||
if (!isMaterialized) this.writeRow(meta)
|
||||
for (const event of events) {
|
||||
const [surfaceSeqs, surfaceOp] = surfaceBindings(event)
|
||||
insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp)
|
||||
const [surfaceSeqs, surfaceOp, ignorable] = envelopeBindings(event)
|
||||
insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp, ignorable)
|
||||
}
|
||||
this.db.prepare('UPDATE sessions SET revision = revision + 1 WHERE id = ?').run(meta.id)
|
||||
this.db.exec('COMMIT')
|
||||
@@ -310,11 +313,11 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
}
|
||||
if (closers.length > 0) {
|
||||
const insertEvent = this.db.prepare(
|
||||
'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)',
|
||||
'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op, ignorable) VALUES (?, ?, ?, ?, ?, ?, ?, ?)',
|
||||
)
|
||||
for (const event of closers) {
|
||||
const [surfaceSeqs, surfaceOp] = surfaceBindings(event)
|
||||
insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp)
|
||||
const [surfaceSeqs, surfaceOp, ignorable] = envelopeBindings(event)
|
||||
insertEvent.run(meta.id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp, ignorable)
|
||||
}
|
||||
}
|
||||
if (tornMarker !== undefined || closers.length > 0) {
|
||||
|
||||
@@ -17,7 +17,7 @@ import type { SessionEvent, SessionId, SessionHeader, SurfaceOp } from '@deepsee
|
||||
* layout; orthogonal to a session's own `version` (which versions the EVENT
|
||||
* vocabulary, stored per session in the `sessions` row).
|
||||
*/
|
||||
export const SCHEMA_VERSION = 14
|
||||
export const SCHEMA_VERSION = 15
|
||||
|
||||
/** SQLite application id protecting unrelated databases from persistence writes. */
|
||||
export const SESSION_PERSISTENCE_SQLITE_APPLICATION_ID = 0x44534850
|
||||
@@ -55,6 +55,8 @@ export interface EventRow {
|
||||
source_event_seqs: string | null
|
||||
/** JSON-encoded `SurfaceOp` — how the event entered the surface, or null. */
|
||||
surface_op: string | null
|
||||
/** `1` iff the event carries the envelope's `ignorable: true` marker, else null. */
|
||||
ignorable: number | null
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -139,6 +141,7 @@ function configureDatabase(db: DatabaseSync, path: string, journalMode: JournalM
|
||||
data TEXT NOT NULL,
|
||||
source_event_seqs TEXT,
|
||||
surface_op TEXT,
|
||||
ignorable INTEGER,
|
||||
PRIMARY KEY (session_id, seq)
|
||||
) STRICT
|
||||
`)
|
||||
@@ -203,12 +206,14 @@ export function rowToEvent(row: EventRow): SessionEvent {
|
||||
...row.source_event_seqs !== null ? { sourceEventSeqs: JSON.parse(row.source_event_seqs) as number[] } : {},
|
||||
...row.surface_op !== null ? { surfaceOp: JSON.parse(row.surface_op) as SurfaceOp } : {},
|
||||
}
|
||||
const ignorableField = row.ignorable === 1 ? { ignorable: true as const } : {}
|
||||
return {
|
||||
type: row.type as SessionEvent['type'],
|
||||
seq: row.seq,
|
||||
time: row.time,
|
||||
data: JSON.parse(row.data) as SessionEvent['data'],
|
||||
...surfaceFields,
|
||||
...ignorableField,
|
||||
} as SessionEvent
|
||||
}
|
||||
|
||||
|
||||
@@ -92,6 +92,7 @@ describe('scanRows', () => {
|
||||
seq: e.seq, type: e.type, time: e.time, data: JSON.stringify(e.data),
|
||||
source_event_seqs: se.sourceEventSeqs !== undefined ? JSON.stringify(se.sourceEventSeqs) : null,
|
||||
surface_op: se.surfaceOp !== undefined ? JSON.stringify(se.surfaceOp) : null,
|
||||
ignorable: e.ignorable === true ? 1 : null,
|
||||
}
|
||||
})
|
||||
|
||||
@@ -142,8 +143,8 @@ describe('scanRows', () => {
|
||||
|
||||
it('throws on an unparsable row inside the committed region', () => {
|
||||
const withCorruptCommitted: EventRow[] = [
|
||||
{ seq: 0, type: 'turn/start', time: 1, data: '{not json', source_event_seqs: null, surface_op: null }, // corrupt, sits before a turn/end
|
||||
{ seq: 1, type: 'turn/end', time: 2, data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }), source_event_seqs: null, surface_op: null },
|
||||
{ seq: 0, type: 'turn/start', time: 1, data: '{not json', source_event_seqs: null, surface_op: null, ignorable: null }, // corrupt, sits before a turn/end
|
||||
{ seq: 1, type: 'turn/end', time: 2, data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }), source_event_seqs: null, surface_op: null, ignorable: null },
|
||||
]
|
||||
expect(() => scanRows(withCorruptCommitted)).toThrow(/unparsable committed event/)
|
||||
})
|
||||
@@ -151,7 +152,7 @@ describe('scanRows', () => {
|
||||
it('tolerates an unparsable torn-tail row after the last turn/end', () => {
|
||||
const withCorruptTail: EventRow[] = [
|
||||
...rows(oneTurnLog()),
|
||||
{ seq: 6, type: 'turn/start', time: 7, data: '{not json', source_event_seqs: null, surface_op: null }, // torn fragment, no committed turn/end after
|
||||
{ seq: 6, type: 'turn/start', time: 7, data: '{not json', source_event_seqs: null, surface_op: null, ignorable: null }, // torn fragment, no committed turn/end after
|
||||
]
|
||||
const { preserved, tornFrom } = scanRows(withCorruptTail)
|
||||
expect(preserved).toEqual(oneTurnLog())
|
||||
@@ -658,7 +659,7 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
|
||||
})
|
||||
|
||||
it('exposes the schema version constant', () => {
|
||||
expect(SCHEMA_VERSION).toBe(14)
|
||||
expect(SCHEMA_VERSION).toBe(15)
|
||||
})
|
||||
|
||||
it('keeps the revision stable for an empty repair hook', async () => {
|
||||
@@ -857,6 +858,7 @@ describe('surface field round-trip', () => {
|
||||
data: JSON.stringify({ turn: 1, step: 1, content: [] }),
|
||||
source_event_seqs: JSON.stringify([3, 5]),
|
||||
surface_op: JSON.stringify('append'),
|
||||
ignorable: null,
|
||||
}
|
||||
const event = rowToEvent(row)
|
||||
expect((event as SurfaceEvent).sourceEventSeqs).toEqual([3, 5])
|
||||
@@ -869,6 +871,7 @@ describe('surface field round-trip', () => {
|
||||
data: JSON.stringify({ turn: 1, step: 1, content: [] }),
|
||||
source_event_seqs: JSON.stringify([0, 1]),
|
||||
surface_op: JSON.stringify({ op: 'replace', start: 0, end: 1 }),
|
||||
ignorable: null,
|
||||
}
|
||||
const event = rowToEvent(row)
|
||||
expect((event as SurfaceEvent).sourceEventSeqs).toEqual([0, 1])
|
||||
@@ -879,10 +882,10 @@ describe('surface field round-trip', () => {
|
||||
const rows: EventRow[] = [
|
||||
{ seq: 0, type: 'user/message', time: 1,
|
||||
data: JSON.stringify({ content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }),
|
||||
source_event_seqs: null, surface_op: '{"op":"replace","start":0,"end":0}' },
|
||||
source_event_seqs: null, surface_op: '{"op":"replace","start":0,"end":0}', ignorable: null },
|
||||
{ seq: 1, type: 'turn/end', time: 2,
|
||||
data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }),
|
||||
source_event_seqs: null, surface_op: null },
|
||||
source_event_seqs: null, surface_op: null, ignorable: 1 },
|
||||
]
|
||||
const { preserved } = scanRows(rows)
|
||||
expect(preserved).toHaveLength(2)
|
||||
|
||||
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write packages/session/session-persistence/README.md
|
||||
README.md: 391548b1b896dca14cbe4f4ae55cf4180c4e0ac2
|
||||
README.zh.md: 7213e1ee71ba418ffacc3685df371dcba33588a7
|
||||
README.md: 7e62360ccf47151f5c450685bfebe6e89bbf187b
|
||||
README.zh.md: 3d819ef0ab4f85c83c2e640f627e36341318ac35
|
||||
|
||||
@@ -14,7 +14,7 @@ The persisted unit IS the existing `SessionEvent` (event-sourced model — the l
|
||||
| `create(meta): Promise<void>` | Register a new session's metadata. MAY defer the physical write until the first `append` (lazy materialization). |
|
||||
| `append(id, events): Promise<void>` | Durably persist a batch. Append-only; first event `seq` == stored next-seq after any repair; rejects non-JSON-serializable data naming the offending type. |
|
||||
| `prepare(id, signal?): Promise<SessionPreparation>` | Reserve the exact unpublished Session used by resume. A coordinator reuses an earlier inspection when available, commits pending recovery, and releases an unpublished reservation back to its bounded cache on disposal. |
|
||||
| `load(id): Promise<{ meta; events }>` | Return an immutable balanced logical log after converting supported older records from the same format version and committing cold recovery. A live load first flushes its snapshot and rejects while its turn is open; a cold load preserves an interrupted final turn and durably closes it with synthetic `tool/result`/`step/end?`/`turn/end {interrupted}` events. Only a torn tail fragment is dropped; committed corruption, malformed records, and unknown `version` reject. |
|
||||
| `load(id): Promise<{ meta; events }>` | Return an immutable balanced logical log after converting supported older records from the same format version and committing cold recovery. A live load first flushes its snapshot and rejects while its turn is open; a cold load preserves an interrupted final turn and durably closes it with synthetic `tool/result`/`step/end?`/`turn/end {interrupted}` events. Only a torn tail fragment is dropped; committed corruption and malformed records reject as `SessionPersistenceCorruptionError`, while an unsupported format `version` or an event type unknown to this build (without the envelope's `ignorable` marker) refuses as `SessionFormatUnsupportedError`, naming the refusal direction and the raw log path when the backend keeps one artifact per session. |
|
||||
| `inspect(id, signal?): Promise<{ meta; events }>` | Return an upgraded, validated, deeply frozen logical view without committing recovery or publishing a Session. A cold view receives in-memory synthetic recovery closers while its physical torn tail remains untouched; an already-live view is its current immutable snapshot and may contain an open turn. Coordinator-backed implementations retain the exact cold unpublished Session in a bounded LRU for later `prepare`, but discard and reload it when the stored revision changes. Same-id inspections share an in-flight read. |
|
||||
| `readFrom(id, fromSeq, signal?): Promise<{ meta; events }>` | Return valid stored events with `seq >= fromSeq` without preparation caching, truncation, closers, or coordinator state. A `fromSeq` at or past the stored end returns an empty event list; a negative or non-safe-integer `fromSeq` rejects. Seek-capable backends (SQLite) read only the suffix unless converting a supported older record requires earlier records; sequential backends (JSONL) parse the whole artifact and skip forward. Intended for checkpoint consumers that apply only events after a stored sequence number. |
|
||||
| `list(signal?): Promise<SessionHeader[]>` | Lightweight listing from metadata, no full-log parse. The optional signal cancels backend listing work. A zero-event lazily-materialized session is absent from `list`. |
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
| `create(meta): Promise<void>` | 注册新会话元数据。可以将物理写入延迟到第一次 `append`(延迟实体化)。 |
|
||||
| `append(id, events): Promise<void>` | 持久保存一个批次。仅追加;任何修复后,第一个事件 `seq` == 已存储 next-seq;非 JSON 可序列化数据会被拒绝,并命名违规类型。 |
|
||||
| `prepare(id, signal?): Promise<SessionPreparation>` | 预留恢复所使用的那个未发布 Session。协调器会尽可能复用之前的检查结果、提交待处理恢复,并在 dispose 时将未发布 reservation 释放回有界缓存。 |
|
||||
| `load(id): Promise<{ meta; events }>` | 转换同一格式版本中受支持的旧记录后,返回不可变、平衡的逻辑日志,并提交冷恢复。实时 load 先 flush 其快照,并在轮次开放时拒绝;冷 load 保留中断的最终轮次,并用合成 `tool/result`/`step/end?`/`turn/end {interrupted}` 事件持久关闭它。只丢弃撕裂尾部碎片;已提交损坏、格式错误的记录和未知 `version` 会被拒绝。 |
|
||||
| `load(id): Promise<{ meta; events }>` | 转换同一格式版本中受支持的旧记录后,返回不可变、平衡的逻辑日志,并提交冷恢复。实时 load 先 flush 其快照,并在轮次开放时拒绝;冷 load 保留中断的最终轮次,并用合成 `tool/result`/`step/end?`/`turn/end {interrupted}` 事件持久关闭它。只丢弃撕裂尾部碎片;已提交损坏和格式错误的记录以 `SessionPersistenceCorruptionError` 拒绝,不支持的格式 `version` 或本构建不认识且信封未带 `ignorable` 标记的事件类型以 `SessionFormatUnsupportedError` 拒绝,消息说明拒绝方向,并在后端为每个会话保留独立文件时给出原始日志路径。 |
|
||||
| `inspect(id, signal?): Promise<{ meta; events }>` | 返回已经升级、验证和深度冻结的逻辑视图,但不提交恢复或发布 Session。冷视图会获得仅存在于内存的合成恢复 closer,物理撕裂尾部保持不变;实时状态下的视图则是当前不可变快照,可能包含开放的轮次。基于协调器的实现会在有界 LRU 中保留该冷状态下未发布的 Session 本身,供后续 `prepare` 使用,但已存储修订值变化后会丢弃并重新读取。同 id 检查共享进行中的读取。 |
|
||||
| `readFrom(id, fromSeq, signal?): Promise<{ meta; events }>` | 返回 `seq >= fromSeq` 的有效已存储事件,不进入 preparation 缓存、不截断、不合成 closer,也不发布协调器状态。`fromSeq` 达到或超过已存储末尾时返回空事件列表;负数或非安全整数 `fromSeq` 会被拒绝。可寻址后端(SQLite)只读后缀,除非转换受支持的旧记录需要读取更早的记录;顺序后端(JSONL)解析整个产物并向前跳过。供 checkpoint 消费方只应用已存序号之后的事件。 |
|
||||
| `list(signal?): Promise<SessionHeader[]>` | 从元数据轻量列出,不解析完整日志。可选信号取消后端列表工作。零事件延迟实体化会话不在 `list` 中。 |
|
||||
|
||||
@@ -9,6 +9,7 @@ import { Context } from '@deepseek-ai/cordis'
|
||||
import {
|
||||
adoptSessionEvent,
|
||||
interruptedTurnClosers,
|
||||
KNOWN_SESSION_EVENT_TYPES,
|
||||
SESSION_FORMAT_VERSION,
|
||||
SessionPreparation,
|
||||
snapshotJsonValue,
|
||||
@@ -16,7 +17,7 @@ import {
|
||||
} from '@deepseek-ai/dsh-session'
|
||||
import type { Session, SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
|
||||
import type { SessionInspection } from './index.ts'
|
||||
import type { SessionInspection, SessionLocation } from './index.ts'
|
||||
import type { SessionPersistenceRevision } from './revision.ts'
|
||||
import { observeQueuedAbort, SessionPreparations } from './preparations.ts'
|
||||
import type { SessionPreparationReservation } from './preparations.ts'
|
||||
@@ -43,6 +44,26 @@ export class SessionPersistenceCorruptionError extends Error {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The stored log is intact but this runtime cannot faithfully interpret it:
|
||||
* the header carries an unsupported format version, or an event's type is
|
||||
* unknown to this build and the event is not marked ignorable. Distinct from
|
||||
* {@link SessionPersistenceCorruptionError} — nothing is damaged; the raw log
|
||||
* remains readable at {@link location} when the backend keeps one artifact
|
||||
* per session.
|
||||
*/
|
||||
export class SessionFormatUnsupportedError extends Error {
|
||||
/**
|
||||
* @param message - stable reason the log cannot be interpreted, already
|
||||
* including the raw-log path when one exists.
|
||||
* @param location - the backend's artifact location, when one exists.
|
||||
*/
|
||||
constructor(message: string, readonly location?: SessionLocation) {
|
||||
super(message)
|
||||
this.name = 'SessionFormatUnsupportedError'
|
||||
}
|
||||
}
|
||||
|
||||
/** Coordinator policy supplied by a concrete persistence backend. */
|
||||
export interface PersistenceCoordinatorOptions {
|
||||
/** Maximum completed unpublished preparations retained for reuse. */
|
||||
@@ -156,6 +177,14 @@ export interface PersistenceBackend<TornMarker = unknown> {
|
||||
*/
|
||||
list(signal?: AbortSignal): Promise<SessionHeader[]>
|
||||
|
||||
/**
|
||||
* Optional side-effect-free artifact locator, used to point refusal
|
||||
* diagnostics ({@link SessionFormatUnsupportedError}) at the raw log.
|
||||
* Backends without one artifact per session omit it or return `undefined`.
|
||||
* @param meta - the header whose artifact is requested.
|
||||
*/
|
||||
locate?(meta: SessionHeader): SessionLocation | undefined
|
||||
|
||||
/**
|
||||
* Optional lifecycle teardown (e.g. close a database handle). Awaited by the
|
||||
* coordinator's dispose effect AFTER the quiescence drain. A stateless file
|
||||
@@ -806,7 +835,9 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
const whole = await this.readStoredPrefix(id, signal)
|
||||
return { meta: whole.meta, events: whole.events.filter(event => event.seq >= fromSeq) }
|
||||
}
|
||||
return { meta: structuredClone(suffix.meta), events: snapshotStoredEvents(suffix.events, id) }
|
||||
const events = snapshotStoredEvents(suffix.events, id)
|
||||
this.assertEventsSupported(suffix.meta, events)
|
||||
return { meta: structuredClone(suffix.meta), events }
|
||||
}
|
||||
const whole = await this.readStoredPrefix(id, signal)
|
||||
// Sequential fallback: contiguous seqs from 0 make the suffix an index slice.
|
||||
@@ -824,9 +855,11 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
if (stored === undefined) throw new Error(`session "${id}" not found`)
|
||||
this.assertStoredId(id, stored.meta)
|
||||
this.assertVersion(stored.meta)
|
||||
const events = snapshotStoredEvents(stored.events, id)
|
||||
this.assertEventsSupported(stored.meta, events)
|
||||
return {
|
||||
meta: structuredClone(stored.meta),
|
||||
events: snapshotStoredEvents(stored.events, id),
|
||||
events,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -839,6 +872,7 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
this.assertStoredId(id, meta)
|
||||
this.assertVersion(meta)
|
||||
const storedEvents = adoptStoredEvents(events, id)
|
||||
this.assertEventsSupported(meta, storedEvents)
|
||||
|
||||
// Preserve complete interrupted events and synthesize only missing closers.
|
||||
const closers = interruptedTurnClosers(storedEvents).map(adoptSessionEvent)
|
||||
@@ -861,6 +895,9 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
closers,
|
||||
}
|
||||
} catch (error: unknown) {
|
||||
// An unsupported format is a refusal over an intact log, not damage —
|
||||
// surface it unwrapped so callers can point at the raw artifact.
|
||||
if (error instanceof SessionFormatUnsupportedError) throw error
|
||||
throw new SessionPersistenceCorruptionError(
|
||||
`stored session "${id}" failed validation: ${String(error)}`,
|
||||
{ cause: error },
|
||||
@@ -982,11 +1019,38 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
}
|
||||
|
||||
private assertVersion(meta: SessionHeader): void {
|
||||
if (meta.version !== SESSION_FORMAT_VERSION) {
|
||||
throw new Error(`unsupported session format version ${meta.version} for "${meta.id}" (only v${SESSION_FORMAT_VERSION} is supported)`)
|
||||
if (meta.version === SESSION_FORMAT_VERSION) return
|
||||
throw this.unsupported(meta, meta.version > SESSION_FORMAT_VERSION
|
||||
? `session "${meta.id}" uses log format v${meta.version}, but this harness reads only v${SESSION_FORMAT_VERSION}: the log was written by a newer harness — upgrade the harness to open it`
|
||||
: `session "${meta.id}" uses log format v${meta.version}, older than the supported v${SESSION_FORMAT_VERSION}, and this build ships no upgrade path for it`)
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuse a log containing an event type this build does not know, unless the
|
||||
* writer marked the event ignorable: an unrecognized required event may
|
||||
* change how the rest of the log must be interpreted, so silently skipping
|
||||
* it would reconstruct a wrong session (the envelope contract on
|
||||
* `SessionEvent.ignorable`). Runs on NORMALIZED events — after
|
||||
* `snapshotStoredEvents`/`adoptStoredEvents` has upgraded the legacy shapes
|
||||
* this build still reads and rejected the ones it does not, so those keep
|
||||
* their specific diagnostics.
|
||||
*/
|
||||
private assertEventsSupported(meta: SessionHeader, events: readonly SessionEvent[]): void {
|
||||
for (const event of events) {
|
||||
if (KNOWN_SESSION_EVENT_TYPES.has(event.type) || event.ignorable === true) continue
|
||||
throw this.unsupported(meta, `session "${meta.id}" contains event type "${event.type}" (seq ${event.seq}) unknown to this harness and not marked ignorable; refusing to interpret the log — it was likely written by a newer harness`)
|
||||
}
|
||||
}
|
||||
|
||||
/** Build a format refusal that points at the raw artifact when the backend has one. */
|
||||
private unsupported(meta: SessionHeader, reason: string): SessionFormatUnsupportedError {
|
||||
const location = this.backend.locate?.(meta)
|
||||
return new SessionFormatUnsupportedError(
|
||||
location === undefined ? reason : `${reason} (raw log: ${location.path})`,
|
||||
location,
|
||||
)
|
||||
}
|
||||
|
||||
/** Reject backend metadata that is not bound to the requested session id. */
|
||||
private assertStoredId(id: SessionId, meta: SessionHeader): void {
|
||||
if (meta.id !== id) {
|
||||
|
||||
@@ -36,6 +36,7 @@ export {
|
||||
DEFAULT_WRITE_BATCH_MAX_DELAY_MS,
|
||||
MAX_WRITE_BATCH_DELAY_MS,
|
||||
PersistenceCoordinator,
|
||||
SessionFormatUnsupportedError,
|
||||
SessionPersistenceCorruptionError,
|
||||
} from './coordinator.ts'
|
||||
export type {
|
||||
|
||||
@@ -706,6 +706,9 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
|
||||
.rejects.toThrow('lacks an identified message')
|
||||
}
|
||||
|
||||
// An out-of-repo event type passes only with the envelope's ignorable
|
||||
// marker (unknown-type refusal otherwise), and its non-object data is
|
||||
// not message-validated.
|
||||
const pluginId = SessionId('non-object-plugin-event')
|
||||
await ctx.sessionPersistence.create(meta(pluginId, WORK))
|
||||
await ctx.sessionPersistence.append(pluginId, [{
|
||||
@@ -713,11 +716,12 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
|
||||
seq: 0,
|
||||
time: 1,
|
||||
data: null,
|
||||
ignorable: true,
|
||||
} as unknown as SessionEvent])
|
||||
await expect(ctx.sessionPersistence.inspect(pluginId))
|
||||
.resolves.toMatchObject({ events: [{ type: 'plugin/test', data: null }] })
|
||||
.resolves.toMatchObject({ events: [{ type: 'plugin/test', data: null, ignorable: true }] })
|
||||
await expect(ctx.sessionPersistence.readFrom(pluginId, 0))
|
||||
.resolves.toMatchObject({ events: [{ type: 'plugin/test', data: null }] })
|
||||
.resolves.toMatchObject({ events: [{ type: 'plugin/test', data: null, ignorable: true }] })
|
||||
|
||||
for (const type of ['user/message', 'assistant/message'] as const) {
|
||||
const missingContentId = SessionId(`invalid-${type}-without-content`)
|
||||
@@ -1321,14 +1325,60 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects an unknown format version on load (assertVersion)', async () => {
|
||||
it('rejects a newer format version on load, naming the upgrade direction', async () => {
|
||||
const fix = await makeFixture()
|
||||
const { ctx, fiber } = await freshCtx(fix)
|
||||
try {
|
||||
const m = { version: 99, id: SessionId('v99'), createdAt: 1, cwd: WORK }
|
||||
await ctx.sessionPersistence.create(m)
|
||||
await ctx.sessionPersistence.append(m.id, oneTurnLog())
|
||||
await expect(ctx.sessionPersistence.load(m.id)).rejects.toThrow(/version/)
|
||||
const failure = await ctx.sessionPersistence.load(m.id).then(() => undefined, (error: unknown) => error as Error)
|
||||
expect(failure?.name).toBe('SessionFormatUnsupportedError')
|
||||
expect(failure?.message).toMatch(/written by a newer harness.*upgrade the harness/)
|
||||
} finally {
|
||||
await fiber.dispose()
|
||||
await fix.cleanup()
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects an older format version on load without claiming an upgrade path', async () => {
|
||||
const fix = await makeFixture()
|
||||
const { ctx, fiber } = await freshCtx(fix)
|
||||
try {
|
||||
const m = { version: -1, id: SessionId('v-older'), createdAt: 1, cwd: WORK }
|
||||
await ctx.sessionPersistence.create(m)
|
||||
await ctx.sessionPersistence.append(m.id, oneTurnLog())
|
||||
const failure = await ctx.sessionPersistence.load(m.id).then(() => undefined, (error: unknown) => error as Error)
|
||||
expect(failure?.name).toBe('SessionFormatUnsupportedError')
|
||||
expect(failure?.message).toMatch(/older than the supported v0.*no upgrade path/)
|
||||
} finally {
|
||||
await fiber.dispose()
|
||||
await fix.cleanup()
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects an unknown event type on load unless the event is marked ignorable', async () => {
|
||||
const fix = await makeFixture()
|
||||
const { ctx, fiber } = await freshCtx(fix)
|
||||
try {
|
||||
const required = meta('unknown-required', WORK)
|
||||
await ctx.sessionPersistence.create(required)
|
||||
await ctx.sessionPersistence.append(required.id, [
|
||||
...oneTurnLog(),
|
||||
{ type: 'future/event', seq: oneTurnLog().length, time: 99, data: { payload: 1 } } as unknown as SessionEvent,
|
||||
])
|
||||
const failure = await ctx.sessionPersistence.load(required.id).then(() => undefined, (error: unknown) => error as Error)
|
||||
expect(failure?.name).toBe('SessionFormatUnsupportedError')
|
||||
expect(failure?.message).toMatch(/event type "future\/event".*not marked ignorable/)
|
||||
|
||||
const skippable = meta('unknown-ignorable', WORK)
|
||||
await ctx.sessionPersistence.create(skippable)
|
||||
await ctx.sessionPersistence.append(skippable.id, [
|
||||
...oneTurnLog(),
|
||||
{ type: 'future/event', seq: oneTurnLog().length, time: 99, data: { payload: 1 }, ignorable: true } as unknown as SessionEvent,
|
||||
])
|
||||
const loaded = await ctx.sessionPersistence.load(skippable.id)
|
||||
expect(loaded.events.some(event => (event.type as string) === 'future/event')).toBe(true)
|
||||
} finally {
|
||||
await fiber.dispose()
|
||||
await fix.cleanup()
|
||||
|
||||
Reference in New Issue
Block a user