mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
refactor: 前端全面替换 any 为 unknown 并统一错误处理,后端用量记录补写请求头/体
- 前端 API 层、stores、conversation 解析器、组件全面替换 any 为 unknown/具体类型 - 错误处理统一使用 parseApiError/getErrorStatus 替代 err.response?.data?.detail 模式 - 后端 handler/TaskService/UsageLifecycle/StreamTracker 链路传递 request_headers/request_body - streaming/pending 状态更新时可补写客户端和提供商的请求头及请求体 - 新增 TaskService 和 UsageService 相关测试
This commit is contained in:
@@ -325,6 +325,8 @@ class BaseMessageHandler:
|
||||
api_format=api_format,
|
||||
endpoint_api_format=endpoint_api_format,
|
||||
has_format_conversion=has_format_conversion,
|
||||
provider_request_headers=ctx.provider_request_headers or None,
|
||||
provider_request_body=ctx.provider_request_body,
|
||||
)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
@@ -563,6 +563,8 @@ class ChatHandlerBase(BaseMessageHandler, ABC):
|
||||
capability_requirements=capability_requirements or None,
|
||||
preferred_key_ids=preferred_key_ids or None,
|
||||
request_body_ref=request_body_ref,
|
||||
request_headers=original_headers,
|
||||
request_body=original_request_body,
|
||||
)
|
||||
stream_generator = exec_result.response
|
||||
provider_name = exec_result.provider_name or "unknown"
|
||||
|
||||
@@ -160,6 +160,8 @@ class ChatSyncExecutor:
|
||||
capability_requirements=capability_requirements or None,
|
||||
preferred_key_ids=preferred_key_ids or None,
|
||||
request_body_ref=request_body_ref,
|
||||
request_headers=original_headers,
|
||||
request_body=original_request_body,
|
||||
)
|
||||
actual_provider_name = exec_result.provider_name or "unknown"
|
||||
ctx.provider_id = exec_result.provider_id
|
||||
|
||||
@@ -445,6 +445,8 @@ class CliEventMixin:
|
||||
api_format=ctx.api_format,
|
||||
endpoint_api_format=ctx.provider_api_format or None,
|
||||
has_format_conversion=ctx.has_format_conversion,
|
||||
provider_request_headers=ctx.provider_request_headers or None,
|
||||
provider_request_body=ctx.provider_request_body,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning("[{}] 同步更新 streaming 状态失败: {}", self.request_id, e)
|
||||
|
||||
@@ -168,6 +168,8 @@ class CliStreamMixin:
|
||||
capability_requirements=capability_requirements or None,
|
||||
preferred_key_ids=preferred_key_ids or None,
|
||||
request_body_ref=request_body_ref,
|
||||
request_headers=original_headers,
|
||||
request_body=original_request_body,
|
||||
)
|
||||
stream_generator = exec_result.response
|
||||
provider_name = exec_result.provider_name or "unknown"
|
||||
|
||||
@@ -464,6 +464,8 @@ class CliSyncMixin:
|
||||
capability_requirements=capability_requirements or None,
|
||||
preferred_key_ids=preferred_key_ids or None,
|
||||
request_body_ref=request_body_ref,
|
||||
request_headers=original_headers,
|
||||
request_body=original_request_body,
|
||||
)
|
||||
result = exec_result.response
|
||||
actual_provider_name = exec_result.provider_name or "unknown"
|
||||
|
||||
@@ -85,6 +85,8 @@ class TaskService:
|
||||
capability_requirements: dict[str, bool] | None = None,
|
||||
preferred_key_ids: list[str] | None = None,
|
||||
request_body_ref: dict[str, Any] | None = None,
|
||||
request_headers: dict[str, Any] | None = None,
|
||||
request_body: dict[str, Any] | None = None,
|
||||
# ASYNC-only (video submit)
|
||||
extract_external_task_id: Any | None = None,
|
||||
supported_auth_types: set[str] | None = None,
|
||||
@@ -180,6 +182,8 @@ class TaskService:
|
||||
capability_requirements=capability_requirements,
|
||||
preferred_key_ids=preferred_key_ids,
|
||||
request_body_ref=request_body_ref,
|
||||
request_headers=request_headers,
|
||||
request_body=request_body,
|
||||
)
|
||||
|
||||
async def _execute_sync_unified(
|
||||
@@ -194,6 +198,8 @@ class TaskService:
|
||||
capability_requirements: dict[str, bool] | None,
|
||||
preferred_key_ids: list[str] | None,
|
||||
request_body_ref: dict[str, Any] | None,
|
||||
request_headers: dict[str, Any] | None,
|
||||
request_body: dict[str, Any] | None,
|
||||
) -> ExecutionResult:
|
||||
"""
|
||||
Unified candidate traversal loop for SYNC.
|
||||
@@ -266,6 +272,8 @@ class TaskService:
|
||||
model=model_name,
|
||||
is_stream=is_stream,
|
||||
api_format=api_format_norm,
|
||||
request_headers=request_headers,
|
||||
request_body=request_body,
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.warning("创建 pending 使用记录失败: {}", str(exc))
|
||||
|
||||
@@ -414,6 +414,10 @@ class UsageQueueConsumer:
|
||||
api_format=data.get("api_format"),
|
||||
endpoint_api_format=data.get("endpoint_api_format"),
|
||||
has_format_conversion=data.get("has_format_conversion"),
|
||||
request_headers=data.get("request_headers"),
|
||||
request_body=data.get("request_body"),
|
||||
provider_request_headers=data.get("provider_request_headers"),
|
||||
provider_request_body=data.get("provider_request_body"),
|
||||
)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
@@ -437,6 +437,10 @@ class UsageLifecycleMixin:
|
||||
endpoint_api_format: str | None = None,
|
||||
has_format_conversion: bool | None = None,
|
||||
status_code: int | None = None,
|
||||
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,
|
||||
) -> Usage | None:
|
||||
"""
|
||||
快速更新使用记录状态
|
||||
@@ -456,6 +460,10 @@ class UsageLifecycleMixin:
|
||||
endpoint_api_format: 端点原生 API 格式(可选)
|
||||
has_format_conversion: 是否发生了格式转换(可选)
|
||||
status_code: HTTP 状态码(可选)
|
||||
request_headers: 客户端请求头(可选,用于补写 pending/streaming 记录)
|
||||
request_body: 客户端请求体(可选,用于补写 pending/streaming 记录)
|
||||
provider_request_headers: 提供商请求头(可选,streaming 时可写入)
|
||||
provider_request_body: 提供商请求体(可选,streaming 时可写入)
|
||||
|
||||
Returns:
|
||||
更新后的 Usage 记录,如果未找到则返回 None
|
||||
@@ -509,6 +517,28 @@ class UsageLifecycleMixin:
|
||||
if status_code is not None:
|
||||
usage.status_code = status_code
|
||||
|
||||
should_log_headers = SystemConfigService.should_log_headers(db)
|
||||
should_log_body = SystemConfigService.should_log_body(db)
|
||||
|
||||
if should_log_headers:
|
||||
if isinstance(request_headers, dict):
|
||||
usage.request_headers = SystemConfigService.mask_sensitive_headers(
|
||||
db, request_headers
|
||||
)
|
||||
if isinstance(provider_request_headers, dict):
|
||||
usage.provider_request_headers = SystemConfigService.mask_sensitive_headers(
|
||||
db, provider_request_headers
|
||||
)
|
||||
if should_log_body:
|
||||
if request_body is not None:
|
||||
usage.request_body = SystemConfigService.truncate_body(
|
||||
db, request_body, is_request=True
|
||||
)
|
||||
if provider_request_body is not None:
|
||||
usage.provider_request_body = SystemConfigService.truncate_body(
|
||||
db, provider_request_body, is_request=True
|
||||
)
|
||||
|
||||
# 结算状态:当请求进入终态时,将 billing_status 标记为 settled
|
||||
# 注意:取消是否应 VOID/部分结算由更高层策略决定;这里默认终态均视为已结算。
|
||||
if status in ("completed", "failed", "cancelled"):
|
||||
|
||||
@@ -121,6 +121,9 @@ class StreamUsageTracker:
|
||||
maxlen=50
|
||||
) # 仅保留最后50个原始chunk(用于错误诊断)
|
||||
|
||||
# 请求体(由 track_stream 设置,初始化为 None 以消除 hasattr 检查)
|
||||
self.request_data: dict[str, Any] | None = None
|
||||
|
||||
# 时间跟踪
|
||||
self.start_time = None
|
||||
self.end_time = None
|
||||
@@ -523,9 +526,12 @@ class StreamUsageTracker:
|
||||
api_format=self.api_format,
|
||||
endpoint_api_format=self.endpoint_api_format,
|
||||
has_format_conversion=self.has_format_conversion,
|
||||
request_headers=self.request_headers,
|
||||
request_body=self.request_data,
|
||||
provider_request_headers=self.provider_request_headers,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"更新使用记录状态为 streaming 失败: {e}")
|
||||
logger.warning("更新使用记录状态为 streaming 失败: {}", e)
|
||||
|
||||
# 解析块以提取内容和使用信息(chunk是原始字节)
|
||||
content, usage = self.parse_stream_chunk(chunk)
|
||||
@@ -773,7 +779,7 @@ class StreamUsageTracker:
|
||||
status_code=self.status_code, # 使用实际的状态码
|
||||
error_message=self.error_message, # 使用实际的错误消息
|
||||
metadata={"stream": True, "content_length": len(self.accumulated_content)},
|
||||
request_body=self.request_data if hasattr(self, "request_data") else None,
|
||||
request_body=self.request_data,
|
||||
request_headers=self.request_headers,
|
||||
provider_request_headers=self.provider_request_headers,
|
||||
response_headers=self.response_headers,
|
||||
@@ -1026,9 +1032,12 @@ class EnhancedStreamUsageTracker(StreamUsageTracker):
|
||||
api_format=self.api_format,
|
||||
endpoint_api_format=self.endpoint_api_format,
|
||||
has_format_conversion=self.has_format_conversion,
|
||||
request_headers=self.request_headers,
|
||||
request_body=self.request_data,
|
||||
provider_request_headers=self.provider_request_headers,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"更新使用记录状态为 streaming 失败: {e}")
|
||||
logger.warning("更新使用记录状态为 streaming 失败: {}", e)
|
||||
|
||||
# 解析块以提取内容和使用信息(chunk是原始字节)
|
||||
content, usage = self.parse_stream_chunk(chunk)
|
||||
|
||||
Reference in New Issue
Block a user