test_workflow_reconciliation.py 5.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187
  1. import pytest
  2. from app.core.orchestration.migration.reconciliation import (
  3. ExecutionSnapshot,
  4. ReconciliationPolicy,
  5. WorkflowReconciler,
  6. )
  7. def snapshot(
  8. *,
  9. engine,
  10. execution_id,
  11. content_hash="content-a",
  12. amount_sum=30.0,
  13. duration=12.0,
  14. partition_duration=None,
  15. cpu_seconds=2.0,
  16. ):
  17. partition_duration = (
  18. duration if partition_duration is None else partition_duration
  19. )
  20. return ExecutionSnapshot.from_dict(
  21. {
  22. "engine": engine,
  23. "execution_id": execution_id,
  24. "status": "success",
  25. "duration_seconds": duration,
  26. "resource_usage": {
  27. "cpu_seconds": cpu_seconds,
  28. "peak_memory_mb": 128.0,
  29. },
  30. "nodes": {
  31. "read_orders": {
  32. "partitions": {
  33. "2026-07-19": {
  34. "row_count": 2,
  35. "primary_keys_hash": "keys-a",
  36. "content_hash": content_hash,
  37. "error_count": 0,
  38. "duration_seconds": partition_duration,
  39. "resource_usage": {"rows_scanned": 2},
  40. "aggregates": {
  41. "amount_sum": amount_sum,
  42. "generated_epoch": 1000,
  43. },
  44. }
  45. }
  46. }
  47. },
  48. }
  49. )
  50. def test_equivalent_runs_pass_all_core_dimensions_without_raw_rows():
  51. report = WorkflowReconciler().compare(
  52. snapshot(engine="n8n", execution_id="n8n-1"),
  53. snapshot(engine="kestra", execution_id="kestra-1"),
  54. ReconciliationPolicy(),
  55. )
  56. assert report["status"] == "passed"
  57. assert report["equivalent"] is True
  58. assert report["difference_count"] == 0
  59. assert report["metrics"]["compared_partitions"] == 1
  60. assert report["metrics"]["row_count_delta"] == 0
  61. assert "rows" not in repr(report).lower()
  62. def test_report_locates_hash_and_aggregate_difference_to_node_partition():
  63. report = WorkflowReconciler().compare(
  64. snapshot(engine="n8n", execution_id="n8n-1"),
  65. snapshot(
  66. engine="kestra",
  67. execution_id="kestra-1",
  68. content_hash="content-b",
  69. amount_sum=32.0,
  70. ),
  71. ReconciliationPolicy(
  72. numeric_tolerances={
  73. "nodes.*.partitions.*.aggregates.amount_sum": 0.5
  74. }
  75. ),
  76. )
  77. assert report["status"] == "failed"
  78. assert report["equivalent"] is False
  79. assert {
  80. (
  81. difference["node_id"],
  82. difference["partition"],
  83. difference["metric"],
  84. )
  85. for difference in report["differences"]
  86. } >= {
  87. ("read_orders", "2026-07-19", "content_hash"),
  88. ("read_orders", "2026-07-19", "aggregates.amount_sum"),
  89. }
  90. def test_explicit_tolerance_and_nondeterministic_ignore_can_pass():
  91. policy = ReconciliationPolicy(
  92. ignored_paths={
  93. "duration_seconds",
  94. "resource_usage.cpu_seconds",
  95. "nodes.*.partitions.*.duration_seconds",
  96. "nodes.*.partitions.*.aggregates.generated_epoch",
  97. },
  98. numeric_tolerances={
  99. "nodes.*.partitions.*.aggregates.amount_sum": 0.1
  100. },
  101. )
  102. report = WorkflowReconciler().compare(
  103. snapshot(engine="n8n", execution_id="n8n-1"),
  104. snapshot(
  105. engine="kestra",
  106. execution_id="kestra-1",
  107. amount_sum=30.05,
  108. duration=30.0,
  109. cpu_seconds=5.0,
  110. ),
  111. policy,
  112. )
  113. assert report["status"] == "passed"
  114. assert report["policy"]["ignored_paths"] == sorted(policy.ignored_paths)
  115. assert report["policy"]["numeric_tolerances"] == {
  116. "nodes.*.partitions.*.aggregates.amount_sum": 0.1
  117. }
  118. @pytest.mark.parametrize(
  119. "path",
  120. [
  121. "status",
  122. "nodes.*.partitions.*.row_count",
  123. "nodes.*.partitions.*.primary_keys_hash",
  124. "nodes.*.partitions.*.content_hash",
  125. "nodes.*.partitions.*.error_count",
  126. "nodes.read_orders.partitions.2026-07-19.row_count",
  127. "nodes.*.partitions.*",
  128. ],
  129. )
  130. def test_policy_cannot_ignore_core_correctness_metrics(path):
  131. with pytest.raises(ValueError, match="core reconciliation metric"):
  132. ReconciliationPolicy(ignored_paths={path})
  133. def test_partition_duration_is_reconciled_independently():
  134. report = WorkflowReconciler().compare(
  135. snapshot(engine="n8n", execution_id="n8n-1"),
  136. snapshot(
  137. engine="kestra",
  138. execution_id="kestra-1",
  139. partition_duration=15.0,
  140. ),
  141. ReconciliationPolicy(),
  142. )
  143. assert {
  144. difference["metric"] for difference in report["differences"]
  145. } == {"duration_seconds"}
  146. assert report["differences"][0]["node_id"] == "read_orders"
  147. def test_missing_partition_is_reported_without_silent_truncation():
  148. baseline = snapshot(engine="n8n", execution_id="n8n-1")
  149. candidate = snapshot(engine="kestra", execution_id="kestra-1")
  150. candidate.nodes["read_orders"]["partitions"] = {}
  151. report = WorkflowReconciler().compare(
  152. baseline,
  153. candidate,
  154. ReconciliationPolicy(),
  155. )
  156. assert report["status"] == "failed"
  157. assert report["differences"] == [
  158. {
  159. "node_id": "read_orders",
  160. "partition": "2026-07-19",
  161. "metric": "partition_presence",
  162. "baseline": True,
  163. "candidate": False,
  164. }
  165. ]
  166. assert report["metrics"]["row_count_delta"] == -2