|
|
@@ -48,10 +48,8 @@ interface MemoryConfig { store?: MemoryStore }
|
|
|
/** Test-only view of the coordinator containers whose retirement is the contract under test. */
|
|
|
interface CoordinatorInternals {
|
|
|
states: Map<unknown, unknown>
|
|
|
- buffers: Map<unknown, unknown>
|
|
|
+ live: Map<unknown, { pending: unknown[]; flush: Promise<void> | undefined }>
|
|
|
chains: Map<unknown, unknown>
|
|
|
- inits: Map<unknown, unknown>
|
|
|
- retirements: Set<Promise<void>>
|
|
|
}
|
|
|
|
|
|
/**
|
|
|
@@ -209,6 +207,75 @@ runCoordinatorContract('memory', async (): Promise<CoordinatorFixture> => {
|
|
|
}
|
|
|
})
|
|
|
|
|
|
+describe('PersistenceCoordinator eager writes', () => {
|
|
|
+ it('starts a follow-up batch for events admitted during an in-flight write', async () => {
|
|
|
+ const ctx = new Context()
|
|
|
+ await ctx.plugin(SessionStore)
|
|
|
+ const backend = new ControlledBackend()
|
|
|
+ const appendGate = Promise.withResolvers<boolean>()
|
|
|
+ backend.beforeAppend = async (attempt) => {
|
|
|
+ if (attempt === 1) await appendGate.promise
|
|
|
+ }
|
|
|
+ const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ new PersistenceCoordinator(inner, backend)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+
|
|
|
+ try {
|
|
|
+ const session = ctx.sessions.create(SessionId('eager-follow-up'))
|
|
|
+ await ctx.sessions.flush(session)
|
|
|
+ session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
|
+ await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
|
|
|
+
|
|
|
+ session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
+ appendGate.resolve(true)
|
|
|
+
|
|
|
+ await vi.waitFor(() => {
|
|
|
+ expect(backend.appendAttempts).toBe(2)
|
|
|
+ expect(backend.store.get(session.id)?.events.map(event => event.seq)).toEqual([0, 1])
|
|
|
+ })
|
|
|
+ } finally {
|
|
|
+ appendGate.resolve(true)
|
|
|
+ await fiber.dispose()
|
|
|
+ await ctx.fiber.dispose()
|
|
|
+ }
|
|
|
+ })
|
|
|
+
|
|
|
+ it('retries a failed overlapping eager write at the explicit flush barrier', async () => {
|
|
|
+ const ctx = new Context()
|
|
|
+ await ctx.plugin(SessionStore)
|
|
|
+ const backend = new ControlledBackend()
|
|
|
+ const appendGate = Promise.withResolvers<boolean>()
|
|
|
+ backend.beforeAppend = async (attempt) => {
|
|
|
+ if (attempt === 1) {
|
|
|
+ await appendGate.promise
|
|
|
+ throw new Error('transient eager failure')
|
|
|
+ }
|
|
|
+ }
|
|
|
+ const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ new PersistenceCoordinator(inner, backend)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+
|
|
|
+ try {
|
|
|
+ const session = ctx.sessions.create(SessionId('eager-flush-retry'))
|
|
|
+ await ctx.sessions.flush(session)
|
|
|
+ session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
|
+ session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
+ await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
|
|
|
+
|
|
|
+ const barriers = [ctx.sessions.flush(session), ctx.sessions.flush(session)]
|
|
|
+ appendGate.resolve(true)
|
|
|
+
|
|
|
+ await expect(Promise.all(barriers)).resolves.toEqual([undefined, undefined])
|
|
|
+ expect(backend.appendAttempts).toBe(2)
|
|
|
+ expect(backend.store.get(session.id)?.events.map(event => event.seq)).toEqual([0, 1])
|
|
|
+ } finally {
|
|
|
+ appendGate.resolve(true)
|
|
|
+ await fiber.dispose()
|
|
|
+ await ctx.fiber.dispose()
|
|
|
+ }
|
|
|
+ })
|
|
|
+})
|
|
|
+
|
|
|
describe('PersistenceCoordinator stored identity', () => {
|
|
|
it('rejects a mismatched backend header before repair or state publication', async () => {
|
|
|
const ctx = new Context()
|
|
|
@@ -237,6 +304,48 @@ describe('PersistenceCoordinator stored identity', () => {
|
|
|
await ctx.fiber.dispose()
|
|
|
}
|
|
|
})
|
|
|
+
|
|
|
+ it('reserves a cold id across asynchronous storage repair', async () => {
|
|
|
+ const ctx = new Context()
|
|
|
+ await ctx.plugin(SessionStore)
|
|
|
+ const backend = new ControlledBackend()
|
|
|
+ const id = SessionId('cold-load-reservation')
|
|
|
+ const header = meta(id)
|
|
|
+ const start: SessionEvent = {
|
|
|
+ type: 'turn/start',
|
|
|
+ seq: 0,
|
|
|
+ time: 1,
|
|
|
+ data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
|
|
+ }
|
|
|
+ backend.store.set(id, { meta: header, events: [start] })
|
|
|
+ const loadGate = Promise.withResolvers<boolean>()
|
|
|
+ backend.beforeLoadStored = async () => { await loadGate.promise }
|
|
|
+ let coordinator!: PersistenceCoordinator<never>
|
|
|
+ const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+
|
|
|
+ try {
|
|
|
+ const loading = coordinator.load(id)
|
|
|
+ await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
|
|
|
+
|
|
|
+ await expect(ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ inner.sessions.create(id, { seed: [start], meta: header })
|
|
|
+ }, { inject: ['sessions'] }))).rejects.toThrow(/persisted history is loading/)
|
|
|
+ expect(ctx.sessions.get(id)).toBeUndefined()
|
|
|
+
|
|
|
+ loadGate.resolve(true)
|
|
|
+ const loaded = await loading
|
|
|
+ expect(loaded.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
|
|
|
+
|
|
|
+ const resumed = ctx.sessions.create(id, { seed: loaded.events, meta: loaded.meta })
|
|
|
+ await expect(ctx.sessions.flush(resumed)).resolves.toBeUndefined()
|
|
|
+ } finally {
|
|
|
+ loadGate.resolve(true)
|
|
|
+ await fiber.dispose()
|
|
|
+ await ctx.fiber.dispose()
|
|
|
+ }
|
|
|
+ })
|
|
|
})
|
|
|
|
|
|
describe('PersistenceCoordinator retirement', () => {
|
|
|
@@ -244,43 +353,76 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
const ctx = new Context()
|
|
|
await ctx.plugin(SessionStore)
|
|
|
const backend = new ControlledBackend()
|
|
|
- let coordinator!: PersistenceCoordinator<never>
|
|
|
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
- coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
+ new PersistenceCoordinator(inner, backend)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
const loadGate = Promise.withResolvers<boolean>()
|
|
|
+ backend.beforeLoadStored = async (attempt) => {
|
|
|
+ if (attempt === 1) await loadGate.promise
|
|
|
+ }
|
|
|
|
|
|
try {
|
|
|
const id = SessionId('retiring-lazy-owner')
|
|
|
+ const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ inner.sessions.create(id)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+ await vi.waitFor(() => { expect(backend.loadAttempts).toBe(1) })
|
|
|
+ await firstFiber.dispose()
|
|
|
+
|
|
|
+ let reuse!: Session
|
|
|
+ await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ reuse = inner.sessions.create(id)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+ const reuseFlush = ctx.sessions.flush(reuse)
|
|
|
+
|
|
|
+ loadGate.resolve(true)
|
|
|
+ await expect(reuseFlush).resolves.toBeUndefined()
|
|
|
+ } finally {
|
|
|
+ loadGate.resolve(true)
|
|
|
+ await backendFiber.dispose()
|
|
|
+ await ctx.fiber.dispose()
|
|
|
+ }
|
|
|
+ })
|
|
|
+
|
|
|
+ it('a replacement queued before retirement cleanup still collides with the live owner', async () => {
|
|
|
+ const ctx = new Context()
|
|
|
+ await ctx.plugin(SessionStore)
|
|
|
+ const backend = new ControlledBackend()
|
|
|
+ const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ new PersistenceCoordinator(inner, backend)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+ const appendGate = Promise.withResolvers<boolean>()
|
|
|
+
|
|
|
+ try {
|
|
|
+ const id = SessionId('retiring-live-owner')
|
|
|
let first!: Session
|
|
|
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
first = inner.sessions.create(id)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
await ctx.sessions.flush(first)
|
|
|
-
|
|
|
- const baselineLoads = backend.loadAttempts
|
|
|
- backend.beforeLoadStored = async () => { await loadGate.promise }
|
|
|
- const blockingLoad = coordinator.load(id)
|
|
|
- await vi.waitFor(() => { expect(backend.loadAttempts).toBe(baselineLoads + 1) })
|
|
|
+ backend.beforeAppend = async () => { await appendGate.promise }
|
|
|
+ first.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
|
+ first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
+ await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
|
|
|
await firstFiber.dispose()
|
|
|
|
|
|
let reuse!: Session
|
|
|
await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
reuse = inner.sessions.create(id)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
- await vi.waitFor(() => { expect(backend.loadAttempts).toBe(baselineLoads + 2) })
|
|
|
+ const reuseFlush = ctx.sessions.flush(reuse)
|
|
|
|
|
|
- loadGate.resolve(true)
|
|
|
- await expect(blockingLoad).rejects.toThrow(/not found/)
|
|
|
- await expect(ctx.sessions.flush(reuse)).resolves.toBeUndefined()
|
|
|
+ appendGate.resolve(true)
|
|
|
+ await expect(reuseFlush).rejects.toThrow(/bound to a different live session/)
|
|
|
+ expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
|
|
|
} finally {
|
|
|
- loadGate.resolve(true)
|
|
|
+ appendGate.resolve(true)
|
|
|
await backendFiber.dispose()
|
|
|
await ctx.fiber.dispose()
|
|
|
}
|
|
|
})
|
|
|
|
|
|
- it('a retiring owner with buffered events still rejects same-id reuse', async () => {
|
|
|
+ it('a racing cold load survives retirement cleanup and rejects same-id reuse', async () => {
|
|
|
const ctx = new Context()
|
|
|
await ctx.plugin(SessionStore)
|
|
|
const backend = new ControlledBackend()
|
|
|
@@ -288,6 +430,7 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
const backendFiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
+ const appendGate = Promise.withResolvers<boolean>()
|
|
|
const loadGate = Promise.withResolvers<boolean>()
|
|
|
|
|
|
try {
|
|
|
@@ -297,27 +440,37 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
first = inner.sessions.create(id)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
await ctx.sessions.flush(first)
|
|
|
+ backend.beforeAppend = async () => { await appendGate.promise }
|
|
|
first.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
|
first.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
-
|
|
|
+ await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
|
|
|
+ await firstFiber.dispose()
|
|
|
const baselineLoads = backend.loadAttempts
|
|
|
backend.beforeLoadStored = async () => { await loadGate.promise }
|
|
|
- const blockingLoad = coordinator.load(id)
|
|
|
+ const coldLoad = coordinator.load(id)
|
|
|
+
|
|
|
+ appendGate.resolve(true)
|
|
|
await vi.waitFor(() => { expect(backend.loadAttempts).toBe(baselineLoads + 1) })
|
|
|
- await firstFiber.dispose()
|
|
|
+
|
|
|
+ await expect(ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ inner.sessions.create(id)
|
|
|
+ }, { inject: ['sessions'] }))).rejects.toThrow(/persisted history is loading/)
|
|
|
+
|
|
|
+ loadGate.resolve(true)
|
|
|
+ await expect(coldLoad).resolves.toMatchObject({
|
|
|
+ events: [{ seq: 0 }, { seq: 1 }],
|
|
|
+ })
|
|
|
|
|
|
let reuse!: Session
|
|
|
await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
reuse = inner.sessions.create(id)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
- await expect(ctx.sessions.flush(reuse)).rejects.toThrow(/bound to a different live session/)
|
|
|
-
|
|
|
- loadGate.resolve(true)
|
|
|
- await expect(blockingLoad).rejects.toThrow(/not found/)
|
|
|
+ await expect(ctx.sessions.flush(reuse)).rejects.toThrow(/id collision/)
|
|
|
await vi.waitFor(() => {
|
|
|
expect(backend.store.get(id)?.events.map(event => event.seq)).toEqual([0, 1])
|
|
|
})
|
|
|
} finally {
|
|
|
+ appendGate.resolve(true)
|
|
|
loadGate.resolve(true)
|
|
|
await backendFiber.dispose()
|
|
|
await ctx.fiber.dispose()
|
|
|
@@ -381,8 +534,9 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
}, { inject: ['sessions'] }))
|
|
|
const internals = coordinator as unknown as CoordinatorInternals
|
|
|
- backend.beforeAppend = async (attempt) => {
|
|
|
- if (attempt === 1) {
|
|
|
+ let retryEnabled = false
|
|
|
+ backend.beforeAppend = async () => {
|
|
|
+ if (!retryEnabled) {
|
|
|
backend.lifecycle.push('append-failed')
|
|
|
throw new Error('transient append failure')
|
|
|
}
|
|
|
@@ -399,17 +553,18 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
await sessionFiber.dispose()
|
|
|
|
|
|
await vi.waitFor(() => {
|
|
|
- expect(backend.appendAttempts).toBe(1)
|
|
|
- expect(internals.retirements.size).toBe(0)
|
|
|
+ expect(backend.appendAttempts).toBeGreaterThanOrEqual(1)
|
|
|
+ expect([...internals.live.values()][0]?.pending).toEqual(expect.arrayContaining([
|
|
|
+ expect.objectContaining({ seq: 0 }),
|
|
|
+ expect.objectContaining({ seq: 1 }),
|
|
|
+ ]))
|
|
|
})
|
|
|
- expect([...internals.buffers.values()]).toEqual([expect.arrayContaining([
|
|
|
- expect.objectContaining({ seq: 0 }),
|
|
|
- expect.objectContaining({ seq: 1 }),
|
|
|
- ])])
|
|
|
|
|
|
+ retryEnabled = true
|
|
|
await backendFiber.dispose()
|
|
|
expect(backend.store.get(SessionId('retry-retirement'))?.events.map(event => event.seq)).toEqual([0, 1])
|
|
|
- expect(backend.lifecycle).toEqual(['append-failed', 'append-committed', 'close'])
|
|
|
+ expect(backend.lifecycle.at(-2)).toBe('append-committed')
|
|
|
+ expect(backend.lifecycle.at(-1)).toBe('close')
|
|
|
} finally {
|
|
|
await backendFiber.dispose()
|
|
|
await ctx.fiber.dispose()
|
|
|
@@ -442,7 +597,8 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
await sessionFiber.dispose()
|
|
|
await vi.waitFor(() => {
|
|
|
expect(backend.appendAttempts).toBe(1)
|
|
|
- expect(internals.retirements.size).toBe(1)
|
|
|
+ expect(internals.live.size).toBe(1)
|
|
|
+ expect([...internals.live.values()][0]?.flush).toBeInstanceOf(Promise)
|
|
|
})
|
|
|
|
|
|
let disposed = false
|
|
|
@@ -461,6 +617,47 @@ describe('PersistenceCoordinator retirement', () => {
|
|
|
await ctx.fiber.dispose()
|
|
|
}
|
|
|
})
|
|
|
+
|
|
|
+ it('backend teardown waits for a detached public append before close', async () => {
|
|
|
+ const ctx = new Context()
|
|
|
+ await ctx.plugin(SessionStore)
|
|
|
+ const backend = new ControlledBackend()
|
|
|
+ let coordinator!: PersistenceCoordinator<never>
|
|
|
+ const fiber = await ctx.plugin(Object.assign((inner: Context) => {
|
|
|
+ coordinator = new PersistenceCoordinator(inner, backend)
|
|
|
+ }, { inject: ['sessions'] }))
|
|
|
+ const appendGate = Promise.withResolvers<boolean>()
|
|
|
+ backend.beforeAppend = async () => {
|
|
|
+ backend.lifecycle.push('append-started')
|
|
|
+ await appendGate.promise
|
|
|
+ backend.lifecycle.push('append-committed')
|
|
|
+ }
|
|
|
+
|
|
|
+ try {
|
|
|
+ const id = SessionId('inflight-public-append')
|
|
|
+ await coordinator.create(meta(id))
|
|
|
+ const append = coordinator.append(id, [{
|
|
|
+ type: 'turn/start',
|
|
|
+ seq: 0,
|
|
|
+ time: 1,
|
|
|
+ data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
|
|
+ }])
|
|
|
+ await vi.waitFor(() => { expect(backend.appendAttempts).toBe(1) })
|
|
|
+
|
|
|
+ let disposed = false
|
|
|
+ const teardown = fiber.dispose().then(() => { disposed = true })
|
|
|
+ await Promise.resolve()
|
|
|
+ expect(disposed).toBe(false)
|
|
|
+
|
|
|
+ appendGate.resolve(true)
|
|
|
+ await Promise.all([append, teardown])
|
|
|
+ expect(backend.lifecycle).toEqual(['append-started', 'append-committed', 'close'])
|
|
|
+ } finally {
|
|
|
+ appendGate.resolve(true)
|
|
|
+ await fiber.dispose()
|
|
|
+ await ctx.fiber.dispose()
|
|
|
+ }
|
|
|
+ })
|
|
|
})
|
|
|
|
|
|
describe('SessionPersistence service registration', () => {
|
|
|
@@ -590,11 +787,9 @@ describe('SessionPersistence service registration', () => {
|
|
|
expect(ctx.sessions.list()).toHaveLength(0)
|
|
|
expect({
|
|
|
states: coordinator.states.size,
|
|
|
- buffers: coordinator.buffers.size,
|
|
|
+ live: coordinator.live.size,
|
|
|
chains: coordinator.chains.size,
|
|
|
- inits: coordinator.inits.size,
|
|
|
- retirements: coordinator.retirements.size,
|
|
|
- }).toEqual({ states: 0, buffers: 0, chains: 0, inits: 0, retirements: 0 })
|
|
|
+ }).toEqual({ states: 0, live: 0, chains: 0 })
|
|
|
})
|
|
|
} finally {
|
|
|
await fiber.dispose()
|