返回源码地图

packages/session/session-persistence-jsonl/src/storage.ts

main snapshot · da00f7f5358f · 正文引用章节 16 / 16;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * The JSONL provider's session storage runtime: its concrete write/read
3 * handle with a per-handle mutation chain and a routed live write-behind
4 * buffer, the in-process bookkeeping that enforces one active writer per
5 * session id, and the backend's live event routing and teardown. Deliberately
6 * provider-local: the persistence seam exposes only the service and handle
7 * contracts, and the shared contract suites pin equivalent observable
8 * behavior across providers.
9 * @module
10 */
11
12import type { Context } from '@deepseek-ai/cordis'
13import { errorChain } from '@deepseek-ai/dsh-llm'
14import type { Session, SessionEvent, SessionHeader, SessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'
15import {
16 assertContiguous,
17 SessionAlreadyExistsError,
18 materializeAppendBatch,
19 SessionAlreadyOwnedError,
20 SessionHandleClosedError,
21 SessionPersistenceNotFoundError,
22 SessionPersistenceRevision,
23 SessionReadOnlyError,
24} from '@deepseek-ai/dsh-session-persistence'
25import type {
26 SessionAccess,
27 SessionHandle,
28 SessionHandleAppendOptions,
29 SessionHandleFlushOptions,
30 SessionHandleReadOptions,
31 SessionHandleReadResult,
32} from '@deepseek-ai/dsh-session-persistence'
33import type { SessionWriteLease } from './lease.ts'
34
35/** Maximum intentional wait before a routed live session batch starts writing. */
36export const LIVE_WRITE_BATCH_MAX_DELAY_MS = 200
37
38/** The file-storage primitives the handle drives on its owning service. */
39export 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): boolean
57 /** 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): void
61}
62
63/** Mutable per-handle log state; a write handle is its session's single mutator. */
64export interface StorageHandleState {
65 /** The stored next-seq (the logical end this handle knows). */
66 cursor: number
67 /** Whether the session has a durable artifact yet. */
68 materialized: boolean
69 /** Torn-tail truncation point, consumed by the first new append. */
70 tornTruncateTo?: number | undefined
71 /** Complete events recovered from the torn final frame; the first mutation rewrites them durably. */
72 recoveredTail?: SessionEvent[] | undefined
73 /** Exact fork-inherited prefix length stored with the log; `0` when unseeded. */
74 inheritedEventCount: SessionLogOffset
75 /** The validated stored prefix from a write open, served to reads until the first append. */
76 primed?: SessionHandleReadResult | undefined
77}
78
79/**
80 * The JSONL session handle. Mutations serialize on a per-handle promise
81 * chain; reads re-scan the artifact on demand and never observe a shorter log
82 * than a prior read on this handle. Routed live events buffer in a bounded
83 * window and drain through the same chain as explicit appends.
84 */
85export class JsonlSessionHandle implements SessionHandle {
86 private chain: Promise<unknown> = Promise.resolve()
87 private closing: Promise<void> | undefined
88 private observedLength = 0
89 /** Routed live events awaiting their batching deadline (persistence-owned copies). */
90 private buffered: SessionEvent[] = []
91 private batchTimer: ReturnType<typeof setTimeout> | undefined
92 /** Set when a drain failed; the automatic timer stays quiet until the next drain. */
93 private drainPaused = false
94 private draining: Promise<void> | undefined
95
96 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 ) {}
105
106 /** Exact fork-inherited prefix length stored with this session's log. */
107 get inheritedEventCount(): SessionLogOffset {
108 return this.state.inheritedEventCount
109 }
110
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 handle
120 // 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: SessionHandleReadResult
130 const primed = this.state.primed
131 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 = undefined
140 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 result
156 }
157
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 }
163
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.length
176 return {
177 eventState: source.eventState,
178 events: source.events.slice(offset, offset + length),
179 }
180 }
181
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 the
190 // 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 }
197
198 /**
199 * Durability barrier; materializes the artifact when nothing has been
200 * 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 resolution
208 await this.ensureLease()
209 await this.storage.persistHeader(this.header, this.state.inheritedEventCount)
210 this.state.materialized = true
211 })
212 }
213
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-open
217 * storage, so backend teardown loses nothing regardless of which fiber
218 * unwinds first; a drain or lock-release failure still frees the in-process
219 * claim, then rejects — both failures together reject as one
220 * `AggregateError`.
221 * @returns settlement of the release.
222 */
223 close(): Promise<void> {
224 return this.closing ??= (async () => {
225 let drainFailure: unknown
226 // Producers on other fibers may still publish while close waits for
227 // in-flight mutations (root disposal is concurrent), so drain again
228 // until a full pass leaves the routed buffer empty. The chain never
229 // rejects because run() swallows each operation's rejection after its
230 // caller observed it.
231 for (;;) {
232 try {
233 await this.drainLive()
234 } catch (error: unknown) {
235 drainFailure = error
236 break
237 }
238 await this.chain
239 if (this.buffered.length === 0) break
240 }
241 // After a drain failure the chain may still hold in-flight mutations.
242 await this.chain
243 // Free the in-process claim no matter how the kernel-lock release
244 // fares: a skipped releaseHandle would wedge the id in this process
245 // 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 }
261
262 /** `await using` support: delegates to {@link close}. */
263 [Symbol.asyncDispose](): Promise<void> {
264 return this.close()
265 }
266
267 /**
268 * Buffer one published live session event and arm the bounded batching
269 * 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 failure
272 * (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) return
277 this.batchTimer = setTimeout(() => {
278 this.batchTimer = undefined
279 this.drainLive().catch(reportBackgroundFailure)
280 }, LIVE_WRITE_BATCH_MAX_DELAY_MS)
281 }
282
283 /**
284 * Durably drain the routed live buffer through the mutation chain;
285 * concurrent callers join one drain, and a failure retains the batch in
286 * order so `session/flush` can retry and reject loudly.
287 */
288 drainLive(): Promise<void> {
289 return this.draining ??= this.drainBuffered().finally(() => {
290 this.draining = undefined
291 })
292 }
293
294 private async drainBuffered(): Promise<void> {
295 if (this.batchTimer !== undefined) {
296 clearTimeout(this.batchTimer)
297 this.batchTimer = undefined
298 }
299 this.drainPaused = false
300 while (this.buffered.length > 0) {
301 // Capture inside the chained turn so events landing while an earlier
302 // 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 the
305 // 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 = true
312 throw error
313 }
314 })
315 }
316 }
317
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) return
322 await this.ensureLease()
323 assertContiguous(this.id, batch, this.state.cursor)
324 // Commit any pending torn-tail repair first, clearing each step's state
325 // only once it lands so a failed step retries on the next mutation:
326 // truncate the torn bytes, then durably rewrite the complete events
327 // 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 = undefined
331 }
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 = undefined
337 }
338 await this.storage.persistBatch(this.header, batch, this.state.materialized, this.state.inheritedEventCount)
339 this.state.materialized = true
340 this.state.cursor += batch.length
341 this.state.primed = undefined
342 this.observedLength = this.state.cursor
343 }
344
345 /**
346 * Hold the cross-process write lock before this session's first durable
347 * write. An open write handle holds it from construction; a create handle
348 * acquires it here — immediately before the first log bytes publish — and
349 * keeps it through close even when materialization then fails, so a
350 * materializing session stays exclusively owned across retries.
351 */
352 private async ensureLease(): Promise<void> {
353 this.lease ??= await this.storage.acquireWriteLease(this.header)
354 }
355
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 next
361 }
362
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 }
371
372 private assertOpen(operation: string): void {
373 if (this.closing !== undefined) throw new SessionHandleClosedError(this.id, operation)
374 }
375}
376
377/** One created-but-unmaterialized session tracked in this process only. */
378export interface PendingSession {
379 readonly header: SessionHeader
380 readonly revision: SessionPersistenceRevision
381 /** Exact fork-inherited prefix length supplied at create. */
382 readonly inheritedEventCount: SessionLogOffset
383}
384
385/**
386 * The JSONL backend's in-process bookkeeping: the single active writer per
387 * session id (doubling as the live event router), the open-handle set the
388 * teardown sweep closes, and the created-but-unmaterialized sessions this
389 * process can already observe.
390 */
391export 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 = 0
398
399 /** @param name - backend label used in in-memory revision tokens and teardown errors. */
400 constructor(private readonly name: string) {}
401
402 /**
403 * Claim write ownership and record the created session as pending, making
404 * it observable to this process before it materializes. Before
405 * materialization this registration is the only guard — session ids do not
406 * collide across processes, and no durable artifact exists for another
407 * process to open; the handle takes the cross-process lock at its first
408 * 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 open
412 * 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 }
423
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 }
433
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 }
441
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 }
450
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 }
459
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 }
467
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 }
475
476 /**
477 * Track one open handle for teardown and, for a write handle, bind it as
478 * 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 handle
486 }
487
488 /**
489 * Release one handle's bookkeeping on close. A write handle drops its
490 * 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') return
498 this.writers.delete(handle.id)
499 if (!materialized) this.pending.delete(handle.id)
500 }
501
502 /**
503 * Drain and flush every active write handle — the service-wide durability
504 * barrier behind `SessionPersistence.flush`.
505 * @throws {AggregateError} naming each session whose flush failed; the
506 * 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 yet
512 try {
513 await writer.drainLive()
514 await writer.flush()
515 } catch (error: unknown) {
516 // A handle closed during the sweep counts as flushed: close itself
517 // drained the routed buffer durably before refusing this flush.
518 if (error instanceof SessionHandleClosedError) continue
519 errors.push(error)
520 }
521 }
522 if (errors.length > 0) throw new AggregateError(errors, `${this.name} flush failed`)
523 }
524
525 /**
526 * Install the backend's live session routing and teardown. Persistence
527 * enforces one active write handle per id, so the listeners route published
528 * sessions' events by id; the teardown effect closes every open handle —
529 * close drains the routed buffer — and aggregates failures. This provider
530 * owns no separate storage connection, so closing handles is the complete
531 * 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 undefined
543 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) return
551 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}