| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187 |
- import pytest
- from app.core.orchestration.migration.reconciliation import (
- ExecutionSnapshot,
- ReconciliationPolicy,
- WorkflowReconciler,
- )
- def snapshot(
- *,
- engine,
- execution_id,
- content_hash="content-a",
- amount_sum=30.0,
- duration=12.0,
- partition_duration=None,
- cpu_seconds=2.0,
- ):
- partition_duration = (
- duration if partition_duration is None else partition_duration
- )
- return ExecutionSnapshot.from_dict(
- {
- "engine": engine,
- "execution_id": execution_id,
- "status": "success",
- "duration_seconds": duration,
- "resource_usage": {
- "cpu_seconds": cpu_seconds,
- "peak_memory_mb": 128.0,
- },
- "nodes": {
- "read_orders": {
- "partitions": {
- "2026-07-19": {
- "row_count": 2,
- "primary_keys_hash": "keys-a",
- "content_hash": content_hash,
- "error_count": 0,
- "duration_seconds": partition_duration,
- "resource_usage": {"rows_scanned": 2},
- "aggregates": {
- "amount_sum": amount_sum,
- "generated_epoch": 1000,
- },
- }
- }
- }
- },
- }
- )
- def test_equivalent_runs_pass_all_core_dimensions_without_raw_rows():
- report = WorkflowReconciler().compare(
- snapshot(engine="n8n", execution_id="n8n-1"),
- snapshot(engine="kestra", execution_id="kestra-1"),
- ReconciliationPolicy(),
- )
- assert report["status"] == "passed"
- assert report["equivalent"] is True
- assert report["difference_count"] == 0
- assert report["metrics"]["compared_partitions"] == 1
- assert report["metrics"]["row_count_delta"] == 0
- assert "rows" not in repr(report).lower()
- def test_report_locates_hash_and_aggregate_difference_to_node_partition():
- report = WorkflowReconciler().compare(
- snapshot(engine="n8n", execution_id="n8n-1"),
- snapshot(
- engine="kestra",
- execution_id="kestra-1",
- content_hash="content-b",
- amount_sum=32.0,
- ),
- ReconciliationPolicy(
- numeric_tolerances={
- "nodes.*.partitions.*.aggregates.amount_sum": 0.5
- }
- ),
- )
- assert report["status"] == "failed"
- assert report["equivalent"] is False
- assert {
- (
- difference["node_id"],
- difference["partition"],
- difference["metric"],
- )
- for difference in report["differences"]
- } >= {
- ("read_orders", "2026-07-19", "content_hash"),
- ("read_orders", "2026-07-19", "aggregates.amount_sum"),
- }
- def test_explicit_tolerance_and_nondeterministic_ignore_can_pass():
- policy = ReconciliationPolicy(
- ignored_paths={
- "duration_seconds",
- "resource_usage.cpu_seconds",
- "nodes.*.partitions.*.duration_seconds",
- "nodes.*.partitions.*.aggregates.generated_epoch",
- },
- numeric_tolerances={
- "nodes.*.partitions.*.aggregates.amount_sum": 0.1
- },
- )
- report = WorkflowReconciler().compare(
- snapshot(engine="n8n", execution_id="n8n-1"),
- snapshot(
- engine="kestra",
- execution_id="kestra-1",
- amount_sum=30.05,
- duration=30.0,
- cpu_seconds=5.0,
- ),
- policy,
- )
- assert report["status"] == "passed"
- assert report["policy"]["ignored_paths"] == sorted(policy.ignored_paths)
- assert report["policy"]["numeric_tolerances"] == {
- "nodes.*.partitions.*.aggregates.amount_sum": 0.1
- }
- @pytest.mark.parametrize(
- "path",
- [
- "status",
- "nodes.*.partitions.*.row_count",
- "nodes.*.partitions.*.primary_keys_hash",
- "nodes.*.partitions.*.content_hash",
- "nodes.*.partitions.*.error_count",
- "nodes.read_orders.partitions.2026-07-19.row_count",
- "nodes.*.partitions.*",
- ],
- )
- def test_policy_cannot_ignore_core_correctness_metrics(path):
- with pytest.raises(ValueError, match="core reconciliation metric"):
- ReconciliationPolicy(ignored_paths={path})
- def test_partition_duration_is_reconciled_independently():
- report = WorkflowReconciler().compare(
- snapshot(engine="n8n", execution_id="n8n-1"),
- snapshot(
- engine="kestra",
- execution_id="kestra-1",
- partition_duration=15.0,
- ),
- ReconciliationPolicy(),
- )
- assert {
- difference["metric"] for difference in report["differences"]
- } == {"duration_seconds"}
- assert report["differences"][0]["node_id"] == "read_orders"
- def test_missing_partition_is_reported_without_silent_truncation():
- baseline = snapshot(engine="n8n", execution_id="n8n-1")
- candidate = snapshot(engine="kestra", execution_id="kestra-1")
- candidate.nodes["read_orders"]["partitions"] = {}
- report = WorkflowReconciler().compare(
- baseline,
- candidate,
- ReconciliationPolicy(),
- )
- assert report["status"] == "failed"
- assert report["differences"] == [
- {
- "node_id": "read_orders",
- "partition": "2026-07-19",
- "metric": "partition_presence",
- "baseline": True,
- "candidate": False,
- }
- ]
- assert report["metrics"]["row_count_delta"] == -2
|