From 871d0ac5cbe412b49037befdfb27add89ab6a376 Mon Sep 17 00:00:00 2001 From: czj <13261895355@163.com> Date: Thu, 16 Jul 2026 12:05:32 +0800 Subject: [PATCH] =?UTF-8?q?=E6=96=B0=E5=A2=9ELLM=E8=BF=BD=E8=B8=AA?= =?UTF-8?q?=E6=AD=A3=E6=96=87=E5=AD=98=E5=82=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../llm/__tests__/trace-payload-store.test.ts | 98 ++++++++++++++++ src/lib/llm/trace-payload-store.ts | 111 ++++++++++++++++++ 2 files changed, 209 insertions(+) create mode 100644 src/lib/llm/__tests__/trace-payload-store.test.ts create mode 100644 src/lib/llm/trace-payload-store.ts diff --git a/src/lib/llm/__tests__/trace-payload-store.test.ts b/src/lib/llm/__tests__/trace-payload-store.test.ts new file mode 100644 index 0000000..8fb2a3a --- /dev/null +++ b/src/lib/llm/__tests__/trace-payload-store.test.ts @@ -0,0 +1,98 @@ +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +import { + createLocalTracePayloadStore, + createR2TracePayloadStore, +} from "../trace-payload-store"; + +describe("LLM trace payload stores", () => { + let tempDir: string; + + beforeEach(() => { + tempDir = mkdtempSync(join(tmpdir(), "geo-llm-payloads-")); + }); + + afterEach(() => { + rmSync(tempDir, { recursive: true, force: true }); + }); + + it("round-trips exact JSON locally and deletes one job prefix", async () => { + const store = createLocalTracePayloadStore(tempDir); + const payload = { + model: "deepseek-v4-pro", + messages: [{ role: "user", content: "原文" }], + }; + const key = "llm-traces/job_1/llmcall_1/request.json"; + + await store.putJson(key, payload); + await expect(store.getJson(key)).resolves.toEqual(payload); + await store.deleteJob("job_1"); + await expect(store.getJson(key)).resolves.toBeNull(); + }); + + it("stores private JSON in R2 without a public URL", async () => { + const put = vi.fn().mockResolvedValue(undefined); + const get = vi.fn().mockResolvedValue({ + json: async () => ({ ok: true }), + }); + const list = vi.fn().mockResolvedValue({ objects: [], truncated: false }); + const deleteObjects = vi.fn().mockResolvedValue(undefined); + const bucket = { + put, + get, + list, + delete: deleteObjects, + } as unknown as R2Bucket; + const store = createR2TracePayloadStore(bucket); + + await store.putJson( + "llm-traces/job_1/llmcall_1/request.json", + { ok: true }, + ); + await expect( + store.getJson("llm-traces/job_1/llmcall_1/request.json"), + ).resolves.toEqual({ ok: true }); + expect(put).toHaveBeenCalledWith( + "llm-traces/job_1/llmcall_1/request.json", + JSON.stringify({ ok: true }), + { httpMetadata: { contentType: "application/json; charset=utf-8" } }, + ); + }); + + it("deletes every paginated R2 object under one job prefix", async () => { + const list = vi + .fn() + .mockResolvedValueOnce({ + objects: [{ key: "llm-traces/job_1/call_1/request.json" }], + truncated: true, + cursor: "next-page", + }) + .mockResolvedValueOnce({ + objects: [{ key: "llm-traces/job_1/call_1/response.json" }], + truncated: false, + }); + const deleteObjects = vi.fn().mockResolvedValue(undefined); + const bucket = { list, delete: deleteObjects } as unknown as R2Bucket; + + await createR2TracePayloadStore(bucket).deleteJob("job_1"); + + expect(list).toHaveBeenNthCalledWith(1, { + prefix: "llm-traces/job_1/", + cursor: undefined, + }); + expect(list).toHaveBeenNthCalledWith(2, { + prefix: "llm-traces/job_1/", + cursor: "next-page", + }); + expect(deleteObjects).toHaveBeenNthCalledWith(1, [ + "llm-traces/job_1/call_1/request.json", + ]); + expect(deleteObjects).toHaveBeenNthCalledWith(2, [ + "llm-traces/job_1/call_1/response.json", + ]); + }); +}); diff --git a/src/lib/llm/trace-payload-store.ts b/src/lib/llm/trace-payload-store.ts new file mode 100644 index 0000000..0136edc --- /dev/null +++ b/src/lib/llm/trace-payload-store.ts @@ -0,0 +1,111 @@ +import { + existsSync, + mkdirSync, + readFileSync, + rmSync, + writeFileSync, +} from "node:fs"; +import { dirname, join, relative, resolve } from "node:path"; + +import { getAppDataDir } from "../db/connection"; +import { getAppCloudflareEnv } from "../runtime/cloudflare"; + +const JSON_CONTENT_TYPE = "application/json; charset=utf-8"; + +export interface LlmTracePayloadStore { + putJson(key: string, value: unknown): Promise; + getJson(key: string): Promise; + deleteJob(jobId: string): Promise; +} + +function assertSafeSegment(value: string, label: string) { + if (!/^[A-Za-z0-9_-]+$/.test(value)) { + throw new Error(`Invalid LLM trace ${label}`); + } +} + +export function tracePayloadKey( + jobId: string, + callId: string, + kind: "request" | "response", +) { + assertSafeSegment(jobId, "job id"); + assertSafeSegment(callId, "call id"); + return `llm-traces/${jobId}/${callId}/${kind}.json`; +} + +function localPathForKey(dataDir: string, key: string) { + const traceRoot = resolve(dataDir, "llm-traces"); + const path = resolve(dataDir, key); + const pathFromTraceRoot = relative(traceRoot, path); + if ( + pathFromTraceRoot.startsWith("..") || + pathFromTraceRoot === "" || + key.startsWith("/") + ) { + throw new Error("Invalid LLM trace payload key"); + } + return path; +} + +export function createLocalTracePayloadStore( + dataDir = getAppDataDir(), +): LlmTracePayloadStore { + return { + async putJson(key, value) { + const path = localPathForKey(dataDir, key); + mkdirSync(dirname(path), { recursive: true }); + writeFileSync(path, JSON.stringify(value), "utf8"); + }, + async getJson(key) { + const path = localPathForKey(dataDir, key); + if (!existsSync(path)) return null; + return JSON.parse(readFileSync(path, "utf8")) as unknown; + }, + async deleteJob(jobId) { + assertSafeSegment(jobId, "job id"); + rmSync(join(dataDir, "llm-traces", jobId), { + recursive: true, + force: true, + }); + }, + }; +} + +export function createR2TracePayloadStore( + bucket: R2Bucket, +): LlmTracePayloadStore { + return { + async putJson(key, value) { + await bucket.put(key, JSON.stringify(value), { + httpMetadata: { contentType: JSON_CONTENT_TYPE }, + }); + }, + async getJson(key) { + const object = await bucket.get(key); + return object ? await object.json() : null; + }, + async deleteJob(jobId) { + assertSafeSegment(jobId, "job id"); + const prefix = `llm-traces/${jobId}/`; + let cursor: string | undefined; + do { + const page = await bucket.list({ prefix, cursor }); + const keys = page.objects.map((object) => object.key); + if (keys.length > 0) await bucket.delete(keys); + cursor = page.truncated ? page.cursor : undefined; + } while (cursor); + }, + }; +} + +export function getLlmTracePayloadStoreFromRuntime(): LlmTracePayloadStore { + if (process.env.APP_RUNTIME === "cloudflare") { + const bucket = getAppCloudflareEnv()?.EXPORT_BUCKET; + if (!bucket) { + throw new Error("Cloudflare R2 binding EXPORT_BUCKET is required"); + } + return createR2TracePayloadStore(bucket); + } + return createLocalTracePayloadStore(getAppDataDir()); +}