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