feat(tui): add safe session resume flow
This commit is contained in:
@@ -6,6 +6,9 @@ The JSONL durable session-persistence backend — a concrete `SessionPersistence
|
||||
|
||||
```
|
||||
<root>/
|
||||
.live/
|
||||
<encoded-id>.lock # PID + nonce cross-process live lease
|
||||
<encoded-id>.lock.reclaim # ephemeral stale-owner takeover guard
|
||||
cwd-<sha256(cwd)[:12]>/ # per-project bucket (or _no-cwd/ when no cwd)
|
||||
<encoded-id>.jsonl.zstd # default: checksummed header frame + append frames
|
||||
<encoded-id>.jsonl # only with compression: 'none'
|
||||
@@ -43,7 +46,7 @@ A root belongs to one encoding. Startup discovery and targeted lookup reject the
|
||||
|
||||
## Write path
|
||||
|
||||
The plugin copies frozen session events into one controller per live session and starts an eager drain. Concurrent events share the current write; events admitted during it form a follow-up batch, while `session/flush` waits until both current and pending batches are durable. A per-session cursor prevents resumed sessions from re-appending stored events, and live sessions are seeded when the plugin loads. The owning backend instance serializes operations for one session; disposal drains every retained controller before teardown.
|
||||
The plugin copies frozen session events into one controller per live session and starts an eager drain. Before a session can flush or resume, the coordinator claims an exclusive `.live/<encoded-id>.lock` containing the process PID and an exec-stable nonce; another live process is rejected, while a dead owner is reclaimed under the separate `.reclaim` guard. Concurrent events share the current write; events admitted during it form a follow-up batch, while `session/flush` waits until both current and pending batches are durable. A per-session cursor prevents resumed sessions from re-appending stored events, and live sessions are seeded when the plugin loads. Disposal drains every retained controller before releasing its lease.
|
||||
|
||||
## Model Experience
|
||||
|
||||
@@ -66,5 +69,6 @@ JSONL storage does not mutate live request prefixes. A resumed loop can reuse pr
|
||||
- **Only the configured encoding and current `SESSION_FORMAT_VERSION` (v0) load** — changing compression requires a separate/fresh root or selecting the legacy raw mode; the pre-release format has no migration.
|
||||
- **Compressed files are not directly line-readable** — use the backend to load them, or select `compression: 'none'` before writing a fresh root when text fixtures or external line readers are required.
|
||||
- **Nothing deletes session files** — logs accumulate under `root` until removed externally (the seam has no deletion surface).
|
||||
- **One live writer per session** — append and repair are coordinated only inside the owning backend instance. Another backend instance or process must not write the same session until that owner reaches quiescent disposal; initial same-id publication remains collision-safe through the POSIX no-overwrite hard link or Windows write-through rename without replacement.
|
||||
- **Lease scope is local-host advisory ownership** — PID liveness prevents two ordinary local Harness processes from resuming the same id, but it is not a distributed lease for shared network filesystems or hostile principals.
|
||||
- **A crash during stale-lease takeover fails closed** — if the reclaiming process itself crashes while holding the short-lived `.reclaim` guard, an operator must remove that guard after confirming no recovery is active.
|
||||
- **POSIX materialization requires hard-link support** — first append uses `link()` so same-id races fail instead of overwriting a committed log; Windows uses write-through rename without replacement.
|
||||
|
||||
@@ -14,8 +14,9 @@ import { dirname, join, resolve } from 'node:path'
|
||||
import { randomBytes } from 'node:crypto'
|
||||
import {
|
||||
SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
|
||||
type PersistenceBackend, type SessionLocation, type SessionPersistenceSnapshot,
|
||||
type StoredPrefix,
|
||||
sessionLeaseProcessIsLive, shareSessionLiveLease,
|
||||
type PersistenceBackend, type SessionLiveLease, type SessionLiveOwner,
|
||||
type SessionLocation, type SessionPersistenceSnapshot, type StoredPrefix,
|
||||
} from '@deepseek-ai/dsh-session-persistence'
|
||||
import type { SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import {
|
||||
@@ -64,6 +65,11 @@ interface JsonlTornMarker {
|
||||
recoveredEvents: SessionEvent[]
|
||||
}
|
||||
|
||||
interface JsonlLiveLeaseRecord {
|
||||
pid: number
|
||||
nonce: string
|
||||
}
|
||||
|
||||
/** Whether a filesystem error means absence; every non-ENOENT failure must surface. */
|
||||
function isENOENT(error: unknown): boolean {
|
||||
return (error as NodeJS.ErrnoException | null)?.code === 'ENOENT'
|
||||
@@ -135,6 +141,14 @@ export class SessionPersistenceJsonl extends SessionPersistence implements Persi
|
||||
return this.coordinator.inspect(id)
|
||||
}
|
||||
|
||||
override claimLive(id: SessionId): Promise<SessionLiveLease> {
|
||||
return this.coordinator.claimLive(id)
|
||||
}
|
||||
|
||||
override isLive(id: SessionId): Promise<boolean> {
|
||||
return this.coordinator.isLive(id)
|
||||
}
|
||||
|
||||
// One method serves both public `list` and the backend hook; delegating it to
|
||||
// the coordinator would call this hook recursively.
|
||||
|
||||
@@ -274,6 +288,110 @@ export class SessionPersistenceJsonl extends SessionPersistence implements Persi
|
||||
return snapshots
|
||||
}
|
||||
|
||||
/** Atomically publish one process lease, reclaiming a crashed owner's record. */
|
||||
async acquireLive(id: SessionId, owner: SessionLiveOwner): Promise<() => Promise<void>> {
|
||||
const path = this.liveLeasePath(id)
|
||||
return shareSessionLiveLease(`jsonl:${path}`, () => this.acquireLiveFile(path, id, owner))
|
||||
}
|
||||
|
||||
private async acquireLiveFile(
|
||||
path: string,
|
||||
id: SessionId,
|
||||
owner: SessionLiveOwner,
|
||||
): Promise<() => Promise<void>> {
|
||||
await mkdir(dirname(path), { recursive: true, mode: 0o700 })
|
||||
for (;;) {
|
||||
try {
|
||||
const handle = await open(path, 'wx', 0o600)
|
||||
try {
|
||||
await handle.writeFile(`${JSON.stringify(owner)}\n`, 'utf8')
|
||||
await handle.sync()
|
||||
} finally {
|
||||
await handle.close()
|
||||
}
|
||||
break
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error
|
||||
const current = await this.readLiveLease(path)
|
||||
if (current !== undefined && current.pid === owner.pid && current.nonce === owner.nonce) break
|
||||
if (current === undefined || sessionLeaseProcessIsLive(current.pid)) {
|
||||
throw new Error(`session "${id}" is occupied by another live process`)
|
||||
}
|
||||
const reclaimPath = `${path}.reclaim`
|
||||
let reclaim: Awaited<ReturnType<typeof open>>
|
||||
try {
|
||||
reclaim = await open(reclaimPath, 'wx', 0o600)
|
||||
} catch (reclaimError) {
|
||||
/* v8 ignore else -- non-contention filesystem failures are propagated verbatim and are not portable to induce */
|
||||
if ((reclaimError as NodeJS.ErrnoException).code === 'EEXIST') {
|
||||
throw new Error(`session "${id}" live-lease reclamation is already in progress`)
|
||||
}
|
||||
/* v8 ignore next -- non-contention filesystem failures are propagated verbatim and are not portable to induce */
|
||||
throw reclaimError
|
||||
}
|
||||
try {
|
||||
/* v8 ignore start -- cross-process revalidation is covered by the two-process race test */
|
||||
const latest = await this.readLiveLease(path)
|
||||
if (latest === undefined) {
|
||||
if (await this.exists(path)) throw new Error(`session "${id}" has an unreadable live-process lease`)
|
||||
} else if (latest.pid !== owner.pid || latest.nonce !== owner.nonce) {
|
||||
if (sessionLeaseProcessIsLive(latest.pid)) {
|
||||
throw new Error(`session "${id}" is occupied by another live process`)
|
||||
}
|
||||
await rm(path, { force: true })
|
||||
}
|
||||
/* v8 ignore stop */
|
||||
} finally {
|
||||
try {
|
||||
await reclaim.close()
|
||||
} finally {
|
||||
await rm(reclaimPath, { force: true })
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return async () => {
|
||||
const current = await this.readLiveLease(path)
|
||||
if (current?.pid === owner.pid && current.nonce === owner.nonce) await rm(path, { force: true })
|
||||
}
|
||||
}
|
||||
|
||||
/** Report one non-stale process lease and clean up a crashed owner's record. */
|
||||
async inspectLive(id: SessionId, owner: SessionLiveOwner): Promise<boolean> {
|
||||
const path = this.liveLeasePath(id)
|
||||
const current = await this.readLiveLease(path)
|
||||
if (current === undefined) return await this.exists(path)
|
||||
if (current.pid === owner.pid && current.nonce === owner.nonce) return true
|
||||
if (sessionLeaseProcessIsLive(current.pid)) return true
|
||||
return false
|
||||
}
|
||||
|
||||
private liveLeasePath(id: SessionId): string {
|
||||
return join(this.root, '.live', `${encodeSegment(id)}.lock`)
|
||||
}
|
||||
|
||||
private async readLiveLease(path: string): Promise<JsonlLiveLeaseRecord | undefined> {
|
||||
let text: string
|
||||
try {
|
||||
text = await readFile(path, 'utf8')
|
||||
} catch (error) {
|
||||
if (isENOENT(error)) return undefined
|
||||
throw error
|
||||
}
|
||||
let value: unknown
|
||||
try {
|
||||
value = JSON.parse(text)
|
||||
} catch {
|
||||
return undefined
|
||||
}
|
||||
if (typeof value !== 'object' || value === null
|
||||
|| !Number.isSafeInteger((value as { pid?: unknown }).pid)
|
||||
|| (value as { pid: number }).pid <= 0
|
||||
|| typeof (value as { nonce?: unknown }).nonce !== 'string'
|
||||
|| (value as { nonce: string }).nonce.length === 0) return undefined
|
||||
return value as JsonlLiveLeaseRecord
|
||||
}
|
||||
|
||||
private async listArtifacts(): Promise<Array<{ header: SessionHeader; path: string }>> {
|
||||
await this.ensureRootEncoding()
|
||||
const artifacts: Array<{ header: SessionHeader; path: string }> = []
|
||||
|
||||
16
packages/session-persistence/session-persistence-jsonl/tests/fixtures/live-lease-child.ts
vendored
Normal file
16
packages/session-persistence/session-persistence-jsonl/tests/fixtures/live-lease-child.ts
vendored
Normal file
@@ -0,0 +1,16 @@
|
||||
/** Child process that holds one JSONL live-session lease until it is killed. */
|
||||
|
||||
import { writeFile } from 'node:fs/promises'
|
||||
import { Context } from 'cordis'
|
||||
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import SessionPersistenceJsonl from '@deepseek-ai/dsh-session-persistence-jsonl'
|
||||
|
||||
const [root, marker] = process.argv.slice(2)
|
||||
if (root === undefined || marker === undefined) throw new Error('usage: live-lease-child.ts <root> <marker>')
|
||||
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
|
||||
await ctx.sessionPersistence.claimLive(SessionId('leased-session'))
|
||||
await writeFile(marker, 'held')
|
||||
await new Promise<never>(() => { setInterval(() => {}, 60_000) })
|
||||
33
packages/session-persistence/session-persistence-jsonl/tests/fixtures/live-lease-race-child.ts
vendored
Normal file
33
packages/session-persistence/session-persistence-jsonl/tests/fixtures/live-lease-race-child.ts
vendored
Normal file
@@ -0,0 +1,33 @@
|
||||
/** Child process competing to reclaim one stale JSONL live-session lease. */
|
||||
|
||||
import { access, writeFile } from 'node:fs/promises'
|
||||
import { Context } from 'cordis'
|
||||
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import SessionPersistenceJsonl from '@deepseek-ai/dsh-session-persistence-jsonl'
|
||||
|
||||
const [root, gate, marker, rawId] = process.argv.slice(2)
|
||||
if (root === undefined || gate === undefined || marker === undefined || rawId === undefined) {
|
||||
throw new Error('usage: live-lease-race-child.ts <root> <gate> <marker> <session-id>')
|
||||
}
|
||||
|
||||
for (;;) {
|
||||
try {
|
||||
await access(gate)
|
||||
break
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error
|
||||
await new Promise(resolve => setTimeout(resolve, 5))
|
||||
}
|
||||
}
|
||||
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
|
||||
try {
|
||||
await ctx.sessionPersistence.claimLive(SessionId(rawId))
|
||||
await writeFile(marker, 'claimed')
|
||||
await new Promise<never>(() => { setInterval(() => {}, 60_000) })
|
||||
} catch (error) {
|
||||
await writeFile(marker, `rejected:${error instanceof Error ? error.message : String(error)}`)
|
||||
await ctx.fiber.dispose()
|
||||
}
|
||||
@@ -1,17 +1,24 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { spawn } from 'node:child_process'
|
||||
import { Context } from 'cordis'
|
||||
import { appendFile, mkdtemp, mkdir, rm, readFile, writeFile, readdir, stat } from 'node:fs/promises'
|
||||
import { access, appendFile, mkdtemp, mkdir, rm, readFile, writeFile, readdir, stat } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { isAbsolute, join, relative, resolve } from 'node:path'
|
||||
import { fileURLToPath } from 'node:url'
|
||||
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import SessionPersistenceJsonl from '@deepseek-ai/dsh-session-persistence-jsonl'
|
||||
import { sessionLiveOwner } from '@deepseek-ai/dsh-session-persistence'
|
||||
import { encodeSegment, eventLines, logPath, scanLog, sessionDir, toHeaderLine } from '../src/format.ts'
|
||||
import { runPersistenceContract, meta, oneTurnLog, appendLog } from '../../session-persistence/tests/contract.ts'
|
||||
import { runCoordinatorContract, type CoordinatorFixture } from '../../session-persistence/tests/coordinator-contract.ts'
|
||||
|
||||
let root: string
|
||||
const dirs: string[] = []
|
||||
const repoRoot = fileURLToPath(new URL('../../../../', import.meta.url))
|
||||
const leaseChild = fileURLToPath(new URL('./fixtures/live-lease-child.ts', import.meta.url))
|
||||
const leaseRaceChild = fileURLToPath(new URL('./fixtures/live-lease-race-child.ts', import.meta.url))
|
||||
const tsxLoader = fileURLToPath(import.meta.resolve('tsx'))
|
||||
|
||||
type MutableSessionHeader = { -readonly [K in keyof SessionHeader]: SessionHeader[K] }
|
||||
|
||||
@@ -142,6 +149,179 @@ describe('SessionPersistenceJsonl: format helpers', () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe('SessionPersistenceJsonl: cross-process live leases', () => {
|
||||
it('reference-counts one physical lease across backend instances in the process', async () => {
|
||||
const dir = await freshRoot()
|
||||
const contexts = [new Context(), new Context()]
|
||||
for (const ctx of contexts) {
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(SessionPersistenceJsonl, { root: dir, compression: 'none' })
|
||||
}
|
||||
try {
|
||||
const first = await contexts[0]!.sessionPersistence.claimLive(SessionId('shared-live'))
|
||||
const second = await contexts[1]!.sessionPersistence.claimLive(SessionId('shared-live'))
|
||||
await first.release()
|
||||
await expect(contexts[1]!.sessionPersistence.isLive(SessionId('shared-live'))).resolves.toBe(true)
|
||||
await second.release()
|
||||
await expect(contexts[1]!.sessionPersistence.isLive(SessionId('shared-live'))).resolves.toBe(false)
|
||||
} finally {
|
||||
await Promise.all(contexts.map(ctx => ctx.fiber.dispose()))
|
||||
}
|
||||
})
|
||||
|
||||
it('disables another live owner and reclaims its lease after the process exits', async () => {
|
||||
const dir = await freshRoot()
|
||||
const marker = join(dir, 'lease-held')
|
||||
const child = spawn(process.execPath, ['--import', tsxLoader, leaseChild, dir, marker], {
|
||||
cwd: repoRoot,
|
||||
env: { ...process.env, TSX_TSCONFIG_PATH: join(repoRoot, 'tsconfig.json') },
|
||||
stdio: ['ignore', 'ignore', 'pipe'],
|
||||
})
|
||||
let stderr = ''
|
||||
child.stderr.setEncoding('utf8')
|
||||
child.stderr.on('data', (chunk: string) => { stderr += chunk })
|
||||
try {
|
||||
await vi.waitFor(() => access(marker), { timeout: 30_000 })
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(SessionPersistenceJsonl, { root: dir, compression: 'none' })
|
||||
try {
|
||||
await expect(ctx.sessionPersistence.isLive(SessionId('leased-session'))).resolves.toBe(true)
|
||||
await expect(ctx.sessionPersistence.claimLive(SessionId('leased-session')))
|
||||
.rejects.toThrow('occupied by another live process')
|
||||
const closed = new Promise<void>(resolve => child.once('close', () => { resolve() }))
|
||||
child.kill()
|
||||
await closed
|
||||
await expect(ctx.sessionPersistence.isLive(SessionId('leased-session'))).resolves.toBe(false)
|
||||
const leasePath = join(dir, '.live', `${encodeSegment('leased-session')}.lock`)
|
||||
await writeFile(leasePath, `${JSON.stringify({ pid: child.pid, nonce: 'dead-owner' })}\n`)
|
||||
const claim = await ctx.sessionPersistence.claimLive(SessionId('leased-session'))
|
||||
await claim.release()
|
||||
} finally {
|
||||
await ctx.fiber.dispose()
|
||||
}
|
||||
} catch (error) {
|
||||
throw new Error(`live-lease child failed: ${stderr}`, { cause: error })
|
||||
} finally {
|
||||
if (child.exitCode === null && child.signalCode === null) child.kill()
|
||||
}
|
||||
}, 40_000)
|
||||
|
||||
it('fails closed on malformed lease records and surfaces lease read errors', async () => {
|
||||
const dir = await freshRoot()
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(SessionPersistenceJsonl, { root: dir, compression: 'none' })
|
||||
const liveDir = join(dir, '.live')
|
||||
await mkdir(liveDir, { recursive: true })
|
||||
try {
|
||||
const malformed = [
|
||||
'not json',
|
||||
JSON.stringify(null),
|
||||
JSON.stringify({ pid: 1.5, nonce: 'x' }),
|
||||
JSON.stringify({ pid: 0, nonce: 'x' }),
|
||||
JSON.stringify({ pid: process.pid, nonce: 1 }),
|
||||
JSON.stringify({ pid: process.pid, nonce: '' }),
|
||||
]
|
||||
for (const [index, content] of malformed.entries()) {
|
||||
const id = SessionId(`malformed-${index}`)
|
||||
const path = join(liveDir, `${encodeSegment(id)}.lock`)
|
||||
await writeFile(path, content)
|
||||
await expect(ctx.sessionPersistence.isLive(id)).resolves.toBe(true)
|
||||
await expect(ctx.sessionPersistence.claimLive(id)).rejects.toThrow('occupied by another live process')
|
||||
}
|
||||
|
||||
const unreadable = SessionId('unreadable-lease')
|
||||
await mkdir(join(liveDir, `${encodeSegment(unreadable)}.lock`))
|
||||
await expect(ctx.sessionPersistence.isLive(unreadable)).rejects.toThrow()
|
||||
|
||||
const replaced = SessionId('replaced-release')
|
||||
const claim = await ctx.sessionPersistence.claimLive(replaced)
|
||||
const replacedPath = join(liveDir, `${encodeSegment(replaced)}.lock`)
|
||||
await writeFile(replacedPath, JSON.stringify({ pid: process.pid, nonce: 'replacement' }))
|
||||
await claim.release()
|
||||
expect(await readFile(replacedPath, 'utf8')).toContain('replacement')
|
||||
|
||||
const inherited = SessionId('inherited-owner')
|
||||
const inheritedPath = join(liveDir, `${encodeSegment(inherited)}.lock`)
|
||||
await writeFile(inheritedPath, JSON.stringify(sessionLiveOwner()))
|
||||
await expect(ctx.sessionPersistence.isLive(inherited)).resolves.toBe(true)
|
||||
const inheritedClaim = await ctx.sessionPersistence.claimLive(inherited)
|
||||
await inheritedClaim.release()
|
||||
|
||||
await expect(ctx.sessionPersistence.claimLive(SessionId('x'.repeat(300))))
|
||||
.rejects.toThrow()
|
||||
|
||||
const guarded = SessionId('guarded-reclaim')
|
||||
const guardedPath = join(liveDir, `${encodeSegment(guarded)}.lock`)
|
||||
await writeFile(guardedPath, JSON.stringify({ pid: 2_147_483_647, nonce: 'dead-owner' }))
|
||||
await writeFile(`${guardedPath}.reclaim`, 'busy')
|
||||
await expect(ctx.sessionPersistence.claimLive(guarded))
|
||||
.rejects.toThrow('reclamation is already in progress')
|
||||
} finally {
|
||||
await ctx.fiber.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('allows exactly one process to reclaim a stale lease', async () => {
|
||||
const dir = await freshRoot()
|
||||
const liveDir = join(dir, '.live')
|
||||
await mkdir(liveDir, { recursive: true })
|
||||
const sessionId = SessionId('reclaim-race')
|
||||
await writeFile(
|
||||
join(liveDir, `${encodeSegment(sessionId)}.lock`),
|
||||
JSON.stringify({ pid: 2_147_483_647, nonce: 'dead-owner' }),
|
||||
)
|
||||
const gate = join(dir, 'race-start')
|
||||
const markers = [join(dir, 'race-a'), join(dir, 'race-b')]
|
||||
const children = markers.map(marker => spawn(
|
||||
process.execPath,
|
||||
['--import', tsxLoader, leaseRaceChild, dir, gate, marker, sessionId],
|
||||
{
|
||||
cwd: repoRoot,
|
||||
env: { ...process.env, TSX_TSCONFIG_PATH: join(repoRoot, 'tsconfig.json') },
|
||||
stdio: ['ignore', 'ignore', 'pipe'],
|
||||
},
|
||||
))
|
||||
const errors = ['', '']
|
||||
children.forEach((child, index) => {
|
||||
child.stderr.setEncoding('utf8')
|
||||
child.stderr.on('data', (chunk: string) => { errors[index] = (errors[index] ?? '') + chunk })
|
||||
})
|
||||
try {
|
||||
await writeFile(gate, 'go')
|
||||
await vi.waitFor(() => Promise.all(markers.map(marker => access(marker))), { timeout: 30_000 })
|
||||
const outcomes = await Promise.all(markers.map(marker => readFile(marker, 'utf8')))
|
||||
expect(outcomes.filter(outcome => outcome === 'claimed')).toHaveLength(1)
|
||||
expect(outcomes.filter(outcome => outcome.startsWith('rejected:'))).toHaveLength(1)
|
||||
|
||||
const winner = children[outcomes.findIndex(outcome => outcome === 'claimed')]!
|
||||
const loser = children[outcomes.findIndex(outcome => outcome.startsWith('rejected:'))]!
|
||||
if (loser.exitCode === null && loser.signalCode === null) {
|
||||
await new Promise<void>(resolve => loser.once('close', () => { resolve() }))
|
||||
}
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(SessionPersistenceJsonl, { root: dir, compression: 'none' })
|
||||
try {
|
||||
await expect(ctx.sessionPersistence.claimLive(sessionId))
|
||||
.rejects.toThrow('occupied by another live process')
|
||||
} finally {
|
||||
await ctx.fiber.dispose()
|
||||
}
|
||||
const closed = new Promise<void>(resolve => winner.once('close', () => { resolve() }))
|
||||
winner.kill()
|
||||
await closed
|
||||
} catch (error) {
|
||||
throw new Error(`live-lease race children failed: ${errors.join('\n')}`, { cause: error })
|
||||
} finally {
|
||||
for (const child of children) {
|
||||
if (child.exitCode === null && child.signalCode === null) child.kill()
|
||||
}
|
||||
}
|
||||
}, 40_000)
|
||||
})
|
||||
|
||||
describe('SessionPersistenceJsonl: durability and crash semantics', () => {
|
||||
let ctx: Context
|
||||
beforeEach(async () => {
|
||||
|
||||
@@ -8,7 +8,7 @@ A SQLite durable session-persistence backend — a second `SessionPersistence` i
|
||||
|
||||
## Storage model
|
||||
|
||||
Each `SessionEvent` maps 1:1 onto a row in an `events` table `(session_id, seq, type, time, data, source_event_seqs, surface_op)` — `data` is the event payload as JSON text, so the row shape is the event verbatim (including `assistant/chunk`, keeping `seq` contiguous). The two `TEXT` columns `source_event_seqs` and `surface_op` are nullable; they store the event's optional surface-metadata fields (see [session surface](../../../.agents/notes/implemented/architecture/2026-06-18-session-surface.md)). Out-of-log metadata (`SessionHeader`), a per-materialization incarnation id, and a monotonic per-log revision live in a `sessions` row; a singleton state row carries the immutable store id. A `sessions` row is written only by the first `append` — its existence is the lazy-materialization signal (`list` reports exactly the sessions that have a row).
|
||||
Each `SessionEvent` maps 1:1 onto a row in an `events` table `(session_id, seq, type, time, data, source_event_seqs, surface_op)` — `data` is the event payload as JSON text, so the row shape is the event verbatim (including `assistant/chunk`, keeping `seq` contiguous). The two `TEXT` columns `source_event_seqs` and `surface_op` are nullable; they store the event's optional surface-metadata fields (see [session surface](../../../.agents/notes/implemented/architecture/2026-06-18-session-surface.md)). Out-of-log metadata (`SessionHeader`), a per-materialization incarnation id, and a monotonic per-log revision live in a `sessions` row; a singleton state row carries the immutable store id, and `live_session_leases` stores one PID and exec-stable nonce per live session. A `sessions` row is written only by the first `append` — its existence is the lazy-materialization signal (`list` reports exactly the sessions that have a row).
|
||||
|
||||
The repository's Node range supports unflagged `node:sqlite`. The database enables foreign keys and uses the configured journal mode (`wal` by default; use a rollback mode where WAL shared-memory files are unsuitable). `PRAGMA user_version` stores the table-layout version; databases with any other version are rejected because this unreleased format has no migrations.
|
||||
|
||||
@@ -33,7 +33,7 @@ interface Config {
|
||||
|
||||
## Write path
|
||||
|
||||
Like the JSONL backend, the plugin copies each frozen `session/event` into one controller per live session and starts an eager drain. Concurrent events share the current transaction; events admitted during it form a follow-up batch, while `session/flush` waits until both current and pending batches are durable. The controller persists a fork's seed once, keeps a write cursor so resume never re-appends stored events, and seeds live sessions on apply because HMR does not replay `session/created`. Dispose drains every retained controller before closing the database.
|
||||
Like the JSONL backend, the plugin copies each frozen `session/event` into one controller per live session and starts an eager drain. A live lease is acquired in a `BEGIN IMMEDIATE` transaction before flush or resume and released after the exact lifecycle retires. Concurrent events share the current transaction; events admitted during it form a follow-up batch, while `session/flush` waits until both current and pending batches are durable. The controller persists a fork's seed once, keeps a write cursor so resume never re-appends stored events, and seeds live sessions on apply because HMR does not replay `session/created`. Dispose drains every retained controller before closing the database.
|
||||
|
||||
## Model Experience
|
||||
|
||||
|
||||
@@ -15,8 +15,9 @@ import { mkdir, open } from 'node:fs/promises'
|
||||
import { dirname, resolve } from 'node:path'
|
||||
import {
|
||||
SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
|
||||
type PersistenceBackend, type SessionLocation, type SessionPersistenceSnapshot,
|
||||
type StoredPrefix,
|
||||
sessionLeaseProcessIsLive, shareSessionLiveLease,
|
||||
type PersistenceBackend, type SessionLiveLease, type SessionLiveOwner,
|
||||
type SessionLocation, type SessionPersistenceSnapshot, type StoredPrefix,
|
||||
} from '@deepseek-ai/dsh-session-persistence'
|
||||
import type { SessionEvent, SurfaceEventType, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import {
|
||||
@@ -161,6 +162,14 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
return this.coordinator.inspect(id)
|
||||
}
|
||||
|
||||
override claimLive(id: SessionId): Promise<SessionLiveLease> {
|
||||
return this.coordinator.claimLive(id)
|
||||
}
|
||||
|
||||
override isLive(id: SessionId): Promise<boolean> {
|
||||
return this.coordinator.isLive(id)
|
||||
}
|
||||
|
||||
// One method serves both public `list` and the backend hook; delegating it to
|
||||
// the coordinator would call this hook recursively.
|
||||
|
||||
@@ -271,6 +280,55 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
}))
|
||||
}
|
||||
|
||||
/** Atomically acquire one SQLite-backed process lease. */
|
||||
async acquireLive(id: SessionId, owner: SessionLiveOwner): Promise<() => Promise<void>> {
|
||||
await this.ready
|
||||
return shareSessionLiveLease(
|
||||
`sqlite:${this.storeIdentity}:${id}`,
|
||||
() => Promise.resolve().then(() => this.acquireLiveRow(id, owner)),
|
||||
)
|
||||
}
|
||||
|
||||
private acquireLiveRow(id: SessionId, owner: SessionLiveOwner): () => Promise<void> {
|
||||
this.db.exec('BEGIN IMMEDIATE')
|
||||
try {
|
||||
const current = this.liveLeaseFor(id)
|
||||
if (current !== undefined
|
||||
&& (current.pid !== owner.pid || current.nonce !== owner.nonce)) {
|
||||
if (sessionLeaseProcessIsLive(current.pid)) {
|
||||
throw new Error(`session "${id}" is occupied by another live process`)
|
||||
}
|
||||
this.db.prepare('DELETE FROM live_session_leases WHERE session_id = ?').run(id)
|
||||
}
|
||||
this.db.prepare(`
|
||||
INSERT INTO live_session_leases (session_id, pid, nonce) VALUES (?, ?, ?)
|
||||
ON CONFLICT(session_id) DO UPDATE SET pid = excluded.pid, nonce = excluded.nonce
|
||||
`).run(id, owner.pid, owner.nonce)
|
||||
this.db.exec('COMMIT')
|
||||
} catch (error) {
|
||||
this.db.exec('ROLLBACK')
|
||||
throw error
|
||||
}
|
||||
return async () => {
|
||||
await this.ready
|
||||
this.db.prepare(
|
||||
'DELETE FROM live_session_leases WHERE session_id = ? AND pid = ? AND nonce = ?',
|
||||
).run(id, owner.pid, owner.nonce)
|
||||
}
|
||||
}
|
||||
|
||||
/** Report a non-stale SQLite lease and remove a crashed owner's row. */
|
||||
async inspectLive(id: SessionId, owner: SessionLiveOwner): Promise<boolean> {
|
||||
await this.ready
|
||||
const current = this.liveLeaseFor(id)
|
||||
if (current === undefined) return false
|
||||
if ((current.pid === owner.pid && current.nonce === owner.nonce)
|
||||
|| sessionLeaseProcessIsLive(current.pid)) return true
|
||||
this.db.prepare('DELETE FROM live_session_leases WHERE session_id = ? AND pid = ? AND nonce = ?')
|
||||
.run(id, current.pid, current.nonce)
|
||||
return false
|
||||
}
|
||||
|
||||
/** Close the database handle (awaited by the coordinator's dispose, post-drain). */
|
||||
async close(): Promise<void> {
|
||||
await this.ready
|
||||
@@ -284,6 +342,11 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
|
||||
return this.db.prepare('SELECT * FROM sessions WHERE id = ?').get(id) as unknown as SessionRow | undefined
|
||||
}
|
||||
|
||||
private liveLeaseFor(id: SessionId): { pid: number; nonce: string } | undefined {
|
||||
return this.db.prepare('SELECT pid, nonce FROM live_session_leases WHERE session_id = ?')
|
||||
.get(id) as { pid: number; nonce: string } | undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Insert-or-replace a session's metadata row. The only caller is the first
|
||||
* materializing `appendBatch`, so writing the row IS the materialization (its
|
||||
|
||||
@@ -17,7 +17,7 @@ import type { SessionEvent, SessionId, SessionHeader, SurfaceOp } from '@deepsee
|
||||
* layout; orthogonal to a session's own `version` (which versions the EVENT
|
||||
* vocabulary, stored per session in the `sessions` row).
|
||||
*/
|
||||
export const SCHEMA_VERSION = 8
|
||||
export const SCHEMA_VERSION = 9
|
||||
|
||||
/**
|
||||
* A row of the `sessions` table — the out-of-log metadata ({@link SessionHeader}).
|
||||
@@ -68,7 +68,7 @@ export type JournalMode = 'wal' | 'delete' | 'truncate' | 'persist'
|
||||
* rather than being migrated in place.
|
||||
* @param path - the SQLite database file to open (created when absent).
|
||||
* @param journalMode - validated journal pragma.
|
||||
* @returns the open handle with pragmas applied and all three tables ensured.
|
||||
* @returns the open handle with pragmas applied and all tables ensured.
|
||||
*/
|
||||
export function openDatabase(path: string, journalMode: JournalMode): DatabaseSync {
|
||||
const db = new DatabaseSync(path)
|
||||
@@ -128,6 +128,13 @@ function configureDatabase(db: DatabaseSync, path: string, journalMode: JournalM
|
||||
PRIMARY KEY (session_id, seq)
|
||||
) STRICT
|
||||
`)
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS live_session_leases (
|
||||
session_id TEXT PRIMARY KEY,
|
||||
pid INTEGER NOT NULL,
|
||||
nonce TEXT NOT NULL
|
||||
) STRICT
|
||||
`)
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -7,6 +7,7 @@ import { dirname, join } from 'node:path'
|
||||
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
||||
import type { Session, SessionEvent, SurfaceEvent, SurfaceEventType } from '@deepseek-ai/dsh-session'
|
||||
import SessionPersistenceSqlite, { SCHEMA_VERSION } from '@deepseek-ai/dsh-session-persistence-sqlite'
|
||||
import { sessionLiveOwner } from '@deepseek-ai/dsh-session-persistence'
|
||||
import { openDatabase, rowToEvent, scanRows, type EventRow } from '../src/schema.ts'
|
||||
import { runPersistenceContract, meta, oneTurnLog, appendLog } from '../../session-persistence/tests/contract.ts'
|
||||
import { runCoordinatorContract, type CoordinatorFixture } from '../../session-persistence/tests/coordinator-contract.ts'
|
||||
@@ -442,7 +443,7 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
|
||||
})
|
||||
|
||||
it('exposes the schema version constant', () => {
|
||||
expect(SCHEMA_VERSION).toBe(8)
|
||||
expect(SCHEMA_VERSION).toBe(9)
|
||||
})
|
||||
|
||||
it('keeps the revision stable for an empty repair hook', async () => {
|
||||
@@ -458,6 +459,38 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
|
||||
})
|
||||
|
||||
describe('SessionPersistenceSqlite: edge cases', () => {
|
||||
it('claims, rejects, reclaims, inspects, and releases SQLite live leases', async () => {
|
||||
const path = await freshDbPath()
|
||||
const b = await backend(path)
|
||||
await b.ctx.sessionPersistence.list()
|
||||
const concrete = b.ctx.sessionPersistence as SessionPersistenceSqlite
|
||||
const owner = sessionLiveOwner()
|
||||
const db = openDatabase(path, 'wal')
|
||||
const insert = db.prepare('INSERT INTO live_session_leases (session_id, pid, nonce) VALUES (?, ?, ?)')
|
||||
insert.run('occupied-lease', process.pid, 'another-owner')
|
||||
insert.run('stale-claim', 2_147_483_647, 'dead-owner')
|
||||
insert.run('stale-inspect', 2_147_483_647, 'dead-owner')
|
||||
insert.run('owned-inspect', owner.pid, owner.nonce)
|
||||
db.close()
|
||||
|
||||
await expect(concrete.acquireLive(SessionId('occupied-lease'), owner))
|
||||
.rejects.toThrow('occupied by another live process')
|
||||
const claim = await concrete.acquireLive(SessionId('stale-claim'), owner)
|
||||
expect(await concrete.inspectLive(SessionId('owned-inspect'), owner)).toBe(true)
|
||||
expect(await concrete.inspectLive(SessionId('stale-inspect'), owner)).toBe(false)
|
||||
expect(await concrete.inspectLive(SessionId('missing-inspect'), owner)).toBe(false)
|
||||
await claim()
|
||||
await b.dispose()
|
||||
|
||||
const memory = new Context()
|
||||
await memory.plugin(SessionStore)
|
||||
await memory.plugin(SessionPersistenceSqlite, { path: ':memory:' })
|
||||
const memoryClaim = await memory.sessionPersistence.claimLive(SessionId('memory-live'))
|
||||
expect(await memory.sessionPersistence.isLive(SessionId('memory-live'))).toBe(true)
|
||||
await memoryClaim.release()
|
||||
await memory.fiber.dispose()
|
||||
})
|
||||
|
||||
it('rejects and closes a current-schema database with an invalid store identity', async () => {
|
||||
const path = await freshDbPath()
|
||||
const db = openDatabase(path, 'wal')
|
||||
|
||||
@@ -15,6 +15,10 @@ The persisted unit IS the existing `SessionEvent` (event-sourced model — the l
|
||||
| `inspect(id): Promise<{ meta; events }>` | Return a detached valid stored prefix without truncating a torn tail, synthesizing recovery closers, or publishing coordinator state. Serialized with same-id writes; intended for read models and other observers that must never recover a log. |
|
||||
| `list(): Promise<SessionHeader[]>` | Lightweight listing from metadata, no full-log parse. A zero-event lazily-materialized session is absent from `list`. |
|
||||
| `listSnapshots(): Promise<SessionPersistenceSnapshot[]>` | Lightweight metadata plus an opaque branded per-log revision, without loading event logs. A revision stays equal while that log and its backing store are unchanged, changes after append or mutating load repair, and cannot collide solely because two stores use the same local counter. |
|
||||
| `claimLive(id): Promise<SessionLiveLease>` | Atomically claim live ownership. First-party backends reject another live process and reclaim a dead owner; release follows quiescence. |
|
||||
| `isLive(id): Promise<boolean>` | Report a current non-stale live lease, including one owned by this process. |
|
||||
|
||||
The abstract base supplies a process-local fallback for lightweight third-party implementations. A backend that needs multi-process safety overrides both live-lease methods.
|
||||
|
||||
## Invariants every backend must honor
|
||||
|
||||
@@ -31,7 +35,7 @@ Each `session/event` copies its event into the session controller and starts an
|
||||
|
||||
Crash repair is cold-only. For a live id, `load(id)` snapshots the authoritative in-memory log, waits for that snapshot to become durable, and returns it with the coordinator's stored header only when balanced; an open live turn rejects instead of receiving synthetic interruption closers. A cold load reserves its id across backend reads and repair writes, so concurrent publication of a same-id live `Session` rejects and rolls back. HMR adoption reads through `loadStored`, applies the coordinator's cwd check, and never closes the active turn.
|
||||
|
||||
When a live session emits `session/disposed`, the coordinator waits for its controller, serializes a final drain, then releases state owned by that exact `Session` object. Failed retirement leaves the controller in the live-session map, so backend teardown can retry it. Backend teardown stops event admission first, flushes every remaining controller, awaits per-id operations, and only then closes the storage handle.
|
||||
When a live session emits `session/disposed`, the coordinator waits for its controller, serializes a final drain, then releases state and the backend-owned live lease for that exact `Session` object. Failed retirement leaves the controller in the live-session map, so backend teardown can retry it. Backend teardown stops event admission first, flushes every remaining controller, releases their leases, awaits per-id operations, and only then closes the storage handle.
|
||||
|
||||
The side-effect-free `locate` and lightweight `listSnapshots` queries remain backend-owned because they describe storage topology and revision identity rather than write orchestration.
|
||||
|
||||
@@ -44,6 +48,8 @@ The `PersistenceBackend<TornMarker>` hooks (the only seam between the coordinato
|
||||
| `appendBatch(meta, events, isMaterialized)` | Durably append a contiguous batch, lazily materializing ATOMICALLY when not yet materialized. |
|
||||
| `commitRepair(meta, tornMarker, closers)` | Make a crash repair durable: truncate the torn tail (iff `tornMarker !== undefined` — a marker may be falsy, e.g. seq/offset `0`) and append `closers`. NOT required to be atomic. Used by load (truncate + closers) and live-adoption (truncate only). |
|
||||
| `list()` | List all stored metadata. |
|
||||
| `acquireLive?(id, owner)` | Atomically acquire a backend-owned cross-process lease and return its physical release. |
|
||||
| `inspectLive?(id, owner)` | Report or reclaim a backend-owned lease without acquiring it. |
|
||||
| `close?()` | Optional lifecycle teardown (e.g. close a db handle), awaited after the dispose drain. |
|
||||
|
||||
The coordinator asserts the stored id and compares stored/live cwd before repair or live adoption. Its `inspect()` path validates and clones the prefix without calling `commitRepair` or publishing write state. The `tornMarker` is fully OPAQUE: the coordinator only tests `!== undefined` and round-trips it to `commitRepair`, never inspecting its value (the JSONL backend uses the byte offset to truncate to, the SQLite backend the seq to delete from). A third-party backend MAY implement the abstract service directly without the coordinator, but it must provide the same non-mutating inspection and trustworthy lightweight snapshot revisions. See [the write-coordinator Agent Note](../../../.agents/notes/implemented/architecture/2026-06-18-shared-persistence-write-coordinator.md).
|
||||
|
||||
@@ -8,6 +8,8 @@
|
||||
import { Context } from 'cordis'
|
||||
import { interruptedTurnClosers, SESSION_FORMAT_VERSION, snapshotJsonValue } from '@deepseek-ai/dsh-session'
|
||||
import type { Session, SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import { sessionLiveOwner } from './lease.ts'
|
||||
import type { SessionLiveLease, SessionLiveOwner } from './lease.ts'
|
||||
|
||||
/**
|
||||
* A stored session's header, valid contiguous event prefix, and optional opaque
|
||||
@@ -63,6 +65,12 @@ export interface PersistenceBackend<TornMarker = unknown> {
|
||||
/** List all stored (materialized) sessions' metadata. */
|
||||
list(): Promise<SessionHeader[]>
|
||||
|
||||
/** Optionally acquire a backend-owned cross-process live-session lease. */
|
||||
acquireLive?(id: SessionId, owner: SessionLiveOwner): Promise<() => Promise<void>>
|
||||
|
||||
/** Optionally inspect and reclaim a backend-owned live-session lease. */
|
||||
inspectLive?(id: SessionId, owner: SessionLiveOwner): Promise<boolean>
|
||||
|
||||
/**
|
||||
* Optional lifecycle teardown (e.g. close a database handle). Awaited by the
|
||||
* coordinator's dispose effect AFTER the quiescence drain. A stateless file
|
||||
@@ -96,6 +104,7 @@ interface LiveSessionState {
|
||||
pending: SessionEvent[]
|
||||
init: Promise<void>
|
||||
flush: Promise<void> | undefined
|
||||
lease?: SessionLiveLease
|
||||
}
|
||||
|
||||
/** Collect the rejection reasons from a set of promises (none-throwing). */
|
||||
@@ -161,6 +170,12 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
* same id, so writes for one session never interleave. Keyed by session id.
|
||||
*/
|
||||
private chains = new Map<SessionId, Promise<unknown>>()
|
||||
/** One backend lease with process-local reference counting per session id. */
|
||||
private liveClaims = new Map<SessionId, {
|
||||
refs: number
|
||||
releaseBackend: () => Promise<void>
|
||||
}>()
|
||||
private readonly liveOwner = sessionLiveOwner()
|
||||
|
||||
constructor(private ctx: Context, private backend: PersistenceBackend<TornMarker>) {
|
||||
this.installWritePath()
|
||||
@@ -273,6 +288,61 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
return this.serialize(id, () => this.inspectCore(id))
|
||||
}
|
||||
|
||||
/**
|
||||
* Acquire one process-local reference to the backend's cross-process lease.
|
||||
* @param id - session identity about to become live.
|
||||
* @returns one idempotent release capability.
|
||||
*/
|
||||
async claimLive(id: SessionId): Promise<SessionLiveLease> {
|
||||
const acquireLive = this.backend.acquireLive?.bind(this.backend)
|
||||
if (acquireLive === undefined) return { release: () => Promise.resolve() }
|
||||
await this.serialize(id, async () => {
|
||||
const existing = this.liveClaims.get(id)
|
||||
if (existing !== undefined) {
|
||||
existing.refs += 1
|
||||
return
|
||||
}
|
||||
const releaseBackend = await acquireLive(id, this.liveOwner)
|
||||
this.liveClaims.set(id, { refs: 1, releaseBackend })
|
||||
})
|
||||
let releaseTask: Promise<void> | undefined
|
||||
return {
|
||||
release: () => {
|
||||
if (releaseTask !== undefined) return releaseTask
|
||||
const task = this.serialize(id, async () => {
|
||||
const claim = this.liveClaims.get(id)
|
||||
/* v8 ignore next -- this capability is returned only after its claim enters the serialized map */
|
||||
if (claim === undefined) return
|
||||
claim.refs -= 1
|
||||
if (claim.refs > 0) return
|
||||
try {
|
||||
await claim.releaseBackend()
|
||||
} catch (error) {
|
||||
claim.refs += 1
|
||||
throw error
|
||||
}
|
||||
this.liveClaims.delete(id)
|
||||
})
|
||||
const wrapped = task.catch((error: unknown) => {
|
||||
releaseTask = undefined
|
||||
throw error
|
||||
})
|
||||
releaseTask = wrapped
|
||||
return wrapped
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Check the backend's current cross-process lease state.
|
||||
* @param id - session identity to inspect.
|
||||
* @returns whether this or another live process owns the session.
|
||||
*/
|
||||
isLive(id: SessionId): Promise<boolean> {
|
||||
if (this.liveClaims.has(id)) return Promise.resolve(true)
|
||||
return this.backend.inspectLive?.(id, this.liveOwner) ?? Promise.resolve(false)
|
||||
}
|
||||
|
||||
private async inspectCore(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
|
||||
const stored = await this.backend.loadStored(id)
|
||||
if (stored === undefined) throw new Error(`session "${id}" not found`)
|
||||
@@ -382,6 +452,9 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
let disposeError: unknown
|
||||
try {
|
||||
const errors = await settledErrors([...this.live.keys()].map(session => this.flush(session)))
|
||||
errors.push(...await settledErrors(
|
||||
[...this.live.values()].flatMap(live => live.lease === undefined ? [] : [live.lease.release()]),
|
||||
))
|
||||
while (this.chains.size > 0) await Promise.allSettled([...this.chains.values()])
|
||||
if (errors.length > 0) {
|
||||
throw new AggregateError(errors, `${this.backend.name} dispose failed`)
|
||||
@@ -441,6 +514,8 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
private async retireCore(session: Session): Promise<void> {
|
||||
await this.flush(session)
|
||||
const id = session.header.id
|
||||
const live = this.live.get(session)
|
||||
await live?.lease?.release()
|
||||
await this.serialize(id, () => {
|
||||
this.live.delete(session)
|
||||
if (this.states.get(id)?.owner === session) this.states.delete(id)
|
||||
@@ -454,7 +529,16 @@ export class PersistenceCoordinator<TornMarker = unknown> {
|
||||
const seed = session.events.map(e => structuredClone(e))
|
||||
const live: LiveSessionState = { pending: [], init: Promise.resolve(), flush: undefined }
|
||||
this.live.set(session, live)
|
||||
live.init = this.serialize(session.header.id, () => this.onCreated(session, seed))
|
||||
live.init = this.claimLive(session.id).then(async (lease) => {
|
||||
live.lease = lease
|
||||
try {
|
||||
await this.serialize(session.header.id, () => this.onCreated(session, seed))
|
||||
} catch (error) {
|
||||
delete live.lease
|
||||
await lease.release()
|
||||
throw error
|
||||
}
|
||||
})
|
||||
live.init.catch(() => { /* observed by flush/dispose through the controller */ })
|
||||
return live
|
||||
}
|
||||
|
||||
@@ -8,10 +8,13 @@
|
||||
import { Context, Service } from 'cordis'
|
||||
import type { SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import type { SessionPersistenceRevision } from './revision.ts'
|
||||
import type { SessionLiveLease } from './lease.ts'
|
||||
|
||||
// Re-export the metadata vocabulary so consumers import it from the seam.
|
||||
export type { SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
export { SessionPersistenceRevision } from './revision.ts'
|
||||
export { sessionLeaseProcessIsLive, sessionLiveOwner, shareSessionLiveLease } from './lease.ts'
|
||||
export type { SessionLiveLease, SessionLiveOwner } from './lease.ts'
|
||||
|
||||
/** Lightweight immutable source identity returned without loading a full log. */
|
||||
export interface SessionPersistenceSnapshot {
|
||||
@@ -50,6 +53,8 @@ export interface SessionLocation {
|
||||
* rewriting committed events.
|
||||
*/
|
||||
export abstract class SessionPersistence extends Service {
|
||||
private readonly localLiveClaims = new Map<SessionId, number>()
|
||||
|
||||
constructor(ctx: Context) {
|
||||
super(ctx, 'sessionPersistence')
|
||||
}
|
||||
@@ -123,6 +128,39 @@ export abstract class SessionPersistence extends Service {
|
||||
* @returns one header and opaque revision per materialized session without loading full logs.
|
||||
*/
|
||||
abstract listSnapshots(): Promise<SessionPersistenceSnapshot[]>
|
||||
|
||||
/**
|
||||
* Atomically acquire this process's live ownership of a session id.
|
||||
* Reentrant claims share one backend lease. First-party backends override
|
||||
* this process-local fallback to reject another live process and reclaim a
|
||||
* dead owner.
|
||||
* @param id - session identity that is about to become live.
|
||||
* @returns a single-release reference owned by the caller.
|
||||
*/
|
||||
claimLive(id: SessionId): Promise<SessionLiveLease> {
|
||||
this.localLiveClaims.set(id, (this.localLiveClaims.get(id) ?? 0) + 1)
|
||||
let released = false
|
||||
return Promise.resolve({
|
||||
release: () => {
|
||||
if (released) return Promise.resolve()
|
||||
released = true
|
||||
const refs = this.localLiveClaims.get(id) as number
|
||||
if (refs <= 1) this.localLiveClaims.delete(id)
|
||||
else this.localLiveClaims.set(id, refs - 1)
|
||||
return Promise.resolve()
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Check whether any process currently owns a live lease for this session.
|
||||
* The base implementation reports only claims on this service instance.
|
||||
* @param id - persisted or prospective session identity.
|
||||
* @returns true while a non-stale lease exists, including this process's lease.
|
||||
*/
|
||||
isLive(id: SessionId): Promise<boolean> {
|
||||
return Promise.resolve(this.localLiveClaims.has(id))
|
||||
}
|
||||
}
|
||||
|
||||
export default SessionPersistence
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
/** Process-backed identity helpers for cross-process live-session leases. */
|
||||
|
||||
import { randomUUID } from 'node:crypto'
|
||||
|
||||
const LIVE_OWNER_ENV = 'DSH_SESSION_LIVE_OWNER'
|
||||
|
||||
/** Process identity stored in backend-owned cross-process live-session leases. */
|
||||
export interface SessionLiveOwner {
|
||||
/** Operating-system process id; retained across an `execve` handoff. */
|
||||
readonly pid: number
|
||||
/** Per-process-start nonce that distinguishes PID reuse. */
|
||||
readonly nonce: string
|
||||
}
|
||||
|
||||
/** Idempotent capability releasing one acquired live-session lease reference. */
|
||||
export interface SessionLiveLease {
|
||||
/** Release this caller's lease reference after its live session reaches quiescence. */
|
||||
release(): Promise<void>
|
||||
}
|
||||
|
||||
/**
|
||||
* Stable owner inherited only by an exec-replaced process, not inferred from a session id.
|
||||
* @returns this process's PID and exec-stable nonce.
|
||||
*/
|
||||
export function sessionLiveOwner(): SessionLiveOwner {
|
||||
const nonce = process.env[LIVE_OWNER_ENV] ?? randomUUID()
|
||||
process.env[LIVE_OWNER_ENV] = nonce
|
||||
return { pid: process.pid, nonce }
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether a lease pid still names a process; permission denial counts as live.
|
||||
* @param pid - positive operating-system process id from a lease record.
|
||||
* @returns true unless the operating system reports that the process is absent.
|
||||
*/
|
||||
export function sessionLeaseProcessIsLive(pid: number): boolean {
|
||||
try {
|
||||
process.kill(pid, 0)
|
||||
return true
|
||||
} catch (error) {
|
||||
return (error as NodeJS.ErrnoException).code !== 'ESRCH'
|
||||
}
|
||||
}
|
||||
|
||||
interface SharedLeaseEntry {
|
||||
refs: number
|
||||
readonly acquired: Promise<() => Promise<void>>
|
||||
}
|
||||
|
||||
const sharedLeases = new Map<string, SharedLeaseEntry>()
|
||||
|
||||
/**
|
||||
* Reference-count one physical lease across backend instances in this process.
|
||||
* @param key - backend-kind plus canonical storage location and session id.
|
||||
* @param acquire - single physical acquisition performed for the first reference.
|
||||
* @returns an idempotent release for this caller's reference.
|
||||
*/
|
||||
export async function shareSessionLiveLease(
|
||||
key: string,
|
||||
acquire: () => Promise<() => Promise<void>>,
|
||||
): Promise<() => Promise<void>> {
|
||||
let entry = sharedLeases.get(key)
|
||||
if (entry === undefined) {
|
||||
entry = { refs: 0, acquired: acquire() }
|
||||
sharedLeases.set(key, entry)
|
||||
void entry.acquired.catch(() => {
|
||||
/* v8 ignore next -- no public operation can replace a still-acquiring module-private entry */
|
||||
if (sharedLeases.get(key) === entry) sharedLeases.delete(key)
|
||||
})
|
||||
}
|
||||
entry.refs += 1
|
||||
try {
|
||||
await entry.acquired
|
||||
} catch (error) {
|
||||
entry.refs -= 1
|
||||
throw error
|
||||
}
|
||||
let releaseTask: Promise<void> | undefined
|
||||
return () => {
|
||||
if (releaseTask !== undefined) return releaseTask
|
||||
const task = (async () => {
|
||||
entry.refs -= 1
|
||||
if (entry.refs > 0 || sharedLeases.get(key) !== entry) return
|
||||
const release = await entry.acquired
|
||||
await release()
|
||||
/* v8 ignore next -- the entry remains installed until this exact final release succeeds */
|
||||
if (sharedLeases.get(key) === entry) sharedLeases.delete(key)
|
||||
})()
|
||||
const wrapped = task.catch((error: unknown) => {
|
||||
entry.refs += 1
|
||||
/* v8 ignore next -- this closure is the sole writer of its releaseTask until settlement */
|
||||
if (releaseTask === wrapped) releaseTask = undefined
|
||||
throw error
|
||||
})
|
||||
releaseTask = wrapped
|
||||
return wrapped
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import {
|
||||
sessionLeaseProcessIsLive,
|
||||
sessionLiveOwner,
|
||||
shareSessionLiveLease,
|
||||
} from '../src/lease.ts'
|
||||
|
||||
const originalOwner = process.env.DSH_SESSION_LIVE_OWNER
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks()
|
||||
if (originalOwner === undefined) delete process.env.DSH_SESSION_LIVE_OWNER
|
||||
else process.env.DSH_SESSION_LIVE_OWNER = originalOwner
|
||||
})
|
||||
|
||||
describe('process live-session lease helpers', () => {
|
||||
it('creates one exec-stable owner identity and classifies process liveness', () => {
|
||||
delete process.env.DSH_SESSION_LIVE_OWNER
|
||||
const first = sessionLiveOwner()
|
||||
expect(first.pid).toBe(process.pid)
|
||||
expect(typeof first.nonce).toBe('string')
|
||||
expect(sessionLiveOwner()).toEqual(first)
|
||||
expect(sessionLeaseProcessIsLive(process.pid)).toBe(true)
|
||||
|
||||
const missing = Object.assign(new Error('missing'), { code: 'ESRCH' })
|
||||
vi.spyOn(process, 'kill').mockImplementationOnce(() => { throw missing })
|
||||
expect(sessionLeaseProcessIsLive(999_999)).toBe(false)
|
||||
const denied = Object.assign(new Error('denied'), { code: 'EPERM' })
|
||||
vi.spyOn(process, 'kill').mockImplementationOnce(() => { throw denied })
|
||||
expect(sessionLeaseProcessIsLive(999_998)).toBe(true)
|
||||
})
|
||||
|
||||
it('shares one physical lease until every process-local reference releases', async () => {
|
||||
const releasePhysical = vi.fn<() => Promise<void>>(() => Promise.resolve())
|
||||
const acquire = vi.fn<() => Promise<() => Promise<void>>>(() => Promise.resolve(releasePhysical))
|
||||
const key = `shared-${randomUUID()}`
|
||||
const first = await shareSessionLiveLease(key, acquire)
|
||||
const second = await shareSessionLiveLease(key, acquire)
|
||||
expect(acquire).toHaveBeenCalledTimes(1)
|
||||
await first()
|
||||
expect(releasePhysical).not.toHaveBeenCalled()
|
||||
await second()
|
||||
await second()
|
||||
expect(releasePhysical).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('removes failed acquisitions and retries a failed physical release', async () => {
|
||||
const key = `retry-${randomUUID()}`
|
||||
await expect(shareSessionLiveLease(key, () => Promise.reject(new Error('claim failed'))))
|
||||
.rejects.toThrow('claim failed')
|
||||
|
||||
let releases = 0
|
||||
const release = await shareSessionLiveLease(key, () => Promise.resolve(async () => {
|
||||
releases += 1
|
||||
if (releases === 1) throw new Error('release failed')
|
||||
}))
|
||||
await expect(release()).rejects.toThrow('release failed')
|
||||
await expect(release()).resolves.toBeUndefined()
|
||||
expect(releases).toBe(2)
|
||||
})
|
||||
})
|
||||
@@ -4,7 +4,7 @@ import SessionStore, { SessionId, isJsonValue } from '@deepseek-ai/dsh-session'
|
||||
import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
|
||||
import {
|
||||
SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
|
||||
type PersistenceBackend, type SessionPersistenceSnapshot, type StoredPrefix,
|
||||
type PersistenceBackend, type SessionLiveOwner, type SessionPersistenceSnapshot, type StoredPrefix,
|
||||
} from '../src/index.ts'
|
||||
import { runPersistenceContract, meta, oneTurnLog } from './contract.ts'
|
||||
import { runCoordinatorContract, type CoordinatorFixture } from './coordinator-contract.ts'
|
||||
@@ -348,6 +348,46 @@ describe('PersistenceCoordinator stored identity', () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe('PersistenceCoordinator live leases', () => {
|
||||
it('degrades without backend hooks and retries a failed final release', async () => {
|
||||
const fallbackCtx = new Context()
|
||||
await fallbackCtx.plugin(SessionStore)
|
||||
const fallback = new PersistenceCoordinator(fallbackCtx, new ControlledBackend())
|
||||
const fallbackClaim = await fallback.claimLive(SessionId('fallback-live'))
|
||||
expect(await fallback.isLive(SessionId('fallback-live'))).toBe(false)
|
||||
await fallbackClaim.release()
|
||||
await fallbackCtx.fiber.dispose()
|
||||
|
||||
class LeaseBackend extends ControlledBackend {
|
||||
releaseAttempts = 0
|
||||
async acquireLive(_id: SessionId, _owner: SessionLiveOwner): Promise<() => Promise<void>> {
|
||||
return async () => {
|
||||
this.releaseAttempts += 1
|
||||
if (this.releaseAttempts === 1) throw new Error('lease release failed')
|
||||
}
|
||||
}
|
||||
inspectLive(): Promise<boolean> {
|
||||
return Promise.resolve(true)
|
||||
}
|
||||
}
|
||||
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const backend = new LeaseBackend()
|
||||
const coordinator = new PersistenceCoordinator(ctx, backend)
|
||||
const first = await coordinator.claimLive(SessionId('leased'))
|
||||
const second = await coordinator.claimLive(SessionId('leased'))
|
||||
expect(await coordinator.isLive(SessionId('leased'))).toBe(true)
|
||||
await first.release()
|
||||
await expect(second.release()).rejects.toThrow('lease release failed')
|
||||
await expect(second.release()).resolves.toBeUndefined()
|
||||
await expect(second.release()).resolves.toBeUndefined()
|
||||
expect(backend.releaseAttempts).toBe(2)
|
||||
expect(await coordinator.isLive(SessionId('leased'))).toBe(true)
|
||||
await ctx.fiber.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
describe('PersistenceCoordinator retirement', () => {
|
||||
it('a retiring unmaterialized owner without buffered events releases its id', async () => {
|
||||
const ctx = new Context()
|
||||
@@ -795,4 +835,20 @@ describe('SessionPersistence service registration', () => {
|
||||
await fiber.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('provides a reference-counted process-local lease fallback', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(MemoryPersistence)
|
||||
const id = SessionId('local-live')
|
||||
const first = await ctx.sessionPersistence.claimLive(id)
|
||||
const second = await ctx.sessionPersistence.claimLive(id)
|
||||
expect(await ctx.sessionPersistence.isLive(id)).toBe(true)
|
||||
await first.release()
|
||||
await first.release()
|
||||
expect(await ctx.sessionPersistence.isLive(id)).toBe(true)
|
||||
await second.release()
|
||||
expect(await ctx.sessionPersistence.isLive(id)).toBe(false)
|
||||
await ctx.fiber.dispose()
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user