#!/usr/bin/env python3 """Run formal acquisition, decode every coarse-hit item, and build payload drafts.""" from __future__ import annotations import argparse import json from dataclasses import asdict, is_dataclass from uuid import UUID from acquisition.repositories.postgres import PostgresAcquisitionRepository from acquisition.runner import DEFAULT_PLATFORMS, run_batch from core.config import CreationDbConfig, Settings from core.db_session import transaction from decode_content.repositories.postgres import PostgresDecodeRepository from decode_content.service import DecodeService from pipeline.decode_runner import run_decode_stage def _model_dump(value): if hasattr(value, "model_dump"): return value.model_dump(mode="json") if is_dataclass(value): return asdict(value) if isinstance(value, dict): return value return dict(value) def _dry_ingest_payloads(repo: PostgresDecodeRepository, outputs: list) -> list[dict]: records: list[dict] = [] for output in outputs: for draft in output.payload_drafts: if draft.id is None: continue repo.mark_payload_draft_ingested(draft.id) record = repo.save_ingest_record( payload_draft_id=draft.id, target_system="dry-run", target_id=str(draft.id), status="ingested", response_payload={ "dry_run": True, "note": "payload generated by formal creation pipeline; external ingest API not called", "payload": draft.payload, }, ) records.append(_model_dump(record)) return records def parse_args(argv: list[str] | None = None) -> argparse.Namespace: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--batch-id", required=True, help="Formal query batch UUID") parser.add_argument( "--platform", action="append", choices=DEFAULT_PLATFORMS, help="Platform to run. Repeat to run multiple platforms. Default: all.", ) parser.add_argument("--search-limit", type=int, default=10) parser.add_argument("--display-limit", type=int, default=5) parser.add_argument("--decode-limit", type=int, default=100) parser.add_argument("--run-key") parser.add_argument("--env-file", default=".env") parser.add_argument("--no-resume", action="store_true") parser.add_argument("--no-skip-done", action="store_true") parser.add_argument("--no-dry-ingest-record", action="store_true") return parser.parse_args(argv) def main(argv: list[str] | None = None) -> int: args = parse_args(argv) settings = Settings.from_env(args.env_file) db_config = CreationDbConfig.from_env(args.env_file) platforms = tuple(args.platform or DEFAULT_PLATFORMS) batch_id = UUID(args.batch_id) with transaction(db_config) as conn: acquisition_repo = PostgresAcquisitionRepository(conn) acquisition = run_batch( acquisition_repo, batch_id=batch_id, settings=settings, platforms=platforms, search_limit=args.search_limit, display_limit=args.display_limit, classify=True, resume=not args.no_resume, skip_done=not args.no_skip_done, run_key=args.run_key, ) with transaction(db_config) as conn: acquisition_repo = PostgresAcquisitionRepository(conn) decode_repo = PostgresDecodeRepository(conn) decode_service = DecodeService(settings=settings, repository=decode_repo) decode = run_decode_stage( candidate_repo=acquisition_repo, decode_service=decode_service, run_id=acquisition.run_id, limit=args.decode_limit, ) ingest_records = [] if args.no_dry_ingest_record else _dry_ingest_payloads(decode_repo, decode.outputs) print(json.dumps( { "batch_id": str(batch_id), "run_id": str(acquisition.run_id), "platforms": list(platforms), "acquisition": _model_dump(acquisition), "decode": _model_dump(decode), "dry_ingest_records": ingest_records, }, ensure_ascii=False, default=str, indent=2, )) return 0 if __name__ == "__main__": raise SystemExit(main())