Files
video-gen/video-gen-api/app/services/generation/log_service.py
T
2026-07-20 13:48:17 +08:00

248 lines
8.8 KiB
Python

from __future__ import annotations
import hashlib
import json
from typing import Any
from app.enums.generation_task import GenerationMode, GenerationOwnerType
from app.models.base import async_session
from app.models.chat_generation_task import ChatGenerationTask
from app.models.chat_generation_task_event import ChatGenerationTaskEvent
from app.models.chat_provider_call_log import ChatProviderCallLog
from app.models.generation_record import GenerationRecord
from app.services.operation_log_service import build_exception_detail, log_operation_event, sanitize_log_value
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(sanitize_log_value(data))
if text is None:
return None
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()
def _owner_fields(
obj: Any | None,
*,
owner_type: str | None,
owner_id: str | None,
task_id: str | None,
record_id: str | None,
generation_attempt_no: int | None,
generation_mode: str | None,
) -> dict[str, Any] | None:
if isinstance(obj, GenerationRecord) or record_id:
resolved_owner_type = GenerationOwnerType.GENERATION_RECORD.value
resolved_owner_id = record_id or owner_id or getattr(obj, "id", None)
resolved_task_id = None
resolved_record_id = resolved_owner_id
resolved_mode = generation_mode or GenerationMode.GENERATION_RECORD.value
elif isinstance(obj, ChatGenerationTask) or task_id:
resolved_owner_type = GenerationOwnerType.CHAT_GENERATION_TASK.value
resolved_owner_id = task_id or owner_id or getattr(obj, "id", None)
resolved_task_id = resolved_owner_id
resolved_record_id = None
resolved_mode = generation_mode or getattr(obj, "generation_mode", GenerationMode.CHATAPI_ASYNC.value)
else:
resolved_owner_type = owner_type or GenerationOwnerType.CHAT_GENERATION_TASK.value
resolved_owner_id = owner_id
if resolved_owner_type == GenerationOwnerType.GENERATION_RECORD.value:
resolved_task_id = None
resolved_record_id = resolved_owner_id
resolved_mode = generation_mode or GenerationMode.GENERATION_RECORD.value
else:
resolved_task_id = resolved_owner_id
resolved_record_id = None
resolved_mode = generation_mode or GenerationMode.CHATAPI_ASYNC.value
if not resolved_owner_id:
return None
return {
"owner_type": resolved_owner_type,
"owner_id": str(resolved_owner_id),
"task_id": str(resolved_task_id) if resolved_task_id else None,
"generation_record_id": str(resolved_record_id) if resolved_record_id else None,
"generation_attempt_no": int(generation_attempt_no or getattr(obj, "generation_attempt_no", 1) or 1),
"generation_mode": str(resolved_mode or ""),
}
def _fallback_log(event_type: str, fields: dict[str, Any] | None, exc: Exception) -> None:
fields = fields or {}
log_operation_event(
domain="generation_pipeline",
event_type="PIPELINE_DB_LOG_FAILED",
event_status="failed",
task_id=fields.get("owner_id"),
message=f"数据库生成日志写入失败: {event_type}",
detail=build_exception_detail(exc, fields),
error=str(exc),
)
async def log_task_event(
task: Any | None = None,
*,
record: Any | None = None,
owner_type: str | None = None,
owner_id: str | None = None,
task_id: str | None = None,
record_id: str | None = None,
generation_attempt_no: int | None = None,
generation_mode: 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 an owner-scoped event in a separate transaction."""
obj = task or record
fields = _owner_fields(
obj,
owner_type=owner_type,
owner_id=owner_id,
task_id=task_id,
record_id=record_id,
generation_attempt_no=generation_attempt_no,
generation_mode=generation_mode,
)
if not fields:
return
upper_event = str(event_type or "").upper()
failed_event = any(marker in upper_event for marker in ("FAILED", "TIMEOUT", "ERROR"))
log_operation_event(
domain="generation_pipeline",
event_type=event_type,
event_status="failed" if failed_event else "success",
source="pipeline",
user_id=str(getattr(obj, "user_id", "") or "") or None,
project_id=str(getattr(obj, "project_id", "") or "") or None,
task_id=fields["owner_id"],
message=message,
detail={
"owner_type": fields["owner_type"],
"owner_id": fields["owner_id"],
"generation_attempt_no": fields["generation_attempt_no"],
"generation_mode": fields["generation_mode"],
"from_status": from_status,
"to_status": to_status,
"from_stage": from_stage,
"to_stage": to_stage,
"detail_excerpt": _excerpt(detail),
},
error=message if failed_event else None,
)
try:
async with async_session() as db:
db.add(ChatGenerationTaskEvent(
id=generate_id(),
owner_type=fields["owner_type"],
task_id=fields["task_id"],
generation_record_id=fields["generation_record_id"],
generation_attempt_no=fields["generation_attempt_no"],
generation_mode=fields["generation_mode"],
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 as exc:
_fallback_log(event_type, fields, exc)
async def log_provider_call(
task: Any | None = None,
*,
record: Any | None = None,
owner_type: str | None = None,
owner_id: str | None = None,
task_id: str | None = None,
record_id: str | None = None,
generation_attempt_no: int | None = None,
generation_mode: 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 an owner-scoped provider call log in a separate transaction."""
obj = task or record
fields = _owner_fields(
obj,
owner_type=owner_type,
owner_id=owner_id,
task_id=task_id,
record_id=record_id,
generation_attempt_no=generation_attempt_no,
generation_mode=generation_mode,
)
if not fields:
return
try:
async with async_session() as db:
db.add(ChatProviderCallLog(
id=generate_id(),
owner_type=fields["owner_type"],
task_id=fields["task_id"],
generation_record_id=fields["generation_record_id"],
generation_attempt_no=fields["generation_attempt_no"],
generation_mode=fields["generation_mode"],
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 as exc:
_fallback_log(api_type, fields, exc)