index.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300
  1. /**
  2. * Persisted projection cache (`ctx.sessionProjectionCache`): durable
  3. * checkpoints of every registered projection unit's state, one record per
  4. * session on the domain data form (`session_projcache` domain — the shipped
  5. * json backend lands it beside `workspace.json`). The cache is a fold
  6. * shortcut, never an authority: a row is possibly stale (its `seq`
  7. * says how stale) but never wrong, so every write path is fail-soft (a lost
  8. * write costs a longer tail replay on the next cold read) and a
  9. * `ver` mismatch discards the row instead of migrating it. Design
  10. * authority: the session-projection RFC
  11. * (.agents/notes/proposed/architecture/2026-07-27-session-projection-and-command-log.md).
  12. * @module @deepseek-ai/dsh-session-projection-cache
  13. */
  14. import { Context, Service } from '@deepseek-ai/cordis'
  15. import z from '@deepseek-ai/schemastery'
  16. import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
  17. import type { Session, SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  18. // Empty type import: applies the package's cordis Context merge
  19. // (`ctx.sessionPersistence`), which this service reads on the cold path.
  20. import type {} from '@deepseek-ai/dsh-session-persistence'
  21. import type { ProjectionCheckpoint, ProjectionSnapshot } from '@deepseek-ai/dsh-session-projection'
  22. import type { KvTable } from '@deepseek-ai/dsh-storage-domain'
  23. import { projectionCacheDomainSpec } from './spec.ts'
  24. import type { CheckpointIdentity, CheckpointRecord } from './spec.ts'
  25. export { checkpointIdentity, checkpointRecord, checkpointRow, projectionCacheDomainSpec } from './spec.ts'
  26. export type { CheckpointIdentity, CheckpointRecord } from './spec.ts'
  27. declare module '@deepseek-ai/cordis' {
  28. interface Context {
  29. sessionProjectionCache: SessionProjectionCache
  30. }
  31. }
  32. /**
  33. * Plugin config. Both throttle triggers are deployment choices with no
  34. * universally correct value, so the composition states them explicitly
  35. * (cordis.yml); the two mandatory write points (`turn/end` and session
  36. * disposal) are policy, not tunables, and always fire.
  37. */
  38. export interface Config {
  39. /** Committed events per session that force a durable checkpoint write between mandatory points. */
  40. writeEveryEvents: number
  41. /** Longest time (milliseconds) a dirty checkpoint may stay unwritten between mandatory points. */
  42. writeIntervalMs: number
  43. }
  44. export const Config: z<Config> = z.object({
  45. writeEveryEvents: z.natural().min(1).required(),
  46. writeIntervalMs: z.natural().min(1).required(),
  47. })
  48. /** Per-session write-behind bookkeeping (live sessions only; dropped at retire). */
  49. interface DirtyState {
  50. /** Committed events since the last durable write. */
  51. pending: number
  52. /** Interval trigger armed at the first dirty event after a clean write. */
  53. timer: ReturnType<typeof setTimeout> | undefined
  54. }
  55. /**
  56. * The persisted projection cache service. Opens the `session_projcache`
  57. * domain at init, checkpoints live sessions on a throttled write-behind
  58. * (count/interval triggers from {@link Config}) plus two mandatory points —
  59. * `turn/end` and session disposal (the live-to-cold moment) — and serves the
  60. * cold-read ladder: cached row, persistence `readFrom` tail, registry
  61. * `restore`, durable write-back. Every durable write is fail-soft: failures
  62. * log a warning and the cache self-heals on the next write or cold read.
  63. */
  64. export class SessionProjectionCache extends Service {
  65. static inject = ['storageDomain', 'sessionProjections', 'sessionPersistence', 'sessions']
  66. static Config: z<Config> = Config
  67. private table?: KvTable<SessionId, CheckpointRecord>
  68. private readonly dirty = new Map<Session, DirtyState>()
  69. constructor(ctx: Context, public config: Config) {
  70. super(ctx, 'sessionProjectionCache')
  71. }
  72. /** Open the domain and install the write-behind listeners. */
  73. protected async [Service.init](): Promise<void> {
  74. const domain = await this.ctx.storageDomain.open(projectionCacheDomainSpec)
  75. this.ctx.effect(() => () => domain.close(), 'sessionProjectionCache.domainClose')
  76. this.table = domain.table('sessions')
  77. this.installWritePath()
  78. }
  79. /**
  80. * The stored record for one session, accepted only when its bound log
  81. * identity matches `expected`. A session id names a slot, not a lifecycle:
  82. * a recreated id or a persistence store swapped under a surviving cache
  83. * must not let an old record seed state folded from an unrelated log.
  84. * Synchronous from the domain's in-memory state.
  85. * @param id - the session whose record is read.
  86. * @param expected - the log identity the caller holds (live or stored header).
  87. * @returns the identity-matching record, or `undefined` (absent or unrelated).
  88. */
  89. private recordFor(id: SessionId, expected: CheckpointIdentity): CheckpointRecord | undefined {
  90. const record = this.requireTable().get(id)
  91. if (record === undefined) return undefined
  92. return identityMatches(record.identity, expected) ? record : undefined
  93. }
  94. /**
  95. * The zero-I/O listing read: whole values viewed straight from the stored
  96. * rows (version-matching keys only), each cut carried with its watermark
  97. * so a client value store can seed under its higher-seq-wins rule — as
  98. * stale as the last durable checkpoint but never wrong, and never from an
  99. * unrelated log (the caller's header is the identity witness). Fresher
  100. * paths (the history tail baseline, {@link coldSnapshot}) supersede these
  101. * values whenever a session is actually opened.
  102. * @param meta - the listed session's header (identity witness; no log read).
  103. * @returns the cut (`asOfSeq` = lowest served-row watermark), or
  104. * `undefined` when no usable row exists for this lifecycle.
  105. */
  106. cachedSnapshot(meta: SessionHeader): ProjectionSnapshot | undefined {
  107. const record = this.recordFor(meta.id, identityOf(meta))
  108. if (record === undefined) return undefined
  109. const values = this.ctx.sessionProjections.viewCheckpoint(record.rows)
  110. const keys = Object.keys(values)
  111. if (keys.length === 0) return undefined
  112. // The block carries ONE cut: the lowest served watermark is the seq every
  113. // value is at least current as of (under-claiming is safe under
  114. // higher-seq-wins; over-claiming would let a stale value outrank pushes).
  115. const asOfSeq = Math.min(...keys.map(key => (record.rows[key] as { seq: number }).seq))
  116. return { asOfSeq, values }
  117. }
  118. /**
  119. * Durably checkpoint one live session NOW (both mandatory points call
  120. * this; tests and carriers may too). The registry cut is snapshotted at
  121. * this boundary (states are live references), then the whole record is
  122. * replaced. NOT fail-soft — callers on the fail-soft paths contain it.
  123. * @param session - the live session to checkpoint.
  124. * @returns resolution after durability and event emission.
  125. */
  126. async write(session: Session): Promise<void> {
  127. const rows = this.ctx.sessionProjections.checkpoint(session)
  128. this.markClean(session)
  129. // Durability barrier: the checkpoint cut was taken above, so flushing
  130. // AFTER it guarantees every event inside the cut is durably logged
  131. // before the cache row lands — a crash can leave the cache behind the
  132. // log (longer tail replay) but never ahead of it (phantom values folded
  133. // from events no stored log contains). At detach the store entry is
  134. // already gone; persistence's own retirement drain covers that path and
  135. // any residual overreach is caught by the cold read's anchored floor.
  136. if (this.ctx.sessions.get(session.id) === session) await this.ctx.sessions.flush(session)
  137. await this.put(session.id, identityOf(session.header), rows)
  138. }
  139. /**
  140. * Cold-read one persisted session's projections with zero full-log load:
  141. * cached rows + a persistence `readFrom` tail from the registry's restore
  142. * floor, refolded by the registry and written back (fail-soft) so the next
  143. * cold read starts closer. A cache row invalidated by a shrunk log
  144. * (crash-repair truncation) triggers one full re-read from seq 0 — the
  145. * ladder's slow rung, still no crash. Rejects when the session has no
  146. * persisted log (`not found` from the persistence seam).
  147. * @param id - the persisted session to read.
  148. * @param signal - optional cancellation for the persistence reads.
  149. * @returns the snapshot cut at the stored log end.
  150. */
  151. async coldSnapshot(id: SessionId, signal?: AbortSignal): Promise<ProjectionSnapshot> {
  152. const record = this.requireTable().get(id)
  153. const cached = record?.rows ?? {}
  154. const floor = this.ctx.sessionProjections.restoreFloor(cached)
  155. const persistence = this.ctx.sessionPersistence
  156. if (floor === undefined) {
  157. // No unit registered: nothing to fold, but the not-found contract must
  158. // hold in this topology too — the probe read rejects for an absent log
  159. // and dates the empty cut for a present one.
  160. const probe = await persistence.readFrom(id, 0, signal)
  161. return { asOfSeq: probe.events.at(-1)?.seq ?? -1, values: {} }
  162. }
  163. let restored: { snapshot: ProjectionSnapshot; checkpoint: ProjectionCheckpoint }
  164. const tail = await persistence.readFrom(id, floor, signal)
  165. // The tail's stored header is the identity witness: a record bound to a
  166. // different lifecycle (recreated id, swapped store) is discarded whole
  167. // before any of its rows can seed a fold.
  168. const related = record === undefined || identityMatches(record.identity, identityOf(tail.meta))
  169. try {
  170. if (!related) throw new Error('unrelated log identity')
  171. restored = this.ctx.sessionProjections.restore(cached, tail.events, floor)
  172. } catch {
  173. // The recoverable restore failures: an unrelated record, or a row
  174. // overreaching the stored log end (or predating the floor). Both imply
  175. // floor > 0 (baseSeq-0 restores never throw and an unrelated record
  176. // still carried a usable watermark), so the full log is a fresh read.
  177. const whole = await persistence.readFrom(id, 0, signal)
  178. restored = this.ctx.sessionProjections.restore({}, whole.events, 0)
  179. }
  180. await this.putSoft(id, identityOf(tail.meta), restored.checkpoint, 'cold-read write-back')
  181. return restored.snapshot
  182. }
  183. // --- write-behind (throttle + mandatory points) ---
  184. private installWritePath(): void {
  185. // Every committed event advances the dirty counter; turn/end is a
  186. // mandatory point (the durable value most reads want is the turn-final
  187. // one), count/interval throttle the in-turn stream.
  188. this.ctx.on('session/event', (session: Session, event: SessionEvent) => {
  189. if (event.type === 'turn/end') {
  190. void this.flushSoft(session, 'turn/end')
  191. return
  192. }
  193. const state = this.dirty.get(session) ?? { pending: 0, timer: undefined }
  194. this.dirty.set(session, state)
  195. state.pending += 1
  196. if (state.pending >= this.config.writeEveryEvents) {
  197. void this.flushSoft(session, 'count threshold')
  198. return
  199. }
  200. state.timer ??= setTimeout(() => {
  201. void this.flushSoft(session, 'interval')
  202. }, this.config.writeIntervalMs)
  203. })
  204. // Detach (the live-to-cold moment): the second mandatory point. After
  205. // this write the cold-read ladder serves the session from the cache.
  206. // flushSoft's synchronous prefix reads and resets the dirty state, so
  207. // dropping it (timer already cleared by markClean) right after is safe.
  208. this.ctx.on('session/disposed', (session: Session) => {
  209. void this.flushSoft(session, 'detach')
  210. this.markClean(session)
  211. this.dirty.delete(session)
  212. })
  213. // Clear pending timers with the plugin (their sessions outlive the cache).
  214. this.ctx.effect(() => () => {
  215. for (const state of this.dirty.values()) {
  216. if (state.timer !== undefined) clearTimeout(state.timer)
  217. }
  218. this.dirty.clear()
  219. }, 'sessionProjectionCache.timers')
  220. }
  221. /**
  222. * One fail-soft durable checkpoint. Every caller has work by construction:
  223. * the throttle triggers only fire dirty (markClean clears the timer with
  224. * the counter) and the two mandatory points write unconditionally.
  225. */
  226. private async flushSoft(session: Session, trigger: string): Promise<void> {
  227. try {
  228. await this.write(session)
  229. } catch (error) {
  230. this.ctx.logger.warn(`session projection cache: ${trigger} write for "${session.id}" failed (cache stays stale): ${String(error)}`)
  231. }
  232. }
  233. /** Reset one session's dirty bookkeeping (its checkpoint is being written). */
  234. private markClean(session: Session): void {
  235. const state = this.dirty.get(session)
  236. if (state === undefined) return
  237. state.pending = 0
  238. if (state.timer !== undefined) {
  239. clearTimeout(state.timer)
  240. state.timer = undefined
  241. }
  242. }
  243. /** Replace one session's stored record with its log identity and a detached snapshot of `rows`. */
  244. private async put(id: SessionId, identity: CheckpointIdentity, rows: ProjectionCheckpoint): Promise<void> {
  245. const detached = snapshotJsonValue(rows)
  246. if (detached === undefined) {
  247. throw new TypeError('projection checkpoint is not losslessly JSON-serializable (a unit state violates the plain-JSON contract)')
  248. }
  249. await this.requireTable().put(id, { identity, rows: detached as CheckpointRecord['rows'] })
  250. }
  251. /** Fail-soft {@link put}: cache writes must never fail their caller's read or event path. */
  252. private async putSoft(id: SessionId, identity: CheckpointIdentity, rows: ProjectionCheckpoint, what: string): Promise<void> {
  253. try {
  254. await this.put(id, identity, rows)
  255. } catch (error) {
  256. this.ctx.logger.warn(`session projection cache: ${what} for "${id}" failed (cache stays stale): ${String(error)}`)
  257. }
  258. }
  259. private requireTable(): KvTable<SessionId, CheckpointRecord> {
  260. /* v8 ignore next -- Service.init assigns the table before the service becomes injectable */
  261. if (this.table === undefined) throw new Error('session projection cache is not initialized')
  262. return this.table
  263. }
  264. }
  265. /** Project a header onto the identity fields a record is bound to. */
  266. function identityOf(header: SessionHeader): CheckpointIdentity {
  267. return { createdAt: header.createdAt, ...header.cwd === undefined ? {} : { cwd: header.cwd } }
  268. }
  269. /** Whether a stored record's bound identity names the caller's lifecycle. */
  270. function identityMatches(stored: CheckpointIdentity, expected: CheckpointIdentity): boolean {
  271. return stored.createdAt === expected.createdAt && stored.cwd === expected.cwd
  272. }
  273. export default SessionProjectionCache