| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364 |
- 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")
- submit.add_argument(
- "--dry-run",
- action="store_true",
- help="execute all local steps but record AIGC effects without dispatching them",
- )
- 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,
- dry_run=args.dry_run,
- ).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()
|