from __future__ import annotations import copy import pytest from app.core.common.identifiers import new_governance_uid from app.core.data_flow.create_reconciliation import DataFlowCreateReconciler from app.core.data_flow.dataflows import DataFlowService class GraphResult: def __init__(self, record=None): self.record = record def single(self): return self.record class GraphSession: def __init__(self, *, existing=None, conflict=False): self.existing = existing self.conflict = conflict self.calls = [] def run(self, query, parameters=None, **kwargs): values = parameters or kwargs self.calls.append((query, values)) if query.startswith("CREATE CONSTRAINT"): return GraphResult() if "WHERE n.uid IS NULL OR n.uid <> $uid" in query: return GraphResult({"uid": "other"}) if self.conflict else GraphResult() if query.startswith("MERGE"): node = copy.deepcopy(self.existing or values["properties"]) self.existing = node return GraphResult({"n": node, "node_id": 73}) raise AssertionError(query) def __enter__(self): return self def __exit__(self, *_args): return None class GraphDriver: def __init__(self, session): self.graph_session = session def session(self): return self.graph_session def close(self): return None def governed_node(): return { "uid": new_governance_uid(), "name_zh": "客户治理生产线", "name_en": "customer_line", "script_type": "governed", "script_requirement": '{"dataflow_spec":{"schema_version":"2.0"}}', "script_path": "", } def test_governed_graph_create_installs_constraint_and_reconciles_same_uid( monkeypatch, ): node = governed_node() graph = GraphSession() monkeypatch.setattr( "app.core.data_flow.dataflows.connect_graph", lambda: GraphDriver(graph), ) first_id, first = DataFlowService._merge_governed_dataflow(node) second_id, second = DataFlowService._merge_governed_dataflow(node) assert first_id == second_id == 73 assert first == second assert any( call[0] == "CREATE CONSTRAINT data_flow_uid IF NOT EXISTS " "FOR (n:DataFlow) REQUIRE n.uid IS UNIQUE" for call in graph.calls ) assert sum(call[0].startswith("MERGE") for call in graph.calls) == 2 def test_governed_graph_reconcile_fails_closed_on_uid_or_name_conflict( monkeypatch, ): node = governed_node() graph = GraphSession(existing={**node, "name_zh": "被篡改的名称"}) monkeypatch.setattr( "app.core.data_flow.dataflows.connect_graph", lambda: GraphDriver(graph), ) with pytest.raises(ValueError, match="dataflow_uid_conflict"): DataFlowService._merge_governed_dataflow(node) name_conflict = GraphSession(conflict=True) monkeypatch.setattr( "app.core.data_flow.dataflows.connect_graph", lambda: GraphDriver(name_conflict), ) with pytest.raises(ValueError, match="dataflow_uid_conflict"): DataFlowService._merge_governed_dataflow(node) def test_completed_saga_replays_result_without_another_neo4j_write(monkeypatch): expected = {"id": 73, **governed_node()} class Repository: def load_published_assets(self, _flow): return {}, {} def begin_dataflow_create( self, _receipt, *, actor_uid, create_request, create_intent, ): assert actor_uid assert create_request["intent"] == create_intent return { "status": "completed", "dataflow_uid": expected["uid"], "result": expected, "request_digest": "a" * 64, } flow = { "schema_version": "2.0", "dataflow_uid": expected["uid"], "name": "客户治理生产线", "input_schema_refs": ["bd:customer:v1"], "output_schema_ref": "bd:customer_clean:v1", "components": [], } # Use the project's validator fixture shape through the repository-level # contract; the replay must happen before any graph write. from tests.core.data_rules.test_contracts import valid_dataflow_spec flow = valid_dataflow_spec() flow["dataflow_uid"] = expected["uid"] envelope = { "dataflow_spec": flow, "dataset_edges": { "source_table": flow["input_schema_refs"], "target_table": flow["output_schema_ref"], }, "migration_metadata": { "status": "migrated", "legacy_fields_present": False, "preserved_for_read_only": True, "governed_semantics": "dataflow_spec", }, } monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda _node: pytest.fail("Neo4j was called during completed replay"), ) result = DataFlowService.create_dataflow( { "name_zh": "客户治理生产线", "describe": "响应丢失重试", "script_type": "governed", "script_requirement": envelope, "draft_reservation": { "reservation_id": new_governance_uid(), "dataflow_uid": expected["uid"], "nonce": "response-loss", }, }, repository=Repository(), actor_uid=new_governance_uid(), ) assert result == expected def test_graph_failure_is_persisted_but_finalize_failure_leaves_reconcilable_lease( monkeypatch, ): from tests.core.data_rules.test_contracts import valid_dataflow_spec from tests.test_legacy_governance_cutover import PublishedAssetRepository flow = valid_dataflow_spec() actor = new_governance_uid() def payload(repository): receipt = repository.reserve_dataflow_draft(actor_uid=actor) receipt["dataflow_uid"] = flow["dataflow_uid"] return { "name_zh": "故障注入生产线", "describe": "故障注入", "script_type": "governed", "script_requirement": { "dataflow_spec": flow, "dataset_edges": { "source_table": flow["input_schema_refs"], "target_table": flow["output_schema_ref"], }, "migration_metadata": { "status": "migrated", "legacy_fields_present": False, "preserved_for_read_only": True, "governed_semantics": "dataflow_spec", }, }, "draft_reservation": { key: receipt[key] for key in ("reservation_id", "dataflow_uid", "nonce") }, } graph_failure = PublishedAssetRepository() failures = [] graph_failure.commit_dataflow_create_failure = ( lambda **kwargs: failures.append(kwargs) ) monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda _node: (_ for _ in ()).throw(RuntimeError("neo4j unavailable")), ) with pytest.raises(RuntimeError, match="neo4j unavailable"): DataFlowService.create_dataflow( payload(graph_failure), repository=graph_failure, actor_uid=actor ) assert failures[0]["error_code"] == "neo4j_create_failed" finalize_failure = PublishedAssetRepository() monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda node: (73, {**node, "id": 73}), ) finalize_failure.complete_dataflow_create = ( lambda **_kwargs: (_ for _ in ()).throw( RuntimeError("postgres finalize failed") ) ) with pytest.raises(RuntimeError, match="postgres finalize failed"): DataFlowService.create_dataflow( payload(finalize_failure), repository=finalize_failure, actor_uid=actor, ) def test_governed_tag_failure_never_finalizes_completed(monkeypatch): from tests.core.data_rules.test_contracts import valid_dataflow_spec from tests.test_legacy_governance_cutover import PublishedAssetRepository flow = valid_dataflow_spec() actor = new_governance_uid() repository = PublishedAssetRepository() receipt = repository.reserve_dataflow_draft(actor_uid=actor) receipt["dataflow_uid"] = flow["dataflow_uid"] failures = [] repository.commit_dataflow_create_failure = ( lambda **kwargs: failures.append(kwargs) ) monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda node: (73, {"id": 73, **node}), ) def fail_tags(_dataflow_id, _tags, *, strict): assert strict is True raise RuntimeError("neo4j tag merge failed") monkeypatch.setattr( DataFlowService, "_handle_tag_relationships", fail_tags ) with pytest.raises(RuntimeError, match="neo4j tag merge failed"): DataFlowService.create_dataflow( { "name_zh": "标签严格生产线", "describe": "严格标签合并验收", "script_type": "governed", "script_requirement": { "dataflow_spec": flow, "dataset_edges": { "source_table": flow["input_schema_refs"], "target_table": flow["output_schema_ref"], }, "migration_metadata": { "status": "migrated", "legacy_fields_present": False, "preserved_for_read_only": True, "governed_semantics": "dataflow_spec", }, }, "tag": [{"id": 81}], "draft_reservation": { key: receipt[key] for key in ("reservation_id", "dataflow_uid", "nonce") }, }, repository=repository, actor_uid=actor, ) assert repository.completed == {} assert failures[0]["error_code"] == "neo4j_tag_merge_failed" def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated( monkeypatch, ): good_uid = new_governance_uid() bad_uid = new_governance_uid() class Session: def __init__(self): self.commits = 0 self.rollbacks = 0 def commit(self): self.commits += 1 def rollback(self): self.rollbacks += 1 class Repository: def __init__(self): self.session = Session() self.claim_calls = 0 self.failures = [] self.completions = [] self.renewals = [] self.claims = [ { "reservation_id": new_governance_uid(), "dataflow_uid": bad_uid, "lease_token": new_governance_uid(), "request_digest": "a" * 64, "create_intent": {}, "attempt": 2, "integrity_valid": False, }, { "reservation_id": new_governance_uid(), "dataflow_uid": good_uid, "lease_token": new_governance_uid(), "request_digest": "b" * 64, "create_intent": {"dataflow_uid": good_uid}, "attempt": 2, "integrity_valid": True, }, ] def list_reconcilable_dataflow_creates(self, *, limit): assert limit == 2 return [{"reservation_id": "preview", "state": "failed"}] def claim_reconcilable_dataflow_creates( self, *, limit, lease_seconds ): assert limit == 1 assert lease_seconds == 60 self.claim_calls += 1 return [self.claims.pop(0)] if self.claims else [] def commit_dataflow_create_claim(self): self.session.commit() def renew_dataflow_create_lease(self, **kwargs): self.renewals.append(kwargs) def commit_dataflow_create_failure(self, **kwargs): self.failures.append(kwargs) self.session.commit() def complete_dataflow_create(self, **kwargs): self.completions.append(kwargs) return {"id": 73, "uid": kwargs["dataflow_uid"]} repository = Repository() reconciler = DataFlowCreateReconciler(repository) preview = reconciler.run(dry_run=True, limit=2) assert preview["candidate_count"] == 1 assert repository.claim_calls == 0 assert repository.session.commits == 0 monkeypatch.setattr( DataFlowService, "validate_governed_create_intent", lambda value, *, repository: { "node": {"uid": value["dataflow_uid"]}, "tags": [], }, ) monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda node: (73, {"id": 73, **node}), ) report = reconciler.run(dry_run=False, limit=2) assert report["reconciled_count"] == 1 assert report["failed_count"] == 1 assert repository.failures[0]["error_code"] == ( "create_intent_integrity_failed" ) assert repository.completions[0]["dataflow_uid"] == good_uid assert len(repository.renewals) == 1 def test_reconciler_validation_is_terminal_and_lost_lease_is_reported( monkeypatch, ): flow_uid = new_governance_uid() class Session: def commit(self): return None def rollback(self): return None class Repository: session = Session() def __init__(self): self.claimed = False self.failure_codes = [] def claim_reconcilable_dataflow_creates(self, **_kwargs): if self.claimed: return [] self.claimed = True return [ { "reservation_id": new_governance_uid(), "dataflow_uid": flow_uid, "lease_token": new_governance_uid(), "request_digest": "c" * 64, "create_intent": {"dataflow_uid": flow_uid}, "attempt": 1, "integrity_valid": True, } ] def commit_dataflow_create_claim(self): return None def commit_dataflow_create_failure(self, **kwargs): self.failure_codes.append(kwargs["error_code"]) repository = Repository() monkeypatch.setattr( DataFlowService, "validate_governed_create_intent", lambda *_args, **_kwargs: (_ for _ in ()).throw( ValueError("published asset is unavailable") ), ) report = DataFlowCreateReconciler(repository).run( dry_run=False, limit=2 ) assert report["items"][0] == { "reservation_id": report["items"][0]["reservation_id"], "dataflow_uid": flow_uid, "attempt": 1, "request_digest": "c" * 64, "status": "failed", "error_code": "create_intent_validation_failed", } assert repository.failure_codes == ["create_intent_validation_failed"] repository = Repository() monkeypatch.setattr( DataFlowService, "validate_governed_create_intent", lambda value, **_kwargs: {"node": value, "tags": []}, ) repository.renew_dataflow_create_lease = lambda **_kwargs: (_ for _ in ()).throw( RuntimeError("dataflow create lease was lost") ) repository.commit_dataflow_create_failure = lambda **_kwargs: ( _ for _ in () ).throw(RuntimeError("dataflow create lease was lost")) monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda _node: pytest.fail("side effect ran after lease loss"), ) report = DataFlowCreateReconciler(repository).run( dry_run=False, limit=1 ) assert report["items"][0]["status"] == "lease_recovery_required" assert report["items"][0]["error_code"] == ( "failure_state_persistence_failed" ) def test_tag_relationship_uses_one_atomic_merge(monkeypatch): calls = [] class Result: def single(self): return {"merged": 1} class Session: def run(self, query, **values): calls.append((query, values)) return Result() def __enter__(self): return self def __exit__(self, *_args): return None class Driver: def session(self): return Session() def close(self): return None monkeypatch.setattr( "app.core.data_flow.dataflows.connect_graph", lambda: Driver() ) DataFlowService._handle_single_tag_relationship(73, 81) DataFlowService._handle_single_tag_relationship(73, 81) assert len(calls) == 2 assert all("MERGE (a)-[:LABEL]->(b)" in query for query, _ in calls) assert all("CREATE (a)-[:LABEL]->(b)" not in query for query, _ in calls) def test_reconciler_tag_failure_does_not_complete_and_retry_succeeds( monkeypatch, ): flow_uid = new_governance_uid() reservation_id = new_governance_uid() digest = "d" * 64 class Session: def commit(self): return None def rollback(self): return None class Repository: session = Session() def __init__(self): self.attempt = 0 self.available = True self.failures = [] self.completions = [] def claim_reconcilable_dataflow_creates(self, **_kwargs): if not self.available: return [] self.available = False self.attempt += 1 return [ { "reservation_id": reservation_id, "dataflow_uid": flow_uid, "lease_token": new_governance_uid(), "request_digest": digest, "create_intent": {"dataflow_uid": flow_uid}, "attempt": self.attempt, "integrity_valid": True, } ] def commit_dataflow_create_claim(self): return None def renew_dataflow_create_lease(self, **_kwargs): return None def commit_dataflow_create_failure(self, **kwargs): self.failures.append(kwargs) self.available = True def complete_dataflow_create(self, **kwargs): self.completions.append(kwargs) return {"id": 73, "uid": flow_uid} repository = Repository() monkeypatch.setattr( DataFlowService, "validate_governed_create_intent", lambda _value, **_kwargs: { "node": {"uid": flow_uid}, "tags": [{"id": 81}], }, ) monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", lambda node: (73, {"id": 73, **node}), ) attempts = {"count": 0} def merge_tags(_node_id, _tags, *, strict): assert strict is True attempts["count"] += 1 if attempts["count"] == 1: raise RuntimeError("neo4j tag write failed") monkeypatch.setattr( DataFlowService, "_handle_tag_relationships", merge_tags ) reconciler = DataFlowCreateReconciler(repository) failed = reconciler.run(dry_run=False, limit=1) assert failed["items"][0]["status"] == "failed" assert repository.completions == [] assert repository.failures[0]["error_code"] == "reconciliation_failed" completed = reconciler.run(dry_run=False, limit=1) assert completed["items"][0]["status"] == "completed" assert len(repository.completions) == 1 assert attempts["count"] == 2 def test_slow_reconciler_does_not_preclaim_next_item_from_second_worker( monkeypatch, ): first_uid = new_governance_uid() second_uid = new_governance_uid() shared = { first_uid: {"state": "failed", "attempt": 0}, second_uid: {"state": "failed", "attempt": 0}, } processed = [] class Session: def commit(self): return None def rollback(self): return None class Repository: session = Session() def __init__(self, worker): self.worker = worker def claim_reconcilable_dataflow_creates( self, *, limit, lease_seconds ): assert limit == 1 assert lease_seconds == 10 available = [ uid for uid, record in shared.items() if record["state"] == "failed" ] if not available: return [] uid = available[0] shared[uid]["state"] = f"creating:{self.worker}" shared[uid]["attempt"] += 1 return [ { "reservation_id": new_governance_uid(), "dataflow_uid": uid, "lease_token": new_governance_uid(), "request_digest": ("a" if uid == first_uid else "b") * 64, "create_intent": {"dataflow_uid": uid}, "attempt": shared[uid]["attempt"], "integrity_valid": True, } ] def commit_dataflow_create_claim(self): return None def renew_dataflow_create_lease(self, **_kwargs): return None def complete_dataflow_create(self, *, dataflow_uid, **_kwargs): assert shared[dataflow_uid]["state"] == f"creating:{self.worker}" shared[dataflow_uid]["state"] = "completed" processed.append((self.worker, dataflow_uid)) return {"id": 73, "uid": dataflow_uid} def commit_dataflow_create_failure(self, **_kwargs): pytest.fail("unexpected reconciliation failure") monkeypatch.setattr( DataFlowService, "validate_governed_create_intent", lambda value, **_kwargs: { "node": {"uid": value["dataflow_uid"]}, "tags": [], }, ) worker_two = DataFlowCreateReconciler(Repository("worker-2")) entered_slow_call = {"value": False} def slow_merge(node): if node["uid"] == first_uid and not entered_slow_call["value"]: entered_slow_call["value"] = True # This models worker 1 crossing its original lease boundary while # processing item A. Item B was never preclaimed, so worker 2 can # safely process B without ever sharing ownership of A. second_report = worker_two.run( dry_run=False, limit=1, lease_seconds=10 ) assert second_report["items"][0]["dataflow_uid"] == second_uid return 73, {"id": 73, **node} monkeypatch.setattr( DataFlowService, "_merge_governed_dataflow", slow_merge ) first_report = DataFlowCreateReconciler(Repository("worker-1")).run( dry_run=False, limit=2, lease_seconds=10 ) assert first_report["reconciled_count"] == 1 assert processed == [ ("worker-2", second_uid), ("worker-1", first_uid), ] assert all(record["attempt"] == 1 for record in shared.values())