| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768 |
- """Pipeline adapter for decoding creation candidate items."""
- from __future__ import annotations
- from dataclasses import dataclass
- from typing import Protocol
- from uuid import UUID
- from acquisition.domain import CandidateItem, MediaAsset
- 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
- 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]
- 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,
- ) -> DecodeBatchResult:
- items = dedupe_candidate_items(candidate_repo.list_creation_candidate_items(run_id=run_id, limit=limit))
- outputs: list[DecodeWorkflowOutput] = []
- decoded = skipped = failed = 0
- for item in items:
- if item.id is None:
- skipped += 1
- continue
- decision = should_decode_item(item, decoded_item_ids=decoded_item_ids)
- if not decision.keep:
- skipped += 1
- continue
- try:
- 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
- except Exception:
- failed += 1
- return DecodeBatchResult(
- total=len(items),
- decoded=decoded,
- skipped=skipped,
- failed=failed,
- outputs=outputs,
- )
|