Files
content-factory/apps/backend/src/application/observability.py
T

240 lines
7.4 KiB
Python

from __future__ import annotations
import re
from datetime import datetime
from typing import Any
from src.domain.contracts import (
AgentJobStatus,
AgentJobSummary,
AgentJobType,
PublishingStatus,
Role,
WorkflowEventSummary,
)
_RETRYABLE_JOB_TYPES = {
AgentJobType.PLAN_GENERATION,
AgentJobType.RESEARCH,
AgentJobType.SECTION_SCAFFOLD,
AgentJobType.SEO_REVIEW,
AgentJobType.LANGUAGE_REVIEW,
AgentJobType.PUBLISH_COMMIT,
AgentJobType.TEST_CODEX,
}
_SENSITIVE_KEY_TOKENS = (
"token",
"password",
"secret",
"authorization",
"cookie",
"api_key",
"apikey",
"session",
)
_SENSITIVE_VALUE_PATTERNS: tuple[tuple[re.Pattern[str], str], ...] = (
(re.compile(r"(?i)(authorization\s*:\s*bearer\s+)[^\s]+"), r"\1[REDACTED]"),
(re.compile(r"(?i)(token\s*[=:]\s*)[^\s,;]+"), r"\1[REDACTED]"),
(re.compile(r"(?i)(password\s*[=:]\s*)[^\s,;]+"), r"\1[REDACTED]"),
(re.compile(r"(?i)(secret\s*[=:]\s*)[^\s,;]+"), r"\1[REDACTED]"),
(re.compile(r"(?i)(api[_-]?key\s*[=:]\s*)[^\s,;]+"), r"\1[REDACTED]"),
(re.compile(r"\b(sk-[a-zA-Z0-9_-]{8,})\b"), "[REDACTED]"),
(re.compile(r"\b(gh[pousr]_[a-zA-Z0-9]{8,})\b"), "[REDACTED]"),
)
def project_job_for_view(
repository: object,
job: AgentJobSummary,
*,
viewer_role: Role,
) -> AgentJobSummary:
retry_eligible, retry_block_reason = evaluate_retry_policy(repository, job)
payload = redact_json_value(job.payload)
error_message = _nullable_redacted(job.error_message)
stdout = redact_text(job.stdout)
stderr = redact_text(job.stderr)
safe_failure_summary = build_safe_failure_summary(job)
if viewer_role != Role.ADMIN:
stdout = ""
stderr = ""
return job.model_copy(
update={
"payload": payload,
"error_message": error_message,
"stdout": stdout,
"stderr": stderr,
"retry_eligible": retry_eligible,
"retry_block_reason": retry_block_reason,
"cancel_eligible": is_cancel_eligible(job),
"safe_failure_summary": safe_failure_summary,
},
deep=True,
)
def evaluate_retry_policy(
repository: object,
job: AgentJobSummary,
) -> tuple[bool, str | None]:
if job.status != AgentJobStatus.FAILED:
return False, "Retry is available only for failed jobs."
if job.job_type not in _RETRYABLE_JOB_TYPES:
return False, f"Retry is not supported for job type {job.job_type.value}."
if job.job_type == AgentJobType.PUBLISH_COMMIT:
if job.article_id is None:
return False, "Publish commit retry requires article context."
existing_publish_commit = repository.publish_commits.latest_for_article_with_statuses(
job.article_id,
statuses={PublishingStatus.PUBLISH_COMMIT_CREATED},
)
if existing_publish_commit is not None:
return (
False,
"Publish commit already exists for this article. Retry is blocked.",
)
return True, None
def is_cancel_eligible(job: AgentJobSummary) -> bool:
return job.status in {AgentJobStatus.QUEUED, AgentJobStatus.RUNNING}
def build_safe_failure_summary(job: AgentJobSummary) -> str | None:
if job.status != AgentJobStatus.FAILED:
return None
category = job.error_category.value if job.error_category is not None else "UNKNOWN"
step = (
_string_value(job.payload.get("last_successful_step"))
or _string_value(job.payload.get("artifact_label"))
or _string_value(job.payload.get("heading"))
)
message = redact_text(job.error_message or "").strip()
parts = [f"{job.job_type.value} failed ({category})."]
if step:
parts.append(f"Last successful step: {step}.")
if message:
parts.append(f"Summary: {message}.")
return " ".join(parts)
def build_observability_timeline(
workflow_events: list[WorkflowEventSummary],
agent_jobs: list[AgentJobSummary],
) -> list[dict[str, Any]]:
timeline: list[dict[str, Any]] = []
for event in workflow_events:
source = "USER" if event.actor_user_id is not None else "SYSTEM"
timeline.append(
{
"id": event.id,
"article_id": event.article_id,
"entry_type": "WORKFLOW_EVENT",
"source": source,
"event_type": event.event_type,
"from_status": event.from_status,
"to_status": event.to_status,
"actor_user_id": event.actor_user_id,
"job_id": None,
"job_type": None,
"job_status": None,
"retry_eligible": False,
"cancel_eligible": False,
"safe_failure_summary": None,
"payload": redact_json_value(event.payload),
"created_at": event.created_at,
}
)
for job in agent_jobs:
created_at = job.finished_at or job.started_at or job.queued_at
timeline.append(
{
"id": job.id,
"article_id": job.article_id,
"entry_type": "AGENT_JOB",
"source": "AGENT",
"event_type": f"{job.job_type.value}_{job.status.value}",
"from_status": None,
"to_status": None,
"actor_user_id": None,
"job_id": job.id,
"job_type": job.job_type,
"job_status": job.status,
"retry_eligible": bool(job.retry_eligible),
"cancel_eligible": bool(job.cancel_eligible),
"safe_failure_summary": job.safe_failure_summary,
"payload": {
"attempt": job.attempt,
"error_category": (
job.error_category.value if job.error_category is not None else None
),
"error_message": job.error_message,
},
"created_at": created_at,
}
)
timeline.sort(
key=lambda item: (
_timeline_dt(item.get("created_at")),
str(item.get("id")),
)
)
return timeline
def redact_json_value(value: Any) -> Any:
if isinstance(value, dict):
redacted: dict[str, Any] = {}
for key, nested in value.items():
if _contains_sensitive_token(key):
redacted[key] = "[REDACTED]"
else:
redacted[key] = redact_json_value(nested)
return redacted
if isinstance(value, list):
return [redact_json_value(item) for item in value]
if isinstance(value, str):
return redact_text(value)
return value
def redact_text(value: str) -> str:
redacted = value
for pattern, replacement in _SENSITIVE_VALUE_PATTERNS:
redacted = pattern.sub(replacement, redacted)
return redacted
def _timeline_dt(value: Any) -> datetime:
if isinstance(value, datetime):
return value
return datetime.min
def _contains_sensitive_token(value: str) -> bool:
normalized = value.strip().lower()
return any(token in normalized for token in _SENSITIVE_KEY_TOKENS)
def _string_value(value: Any) -> str:
if isinstance(value, str):
return value.strip()
return ""
def _nullable_redacted(value: str | None) -> str | None:
if value is None:
return None
redacted = redact_text(value).strip()
if not redacted:
return None
return redacted