| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766 |
- 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
- def consume(self):
- return None
- class GraphSession:
- def __init__(
- self, *, existing=None, conflict=False, constraint_error=None
- ):
- self.existing = existing
- self.conflict = conflict
- self.constraint_error = constraint_error
- self.calls = []
- def run(self, query, parameters=None, **kwargs):
- values = parameters or kwargs
- self.calls.append((query, values))
- if query.startswith("CREATE CONSTRAINT"):
- if (
- self.constraint_error is not None
- and "data_flow_name_zh" in query
- ):
- raise self.constraint_error
- 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 any(
- call[0]
- == "CREATE CONSTRAINT data_flow_name_zh IF NOT EXISTS "
- "FOR (n:DataFlow) REQUIRE n.name_zh 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_governed_graph_fails_closed_when_name_constraint_finds_duplicates(
- monkeypatch,
- ):
- class ConstraintFailure(RuntimeError):
- code = "Neo.ClientError.Schema.ConstraintValidationFailed"
- graph = GraphSession(constraint_error=ConstraintFailure("duplicates"))
- monkeypatch.setattr(
- "app.core.data_flow.dataflows.connect_graph",
- lambda: GraphDriver(graph),
- )
- with pytest.raises(ValueError, match="dataflow_uid_conflict"):
- DataFlowService._merge_governed_dataflow(governed_node())
- assert not any(call[0].startswith("MERGE") for call in graph.calls)
- 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())
|