import { createHash } from 'node:crypto' import { createReadStream } from 'node:fs' import { open as fsOpen, mkdir, mkdtemp, readdir, readFile, rename, rm, writeFile, } from 'node:fs/promises' import { tmpdir } from 'node:os' import { dirname, join } from 'node:path' import { Writable } from 'node:stream' import { pipeline } from 'node:stream/promises' import { createZstdCompress, createZstdDecompress } from 'node:zlib' import { ARCHIVE_FORMAT_VERSION, type ArchivedRecord, type SegmentSidecar } from './retention.ts' /** * Versioned zstd JSONL archive segments with immutable sidecar indexes. * Publication is crash-safe: write tmp -> fsync -> rename; reconcile handles orphans. */ export interface PublishedSegment { sidecar: SegmentSidecar zstPath: string } export async function sha256File(path: string): Promise { const hash = createHash('sha256') await pipeline(createReadStream(path), hash) return hash.digest('hex') } export function zstdCompressBuffer(data: Buffer): Promise { const chunks: Buffer[] = [] const stream = createZstdCompress() stream.write(data) stream.end() return new Promise((resolve, reject) => { stream.on('data', (c: Buffer) => chunks.push(c)) stream.on('error', reject) stream.on('end', () => resolve(Buffer.concat(chunks))) }) } export async function zstdDecompressFile(path: string): Promise { const chunks: Buffer[] = [] const sink = new Writable({ write(chunk, _enc, callback) { chunks.push(Buffer.from(chunk)) callback() }, }) await pipeline(createReadStream(path), createZstdDecompress(), sink) return Buffer.concat(chunks) } /** * Write one segment: records as JSONL, zstd-compressed, atomically renamed into place, * with a sidecar index written after the data file so a sidecar implies a complete segment. */ export async function writeSegment( archiveDir: string, kind: 'audit' | 'chat', month: string, records: ArchivedRecord[], segmentName: string, ): Promise { const dir = join(archiveDir, kind, month) const tmp = await mkdtemp(join(tmpdir(), 'will-seg-')) try { const lines = records.map((r) => JSON.stringify({ id: r.id, ts: r.ts, record: r.record })) const zst = await zstdCompressBuffer( Buffer.from(lines.join('\n') + (lines.length ? '\n' : ''), 'utf8'), ) const zstTmp = join(tmp, 'segment.jsonl.zst') const fh = await fsOpen(zstTmp, 'w') try { await fh.writeFile(zst) await fh.sync() } finally { await fh.close() } const zstPath = join(dir, `${segmentName}.jsonl.zst`) const indexPath = join(dir, `${segmentName}.index.json`) await mkdir(dir, { recursive: true }) await rename(zstTmp, zstPath) const first = records[0] const last = records[records.length - 1] if (!first || !last) throw new Error('cannot archive an empty segment') const sidecar: SegmentSidecar = { version: ARCHIVE_FORMAT_VERSION, kind, month, segment: segmentName, count: records.length, firstId: first.id, lastId: last.id, firstTs: first.ts, lastTs: last.ts, sha256: createHash('sha256').update(zst).digest('hex'), bytes: zst.length, } const indexTmp = join(tmp, 'index.json') await writeFile(indexTmp, JSON.stringify(sidecar), 'utf8') await rename(indexTmp, indexPath) return { sidecar, zstPath } } finally { await rm(tmp, { recursive: true, force: true }) } } /** Read all records from a segment file (cold path; slowness accepted by spec). */ export async function readSegment( zstPath: string, filter: (r: ArchivedRecord) => boolean, limit: number, ): Promise { const buf = await zstdDecompressFile(zstPath) const out: ArchivedRecord[] = [] for (const line of buf.toString('utf8').split('\n')) { if (line === '') continue const parsed = JSON.parse(line) as ArchivedRecord if (filter(parsed)) { out.push(parsed) if (out.length >= limit) break } } return out } /** Load and validate every sidecar under archive//. */ export async function listSidecars( archiveDir: string, kind: 'audit' | 'chat', ): Promise<{ sidecar: SegmentSidecar; zstPath: string }[]> { const root = join(archiveDir, kind) const out: { sidecar: SegmentSidecar; zstPath: string }[] = [] let months: string[] try { months = await readdir(root) } catch { return out } for (const month of months) { const files = await readdir(join(root, month)) for (const f of files) { if (!f.endsWith('.index.json')) continue const raw = await readFile(join(root, month, f), 'utf8') const sidecar = JSON.parse(raw) as SegmentSidecar if (sidecar.version !== ARCHIVE_FORMAT_VERSION) { throw new Error(`unsupported segment version ${sidecar.version} in ${month}/${f}`) } out.push({ sidecar, zstPath: join(root, month, f.replace('.index.json', '.jsonl.zst')) }) } } return out } /** Idempotent atomic file publication helper used by the CAS media store too. */ export async function atomicWriteFile(path: string, data: Buffer): Promise { const dir = dirname(path) await mkdir(dir, { recursive: true }) const tmp = join(dir, `.tmp-${process.pid}-${Date.now()}`) const fh = await fsOpen(tmp, 'w') try { await fh.writeFile(data) await fh.sync() } finally { await fh.close() } await rename(tmp, path) }