/** * Raw body archive on S3/MinIO. Keys: raw///..zst * Bodies are compressed with zstd (node:zlib, Node ≥ 22.15) — gzip fallback (".gz") when unavailable. * The same content hash is never uploaded twice (HEAD before PUT). */ import { S3Client, HeadBucketCommand, CreateBucketCommand, PutObjectCommand, GetObjectCommand, HeadObjectCommand } from "@aws-sdk/client-s3"; import * as zlib from "node:zlib"; import { getEnv } from "./env.js"; type Codec = { ext: "zst" | "gz"; compress(b: Buffer): Buffer; decompress(b: Buffer): Buffer }; const zstd = (zlib as unknown as { zstdCompressSync?: (b: Buffer) => Buffer; zstdDecompressSync?: (b: Buffer) => Buffer }); export const codec: Codec = typeof zstd.zstdCompressSync === "function" && typeof zstd.zstdDecompressSync === "function" ? { ext: "zst", compress: (b) => zstd.zstdCompressSync!(b), decompress: (b) => zstd.zstdDecompressSync!(b) } : { ext: "gz", compress: (b) => zlib.gzipSync(b), decompress: (b) => zlib.gunzipSync(b) }; /** Decompress by key suffix so gz-era objects stay readable after a zstd upgrade. */ export function decompressByKey(key: string, body: Buffer): Buffer { if (key.endsWith(".zst")) { if (!zstd.zstdDecompressSync) throw new Error("zstd not available in this Node build"); return zstd.zstdDecompressSync(body); } if (key.endsWith(".gz")) return zlib.gunzipSync(body); return body; } let _client: S3Client | null = null; export function getS3(): S3Client { if (_client) return _client; const env = getEnv(); _client = new S3Client({ endpoint: env.s3Endpoint, region: env.s3Region, forcePathStyle: true, credentials: { accessKeyId: env.s3AccessKey, secretAccessKey: env.s3SecretKey }, requestHandler: { requestTimeout: 30_000, connectionTimeout: 5_000 }, }); return _client; } export function bucketName(): string { return getEnv().s3Bucket; } let bucketReady = false; export async function ensureBucket(): Promise { if (bucketReady) return; const s3 = getS3(); const Bucket = bucketName(); try { await s3.send(new HeadBucketCommand({ Bucket })); } catch (e) { const code = (e as { name?: string; $metadata?: { httpStatusCode?: number } }); if (code.name === "NotFound" || code.name === "NoSuchBucket" || code.$metadata?.httpStatusCode === 404) { await s3.send(new CreateBucketCommand({ Bucket })); } else throw e; } bucketReady = true; } /** Map a content type (and URL hint) to a short file extension. */ export function extensionFor(contentType: string | null | undefined, url?: string): string { const ct = (contentType ?? "").toLowerCase(); if (/pdf/.test(ct) || /\.pdf(\?|$)/i.test(url ?? "")) return "pdf"; if (/json/.test(ct)) return "json"; if (/xml/.test(ct)) return "xml"; if (/html/.test(ct)) return "html"; if (/markdown/.test(ct)) return "md"; if (/csv/.test(ct)) return "csv"; if (/plain/.test(ct)) return "txt"; if (/javascript/.test(ct)) return "js"; if (ct.startsWith("image/")) return ct.split("/")[1]?.split(";")[0] ?? "bin"; return "bin"; } export function rawKey(connectorId: string, docId: string, contentHash: string, ext: string): string { return `raw/${connectorId}/${docId}/${contentHash}.${ext}.${codec.ext}`; } export interface PutResult { key: string; stored: boolean; bytes: number; compressedBytes: number } /** Store a raw body once per content hash. Returns the key (existing or new). */ export async function putRaw(connectorId: string, docId: string, contentHash: string, body: Buffer, contentType: string | null | undefined, url?: string): Promise { await ensureBucket(); const key = rawKey(connectorId, docId, contentHash, extensionFor(contentType, url)); const existing = await headRaw(key); if (existing) return { key, stored: false, bytes: body.length, compressedBytes: existing.size }; const compressed = codec.compress(body); await getS3().send( new PutObjectCommand({ Bucket: bucketName(), Key: key, Body: compressed, ContentType: "application/octet-stream", ContentEncoding: codec.ext === "zst" ? "zstd" : "gzip", Metadata: { "content-type": (contentType ?? "application/octet-stream").slice(0, 200), "content-hash": contentHash, "original-bytes": String(body.length) }, }), ); return { key, stored: true, bytes: body.length, compressedBytes: compressed.length }; } export async function headRaw(key: string): Promise<{ size: number; contentType: string | null } | null> { try { const r = await getS3().send(new HeadObjectCommand({ Bucket: bucketName(), Key: key })); return { size: Number(r.ContentLength ?? 0), contentType: r.Metadata?.["content-type"] ?? null }; } catch (e) { const err = e as { name?: string; $metadata?: { httpStatusCode?: number } }; if (err.name === "NotFound" || err.name === "NoSuchKey" || err.$metadata?.httpStatusCode === 404) return null; throw e; } } /** Fetch and decompress a stored raw body. Returns null when the key does not exist. */ export async function getRaw(key: string): Promise<{ body: Buffer; contentType: string | null } | null> { try { const r = await getS3().send(new GetObjectCommand({ Bucket: bucketName(), Key: key })); const bytes = r.Body ? Buffer.from(await r.Body.transformToByteArray()) : Buffer.alloc(0); return { body: decompressByKey(key, bytes), contentType: r.Metadata?.["content-type"] ?? null }; } catch (e) { const err = e as { name?: string; $metadata?: { httpStatusCode?: number } }; if (err.name === "NoSuchKey" || err.name === "NotFound" || err.$metadata?.httpStatusCode === 404) return null; throw e; } } /** Lightweight connectivity probe for `dci doctor`. */ export async function storageProbe(): Promise<{ ok: boolean; message: string }> { try { await ensureBucket(); return { ok: true, message: `bucket ${bucketName()} reachable at ${getEnv().s3Endpoint} (${codec.ext})` }; } catch (e) { return { ok: false, message: (e as Error).message }; } }