diff --git a/README.md b/README.md index 422344a..045b23c 100644 --- a/README.md +++ b/README.md @@ -217,6 +217,7 @@ gateway 运行时只负责 API 与静态文件服务,不再把复杂 UI 硬写 - `516` 占比 - 看最近请求记录 - 请求时间戳、首字耗时、总耗时、请求体大小、路径、模型、状态码 + - 相同重发请求会带相同的 `request_id` - `usage` 中的 input / output / total / reasoning tokens - 管理 profiles - 新建 / 编辑 profile env diff --git a/gateway.mjs b/gateway.mjs index cad1eb6..1915953 100644 --- a/gateway.mjs +++ b/gateway.mjs @@ -2,6 +2,7 @@ import http from "node:http"; import { spawn } from "node:child_process"; +import { createHash } from "node:crypto"; import { chmod, copyFile, mkdir, readFile, readdir, rm, writeFile } from "node:fs/promises"; import fs from "node:fs"; import path from "node:path"; @@ -456,6 +457,7 @@ function openRequestsDatabase(dbPath) { PRAGMA journal_mode = WAL; CREATE TABLE IF NOT EXISTS requests ( seq INTEGER PRIMARY KEY, + request_id TEXT, started_at TEXT, finished_at TEXT, duration_ms INTEGER, @@ -489,6 +491,12 @@ function openRequestsDatabase(dbPath) { CREATE INDEX IF NOT EXISTS idx_requests_matched ON requests(matched); CREATE INDEX IF NOT EXISTS idx_requests_response_stream ON requests(response_stream); `); + const requestColumns = db.prepare("PRAGMA table_info(requests)").all(); + const requestColumnNames = new Set(requestColumns.map((column) => column.name)); + if (!requestColumnNames.has("request_id")) { + db.exec("ALTER TABLE requests ADD COLUMN request_id TEXT"); + } + db.exec("CREATE INDEX IF NOT EXISTS idx_requests_request_id ON requests(request_id)"); return db; } @@ -507,10 +515,19 @@ function buildPersistedRequestPayload(entry) { ); } +function computeRequestId(pathname, rawBody) { + const hash = createHash("sha256"); + hash.update(normalizePath(pathname)); + hash.update("\n"); + hash.update(Buffer.isBuffer(rawBody) ? rawBody : Buffer.from(rawBody || "")); + return `req_${hash.digest("hex").slice(0, 16)}`; +} + function requestRowFromEntry(entry) { const payload = buildPersistedRequestPayload(entry); return { seq: entry.seq, + request_id: entry.request_id || null, started_at: entry.started_at || null, finished_at: entry.finished_at || null, duration_ms: Number.isInteger(entry.duration_ms) ? entry.duration_ms : null, @@ -543,13 +560,13 @@ function requestRowFromEntry(entry) { function insertRequestRow(db, row) { db.prepare(` INSERT OR REPLACE INTO requests ( - seq, started_at, finished_at, duration_ms, profile_name, method, path, model, + seq, request_id, started_at, finished_at, duration_ms, profile_name, method, path, model, requested_model, forwarded_model, request_stream, response_stream, inspected, matched, status_code, upstream_status_code, reasoning_tokens, input_tokens, output_tokens, total_tokens, cached_tokens, error, upstream_origin, upstream_path, upstream_auth_mode, upstream_auth_source, payload_json ) VALUES ( - @seq, @started_at, @finished_at, @duration_ms, @profile_name, @method, @path, @model, + @seq, @request_id, @started_at, @finished_at, @duration_ms, @profile_name, @method, @path, @model, @requested_model, @forwarded_model, @request_stream, @response_stream, @inspected, @matched, @status_code, @upstream_status_code, @reasoning_tokens, @input_tokens, @output_tokens, @total_tokens, @cached_tokens, @error, @upstream_origin, @upstream_path, @@ -609,6 +626,7 @@ function buildRequestQueryFilters({ query, filter }) { lower(coalesce(profile_name, '')) LIKE @query OR lower(coalesce(method, '')) LIKE @query OR lower(coalesce(path, '')) LIKE @query OR + lower(coalesce(request_id, '')) LIKE @query OR lower(coalesce(model, '')) LIKE @query OR lower(coalesce(requested_model, '')) LIKE @query OR lower(coalesce(forwarded_model, '')) LIKE @query OR @@ -1203,6 +1221,7 @@ async function buildPersistentRequestsSnapshot(runtime, { limit = 50, offset = 0 function buildRequestEntry({ seq, startedAt, startedMs, req, pathname, requestJson, profileName }) { return { seq, + request_id: null, lifecycle_state: "sent", started_at: startedAt.toISOString(), first_response_at: null, @@ -3060,6 +3079,7 @@ async function proxyRequest(runtime, req, res) { const requestIsStream = Boolean(requestJson?.stream); let totalUpstreamAttempts = 0; requestEntry.request_body_bytes = rawRequestBody.length; + requestEntry.request_id = computeRequestId(pathname, rawRequestBody); requestEntry.model = requestJson?.model || null; requestEntry.requested_model = parsedRequestJson?.model || null; requestEntry.forwarded_model = forwardedModel || parsedRequestJson?.model || null; diff --git a/scripts/test-gateway-e2e.mjs b/scripts/test-gateway-e2e.mjs index 9a67a18..2d952d3 100644 --- a/scripts/test-gateway-e2e.mjs +++ b/scripts/test-gateway-e2e.mjs @@ -459,6 +459,77 @@ async function run() { recoveredEntry?.request_body_bytes === Buffer.byteLength(recoveredPayload), `请求体大小记录异常: ${recoveredEntry?.request_body_bytes}`, ); + assert(recoveredEntry?.request_id, "请求记录未生成 request_id"); + + const sameRequestPayload = JSON.stringify({ test_reasoning_tokens: 128, test_request_id_marker: "same" }); + const sameRequestFirstResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: sameRequestPayload, + }); + assert(sameRequestFirstResponse.status === 200, `相同请求首次发送失败: ${sameRequestFirstResponse.status}`); + await sameRequestFirstResponse.json(); + + const sameRequestSecondResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: sameRequestPayload, + }); + assert(sameRequestSecondResponse.status === 200, `相同请求二次发送失败: ${sameRequestSecondResponse.status}`); + await sameRequestSecondResponse.json(); + + const differentBodyPayload = JSON.stringify({ test_reasoning_tokens: 128, test_request_id_marker: "different" }); + const differentBodyResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: differentBodyPayload, + }); + assert(differentBodyResponse.status === 200, `不同请求体发送失败: ${differentBodyResponse.status}`); + await differentBodyResponse.json(); + + const differentPathResponse = await fetch(`http://127.0.0.1:${gatewayPort}/v1/responses`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: sameRequestPayload, + }); + assert(differentPathResponse.status === 200, `不同路径发送失败: ${differentPathResponse.status}`); + await differentPathResponse.json(); + + const requestIdRequestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=40`); + const requestIdRequestsPayload = await requestIdRequestsResponse.json(); + const sameRequestEntries = (requestIdRequestsPayload?.entries || []).filter( + (entry) => entry.path === "/responses" && entry.request_body_bytes === Buffer.byteLength(sameRequestPayload), + ); + assert(sameRequestEntries.length >= 2, "未找到两条相同请求记录"); + const [sameRequestEntryA, sameRequestEntryB] = sameRequestEntries; + assert(sameRequestEntryA.request_id, "相同请求记录缺少 request_id"); + assert(sameRequestEntryA.request_id === sameRequestEntryB.request_id, "相同请求未复用 request_id"); + assert(sameRequestEntryA.seq !== sameRequestEntryB.seq, "相同请求不应复用 seq"); + + const differentBodyEntry = (requestIdRequestsPayload?.entries || []).find( + (entry) => + entry.path === "/responses" && + entry.request_body_bytes !== Buffer.byteLength(sameRequestPayload) && + entry.request_body_bytes === Buffer.byteLength(differentBodyPayload), + ); + assert(differentBodyEntry?.request_id, "不同请求体记录缺少 request_id"); + assert(differentBodyEntry.request_id !== sameRequestEntryA.request_id, "不同请求体错误复用 request_id"); + + const differentPathEntry = (requestIdRequestsPayload?.entries || []).find( + (entry) => entry.path === "/v1/responses" && entry.request_body_bytes === Buffer.byteLength(sameRequestPayload), + ); + assert(differentPathEntry?.request_id, "不同路径记录缺少 request_id"); + assert(differentPathEntry.request_id !== sameRequestEntryA.request_id, "不同路径错误复用 request_id"); + + const requestIdQueryResponse = await fetch( + `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(sameRequestEntryA.request_id)}`, + ); + const requestIdQueryPayload = await requestIdQueryResponse.json(); + assert(requestIdQueryResponse.status === 200, `request_id 搜索失败: ${requestIdQueryResponse.status}`); + assert( + (requestIdQueryPayload?.entries || []).some((entry) => entry.request_id === sameRequestEntryA.request_id), + "request_id 搜索未命中对应记录", + ); const capacityResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", @@ -510,6 +581,7 @@ async function run() { (entry) => entry.path === "/responses" && entry.status_code === 200 && entry.upstream_attempt_count >= 3, ); assert(capacityRecoveredEntry, "capacity 抖动恢复后的请求记录未保留重试次数"); + assert(capacityRecoveredEntry.request_id, "capacity 恢复后的请求记录缺少 request_id"); const capacityStatus200RecoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", diff --git a/ui-src/src/App.tsx b/ui-src/src/App.tsx index 588b9f7..916caf7 100644 --- a/ui-src/src/App.tsx +++ b/ui-src/src/App.tsx @@ -70,6 +70,7 @@ type Usage = { type RequestEntry = { seq: number; + request_id?: string | null; lifecycle_state?: string | null; started_at?: string; first_response_at?: string | null; @@ -1104,6 +1105,7 @@ export default function App() {
{`${entry.method || "-"} ${entry.path || "-"}`}
{entry.response_stream ? "stream" : "non-stream"}
+ {entry.request_id || "-"}