reconcile_dual_run.py 1.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960
  1. """Compare credential-safe n8n and Kestra execution snapshots."""
  2. from __future__ import annotations
  3. import argparse
  4. import json
  5. from pathlib import Path
  6. from app.core.orchestration.migration.reconciliation import (
  7. ExecutionSnapshot,
  8. ReconciliationPolicy,
  9. WorkflowReconciler,
  10. )
  11. def reconcile_exports(baseline, candidate, policy=None):
  12. policy = policy or {}
  13. if not isinstance(policy, dict):
  14. raise ValueError("reconciliation policy must be an object")
  15. return WorkflowReconciler().compare(
  16. ExecutionSnapshot.from_dict(baseline),
  17. ExecutionSnapshot.from_dict(candidate),
  18. ReconciliationPolicy(
  19. ignored_paths=policy.get("ignored_paths"),
  20. numeric_tolerances=policy.get("numeric_tolerances"),
  21. ),
  22. )
  23. def main() -> int:
  24. parser = argparse.ArgumentParser(
  25. description="Reconcile bounded execution metrics without raw rows."
  26. )
  27. parser.add_argument("--n8n-snapshot", required=True)
  28. parser.add_argument("--kestra-snapshot", required=True)
  29. parser.add_argument("--policy-json")
  30. parser.add_argument("--output", required=True)
  31. args = parser.parse_args()
  32. baseline = json.loads(
  33. Path(args.n8n_snapshot).read_text(encoding="utf-8")
  34. )
  35. candidate = json.loads(
  36. Path(args.kestra_snapshot).read_text(encoding="utf-8")
  37. )
  38. policy = (
  39. json.loads(Path(args.policy_json).read_text(encoding="utf-8"))
  40. if args.policy_json
  41. else {}
  42. )
  43. report = reconcile_exports(baseline, candidate, policy)
  44. Path(args.output).write_text(
  45. json.dumps(report, ensure_ascii=False, indent=2) + "\n",
  46. encoding="utf-8",
  47. )
  48. return 0 if report["status"] == "passed" else 2
  49. if __name__ == "__main__":
  50. raise SystemExit(main())