|
|
|
|
@@ -91,12 +91,17 @@ class MemoryPersistence extends SessionPersistence implements PersistenceBackend
|
|
|
|
|
return this.coordinator.append(id, events)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
override prepare(id: SessionId, signal?: AbortSignal): ReturnType<PersistenceCoordinator['prepare']> {
|
|
|
|
|
return this.coordinator.prepare(id, signal)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
|
|
|
|
|
return this.coordinator.load(id)
|
|
|
|
|
return this.coordinator.load(id).then(loaded => ({ meta: loaded.meta, events: [...loaded.events] }))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
|
|
|
|
|
return this.coordinator.inspect(id, signal)
|
|
|
|
|
.then(loaded => ({ meta: loaded.meta, events: [...loaded.events] }))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
readFrom(id: SessionId, fromSeq: number, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
|
|
|
|
|
@@ -170,7 +175,8 @@ class ControlledBackend implements PersistenceBackend<never> {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<never> | undefined> {
|
|
|
|
|
await this.beforeLoadStored?.(++this.loadAttempts, signal)
|
|
|
|
|
const attempt = ++this.loadAttempts
|
|
|
|
|
await this.beforeLoadStored?.(attempt, signal)
|
|
|
|
|
const entry = this.store.get(id)
|
|
|
|
|
if (entry === undefined) return undefined
|
|
|
|
|
return { meta: structuredClone(entry.meta), events: structuredClone(entry.events) }
|
|
|
|
|
@@ -187,8 +193,10 @@ class ControlledBackend implements PersistenceBackend<never> {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async commitRepair(_m: SessionHeader, _tornMarker: undefined, _closers: readonly SessionEvent[]): Promise<void> {
|
|
|
|
|
async commitRepair(m: SessionHeader, _tornMarker: undefined, closers: readonly SessionEvent[]): Promise<void> {
|
|
|
|
|
this.repairAttempts += 1
|
|
|
|
|
const entry = this.store.get(m.id)
|
|
|
|
|
if (entry !== undefined) entry.events.push(...structuredClone(closers) as SessionEvent[])
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async list(): Promise<SessionHeader[]> {
|
|
|
|
|
@@ -345,7 +353,7 @@ describe('PersistenceCoordinator stored identity', () => {
|
|
|
|
|
|
|
|
|
|
await expect(ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
inner.sessions.create(id, { seed: [start], meta: header })
|
|
|
|
|
}, { inject: ['sessions'] }))).rejects.toThrow(/persisted history is loading/)
|
|
|
|
|
}, { inject: ['sessions'] }))).rejects.toThrow(/persisted preparation exists/)
|
|
|
|
|
expect(ctx.sessions.get(id)).toBeUndefined()
|
|
|
|
|
|
|
|
|
|
loadGate.resolve(true)
|
|
|
|
|
@@ -362,6 +370,170 @@ describe('PersistenceCoordinator stored identity', () => {
|
|
|
|
|
})
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
describe('PersistenceCoordinator session preparations', () => {
|
|
|
|
|
it('reuses the exact Session from inspect through repeated unpublished prepare calls', async () => {
|
|
|
|
|
const ctx = new Context()
|
|
|
|
|
await ctx.plugin(SessionStore)
|
|
|
|
|
const backend = new ControlledBackend()
|
|
|
|
|
const id = SessionId('inspect-prepare-reuse')
|
|
|
|
|
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
|
|
|
|
|
let coordinator!: PersistenceCoordinator<never>
|
|
|
|
|
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
|
|
}, { inject: ['sessions'] }))
|
|
|
|
|
let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
const inspected = await coordinator.inspect(id)
|
|
|
|
|
first = await coordinator.prepare(id)
|
|
|
|
|
|
|
|
|
|
expect(backend.loadAttempts).toBe(1)
|
|
|
|
|
expect(first.session.events[0]).toBe(inspected.events[0])
|
|
|
|
|
|
|
|
|
|
first[Symbol.dispose]()
|
|
|
|
|
second = await coordinator.prepare(id)
|
|
|
|
|
expect(second.session).toBe(first.session)
|
|
|
|
|
expect(backend.loadAttempts).toBe(1)
|
|
|
|
|
} finally {
|
|
|
|
|
second?.[Symbol.dispose]()
|
|
|
|
|
first?.[Symbol.dispose]()
|
|
|
|
|
await fiber.dispose()
|
|
|
|
|
await ctx.fiber.dispose()
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
it('keeps synthetic recovery in memory during inspect and commits it only once on prepare', async () => {
|
|
|
|
|
const ctx = new Context()
|
|
|
|
|
await ctx.plugin(SessionStore)
|
|
|
|
|
const backend = new ControlledBackend()
|
|
|
|
|
const id = SessionId('inspect-repair-commit')
|
|
|
|
|
backend.store.set(id, {
|
|
|
|
|
meta: meta(id),
|
|
|
|
|
events: [{
|
|
|
|
|
type: 'turn/start',
|
|
|
|
|
seq: 0,
|
|
|
|
|
time: 1,
|
|
|
|
|
data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
|
|
|
|
}],
|
|
|
|
|
})
|
|
|
|
|
let coordinator!: PersistenceCoordinator<never>
|
|
|
|
|
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
|
|
}, { inject: ['sessions'] }))
|
|
|
|
|
let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
const inspected = await coordinator.inspect(id)
|
|
|
|
|
expect(inspected.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
|
|
|
|
|
expect(backend.store.get(id)?.events.map(event => event.type)).toEqual(['turn/start'])
|
|
|
|
|
expect(backend.repairAttempts).toBe(0)
|
|
|
|
|
|
|
|
|
|
first = await coordinator.prepare(id)
|
|
|
|
|
expect(backend.repairAttempts).toBe(1)
|
|
|
|
|
expect(backend.store.get(id)?.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
|
|
|
|
|
first[Symbol.dispose]()
|
|
|
|
|
|
|
|
|
|
second = await coordinator.prepare(id)
|
|
|
|
|
expect(second.session).toBe(first.session)
|
|
|
|
|
expect(backend.loadAttempts).toBe(1)
|
|
|
|
|
expect(backend.repairAttempts).toBe(1)
|
|
|
|
|
} finally {
|
|
|
|
|
second?.[Symbol.dispose]()
|
|
|
|
|
first?.[Symbol.dispose]()
|
|
|
|
|
await fiber.dispose()
|
|
|
|
|
await ctx.fiber.dispose()
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
it('waits for an existing reservation and reuses it after release', async () => {
|
|
|
|
|
const ctx = new Context()
|
|
|
|
|
await ctx.plugin(SessionStore)
|
|
|
|
|
const backend = new ControlledBackend()
|
|
|
|
|
const id = SessionId('prepare-reservation-wait')
|
|
|
|
|
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
|
|
|
|
|
let coordinator!: PersistenceCoordinator<never>
|
|
|
|
|
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
|
|
}, { inject: ['sessions'] }))
|
|
|
|
|
let first: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
let second: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
first = await coordinator.prepare(id)
|
|
|
|
|
let secondResolved = false
|
|
|
|
|
const waiting = coordinator.prepare(id).then((preparation) => {
|
|
|
|
|
secondResolved = true
|
|
|
|
|
return preparation
|
|
|
|
|
})
|
|
|
|
|
await Promise.resolve()
|
|
|
|
|
expect(secondResolved).toBe(false)
|
|
|
|
|
|
|
|
|
|
first[Symbol.dispose]()
|
|
|
|
|
second = await waiting
|
|
|
|
|
expect(second.session).toBe(first.session)
|
|
|
|
|
expect(backend.loadAttempts).toBe(1)
|
|
|
|
|
} finally {
|
|
|
|
|
second?.[Symbol.dispose]()
|
|
|
|
|
first?.[Symbol.dispose]()
|
|
|
|
|
await fiber.dispose()
|
|
|
|
|
await ctx.fiber.dispose()
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
it('evicts only ready preparations by LRU capacity', async () => {
|
|
|
|
|
const ctx = new Context()
|
|
|
|
|
await ctx.plugin(SessionStore)
|
|
|
|
|
const backend = new ControlledBackend()
|
|
|
|
|
const firstId = SessionId('preparation-lru-first')
|
|
|
|
|
const secondId = SessionId('preparation-lru-second')
|
|
|
|
|
backend.store.set(firstId, { meta: meta(firstId), events: oneTurnLog() })
|
|
|
|
|
backend.store.set(secondId, { meta: meta(secondId), events: oneTurnLog() })
|
|
|
|
|
let coordinator!: PersistenceCoordinator<never>
|
|
|
|
|
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
coordinator = new PersistenceCoordinator(inner, backend, { preparedSessionCacheSize: 1 })
|
|
|
|
|
}, { inject: ['sessions'] }))
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
await coordinator.inspect(firstId)
|
|
|
|
|
await coordinator.inspect(secondId)
|
|
|
|
|
await coordinator.inspect(firstId)
|
|
|
|
|
expect(backend.loadAttempts).toBe(3)
|
|
|
|
|
} finally {
|
|
|
|
|
await fiber.dispose()
|
|
|
|
|
await ctx.fiber.dispose()
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
it('rejects append while an unpublished preparation owns the persisted cursor', async () => {
|
|
|
|
|
const ctx = new Context()
|
|
|
|
|
await ctx.plugin(SessionStore)
|
|
|
|
|
const backend = new ControlledBackend()
|
|
|
|
|
const id = SessionId('reserved-append')
|
|
|
|
|
backend.store.set(id, { meta: meta(id), events: oneTurnLog() })
|
|
|
|
|
let coordinator!: PersistenceCoordinator<never>
|
|
|
|
|
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
|
|
}, { inject: ['sessions'] }))
|
|
|
|
|
let preparation: Awaited<ReturnType<typeof coordinator.prepare>> | undefined
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
preparation = await coordinator.prepare(id)
|
|
|
|
|
await expect(coordinator.append(id, [{
|
|
|
|
|
type: 'turn/start',
|
|
|
|
|
seq: oneTurnLog().length,
|
|
|
|
|
time: 7,
|
|
|
|
|
data: { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } },
|
|
|
|
|
}])).rejects.toThrow(/persisted preparation is reserved/)
|
|
|
|
|
} finally {
|
|
|
|
|
preparation?.[Symbol.dispose]()
|
|
|
|
|
await fiber.dispose()
|
|
|
|
|
await ctx.fiber.dispose()
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
describe('PersistenceCoordinator observation cancellation', () => {
|
|
|
|
|
it('promptly rejects a queued inspect without invoking it and keeps the same-id chain healthy', async () => {
|
|
|
|
|
const ctx = new Context()
|
|
|
|
|
@@ -400,7 +572,7 @@ describe('PersistenceCoordinator observation cancellation', () => {
|
|
|
|
|
await expect(prior).resolves.toMatchObject({ meta: { id } })
|
|
|
|
|
await observedAbort
|
|
|
|
|
await expect(subsequent).resolves.toMatchObject({ meta: { id } })
|
|
|
|
|
expect(backend.loadAttempts).toBe(2)
|
|
|
|
|
expect(backend.loadAttempts).toBe(1)
|
|
|
|
|
await vi.waitFor(() => {
|
|
|
|
|
expect((coordinator as unknown as CoordinatorInternals).chains.size).toBe(0)
|
|
|
|
|
})
|
|
|
|
|
@@ -540,6 +712,7 @@ describe('PersistenceCoordinator observation cancellation', () => {
|
|
|
|
|
// retirement promise stays pending in the coordinator.
|
|
|
|
|
await sessionFiber.dispose()
|
|
|
|
|
await vi.waitFor(() => { expect(internals.retirements.has(id)).toBe(true) })
|
|
|
|
|
const baselineLoads = backend.loadAttempts
|
|
|
|
|
|
|
|
|
|
const controller = new AbortController()
|
|
|
|
|
const reason = new Error('inspect cancelled during retirement')
|
|
|
|
|
@@ -552,7 +725,7 @@ describe('PersistenceCoordinator observation cancellation', () => {
|
|
|
|
|
// backend read.
|
|
|
|
|
controller.abort(reason)
|
|
|
|
|
await vi.waitFor(() => { expect(observedReason).toBe(reason) })
|
|
|
|
|
expect(backend.loadAttempts).toBe(0)
|
|
|
|
|
expect(backend.loadAttempts).toBe(baselineLoads)
|
|
|
|
|
|
|
|
|
|
appendGate.resolve(true)
|
|
|
|
|
await observed
|
|
|
|
|
@@ -621,15 +794,17 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
|
|
}, { inject: ['sessions'] }))
|
|
|
|
|
await ctx.sessions.flush(first)
|
|
|
|
|
|
|
|
|
|
// Occupy the per-id serialize chain with a gated read: everything the
|
|
|
|
|
// two retirements queue stays pending behind it. (Attempt counting
|
|
|
|
|
// starts here — an absent beforeLoadStored short-circuits the optional
|
|
|
|
|
// call without evaluating its ++ argument.)
|
|
|
|
|
backend.beforeLoadStored = async (attempt) => {
|
|
|
|
|
if (attempt === 1) await readGate.promise
|
|
|
|
|
// Occupy the per-id serialize chain with a gated physical read:
|
|
|
|
|
// inspect() correctly borrows the still-live Session without entering
|
|
|
|
|
// the backend chain, while both retirements must queue behind readFrom().
|
|
|
|
|
const readEntered = Promise.withResolvers<undefined>()
|
|
|
|
|
backend.seekHook = async () => {
|
|
|
|
|
readEntered.resolve(undefined)
|
|
|
|
|
await readGate.promise
|
|
|
|
|
return undefined
|
|
|
|
|
}
|
|
|
|
|
const parked = coordinator.inspect(id).catch((error: unknown) => error)
|
|
|
|
|
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
|
|
|
|
|
const parked = coordinator.readFrom(id, 0).catch((error: unknown) => error)
|
|
|
|
|
await readEntered.promise
|
|
|
|
|
|
|
|
|
|
// First retirement queues behind the gate and stays pending.
|
|
|
|
|
await firstFiber.dispose()
|
|
|
|
|
@@ -730,7 +905,7 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
|
|
|
|
|
|
|
await expect(ctx.plugin(Object.assign((inner: Context) => {
|
|
|
|
|
inner.sessions.create(id)
|
|
|
|
|
}, { inject: ['sessions'] }))).rejects.toThrow(/persisted history is loading/)
|
|
|
|
|
}, { inject: ['sessions'] }))).rejects.toThrow(/persisted preparation exists/)
|
|
|
|
|
|
|
|
|
|
loadGate.resolve(true)
|
|
|
|
|
await expect(coldLoad).resolves.toMatchObject({
|
|
|
|
|
|