index.ts 16 KB

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