fix(persistence): restore pre-identity sessions
This commit is contained in:
18 files changed
+457
-33
No files matched your search
@@ -146,10 +146,142 @@ function assertSupportedEvents(events: readonly SessionEvent[], id: SessionId):
|
||||
}
|
||||
}
|
||||
|
||||
/** Materialize stored events as validated snapshots with immutable messages. */
|
||||
/** Return an object record without widening arrays into message payloads. */
|
||||
function asRecord(value: unknown): Record<string, unknown> | undefined {
|
||||
return typeof value === 'object' && value !== null && !Array.isArray(value)
|
||||
? value as Record<string, unknown>
|
||||
: undefined
|
||||
}
|
||||
|
||||
type PersistedMessageId = SessionEvent<'user/message'>['data']['id']
|
||||
|
||||
/** Mint the stable import identity for a message persisted before identities existed. */
|
||||
function legacyMessageId(id: SessionId, seq: number): PersistedMessageId {
|
||||
return `legacy-message:${id}:${seq}` as PersistedMessageId
|
||||
}
|
||||
|
||||
/** Read a replacement target while leaving malformed surface metadata to the session validator. */
|
||||
function replacementStart(event: SessionEvent): number | undefined {
|
||||
const op = asRecord((event as SessionEvent & { surfaceOp?: unknown }).surfaceOp)
|
||||
return op?.['op'] === 'replace' && typeof op['start'] === 'number'
|
||||
? op['start']
|
||||
: undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Upgrade one pre-identity message event into the current wrapper shape.
|
||||
* Current-looking malformed events remain untouched so validation rejects them
|
||||
* instead of disguising corruption as legacy data.
|
||||
*/
|
||||
function migrateLegacyMessageEvent(
|
||||
event: SessionEvent,
|
||||
id: SessionId,
|
||||
messageIds: ReadonlyMap<number, PersistedMessageId>,
|
||||
): SessionEvent {
|
||||
const data = asRecord(event.data)
|
||||
if (data === undefined) return event
|
||||
switch (event.type) {
|
||||
case 'user/message': {
|
||||
if (Object.hasOwn(data, 'id') || Object.hasOwn(data, 'role')
|
||||
|| Object.hasOwn(data, 'message')
|
||||
|| !Object.hasOwn(data, 'content') || !Object.hasOwn(data, 'source')) return event
|
||||
return {
|
||||
...event,
|
||||
data: {
|
||||
...data,
|
||||
id: legacyMessageId(id, event.seq),
|
||||
role: 'user',
|
||||
},
|
||||
} as SessionEvent
|
||||
}
|
||||
case 'assistant/message': {
|
||||
if (Object.hasOwn(data, 'message')
|
||||
|| !Object.hasOwn(data, 'content') || !Object.hasOwn(data, 'provenance')) return event
|
||||
const { content, provenance, ...eventData } = data
|
||||
return {
|
||||
...event,
|
||||
data: {
|
||||
...eventData,
|
||||
message: {
|
||||
id: legacyMessageId(id, event.seq),
|
||||
role: 'assistant',
|
||||
content,
|
||||
source: {
|
||||
...asRecord(provenance),
|
||||
kind: 'model',
|
||||
},
|
||||
},
|
||||
},
|
||||
} as SessionEvent
|
||||
}
|
||||
case 'tool/result': {
|
||||
if (Object.hasOwn(data, 'message')
|
||||
|| !Object.hasOwn(data, 'callId') || !Object.hasOwn(data, 'content')
|
||||
|| !Object.hasOwn(data, 'isError')) return event
|
||||
const { callId, content, isError, ...eventData } = data
|
||||
const inheritedId = replacementStart(event)
|
||||
return {
|
||||
...event,
|
||||
data: {
|
||||
...eventData,
|
||||
message: {
|
||||
id: inheritedId === undefined
|
||||
? legacyMessageId(id, event.seq)
|
||||
: messageIds.get(inheritedId),
|
||||
role: 'user',
|
||||
content: [{
|
||||
type: 'tool-result',
|
||||
toolCallId: callId,
|
||||
content,
|
||||
isError,
|
||||
}],
|
||||
source: {
|
||||
kind: 'tool',
|
||||
callId,
|
||||
},
|
||||
},
|
||||
},
|
||||
} as SessionEvent
|
||||
}
|
||||
case 'steering/message': {
|
||||
if (Object.hasOwn(data, 'message')
|
||||
|| !Object.hasOwn(data, 'content') || !Object.hasOwn(data, 'source')) return event
|
||||
const { content, source, ...eventData } = data
|
||||
return {
|
||||
...event,
|
||||
data: {
|
||||
...eventData,
|
||||
message: {
|
||||
id: legacyMessageId(id, event.seq),
|
||||
role: 'user',
|
||||
content,
|
||||
source,
|
||||
},
|
||||
},
|
||||
} as SessionEvent
|
||||
}
|
||||
default:
|
||||
return event
|
||||
}
|
||||
}
|
||||
|
||||
/** Read the identified message carried by one validated current event. */
|
||||
function eventMessageId(event: SessionEvent): PersistedMessageId | undefined {
|
||||
const data = asRecord(event.data)
|
||||
const message = event.type === 'user/message' ? data : asRecord(data?.['message'])
|
||||
return typeof message?.['id'] === 'string' ? message['id'] as PersistedMessageId : undefined
|
||||
}
|
||||
|
||||
/** Materialize stored events as upgraded, validated snapshots with immutable messages. */
|
||||
function snapshotStoredEvents(events: readonly SessionEvent[], id: SessionId): SessionEvent[] {
|
||||
assertSupportedEvents(events, id)
|
||||
return events.map(snapshotSessionEvent)
|
||||
const messageIds = new Map<number, PersistedMessageId>()
|
||||
return events.map((event) => {
|
||||
const snapshot = snapshotSessionEvent(migrateLegacyMessageEvent(event, id, messageIds))
|
||||
const messageId = eventMessageId(snapshot)
|
||||
if (messageId !== undefined) messageIds.set(snapshot.seq, messageId)
|
||||
return snapshot
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -526,7 +658,7 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
/* v8 ignore next -- a cursor > 0 means the session was materialized, so it exists */
|
||||
if (stored === undefined) return false
|
||||
this.assertStoredId(id, stored.meta)
|
||||
return seedCoversPrefix(seed, stored.events.slice(0, cursor))
|
||||
return seedCoversPrefix(seed, snapshotStoredEvents(stored.events, id).slice(0, cursor))
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -614,19 +746,19 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
throw new Error(`session "${session.header.id}" is already persisted at a different cwd (persisted: ${String(meta.cwd)}, live: ${String(session.header.cwd)}) (id collision)`)
|
||||
}
|
||||
this.assertVersion(meta)
|
||||
assertSupportedEvents(events, session.header.id)
|
||||
if (!seedCoversPrefix(seed, events)) {
|
||||
const storedEvents = snapshotStoredEvents(events, session.header.id)
|
||||
if (!seedCoversPrefix(seed, storedEvents)) {
|
||||
throw new Error(`session "${session.header.id}" already has a persisted log on disk that does not match this live session (id collision)`)
|
||||
}
|
||||
// Truncate-only repair (no closers): the open turn is NOT closed here.
|
||||
if (tornMarker !== undefined) await this.backend.commitRepair(meta, tornMarker, [])
|
||||
this.states.set(session.header.id, {
|
||||
meta: { ...meta },
|
||||
cursor: events.length,
|
||||
cursor: storedEvents.length,
|
||||
materialized: true,
|
||||
owner: session,
|
||||
})
|
||||
const suffix = seed.slice(events.length)
|
||||
const suffix = seed.slice(storedEvents.length)
|
||||
if (suffix.length > 0) await this.appendCore(session.header.id, suffix)
|
||||
}
|
||||
|
||||
|
||||
@@ -93,7 +93,9 @@ export abstract class SessionPersistence extends Service {
|
||||
* A coordinator-backed cold load reserves the identity across storage awaits,
|
||||
* so concurrent publication of a same-id live Session rejects.
|
||||
* Returned events are detached, and every identified message is deeply
|
||||
* frozen; malformed identified messages reject before any stored event is returned.
|
||||
* frozen. Coordinator-backed implementations upgrade supported pre-identity
|
||||
* message events before validation; other malformed messages reject before
|
||||
* any stored event is returned.
|
||||
* @param id - the persisted session to reload.
|
||||
* @returns the header and a log ending on a balanced `turn/end`.
|
||||
*/
|
||||
@@ -103,8 +105,9 @@ export abstract class SessionPersistence extends Service {
|
||||
* Inspect a header and its valid contiguous stored prefix without repairing
|
||||
* a torn tail, closing an interrupted turn, or publishing coordinator state.
|
||||
* This read is serialized with writes for the same id and returns detached
|
||||
* values with deeply frozen identified messages, so observers cannot mutate message
|
||||
* identity/content or backend-owned state. Malformed identified messages reject.
|
||||
* values with upgraded, deeply frozen identified messages, so observers
|
||||
* cannot mutate message identity/content or backend-owned state. Other
|
||||
* malformed messages reject.
|
||||
* @param id - the persisted session to inspect.
|
||||
* @param signal - optional cancellation for queued and backend read work.
|
||||
* @returns the header and valid stored event prefix exactly as observed.
|
||||
|
||||
Reference in New Issue
Block a user