| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165 |
- """Pipeline adapter for decoding creation candidate items."""
- from __future__ import annotations
- from dataclasses import dataclass, field
- from typing import Protocol
- from uuid import UUID
- from acquisition.domain import CandidateItem, MediaAsset
- from core.text_limits import ERROR_MESSAGE_MAX_CHARS, clip_text
- from decode_content.readers.service import post_from_candidate_item
- from decode_content.service import DecodeService, DecodeWorkflowOutput
- from pipeline.dedupe import dedupe_candidate_items, should_decode_item
- from pipeline.tracing import NoopTraceWriter, TraceContext, TraceWriter
- class DecodeCandidateRepository(Protocol):
- def list_creation_candidate_items(
- self,
- *,
- run_id: UUID | None = None,
- limit: int = 100,
- ) -> list[CandidateItem]:
- ...
- def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
- ...
- @dataclass(frozen=True)
- class DecodeBatchResult:
- total: int
- decoded: int
- skipped: int
- failed: int
- outputs: list[DecodeWorkflowOutput]
- failures: list[dict[str, str]] = field(default_factory=list)
- def _record_failure(decode_service: DecodeService, item: CandidateItem, exc: Exception) -> dict[str, str]:
- message = clip_text(str(exc) or exc.__class__.__name__, ERROR_MESSAGE_MAX_CHARS)
- failure = {
- "item_id": str(item.id),
- "platform": item.platform,
- "title": item.title or "",
- "error": message,
- }
- repo = getattr(decode_service, "repository", None)
- if repo is not None and item.id is not None:
- mark_jobs = getattr(repo, "mark_running_decode_jobs_failed", None)
- if mark_jobs is not None:
- mark_jobs(item.id, message)
- save_result = getattr(repo, "save_decode_result", None)
- if save_result is not None:
- save_result(
- item_id=item.id,
- read_result={"is_empty": True, "text": "", "metadata": {"decode_error": message}},
- gate_result={"passed": False, "reason": "decode_failed", "details": {"error": message}},
- framing_result={"error": message},
- status="failed",
- )
- return failure
- def run_decode_stage(
- *,
- candidate_repo: DecodeCandidateRepository,
- decode_service: DecodeService,
- run_id: UUID | None = None,
- limit: int = 100,
- decoded_item_ids: set[str] | None = None,
- trace_writer: TraceWriter | None = None,
- trace_context: TraceContext | None = None,
- ) -> DecodeBatchResult:
- items = dedupe_candidate_items(candidate_repo.list_creation_candidate_items(run_id=run_id, limit=limit))
- trace_writer = trace_writer or NoopTraceWriter()
- base_context = trace_context or TraceContext(acquisition_run_id=run_id, stage="decode")
- trace_writer.event(
- context=base_context,
- stage="decode",
- event_type="decode_stage_started",
- status="running",
- payload={"run_id": str(run_id) if run_id else None, "candidate_count": len(items), "limit": limit},
- )
- outputs: list[DecodeWorkflowOutput] = []
- failures: list[dict[str, str]] = []
- decoded = skipped = failed = 0
- for item in items:
- if item.id is None:
- skipped += 1
- trace_writer.event(
- context=base_context.child(platform=item.platform),
- stage="decode",
- event_type="decode_item_skipped",
- status="skipped",
- payload={"reason": "missing_item_id", "title": item.title or ""},
- )
- continue
- item_context = base_context.child(
- item_id=item.id,
- acquisition_job_id=item.job_id,
- query_id=item.query_id,
- platform=item.platform,
- )
- decision = should_decode_item(item, decoded_item_ids=decoded_item_ids)
- if not decision.keep:
- skipped += 1
- trace_writer.event(
- context=item_context,
- stage="decode",
- event_type="decode_item_skipped",
- status="skipped",
- target_table="candidate_items",
- target_id=item.id,
- payload={"reason": decision.reason},
- )
- continue
- try:
- trace_writer.event(
- context=item_context,
- stage="decode",
- event_type="decode_item_started",
- status="running",
- target_table="candidate_items",
- target_id=item.id,
- )
- media = candidate_repo.list_media_assets_for_item(item.id)
- post = post_from_candidate_item(item, media)
- outputs.append(decode_service.decode_post(item_id=item.id, post=post))
- decoded += 1
- trace_writer.event(
- context=item_context,
- stage="decode",
- event_type="decode_item_finished",
- status="done",
- target_table="candidate_items",
- target_id=item.id,
- )
- except Exception as exc:
- failed += 1
- failures.append(_record_failure(decode_service, item, exc))
- trace_writer.event(
- context=item_context,
- stage="decode",
- event_type="decode_item_failed",
- status="failed",
- severity="error",
- target_table="candidate_items",
- target_id=item.id,
- error_message=str(exc),
- )
- trace_writer.event(
- context=base_context,
- stage="decode",
- event_type="decode_stage_finished",
- status="done" if failed == 0 else "partial",
- payload={"total": len(items), "decoded": decoded, "skipped": skipped, "failed": failed},
- )
- return DecodeBatchResult(
- total=len(items),
- decoded=decoded,
- skipped=skipped,
- failed=failed,
- outputs=outputs,
- failures=failures,
- )
|