feat: 拆分 usage 记录的请求体/响应体为客户端侧与提供商侧

将 request_body/response_body 语义明确为客户端原始请求体和提供商原始响应体,
新增 provider_request_body(格式转换后发给提供商的请求体)和 client_response_body
(格式转换后返回给客户端的响应体),支持跨格式转换场景下分别查看两侧数据。

- 数据库新增 provider_request_body/client_response_body 及对应压缩字段
- 全链路(telemetry/recording/handler/stream_context)传递新字段
- 维护调度器同步支持新字段的压缩与清理
- 前端请求详情抽屉支持请求体/响应体/响应头的客户端/提供商视图切换
This commit is contained in:
fawney19
2026-02-21 01:56:28 +08:00
parent 4c5dac603f
commit 314e4a497d
23 changed files with 449 additions and 106 deletions

View File

@@ -928,9 +928,20 @@ class MaintenanceScheduler:
# 1. 查询需要压缩的记录
# 注意:排除已经是 NULL 或 JSON null 的记录
records = (
batch_db.query(Usage.id, Usage.request_body, Usage.response_body)
batch_db.query(
Usage.id,
Usage.request_body,
Usage.response_body,
Usage.provider_request_body,
Usage.client_response_body,
)
.filter(Usage.created_at < cutoff_time)
.filter((Usage.request_body.isnot(None)) | (Usage.response_body.isnot(None)))
.filter(
(Usage.request_body.isnot(None))
| (Usage.response_body.isnot(None))
| (Usage.provider_request_body.isnot(None))
| (Usage.client_response_body.isnot(None))
)
.limit(batch_size)
.all()
)
@@ -940,9 +951,12 @@ class MaintenanceScheduler:
# 过滤掉实际值为 None 的记录JSON null 被解析为 Python None
valid_records = [
(rid, req, resp)
for rid, req, resp in records
if req is not None or resp is not None
r
for r in records
if r.request_body is not None
or r.response_body is not None
or r.provider_request_body is not None
or r.client_response_body is not None
]
if not valid_records:
@@ -950,17 +964,22 @@ class MaintenanceScheduler:
logger.warning(
f"检测到 {len(records)} 条记录的 body 字段为 JSON null进行清理"
)
for record_id, _, _ in records:
for r in records:
batch_db.execute(
update(Usage)
.where(Usage.id == record_id)
.values(request_body=null(), response_body=null())
.where(Usage.id == r.id)
.values(
request_body=null(),
response_body=null(),
provider_request_body=null(),
client_response_body=null(),
)
)
batch_db.commit()
continue
# 检测是否有重复的 ID说明更新未生效
current_ids = {r[0] for r in valid_records}
current_ids = {r.id for r in valid_records}
repeated_ids = current_ids & processed_ids
if repeated_ids:
logger.error(
@@ -972,28 +991,40 @@ class MaintenanceScheduler:
batch_success = 0
# 2. 逐条更新(确保每条都正确处理)
for record_id, req_body, resp_body in valid_records:
for r in valid_records:
try:
# 使用 null() 确保设置的是 SQL NULL 而不是 JSON null
result = batch_db.execute(
update(Usage)
.where(Usage.id == record_id)
.where(Usage.id == r.id)
.values(
request_body=null(),
response_body=null(),
provider_request_body=null(),
client_response_body=null(),
request_body_compressed=(
compress_json(req_body) if req_body else None
compress_json(r.request_body) if r.request_body else None
),
response_body_compressed=(
compress_json(resp_body) if resp_body else None
compress_json(r.response_body) if r.response_body else None
),
provider_request_body_compressed=(
compress_json(r.provider_request_body)
if r.provider_request_body
else None
),
client_response_body_compressed=(
compress_json(r.client_response_body)
if r.client_response_body
else None
),
)
)
if result.rowcount > 0:
batch_success += 1
processed_ids.add(record_id)
processed_ids.add(r.id)
except Exception as e:
logger.warning(f"压缩记录 {record_id} 失败: {e}")
logger.warning(f"压缩记录 {r.id} 失败: {e}")
continue
batch_db.commit()
@@ -1050,6 +1081,8 @@ class MaintenanceScheduler:
.filter(
(Usage.request_body_compressed.isnot(None))
| (Usage.response_body_compressed.isnot(None))
| (Usage.provider_request_body_compressed.isnot(None))
| (Usage.client_response_body_compressed.isnot(None))
)
.limit(batch_size)
.all()
@@ -1067,6 +1100,8 @@ class MaintenanceScheduler:
.values(
request_body_compressed=null(),
response_body_compressed=null(),
provider_request_body_compressed=null(),
client_response_body_compressed=null(),
)
)

View File

@@ -146,9 +146,11 @@ class UsageBillingIntegrationMixin:
request_headers=params.request_headers,
request_body=params.request_body,
provider_request_headers=params.provider_request_headers,
provider_request_body=params.provider_request_body,
response_headers=params.response_headers,
client_response_headers=params.client_response_headers,
response_body=params.response_body,
client_response_body=params.client_response_body,
request_id=params.request_id,
provider_id=params.provider_id,
provider_endpoint_id=params.provider_endpoint_id,

View File

@@ -56,9 +56,11 @@ def build_usage_params(
request_headers: dict[str, Any] | None,
request_body: Any | None,
provider_request_headers: dict[str, Any] | None,
provider_request_body: Any | None,
response_headers: dict[str, Any] | None,
client_response_headers: dict[str, Any] | None,
response_body: Any | None,
client_response_body: Any | None,
request_id: str,
provider_id: str | None,
provider_endpoint_id: str | None,
@@ -103,16 +105,26 @@ def build_usage_params(
# 处理请求体和响应体(可能需要截断)
processed_request_body = None
processed_provider_request_body = None
processed_response_body = None
processed_client_response_body = None
if should_log_body:
if request_body:
processed_request_body = SystemConfigService.truncate_body(
db, request_body, is_request=True
)
if provider_request_body:
processed_provider_request_body = SystemConfigService.truncate_body(
db, provider_request_body, is_request=True
)
if response_body:
processed_response_body = SystemConfigService.truncate_body(
db, response_body, is_request=False
)
if client_response_body:
processed_client_response_body = SystemConfigService.truncate_body(
db, client_response_body, is_request=False
)
# 处理响应头
processed_response_headers = None
@@ -197,9 +209,11 @@ def build_usage_params(
"request_headers": processed_request_headers,
"request_body": processed_request_body,
"provider_request_headers": processed_provider_request_headers,
"provider_request_body": processed_provider_request_body,
"response_headers": processed_response_headers,
"client_response_headers": processed_client_response_headers,
"response_body": processed_response_body,
"client_response_body": processed_client_response_body,
}
@@ -231,9 +245,12 @@ def update_existing_usage(
existing_usage.request_body = usage_params["request_body"]
if usage_params["provider_request_headers"] is not None:
existing_usage.provider_request_headers = usage_params["provider_request_headers"]
if usage_params["provider_request_body"] is not None:
existing_usage.provider_request_body = usage_params["provider_request_body"]
existing_usage.response_body = usage_params["response_body"]
existing_usage.response_headers = usage_params["response_headers"]
existing_usage.client_response_headers = usage_params["client_response_headers"]
existing_usage.client_response_body = usage_params["client_response_body"]
# 更新 token 和费用信息
existing_usage.input_tokens = usage_params["input_tokens"]

View File

@@ -34,9 +34,11 @@ class UsageRecordParams:
request_headers: dict[str, Any] | None
request_body: Any | None
provider_request_headers: dict[str, Any] | None
provider_request_body: Any | None
response_headers: dict[str, Any] | None
client_response_headers: dict[str, Any] | None
response_body: Any | None
client_response_body: Any | None
request_id: str
provider_id: str | None
provider_endpoint_id: str | None

View File

@@ -81,9 +81,11 @@ def _event_to_record(event: UsageEvent) -> dict[str, Any]:
"request_headers": data.get("request_headers"),
"request_body": _parse_body(data.get("request_body")),
"provider_request_headers": data.get("provider_request_headers"),
"provider_request_body": _parse_body(data.get("provider_request_body")),
"response_headers": data.get("response_headers"),
"client_response_headers": data.get("client_response_headers"),
"response_body": _parse_body(data.get("response_body")),
"client_response_body": _parse_body(data.get("client_response_body")),
"provider_id": data.get("provider_id"),
"provider_endpoint_id": data.get("provider_endpoint_id"),
"provider_api_key_id": data.get("provider_api_key_id"),
@@ -464,9 +466,11 @@ class UsageQueueConsumer:
request_headers=data.get("request_headers"),
request_body=_parse_body(data.get("request_body")),
provider_request_headers=data.get("provider_request_headers"),
provider_request_body=_parse_body(data.get("provider_request_body")),
response_headers=data.get("response_headers"),
client_response_headers=data.get("client_response_headers"),
response_body=_parse_body(data.get("response_body")),
client_response_body=_parse_body(data.get("client_response_body")),
request_id=event.request_id,
provider_id=data.get("provider_id"),
provider_endpoint_id=data.get("provider_endpoint_id"),

View File

@@ -78,9 +78,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers: dict[str, Any] | None = None,
request_body: Any | None = None,
provider_request_headers: dict[str, Any] | None = None,
provider_request_body: Any | None = None,
response_headers: dict[str, Any] | None = None,
client_response_headers: dict[str, Any] | None = None,
response_body: Any | None = None,
client_response_body: Any | None = None,
request_id: str | None = None,
provider_id: str | None = None,
provider_endpoint_id: str | None = None,
@@ -125,9 +127,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers=request_headers,
request_body=request_body,
provider_request_headers=provider_request_headers,
provider_request_body=provider_request_body,
response_headers=response_headers,
client_response_headers=client_response_headers,
response_body=response_body,
client_response_body=client_response_body,
request_id=request_id,
provider_id=provider_id,
provider_endpoint_id=provider_endpoint_id,
@@ -196,9 +200,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers: dict[str, Any] | None = None,
request_body: Any | None = None,
provider_request_headers: dict[str, Any] | None = None,
provider_request_body: Any | None = None,
response_headers: dict[str, Any] | None = None,
client_response_headers: dict[str, Any] | None = None,
response_body: Any | None = None,
client_response_body: Any | None = None,
request_id: str | None = None,
provider_id: str | None = None,
provider_endpoint_id: str | None = None,
@@ -245,9 +251,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers=request_headers,
request_body=request_body,
provider_request_headers=provider_request_headers,
provider_request_body=provider_request_body,
response_headers=response_headers,
client_response_headers=client_response_headers,
response_body=response_body,
client_response_body=client_response_body,
request_id=request_id,
provider_id=provider_id,
provider_endpoint_id=provider_endpoint_id,
@@ -383,9 +391,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers: dict[str, Any] | None = None,
request_body: Any | None = None,
provider_request_headers: dict[str, Any] | None = None,
provider_request_body: Any | None = None,
response_headers: dict[str, Any] | None = None,
client_response_headers: dict[str, Any] | None = None,
response_body: Any | None = None,
client_response_body: Any | None = None,
request_id: str | None = None,
provider_id: str | None = None,
provider_endpoint_id: str | None = None,
@@ -444,9 +454,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers=request_headers,
request_body=request_body,
provider_request_headers=provider_request_headers,
provider_request_body=provider_request_body,
response_headers=response_headers,
client_response_headers=client_response_headers,
response_body=response_body,
client_response_body=client_response_body,
request_id=request_id,
provider_id=provider_id,
provider_endpoint_id=provider_endpoint_id,
@@ -743,9 +755,11 @@ class UsageRecordingMixin(UsageBillingIntegrationMixin):
request_headers=record.get("request_headers"),
request_body=record.get("request_body"),
provider_request_headers=record.get("provider_request_headers"),
provider_request_body=record.get("provider_request_body"),
response_headers=record.get("response_headers"),
client_response_headers=record.get("client_response_headers"),
response_body=record.get("response_body"),
client_response_body=record.get("client_response_body"),
request_id=request_id,
provider_id=record.get("provider_id"),
provider_endpoint_id=record.get("provider_endpoint_id"),

View File

@@ -72,6 +72,8 @@ class MessageTelemetry:
cache_read_tokens: int = 0,
is_stream: bool = False,
provider_request_headers: dict[str, Any] | None = None,
provider_request_body: Any | None = None,
client_response_body: Any | None = None,
# 时间指标
first_byte_time_ms: int | None = None, # 首字时间/TTFB
# Provider 侧追踪信息(用于记录真实成本)
@@ -117,9 +119,11 @@ class MessageTelemetry:
request_headers=request_headers,
request_body=request_body,
provider_request_headers=provider_request_headers or {},
provider_request_body=provider_request_body,
response_headers=response_headers,
client_response_headers=client_response_headers,
response_body=response_body,
client_response_body=client_response_body,
request_id=self.request_id,
# Provider 侧追踪信息(用于记录真实成本)
provider_id=provider_id,
@@ -164,6 +168,7 @@ class MessageTelemetry:
is_stream: bool,
api_format: str | None = None,
provider_request_headers: dict[str, Any] | None = None,
provider_request_body: Any | None = None,
# 预估 token 信息(来自 message_start 事件,用于中断请求的成本估算)
input_tokens: int = 0,
output_tokens: int = 0,
@@ -172,6 +177,7 @@ class MessageTelemetry:
response_body: dict[str, Any] | None = None,
response_headers: dict[str, Any] | None = None,
client_response_headers: dict[str, Any] | None = None,
client_response_body: Any | None = None,
# Provider 侧追踪信息(用于 curl 复现等场景)
provider_id: str | None = None,
provider_endpoint_id: str | None = None,
@@ -225,9 +231,11 @@ class MessageTelemetry:
request_headers=request_headers,
request_body=request_body,
provider_request_headers=provider_request_headers or {},
provider_request_body=provider_request_body,
response_headers=response_headers or {},
client_response_headers=client_response_headers,
response_body=response_body or {"error": error_message},
client_response_body=client_response_body,
request_id=self.request_id,
# Provider 侧追踪信息
provider_id=provider_id,
@@ -252,6 +260,7 @@ class MessageTelemetry:
is_stream: bool,
api_format: str | None = None,
provider_request_headers: dict[str, Any] | None = None,
provider_request_body: Any | None = None,
input_tokens: int = 0,
output_tokens: int = 0,
cache_creation_tokens: int = 0,
@@ -259,6 +268,7 @@ class MessageTelemetry:
response_body: dict[str, Any] | None = None,
response_headers: dict[str, Any] | None = None,
client_response_headers: dict[str, Any] | None = None,
client_response_body: Any | None = None,
# Provider 侧追踪信息
provider_id: str | None = None,
provider_endpoint_id: str | None = None,
@@ -299,9 +309,11 @@ class MessageTelemetry:
request_headers=request_headers,
request_body=request_body,
provider_request_headers=provider_request_headers or {},
provider_request_body=provider_request_body,
response_headers=response_headers or {},
client_response_headers=client_response_headers,
response_body=response_body or {},
client_response_body=client_response_body,
request_id=self.request_id,
# Provider 侧追踪信息
provider_id=provider_id,

View File

@@ -268,14 +268,28 @@ class QueueTelemetryWriter(TelemetryWriter):
max_size=self._max_request_body_size,
is_request=True,
)
provider_request_body = self._truncate_body(
kwargs.get("provider_request_body"),
max_size=self._max_request_body_size,
is_request=True,
)
response_body = self._truncate_body(
kwargs.get("response_body"),
max_size=self._max_response_body_size,
is_request=False,
)
client_response_body = self._truncate_body(
kwargs.get("client_response_body"),
max_size=self._max_response_body_size,
is_request=False,
)
if request_body is not None:
data["request_body"] = request_body
if provider_request_body is not None:
data["provider_request_body"] = provider_request_body
if response_body is not None:
data["response_body"] = response_body
if client_response_body is not None:
data["client_response_body"] = client_response_body
return data