Merge remote-tracking branch 'origin/master' into codex/sqlite-metadata-contracts

This commit is contained in:
Hypatia May
2026-07-24 09:48:06 +08:00
228 changed files with 6093 additions and 2470 deletions

View File

@@ -88,6 +88,84 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
}
})
it('rejects crash-repair load while a live session owns the persisted prefix', async () => {
const fix = await makeFixture()
const { ctx, fiber } = await freshCtx(fix)
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(SessionId('live-load'), { meta: { cwd: WORK } })
}, { inject: ['sessions'] }))
try {
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
await ctx.sessions.flush(session)
await expect(ctx.sessionPersistence.load(session.id))
.rejects.toThrow(`cannot load session "${session.id}" while its live turn is open`)
send(session, oneTurnLog().slice(1))
await ctx.sessions.flush(session)
await sessionFiber.dispose()
await vi.waitFor(async () => {
const loaded = await ctx.sessionPersistence.load(session.id)
expect(loaded.events.map(event => event.type)).toEqual(oneTurnLog().map(event => event.type))
expect(loaded.events.at(-1)).toMatchObject({
type: 'turn/end',
data: { reason: { kind: 'completed' } },
})
})
} finally {
await sessionFiber.dispose()
await fiber.dispose()
await fix.cleanup()
}
})
it('rechecks live ownership after a cold load enters the per-id chain', async () => {
const fix = await makeFixture()
const { ctx, fiber } = await freshCtx(fix)
try {
const id = SessionId('queued-load-live-race')
const header = meta(id, WORK)
const start: SessionEvent = {
type: 'turn/start',
seq: 0,
time: 1,
data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
}
await ctx.sessionPersistence.create(header)
await ctx.sessionPersistence.append(id, [start])
const loading = ctx.sessionPersistence.load(id)
const live = ctx.sessions.create(id, { seed: [start], meta: header })
await expect(loading).rejects.toThrow(/live turn is open/)
live.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
await ctx.sessions.flush(live)
const loaded = await ctx.sessionPersistence.load(id)
expect(loaded.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
expect(loaded.events.at(-1)).toMatchObject({
type: 'turn/end',
data: { reason: { kind: 'completed' } },
})
} finally {
await fiber.dispose()
await fix.cleanup()
}
})
it('does not load an unmaterialized empty live session', async () => {
const fix = await makeFixture()
const { ctx, fiber } = await freshCtx(fix)
try {
const session = ctx.sessions.create(SessionId('empty-live'), { meta: { cwd: WORK } })
await expect(ctx.sessionPersistence.load(session.id)).rejects.toThrow(/not found/)
} finally {
await fiber.dispose()
await fix.cleanup()
}
})
it('round-trips the seed boundary (seedLength) through persistence', async () => {
// A forked child records how many leading events were inherited via the seed; the
// boundary must survive a reload (so a resume/replay can tell the inherited prefix from
@@ -95,9 +173,13 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
const fix = await makeFixture()
const { ctx, fiber } = await freshCtx(fix)
try {
const session = ctx.sessions.create(SessionId('forked-child'), { meta: { cwd: WORK, seedLength: 3 } })
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(SessionId('forked-child'), { meta: { cwd: WORK, seedLength: 3 } })
}, { inject: ['sessions'] }))
send(session, oneTurnLog())
await ctx.sessions.flush(session)
await sessionFiber.dispose()
const loaded = await ctx.sessionPersistence.load(SessionId('forked-child'))
expect(loaded.meta.seedLength).toBe(3)
@@ -114,11 +196,15 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
const fix = await makeFixture()
const { ctx, fiber } = await freshCtx(fix)
try {
const session = ctx.sessions.create(SessionId('delegated-child'), {
meta: { cwd: WORK, parentSession: SessionId('root'), delegationDepth: 2 },
})
let session!: Session
const sessionFiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create(SessionId('delegated-child'), {
meta: { cwd: WORK, parentSession: SessionId('root'), delegationDepth: 2 },
})
}, { inject: ['sessions'] }))
send(session, oneTurnLog())
await ctx.parallel('session/flush', session)
await sessionFiber.dispose()
const loaded = await ctx.sessionPersistence.load(SessionId('delegated-child'))
expect(loaded.meta.delegationDepth).toBe(2)
@@ -526,20 +612,31 @@ export function runCoordinatorContract(name: string, makeFixture: () => Promise<
const { ctx, fiber } = await freshCtx(fix)
try {
// Materialize and load (ownerless, cursor = 6).
await ctx.sessionPersistence.create(meta('claim', WORK))
const storedMeta = meta('claim', WORK)
await ctx.sessionPersistence.create(storedMeta)
await ctx.sessionPersistence.append(SessionId('claim'), oneTurnLog())
const { events } = await ctx.sessionPersistence.load(SessionId('claim'))
const { events, meta: durableMeta } = 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(SessionId('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 } })
let cont!: Session
const contFiber = await ctx.plugin(Object.assign((inner: Context) => {
cont = inner.sessions.create(SessionId('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, createdAt: 2000 } })
}, { inject: ['sessions'] }))
await ctx.sessions.flush(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])
expect(loaded.meta).toEqual(durableMeta)
expect(loaded.meta.createdAt).toBe(1000)
await contFiber.dispose()
await vi.waitFor(async () => {
expect((await ctx.sessionPersistence.load(SessionId('claim'))).meta).toEqual(durableMeta)
})
} finally {
await fiber.dispose()
await fix.cleanup()

View File

@@ -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,35 +353,30 @@ 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')
let first!: Session
const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
first = inner.sessions.create(id)
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) })
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'] }))
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()
await expect(reuseFlush).resolves.toBeUndefined()
} finally {
loadGate.resolve(true)
await backendFiber.dispose()
@@ -280,7 +384,45 @@ describe('PersistenceCoordinator retirement', () => {
}
})
it('a retiring owner with buffered events still rejects same-id reuse', async () => {
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)
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'] }))
const reuseFlush = ctx.sessions.flush(reuse)
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 {
appendGate.resolve(true)
await backendFiber.dispose()
await ctx.fiber.dispose()
}
})
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()