repositories / will
will
owned by admin
src/shared/archive.ts
Rawimport { 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<string> {
const hash = createHash('sha256')
await pipeline(createReadStream(path), hash)
return hash.digest('hex')
}
export function zstdCompressBuffer(data: Buffer): Promise<Buffer> {
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<Buffer> {
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<PublishedSegment> {
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<ArchivedRecord[]> {
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/<kind>/. */
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<void> {
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)
}