spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
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