# TokenThief newapi 反向代理,采集请求/响应(包括 SSE/chunked 流式响应)并异步批量写入 ClickHouse。 ## 特性 - 透明反代任意 HTTP 后端(默认目标为 newapi)。 - 流式响应(SSE、chunked、OpenAI 兼容 chat completions)边转发边缓冲,结束后整体入库。 - 异步队列 + 批量写入,主路径零阻塞。 - **DB 故障不影响代理服务**:连接失败时丢弃日志,后台持续重连。 - 通过 yaml 配置 glob 风格的黑白名单。 - 原样存储 headers/body(String 字段保存 JSON 文本与 body 内容)。 ## 快速开始 ### 本地运行 ```bash cp .env.example .env # 编辑 .env 设置 UPSTREAM_URL / CLICKHOUSE_URL set -a; source .env; set +a go run . ``` ### Docker ```bash docker build -t tokenthief . docker run --rm -p 8080:8080 \ -e UPSTREAM_URL=http://newapi:3000 \ -e CLICKHOUSE_URL='clickhouses://tokenthief:@clickhouse:9440/tokenthief' \ -v $(pwd)/filter.yaml:/etc/tokenthief/filter.yaml \ tokenthief ``` 镜像内二进制路径为 `/usr/local/bin/TokenThief`。 ### Docker Compose `compose.yml` 已包含 TokenThief 与 ClickHouse: ```bash cp .env.example .env # 编辑 .env,设置 UPSTREAM_URL、强密码 CLICKHOUSE_PASSWORD 和对应的 CLICKHOUSE_URL docker compose -f compose.yml up -d --build ``` 默认端口: | 服务 | 地址 | |---|---| | TokenThief | `http://localhost:8080` | | ClickHouse | 仅 Compose 内部网络,不默认发布宿主机端口 | Compose 在渲染配置时要求显式提供 `UPSTREAM_URL`、应用 DSN 和非空 ClickHouse 密码。密码中的 URL 特殊字符必须编码: ```text clickhouse://tokenthief:@thief_clickhouse:9000/tokenthief ``` ### 运行测试 测试代码主要集中在 `tests/` 目录,数据库包还包含同包测试。 ```bash gofmt -w . go test ./... go vet ./... ``` 构建二进制时不会引入测试内容:`_test.go` 不参与 `go build`,`tests/scripts/` 也不被主程序 import;Docker 构建时 `.dockerignore` 会把整个 `tests/` 目录排除在 build context 之外。 ## 环境变量 | 变量 | 说明 | 默认 | |---|---|---| | `LISTEN_ADDR` | 监听地址 | `:8080` | | `UPSTREAM_URL` | 后端地址(必填) | - | | `UPSTREAM_TLS_INSECURE_SKIP_VERIFY` | 跳过上游 HTTPS 证书校验;仅开发/可信内网自签证书场景使用 | `false` | | `CLICKHOUSE_URL` | ClickHouse 原生协议地址(必填);支持严格 `host:port`、`clickhouse://` 或 TLS `clickhouses://` | - | | `MAX_BODY_BYTES` | 单个请求体和响应体的记录上限 | `1048576` | | `MAX_REQUEST_BYTES` | 单个代理请求体硬上限,超出时返回 413 | `16777216` | | `LOG_QUEUE_SIZE` | 异步队列条数上限 | `256` | | `LOG_QUEUE_BYTES` | 异步队列总字节预算 | `67108864` | | `LOG_BATCH_SIZE` | 批量写入条数 | `50` | | `LOG_BATCH_INTERVAL` | 批量刷新间隔 | `2s` | | `LOG_WORKERS` | worker 数 | `2` | | `FILTER_FILE` | 黑白名单文件 | `./filter.yaml` | | `DB_RECONNECT_INTERVAL` | DB 重连间隔 | `10s` | | `READ_TIMEOUT` | 请求读取总超时 | `30s` | | `WRITE_TIMEOUT` | 响应写入总超时 | `10m` | | `IDLE_TIMEOUT` | HTTP keep-alive 空闲超时 | `5m` | | `UPSTREAM_TIMEOUT` | 上游连接、TLS 握手、响应头等待超时 | `30s` | | `UPSTREAM_RESPONSE_TIMEOUT` | 普通上游响应体总超时 | `30s` | | `UPSTREAM_STREAM_IDLE_TIMEOUT` | SSE 上游响应体空闲超时,每次成功读取后重置 | `2m` | | `TRUSTED_PROXIES` | 可信反向代理 IP/CIDR,逗号分隔;为空时忽略转发头 | 空 | ## 黑白名单(filter.yaml) `CLICKHOUSE_URL` 中的账号、密码和数据库名会传给 ClickHouse 原生协议连接。密码如果包含 `@`、`:`、`/`、`#` 等 URL 特殊字符,需要先做 URL encode。仅填写 `host:port` 时会使用 ClickHouse 默认用户、空密码和默认数据库。 默认 `filter.yaml` 已按 [newapi 官方文档](https://docs.newapi.pro/zh/docs/api) 列出全部 AI 模型接口(chat、completions、embeddings、moderations、rerank、realtime、audio、images、videos、Claude、Gemini 等)。 ```yaml mode: whitelist # whitelist | blacklist | disabled patterns: - /v1/chat/completions - /v1/audio/* # 单段通配 - /v1/videos/** # 跨段通配 - /v1beta/models/*:generateContent # 段内通配 ``` Pattern 语法: | 通配符 | 含义 | |---|---| | `*` | 匹配单个路径段内除 `/` 之外的任意字符 | | `**` | 跨段匹配任意字符(包括 `/`) | | `?` | 匹配单个非 `/` 字符 | ## 端点 - `/healthz` — 始终返回 200,不参与反代与日志。 - 其余路径 — 全部反代到 `UPSTREAM_URL`。 ## 数据库表 启动时自动创建单张 ClickHouse MergeTree 表 `proxy_logs`(详见 `db/migrate.go`),按 `started_at` 月份分区。DDL 如下: ```sql CREATE TABLE IF NOT EXISTS proxy_logs ( request_id String, method String, path String, query String, client_ip String, request_headers String, request_body String, request_truncated Bool DEFAULT false, status_code Int32, response_headers String, response_body String, response_truncated Bool DEFAULT false, is_stream Bool DEFAULT false, latency_ms Int64, started_at DateTime64(3), finished_at DateTime64(3), error String ) ENGINE = MergeTree PARTITION BY toYYYYMM(started_at) ORDER BY (started_at, request_id); ``` 完整字段说明如下。 ### 字段一览 | 字段 | 类型 | 可空 | 说明 | |---|---|---|---| | `request_id` | `String` | 否 | 每次请求由代理生成的 16 字节随机 hex(32 个字符),同时写入响应头 `X-Request-Id`,便于客户端日志关联 | | `method` | `String` | 否 | HTTP 方法(`GET` / `POST` / …),与客户端实际发送一致 | | `path` | `String` | 否 | 请求路径,不含 query string,例如 `/v1/chat/completions` | | `query` | `String` | 否 | 原始 query string(不含 `?`),例如 `model=gpt-4&stream=true`;无 query 时为空串 | | `client_ip` | `String` | 否 | 默认记录 TCP 对端;仅直接对端命中 `TRUSTED_PROXIES` 时,从右向左解析 `X-Forwarded-For`,无 XFF 时使用有效的 `X-Real-IP` | | `request_headers` | `String` | 否 | 完整请求头,序列化为 `{"Header-Name": ["value1", "value2"], ...}` 的 JSON 对象;**注意 Authorization、Cookie、API-Key 等敏感头未脱敏**,按设计原样存储 | | `request_body` | `String` | 否 | 请求体内容。最多保留 `MAX_BODY_BYTES`(默认 1 MiB)字节,超出部分丢弃;读取失败时请求不会转发 | | `request_truncated` | `Bool` | 否 | 请求体是否被 `MAX_BODY_BYTES` 截断。截断只影响数据库记录;请求体超过 `MAX_REQUEST_BYTES` 时不会转发并返回 413 | | `status_code` | `Int32` | 否 | 上游返回的 HTTP 状态码。`502` 通常意味着上游连接失败;WebSocket 成功升级记录为 `101` | | `response_headers` | `String` | 否 | 响应头,结构同 `request_headers`。对于 SSE,会包含 `Content-Type: text/event-stream` 等 | | `response_body` | `String` | 否 | 响应体内容。对于未截断的 SSE,代理会尝试将 OpenAI Completions、Chat Completions、Responses、Anthropic 或 Gemini 事件组装为单个 JSON,并以该 JSON 替换原始 SSE 字节流;无法识别或组装失败时保留原始 SSE。超过 `MAX_BODY_BYTES` 时不组装,仅保留原始字节流的前 `MAX_BODY_BYTES` 字节 | | `response_truncated` | `Bool` | 否 | 响应体是否被截断。截断只影响数据库存储,客户端始终收到完整数据 | | `is_stream` | `Bool` | 否 | 是否 SSE 响应,仅接受媒体类型 `text/event-stream` | | `latency_ms` | `Int64` | 否 | 端到端耗时(毫秒),从代理接收到请求到响应完成。对流式响应 = 从首请求到最后一个 chunk 发出 | | `started_at` | `DateTime64(3)` | 否 | 代理接收到请求的时刻;ClickHouse 按 `toYYYYMM(started_at)` 月度分区 | | `finished_at` | `DateTime64(3)` | 否 | 响应完全写回客户端(包括所有 chunk)的时刻 | | `error` | `String` | 否 | 仅在反代过程中出现错误时填充。常见值:上游不可达、超时、读取请求体失败等 | ### 常用查询示例 ```sql -- 查看最近 20 次失败请求 SELECT started_at, path, status_code, error FROM proxy_logs WHERE status_code >= 400 OR error != '' ORDER BY started_at DESC LIMIT 20; -- 查看某次请求的完整内容 SELECT request_id, method, path, request_body AS req_text, response_body AS resp_text, latency_ms, is_stream FROM proxy_logs WHERE request_id = '0123456789abcdef0123456789abcdef'; -- 按模型统计调用量(从请求体里提取 JSON 字段) SELECT JSONExtractString(request_body, 'model') AS model, count(*) AS calls, toInt32(avg(latency_ms)) AS avg_ms FROM proxy_logs WHERE path = '/v1/chat/completions' AND started_at > now() - INTERVAL 1 DAY GROUP BY 1 ORDER BY calls DESC; -- 查 Authorization(注意:敏感信息) SELECT JSONExtractRaw(request_headers, 'Authorization') FROM proxy_logs LIMIT 5; ``` ### 注意事项 - **敏感信息**:请求头中的 `Authorization`、`Cookie`、`X-Api-Key` 等**未脱敏**。如需脱敏请在 `proxy/capture.go` 的 `headersJSON` 中改造,或对数据库做列级权限控制。 - **body 编码**:ClickHouse 以 `String` 保存 body 内容;文本接口可直接查询,二进制或压缩内容需按业务格式离线解析。 - **WebSocket**:只记录 101 握手元数据,`response_body` 为空;WS 帧内容不采集,进程关闭时会关闭受管升级连接。 - **截断**:`request_truncated` / `response_truncated` 为 `true` 时,对应 `*_body` 仅包含前 `MAX_BODY_BYTES` 字节。需保留完整内容请调高 `MAX_BODY_BYTES`,但要警惕数据库膨胀。 ## 设计要点 - 转发前最多读取 `MAX_REQUEST_BYTES + 1` 字节,超限返回 413,读取或关闭失败时不向上游发送请求;数据库仅记录前 `MAX_BODY_BYTES` 字节,随后通过 `bytes.Reader` 重放完整请求体。 - `httputil.ReverseProxy` + `FlushInterval = -1`,自定义 `ResponseWriter` 同时实现 `Flusher`/`Hijacker`,写入时先转发再缓冲,保证流式实时性。 - 日志通过非阻塞 channel 投递,队列满或 DB 不健康时直接丢弃(每 30 秒打印 metrics)。 - DB 健康状态机:写入失败立即标记 unhealthy,后台 ping 恢复后重新启用。 - ClickHouse `PrepareBatch` 失败会保留整批重试;逐项 `Append` 失败只保留明确失败项,已成功追加项继续发送。`Send` 返回错误时提交结果可能不明,系统不会自动重发该批,避免静默重复,并通过 `ambiguous_send` 计数暴露可能丢失;进程内 best-effort 队列不承诺分布式 exactly-once。 - 启动会校验现有 `proxy_logs` 的列、引擎、分区键和排序键;不兼容 schema 会保持数据库 unhealthy,不自动重建数据表。 - 所有 batch 路径都会执行清理;`Abort`/关闭失败会与原始错误合并记录,但不能使模糊提交变得可判定。