Luigit
repositories / will

will

owned by admin

src/shared/archive.ts

Raw
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<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)
}