135 lines
4.1 KiB
Python
135 lines
4.1 KiB
Python
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
from typing import Any
|
|
|
|
from app.models.base import async_session
|
|
from app.models.chat_generation_task_event import ChatGenerationTaskEvent
|
|
from app.models.chat_provider_call_log import ChatProviderCallLog
|
|
from app.utils.id_gen import generate_id
|
|
|
|
MAX_EXCERPT_CHARS = 2000
|
|
|
|
|
|
def _safe_json(data: Any) -> str | None:
|
|
if data is None:
|
|
return None
|
|
try:
|
|
return json.dumps(data, ensure_ascii=False, default=str)
|
|
except Exception:
|
|
return str(data)
|
|
|
|
|
|
def _excerpt(data: Any, limit: int = MAX_EXCERPT_CHARS) -> str | None:
|
|
text = _safe_json(data)
|
|
if text is None:
|
|
return None
|
|
# Avoid storing secrets in logs.
|
|
text = text.replace("Authorization", "Authorization-REDACTED")
|
|
text = text.replace("api_key", "api_key_REDACTED")
|
|
if len(text) > limit:
|
|
return text[:limit] + "...[truncated]"
|
|
return text
|
|
|
|
|
|
def _hash(data: Any) -> str | None:
|
|
text = _safe_json(data)
|
|
if text is None:
|
|
return None
|
|
return hashlib.sha256(text.encode("utf-8")).hexdigest()
|
|
|
|
|
|
async def log_task_event(
|
|
task: Any | None = None,
|
|
*,
|
|
record: Any | None = None,
|
|
task_id: str | None = None,
|
|
record_id: str | None = None,
|
|
event_type: str,
|
|
from_status: str | None = None,
|
|
to_status: str | None = None,
|
|
from_stage: str | None = None,
|
|
to_stage: str | None = None,
|
|
message: str | None = None,
|
|
detail: Any = None,
|
|
) -> None:
|
|
"""Write task event in a separate transaction; failure must not affect main flow."""
|
|
try:
|
|
obj = task or record
|
|
tid = task_id or record_id or (obj.id if obj else None)
|
|
if not tid:
|
|
return
|
|
async with async_session() as db:
|
|
db.add(ChatGenerationTaskEvent(
|
|
id=generate_id(),
|
|
task_id=tid,
|
|
generation_mode=getattr(obj, "generation_mode", "chatapi_async"),
|
|
event_type=event_type,
|
|
from_status=from_status,
|
|
to_status=to_status,
|
|
from_stage=from_stage,
|
|
to_stage=to_stage,
|
|
message=message,
|
|
detail_json=_excerpt(detail),
|
|
))
|
|
await db.commit()
|
|
except Exception:
|
|
return
|
|
|
|
|
|
async def log_provider_call(
|
|
task: Any | None = None,
|
|
*,
|
|
record: Any | None = None,
|
|
task_id: str | None = None,
|
|
record_id: str | None = None,
|
|
provider: str | None,
|
|
api_type: str,
|
|
model: str | None = None,
|
|
engine_id: str | None = None,
|
|
status: str,
|
|
latency_ms: int | None = None,
|
|
http_status: int | None = None,
|
|
provider_task_id: str | None = None,
|
|
request_data: Any = None,
|
|
response_data: Any = None,
|
|
prompt_tokens: int = 0,
|
|
completion_tokens: int = 0,
|
|
total_tokens: int = 0,
|
|
error_code: str | None = None,
|
|
error_message: str | None = None,
|
|
) -> None:
|
|
"""Write provider call log in a separate transaction; failure must not affect main flow."""
|
|
try:
|
|
obj = task or record
|
|
tid = task_id or record_id or (obj.id if obj else None)
|
|
if not tid:
|
|
return
|
|
async with async_session() as db:
|
|
db.add(ChatProviderCallLog(
|
|
id=generate_id(),
|
|
task_id=tid,
|
|
generation_mode=getattr(obj, "generation_mode", "chatapi_async"),
|
|
provider=provider,
|
|
api_type=api_type,
|
|
model=model,
|
|
engine_id=engine_id,
|
|
status=status,
|
|
latency_ms=latency_ms,
|
|
http_status=http_status,
|
|
provider_task_id=provider_task_id,
|
|
request_hash=_hash(request_data),
|
|
response_hash=_hash(response_data),
|
|
request_excerpt=_excerpt(request_data),
|
|
response_excerpt=_excerpt(response_data),
|
|
prompt_tokens=prompt_tokens or 0,
|
|
completion_tokens=completion_tokens or 0,
|
|
total_tokens=total_tokens or 0,
|
|
error_code=error_code,
|
|
error_message=error_message,
|
|
))
|
|
await db.commit()
|
|
except Exception:
|
|
return
|