From f15e18cf05d6ec3b47920cf46d59d3e2e3de132f Mon Sep 17 00:00:00 2001 From: redshift <213178690+1ftredsh@users.noreply.github.com> Date: Sat, 20 Jun 2026 17:00:52 +0800 Subject: [PATCH] Adapt request response logging to SDK sink --- src/daemon/http/index.ts | 9 +- src/daemon/index.ts | 9 +- src/daemon/request-response-log-sink.ts | 269 ++++++++++++++++++++++++ 3 files changed, 282 insertions(+), 5 deletions(-) create mode 100644 src/daemon/request-response-log-sink.ts diff --git a/src/daemon/http/index.ts b/src/daemon/http/index.ts index 2cfc503..b66c3fa 100644 --- a/src/daemon/http/index.ts +++ b/src/daemon/http/index.ts @@ -7,6 +7,7 @@ import { ProviderManager, } from "@routstr/sdk"; import type { UsageTrackingDriver, SdkLogger } from "@routstr/sdk"; +import type { RequestResponseLogSink } from "../request-response-log-sink"; import { logger } from "../../utils/logger"; import { loadDaemonConfig, saveDaemonConfig } from "../config-store"; import { @@ -48,7 +49,7 @@ type DaemonDeps = { routstrPubkey?: string; providerManager: ProviderManager; refundClient: any; - requestResponseLogDir?: string; + requestResponseLogSink?: RequestResponseLogSink; }; /** @@ -304,7 +305,7 @@ export function createDaemonRequestHandler(deps: { usageTrackingDriver: UsageTrackingDriver; providerManager: ProviderManager; refundClient: any; - requestResponseLogDir?: string; + requestResponseLogSink?: RequestResponseLogSink; }) { return async function handler(req: IncomingMessage, res: ServerResponse) { const host = req.headers.host || "localhost"; @@ -1276,8 +1277,8 @@ export function createDaemonRequestHandler(deps: { sdkStore: deps.store, providerManager: deps.providerManager, logger: reqLogger, - ...(deps.requestResponseLogDir - ? { requestResponseLogDir: deps.requestResponseLogDir } + ...(deps.requestResponseLogSink + ? { requestResponseLogSink: deps.requestResponseLogSink } : {}), ...(deps.routstrPubkey ? { routstrPubkey: deps.routstrPubkey } : {}), }); diff --git a/src/daemon/index.ts b/src/daemon/index.ts index f7538a7..9875400 100644 --- a/src/daemon/index.ts +++ b/src/daemon/index.ts @@ -46,6 +46,7 @@ import type { AutoRefillConfig } from "./wallet/auto-refill"; import { createCocodClient } from "./wallet/cocod-client"; import { createModelService } from "./models"; import { createDaemonRequestHandler } from "./http"; +import { FileRequestResponseLogSink } from "./request-response-log-sink"; import { refreshModelsAndIntegrations } from "../integrations"; import { RoutstrClient } from "@routstr/sdk"; @@ -60,6 +61,12 @@ async function main(): Promise { (config.requestResponseLogging?.enabled ? config.requestResponseLogging.dir || REQUEST_RESPONSE_LOGS_DIR : undefined); + const requestResponseLogSink = requestResponseLogDir + ? new FileRequestResponseLogSink({ + dir: requestResponseLogDir, + logger: daemonSdkLogger.child("request-response-log"), + }) + : undefined; await ensureDirs(); @@ -144,7 +151,7 @@ async function main(): Promise { usageTrackingDriver, providerManager, refundClient, - requestResponseLogDir, + requestResponseLogSink, }), ); diff --git a/src/daemon/request-response-log-sink.ts b/src/daemon/request-response-log-sink.ts new file mode 100644 index 0000000..19350b3 --- /dev/null +++ b/src/daemon/request-response-log-sink.ts @@ -0,0 +1,269 @@ +import { createWriteStream, mkdirSync, type WriteStream } from "fs"; +import { writeFile } from "fs/promises"; +import { join } from "path"; +import type { SdkLogger } from "@routstr/sdk"; + +export interface RequestResponseLogRequestInput { + method: string; + url: string; + path: string; + baseUrl: string; + headers: Record; + body?: unknown; + rawBody?: string; +} + +export interface RequestResponseLogSink { + logRequest?(input: RequestResponseLogRequestInput): string | undefined | Promise; + logResponseStart?(id: string | undefined, response: Response): void | Promise; + logResponseChunk?(id: string | undefined, sequence: number, text: string): void | Promise; + logResponseEnd?(id: string | undefined): void | Promise; + logResponseError?(id: string | undefined, error: unknown): void | Promise; + logResponseBody?(id: string | undefined, response: Response): void | Promise; +} + +interface ActiveResponseLog { + stream: WriteStream; + pending: Promise; +} + +export interface FileRequestResponseLogSinkOptions { + dir: string; + logger?: SdkLogger; +} + +const SENSITIVE_HEADER_NAMES = new Set([ + "authorization", + "x-cashu", + "cookie", + "set-cookie", + "proxy-authorization", +]); + +const SENSITIVE_BODY_FIELD_NAMES = new Set([ + "authorization", + "api_key", + "apikey", + "apiKey", + "access_token", + "accessToken", + "bearer", + "cashu", + "cookie", + "key", + "password", + "secret", + "token", + "x-cashu", +]); + +const REDACTED = "[REDACTED]"; + +const sanitizeForFilename = (value: string): string => + value.replace(/[^a-zA-Z0-9._-]+/g, "-").replace(/^-+|-+$/g, ""); + +const makeId = (): string => { + const timestamp = new Date().toISOString().replace(/[:.]/g, "-"); + const random = crypto.randomUUID().slice(0, 8); + return `${timestamp}-${random}`; +}; + +const headersToObject = (headers: Headers): Record => { + const out: Record = {}; + headers.forEach((value, key) => { + out[key] = value; + }); + return out; +}; + +const redactHeaders = (headers: Record): Record => + Object.fromEntries( + Object.entries(headers).map(([key, value]) => [ + key, + SENSITIVE_HEADER_NAMES.has(key.toLowerCase()) ? REDACTED : value, + ]), + ); + +const redactBody = (value: unknown): unknown => { + if (!value || typeof value !== "object") return value; + if (Array.isArray(value)) return value.map(redactBody); + + return Object.fromEntries( + Object.entries(value as Record).map(([key, entry]) => [ + key, + SENSITIVE_BODY_FIELD_NAMES.has(key) || SENSITIVE_BODY_FIELD_NAMES.has(key.toLowerCase()) + ? REDACTED + : redactBody(entry), + ]), + ); +}; + +const redactRawBody = (rawBody: string | undefined): string | undefined => { + if (!rawBody) return undefined; + try { + return JSON.stringify(redactBody(JSON.parse(rawBody))); + } catch { + return rawBody; + } +}; + +export class FileRequestResponseLogSink implements RequestResponseLogSink { + private requestsDir: string; + private responsesDir: string; + private activeResponses = new Map(); + + constructor(private options: FileRequestResponseLogSinkOptions) { + this.requestsDir = join(options.dir, "requests"); + this.responsesDir = join(options.dir, "responses"); + mkdirSync(this.requestsDir, { recursive: true }); + mkdirSync(this.responsesDir, { recursive: true }); + } + + async logRequest(input: RequestResponseLogRequestInput): Promise { + try { + const id = makeId(); + const filePath = join(this.requestsDir, `${sanitizeForFilename(id)}.json`); + await writeFile( + filePath, + JSON.stringify( + { + id, + timestamp: new Date().toISOString(), + method: input.method, + url: input.url, + path: input.path, + baseUrl: input.baseUrl, + headers: redactHeaders(input.headers), + body: redactBody(input.body), + rawBody: redactRawBody(input.rawBody), + }, + null, + 2, + ), + ); + return id; + } catch (error) { + this.options.logger?.error?.("[request-response-log] failed to log request:", error); + return undefined; + } + } + + async logResponseStart(id: string | undefined, response: Response): Promise { + if (!id) return; + await this.append(id, { + type: "response_start", + status: response.status, + statusText: response.statusText, + headers: redactHeaders(headersToObject(response.headers)), + }); + } + + logResponseChunk(id: string | undefined, sequence: number, text: string): void { + if (!id) return; + void this.append(id, { + type: "chunk", + sequence, + text, + }); + } + + async logResponseEnd(id: string | undefined): Promise { + if (!id) return; + await this.append(id, { type: "end" }); + await this.close(id); + } + + async logResponseError(id: string | undefined, error: unknown): Promise { + if (!id) return; + await this.append(id, { + type: "error", + error: error instanceof Error ? { message: error.message, stack: error.stack } : String(error), + }); + await this.close(id); + } + + async logResponseBody(id: string | undefined, response: Response): Promise { + if (!id) return; + try { + if (!response.body) { + await this.logResponseEnd(id); + return; + } + + const reader = response.body.getReader(); + const decoder = new TextDecoder("utf-8"); + let sequence = 0; + + while (true) { + const { done, value } = await reader.read(); + if (done) break; + if (value && value.byteLength > 0) { + await this.append(id, { + type: "chunk", + sequence: sequence++, + text: decoder.decode(value, { stream: true }), + }); + } + } + + const tail = decoder.decode(); + if (tail) { + await this.append(id, { + type: "chunk", + sequence: sequence++, + text: tail, + }); + } + + await this.logResponseEnd(id); + } catch (error) { + await this.logResponseError(id, error); + } + } + + private getOrCreate(id: string): ActiveResponseLog { + const existing = this.activeResponses.get(id); + if (existing) return existing; + + const filePath = join(this.responsesDir, `${sanitizeForFilename(id)}.jsonl`); + const stream = createWriteStream(filePath, { flags: "a" }); + const active: ActiveResponseLog = { + stream, + pending: Promise.resolve(), + }; + + stream.on("error", (error) => { + this.options.logger?.error?.("[request-response-log] response log stream error:", error); + }); + + this.activeResponses.set(id, active); + return active; + } + + private async append(id: string, event: Record): Promise { + try { + const active = this.getOrCreate(id); + const line = JSON.stringify({ requestLogId: id, timestamp: new Date().toISOString(), ...event }) + "\n"; + active.pending = active.pending.then( + () => + new Promise((resolve, reject) => { + active.stream.write(line, (error) => { + if (error) reject(error); + else resolve(); + }); + }), + ); + await active.pending; + } catch (error) { + this.options.logger?.error?.("[request-response-log] failed to append response event:", error); + } + } + + private async close(id: string): Promise { + const active = this.activeResponses.get(id); + if (!active) return; + this.activeResponses.delete(id); + await active.pending; + await new Promise((resolve) => active.stream.end(resolve)); + } +}