index.ts 65 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647
  1. /**
  2. * JSONL durable session-persistence backend. It stores a header and contiguous
  3. * events in immutable generation files under one directory per session and serves the handle-based
  4. * `SessionPersistence` API: `create`/`open` return per-session handles, and
  5. * every read validates the same fail-closed storage contract.
  6. * @module @deepseek-ai/dsh-session-persistence-jsonl
  7. */
  8. import { Context } from '@deepseek-ai/cordis'
  9. import z from '@deepseek-ai/schemastery'
  10. import {
  11. SessionFormatUnsupportedMigrationError,
  12. sessionFormatCatalog,
  13. } from '@deepseek-ai/dsh-session-format-catalog'
  14. import { readdirSync, type Dirent } from 'node:fs'
  15. import { open, mkdir, readdir, realpath, link, rm, stat, truncate } from 'node:fs/promises'
  16. import { dirname, join, resolve } from 'node:path'
  17. import { performance } from 'node:perf_hooks'
  18. import { scheduler } from 'node:timers/promises'
  19. import { randomBytes } from 'node:crypto'
  20. import {
  21. SessionPersistence, SessionPersistenceRevision, SessionFormatUnsupportedError,
  22. SessionPersistenceCorruptionError,
  23. SessionAlreadyExistsError, SessionPersistenceNotFoundError,
  24. assertStoredId, materializeCreateHeader, sessionFormatVersionRefusal, validateStoredEvents,
  25. type SessionAccess, type SessionHandle,
  26. type SessionHandleReadResult,
  27. type SessionLocation, type SessionPersistenceCreateOptions,
  28. type SessionPersistenceListOptions, type SessionPersistenceOpenOptions,
  29. type SessionPersistenceSnapshot, type SessionPersistenceStatOptions,
  30. type SessionPersistenceRevision as PersistenceRevision,
  31. } from '@deepseek-ai/dsh-session-persistence'
  32. import { JsonlBackendTracker, JsonlSessionHandle, type StorageHandleState } from './storage.ts'
  33. import { SessionWriteLease } from './lease.ts'
  34. import { SESSION_FORMAT_VERSION, SessionId as makeSessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'
  35. import type { SessionEvent, SessionId, SessionHeader, SessionLogOffset as SessionLogOffsetType } from '@deepseek-ai/dsh-session'
  36. import {
  37. assertNoRetiredHeaderFields, encodeSegment, eventLines, generationLogFilename, generationLogPath, logPath, logSuffix,
  38. parseGenerationLogFilename, projectDir, scanLog, sessionDir, SessionLogScanner, toHeaderLine,
  39. type JsonlCompression,
  40. } from './format.ts'
  41. import {
  42. compressZstdFrame, createZstdFrameDecoder, decompressZstdFrame, decompressZstdPrefix, scanZstdFrames,
  43. } from './zstd.ts'
  44. import { ensureDurableDirectoryWin32, publishNewFileWin32 } from './win32.ts'
  45. import { verifyCurrentGenerationInWorker } from './migration-verifier.ts'
  46. import {
  47. JsonlGenerationSourceChangedError,
  48. JsonlGenerationUnsupportedMigrationError,
  49. prepareJsonlMigration,
  50. readStableJsonlFile,
  51. type JsonlGenerationFormatAdapter,
  52. type JsonlPhysicalIdentity,
  53. type PreparedJsonlMigration,
  54. } from './generation.ts'
  55. export type { JsonlCompression } from './format.ts'
  56. /**
  57. * Internal handoff-reuse policy, not deployment configuration: a cold
  58. * observation and the resume that immediately follows it reuse one parsed
  59. * log, so the memo only needs the sessions in flight between those steps.
  60. */
  61. const COLD_LOG_MEMO_MAX_ENTRIES = 2
  62. const DEFAULT_COMPRESSION: JsonlCompression = 'zstd'
  63. /**
  64. * Internal scheduling constant, not deployment configuration: balance
  65. * frame-boundary event-loop yields against `setImmediate` overhead. One frame
  66. * remains an indivisible synchronous decode.
  67. */
  68. const ZSTD_DECODE_YIELD_INTERVAL_MS = 500
  69. /** Assert that the independently decodable first frame contains only the header record. */
  70. function assertZstdHeaderFrame(plaintext: Buffer): void {
  71. if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {
  72. throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
  73. }
  74. }
  75. /** Loader schema for the JSONL artifact's physical encoding. */
  76. export const JsonlCompressionSchema: z<JsonlCompression> = z.union([
  77. z.const('zstd'),
  78. z.const('none'),
  79. ]).default(DEFAULT_COMPRESSION)
  80. /** Plugin config for the JSONL backend's root and physical encoding. */
  81. export interface Config {
  82. /**
  83. * Root directory for all session files. Required (no default): a default of
  84. * `process.cwd()` would scatter session files as the process's cwd changes
  85. * (bash calls, subprocesses). Sessions group under human-readable project
  86. * directories, then per-session directories. An existing root must be a
  87. * readable directory; an absent root is created on first materialization.
  88. */
  89. root: string
  90. /** Physical encoding; defaults to checksummed Zstandard frames. */
  91. compression?: JsonlCompression
  92. }
  93. /** One stored event graph whose producer has established immutable sharing. */
  94. interface FrozenStoredEvents extends SessionHandleReadResult {
  95. readonly eventState: 'shared-frozen'
  96. }
  97. /** State shared by prepared historical and published current logs. */
  98. interface StoredLogBase extends FrozenStoredEvents {
  99. readonly meta: SessionHeader
  100. readonly tornTruncateTo: number | undefined
  101. /** Complete events recovered from the torn final frame; the write path rewrites them durably. */
  102. readonly recoveredTail: SessionEvent[]
  103. /** Exact fork-inherited prefix length stored in the header line. */
  104. readonly inheritedEventCount: SessionLogOffsetType
  105. readonly revision: PersistenceRevision
  106. }
  107. /** A decoded current generation that is already durable. */
  108. interface CurrentStoredLog extends StoredLogBase {
  109. readonly status: 'current'
  110. }
  111. /** A migrated historical generation retained until an explicit write open publishes it. */
  112. interface PreparedStoredLog extends StoredLogBase {
  113. readonly status: 'prepared'
  114. readonly publication: {
  115. readonly source: ResolvedJsonlGeneration
  116. readonly value: PreparedJsonlMigration
  117. }
  118. }
  119. /** A validated logical log, either durable current state or prepared historical state. */
  120. type StoredLog = CurrentStoredLog | PreparedStoredLog
  121. /** Deep-freeze one acyclic stored JSON event without recursive calls. */
  122. function freezeStoredEvent(event: SessionEvent): void {
  123. const pending: object[] = [event]
  124. while (pending.length > 0) {
  125. // The non-empty check proves an object remains to visit.
  126. // oxlint-disable-next-line typescript/no-non-null-assertion
  127. const current = pending.pop()!
  128. Object.freeze(current)
  129. for (const key in current) {
  130. const child = (current as Record<string, unknown>)[key]
  131. if (child !== null && typeof child === 'object') pending.push(child)
  132. }
  133. }
  134. }
  135. /** Establish immutable sharing for one decoded event graph and report that state. */
  136. function freezeStoredEvents(events: SessionEvent[]): FrozenStoredEvents {
  137. for (const event of events) freezeStoredEvent(event)
  138. Object.freeze(events)
  139. return { eventState: 'shared-frozen', events }
  140. }
  141. /** One authoritative immutable generation selected from a Session directory. */
  142. interface ResolvedJsonlGeneration {
  143. readonly sourcePath: string
  144. readonly sourceVersion: number
  145. readonly currentPath: string
  146. }
  147. /** One backend-owned historical preparation shared by its current callers. */
  148. interface MigrationPreparation {
  149. readonly sourcePath: string
  150. readonly sourceRevision: PersistenceRevision
  151. readonly controller: AbortController
  152. readonly promise: Promise<PreparedStoredLog>
  153. settled: boolean
  154. waiters: number
  155. }
  156. /** Build the stat-derived best-effort change token shared by full and lightweight reads. */
  157. function fileRevision(identity: JsonlPhysicalIdentity): PersistenceRevision {
  158. return SessionPersistenceRevision([
  159. identity.dev,
  160. identity.ino,
  161. identity.size,
  162. identity.mtimeNs,
  163. identity.ctimeNs,
  164. ].join(':'))
  165. }
  166. /** Whether a filesystem error means absence; every non-ENOENT failure must surface. */
  167. function isENOENT(error: unknown): boolean {
  168. return (error as NodeJS.ErrnoException | null)?.code === 'ENOENT'
  169. }
  170. /** Whether a filesystem-owned failure should retain its original errno and path. */
  171. function isErrnoException(error: unknown): error is NodeJS.ErrnoException {
  172. return typeof (error as NodeJS.ErrnoException | null)?.code === 'string'
  173. }
  174. /** Preserve an Error abort reason and normalize hostile non-Error reasons. */
  175. function abortError(signal: AbortSignal): Error {
  176. return signal.reason instanceof Error
  177. ? signal.reason
  178. : new Error('session migration preparation aborted', { cause: signal.reason })
  179. }
  180. /** Let one caller stop waiting without transferring cancellation ownership to shared work. */
  181. function waitWithAbort<T>(operation: Promise<T>, signal?: AbortSignal): Promise<T> {
  182. if (signal === undefined) return operation
  183. /* v8 ignore next -- requireStoredLog synchronously rechecks the signal immediately before waiting. */
  184. if (signal.aborted) return Promise.reject(abortError(signal))
  185. return new Promise<T>((resolve, reject) => {
  186. const stopWaiting = (): void => {
  187. reject(abortError(signal))
  188. }
  189. signal.addEventListener('abort', stopWaiting, { once: true })
  190. void operation.then(
  191. (value) => {
  192. signal.removeEventListener('abort', stopWaiting)
  193. resolve(value)
  194. },
  195. (error: unknown) => {
  196. signal.removeEventListener('abort', stopWaiting)
  197. /* v8 ignore else -- the preparation owner normalizes every rejection before this waiter sees it. */
  198. if (error instanceof Error) {
  199. reject(error)
  200. } else {
  201. reject(new Error('session migration preparation failed', { cause: error }))
  202. }
  203. },
  204. )
  205. })
  206. }
  207. /**
  208. * The JSONL persistence backend. Load as a plugin; it registers as
  209. * `ctx.sessionPersistence`. Sessions materialize lazily: a created session is
  210. * visible to this process immediately, reaches disk on its first append or
  211. * flush, and never existed if the process crashes before that.
  212. */
  213. class JsonlSessionPersistence extends SessionPersistence {
  214. static Config: z<Config> = z.object({
  215. root: z.string().required(),
  216. compression: JsonlCompressionSchema,
  217. })
  218. /** Backend label for diagnostics and effects; shadows `Service.name` without changing the service key. */
  219. override readonly name = 'session-persistence-jsonl'
  220. private root: string
  221. private compression: JsonlCompression
  222. private rootEncodingCheck: Promise<void> | undefined
  223. private readonly tracker = new JsonlBackendTracker(this.name)
  224. private readonly generationFormat: JsonlGenerationFormatAdapter
  225. /**
  226. * Bounded LRU of parsed, validated stored logs keyed by session id and
  227. * guarded by the stat-derived revision, so an immediate cold-read handoff
  228. * (observation then resume) parses the artifact once. Every local mutation
  229. * for an id invalidates its entry; a foreign write misses through the
  230. * revision guard.
  231. */
  232. private readonly coldLogMemo = new Map<SessionId, StoredLog>()
  233. /** One joinable decode/migration operation per selected historical Session file revision. */
  234. private readonly migrationPreparations = new Map<SessionId, MigrationPreparation>()
  235. constructor(ctx: Context, public config: Config) {
  236. super(ctx)
  237. /* v8 ignore next 5 -- generated catalog and Session source share one build-time version owner. */
  238. if (sessionFormatCatalog.currentVersion !== SESSION_FORMAT_VERSION) {
  239. throw new Error(
  240. `session-persistence-jsonl: format catalog v${sessionFormatCatalog.currentVersion} `
  241. + `does not match Session v${SESSION_FORMAT_VERSION}`,
  242. )
  243. }
  244. // Resolve once so later process.cwd() changes cannot split one backend across roots.
  245. this.root = resolve(config.root)
  246. this.compression = config.compression ?? DEFAULT_COMPRESSION
  247. this.generationFormat = {
  248. currentVersion: sessionFormatCatalog.currentVersion,
  249. createRestore: header => sessionFormatCatalog.createRestore(header, {
  250. recovery: 'recoverable',
  251. validation: 'transformed',
  252. }),
  253. encodeHeader: (header, inheritedEventCount) =>
  254. sessionFormatCatalog.encodeCurrentHeader(header, inheritedEventCount),
  255. encodeEvent: event => sessionFormatCatalog.encodeCurrentEvent(event),
  256. isUnsupportedMigrationError: (error): error is SessionFormatUnsupportedMigrationError =>
  257. error instanceof SessionFormatUnsupportedMigrationError,
  258. }
  259. this.assertUsableRoot()
  260. this.tracker.install(ctx)
  261. }
  262. /**
  263. * Refusal-diagnostics hook: the absolute target path, without touching the filesystem.
  264. * @param meta - the stored header naming the session and its cwd.
  265. * @returns the artifact kind and absolute path.
  266. */
  267. private locate(meta: SessionHeader): SessionLocation {
  268. return { kind: 'jsonl', path: logPath(this.root, meta.cwd, meta.id, this.compression) }
  269. }
  270. // --- SessionPersistence service API ---
  271. /**
  272. * Create a new stored session and take its write ownership. The session is
  273. * visible to this process immediately; the physical artifact appears on the
  274. * first append or flush.
  275. * @param header - the immutable header to store; must be losslessly
  276. * JSON-serializable with a non-negative safe-integer `createdAt`.
  277. * @param options - optional cancellation.
  278. * @returns the owned write handle.
  279. */
  280. async create(header: SessionHeader, options?: SessionPersistenceCreateOptions): Promise<SessionHandle> {
  281. options?.signal?.throwIfAborted()
  282. const snapshot = materializeCreateHeader(header)
  283. // Fail fast on a seeded/cut mismatch with the exact refusal the header
  284. // line encoder enforces at materialization.
  285. toHeaderLine(snapshot, options?.inheritedEventCount)
  286. const inheritedEventCount = SessionLogOffset(options?.inheritedEventCount ?? 0)
  287. await this.ensureRootEncoding()
  288. options?.signal?.throwIfAborted()
  289. if (this.tracker.hasPending(snapshot.id) || await this.findLog(snapshot.id, options?.signal) !== undefined) {
  290. throw new SessionAlreadyExistsError(snapshot.id)
  291. }
  292. options?.signal?.throwIfAborted()
  293. // No lock yet: before materialization there is no durable artifact for
  294. // another process to contend over, so the handle acquires the lock right
  295. // before its first log bytes publish (ensureLease); an unmaterialized
  296. // session leaves no filesystem footprint at all.
  297. this.tracker.registerCreated(snapshot, inheritedEventCount)
  298. return this.tracker.adopt(new JsonlSessionHandle(this, snapshot.id, snapshot, 'write', { cursor: 0, materialized: false, inheritedEventCount }))
  299. }
  300. /**
  301. * Open an existing stored session for `read` or single-writer `write`.
  302. * @param id - the stored session to open.
  303. * @param access - `read` (no ownership) or `write` (atomic in-process claim).
  304. * @param options - optional cancellation.
  305. * @returns the open handle.
  306. */
  307. async open(id: SessionId, access: SessionAccess, options?: SessionPersistenceOpenOptions): Promise<SessionHandle> {
  308. options?.signal?.throwIfAborted()
  309. await this.ensureRootEncoding()
  310. options?.signal?.throwIfAborted()
  311. const pending = this.tracker.pendingOf(id)
  312. if (access === 'read') {
  313. if (pending !== undefined) {
  314. return this.tracker.adopt(new JsonlSessionHandle(this, id, pending.header, 'read', { cursor: 0, materialized: false, inheritedEventCount: pending.inheritedEventCount }))
  315. }
  316. const stored = await this.requireStoredLog(id, options?.signal)
  317. let state: StorageHandleState
  318. if (stored.status === 'prepared') {
  319. state = {
  320. cursor: 0,
  321. materialized: true,
  322. inheritedEventCount: stored.inheritedEventCount,
  323. primed: stored,
  324. }
  325. } else {
  326. state = {
  327. cursor: 0,
  328. materialized: true,
  329. inheritedEventCount: stored.inheritedEventCount,
  330. }
  331. }
  332. return this.tracker.adopt(new JsonlSessionHandle(this, id, stored.meta, 'read', state))
  333. }
  334. // A pending entry always belongs to an ACTIVE creator handle (close erases
  335. // it), so the claim below rejects that case as already owned.
  336. this.tracker.claimWrite(id)
  337. let lease: SessionWriteLease | undefined
  338. try {
  339. const resolved = await this.findLog(id, options?.signal)
  340. if (resolved === undefined) throw new SessionPersistenceNotFoundError(id)
  341. lease = await this.acquireLease(id, undefined, dirname(resolved.currentPath))
  342. const prepared = await this.requireStoredLog(id, options?.signal)
  343. options?.signal?.throwIfAborted()
  344. let stored: CurrentStoredLog
  345. if (prepared.status === 'prepared') {
  346. stored = await this.publishStoredMigration(id, prepared)
  347. } else {
  348. stored = prepared
  349. }
  350. options?.signal?.throwIfAborted()
  351. return this.tracker.adopt(new JsonlSessionHandle(this, id, stored.meta, 'write', {
  352. cursor: stored.events.length,
  353. materialized: true,
  354. tornTruncateTo: stored.tornTruncateTo,
  355. recoveredTail: stored.recoveredTail,
  356. inheritedEventCount: stored.inheritedEventCount,
  357. primed: stored,
  358. }, lease))
  359. } catch (error) {
  360. // Free the in-process claim no matter how the kernel-lock release
  361. // fares, and keep the original diagnostic: a release failure joins it
  362. // instead of replacing it.
  363. /* v8 ignore next -- typed backends and fs reject with Error */
  364. const failure = error instanceof Error ? error : new Error(String(error))
  365. let releaseFailure: Error | undefined
  366. try {
  367. await lease?.release()
  368. } catch (raw: unknown) {
  369. /* v8 ignore next -- lock releases reject with Error */
  370. releaseFailure = raw instanceof Error ? raw : new Error(String(raw))
  371. }
  372. this.tracker.releaseClaim(id)
  373. if (releaseFailure !== undefined) {
  374. throw new AggregateError([failure, releaseFailure], `session "${id}": write open failed and its lock release failed`)
  375. }
  376. throw failure
  377. }
  378. }
  379. /**
  380. * Flush every active write handle in one durability barrier; see the seam
  381. * contract.
  382. * @returns resolution once every write handle active at the call has flushed.
  383. */
  384. flush(): Promise<void> {
  385. return this.tracker.flushAll()
  386. }
  387. /**
  388. * Observe one stored session without reading its event log.
  389. * @param id - the stored session to observe.
  390. * @param options - optional cancellation.
  391. * @returns the snapshot (`sizeBytes` carries the physical artifact size), or
  392. * `undefined` when the session does not exist.
  393. */
  394. async stat(
  395. id: SessionId,
  396. options?: SessionPersistenceStatOptions,
  397. ): Promise<SessionPersistenceSnapshot | undefined> {
  398. options?.signal?.throwIfAborted()
  399. await this.ensureRootEncoding()
  400. options?.signal?.throwIfAborted()
  401. const pending = this.tracker.pendingOf(id)
  402. if (pending !== undefined) {
  403. return { header: pending.header, revision: pending.revision }
  404. }
  405. const selected = await this.findLog(id, options?.signal)
  406. if (selected === undefined) return undefined
  407. const header = await this.readGenerationHeader(selected, id, options?.signal)
  408. if (header === undefined) return undefined
  409. try {
  410. const identity = await stat(selected.sourcePath, { bigint: true })
  411. options?.signal?.throwIfAborted()
  412. return {
  413. header,
  414. revision: fileRevision(identity),
  415. sizeBytes: Number(identity.size),
  416. }
  417. } catch (error: unknown) {
  418. options?.signal?.throwIfAborted()
  419. if (isENOENT(error)) return undefined
  420. throw error
  421. }
  422. }
  423. /**
  424. * List every stored session visible to this process: materialized artifacts
  425. * plus this process's created-but-unmaterialized sessions.
  426. * @param options - optional cancellation.
  427. * @returns one snapshot per session, in no promised order.
  428. */
  429. async list(options?: SessionPersistenceListOptions): Promise<readonly SessionPersistenceSnapshot[]> {
  430. const signal = options?.signal
  431. const snapshots: SessionPersistenceSnapshot[] = []
  432. const listed = new Set<SessionId>()
  433. // Snapshot pending entries BEFORE scanning storage: a session whose first
  434. // append lands mid-scan is then still in this snapshot (its artifact may
  435. // predate the scan), so create-to-list visibility never has a hole.
  436. const pending = [...this.tracker.pendingEntries()]
  437. for (const artifact of await this.listArtifacts(signal)) {
  438. signal?.throwIfAborted()
  439. try {
  440. const identity = await stat(artifact.path, { bigint: true })
  441. signal?.throwIfAborted()
  442. listed.add(artifact.header.id)
  443. snapshots.push({
  444. header: artifact.header,
  445. revision: fileRevision(identity),
  446. sizeBytes: Number(identity.size),
  447. })
  448. } catch (error: unknown) {
  449. signal?.throwIfAborted()
  450. if (!isENOENT(error)) throw error
  451. }
  452. }
  453. for (const [id, entry] of pending) {
  454. if (!listed.has(id)) snapshots.push({ header: entry.header, revision: entry.revision })
  455. }
  456. signal?.throwIfAborted()
  457. return snapshots
  458. }
  459. // --- handle-facing storage internals (package-private via the handle class below) ---
  460. /** Resolve and read one stored log, refusing loudly when the artifact is absent. */
  461. private async requireStoredLog(id: SessionId, signal?: AbortSignal): Promise<StoredLog> {
  462. const selected = await this.findLog(id, signal)
  463. if (selected === undefined) throw new SessionPersistenceNotFoundError(id)
  464. if (selected.sourceVersion < SESSION_FORMAT_VERSION) {
  465. const sourceRevision = fileRevision(await stat(selected.sourcePath, { bigint: true }))
  466. signal?.throwIfAborted()
  467. let preparation = this.migrationPreparations.get(id)
  468. if (preparation === undefined
  469. || preparation.sourcePath !== selected.sourcePath
  470. || preparation.sourceRevision !== sourceRevision) {
  471. const controller = new AbortController()
  472. const promise = this.loadStoredMigration(id, selected, sourceRevision, controller.signal)
  473. preparation = {
  474. sourcePath: selected.sourcePath,
  475. sourceRevision,
  476. controller,
  477. promise,
  478. settled: false,
  479. waiters: 0,
  480. }
  481. this.migrationPreparations.set(id, preparation)
  482. const created = preparation
  483. const release = (): void => {
  484. created.settled = true
  485. if (this.migrationPreparations.get(id) === created) {
  486. this.migrationPreparations.delete(id)
  487. }
  488. }
  489. void promise.then(release, release)
  490. }
  491. signal?.throwIfAborted()
  492. return this.waitForPreparation(id, preparation, signal)
  493. }
  494. if (selected.sourceVersion > SESSION_FORMAT_VERSION) {
  495. const header = await this.readGenerationHeader(selected, id, signal)
  496. /* v8 ignore else -- a readable future header is rejected inside readGenerationHeader. */
  497. if (header === undefined) {
  498. throw new SessionPersistenceCorruptionError(
  499. `session "${id}": stored log has a malformed header (raw log: ${selected.sourcePath})`,
  500. { cause: new Error('malformed Session header') },
  501. )
  502. }
  503. /* v8 ignore next -- readGenerationHeader rejects every future version. */
  504. throw new SessionFormatUnsupportedError(
  505. `${sessionFormatVersionRefusal(id, selected.sourceVersion)} (raw log: ${selected.sourcePath})`,
  506. { kind: 'jsonl', path: selected.sourcePath },
  507. )
  508. }
  509. const probe = fileRevision(await stat(selected.sourcePath, { bigint: true }))
  510. const memoized = this.coldLogMemo.get(id)
  511. if (memoized?.status === 'current' && memoized.revision === probe) {
  512. this.coldLogMemo.delete(id)
  513. this.coldLogMemo.set(id, memoized)
  514. return memoized
  515. }
  516. const current = await readStableJsonlFile(selected.sourcePath, signal)
  517. return this.decodeStoredLog(
  518. selected.sourcePath,
  519. id,
  520. current.bytes,
  521. fileRevision(current.identity),
  522. signal,
  523. )
  524. }
  525. /** Probe the memo and otherwise decode one historical generation under backend cancellation. */
  526. private async loadStoredMigration(
  527. id: SessionId,
  528. selected: ResolvedJsonlGeneration,
  529. sourceRevision: PersistenceRevision,
  530. signal: AbortSignal,
  531. ): Promise<PreparedStoredLog> {
  532. signal.throwIfAborted()
  533. const memoized = this.coldLogMemo.get(id)
  534. if (memoized?.status === 'prepared' && memoized.revision === sourceRevision) {
  535. this.coldLogMemo.delete(id)
  536. this.coldLogMemo.set(id, memoized)
  537. return memoized
  538. }
  539. return this.prepareStoredMigration(id, selected, signal)
  540. }
  541. /** Await shared preparation for one caller and abort it only after its last waiter leaves. */
  542. private async waitForPreparation(
  543. id: SessionId,
  544. preparation: MigrationPreparation,
  545. signal?: AbortSignal,
  546. ): Promise<PreparedStoredLog> {
  547. preparation.waiters += 1
  548. try {
  549. return await waitWithAbort(preparation.promise, signal)
  550. } finally {
  551. preparation.waiters -= 1
  552. if (preparation.waiters === 0 && !preparation.settled) {
  553. /* v8 ignore else -- a newer selected source may already own this id's preparation slot. */
  554. if (this.migrationPreparations.get(id) === preparation) {
  555. this.migrationPreparations.delete(id)
  556. }
  557. preparation.controller.abort()
  558. }
  559. }
  560. }
  561. /** Decode one historical generation without publishing a successor. */
  562. private async prepareStoredMigration(
  563. id: SessionId,
  564. selected: ResolvedJsonlGeneration,
  565. signal: AbortSignal,
  566. ): Promise<PreparedStoredLog> {
  567. let prepared: Awaited<ReturnType<typeof prepareJsonlMigration>>
  568. try {
  569. prepared = await prepareJsonlMigration({
  570. sourcePath: selected.sourcePath,
  571. sourceVersion: selected.sourceVersion,
  572. currentPath: selected.currentPath,
  573. compression: this.compression,
  574. format: this.generationFormat,
  575. verifyCurrentFile: verifyCurrentGenerationInWorker,
  576. validateHistoricalHeader: headerValue => this.validateSourceIdentity(
  577. selected,
  578. headerValue,
  579. id,
  580. signal,
  581. ),
  582. signal,
  583. })
  584. } catch (error: unknown) {
  585. throw this.generationFailure(id, selected, error)
  586. }
  587. const meta = this.currentHeader(prepared.artifact.header)
  588. assertStoredId(id, meta)
  589. const events = prepared.artifact.events as SessionEvent[]
  590. validateStoredEvents(meta, events, { kind: 'jsonl', path: selected.sourcePath })
  591. const stored: PreparedStoredLog = {
  592. status: 'prepared',
  593. meta,
  594. ...freezeStoredEvents(events),
  595. tornTruncateTo: undefined,
  596. recoveredTail: [],
  597. inheritedEventCount: SessionLogOffset(prepared.artifact.inheritedEventCount),
  598. revision: fileRevision(prepared.sourceIdentity),
  599. publication: { source: selected, value: prepared },
  600. }
  601. this.memoizeStoredLog(id, stored)
  602. return stored
  603. }
  604. /** Publish a prepared historical log before granting write access. */
  605. private async publishStoredMigration(id: SessionId, stored: PreparedStoredLog): Promise<CurrentStoredLog> {
  606. const migration = stored.publication
  607. let identity: JsonlPhysicalIdentity
  608. try {
  609. identity = await migration.value.publish()
  610. } catch (error: unknown) {
  611. /* v8 ignore else -- a newer preparation may have replaced this stale cache entry. */
  612. if (this.coldLogMemo.get(id) === stored) this.coldLogMemo.delete(id)
  613. throw this.generationFailure(id, migration.source, error)
  614. }
  615. const published: CurrentStoredLog = {
  616. status: 'current',
  617. meta: stored.meta,
  618. eventState: stored.eventState,
  619. events: stored.events,
  620. tornTruncateTo: stored.tornTruncateTo,
  621. recoveredTail: stored.recoveredTail,
  622. inheritedEventCount: stored.inheritedEventCount,
  623. revision: fileRevision(identity),
  624. }
  625. this.memoizeStoredLog(id, published)
  626. return published
  627. }
  628. /** Translate generation-layer failures into the persistence seam's error vocabulary. */
  629. private generationFailure(
  630. id: SessionId,
  631. selected: ResolvedJsonlGeneration,
  632. error: unknown,
  633. ): Error {
  634. if (error instanceof JsonlGenerationUnsupportedMigrationError) {
  635. return new SessionFormatUnsupportedError(
  636. `${error.message}; source v${error.fromVersion} artifact remains unchanged (raw log: ${selected.sourcePath})`,
  637. { kind: 'jsonl', path: selected.sourcePath },
  638. )
  639. }
  640. if (error instanceof JsonlGenerationSourceChangedError) return error
  641. if (error instanceof SessionFormatUnsupportedError
  642. || error instanceof SessionPersistenceCorruptionError
  643. || isErrnoException(error)
  644. || error instanceof DOMException && error.name === 'AbortError') return error
  645. return new SessionPersistenceCorruptionError(
  646. `session "${id}": stored log is corrupt: ${String(error)} (raw log: ${selected.sourcePath})`,
  647. { cause: error },
  648. )
  649. }
  650. /**
  651. * Read, parse, and validate one stored log as the current logical prefix.
  652. * @param path - the artifact file to read.
  653. * @param expectedId - the session identity the artifact must carry.
  654. * @param signal - optional cancellation for the stat/read/decode work.
  655. * @returns the validated stored log with any torn-tail truncation point.
  656. */
  657. async readStoredLog(path: string, expectedId: SessionId, signal?: AbortSignal): Promise<CurrentStoredLog> {
  658. signal?.throwIfAborted()
  659. const probe = fileRevision(await stat(path, { bigint: true }))
  660. const memoized = this.coldLogMemo.get(expectedId)
  661. if (memoized?.status === 'current' && memoized.revision === probe) {
  662. this.coldLogMemo.delete(expectedId)
  663. this.coldLogMemo.set(expectedId, memoized)
  664. return memoized
  665. }
  666. const { bytes, identity } = await readStableJsonlFile(path, signal)
  667. return this.decodeStoredLog(path, expectedId, bytes, fileRevision(identity), signal)
  668. }
  669. /** Decode and memoize one already-stable current physical snapshot. */
  670. private async decodeStoredLog(
  671. path: string,
  672. expectedId: SessionId,
  673. buffer: Buffer,
  674. revision: PersistenceRevision,
  675. signal?: AbortSignal,
  676. ): Promise<CurrentStoredLog> {
  677. let parsed: {
  678. meta: SessionHeader
  679. inheritedEventCount: SessionLogOffsetType
  680. events: SessionEvent[]
  681. tornTruncateTo: number | undefined
  682. recoveredTail: SessionEvent[]
  683. }
  684. try {
  685. if (this.compression === 'zstd') {
  686. parsed = await this.readZstdPrefix(buffer, signal)
  687. } else {
  688. signal?.throwIfAborted()
  689. const { meta, inheritedEventCount, events, committedBytes } = scanLog(buffer)
  690. signal?.throwIfAborted()
  691. parsed = {
  692. meta,
  693. inheritedEventCount,
  694. events,
  695. tornTruncateTo: committedBytes < buffer.byteLength ? committedBytes : undefined,
  696. // A torn raw tail is one incomplete JSONL line; it holds no complete
  697. // record to recover.
  698. recoveredTail: [],
  699. }
  700. }
  701. } catch (error: unknown) {
  702. signal?.throwIfAborted()
  703. // A parse-time format refusal predates any SessionHeader, so attach the
  704. // artifact this read actually refused; every other parse failure is
  705. // committed bytes the decoder cannot interpret — damage, classified for
  706. // the seam's stable error vocabulary.
  707. if (error instanceof SessionFormatUnsupportedError) {
  708. throw new SessionFormatUnsupportedError(`${error.message} (raw log: ${path})`, { kind: 'jsonl', path })
  709. }
  710. throw new SessionPersistenceCorruptionError(`session "${expectedId}": stored log is corrupt: ${String(error)} (raw log: ${path})`, { cause: error })
  711. }
  712. signal?.throwIfAborted()
  713. await this.assertStoredIdentity(path, SESSION_FORMAT_VERSION, parsed.meta, expectedId, signal)
  714. signal?.throwIfAborted()
  715. assertStoredId(expectedId, parsed.meta)
  716. const location = this.locate(parsed.meta)
  717. validateStoredEvents(parsed.meta, parsed.events, location)
  718. const { events, ...rest } = parsed
  719. const stored: CurrentStoredLog = {
  720. status: 'current',
  721. ...rest,
  722. ...freezeStoredEvents(events),
  723. revision,
  724. }
  725. this.memoizeStoredLog(expectedId, stored)
  726. return stored
  727. }
  728. /** Insert one parsed log into the bounded handoff cache. */
  729. private memoizeStoredLog(id: SessionId, stored: StoredLog): void {
  730. this.coldLogMemo.delete(id)
  731. this.coldLogMemo.set(id, stored)
  732. for (const oldest of this.coldLogMemo.keys()) {
  733. if (this.coldLogMemo.size <= COLD_LOG_MEMO_MAX_ENTRIES) break
  734. this.coldLogMemo.delete(oldest)
  735. }
  736. }
  737. /**
  738. * Resolve a session's current-generation log path.
  739. * @param id - the stored session to locate.
  740. * @param signal - optional cancellation for the directory scans.
  741. * @returns the current artifact path, or `undefined` while only a historical generation exists.
  742. */
  743. async resolveCurrentLog(id: SessionId, signal?: AbortSignal): Promise<string | undefined> {
  744. await this.ensureRootEncoding()
  745. signal?.throwIfAborted()
  746. const selected = await this.findLog(id, signal)
  747. if (selected === undefined) return undefined
  748. if (selected.sourceVersion === SESSION_FORMAT_VERSION) return selected.sourcePath
  749. if (selected.sourceVersion < SESSION_FORMAT_VERSION) return undefined
  750. const reason = sessionFormatVersionRefusal(id, selected.sourceVersion)
  751. throw new SessionFormatUnsupportedError(
  752. `${reason} (raw log: ${selected.sourcePath})`,
  753. { kind: 'jsonl', path: selected.sourcePath },
  754. )
  755. }
  756. /**
  757. * Durably append one validated batch; lazily materializes on the first write.
  758. * @param header - the session's stored header.
  759. * @param events - the validated contiguous batch, in seq order.
  760. * @param isMaterialized - whether the session already has a durable artifact.
  761. * @param inheritedEventCount - the exact fork-inherited prefix length written into a materializing header line.
  762. */
  763. async persistBatch(
  764. header: SessionHeader,
  765. events: readonly SessionEvent[],
  766. isMaterialized: boolean,
  767. inheritedEventCount: SessionLogOffsetType,
  768. ): Promise<void> {
  769. this.coldLogMemo.delete(header.id)
  770. await this.ensureRootEncoding()
  771. if (isMaterialized) {
  772. await this.appendLines(header, events)
  773. } else {
  774. await this.materialize(header, inheritedEventCount, events)
  775. this.tracker.materialized(header.id)
  776. }
  777. }
  778. /**
  779. * Materialize a header-only artifact for an explicitly durable empty session.
  780. * @param header - the session's stored header.
  781. * @param inheritedEventCount - the exact fork-inherited prefix length written into the header line.
  782. */
  783. async persistHeader(header: SessionHeader, inheritedEventCount: SessionLogOffsetType): Promise<void> {
  784. this.coldLogMemo.delete(header.id)
  785. await this.ensureRootEncoding()
  786. await this.materialize(header, inheritedEventCount, [])
  787. this.tracker.materialized(header.id)
  788. }
  789. /**
  790. * Truncate a torn physical tail durably before this session's first new append.
  791. * @param header - the session's stored header.
  792. * @param truncateTo - the byte offset the artifact is truncated to.
  793. */
  794. async truncateTornTail(header: SessionHeader, truncateTo: number): Promise<void> {
  795. this.coldLogMemo.delete(header.id)
  796. await this.repair(header, truncateTo)
  797. this.ctx.logger.warn(`${this.name}: session "${header.id}" recovered from a torn tail; incomplete tail bytes were discarded`)
  798. }
  799. /**
  800. * Whether this process still tracks a created-but-unmaterialized session.
  801. * @param id - the session to test.
  802. * @returns true while the pending entry exists.
  803. */
  804. hasPendingSession(id: SessionId): boolean {
  805. return this.tracker.hasPending(id)
  806. }
  807. /**
  808. * Release one handle's backend bookkeeping on close.
  809. * @param handle - the closing handle.
  810. * @param materialized - whether the session reached durable storage.
  811. */
  812. releaseHandle(handle: JsonlSessionHandle, materialized: boolean): void {
  813. this.tracker.release(handle, materialized)
  814. }
  815. /**
  816. * Acquire the session directory's kernel write lock; the kernel holds it
  817. * until the handle's close releases the descriptor, including on process death.
  818. * @param id - the session the lock guards.
  819. * @param cwd - header cwd used to derive the directory for a fresh session.
  820. * @param dir - the resolved directory of an existing artifact, when known.
  821. * @returns the held lock.
  822. */
  823. private acquireLease(id: SessionId, cwd: string | undefined, dir = sessionDir(this.root, cwd, id)): Promise<SessionWriteLease> {
  824. return SessionWriteLease.acquire(dir, id)
  825. }
  826. /**
  827. * Acquire the cross-process write lock for a materializing created session,
  828. * called by its handle immediately before the first log bytes publish.
  829. * @param header - the session's stored header (its cwd derives the directory).
  830. * @returns the held lock.
  831. */
  832. async acquireWriteLease(header: SessionHeader): Promise<SessionWriteLease> {
  833. // Refuse an opposite-encoding artifact before the lock's mkdir publishes
  834. // the session directory — the last moment the directory can be absent.
  835. await this.rejectOppositeArtifact(header.cwd, header.id)
  836. return this.acquireLease(header.id, header.cwd)
  837. }
  838. /** Decode complete frames and retain complete JSONL records from a torn final frame. */
  839. private async readZstdPrefix(
  840. buffer: Buffer,
  841. signal?: AbortSignal,
  842. ): Promise<{
  843. meta: SessionHeader
  844. inheritedEventCount: SessionLogOffsetType
  845. events: SessionEvent[]
  846. tornTruncateTo: number | undefined
  847. recoveredTail: SessionEvent[]
  848. }> {
  849. signal?.throwIfAborted()
  850. const { frames, tornStart } = scanZstdFrames(buffer)
  851. signal?.throwIfAborted()
  852. if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')
  853. const decoder = createZstdFrameDecoder()
  854. let yieldDeadline = performance.now() + ZSTD_DECODE_YIELD_INTERVAL_MS
  855. try {
  856. const decodedFrames = decoder.decode(buffer, frames)
  857. signal?.throwIfAborted()
  858. const headerFrame = decodedFrames.next()
  859. signal?.throwIfAborted()
  860. /* v8 ignore next -- a non-empty structural frame list makes the decoder yield its first frame or throw. */
  861. if (headerFrame.done) throw new Error('empty or header-less Zstandard session log')
  862. assertZstdHeaderFrame(headerFrame.value)
  863. const scanner = new SessionLogScanner(headerFrame.value)
  864. let remainingFrames = frames.length - 1
  865. for (const plaintext of decodedFrames) {
  866. signal?.throwIfAborted()
  867. scanner.write(plaintext)
  868. remainingFrames -= 1
  869. if (remainingFrames > 0 && performance.now() >= yieldDeadline) {
  870. await scheduler.yield()
  871. signal?.throwIfAborted()
  872. yieldDeadline = performance.now() + ZSTD_DECODE_YIELD_INTERVAL_MS
  873. }
  874. }
  875. signal?.throwIfAborted()
  876. const complete = scanner.checkpoint()
  877. if (complete.committedBytes !== complete.inputBytes) {
  878. throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')
  879. }
  880. if (tornStart === undefined) {
  881. const prefix = scanner.finish()
  882. return {
  883. meta: prefix.meta,
  884. inheritedEventCount: prefix.inheritedEventCount,
  885. events: prefix.events,
  886. tornTruncateTo: undefined,
  887. recoveredTail: [],
  888. }
  889. }
  890. // A torn final frame's append never resolved, but complete JSONL records
  891. // already flushed into it are real emitted events: recover them, and let
  892. // the write path truncate the torn bytes and rewrite them durably.
  893. let recoveredPlaintext: Buffer = Buffer.alloc(0)
  894. try {
  895. signal?.throwIfAborted()
  896. recoveredPlaintext = await decompressZstdPrefix(buffer.subarray(tornStart))
  897. } catch {
  898. /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
  899. if (signal?.aborted) signal.throwIfAborted()
  900. // A structurally incomplete final frame may end before Node's decoder
  901. // can emit any plaintext; the complete prior frames remain recoverable.
  902. }
  903. signal?.throwIfAborted()
  904. scanner.write(recoveredPlaintext)
  905. const prefix = scanner.finish()
  906. return {
  907. meta: prefix.meta,
  908. inheritedEventCount: prefix.inheritedEventCount,
  909. events: prefix.events,
  910. tornTruncateTo: tornStart,
  911. recoveredTail: prefix.events.slice(complete.eventCount),
  912. }
  913. } catch (error) {
  914. /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
  915. if (signal?.aborted) signal.throwIfAborted()
  916. throw error
  917. } finally {
  918. decoder.close()
  919. }
  920. }
  921. private async listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>> {
  922. signal?.throwIfAborted()
  923. await this.ensureRootEncoding()
  924. signal?.throwIfAborted()
  925. const artifacts: Array<{ header: SessionHeader; path: string }> = []
  926. const ids = new Set<SessionId>()
  927. for (const project of await this.listProjectDirs(signal)) {
  928. signal?.throwIfAborted()
  929. for (const dir of await this.listSessionDirs(project, signal)) {
  930. signal?.throwIfAborted()
  931. const selected = await this.resolveGenerationInDirectory(dir, signal)
  932. if (selected === undefined) continue
  933. let header: SessionHeader | undefined
  934. try {
  935. header = await this.readGenerationHeader(selected, undefined, signal)
  936. } catch (error: unknown) {
  937. // Listing skips a foreign format while opening its id still refuses
  938. // with the selected physical location.
  939. if (error instanceof SessionFormatUnsupportedError) continue
  940. throw error
  941. }
  942. if (header === undefined) continue
  943. if (ids.has(header.id)) {
  944. throw new Error(`duplicate JSONL session id "${header.id}" appears in multiple project directories`)
  945. }
  946. ids.add(header.id)
  947. artifacts.push({ header, path: selected.sourcePath })
  948. }
  949. }
  950. signal?.throwIfAborted()
  951. return artifacts
  952. }
  953. /** Read and translate one selected generation header without inspecting its body. */
  954. private async readGenerationHeader(
  955. selected: ResolvedJsonlGeneration,
  956. expectedId?: SessionId,
  957. signal?: AbortSignal,
  958. ): Promise<SessionHeader | undefined> {
  959. let first: string | undefined
  960. try {
  961. first = this.compression === 'zstd'
  962. ? await this.readFirstZstdLine(selected.sourcePath, signal)
  963. : await this.readFirstLine(selected.sourcePath, signal)
  964. } catch (error: unknown) {
  965. signal?.throwIfAborted()
  966. if (isENOENT(error)) return undefined
  967. throw error
  968. }
  969. signal?.throwIfAborted()
  970. if (first === undefined) return undefined
  971. let value: unknown
  972. try {
  973. value = JSON.parse(first)
  974. } catch {
  975. return undefined
  976. }
  977. assertNoRetiredHeaderFields(value)
  978. const result = sessionFormatCatalog.readHeader(value)
  979. if ('storedVersion' in result && result.storedVersion !== selected.sourceVersion) {
  980. throw new Error(
  981. `session generation filename identifies v${selected.sourceVersion}, `
  982. + `but its header identifies v${result.storedVersion}`,
  983. )
  984. }
  985. if (result.status === 'unsupported') {
  986. const physicalId = String((value as { id?: unknown }).id)
  987. let reason = result.reason
  988. /* v8 ignore else -- released historical header migrations cannot refuse after physical decoding. */
  989. if (result.storedVersion > SESSION_FORMAT_VERSION) {
  990. reason = sessionFormatVersionRefusal(physicalId, result.storedVersion)
  991. }
  992. throw new SessionFormatUnsupportedError(
  993. `${reason} (raw log: ${selected.sourcePath})`,
  994. { kind: 'jsonl', path: selected.sourcePath },
  995. )
  996. }
  997. if (result.status === 'malformed') return undefined
  998. const header = this.currentHeader(result.header)
  999. await this.assertStoredIdentity(
  1000. selected.sourcePath,
  1001. selected.sourceVersion,
  1002. header,
  1003. expectedId,
  1004. signal,
  1005. )
  1006. return header
  1007. }
  1008. /** Convert format-catalog string identities to current branded Session metadata. */
  1009. private currentHeader(header: {
  1010. readonly version: number
  1011. readonly id: string
  1012. readonly createdAt: number
  1013. readonly cwd?: string
  1014. readonly parentSession?: string
  1015. readonly isSeeded: boolean
  1016. readonly origin?: 'subagent'
  1017. readonly delegationDepth: number
  1018. readonly agentPreset?: string
  1019. }): SessionHeader {
  1020. /* v8 ignore next 3 -- readable catalog results are restored to its configured current version. */
  1021. if (header.version !== SESSION_FORMAT_VERSION) {
  1022. throw new Error(`format catalog returned non-current logical header v${header.version}`)
  1023. }
  1024. return {
  1025. version: SESSION_FORMAT_VERSION,
  1026. id: makeSessionId(header.id),
  1027. createdAt: header.createdAt,
  1028. ...(header.cwd === undefined ? {} : { cwd: header.cwd }),
  1029. ...(header.parentSession === undefined
  1030. ? {}
  1031. : { parentSession: makeSessionId(header.parentSession) }),
  1032. isSeeded: header.isSeeded,
  1033. ...(header.origin === undefined ? {} : { origin: header.origin }),
  1034. delegationDepth: header.delegationDepth,
  1035. ...(header.agentPreset === undefined ? {} : { agentPreset: header.agentPreset }),
  1036. }
  1037. }
  1038. // --- materialization / append / repair (file mechanics) ---
  1039. /** Atomically write the header line + first batch (temp-write, fsync, publish). */
  1040. private async materialize(
  1041. meta: SessionHeader,
  1042. inheritedEventCount: SessionLogOffsetType,
  1043. events: readonly SessionEvent[],
  1044. ): Promise<void> {
  1045. const project = projectDir(this.root, meta.cwd)
  1046. const dir = sessionDir(this.root, meta.cwd, meta.id)
  1047. const finalPath = logPath(this.root, meta.cwd, meta.id, this.compression)
  1048. await this.rejectOppositeArtifact(meta.cwd, meta.id)
  1049. const content = await this.encodeMaterialization(meta, inheritedEventCount, events)
  1050. /* v8 ignore next -- native Windows coverage exercises this platform dispatch; Linux covers the POSIX peer */
  1051. if (process.platform === 'win32') {
  1052. await this.materializeWin32(project, dir, finalPath, meta.id, content)
  1053. } else {
  1054. await this.materializePosix(project, dir, finalPath, meta.id, content)
  1055. }
  1056. }
  1057. /* v8 ignore start -- Windows uses the Win32 durable-publish path; POSIX coverage exercises this peer. */
  1058. private async materializePosix(
  1059. project: string,
  1060. dir: string,
  1061. finalPath: string,
  1062. id: SessionId,
  1063. content: Buffer | string,
  1064. ): Promise<void> {
  1065. await mkdir(this.root, { recursive: true, mode: 0o700 })
  1066. await this.syncDirPosix(dirname(this.root))
  1067. await mkdir(project, { recursive: true, mode: 0o700 })
  1068. await this.syncDirPosix(this.root)
  1069. await mkdir(dir, { recursive: true, mode: 0o700 })
  1070. await this.syncDirPosix(project)
  1071. await this.rejectExistingLog(finalPath, id)
  1072. const tmp = await this.writeSyncedTempFile(finalPath, content)
  1073. // Publish via link()+unlink(), NOT rename(): link fails with EEXIST if the
  1074. // final path already exists, so two processes materializing the same id
  1075. // concurrently cannot clobber each other. rename() would silently overwrite.
  1076. let linked = false
  1077. try {
  1078. await link(tmp, finalPath)
  1079. linked = true
  1080. } finally {
  1081. // Remove an unpublished temp on failure. After publication, defer cleanup
  1082. // until the directory entry is durable so cleanup cannot reject a live log.
  1083. /* v8 ignore next -- link failure is the TOCTOU/IO race guarded above; not reachable in test */
  1084. if (!linked) await rm(tmp, { force: true })
  1085. }
  1086. // link() succeeded — the log is published. fsync the directory so the new
  1087. // entry survives a power loss: the new link is not crash-durable until the
  1088. // parent directory's metadata is synced.
  1089. await this.syncDirPosix(dir)
  1090. // Best-effort temp cleanup: the log is already published and durable, so a
  1091. // failure to remove the redundant temp hard link must NOT reject the
  1092. // append. Swallow only the rm failure; nothing else of consequence runs here.
  1093. try {
  1094. await rm(tmp, { force: true })
  1095. } catch {
  1096. /* v8 ignore next -- redundant temp link; publish already durable, rm failure is an unreachable IO edge */
  1097. }
  1098. }
  1099. /* v8 ignore stop */
  1100. /* v8 ignore start -- native Windows coverage exercises this integration path */
  1101. private async materializeWin32(
  1102. project: string,
  1103. dir: string,
  1104. finalPath: string,
  1105. id: SessionId,
  1106. content: Buffer | string,
  1107. ): Promise<void> {
  1108. await ensureDurableDirectoryWin32(this.root)
  1109. await ensureDurableDirectoryWin32(project)
  1110. await ensureDurableDirectoryWin32(dir)
  1111. await this.rejectExistingLog(finalPath, id)
  1112. const tmp = await this.writeSyncedTempFile(finalPath, content)
  1113. try {
  1114. await publishNewFileWin32(tmp, finalPath)
  1115. } catch (error) {
  1116. await rm(tmp, { force: true })
  1117. throw error
  1118. }
  1119. }
  1120. /* v8 ignore stop */
  1121. private async rejectExistingLog(finalPath: string, id: SessionId): Promise<void> {
  1122. // Never publish over an existing committed log: materialize is the first
  1123. // write of a session the backend believes is new. A file here means a
  1124. // different session shares this id on disk — reject loudly. (create already
  1125. // guards the create path, so this is unreachable-in-practice TOCTOU
  1126. // defense.)
  1127. /* v8 ignore next 3 -- create guards collisions before materialize; this is a TOCTOU backstop */
  1128. if (await this.resolveGenerationInDirectory(dirname(finalPath)) !== undefined) {
  1129. throw new Error(`refusing to materialize "${id}": a log already exists on disk (open it instead)`)
  1130. }
  1131. }
  1132. private async writeSyncedTempFile(finalPath: string, content: Buffer | string): Promise<string> {
  1133. const tmp = `${finalPath}.${randomBytes(6).toString('hex')}.tmp`
  1134. const handle = await open(tmp, 'wx', 0o600)
  1135. try {
  1136. await handle.writeFile(content)
  1137. await handle.sync()
  1138. } finally {
  1139. await handle.close()
  1140. }
  1141. return tmp
  1142. }
  1143. /** Encode the header and first batch without combining their frame boundaries. */
  1144. private async encodeMaterialization(
  1145. meta: SessionHeader,
  1146. inheritedEventCount: SessionLogOffsetType,
  1147. events: readonly SessionEvent[],
  1148. ): Promise<Buffer | string> {
  1149. const header = JSON.stringify(toHeaderLine(meta, meta.isSeeded ? inheritedEventCount : undefined)) + '\n'
  1150. if (events.length === 0) {
  1151. return this.compression === 'none' ? header : compressZstdFrame(header)
  1152. }
  1153. const body = eventLines(events) + '\n'
  1154. if (this.compression === 'none') return header + body
  1155. const headerFrame = await compressZstdFrame(header)
  1156. const eventFrame = await compressZstdFrame(body)
  1157. return Buffer.concat([headerFrame, eventFrame])
  1158. }
  1159. /** Encode one durable append batch in the configured physical representation. */
  1160. private async encodeEventBatch(events: readonly SessionEvent[]): Promise<Buffer | string> {
  1161. const body = eventLines(events) + '\n'
  1162. return this.compression === 'zstd' ? compressZstdFrame(body) : body
  1163. }
  1164. /** fsync a POSIX directory so a just-created/renamed entry is crash-durable. */
  1165. /* v8 ignore start -- Windows uses write-through namespace operations; POSIX coverage exercises directory fsync. */
  1166. private async syncDirPosix(dir: string): Promise<void> {
  1167. const handle = await open(dir, 'r')
  1168. try {
  1169. await handle.sync()
  1170. } finally {
  1171. await handle.close()
  1172. }
  1173. }
  1174. /* v8 ignore stop */
  1175. /**
  1176. * Append and fsync event lines. On a partial write or sync failure, restore the
  1177. * previous size before rethrowing because the unchanged cursor will retry the
  1178. * batch; leaving partial bytes would create duplicate sequence numbers.
  1179. */
  1180. private async appendLines(meta: SessionHeader, events: readonly SessionEvent[]): Promise<void> {
  1181. const content = await this.encodeEventBatch(events)
  1182. const path = logPath(this.root, meta.cwd, meta.id, this.compression)
  1183. const handle = await open(path, 'a')
  1184. let closed = false
  1185. const closeAppendHandle = async (): Promise<void> => {
  1186. if (closed) return
  1187. closed = true
  1188. await handle.close()
  1189. }
  1190. try {
  1191. const { size: before } = await handle.stat()
  1192. try {
  1193. await handle.writeFile(content)
  1194. await handle.sync()
  1195. } catch (error) {
  1196. try {
  1197. await closeAppendHandle()
  1198. await this.rollbackAppend(path, before)
  1199. } catch (rollbackError) {
  1200. throw new AggregateError([error, rollbackError], `failed to roll back append to "${path}"`)
  1201. }
  1202. throw error
  1203. }
  1204. } finally {
  1205. await closeAppendHandle()
  1206. }
  1207. }
  1208. private async rollbackAppend(path: string, size: number): Promise<void> {
  1209. const handle = await open(path, 'r+')
  1210. try {
  1211. await handle.truncate(size)
  1212. await handle.sync()
  1213. } finally {
  1214. await handle.close()
  1215. }
  1216. }
  1217. /** Truncate the log file to `offset` bytes and fsync (discard the crash tail). */
  1218. private async repair(meta: SessionHeader, offset: number): Promise<void> {
  1219. const path = logPath(this.root, meta.cwd, meta.id, this.compression)
  1220. await truncate(path, offset)
  1221. const handle = await open(path, 'r+')
  1222. try {
  1223. await handle.sync()
  1224. } finally {
  1225. await handle.close()
  1226. }
  1227. }
  1228. // --- discovery helpers ---
  1229. /**
  1230. * Read the first newline-terminated line of a file without loading the whole
  1231. * file. Returns undefined if the file is empty or has no complete first line.
  1232. * Reads in bounded chunks so a huge log costs only the header read.
  1233. */
  1234. private async readFirstLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
  1235. signal?.throwIfAborted()
  1236. const handle = await open(path, 'r')
  1237. try {
  1238. signal?.throwIfAborted()
  1239. const chunks: Buffer[] = []
  1240. const buf = Buffer.alloc(8192)
  1241. for (;;) {
  1242. signal?.throwIfAborted()
  1243. const { bytesRead } = await handle.read(buf, 0, buf.length, null)
  1244. signal?.throwIfAborted()
  1245. if (bytesRead === 0) return undefined // EOF with no newline → no complete line
  1246. const slice = buf.subarray(0, bytesRead)
  1247. const nl = slice.indexOf(0x0a)
  1248. if (nl !== -1) {
  1249. chunks.push(slice.subarray(0, nl))
  1250. signal?.throwIfAborted()
  1251. return Buffer.concat(chunks).toString('utf8')
  1252. }
  1253. chunks.push(Buffer.from(slice))
  1254. }
  1255. } finally {
  1256. await handle.close()
  1257. }
  1258. }
  1259. /** Read and validate only the independently compressed header frame. */
  1260. private async readFirstZstdLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
  1261. signal?.throwIfAborted()
  1262. const handle = await open(path, 'r')
  1263. try {
  1264. signal?.throwIfAborted()
  1265. let content = Buffer.alloc(0)
  1266. const chunk = Buffer.alloc(8192)
  1267. for (;;) {
  1268. signal?.throwIfAborted()
  1269. const { bytesRead } = await handle.read(chunk, 0, chunk.length, null)
  1270. signal?.throwIfAborted()
  1271. if (bytesRead === 0) return undefined
  1272. signal?.throwIfAborted()
  1273. content = Buffer.concat([content, chunk.subarray(0, bytesRead)])
  1274. signal?.throwIfAborted()
  1275. const first = scanZstdFrames(content, 1).frames[0]
  1276. signal?.throwIfAborted()
  1277. if (first === undefined) continue
  1278. let plaintext: Buffer
  1279. try {
  1280. signal?.throwIfAborted()
  1281. plaintext = await decompressZstdFrame(content.subarray(first.start, first.end))
  1282. } catch (error) {
  1283. /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
  1284. if (signal?.aborted) signal.throwIfAborted()
  1285. throw new Error('corrupt Zstandard session log: header frame failed validation', { cause: error })
  1286. }
  1287. signal?.throwIfAborted()
  1288. assertZstdHeaderFrame(plaintext)
  1289. return plaintext.subarray(0, -1).toString('utf8')
  1290. }
  1291. } finally {
  1292. await handle.close()
  1293. }
  1294. }
  1295. /** Select the numerically highest canonical generation in one Session directory. */
  1296. private async resolveGenerationInDirectory(
  1297. dir: string,
  1298. signal?: AbortSignal,
  1299. ): Promise<ResolvedJsonlGeneration | undefined> {
  1300. signal?.throwIfAborted()
  1301. let entries: Dirent[]
  1302. try {
  1303. entries = await readdir(dir, { withFileTypes: true })
  1304. } catch (error: unknown) {
  1305. if (isENOENT(error)) return undefined
  1306. throw error
  1307. }
  1308. signal?.throwIfAborted()
  1309. const generations: Array<{ readonly path: string; readonly version: number }> = []
  1310. const opposite: string[] = []
  1311. for (const entry of entries) {
  1312. const version = parseGenerationLogFilename(entry.name, this.compression)
  1313. if (version !== undefined) {
  1314. generations.push({ path: join(dir, entry.name), version })
  1315. continue
  1316. }
  1317. if (parseGenerationLogFilename(entry.name, this.oppositeCompression()) !== undefined) {
  1318. opposite.push(join(dir, entry.name))
  1319. }
  1320. }
  1321. if (opposite.length > 0) throw this.encodingMismatch(opposite[0] as string)
  1322. const latest = generations.sort((left, right) => right.version - left.version)[0]
  1323. if (latest === undefined) return undefined
  1324. return {
  1325. sourcePath: latest.path,
  1326. sourceVersion: latest.version,
  1327. currentPath: join(
  1328. dir,
  1329. generationLogFilename(sessionFormatCatalog.currentVersion, this.compression),
  1330. ),
  1331. }
  1332. }
  1333. /** Find the unique authoritative generation for an id across project directories. */
  1334. private async findLog(id: SessionId, signal?: AbortSignal): Promise<ResolvedJsonlGeneration | undefined> {
  1335. const matches: ResolvedJsonlGeneration[] = []
  1336. for (const project of await this.listProjectDirs(signal)) {
  1337. signal?.throwIfAborted()
  1338. await this.rejectLegacyFlatArtifact(project, id, signal)
  1339. signal?.throwIfAborted()
  1340. const dir = join(project, encodeSegment(id))
  1341. const selected = await this.resolveGenerationInDirectory(dir, signal)
  1342. if (selected !== undefined) matches.push(selected)
  1343. }
  1344. if (matches.length > 1) {
  1345. throw new Error(`duplicate JSONL session id "${id}" appears in multiple project directories`)
  1346. }
  1347. signal?.throwIfAborted()
  1348. return matches[0]
  1349. }
  1350. /** Require an existing configured root to be a readable directory. */
  1351. private assertUsableRoot(): void {
  1352. try {
  1353. readdirSync(this.root)
  1354. } catch (error) {
  1355. if (isENOENT(error)) return
  1356. throw error
  1357. }
  1358. }
  1359. /** Reject metadata that does not identify the selected physical log. */
  1360. private async assertStoredIdentity(
  1361. path: string,
  1362. storedVersion: number,
  1363. meta: SessionHeader,
  1364. expectedId?: SessionId,
  1365. signal?: AbortSignal,
  1366. ): Promise<void> {
  1367. signal?.throwIfAborted()
  1368. if (expectedId !== undefined && meta.id !== expectedId) {
  1369. throw new Error(`corrupt session log "${path}": requested id "${expectedId}" does not match header id "${meta.id}"`)
  1370. }
  1371. let expectedPath: string
  1372. try {
  1373. expectedPath = generationLogPath(
  1374. this.root,
  1375. meta.cwd,
  1376. meta.id,
  1377. storedVersion,
  1378. this.compression,
  1379. )
  1380. } catch (error) {
  1381. throw new Error(`corrupt session log "${path}": header id cannot name a storage path`, { cause: error })
  1382. }
  1383. if (path !== expectedPath && !await this.sameFile(path, expectedPath, signal)) {
  1384. throw new Error(`corrupt session log "${path}": header id "${meta.id}" and cwd identify "${expectedPath}"`)
  1385. }
  1386. signal?.throwIfAborted()
  1387. }
  1388. /** Validate a supported historical header against the selected source path. */
  1389. private validateSourceIdentity(
  1390. selected: ResolvedJsonlGeneration,
  1391. headerValue: Readonly<Record<string, unknown>>,
  1392. expectedId: SessionId,
  1393. signal?: AbortSignal,
  1394. ): void | Promise<void> {
  1395. const result = sessionFormatCatalog.readHeader(headerValue)
  1396. if (result.status !== 'current' && result.status !== 'migration-required') return
  1397. return this.assertStoredIdentity(
  1398. selected.sourcePath,
  1399. selected.sourceVersion,
  1400. this.currentHeader(result.header),
  1401. expectedId,
  1402. signal,
  1403. )
  1404. }
  1405. /**
  1406. * Whether two path spellings resolve to the same physical file. This admits
  1407. * case aliases on case-insensitive filesystems without weakening identity
  1408. * checks on case-sensitive stores.
  1409. */
  1410. private async sameFile(path: string, expectedPath: string, signal?: AbortSignal): Promise<boolean> {
  1411. signal?.throwIfAborted()
  1412. try {
  1413. const [actual, expected] = await Promise.all([realpath(path), realpath(expectedPath)])
  1414. signal?.throwIfAborted()
  1415. return actual === expected
  1416. } catch (error) {
  1417. signal?.throwIfAborted()
  1418. /* v8 ignore else -- non-ENOENT realpath failures require an external permission or I/O fault */
  1419. if (isENOENT(error)) return false
  1420. /* v8 ignore next -- non-ENOENT realpath failures are external I/O faults, propagated unchanged */
  1421. throw error
  1422. }
  1423. }
  1424. /** The human-readable project directories under the configured root. */
  1425. private async listProjectDirs(signal?: AbortSignal): Promise<string[]> {
  1426. try {
  1427. signal?.throwIfAborted()
  1428. const entries = await readdir(this.root, { withFileTypes: true })
  1429. signal?.throwIfAborted()
  1430. return entries.filter(e => e.isDirectory()).map(e => join(this.root, e.name))
  1431. } catch (error) {
  1432. // Only an absent root means no sessions; rethrow every other I/O failure.
  1433. if (isENOENT(error)) return []
  1434. throw error
  1435. }
  1436. }
  1437. /** List session-owned directories and reject the obsolete flat-file layout. */
  1438. private async listSessionDirs(project: string, signal?: AbortSignal): Promise<string[]> {
  1439. signal?.throwIfAborted()
  1440. const entries = await readdir(project, { withFileTypes: true })
  1441. signal?.throwIfAborted()
  1442. const legacy = entries.find(entry =>
  1443. entry.isFile() && (entry.name.endsWith('.jsonl') || entry.name.endsWith('.jsonl.zstd')))
  1444. if (legacy !== undefined) throw this.legacyLayout(join(project, legacy.name))
  1445. return entries.filter(entry => entry.isDirectory()).map(entry => join(project, entry.name))
  1446. }
  1447. /** Reject a root that already belongs to the other physical encoding. */
  1448. private ensureRootEncoding(): Promise<void> {
  1449. this.rootEncodingCheck ??= this.checkRootEncoding()
  1450. return this.rootEncodingCheck
  1451. }
  1452. private async checkRootEncoding(): Promise<void> {
  1453. for (const project of await this.listProjectDirs()) {
  1454. for (const dir of await this.listSessionDirs(project)) {
  1455. const incompatible = await this.findOppositeGenerationInDirectory(dir)
  1456. if (incompatible !== undefined) throw this.encodingMismatch(incompatible)
  1457. }
  1458. }
  1459. }
  1460. private async rejectLegacyFlatArtifact(
  1461. project: string,
  1462. id: SessionId,
  1463. signal?: AbortSignal,
  1464. ): Promise<void> {
  1465. signal?.throwIfAborted()
  1466. const encoded = encodeSegment(id)
  1467. for (const compression of ['zstd', 'none'] as const) {
  1468. const path = join(project, encoded + logSuffix(compression))
  1469. const artifactExists = await this.exists(path)
  1470. signal?.throwIfAborted()
  1471. if (artifactExists) throw this.legacyLayout(path)
  1472. }
  1473. }
  1474. private async rejectOppositeArtifact(cwd: string | undefined, id: SessionId): Promise<void> {
  1475. const path = await this.findOppositeGenerationInDirectory(sessionDir(this.root, cwd, id))
  1476. if (path !== undefined) throw this.encodingMismatch(path)
  1477. }
  1478. /** Return the highest canonical generation encoded with the other configured suffix. */
  1479. private async findOppositeGenerationInDirectory(dir: string): Promise<string | undefined> {
  1480. let entries: Dirent[]
  1481. try {
  1482. entries = await readdir(dir, { withFileTypes: true })
  1483. } catch (error: unknown) {
  1484. if (isENOENT(error)) return undefined
  1485. throw error
  1486. }
  1487. const generations: Array<{ readonly name: string; readonly version: number }> = []
  1488. for (const entry of entries) {
  1489. const version = parseGenerationLogFilename(entry.name, this.oppositeCompression())
  1490. if (version !== undefined) generations.push({ name: entry.name, version })
  1491. }
  1492. const latest = generations.sort((left, right) => right.version - left.version)[0]
  1493. return latest === undefined ? undefined : join(dir, latest.name)
  1494. }
  1495. private oppositeCompression(): JsonlCompression {
  1496. return this.compression === 'zstd' ? 'none' : 'zstd'
  1497. }
  1498. private encodingMismatch(path: string): Error {
  1499. return new Error(
  1500. `session artifact ${JSON.stringify(path)} uses ${logSuffix(this.oppositeCompression())}, `
  1501. + `but this backend is configured for compression ${JSON.stringify(this.compression)}; `
  1502. + 'use a separate root or select the matching compression mode',
  1503. )
  1504. }
  1505. private legacyLayout(path: string): Error {
  1506. return new Error(
  1507. `session artifact ${JSON.stringify(path)} uses the unsupported flat-file layout; `
  1508. + 'use a separate root or move it into a project/session directory before loading',
  1509. )
  1510. }
  1511. private async exists(path: string): Promise<boolean> {
  1512. try {
  1513. const handle = await open(path, 'r')
  1514. await handle.close()
  1515. return true
  1516. } catch (error) {
  1517. // Only ENOENT means absent. A permission/I/O error must surface rather
  1518. // than letting load or collision checks proceed under false absence.
  1519. /* v8 ignore else -- Windows reports file-valued parents as ENOENT; POSIX covers direct ENOTDIR. */
  1520. if (isENOENT(error)) {
  1521. // Windows reports ENOENT, not ENOTDIR, for `regular-file/child`, so it
  1522. // alone verifies the immediate parent to keep a blocked session
  1523. // directory a storage fault. POSIX open already reported ENOTDIR before
  1524. // this point, where the extra stat would only cost a syscall per probe.
  1525. /* v8 ignore next -- native Windows coverage exercises this platform dispatch; POSIX reports ENOTDIR from open */
  1526. if (process.platform === 'win32') await this.assertLogParentAllowsAbsence(path)
  1527. return false
  1528. }
  1529. /* v8 ignore next -- Windows repairs ENOTDIR from ENOENT above; POSIX covers direct ENOTDIR. */
  1530. throw error
  1531. }
  1532. }
  1533. /* v8 ignore start -- native Windows coverage exercises this repair; POSIX open reports ENOTDIR before this point. */
  1534. private async assertLogParentAllowsAbsence(path: string): Promise<void> {
  1535. try {
  1536. const parent = dirname(path)
  1537. const info = await stat(parent)
  1538. if (info.isDirectory()) return
  1539. const error = new Error(`ENOTDIR: parent path exists but is not a directory: ${parent}`) as NodeJS.ErrnoException
  1540. error.code = 'ENOTDIR'
  1541. error.path = parent
  1542. throw error
  1543. } catch (error) {
  1544. if (isENOENT(error)) return
  1545. throw error
  1546. }
  1547. }
  1548. /* v8 ignore stop */
  1549. }
  1550. /**
  1551. * One open channel onto a JSONL-stored session: the shared storage-handle
  1552. * scaffolding over this backend's file primitives. Reads re-scan the artifact
  1553. * under the stable-read loop.
  1554. */
  1555. export default JsonlSessionPersistence