| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120 |
- #!/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, IngestApiConfig, Settings
- from core.db_session import transaction
- from decode_content.ingest import KnowledgeIngestClient, ingest_payload_draft
- 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 _ingest_payloads(repo: PostgresDecodeRepository, outputs: list, *, dry_run: bool, env_file: str) -> list[dict]:
- records: list[dict] = []
- client = None if dry_run else KnowledgeIngestClient(IngestApiConfig.from_env(env_file))
- for output in outputs:
- for draft in output.payload_drafts:
- if draft.id is None:
- continue
- record = ingest_payload_draft(repo, draft, dry_run=dry_run, client=client)
- 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")
- parser.add_argument("--real-ingest", action="store_true", help="Call CK_INGEST_API_URL instead of dry-run records.")
- 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 _ingest_payloads(
- decode_repo,
- decode.outputs,
- dry_run=not args.real_ingest,
- env_file=args.env_file,
- )
- 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),
- "ingest_records": ingest_records,
- "dry_ingest_records": ingest_records if not args.real_ingest else [],
- },
- ensure_ascii=False,
- default=str,
- indent=2,
- ))
- return 0
- if __name__ == "__main__":
- raise SystemExit(main())
|