index.ts 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161
  1. /**
  2. * The durable session-persistence seam (`ctx.sessionPersistence`): an abstract
  3. * service defining WHAT a persistence backend does — durably store, reload,
  4. * and list sessions — without saying HOW. Implementations subclass
  5. * {@link SessionPersistence} and register themselves as the
  6. * `sessionPersistence` service; `@deepseek-ai/dsh-session-persistence-jsonl`
  7. * (an append-only JSONL log per session) is the first and
  8. * `@deepseek-ai/dsh-session-persistence-sqlite` (`node:sqlite`, one row per
  9. * event) is a second that validates the seam is backend-agnostic by passing
  10. * the same `runPersistenceContract` suite. Further backends swap in an object
  11. * store or a remote service without touching the consumers (the write-path
  12. * plugin, the agent-loop resume seam).
  13. *
  14. * The persisted unit IS the existing {@link SessionEvent} — there is no
  15. * parallel "persisted message" type the log must be converted to and from
  16. * (faithful to the event-sourced model: the log is the single source of
  17. * truth). Metadata that is NOT replayable conversation state (format version,
  18. * cwd, lineage, seed boundary) travels separately as {@link SessionHeader},
  19. * which is owned by `dsh-session` and re-exported here.
  20. *
  21. * @module @deepseek-ai/dsh-session-persistence
  22. */
  23. import { Context, Service } from 'cordis'
  24. import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
  25. import type { SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
  26. // Re-export the metadata vocabulary so consumers import it from the seam.
  27. export type { SessionHeader } from '@deepseek-ai/dsh-session'
  28. // The backend-agnostic write-path orchestration first-party backends compose.
  29. export { PersistenceCoordinator } from './coordinator.ts'
  30. export type { PersistenceBackend, StoredPrefix } from './coordinator.ts'
  31. declare module 'cordis' {
  32. interface Context {
  33. sessionPersistence: SessionPersistence
  34. }
  35. }
  36. /**
  37. * Whether a live session's seed reproduces a persisted prefix exactly. Backends
  38. * use this collision check to distinguish a legitimate resume/HMR rebind from a
  39. * different live session reusing an existing session id.
  40. *
  41. * The comparison includes the full event payload, not just seq/type/time, so a
  42. * mutated seed cannot be grafted onto a durable log with the same envelope.
  43. * @param seed - the live session's creation-time event snapshot.
  44. * @param prefix - the persisted prefix the seed must reproduce.
  45. * @returns `true` when the prefix fits within the seed and every event matches by JSON text.
  46. */
  47. export function seedCoversPrefix(seed: readonly SessionEvent[], prefix: readonly SessionEvent[]): boolean {
  48. return prefix.length <= seed.length
  49. && prefix.every((event, index) => {
  50. const seedEvent = seed[index]
  51. return seedEvent !== undefined && JSON.stringify(seedEvent) === JSON.stringify(event)
  52. })
  53. }
  54. /**
  55. * Reject a batch that is not wholly losslessly JSON-serializable. Live session
  56. * appends already enforce this; persistence append paths also accept replay or
  57. * direct batches that may bypass a live session instance. Validation uses the
  58. * same one-pass materializer as the coordinator, so getters are read once.
  59. * @param events - the complete event batch to validate.
  60. */
  61. export function assertSerializable(events: readonly SessionEvent[]): void {
  62. const snapshot = snapshotJsonValue(events)
  63. if (snapshot === undefined) {
  64. throw new Error('session event batch is not losslessly JSON-serializable because it contains non-JSON-serializable data')
  65. }
  66. }
  67. /**
  68. * Abstract durable session-persistence service. Subclass, implement the
  69. * abstract methods, and load the subclass as a plugin — it registers as
  70. * `ctx.sessionPersistence` (one implementation per context; loading a second
  71. * throws, cordis' standard duplicate-service behavior).
  72. *
  73. * Contracts every implementation MUST honor (a DB backend asserts them inside
  74. * a transaction; a file backend appends at EOF):
  75. *
  76. * - **Append-only; a crashed turn is closed, not truncated.** Committed events
  77. * — those at or below a flushed `turn/end` — are never rewritten. A crash can
  78. * leave an unclosed final turn whose events are real (and possibly large);
  79. * {@link load} preserves them and closes the orphaned turn with synthetic
  80. * boundary events (see {@link load}). Only a never-fully-written torn tail
  81. * fragment is discarded.
  82. * - **Contiguous seq.** A persisted log is contiguous: `events[i].seq === i`.
  83. * {@link load} rejects a parse error or a `seq` gap in the COMMITTED region
  84. * (unloadable); {@link append}'s first event `seq` MUST equal the backend's
  85. * stored next-seq (after `load` has balanced any interrupted turn).
  86. * - **JSON-serializable events.** `SessionEventMap` is merge-extensible, so
  87. * {@link append} materializes each complete batch through the shared
  88. * lossless-JSON boundary before buffering it. The public `session.events`
  89. * view is immutable, but persistence still snapshots direct/replay callers at
  90. * this independent trust boundary.
  91. * - **Durability.** {@link append} returns only once the batch is durable
  92. * (the file backend fsyncs; a DB commits). {@link create} MAY defer the
  93. * physical write until the first {@link append} (lazy materialization).
  94. */
  95. export abstract class SessionPersistence extends Service {
  96. constructor(ctx: Context) {
  97. super(ctx, 'sessionPersistence')
  98. }
  99. /**
  100. * Register a new session's metadata. A backend MAY defer the physical write
  101. * until the first {@link append} (lazy materialization), in which case a
  102. * created-but-never-appended session is absent from {@link list}
  103. * — abandoned sessions leave nothing behind.
  104. * @param meta - the immutable header (id, version, cwd, lineage) to record.
  105. */
  106. abstract create(meta: SessionHeader): Promise<void>
  107. /**
  108. * Durably persist a batch of events (called from the write-behind drain at
  109. * the `session/flush` checkpoint). Honors the append-only and contiguous-seq
  110. * contracts: the first event's `seq` MUST equal the stored next-seq (after
  111. * `load` has durably closed any interrupted turn). Rejects non-JSON-
  112. * serializable `event.data` with an error naming the offending event type.
  113. * @param id - the session the batch belongs to.
  114. * @param events - the contiguous batch to persist, in seq order.
  115. */
  116. abstract append(id: SessionId, events: readonly SessionEvent[]): Promise<void>
  117. /**
  118. * Reload a session: its {@link SessionHeader} plus the event log up to the last
  119. * durable checkpoint. Returns `meta` AND `events` so the live session is
  120. * reconstructed with its `cwd`/lineage, not just its log.
  121. *
  122. * The loop only flushes at `turn/end`, so a crash can leave a durable log
  123. * whose final turn never closed: real, fully-written events sit after the last
  124. * `turn/end`. Those events are PRESERVED — a single turn can be huge in a
  125. * long-horizon task, so truncating it would destroy real work — and `load`
  126. * CLOSES the orphaned turn by durably appending the minimal synthetic boundary
  127. * events: an error `tool/result` for every `tool-call` the crash left
  128. * unanswered (so the rehydrated history is a valid provider transcript — a
  129. * dangling assistant tool-call is otherwise rejected), then a `step/end` if a
  130. * step was open, then a `turn/end` carrying the `{ kind: 'interrupted' }`
  131. * reason. The returned `events` therefore end on a balanced `turn/end` and are
  132. * immediately usable as a session seed. Only a never-fully-written TORN tail
  133. * fragment (a half-written final record) is discarded. Returned events are
  134. * contiguous (`events[i].seq === i`); a parse error or a `seq` gap in the
  135. * COMMITTED region (at or before the last real `turn/end`) makes the session
  136. * unloadable (reject). Rejects an unknown format `version`. See the session-persistence RFC for
  137. * the crash-recovery contract.
  138. * @param id - the persisted session to reload.
  139. * @returns the header plus the event log, ending on a balanced `turn/end` —
  140. * immediately usable as a session seed.
  141. */
  142. abstract load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }>
  143. /**
  144. * Lightweight listing from metadata, without a full-log parse.
  145. * @returns one header per materialized session.
  146. */
  147. abstract list(): Promise<SessionHeader[]>
  148. }
  149. export default SessionPersistence