Add durable reconciliation row evidence
This commit is contained in:
@@ -30,7 +30,7 @@ nodes by purpose:
|
|||||||
| Combine | Append rows, inner/outer/semi/anti joins |
|
| Combine | Append rows, inner/outer/semi/anti joins |
|
||||||
| Filter | Column filter, typed-expression filter, remove duplicates |
|
| Filter | Column filter, typed-expression filter, remove duplicates |
|
||||||
| Transform | Select, derive, ordered calculations, typed expression, conversion, replacement, aggregate, partitioned rank, sort, limit, reusable subflow |
|
| Transform | Select, derive, ordered calculations, typed expression, conversion, replacement, aggregate, partitioned rank, sort, limit, reusable subflow |
|
||||||
| Quality | Quality rules, keyed reconciliation |
|
| Quality | Quality rules; keyed reconciliation with stable row identity, explicit before/after evidence, and rerun invalidation hashes |
|
||||||
| Output | Preview output |
|
| Output | Preview output |
|
||||||
|
|
||||||
Join nodes have explicit left and right ports. Append nodes accept two or more
|
Join nodes have explicit left and right ports. Append nodes accept two or more
|
||||||
|
|||||||
@@ -7,7 +7,23 @@
|
|||||||
"monthly_status": "open",
|
"monthly_status": "open",
|
||||||
"monthly_amount": 10,
|
"monthly_amount": 10,
|
||||||
"_reconciliation_status": "match",
|
"_reconciliation_status": "match",
|
||||||
"_reconciliation_differences": []
|
"_reconciliation_differences": [],
|
||||||
|
"_reconciliation_changes": [],
|
||||||
|
"_reconciliation_key": [
|
||||||
|
"A-1"
|
||||||
|
],
|
||||||
|
"_reconciliation_key_hash": "04c89a1f5d4765e5c245cbb4bb2964de94b34494bd5af928428042b450d16c09",
|
||||||
|
"_reconciliation_input_hash": "34e8f55e551b289229822bb93bdd30c9ea42832ed67ad67d99050a3ef9fbdb18",
|
||||||
|
"_reconciliation_before": {
|
||||||
|
"case_id": "A-1",
|
||||||
|
"status": "open",
|
||||||
|
"amount": 10
|
||||||
|
},
|
||||||
|
"_reconciliation_after": {
|
||||||
|
"case_id": "A-1",
|
||||||
|
"status": "open",
|
||||||
|
"amount": 10
|
||||||
|
}
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"case_id": "A-2",
|
"case_id": "A-2",
|
||||||
@@ -17,13 +33,50 @@
|
|||||||
"monthly_status": "closed",
|
"monthly_status": "closed",
|
||||||
"monthly_amount": 25,
|
"monthly_amount": 25,
|
||||||
"_reconciliation_status": "changed",
|
"_reconciliation_status": "changed",
|
||||||
"_reconciliation_differences": ["amount"]
|
"_reconciliation_differences": [
|
||||||
|
"amount"
|
||||||
|
],
|
||||||
|
"_reconciliation_changes": [
|
||||||
|
{
|
||||||
|
"expected_field": "amount",
|
||||||
|
"observed_field": "amount",
|
||||||
|
"expected": 20,
|
||||||
|
"observed": 25
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"_reconciliation_key": [
|
||||||
|
"A-2"
|
||||||
|
],
|
||||||
|
"_reconciliation_key_hash": "20f1552752978558506b159b92bc738d34f484d3d76068173b0c96e1fc4a8c2d",
|
||||||
|
"_reconciliation_input_hash": "182de78bc3297d90cd4b02392bfb4a53fc9d9c29851d567f68e10be477a57dd8",
|
||||||
|
"_reconciliation_before": {
|
||||||
|
"case_id": "A-2",
|
||||||
|
"status": "closed",
|
||||||
|
"amount": 20
|
||||||
|
},
|
||||||
|
"_reconciliation_after": {
|
||||||
|
"case_id": "A-2",
|
||||||
|
"status": "closed",
|
||||||
|
"amount": 25
|
||||||
|
}
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"monthly_case_id": "A-3",
|
"monthly_case_id": "A-3",
|
||||||
"monthly_status": "open",
|
"monthly_status": "open",
|
||||||
"monthly_amount": 30,
|
"monthly_amount": 30,
|
||||||
"_reconciliation_status": "missing_expected",
|
"_reconciliation_status": "missing_expected",
|
||||||
"_reconciliation_differences": []
|
"_reconciliation_differences": [],
|
||||||
|
"_reconciliation_changes": [],
|
||||||
|
"_reconciliation_key": [
|
||||||
|
"A-3"
|
||||||
|
],
|
||||||
|
"_reconciliation_key_hash": "fedd5f48e5721d007a161015a4bab211555795d54a9492afa1b018050d3767f2",
|
||||||
|
"_reconciliation_input_hash": "fa82f93e97badd67750af911c545d703bb263a336bf8ba6d995eea5377abc2a0",
|
||||||
|
"_reconciliation_before": null,
|
||||||
|
"_reconciliation_after": {
|
||||||
|
"case_id": "A-3",
|
||||||
|
"status": "open",
|
||||||
|
"amount": 30
|
||||||
|
}
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -878,12 +878,32 @@ def _reconcile_rows(
|
|||||||
right_keys = [str(item) for item in config["right_keys"]]
|
right_keys = [str(item) for item in config["right_keys"]]
|
||||||
right_prefix = str(config.get("right_prefix", "observed_"))
|
right_prefix = str(config.get("right_prefix", "observed_"))
|
||||||
comparisons = _comparison_fields(config.get("compare_columns"))
|
comparisons = _comparison_fields(config.get("compare_columns"))
|
||||||
left_index = _unique_row_index(left_rows, left_keys, node_id=node_id)
|
left_index = _unique_row_index(
|
||||||
right_index = _unique_row_index(right_rows, right_keys, node_id=node_id)
|
left_rows,
|
||||||
|
left_keys,
|
||||||
|
node_id=node_id,
|
||||||
|
input_label="expected",
|
||||||
|
)
|
||||||
|
right_index = _unique_row_index(
|
||||||
|
right_rows,
|
||||||
|
right_keys,
|
||||||
|
node_id=node_id,
|
||||||
|
input_label="observed",
|
||||||
|
)
|
||||||
output: list[dict[str, Any]] = []
|
output: list[dict[str, Any]] = []
|
||||||
for key in dict.fromkeys((*left_index, *right_index)):
|
for key in dict.fromkeys((*left_index, *right_index)):
|
||||||
left = left_index.get(key)
|
left = left_index.get(key)
|
||||||
right = right_index.get(key)
|
right = right_index.get(key)
|
||||||
|
key_values = [
|
||||||
|
(left or {}).get(left_name)
|
||||||
|
if left is not None
|
||||||
|
else (right or {}).get(right_name)
|
||||||
|
for left_name, right_name in zip(
|
||||||
|
left_keys,
|
||||||
|
right_keys,
|
||||||
|
strict=True,
|
||||||
|
)
|
||||||
|
]
|
||||||
item = dict(left or {})
|
item = dict(left or {})
|
||||||
if right is not None:
|
if right is not None:
|
||||||
item.update({f"{right_prefix}{name}": value for name, value in right.items()})
|
item.update({f"{right_prefix}{name}": value for name, value in right.items()})
|
||||||
@@ -895,18 +915,55 @@ def _reconcile_rows(
|
|||||||
differences = []
|
differences = []
|
||||||
else:
|
else:
|
||||||
fields = comparisons or tuple((name, name) for name in left if name not in left_keys)
|
fields = comparisons or tuple((name, name) for name in left if name not in left_keys)
|
||||||
differences = [
|
changes = [
|
||||||
left_name
|
{
|
||||||
|
"expected_field": left_name,
|
||||||
|
"observed_field": right_name,
|
||||||
|
"expected": left.get(left_name),
|
||||||
|
"observed": right.get(right_name),
|
||||||
|
}
|
||||||
for left_name, right_name in fields
|
for left_name, right_name in fields
|
||||||
if left.get(left_name) != right.get(right_name)
|
if left.get(left_name) != right.get(right_name)
|
||||||
]
|
]
|
||||||
|
differences = [str(change["expected_field"]) for change in changes]
|
||||||
status = "changed" if differences else "match"
|
status = "changed" if differences else "match"
|
||||||
|
if left is None or right is None:
|
||||||
|
changes = []
|
||||||
item["_reconciliation_status"] = status
|
item["_reconciliation_status"] = status
|
||||||
item["_reconciliation_differences"] = differences
|
item["_reconciliation_differences"] = differences
|
||||||
|
item["_reconciliation_changes"] = changes
|
||||||
|
item["_reconciliation_key"] = key_values
|
||||||
|
item["_reconciliation_key_hash"] = _reconciliation_hash(
|
||||||
|
{
|
||||||
|
"left_keys": left_keys,
|
||||||
|
"right_keys": right_keys,
|
||||||
|
"values": key_values,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
item["_reconciliation_input_hash"] = _reconciliation_hash(
|
||||||
|
{
|
||||||
|
"key": key_values,
|
||||||
|
"expected": left,
|
||||||
|
"observed": right,
|
||||||
|
"comparisons": comparisons,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
item["_reconciliation_before"] = dict(left) if left is not None else None
|
||||||
|
item["_reconciliation_after"] = dict(right) if right is not None else None
|
||||||
output.append(item)
|
output.append(item)
|
||||||
return output
|
return output
|
||||||
|
|
||||||
|
|
||||||
|
def _reconciliation_hash(value: object) -> str:
|
||||||
|
encoded = json.dumps(
|
||||||
|
value,
|
||||||
|
sort_keys=True,
|
||||||
|
separators=(",", ":"),
|
||||||
|
default=str,
|
||||||
|
).encode("utf-8")
|
||||||
|
return hashlib.sha256(encoded).hexdigest()
|
||||||
|
|
||||||
|
|
||||||
def _comparison_fields(value: object) -> tuple[tuple[str, str], ...]:
|
def _comparison_fields(value: object) -> tuple[tuple[str, str], ...]:
|
||||||
if not isinstance(value, list):
|
if not isinstance(value, list):
|
||||||
return ()
|
return ()
|
||||||
@@ -927,13 +984,14 @@ def _unique_row_index(
|
|||||||
columns: list[str],
|
columns: list[str],
|
||||||
*,
|
*,
|
||||||
node_id: str,
|
node_id: str,
|
||||||
|
input_label: str,
|
||||||
) -> dict[tuple[Any, ...], dict[str, Any]]:
|
) -> dict[tuple[Any, ...], dict[str, Any]]:
|
||||||
index: dict[tuple[Any, ...], dict[str, Any]] = {}
|
index: dict[tuple[Any, ...], dict[str, Any]] = {}
|
||||||
for row in rows:
|
for row in rows:
|
||||||
key = tuple(_hashable(row.get(column)) for column in columns)
|
key = tuple(_hashable(row.get(column)) for column in columns)
|
||||||
if key in index:
|
if key in index:
|
||||||
raise PipelineExecutionError(
|
raise PipelineExecutionError(
|
||||||
f"Reconciliation keys are not unique: {key!r}.",
|
f"Reconciliation keys are not unique in the {input_label} input.",
|
||||||
node_id=node_id,
|
node_id=node_id,
|
||||||
)
|
)
|
||||||
index[key] = row
|
index[key] = row
|
||||||
|
|||||||
@@ -200,6 +200,7 @@ DOCUMENTATION = (
|
|||||||
"Every graph node declares typed inputs, configuration, output schema, and validation rules. "
|
"Every graph node declares typed inputs, configuration, output schema, and validation rules. "
|
||||||
"Source nodes pin inline content or governed Datasource references; combine, filter, transform, "
|
"Source nodes pin inline content or governed Datasource references; combine, filter, transform, "
|
||||||
"quality, reconciliation, reusable-subflow, and output nodes remain explicit in the canonical graph. "
|
"quality, reconciliation, reusable-subflow, and output nodes remain explicit in the canonical graph. "
|
||||||
|
"Reconciliation rows expose stable key hashes, explicit before/after values, and input hashes so a later governed human decision can be replayed only while its inputs still match. "
|
||||||
"Expressions use the typed Dataflow expression language and never execute arbitrary host or database "
|
"Expressions use the typed Dataflow expression language and never execute arbitrary host or database "
|
||||||
"code. Selecting a node may request a bounded intermediate preview; preview rows are transient, "
|
"code. Selecting a node may request a bounded intermediate preview; preview rows are transient, "
|
||||||
"privacy-filtered for the actor, and are not retained as run output. SQL editing compiles into the same "
|
"privacy-filtered for the actor, and are not retained as run output. SQL editing compiles into the same "
|
||||||
|
|||||||
@@ -510,7 +510,7 @@ _NODE_TYPES = (
|
|||||||
type="reconcile.compare",
|
type="reconcile.compare",
|
||||||
category="quality",
|
category="quality",
|
||||||
label="Reconcile tables",
|
label="Reconcile tables",
|
||||||
description="Compare keyed rows and expose missing, matching, and changed records.",
|
description="Compare keyed rows with stable identity, before/after evidence, and decision invalidation fingerprints.",
|
||||||
icon="scan-search",
|
icon="scan-search",
|
||||||
input_ports=(
|
input_ports=(
|
||||||
NodePortDefinition(id="left", label="Expected"),
|
NodePortDefinition(id="left", label="Expected"),
|
||||||
|
|||||||
@@ -612,6 +612,12 @@ def _reconcile(context: SchemaPropagationContext) -> SchemaPropagationResult:
|
|||||||
(
|
(
|
||||||
"_reconciliation_status",
|
"_reconciliation_status",
|
||||||
"_reconciliation_differences",
|
"_reconciliation_differences",
|
||||||
|
"_reconciliation_changes",
|
||||||
|
"_reconciliation_key",
|
||||||
|
"_reconciliation_key_hash",
|
||||||
|
"_reconciliation_input_hash",
|
||||||
|
"_reconciliation_before",
|
||||||
|
"_reconciliation_after",
|
||||||
)
|
)
|
||||||
),
|
),
|
||||||
open=left_state.open or right_state.open,
|
open=left_state.open or right_state.open,
|
||||||
@@ -623,6 +629,12 @@ def _reconcile(context: SchemaPropagationContext) -> SchemaPropagationResult:
|
|||||||
},
|
},
|
||||||
"_reconciliation_status": "string",
|
"_reconciliation_status": "string",
|
||||||
"_reconciliation_differences": "array",
|
"_reconciliation_differences": "array",
|
||||||
|
"_reconciliation_changes": "array",
|
||||||
|
"_reconciliation_key": "array",
|
||||||
|
"_reconciliation_key_hash": "string",
|
||||||
|
"_reconciliation_input_hash": "string",
|
||||||
|
"_reconciliation_before": "object",
|
||||||
|
"_reconciliation_after": "object",
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
return SchemaPropagationResult(state, tuple(diagnostics))
|
return SchemaPropagationResult(state, tuple(diagnostics))
|
||||||
|
|||||||
+123
-1
@@ -2,7 +2,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import unittest
|
import unittest
|
||||||
|
|
||||||
from govoplan_dataflow.backend.executor import execute_preview
|
from govoplan_dataflow.backend.executor import PipelineExecutionError, execute_preview
|
||||||
from govoplan_dataflow.backend.expressions import (
|
from govoplan_dataflow.backend.expressions import (
|
||||||
ExpressionError,
|
ExpressionError,
|
||||||
evaluate_expression,
|
evaluate_expression,
|
||||||
@@ -389,6 +389,128 @@ class DataflowOperatorTests(unittest.TestCase):
|
|||||||
["changed", "missing_observed", "missing_expected"],
|
["changed", "missing_observed", "missing_expected"],
|
||||||
[item["_reconciliation_status"] for item in result.rows],
|
[item["_reconciliation_status"] for item in result.rows],
|
||||||
)
|
)
|
||||||
|
changed = result.rows[0]
|
||||||
|
self.assertEqual(["1"], changed["_reconciliation_key"])
|
||||||
|
self.assertEqual(
|
||||||
|
{"id": "1", "amount": 10},
|
||||||
|
changed["_reconciliation_before"],
|
||||||
|
)
|
||||||
|
self.assertEqual(
|
||||||
|
{"id": "1", "amount": 11},
|
||||||
|
changed["_reconciliation_after"],
|
||||||
|
)
|
||||||
|
self.assertEqual(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"expected_field": "amount",
|
||||||
|
"observed_field": "amount",
|
||||||
|
"expected": 10,
|
||||||
|
"observed": 11,
|
||||||
|
}
|
||||||
|
],
|
||||||
|
changed["_reconciliation_changes"],
|
||||||
|
)
|
||||||
|
replay = execute_preview(graph, row_limit=100).rows[0]
|
||||||
|
self.assertEqual(
|
||||||
|
changed["_reconciliation_key_hash"],
|
||||||
|
replay["_reconciliation_key_hash"],
|
||||||
|
)
|
||||||
|
self.assertEqual(
|
||||||
|
changed["_reconciliation_input_hash"],
|
||||||
|
replay["_reconciliation_input_hash"],
|
||||||
|
)
|
||||||
|
|
||||||
|
changed_nodes = [
|
||||||
|
(
|
||||||
|
graph_node.model_copy(
|
||||||
|
update={
|
||||||
|
"config": {
|
||||||
|
**graph_node.config,
|
||||||
|
"rows": [
|
||||||
|
{"id": "1", "amount": 12},
|
||||||
|
{"id": "3", "amount": 30},
|
||||||
|
],
|
||||||
|
}
|
||||||
|
}
|
||||||
|
)
|
||||||
|
if graph_node.id == "observed"
|
||||||
|
else graph_node
|
||||||
|
)
|
||||||
|
for graph_node in graph.nodes
|
||||||
|
]
|
||||||
|
rerun = execute_preview(
|
||||||
|
graph.model_copy(update={"nodes": changed_nodes}),
|
||||||
|
row_limit=100,
|
||||||
|
).rows[0]
|
||||||
|
self.assertEqual(
|
||||||
|
changed["_reconciliation_key_hash"],
|
||||||
|
rerun["_reconciliation_key_hash"],
|
||||||
|
)
|
||||||
|
self.assertNotEqual(
|
||||||
|
changed["_reconciliation_input_hash"],
|
||||||
|
rerun["_reconciliation_input_hash"],
|
||||||
|
)
|
||||||
|
self.assertIsNone(result.rows[1]["_reconciliation_after"])
|
||||||
|
self.assertIsNone(result.rows[2]["_reconciliation_before"])
|
||||||
|
|
||||||
|
def test_reconciliation_rejects_ambiguous_keys_without_echoing_values(self) -> None:
|
||||||
|
graph = PipelineGraph(
|
||||||
|
nodes=[
|
||||||
|
node(
|
||||||
|
"expected",
|
||||||
|
"source.inline",
|
||||||
|
{
|
||||||
|
"source_name": "expected",
|
||||||
|
"rows": [
|
||||||
|
{"id": "private-key", "amount": 10},
|
||||||
|
{"id": "private-key", "amount": 11},
|
||||||
|
],
|
||||||
|
},
|
||||||
|
x=0,
|
||||||
|
),
|
||||||
|
node(
|
||||||
|
"observed",
|
||||||
|
"source.inline",
|
||||||
|
{"source_name": "observed", "rows": []},
|
||||||
|
x=0,
|
||||||
|
),
|
||||||
|
node(
|
||||||
|
"reconcile",
|
||||||
|
"reconcile.compare",
|
||||||
|
{
|
||||||
|
"left_keys": ["id"],
|
||||||
|
"right_keys": ["id"],
|
||||||
|
"compare_columns": ["amount"],
|
||||||
|
"right_prefix": "observed_",
|
||||||
|
},
|
||||||
|
x=300,
|
||||||
|
),
|
||||||
|
node("output", "output", {}, x=600),
|
||||||
|
],
|
||||||
|
edges=[
|
||||||
|
GraphEdge(
|
||||||
|
id="e1",
|
||||||
|
source="expected",
|
||||||
|
target="reconcile",
|
||||||
|
target_port="left",
|
||||||
|
),
|
||||||
|
GraphEdge(
|
||||||
|
id="e2",
|
||||||
|
source="observed",
|
||||||
|
target="reconcile",
|
||||||
|
target_port="right",
|
||||||
|
),
|
||||||
|
GraphEdge(id="e3", source="reconcile", target="output"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
|
||||||
|
with self.assertRaisesRegex(
|
||||||
|
PipelineExecutionError,
|
||||||
|
"not unique in the expected input",
|
||||||
|
) as raised:
|
||||||
|
execute_preview(graph, row_limit=100)
|
||||||
|
|
||||||
|
self.assertNotIn("private-key", str(raised.exception))
|
||||||
|
|
||||||
def test_parameterized_subflow_runs_a_pinned_graph(self) -> None:
|
def test_parameterized_subflow_runs_a_pinned_graph(self) -> None:
|
||||||
nested = {
|
nested = {
|
||||||
|
|||||||
Reference in New Issue
Block a user