Reword the suite's module doc to describe what it IS — each scenario lives here once and runs per backend through the fixture — rather than narrating that the scenarios were previously duplicated in the per-backend specs. Per the repo doc-current-state convention (no process/history in comments).
787 lines
37 KiB
TypeScript
787 lines
37 KiB
TypeScript
/**
|
|
* Reusable ORCHESTRATION suite for any backend that composes a
|
|
* {@link PersistenceCoordinator}. Where {@link runPersistenceContract} (in
|
|
* contract.ts) pins the public read/write SEMANTICS, this suite pins the
|
|
* coordinator's WRITE-PATH ORCHESTRATION — the behavior that is identical across
|
|
* every first-party backend because it lives in the shared coordinator, not in
|
|
* the storage primitives: the `session/created` → `session/event` →
|
|
* `session/flush` → dispose drain, lazy materialization, fork-seed persistence,
|
|
* the four `onCreated` adoption cases (new / HMR-adopt / collision /
|
|
* ownerless-claim), crash-tail repair on load, and dispose-time quiescence.
|
|
*
|
|
* A backend imports {@link runCoordinatorContract} and calls it with a
|
|
* {@link CoordinatorFixture} factory that knows how to (a) mount the REAL
|
|
* backend plugin on a {@link Context} over a SHARED storage scope (so HMR/reload
|
|
* tests can dispose one instance and mount another over the same bytes/rows),
|
|
* and (b) inject a never-committed torn tail for one session
|
|
* ({@link CoordinatorFixture.corruptTail}) so the through-coordinator torn-tail
|
|
* repair branch is exercised against real storage. The suite drives everything
|
|
* through the PUBLIC {@link SessionPersistence} API + the cordis SessionStore
|
|
* write path — never the storage primitives directly — so it runs unchanged for
|
|
* every backend (memory / jsonl / sqlite).
|
|
*
|
|
* Each scenario lives here once and runs once per backend through the fixture;
|
|
* the per-backend specs keep ONLY their storage-mechanics tests.
|
|
*
|
|
* @module @deepseek-ai/dsh-session-persistence/tests/coordinator-contract
|
|
*/
|
|
|
|
import { describe, expect, it } from 'vitest'
|
|
import { Context, type Fiber } from 'cordis'
|
|
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
|
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
|
|
import type { SessionPersistence } from '../src/index.ts'
|
|
import { meta, oneTurnLog } from './contract.ts'
|
|
|
|
/**
|
|
* The backend-specific capabilities the orchestration suite needs beyond the
|
|
* public service API. A fresh fixture is created per test (isolated storage);
|
|
* the suite mounts/disposes backend instances on it and cleans it up at the end.
|
|
*/
|
|
export interface CoordinatorFixture {
|
|
/**
|
|
* Mount the REAL backend plugin (via `ctx.plugin`, the Loader path) on `ctx`,
|
|
* over THIS fixture's shared storage scope. Returns the plugin fiber so the
|
|
* suite can dispose a single instance (HMR/reload) while the storage — and any
|
|
* still-live session in another fiber — survives. The caller has already
|
|
* mounted `SessionStore` on `ctx`.
|
|
*/
|
|
mount: (ctx: Context) => Promise<Fiber>
|
|
|
|
/**
|
|
* Inject a NEVER-COMMITTED torn tail into the backend's storage for `id` at
|
|
* the given `cwd` (the cwd the session was created with): a half-written
|
|
* record past the committed region (JSONL: a partial line with no newline;
|
|
* SQLite: a row with invalid `data` JSON past the committed seq). This drives
|
|
* the coordinator's `loadCore` `tornMarker !== undefined` → `commitRepair`
|
|
* branch against real storage.
|
|
*
|
|
* OMITTED by a backend that structurally has no torn tails (memory): the
|
|
* torn-tail scenario then self-skips (asserted explicitly in the suite).
|
|
*/
|
|
corruptTail?: (id: SessionId, cwd: string | undefined) => Promise<void>
|
|
|
|
/** Tear down the storage scope (remove the temp dir / file). */
|
|
cleanup: () => Promise<void>
|
|
}
|
|
|
|
/** A constant absolute cwd; jsonl keys directories off it, memory/sqlite ignore it. */
|
|
const WORK = '/w'
|
|
const OTHER = '/other'
|
|
|
|
/** The per-session init map a backend exposes for white-box init awaits. */
|
|
function inits(persistence: SessionPersistence): Map<Session, Promise<void>> {
|
|
return (persistence as unknown as { inits: Map<Session, Promise<void>> }).inits
|
|
}
|
|
|
|
/** Append a whole event log to a live session, event by event (drives session/event). */
|
|
function send(session: Session, events: readonly SessionEvent[]): void {
|
|
for (const e of events) session.append(e.type, e.data)
|
|
}
|
|
|
|
/** A live session created inside its OWN fiber, so it survives a backend reload. */
|
|
async function liveSessionInFiber(
|
|
ctx: Context, id: string, cwd: string | undefined,
|
|
): Promise<Session> {
|
|
let session!: Session
|
|
await ctx.plugin(Object.assign((inner: Context) => {
|
|
session = inner.sessions.create(id, cwd !== undefined ? { meta: { cwd } } : undefined)
|
|
}, { inject: ['sessions'] }))
|
|
return session
|
|
}
|
|
|
|
/**
|
|
* Run the coordinator orchestration suite against a backend. `makeFixture()`
|
|
* MUST return a fresh fixture (isolated storage) each call.
|
|
*/
|
|
export function runCoordinatorContract(name: string, makeFixture: () => Promise<CoordinatorFixture>): void {
|
|
describe(`PersistenceCoordinator orchestration: ${name}`, () => {
|
|
/** Mount SessionStore + a backend instance on a fresh context over the fixture's storage. */
|
|
async function freshCtx(fix: CoordinatorFixture): Promise<{ ctx: Context; fiber: Fiber }> {
|
|
const ctx = new Context()
|
|
await ctx.plugin(SessionStore)
|
|
const fiber = await fix.mount(ctx)
|
|
return { ctx, fiber }
|
|
}
|
|
|
|
// --- write path: live session → flush → reload ---
|
|
|
|
it('persists a live session driven through the store, surviving reload', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const session = ctx.sessions.create('live', { meta: { cwd: WORK } })
|
|
send(session, oneTurnLog())
|
|
await ctx.parallel('session/flush', session)
|
|
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('live'))
|
|
expect(loaded.events).toHaveLength(6)
|
|
expect(loaded.meta.cwd).toBe(WORK)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('snapshot-on-buffer: mutating an event after session/event does not corrupt the persisted copy', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const session = ctx.sessions.create('mutate', { meta: { cwd: WORK } })
|
|
const ev = session.append('user/message', { content: [{ type: 'text', text: 'original' }], source: { kind: 'user' } })
|
|
// Mutate the live event object AFTER it was buffered by session/event.
|
|
;(ev.data as { content: { type: 'text'; text: string }[] }).content[0]!.text = 'HACKED'
|
|
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', session)
|
|
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('mutate'))
|
|
const first = loaded.events[0]
|
|
expect(first?.type === 'user/message' && (first.data.content[0] as { text: string }).text).toBe('original')
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('append snapshots the batch: mutating the caller array/events after the call is ignored', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const m = meta('snapshot', WORK)
|
|
await ctx.sessionPersistence.create(m)
|
|
const events = oneTurnLog() // seqs 0..5
|
|
const userMsg = events[1] // the user/message event
|
|
const p = ctx.sessionPersistence.append(m.id, events)
|
|
// Mutate the caller's array AND an event object after the call but before
|
|
// the queued op runs: the snapshot taken at call time must shield the copy.
|
|
events.push({ type: 'turn/start', seq: 6, time: 99, data: { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } } })
|
|
if (userMsg?.type === 'user/message') userMsg.data.content = [{ type: 'text', text: 'MUTATED' }]
|
|
await p
|
|
const loaded = await ctx.sessionPersistence.load(m.id)
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5]) // not 0..6
|
|
const persisted = JSON.stringify(loaded.events)
|
|
expect(persisted).toContain('hi') // original content
|
|
expect(persisted).not.toContain('MUTATED')
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- fork / resume ---
|
|
|
|
it('fork: a seeded new session persists its seed once (no double-write on a no-op flush)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const seed = oneTurnLog()
|
|
// A fork: a brand-new id whose seed came from elsewhere.
|
|
const forked = ctx.sessions.create('forked', { seed, meta: { cwd: WORK } })
|
|
await inits(ctx.sessionPersistence).get(forked) // onCreated persisted the seed
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('forked'))
|
|
expect(loaded.events).toEqual(seed)
|
|
// A flush with no NEW events must not double-write.
|
|
await ctx.parallel('session/flush', forked)
|
|
const reloaded = await ctx.sessionPersistence.load(SessionId('forked'))
|
|
expect(reloaded.events).toEqual(seed)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('resume: a re-created session seeded with the loaded log does not re-append its seed and continues the seq', async () => {
|
|
const fix = await makeFixture()
|
|
const first = await freshCtx(fix)
|
|
try {
|
|
// First lifecycle: persist a session through the store.
|
|
const s1 = first.ctx.sessions.create('resumed', { meta: { cwd: WORK } })
|
|
send(s1, oneTurnLog())
|
|
await first.ctx.parallel('session/flush', s1)
|
|
} finally {
|
|
await first.fiber.dispose()
|
|
}
|
|
|
|
// Second lifecycle: a NEW backend instance + a session re-created with the
|
|
// same id SEEDED with the loaded events. onCreated adopts the stored log
|
|
// (does not re-persist the seed); a new turn appends at seq 6.
|
|
const second = await freshCtx(fix)
|
|
try {
|
|
const loaded = await second.ctx.sessionPersistence.load(SessionId('resumed'))
|
|
const s2 = second.ctx.sessions.create('resumed', { seed: loaded.events, meta: { cwd: WORK } })
|
|
await inits(second.ctx.sessionPersistence).get(s2) // let onCreated adopt
|
|
s2.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
s2.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
|
|
await second.ctx.parallel('session/flush', s2)
|
|
|
|
const reloaded = await second.ctx.sessionPersistence.load(SessionId('resumed'))
|
|
// 6 original + 2 new, contiguous, no duplicated seed.
|
|
expect(reloaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
|
|
} finally {
|
|
await second.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- HMR ---
|
|
|
|
it('HMR: applying the plugin seeds existing live sessions', async () => {
|
|
const fix = await makeFixture()
|
|
const ctx = new Context()
|
|
await ctx.plugin(SessionStore)
|
|
// A session exists BEFORE the persistence plugin is applied.
|
|
const session = ctx.sessions.create('pre-existing', { meta: { cwd: WORK } })
|
|
session.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
|
const fiber = await fix.mount(ctx)
|
|
try {
|
|
// The plugin seeded it on apply; a subsequent flush persists its events.
|
|
await ctx.parallel('session/flush', session)
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('pre-existing'))
|
|
expect(loaded.events.length).toBeGreaterThanOrEqual(2)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('HMR: dispose drains remaining buffers', async () => {
|
|
const fix = await makeFixture()
|
|
const ctx = new Context()
|
|
await ctx.plugin(SessionStore)
|
|
const fiber = await fix.mount(ctx)
|
|
const session = await liveSessionInFiber(ctx, 'drain', WORK)
|
|
session.append('user/message', { content: [{ type: 'text', text: 'buffered' }], source: { kind: 'user' } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
// No explicit flush — dispose must drain.
|
|
await fiber.dispose()
|
|
|
|
// A fresh backend instance reads what the disposed one drained.
|
|
const second = await freshCtx(fix)
|
|
try {
|
|
const loaded = await second.ctx.sessionPersistence.load(SessionId('drain'))
|
|
expect(loaded.events.length).toBeGreaterThanOrEqual(2)
|
|
} finally {
|
|
await second.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('HMR: reloading the backend adopts a still-live, already-materialized session', async () => {
|
|
const fix = await makeFixture()
|
|
const ctx = new Context()
|
|
await ctx.plugin(SessionStore)
|
|
// The session lives in its OWN fiber so it survives the backend reload.
|
|
const session = await liveSessionInFiber(ctx, 'hmr-adopt', WORK)
|
|
try {
|
|
// Backend instance 1 materializes the session.
|
|
const backend1 = await fix.mount(ctx)
|
|
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
session.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', session)
|
|
|
|
// Hot-reload: dispose instance 1, mount instance 2 over the SAME storage
|
|
// while the session stays live. Instance 2 has an empty states map but the
|
|
// log is materialized and is a prefix of the live events — it must ADOPT
|
|
// (not reject). A second turn appended after reload then persists.
|
|
await backend1.dispose()
|
|
await fix.mount(ctx)
|
|
session.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
session.append('user/message', { content: [{ type: 'text', text: 'again' }], source: { kind: 'user' } })
|
|
session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
|
|
await expect(ctx.parallel('session/flush', session)).resolves.not.toThrow()
|
|
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('hmr-adopt'))
|
|
expect(loaded.events.filter(e => e.type === 'turn/start')).toHaveLength(2)
|
|
} finally {
|
|
await ctx.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('HMR: adoption persists the live SUFFIX that was ahead of the stored prefix', async () => {
|
|
const fix = await makeFixture()
|
|
const ctx = new Context()
|
|
await ctx.plugin(SessionStore)
|
|
const session = await liveSessionInFiber(ctx, 'hmr-suffix', WORK)
|
|
try {
|
|
// Instance 1 flushes turn 1.
|
|
const backend1 = await fix.mount(ctx)
|
|
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', session)
|
|
|
|
// Append turn 2 to the LIVE session, then dispose instance 1 WITHOUT
|
|
// flushing turn 2: it is now ONLY in the live session's events; the new
|
|
// backend never buffered it via session/event.
|
|
await backend1.dispose()
|
|
session.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
|
|
|
|
// Instance 2 adopts the stored prefix (turn 1) and MUST also persist the
|
|
// live suffix (turn 2) carried in the session's events.
|
|
await fix.mount(ctx)
|
|
await ctx.parallel('session/flush', session)
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('hmr-suffix'))
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3])
|
|
expect(loaded.events.filter(e => e.type === 'turn/start')).toHaveLength(2)
|
|
} finally {
|
|
await ctx.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('HMR adoption does NOT crash-repair an active open turn as interrupted (truncate without closers)', async () => {
|
|
const fix = await makeFixture()
|
|
const ctx = new Context()
|
|
await ctx.plugin(SessionStore)
|
|
const session = await liveSessionInFiber(ctx, 'hmr-open', WORK)
|
|
try {
|
|
const first = await fix.mount(ctx)
|
|
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
session.append('step/start', { turn: 1, step: 1 })
|
|
await ctx.parallel('session/flush', session)
|
|
|
|
// Crash-tail a torn fragment past the (open) committed turn, then reload.
|
|
await first.dispose()
|
|
if (fix.corruptTail) await fix.corruptTail(SessionId('hmr-open'), WORK)
|
|
const second = await fix.mount(ctx)
|
|
// The live session is still the authority: it appends the REAL step/turn
|
|
// end. Adoption must truncate the torn tail but NOT synthesize closers.
|
|
session.append('step/end', { turn: 1, step: 1 })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', session)
|
|
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('hmr-open'))
|
|
expect(loaded.events.map(e => e.type)).toEqual(['turn/start', 'step/start', 'step/end', 'turn/end'])
|
|
expect(loaded.events.at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'completed' } } })
|
|
await second.dispose()
|
|
} finally {
|
|
await ctx.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- collision / id reuse ---
|
|
|
|
it('a NEW live session colliding on a persisted id is rejected, not silently adopted', async () => {
|
|
const fix = await makeFixture()
|
|
const first = await freshCtx(fix)
|
|
try {
|
|
const s1 = first.ctx.sessions.create('collide', { meta: { cwd: WORK } })
|
|
send(s1, oneTurnLog())
|
|
await first.ctx.parallel('session/flush', s1)
|
|
} finally {
|
|
await first.fiber.dispose()
|
|
}
|
|
|
|
// A FRESH backend + a NEW live session with the same id but NO explicit
|
|
// resume. onCreated treats it as new; create() rejects because a log already
|
|
// exists. The rejection surfaces via the init promise (flush awaits it).
|
|
const second = await freshCtx(fix)
|
|
try {
|
|
const s2 = second.ctx.sessions.create('collide', { meta: { cwd: WORK } })
|
|
s2.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
await expect(inits(second.ctx.sessionPersistence).get(s2))
|
|
.rejects.toThrow(/already has a persisted log|id collision/)
|
|
} finally {
|
|
await second.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('an abandoned lazy session (never materialized) releases its id for reuse', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// A live session created then disposed BEFORE its first append: cursor 0,
|
|
// never materialized. A new live session reusing the id must reclaim it.
|
|
let firstSession!: Session
|
|
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
firstSession = inner.sessions.create('abandoned', { meta: { cwd: WORK } })
|
|
}, { inject: ['sessions'] }))
|
|
await inits(ctx.sessionPersistence).get(firstSession) // register the lazy state
|
|
await firstFiber.dispose() // disposed before any append → never materialized
|
|
|
|
let reuse!: Session
|
|
await ctx.plugin(Object.assign((inner: Context) => {
|
|
reuse = inner.sessions.create('abandoned', { meta: { cwd: WORK } })
|
|
}, { inject: ['sessions'] }))
|
|
await expect(inits(ctx.sessionPersistence).get(reuse)).resolves.toBeUndefined()
|
|
reuse.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
reuse.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', reuse)
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('abandoned'))
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1])
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('does NOT reclaim an id whose abandoned owner still has buffered (unflushed) events', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
let first!: Session
|
|
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
first = inner.sessions.create('buffered', { meta: { cwd: WORK } })
|
|
}, { inject: ['sessions'] }))
|
|
await inits(ctx.sessionPersistence).get(first)
|
|
// Append a turn but do NOT flush — events sit in the write-behind buffer.
|
|
first.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await firstFiber.dispose() // disposed before flush; not materialized, buffer pending
|
|
|
|
let reuse!: Session
|
|
await ctx.plugin(Object.assign((inner: Context) => {
|
|
reuse = inner.sessions.create('buffered', { meta: { cwd: WORK } })
|
|
}, { inject: ['sessions'] }))
|
|
await expect(inits(ctx.sessionPersistence).get(reuse)).rejects.toThrow(/already bound to a different live session/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('initFor is idempotent: re-emitting session/created does not re-initialize', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const session = ctx.sessions.create('idem', { meta: { cwd: WORK } })
|
|
session.append('user/message', { content: [{ type: 'text', text: 'x' }], source: { kind: 'user' } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', session)
|
|
// Re-emit session/created for the SAME live session (idempotent initFor).
|
|
ctx.emit('session/created', session)
|
|
await ctx.parallel('session/flush', session)
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('idem'))
|
|
expect(loaded.events).toHaveLength(2) // not doubled
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- ownerless-state claim (public create()/load() then a live session arrives) ---
|
|
|
|
it('a live session claims cursor-0 ownerless state created via the public API and persists its seed', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// create() registers ownerless state with cursor 0 (lazy, nothing persisted).
|
|
await ctx.sessionPersistence.create(meta('lazy-claim', WORK))
|
|
// A live session with that id arrives and claims it (cursor 0 matches
|
|
// trivially), persisting its seed.
|
|
const live = ctx.sessions.create('lazy-claim', { seed: oneTurnLog(), meta: { cwd: WORK } })
|
|
await expect(inits(ctx.sessionPersistence).get(live)).resolves.toBeUndefined()
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('lazy-claim'))
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5])
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('a fresh session reusing a previously-loaded id is rejected (ownerless guard)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// Materialize a log, then load() it WITHOUT a live session — ownerless
|
|
// state, cursor at the persisted length.
|
|
await ctx.sessionPersistence.create(meta('preview', WORK))
|
|
await ctx.sessionPersistence.append(SessionId('preview'), oneTurnLog())
|
|
await ctx.sessionPersistence.load(SessionId('preview'))
|
|
|
|
// A FRESH (empty-seed) live session reusing that id must be rejected: its
|
|
// seq 0..cursor-1 events would otherwise be filtered as already-persisted.
|
|
let fresh!: Session
|
|
await ctx.plugin(Object.assign((inner: Context) => {
|
|
fresh = inner.sessions.create('preview', { meta: { cwd: WORK } })
|
|
}, { inject: ['sessions'] }))
|
|
await expect(inits(ctx.sessionPersistence).get(fresh))
|
|
.rejects.toThrow(/do not match this live session|already has a persisted log|id collision/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('a live session whose seed matches the loaded prefix claims ownerless state and persists the suffix', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// Materialize and load (ownerless, cursor = 6).
|
|
await ctx.sessionPersistence.create(meta('claim', WORK))
|
|
await ctx.sessionPersistence.append(SessionId('claim'), oneTurnLog())
|
|
const { events } = await ctx.sessionPersistence.load(SessionId('claim'))
|
|
|
|
// A live session SEEDED with the loaded log PLUS a new turn claims the
|
|
// ownerless state and persists only the suffix.
|
|
const cont = ctx.sessions.create('claim', { seed: [
|
|
...events,
|
|
{ type: 'turn/start', seq: 6, time: 7, data: { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } } },
|
|
{ type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
|
|
], meta: { cwd: WORK } })
|
|
await inits(ctx.sessionPersistence).get(cont)
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('claim'))
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('a live session at a DIFFERENT cwd cannot claim cursor-0 ownerless state (cwd scope)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// create() registers ownerless state at cwd /a (cursor 0 — claims would
|
|
// otherwise match trivially on the seed).
|
|
await ctx.sessionPersistence.create(meta('wrong-cwd-claim', OTHER))
|
|
// A live session reusing the id but at cwd WORK must NOT claim it — the
|
|
// cwd scope is the fence (without it, WORK events would append under the
|
|
// OTHER header). Rejected as a collision.
|
|
const live = ctx.sessions.create('wrong-cwd-claim', { seed: oneTurnLog(), meta: { cwd: WORK } })
|
|
await expect(inits(ctx.sessionPersistence).get(live)).rejects.toThrow(/different cwd|id collision/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('a live session at a DIFFERENT cwd cannot claim loaded-prefix ownerless state (cwd scope)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// Materialize + load at cwd OTHER (ownerless, cursor = 6).
|
|
await ctx.sessionPersistence.create(meta('wrong-cwd-load', OTHER))
|
|
await ctx.sessionPersistence.append(SessionId('wrong-cwd-load'), oneTurnLog())
|
|
const { events } = await ctx.sessionPersistence.load(SessionId('wrong-cwd-load'))
|
|
// A live session whose SEED matches the loaded prefix but whose cwd is
|
|
// WORK must still be rejected — the cwd guard runs before the seed check.
|
|
const live = ctx.sessions.create('wrong-cwd-load', { seed: events, meta: { cwd: WORK } })
|
|
await expect(inits(ctx.sessionPersistence).get(live)).rejects.toThrow(/different cwd|id collision/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('a no-cwd ownerless state cannot be claimed by a live session WITH a cwd (cwd scope, undefined side)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// Ownerless state created WITHOUT a cwd (the no-cwd bucket).
|
|
await ctx.sessionPersistence.create(meta('no-cwd-state'))
|
|
// A live session reusing the id but WITH cwd WORK is a cwd mismatch
|
|
// (undefined vs WORK) and must be rejected.
|
|
const live = ctx.sessions.create('no-cwd-state', { seed: oneTurnLog(), meta: { cwd: WORK } })
|
|
await expect(inits(ctx.sessionPersistence).get(live)).rejects.toThrow(/different cwd|id collision/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- append adopts a storage-only session (fresh instance, no prior create/load) ---
|
|
|
|
it('append adopts a storage-only session (fresh instance) and continues the seq', async () => {
|
|
const fix = await makeFixture()
|
|
const first = await freshCtx(fix)
|
|
try {
|
|
const m = meta('adopt-append', WORK)
|
|
await first.ctx.sessionPersistence.create(m)
|
|
await first.ctx.sessionPersistence.append(m.id, oneTurnLog())
|
|
} finally {
|
|
await first.fiber.dispose()
|
|
}
|
|
|
|
// A fresh instance appends a second turn WITHOUT a prior create/load: append
|
|
// must adopt the stored session (cursor = stored length) and continue.
|
|
const second = await freshCtx(fix)
|
|
try {
|
|
await second.ctx.sessionPersistence.append(SessionId('adopt-append'), [
|
|
{ type: 'turn/start', seq: 6, time: 7, data: { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } } },
|
|
{ type: 'turn/end', seq: 7, time: 8, data: { turn: 2, reason: { kind: 'completed' } } },
|
|
])
|
|
const loaded = await second.ctx.sessionPersistence.load(SessionId('adopt-append'))
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
|
|
} finally {
|
|
await second.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- small public-API edges that the coordinator owns uniformly ---
|
|
|
|
it('append of an empty batch is a no-op (stays lazy)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const m = meta('empty-batch', WORK)
|
|
await ctx.sessionPersistence.create(m)
|
|
await ctx.sessionPersistence.append(m.id, [])
|
|
expect(await ctx.sessionPersistence.has(m.id)).toBe(false)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('load rejects a missing session', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
await expect(ctx.sessionPersistence.load(SessionId('nope'))).rejects.toThrow(/not found/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('delete of a non-existent session is a no-op', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
await expect(ctx.sessionPersistence.delete(SessionId('ghost'))).resolves.toBeUndefined()
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('create rejects a duplicate id (in memory and on a persisted log)', async () => {
|
|
const fix = await makeFixture()
|
|
const first = await freshCtx(fix)
|
|
try {
|
|
const m = meta('dup', WORK)
|
|
await first.ctx.sessionPersistence.create(m)
|
|
// Same in-memory state.
|
|
await expect(first.ctx.sessionPersistence.create(m)).rejects.toThrow(/already exists in this backend/)
|
|
await first.ctx.sessionPersistence.append(m.id, oneTurnLog())
|
|
} finally {
|
|
await first.fiber.dispose()
|
|
}
|
|
|
|
// A fresh instance over the same storage sees the persisted log.
|
|
const second = await freshCtx(fix)
|
|
try {
|
|
await expect(second.ctx.sessionPersistence.create(meta('dup', WORK)))
|
|
.rejects.toThrow(/already has a persisted log on disk/)
|
|
} finally {
|
|
await second.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('rejects an unknown format version on load (assertVersion)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const m = { version: 2, id: SessionId('v2'), 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/)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('round-trips a header with parentSession (fork lineage)', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
const m = { version: 1, id: SessionId('forked-child'), createdAt: 1, cwd: WORK, parentSession: SessionId('the-parent') }
|
|
await ctx.sessionPersistence.create(m)
|
|
await ctx.sessionPersistence.append(m.id, oneTurnLog())
|
|
const loaded = await ctx.sessionPersistence.load(m.id)
|
|
expect(loaded.meta.parentSession).toBe('the-parent')
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
it('flush before init resolves uses cursor 0', async () => {
|
|
const fix = await makeFixture()
|
|
const { ctx, fiber } = await freshCtx(fix)
|
|
try {
|
|
// Append directly to a live session and flush IMMEDIATELY, before the
|
|
// async onCreated init has necessarily set state (exercises the
|
|
// state-undefined cursor path).
|
|
const session = ctx.sessions.create('flush-nostate', { meta: { cwd: WORK } })
|
|
session.append('user/message', { content: [{ type: 'text', text: 'q' }], source: { kind: 'user' } })
|
|
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
await ctx.parallel('session/flush', session)
|
|
const loaded = await ctx.sessionPersistence.load(SessionId('flush-nostate'))
|
|
expect(loaded.events).toHaveLength(2)
|
|
} finally {
|
|
await fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
|
|
// --- crash-tail repair THROUGH the coordinator (real storage torn tail) ---
|
|
|
|
it('torn-tail load: a never-committed tail is truncated and the open turn closed during load (commitRepair w/ tornMarker)', async () => {
|
|
const fix = await makeFixture()
|
|
if (!fix.corruptTail) {
|
|
// A memory-style store has no torn tails (every write is atomic in RAM),
|
|
// so there is no tornMarker path to exercise. Assert that explicitly
|
|
// instead of silently skipping, then bail.
|
|
expect(fix.corruptTail).toBeUndefined()
|
|
await fix.cleanup()
|
|
return
|
|
}
|
|
const first = await freshCtx(fix)
|
|
try {
|
|
const m = meta('torn', WORK)
|
|
await first.ctx.sessionPersistence.create(m)
|
|
await first.ctx.sessionPersistence.append(m.id, oneTurnLog()) // committed 0..5 (balanced)
|
|
// A second turn whose real events are durable but never closed (open turn).
|
|
await first.ctx.sessionPersistence.append(m.id, [
|
|
{ type: 'turn/start', seq: 6, time: 7, data: { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } } },
|
|
{ type: 'step/start', seq: 7, time: 8, data: { turn: 2, step: 1 } },
|
|
])
|
|
} finally {
|
|
await first.fiber.dispose()
|
|
}
|
|
// Inject a torn fragment past the committed region (never-committed tail).
|
|
await fix.corruptTail(SessionId('torn'), WORK)
|
|
|
|
// A FRESH instance loads: the torn tail is truncated (tornMarker !==
|
|
// undefined) AND the open turn 2 is closed with synthetic step/end +
|
|
// turn/end {interrupted} — commitRepair runs with BOTH a torn marker and
|
|
// closers. The preserved real events (0..7) are never truncated.
|
|
const second = await freshCtx(fix)
|
|
try {
|
|
const loaded = await second.ctx.sessionPersistence.load(SessionId('torn'))
|
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8, 9])
|
|
expect(loaded.events.map(e => e.type)).toEqual([
|
|
'turn/start', 'user/message', 'step/start', 'assistant/message', 'step/end', 'turn/end', // turn 1
|
|
'turn/start', 'step/start', 'step/end', 'turn/end', // turn 2: real + synthetic closers
|
|
])
|
|
const last = loaded.events.at(-1)!
|
|
expect(last.type === 'turn/end' && last.data.reason).toEqual({ kind: 'interrupted' })
|
|
|
|
// The repair is durable: the next append continues at the balanced length
|
|
// (seq 10) and a reload round-trips identically.
|
|
await second.ctx.sessionPersistence.append(SessionId('torn'), [
|
|
{ type: 'turn/start', seq: 10, time: 9, data: { turn: 3, trigger: { kind: 'message', source: { kind: 'user' } } } },
|
|
{ type: 'turn/end', seq: 11, time: 10, data: { turn: 3, reason: { kind: 'completed' } } },
|
|
])
|
|
const reloaded = await second.ctx.sessionPersistence.load(SessionId('torn'))
|
|
expect(reloaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11])
|
|
} finally {
|
|
await second.fiber.dispose()
|
|
await fix.cleanup()
|
|
}
|
|
})
|
|
})
|
|
}
|