decode_runner.py 2.0 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768
  1. """Pipeline adapter for decoding creation candidate items."""
  2. from __future__ import annotations
  3. from dataclasses import dataclass
  4. from typing import Protocol
  5. from uuid import UUID
  6. from acquisition.domain import CandidateItem, MediaAsset
  7. from decode_content.readers.service import post_from_candidate_item
  8. from decode_content.service import DecodeService, DecodeWorkflowOutput
  9. from pipeline.dedupe import dedupe_candidate_items, should_decode_item
  10. class DecodeCandidateRepository(Protocol):
  11. def list_creation_candidate_items(
  12. self,
  13. *,
  14. run_id: UUID | None = None,
  15. limit: int = 100,
  16. ) -> list[CandidateItem]:
  17. ...
  18. def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
  19. ...
  20. @dataclass(frozen=True)
  21. class DecodeBatchResult:
  22. total: int
  23. decoded: int
  24. skipped: int
  25. failed: int
  26. outputs: list[DecodeWorkflowOutput]
  27. def run_decode_stage(
  28. *,
  29. candidate_repo: DecodeCandidateRepository,
  30. decode_service: DecodeService,
  31. run_id: UUID | None = None,
  32. limit: int = 100,
  33. decoded_item_ids: set[str] | None = None,
  34. ) -> DecodeBatchResult:
  35. items = dedupe_candidate_items(candidate_repo.list_creation_candidate_items(run_id=run_id, limit=limit))
  36. outputs: list[DecodeWorkflowOutput] = []
  37. decoded = skipped = failed = 0
  38. for item in items:
  39. if item.id is None:
  40. skipped += 1
  41. continue
  42. decision = should_decode_item(item, decoded_item_ids=decoded_item_ids)
  43. if not decision.keep:
  44. skipped += 1
  45. continue
  46. try:
  47. media = candidate_repo.list_media_assets_for_item(item.id)
  48. post = post_from_candidate_item(item, media)
  49. outputs.append(decode_service.decode_post(item_id=item.id, post=post))
  50. decoded += 1
  51. except Exception:
  52. failed += 1
  53. return DecodeBatchResult(
  54. total=len(items),
  55. decoded=decoded,
  56. skipped=skipped,
  57. failed=failed,
  58. outputs=outputs,
  59. )