SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
5.9 KB · 131 lines typescript
Raw Blame History
1/**2 * Raw body archive on S3/MinIO. Keys: raw/<connector>/<docId>/<contentHash>.<ext>.zst3 * Bodies are compressed with zstd (node:zlib, Node ≥ 22.15) — gzip fallback (".gz") when unavailable.4 * The same content hash is never uploaded twice (HEAD before PUT).5 */6import { S3Client, HeadBucketCommand, CreateBucketCommand, PutObjectCommand, GetObjectCommand, HeadObjectCommand } from "@aws-sdk/client-s3";7import * as zlib from "node:zlib";8import { getEnv } from "./env.js";910type Codec = { ext: "zst" | "gz"; compress(b: Buffer): Buffer; decompress(b: Buffer): Buffer };1112const zstd = (zlib as unknown as { zstdCompressSync?: (b: Buffer) => Buffer; zstdDecompressSync?: (b: Buffer) => Buffer });13export const codec: Codec =14  typeof zstd.zstdCompressSync === "function" && typeof zstd.zstdDecompressSync === "function"15    ? { ext: "zst", compress: (b) => zstd.zstdCompressSync!(b), decompress: (b) => zstd.zstdDecompressSync!(b) }16    : { ext: "gz", compress: (b) => zlib.gzipSync(b), decompress: (b) => zlib.gunzipSync(b) };1718/** Decompress by key suffix so gz-era objects stay readable after a zstd upgrade. */19export function decompressByKey(key: string, body: Buffer): Buffer {20  if (key.endsWith(".zst")) { if (!zstd.zstdDecompressSync) throw new Error("zstd not available in this Node build"); return zstd.zstdDecompressSync(body); }21  if (key.endsWith(".gz")) return zlib.gunzipSync(body);22  return body;23}2425let _client: S3Client | null = null;26export function getS3(): S3Client {27  if (_client) return _client;28  const env = getEnv();29  _client = new S3Client({30    endpoint: env.s3Endpoint,31    region: env.s3Region,32    forcePathStyle: true,33    credentials: { accessKeyId: env.s3AccessKey, secretAccessKey: env.s3SecretKey },34    requestHandler: { requestTimeout: 30_000, connectionTimeout: 5_000 },35  });36  return _client;37}3839export function bucketName(): string { return getEnv().s3Bucket; }4041let bucketReady = false;42export async function ensureBucket(): Promise<void> {43  if (bucketReady) return;44  const s3 = getS3();45  const Bucket = bucketName();46  try {47    await s3.send(new HeadBucketCommand({ Bucket }));48  } catch (e) {49    const code = (e as { name?: string; $metadata?: { httpStatusCode?: number } });50    if (code.name === "NotFound" || code.name === "NoSuchBucket" || code.$metadata?.httpStatusCode === 404) {51      await s3.send(new CreateBucketCommand({ Bucket }));52    } else throw e;53  }54  bucketReady = true;55}5657/** Map a content type (and URL hint) to a short file extension. */58export function extensionFor(contentType: string | null | undefined, url?: string): string {59  const ct = (contentType ?? "").toLowerCase();60  if (/pdf/.test(ct) || /\.pdf(\?|$)/i.test(url ?? "")) return "pdf";61  if (/json/.test(ct)) return "json";62  if (/xml/.test(ct)) return "xml";63  if (/html/.test(ct)) return "html";64  if (/markdown/.test(ct)) return "md";65  if (/csv/.test(ct)) return "csv";66  if (/plain/.test(ct)) return "txt";67  if (/javascript/.test(ct)) return "js";68  if (ct.startsWith("image/")) return ct.split("/")[1]?.split(";")[0] ?? "bin";69  return "bin";70}7172export function rawKey(connectorId: string, docId: string, contentHash: string, ext: string): string {73  return `raw/${connectorId}/${docId}/${contentHash}.${ext}.${codec.ext}`;74}7576export interface PutResult { key: string; stored: boolean; bytes: number; compressedBytes: number }7778/** Store a raw body once per content hash. Returns the key (existing or new). */79export async function putRaw(connectorId: string, docId: string, contentHash: string, body: Buffer, contentType: string | null | undefined, url?: string): Promise<PutResult> {80  await ensureBucket();81  const key = rawKey(connectorId, docId, contentHash, extensionFor(contentType, url));82  const existing = await headRaw(key);83  if (existing) return { key, stored: false, bytes: body.length, compressedBytes: existing.size };84  const compressed = codec.compress(body);85  await getS3().send(86    new PutObjectCommand({87      Bucket: bucketName(),88      Key: key,89      Body: compressed,90      ContentType: "application/octet-stream",91      ContentEncoding: codec.ext === "zst" ? "zstd" : "gzip",92      Metadata: { "content-type": (contentType ?? "application/octet-stream").slice(0, 200), "content-hash": contentHash, "original-bytes": String(body.length) },93    }),94  );95  return { key, stored: true, bytes: body.length, compressedBytes: compressed.length };96}9798export async function headRaw(key: string): Promise<{ size: number; contentType: string | null } | null> {99  try {100    const r = await getS3().send(new HeadObjectCommand({ Bucket: bucketName(), Key: key }));101    return { size: Number(r.ContentLength ?? 0), contentType: r.Metadata?.["content-type"] ?? null };102  } catch (e) {103    const err = e as { name?: string; $metadata?: { httpStatusCode?: number } };104    if (err.name === "NotFound" || err.name === "NoSuchKey" || err.$metadata?.httpStatusCode === 404) return null;105    throw e;106  }107}108109/** Fetch and decompress a stored raw body. Returns null when the key does not exist. */110export async function getRaw(key: string): Promise<{ body: Buffer; contentType: string | null } | null> {111  try {112    const r = await getS3().send(new GetObjectCommand({ Bucket: bucketName(), Key: key }));113    const bytes = r.Body ? Buffer.from(await r.Body.transformToByteArray()) : Buffer.alloc(0);114    return { body: decompressByKey(key, bytes), contentType: r.Metadata?.["content-type"] ?? null };115  } catch (e) {116    const err = e as { name?: string; $metadata?: { httpStatusCode?: number } };117    if (err.name === "NoSuchKey" || err.name === "NotFound" || err.$metadata?.httpStatusCode === 404) return null;118    throw e;119  }120}121122/** Lightweight connectivity probe for `dci doctor`. */123export async function storageProbe(): Promise<{ ok: boolean; message: string }> {124  try {125    await ensureBucket();126    return { ok: true, message: `bucket ${bucketName()} reachable at ${getEnv().s3Endpoint} (${codec.ext})` };127  } catch (e) {128    return { ok: false, message: (e as Error).message };129  }130}131