Files
video-gen/video-gen-api/app/services/generation/log_service.py
T

302 lines
11 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.generation_record import GenerationRecord
from app.services.operation_log_service import build_exception_detail, log_ai_model_event, 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:
inferred_mode = generation_mode or getattr(obj, "generation_mode", None)
inferred_owner_type = (
GenerationOwnerType.GENERATION_RECORD.value
if inferred_mode == GenerationMode.GENERATION_RECORD.value
else GenerationOwnerType.CHAT_GENERATION_TASK.value
)
resolved_owner_type = owner_type or inferred_owner_type
resolved_owner_id = owner_id or getattr(obj, "id", None)
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,
call_id: str | None = None,
module: str | None = None,
step_code: str | None = None,
) -> str | None:
"""Write provider audit events to the AiModel file log only.
``ChatProviderCallLog`` is intentionally no longer written. The model and
historical table remain registered for backward compatibility, so no schema
migration is required.
"""
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 None
resolved_call_id = call_id or generate_id()
resolved_module = module or fields.get("generation_mode") or "generation_pipeline"
resolved_step = step_code or api_type
common = {
"module": resolved_module,
"step_code": resolved_step,
"call_id": resolved_call_id,
"source": "app.services.generation.log_service",
"task_id": fields.get("owner_id"),
"owner_type": fields.get("owner_type"),
"owner_id": fields.get("owner_id"),
"generation_attempt_no": fields.get("generation_attempt_no"),
"remote_action": api_type,
"remote_request_id": provider_task_id,
"model_config_id": engine_id,
"model_config_name": engine_id,
"model_name": model,
"provider": provider,
"http_status": http_status,
}
detail = {
"generation_mode": fields.get("generation_mode"),
"provider_task_id": provider_task_id,
"error_code": error_code,
}
token_usage = {
"prompt_tokens": int(prompt_tokens or 0),
"completion_tokens": int(completion_tokens or 0),
"total_tokens": int(total_tokens or 0),
}
if request_data is not None or str(status).lower() == "request":
log_ai_model_event(
event_type="REQUEST",
event_phase="REQUEST",
event_status="started",
request=request_data if request_data is not None else {},
detail=detail,
**common,
)
normalized_status = str(status or "").lower()
if response_data is not None or normalized_status in {"success", "completed", "succeeded"}:
log_ai_model_event(
event_type="RESPONSE",
event_phase="RESPONSE",
event_status="success" if normalized_status not in {"failed", "error"} else "failed",
latency_ms=latency_ms,
response=response_data,
token_usage=token_usage,
detail=detail,
error=error_message if normalized_status in {"failed", "error"} else None,
**common,
)
if normalized_status in {"failed", "error"} or error_message:
log_ai_model_event(
event_type="ERROR",
event_phase="ERROR",
event_status="failed",
latency_ms=latency_ms,
token_usage=token_usage,
detail=build_exception_detail(
RuntimeError(error_message or "provider call failed"),
detail,
),
error=error_message or "provider call failed",
**common,
)
return resolved_call_id