| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245 |
- 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,
- )
|