Files
planet/backend/app/services/persistent_logs.py
linkong acbbfdf9e2
Some checks failed
ci / backend (push) Has been cancelled
ci / frontend (push) Has been cancelled
ci / delivery (push) Has been cancelled
release / images (push) Has been cancelled
release: bump version to 0.69.0
2026-06-03 17:27:00 +08:00

275 lines
9.2 KiB
Python

from __future__ import annotations
import hashlib
import re
from datetime import UTC, datetime
from typing import Any
from app.core.logging import get_logger, sanitize_log_value
from app.core.request_context import get_request_id
from app.db.session import async_session_factory
from app.models.system_log import AuditLog, ObservabilityEvent, ObservabilityEventGroup, SystemLog
logger = get_logger(__name__)
HLS_TRANSIENT_RE = re.compile(r"(index|chunk|segment)[_-]?\d+(?:_\d+)?\.(?:ts|m4s|vtt)", re.IGNORECASE)
QUERY_RE = re.compile(r"([?&](?:m|t|token|expires|signature|X-Amz-[^=]+)=[^&\\s]+)", re.IGNORECASE)
UUID_RE = re.compile(r"\b[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\b", re.IGNORECASE)
CONNECTION_RE = re.compile(r"\bconn_[A-Za-z0-9:._-]+\b")
NUMBER_RE = re.compile(r"\b\d{5,}\b")
def normalize_observability_text(value: Any) -> str:
text = str(sanitize_log_value(value or "")).strip()
text = QUERY_RE.sub("", text)
text = HLS_TRANSIENT_RE.sub("<hls-fragment>", text)
text = UUID_RE.sub("<uuid>", text)
text = CONNECTION_RE.sub("<connection>", text)
text = NUMBER_RE.sub("<number>", text)
return re.sub(r"\s+", " ", text).strip()
def build_observability_fingerprint(
*,
source: str,
service: str | None = None,
module: str | None = None,
category: str | None = None,
event: str | None = None,
message: str,
context: dict[str, Any] | None = None,
) -> str:
context = context or {}
stable_context = {
key: context.get(key)
for key in (
"task_type",
"source_id",
"source",
"provider",
"status_code",
"error_type",
"details",
)
if context.get(key) not in (None, "")
}
raw = "|".join(
[
normalize_observability_text(source),
normalize_observability_text(service),
normalize_observability_text(module),
normalize_observability_text(category),
normalize_observability_text(event),
normalize_observability_text(message),
normalize_observability_text(stable_context),
]
)
return hashlib.sha1(raw.encode("utf-8", errors="replace")).hexdigest()
def _context_text(context: dict[str, Any] | None, key: str) -> str | None:
value = (context or {}).get(key)
if value in (None, ""):
return None
return str(value)
async def record_observability_event(
*,
source: str,
level: str,
message: str,
service: str | None = None,
module: str | None = None,
event: str | None = None,
request_id: str | None = None,
trace_id: str | None = None,
user_id: int | None = None,
category: str | None = None,
context: dict[str, Any] | None = None,
fingerprint: str | None = None,
occurred_at: datetime | None = None,
occurrence_count: int = 1,
) -> None:
normalized_context = sanitize_log_value(context or {})
if not isinstance(normalized_context, dict):
normalized_context = {"value": normalized_context}
safe_message = str(sanitize_log_value(message))
normalized_level = str(level or "info").lower()
count = max(1, int(occurrence_count or 1))
event_time = occurred_at or datetime.now(UTC)
event_fingerprint = fingerprint or build_observability_fingerprint(
source=source,
service=service,
module=module,
category=category,
event=event,
message=safe_message,
context=normalized_context,
)
detail = _context_text(normalized_context, "detail") or _context_text(normalized_context, "error")
affected_sources = sorted(
{
item
for item in (
source,
service,
module,
_context_text(normalized_context, "source_id"),
_context_text(normalized_context, "source"),
)
if item
}
)
try:
async with async_session_factory() as session:
session.add(
ObservabilityEvent(
source=source,
service=service,
module=module,
category=category,
event=event,
level=normalized_level,
message=safe_message,
fingerprint=event_fingerprint,
occurred_at=event_time,
request_id=request_id or get_request_id(),
trace_id=trace_id,
user_id=user_id,
task_id=_context_text(normalized_context, "task_id"),
source_ref_id=_context_text(normalized_context, "source_id") or _context_text(normalized_context, "source"),
provider=_context_text(normalized_context, "provider"),
context=normalized_context,
occurrence_count=count,
)
)
group = await session.get(ObservabilityEventGroup, event_fingerprint)
if group is None:
session.add(
ObservabilityEventGroup(
fingerprint=event_fingerprint,
source=source,
service=service,
module=module,
category=category,
event=event,
last_level=normalized_level,
sample_message=safe_message,
sample_detail=detail,
affected_sources=affected_sources,
count=count,
first_seen_at=event_time,
last_seen_at=event_time,
)
)
else:
group.count = int(group.count or 0) + count
group.last_seen_at = event_time
group.last_level = normalized_level
group.sample_message = safe_message
group.sample_detail = detail
merged_sources = sorted(set(group.affected_sources or []) | set(affected_sources))
group.affected_sources = merged_sources
await session.commit()
except Exception:
logger.exception_event(
"Failed to persist observability event",
event="observability_event.persist.failed",
context={"event_name": event, "source": source},
)
async def record_system_log(
*,
source: str,
level: str,
message: str,
service: str | None = None,
module: str | None = None,
event: str | None = None,
request_id: str | None = None,
trace_id: str | None = None,
user_id: int | None = None,
category: str | None = None,
context: dict[str, Any] | None = None,
fingerprint: str | None = None,
occurrence_count: int = 1,
) -> None:
try:
async with async_session_factory() as session:
session.add(
SystemLog(
source=source,
service=service,
module=module,
event=event,
level=level.lower(),
message=str(sanitize_log_value(message)),
request_id=request_id or get_request_id(),
trace_id=trace_id,
user_id=user_id,
category=category,
context=sanitize_log_value(context or {}),
)
)
await session.commit()
except Exception:
logger.exception_event(
"Failed to persist system log",
event="system_log.persist.failed",
context={"event_name": event, "source": source},
)
await record_observability_event(
source=source,
service=service,
module=module,
event=event,
level=level,
message=message,
request_id=request_id,
trace_id=trace_id,
user_id=user_id,
category=category,
context=context,
fingerprint=fingerprint,
occurrence_count=occurrence_count,
)
async def record_audit_log(
*,
action: str,
actor_id: int | None = None,
actor_name: str | None = None,
target_type: str | None = None,
target_id: str | None = None,
result: str | None = None,
request_id: str | None = None,
ip: str | None = None,
details: dict[str, Any] | None = None,
) -> None:
try:
async with async_session_factory() as session:
session.add(
AuditLog(
actor_id=actor_id,
actor_name=actor_name,
action=action,
target_type=target_type,
target_id=target_id,
result=result,
request_id=request_id or get_request_id(),
ip=ip,
details=sanitize_log_value(details or {}),
)
)
await session.commit()
except Exception:
logger.exception_event(
"Failed to persist audit log",
event="audit_log.persist.failed",
context={"action": action},
)