creation_pipeline.py 2.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566
  1. """Formal query -> acquisition -> decode orchestration."""
  2. from __future__ import annotations
  3. from dataclasses import dataclass
  4. from uuid import UUID
  5. from acquisition.repositories.base import AcquisitionRepository
  6. from core.config import Settings
  7. from decode_content.service import DecodeService
  8. from pipeline.acquisition_runner import AcquisitionStageResult, run_acquisition_stage
  9. from pipeline.decode_runner import DecodeBatchResult, run_decode_stage
  10. from pipeline.models import PipelineRun, ResumePolicy
  11. from pipeline.repository import PipelineRepository
  12. @dataclass(frozen=True)
  13. class CreationPipelineResult:
  14. pipeline_run: PipelineRun | None
  15. acquisition: AcquisitionStageResult
  16. decode: DecodeBatchResult | None = None
  17. def run_creation_pipeline(
  18. *,
  19. acquisition_repo: AcquisitionRepository,
  20. batch_id: UUID,
  21. settings: Settings,
  22. pipeline_repo: PipelineRepository | None = None,
  23. decode_service: DecodeService | None = None,
  24. run_key: str | None = None,
  25. resume_policy: ResumePolicy | None = None,
  26. decode_limit: int = 100,
  27. ) -> CreationPipelineResult:
  28. resume_policy = resume_policy or ResumePolicy()
  29. pipeline_run = (
  30. pipeline_repo.create_pipeline_run(
  31. run_key=run_key or f"creation-pipeline:{batch_id}",
  32. batch_id=batch_id,
  33. status="running",
  34. metadata=resume_policy.model_dump(),
  35. )
  36. if pipeline_repo
  37. else None
  38. )
  39. acquisition_result = run_acquisition_stage(
  40. acquisition_repo=acquisition_repo,
  41. batch_id=batch_id,
  42. settings=settings,
  43. pipeline_repo=pipeline_repo,
  44. pipeline_run=pipeline_run,
  45. resume=resume_policy.resume,
  46. skip_done=resume_policy.skip_done,
  47. )
  48. decode_result = None
  49. if decode_service is not None:
  50. decode_result = run_decode_stage(
  51. candidate_repo=acquisition_repo,
  52. decode_service=decode_service,
  53. run_id=acquisition_result.acquisition.run_id,
  54. limit=decode_limit,
  55. )
  56. return CreationPipelineResult(
  57. pipeline_run=pipeline_run,
  58. acquisition=acquisition_result,
  59. decode=decode_result,
  60. )