from __future__ import annotations import copy import pytest from app.core.common.identifiers import new_governance_uid 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): assert actor_uid return { "status": "completed", "dataflow_uid": expected["uid"], "result": expected, } 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, )