1
/**2
* The JSONL provider's session storage runtime: its concrete write/read3
* handle with a per-handle mutation chain and a routed live write-behind4
* buffer, the in-process bookkeeping that enforces one active writer per5
* session id, and the backend's live event routing and teardown. Deliberately6
* provider-local: the persistence seam exposes only the service and handle7
* contracts, and the shared contract suites pin equivalent observable8
* behavior across providers.9
* @module10
*/12
import type { Context } from '@deepseek-ai/cordis'13
import { errorChain } from '@deepseek-ai/dsh-llm'14
import type { Session, SessionEvent, SessionHeader, SessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'15
import {16
assertContiguous,17
SessionAlreadyExistsError,18
materializeAppendBatch,19
SessionAlreadyOwnedError,20
SessionHandleClosedError,21
SessionPersistenceNotFoundError,22
SessionPersistenceRevision,23
SessionReadOnlyError,24
} from '@deepseek-ai/dsh-session-persistence'25
import type {26
SessionAccess,27
SessionHandle,28
SessionHandleAppendOptions,29
SessionHandleFlushOptions,30
SessionHandleReadOptions,31
SessionHandleReadResult,32
} from '@deepseek-ai/dsh-session-persistence'33
import type { SessionWriteLease } from './lease.ts'35
/** Maximum intentional wait before a routed live session batch starts writing. */36
export const LIVE_WRITE_BATCH_MAX_DELAY_MS = 20038
/** The file-storage primitives the handle drives on its owning service. */39
export interface JsonlHandleStorage {40
/** Append encoded lines; `isMaterialized` selects create-vs-extend publication. */41
persistBatch(42
header: SessionHeader,43
events: readonly SessionEvent[],44
isMaterialized: boolean,45
inheritedEventCount: SessionLogOffset,46
): Promise<void>47
/** Materialize the header-only artifact for an explicitly flushed empty session. */48
persistHeader(header: SessionHeader, inheritedEventCount: SessionLogOffset): Promise<void>49
/** Truncate a torn physical tail before the first new append lands. */50
truncateTornTail(header: SessionHeader, truncateTo: number): Promise<void>51
/** Resolve the current-generation artifact path, or `undefined` when absent. */52
resolveCurrentLog(id: SessionId, signal?: AbortSignal): Promise<string | undefined>53
/** Read and validate the stored log at `path`, including its established event aliasing state. */54
readStoredLog(path: string, expectedId: SessionId, signal?: AbortSignal): Promise<SessionHandleReadResult>55
/** Whether the id is still a created-but-unmaterialized session here. */56
hasPendingSession(id: SessionId): boolean57
/** Acquire the session's cross-process write lock in its artifact directory. */58
acquireWriteLease(header: SessionHeader): Promise<SessionWriteLease>59
/** Drop the handle's bookkeeping on close. */60
releaseHandle(handle: JsonlSessionHandle, materialized: boolean): void61
}63
/** Mutable per-handle log state; a write handle is its session's single mutator. */64
export interface StorageHandleState {65
/** The stored next-seq (the logical end this handle knows). */66
cursor: number67
/** Whether the session has a durable artifact yet. */68
materialized: boolean69
/** Torn-tail truncation point, consumed by the first new append. */70
tornTruncateTo?: number | undefined71
/** Complete events recovered from the torn final frame; the first mutation rewrites them durably. */72
recoveredTail?: SessionEvent[] | undefined73
/** Exact fork-inherited prefix length stored with the log; `0` when unseeded. */74
inheritedEventCount: SessionLogOffset75
/** The validated stored prefix from a write open, served to reads until the first append. */76
primed?: SessionHandleReadResult | undefined77
}79
/**80
* The JSONL session handle. Mutations serialize on a per-handle promise81
* chain; reads re-scan the artifact on demand and never observe a shorter log82
* than a prior read on this handle. Routed live events buffer in a bounded83
* window and drain through the same chain as explicit appends.84
*/85
export class JsonlSessionHandle implements SessionHandle {86
private chain: Promise<unknown> = Promise.resolve()87
private closing: Promise<void> | undefined88
private observedLength = 089
/** Routed live events awaiting their batching deadline (persistence-owned copies). */90
private buffered: SessionEvent[] = []91
private batchTimer: ReturnType<typeof setTimeout> | undefined92
/** Set when a drain failed; the automatic timer stays quiet until the next drain. */93
private drainPaused = false94
private draining: Promise<void> | undefined96
constructor(97
private readonly storage: JsonlHandleStorage,98
readonly id: SessionId,99
readonly header: SessionHeader,100
readonly access: SessionAccess,101
private readonly state: StorageHandleState,102
/** The cross-process write lock; a create handle acquires it lazily at first materialization. */103
private lease?: SessionWriteLease,104
) {}106
/** Exact fork-inherited prefix length stored with this session's log. */107
get inheritedEventCount(): SessionLogOffset {108
return this.state.inheritedEventCount109
}111
/**112
* Read a slice of the valid contiguous logical log; see the seam contract.113
* @param offset - first logical seq to include (default 0).114
* @param length - maximum events returned (default: the rest).115
* @param options - optional cancellation.116
* @returns a slice carrying the aliasing state established by its producer.117
*/118
async read(offset = 0, length = Number.MAX_SAFE_INTEGER, options?: SessionHandleReadOptions): Promise<SessionHandleReadResult> {119
// Closed-handle refusal precedes argument validation: a closed handle120
// rejects SessionHandleClosedError regardless of the arguments.121
this.assertOpen('read')122
if (!Number.isSafeInteger(offset) || offset < 0) {123
throw new TypeError(`read offset must be a non-negative safe integer, got ${String(offset)}`)124
}125
if (!Number.isSafeInteger(length) || length < 0) {126
throw new TypeError(`read length must be a non-negative safe integer, got ${String(length)}`)127
}128
options?.signal?.throwIfAborted()129
let result: SessionHandleReadResult130
const primed = this.state.primed131
if (primed !== undefined) {132
if (this.access === 'write') {133
result = this.readPrimed(primed, offset, length)134
} else {135
const currentPath = await this.storage.resolveCurrentLog(this.id, options?.signal)136
if (currentPath === undefined) {137
result = this.readPrimed(primed, offset, length)138
} else {139
this.state.primed = undefined140
result = await this.readCurrent(currentPath, offset, length, options?.signal)141
}142
}143
} else if (this.access === 'write' && !this.state.materialized) {144
result = { eventState: 'detached', events: [] }145
} else {146
const currentPath = await this.storage.resolveCurrentLog(this.id, options?.signal)147
if (currentPath !== undefined) {148
result = await this.readCurrent(currentPath, offset, length, options?.signal)149
} else if (this.storage.hasPendingSession(this.id)) {150
result = { eventState: 'detached', events: [] }151
} else {152
throw new SessionPersistenceNotFoundError(this.id)153
}154
}155
return result156
}158
/** Read one slice from the prepared historical prefix retained by this handle. */159
private readPrimed(source: SessionHandleReadResult, offset: number, length: number): SessionHandleReadResult {160
this.observedLength = Math.max(this.observedLength, source.events.length)161
return { eventState: source.eventState, events: source.events.slice(offset, offset + length) }162
}164
/** Read one current physical generation and enforce this handle's monotonic view. */165
private async readCurrent(166
path: string,167
offset: number,168
length: number,169
signal?: AbortSignal,170
): Promise<SessionHandleReadResult> {171
const source = await this.storage.readStoredLog(path, this.id, signal)172
if (source.events.length < this.observedLength) {173
throw new Error(`session "${this.id}": stored log shrank below a previously observed prefix (${source.events.length} < ${this.observedLength})`)174
}175
this.observedLength = source.events.length176
return {177
eventState: source.eventState,178
events: source.events.slice(offset, offset + length),179
}180
}182
/**183
* Durably append a contiguous batch; see the seam contract.184
* @param events - the contiguous batch in seq order.185
* @param options - optional cancellation observed before the write starts.186
*/187
async append(events: readonly SessionEvent[], options?: SessionHandleAppendOptions): Promise<void> {188
this.assertOpen('append')189
// Validate and deep-snapshot the batch HERE, before queueing behind the190
// chain, so the checked value is exactly the value persisted.191
const batch = materializeAppendBatch(events)192
return this.run('append', async () => {193
options?.signal?.throwIfAborted()194
await this.persistContiguous(batch)195
})196
}198
/**199
* Durability barrier; materializes the artifact when nothing has been200
* appended yet, so an explicitly flushed empty session survives this process.201
* @param options - optional cancellation observed before the barrier starts.202
*/203
flush(options?: SessionHandleFlushOptions): Promise<void> {204
return this.run('flush', async () => {205
options?.signal?.throwIfAborted()206
if (this.access !== 'write') throw new SessionReadOnlyError(this.id, 'flush')207
if (this.state.materialized) return // appends are durable on resolution208
await this.ensureLease()209
await this.storage.persistHeader(this.header, this.state.inheritedEventCount)210
this.state.materialized = true211
})212
}214
/**215
* Release the handle; see the seam contract. Idempotent and uncancellable.216
* A write handle first drains its routed live buffer through the still-open217
* storage, so backend teardown loses nothing regardless of which fiber218
* unwinds first; a drain or lock-release failure still frees the in-process219
* claim, then rejects — both failures together reject as one220
* `AggregateError`.221
* @returns settlement of the release.222
*/223
close(): Promise<void> {224
return this.closing ??= (async () => {225
let drainFailure: unknown226
// Producers on other fibers may still publish while close waits for227
// in-flight mutations (root disposal is concurrent), so drain again228
// until a full pass leaves the routed buffer empty. The chain never229
// rejects because run() swallows each operation's rejection after its230
// caller observed it.231
for (;;) {232
try {233
await this.drainLive()234
} catch (error: unknown) {235
drainFailure = error236
break237
}238
await this.chain239
if (this.buffered.length === 0) break240
}241
// After a drain failure the chain may still hold in-flight mutations.242
await this.chain243
// Free the in-process claim no matter how the kernel-lock release244
// fares: a skipped releaseHandle would wedge the id in this process245
// behind a lock the kernel may already have dropped.246
const failures: Error[] = []247
if (drainFailure !== undefined) {248
failures.push(drainFailure instanceof Error ? drainFailure : new Error(errorChain(drainFailure)))249
}250
try {251
await this.lease?.release()252
} catch (releaseFailure: unknown) {253
/* v8 ignore next -- lock releases reject with Error */254
failures.push(releaseFailure instanceof Error ? releaseFailure : new Error(errorChain(releaseFailure)))255
}256
this.storage.releaseHandle(this, this.state.materialized)257
if (failures.length > 1) throw new AggregateError(failures, `session "${this.id}": close failed to drain and to release its write lock`)258
if (failures[0] !== undefined) throw failures[0]259
})()260
}262
/** `await using` support: delegates to {@link close}. */263
[Symbol.asyncDispose](): Promise<void> {264
return this.close()265
}267
/**268
* Buffer one published live session event and arm the bounded batching269
* window when it is idle. The routing installer is the only caller.270
* @param event - the live event, retained as a persistence-owned copy.271
* @param reportBackgroundFailure - observes a deadline-driven drain failure272
* (the events stay buffered; the next {@link drainLive} retries loudly).273
*/274
enqueueLive(event: SessionEvent, reportBackgroundFailure: (error: unknown) => void): void {275
this.buffered.push(structuredClone(event))276
if (this.batchTimer !== undefined || this.drainPaused) return277
this.batchTimer = setTimeout(() => {278
this.batchTimer = undefined279
this.drainLive().catch(reportBackgroundFailure)280
}, LIVE_WRITE_BATCH_MAX_DELAY_MS)281
}283
/**284
* Durably drain the routed live buffer through the mutation chain;285
* concurrent callers join one drain, and a failure retains the batch in286
* order so `session/flush` can retry and reject loudly.287
*/288
drainLive(): Promise<void> {289
return this.draining ??= this.drainBuffered().finally(() => {290
this.draining = undefined291
})292
}294
private async drainBuffered(): Promise<void> {295
if (this.batchTimer !== undefined) {296
clearTimeout(this.batchTimer)297
this.batchTimer = undefined298
}299
this.drainPaused = false300
while (this.buffered.length > 0) {301
// Capture inside the chained turn so events landing while an earlier302
// batch writes coalesce into the next one, in order.303
await this.enqueueChain(async () => {304
// Only this single-flight drain splices the buffer, so the batch the305
// while-guard saw is still here when the chained turn runs.306
const batch = this.buffered.splice(0)307
try {308
await this.persistContiguous(materializeAppendBatch(batch))309
} catch (error: unknown) {310
this.buffered = batch.concat(this.buffered)311
this.drainPaused = true312
throw error313
}314
})315
}316
}318
/** The shared durable-append body: contiguity, ownership, torn-tail repair, storage write, state advance. */319
private async persistContiguous(batch: readonly SessionEvent[]): Promise<void> {320
if (this.access !== 'write') throw new SessionReadOnlyError(this.id, 'append')321
if (batch.length === 0) return322
await this.ensureLease()323
assertContiguous(this.id, batch, this.state.cursor)324
// Commit any pending torn-tail repair first, clearing each step's state325
// only once it lands so a failed step retries on the next mutation:326
// truncate the torn bytes, then durably rewrite the complete events327
// recovered from them (already counted in the primed cursor).328
if (this.state.tornTruncateTo !== undefined) {329
await this.storage.truncateTornTail(this.header, this.state.tornTruncateTo)330
this.state.tornTruncateTo = undefined331
}332
if (this.state.recoveredTail !== undefined) {333
if (this.state.recoveredTail.length > 0) {334
await this.storage.persistBatch(this.header, this.state.recoveredTail, this.state.materialized, this.state.inheritedEventCount)335
}336
this.state.recoveredTail = undefined337
}338
await this.storage.persistBatch(this.header, batch, this.state.materialized, this.state.inheritedEventCount)339
this.state.materialized = true340
this.state.cursor += batch.length341
this.state.primed = undefined342
this.observedLength = this.state.cursor343
}345
/**346
* Hold the cross-process write lock before this session's first durable347
* write. An open write handle holds it from construction; a create handle348
* acquires it here — immediately before the first log bytes publish — and349
* keeps it through close even when materialization then fails, so a350
* materializing session stays exclusively owned across retries.351
*/352
private async ensureLease(): Promise<void> {353
this.lease ??= await this.storage.acquireWriteLease(this.header)354
}356
/** Serialize one operation onto the chain without the closed-handle refusal (drain-from-close). */357
private enqueueChain(op: () => Promise<void>): Promise<void> {358
const next = this.chain.then(op)359
this.chain = next.catch(() => {})360
return next361
}363
/** Serialize one public mutating operation onto this handle's chain. */364
private async run(operation: string, op: () => Promise<void>): Promise<void> {365
this.assertOpen(operation)366
return this.enqueueChain(async () => {367
this.assertOpen(operation)368
return op()369
})370
}372
private assertOpen(operation: string): void {373
if (this.closing !== undefined) throw new SessionHandleClosedError(this.id, operation)374
}375
}377
/** One created-but-unmaterialized session tracked in this process only. */378
export interface PendingSession {379
readonly header: SessionHeader380
readonly revision: SessionPersistenceRevision381
/** Exact fork-inherited prefix length supplied at create. */382
readonly inheritedEventCount: SessionLogOffset383
}385
/**386
* The JSONL backend's in-process bookkeeping: the single active writer per387
* session id (doubling as the live event router), the open-handle set the388
* teardown sweep closes, and the created-but-unmaterialized sessions this389
* process can already observe.390
*/391
export class JsonlBackendTracker {392
/** Every open handle; teardown closes what remains. */393
readonly openHandles = new Set<SessionHandle>()394
/** `null` marks a claim whose handle is still being constructed. */395
private readonly writers = new Map<SessionId, JsonlSessionHandle | null>()396
private readonly pending = new Map<SessionId, PendingSession>()397
private counter = 0399
/** @param name - backend label used in in-memory revision tokens and teardown errors. */400
constructor(private readonly name: string) {}402
/**403
* Claim write ownership and record the created session as pending, making404
* it observable to this process before it materializes. Before405
* materialization this registration is the only guard — session ids do not406
* collide across processes, and no durable artifact exists for another407
* process to open; the handle takes the cross-process lock at its first408
* materializing write.409
* @param header - the validated detached header.410
* @param inheritedEventCount - the exact fork-inherited prefix length.411
* @throws {SessionAlreadyExistsError} when a concurrent create or an open412
* write handle holds the id — for create, the duplicate is the fact.413
*/414
registerCreated(header: SessionHeader, inheritedEventCount: SessionLogOffset): void {415
if (this.writers.has(header.id)) throw new SessionAlreadyExistsError(header.id)416
this.writers.set(header.id, null)417
this.pending.set(header.id, {418
header,419
revision: SessionPersistenceRevision(`memory:${this.name}:${++this.counter}`),420
inheritedEventCount,421
})422
}424
/**425
* Claim write ownership for an existing session.426
* @param id - the session to claim.427
* @throws {SessionAlreadyOwnedError} when an active write handle exists.428
*/429
claimWrite(id: SessionId): void {430
if (this.writers.has(id)) throw new SessionAlreadyOwnedError(id)431
this.writers.set(id, null)432
}434
/**435
* Roll a failed write open back.436
* @param id - the session whose claim is dropped.437
*/438
releaseClaim(id: SessionId): void {439
this.writers.delete(id)440
}442
/**443
* The pending entry for a created-but-unmaterialized session, if any.444
* @param id - the session to look up.445
* @returns the pending header and in-memory revision.446
*/447
pendingOf(id: SessionId): PendingSession | undefined {448
return this.pending.get(id)449
}451
/**452
* Whether this process still tracks a created-but-unmaterialized session.453
* @param id - the session to test.454
* @returns true while the pending entry exists.455
*/456
hasPending(id: SessionId): boolean {457
return this.pending.has(id)458
}460
/**461
* Iterate the pending sessions for listing.462
* @returns the pending entries, keyed by session id.463
*/464
pendingEntries(): IterableIterator<[SessionId, PendingSession]> {465
return this.pending.entries()466
}468
/**469
* Drop a pending entry once the session materialized durably.470
* @param id - the session that reached durable storage.471
*/472
materialized(id: SessionId): void {473
this.pending.delete(id)474
}476
/**477
* Track one open handle for teardown and, for a write handle, bind it as478
* the session's live event route.479
* @param handle - the just-constructed handle.480
* @returns the same handle, for construction-site chaining.481
*/482
adopt(handle: JsonlSessionHandle): JsonlSessionHandle {483
this.openHandles.add(handle)484
if (handle.access === 'write') this.writers.set(handle.id, handle)485
return handle486
}488
/**489
* Release one handle's bookkeeping on close. A write handle drops its490
* ownership claim; a creator that never materialized leaves nothing behind —491
* the session never existed.492
* @param handle - the closing handle.493
* @param materialized - whether the session reached durable storage.494
*/495
release(handle: JsonlSessionHandle, materialized: boolean): void {496
this.openHandles.delete(handle)497
if (handle.access !== 'write') return498
this.writers.delete(handle.id)499
if (!materialized) this.pending.delete(handle.id)500
}502
/**503
* Drain and flush every active write handle — the service-wide durability504
* barrier behind `SessionPersistence.flush`.505
* @throws {AggregateError} naming each session whose flush failed; the506
* remaining handles still flush.507
*/508
async flushAll(): Promise<void> {509
const errors: unknown[] = []510
for (const writer of [...this.writers.values()]) {511
if (writer === null) continue // a claim mid-construction routes nothing yet512
try {513
await writer.drainLive()514
await writer.flush()515
} catch (error: unknown) {516
// A handle closed during the sweep counts as flushed: close itself517
// drained the routed buffer durably before refusing this flush.518
if (error instanceof SessionHandleClosedError) continue519
errors.push(error)520
}521
}522
if (errors.length > 0) throw new AggregateError(errors, `${this.name} flush failed`)523
}525
/**526
* Install the backend's live session routing and teardown. Persistence527
* enforces one active write handle per id, so the listeners route published528
* sessions' events by id; the teardown effect closes every open handle —529
* close drains the routed buffer — and aggregates failures. This provider530
* owns no separate storage connection, so closing handles is the complete531
* teardown. Registrations are effects of the current fiber.532
* @param ctx - the backend's context.533
*/534
install(ctx: Context): void {535
ctx.on('session/event', (session: Session, event) => {536
this.writers.get(session.id)?.enqueueLive(event, (error) => {537
ctx.logger.warn(`session-persistence: background write for session "${session.id}" failed (buffered events retained): ${String(error)}`)538
})539
})540
ctx.on('session/flush', (session: Session) => {541
const writer = this.writers.get(session.id)542
if (writer === null || writer === undefined) return undefined543
return (async () => {544
await writer.drainLive()545
await writer.flush()546
})()547
})548
ctx.on('session/disposed', (session: Session) => {549
const writer = this.writers.get(session.id)550
if (writer === null || writer === undefined) return551
writer.close().catch((error: unknown) => {552
ctx.logger.warn(`session-persistence: final drain for session "${session.id}" failed: ${String(error)}`)553
})554
})555
ctx.effect(() => async () => {556
const errors: unknown[] = []557
for (const handle of [...this.openHandles]) {558
try {559
await handle.close()560
} catch (error: unknown) {561
errors.push(error)562
}563
}564
if (errors.length > 0) throw new AggregateError(errors, `${this.name} dispose failed`)565
}, `${this.name} open handles`)566
}567
}