from __future__ import annotations import argparse import json import logging from supply_infra.pipeline.reconciler import PipelineReconciler, reconcile_once from supply_infra.pipeline.run_service import ( get_pipeline_run, submit_pipeline_run, ) from supply_infra.pipeline.worker import PipelineWorker def _parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(description="SupplyAgent durable pipeline control") sub = parser.add_subparsers(dest="command", required=True) submit = sub.add_parser("submit") submit.add_argument("--biz-dt") submit.add_argument("--reason") inspect = sub.add_parser("inspect") inspect.add_argument("--run-id", required=True) sub.add_parser("worker") sub.add_parser("reconciler") sub.add_parser("reconcile-once") return parser def main() -> None: logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) args = _parser().parse_args() if args.command == "submit": result = submit_pipeline_run( biz_dt=args.biz_dt, trigger_type="cli", trigger_source="cli", trigger_reason=args.reason, ).to_dict() print(json.dumps(result, ensure_ascii=False)) return if args.command == "inspect": result = get_pipeline_run(args.run_id) print(json.dumps(result, ensure_ascii=False, default=str)) raise SystemExit(0 if result is not None else 1) if args.command == "worker": PipelineWorker().run_forever() return if args.command == "reconciler": PipelineReconciler().run_forever() return print(json.dumps(reconcile_once(), ensure_ascii=False)) if __name__ == "__main__": main()