"""Formal query -> acquisition -> decode orchestration.""" from __future__ import annotations from dataclasses import dataclass from uuid import UUID from acquisition.repositories.base import AcquisitionRepository from core.config import Settings from decode_content.service import DecodeService from pipeline.acquisition_runner import AcquisitionStageResult, run_acquisition_stage from pipeline.decode_runner import DecodeBatchResult, run_decode_stage from pipeline.models import PipelineRun, ResumePolicy from pipeline.repository import PipelineRepository @dataclass(frozen=True) class CreationPipelineResult: pipeline_run: PipelineRun | None acquisition: AcquisitionStageResult decode: DecodeBatchResult | None = None def run_creation_pipeline( *, acquisition_repo: AcquisitionRepository, batch_id: UUID, settings: Settings, pipeline_repo: PipelineRepository | None = None, decode_service: DecodeService | None = None, run_key: str | None = None, resume_policy: ResumePolicy | None = None, decode_limit: int = 100, ) -> CreationPipelineResult: resume_policy = resume_policy or ResumePolicy() pipeline_run = ( pipeline_repo.create_pipeline_run( run_key=run_key or f"creation-pipeline:{batch_id}", batch_id=batch_id, status="running", metadata=resume_policy.model_dump(), ) if pipeline_repo else None ) acquisition_result = run_acquisition_stage( acquisition_repo=acquisition_repo, batch_id=batch_id, settings=settings, pipeline_repo=pipeline_repo, pipeline_run=pipeline_run, resume=resume_policy.resume, skip_done=resume_policy.skip_done, ) decode_result = None if decode_service is not None: decode_result = run_decode_stage( candidate_repo=acquisition_repo, decode_service=decode_service, run_id=acquisition_result.acquisition.run_id, limit=decode_limit, ) return CreationPipelineResult( pipeline_run=pipeline_run, acquisition=acquisition_result, decode=decode_result, )