| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524 |
- from __future__ import annotations
- import hashlib
- import json
- import uuid
- from datetime import UTC, datetime, timedelta
- import pytest
- from app.core.common.identifiers import new_governance_uid
- from app.core.system.tokens import decode_access_token, issue_access_token
- from tests.core.data_rules.test_contracts import valid_dataflow_spec
- class PublishedAssetRepository:
- def __init__(self, *, published=True):
- self.published = published
- self.rule_calls = []
- self.dataflow_calls = []
- self.consumed = set()
- self.completed = {}
- def require_published_rule_version(self, rule_version_id):
- self.rule_calls.append(rule_version_id)
- if not self.published:
- raise ValueError("rule version is not published")
- return {"id": rule_version_id, "status": "published"}
- def load_published_assets(self, dataflow_spec):
- self.dataflow_calls.append(dataflow_spec)
- if not self.published:
- raise ValueError("dataflow references unpublished assets")
- return {}, {}
- def reserve_dataflow_draft(self, *, actor_uid):
- return {
- "reservation_id": new_governance_uid(),
- "dataflow_uid": new_governance_uid(),
- "nonce": "single-use-test-nonce",
- "expires_at": "2099-01-01T00:00:00+00:00",
- }
- def begin_dataflow_create(
- self, receipt, *, actor_uid, create_request, create_intent
- ):
- key = (receipt["reservation_id"], actor_uid)
- request_digest = hashlib.sha256(
- json.dumps(
- create_request,
- sort_keys=True,
- separators=(",", ":"),
- ensure_ascii=False,
- ).encode("utf-8")
- ).hexdigest()
- if key in self.completed:
- return {
- "status": "completed",
- "dataflow_uid": receipt["dataflow_uid"],
- "result": self.completed[key],
- "request_digest": request_digest,
- }
- self.consumed.add(key)
- return {
- "status": "claimed",
- "dataflow_uid": receipt["dataflow_uid"],
- "lease_token": new_governance_uid(),
- "attempt": 1,
- "create_intent": create_intent,
- "request_digest": request_digest,
- }
- def commit_dataflow_create_claim(self):
- return None
- def commit_dataflow_create_failure(self, **_kwargs):
- return None
- def complete_dataflow_create(self, *, reservation_id, result, **_kwargs):
- key = next(key for key in self.consumed if key[0] == reservation_id)
- self.completed[key] = dict(result)
- return dict(result)
- def _headers(app, role="editor"):
- token = issue_access_token(
- user_id=new_governance_uid(),
- roles=[role],
- secret=app.config["SECRET_KEY"],
- now=datetime.now(UTC),
- lifetime=timedelta(minutes=10),
- )
- return {"Authorization": f"Bearer {token}"}
- def _use_token_identity(monkeypatch):
- def load(token, *, secret):
- claims = decode_access_token(token, secret=secret)
- return {
- "id": claims["sub"],
- "username": "cutover-test",
- "display_name": "Cutover Test",
- "roles": claims["roles"],
- }
- monkeypatch.setattr("app.core.system.auth.load_identity_from_token", load)
- def test_production_line_draft_identity_is_server_owned_closed_and_governed(
- monkeypatch,
- ):
- from app import create_app
- app = create_app()
- app.config["TESTING"] = True
- _use_token_identity(monkeypatch)
- app.extensions["data_rule_repository"] = PublishedAssetRepository()
- client = app.test_client()
- response = client.post(
- "/api/rules/production-lines/draft-identity",
- json={},
- headers=_headers(app),
- )
- assert response.status_code == 201
- receipt = response.get_json()["data"]
- value = receipt["dataflow_uid"]
- parsed = uuid.UUID(value)
- assert parsed.version == 7
- assert parsed.variant == uuid.RFC_4122
- assert (
- client.post(
- "/api/rules/production-lines/draft-identity",
- json={"dataflow_uid": value},
- headers=_headers(app),
- ).status_code
- == 400
- )
- assert (
- client.post(
- "/api/rules/production-lines/draft-identity",
- json={},
- headers=_headers(app, "viewer"),
- ).status_code
- == 403
- )
- def test_legacy_standard_code_generation_is_closed_and_code_cannot_be_written(
- monkeypatch,
- ):
- from app import create_app
- app = create_app()
- app.config["TESTING"] = True
- _use_token_identity(monkeypatch)
- client = app.test_client()
- monkeypatch.setattr(
- "app.api.data_interface.routes.create_or_get_node",
- lambda *_args, **_kwargs: pytest.fail("legacy code was persisted"),
- )
- generated = client.post(
- "/api/interface/data/standard/code",
- json={"input": [], "describe": "生成代码", "output": []},
- headers=_headers(app),
- )
- assert generated.status_code == 410
- assert generated.get_json()["data"]["semantics"] == "read_only_migration"
- added = client.post(
- "/api/interface/data/standard/add",
- json={
- "name_zh": "旧代码标准",
- "tag": [],
- "code": "print('must not persist')",
- },
- headers=_headers(app),
- )
- assert added.status_code == 400
- updated = client.post(
- "/api/interface/data/standard/update",
- json={
- "name_zh": "旧代码标准",
- "tag": [],
- "code": "print('must not persist')",
- },
- headers=_headers(app),
- )
- assert updated.status_code == 400
- missing_rule = client.post(
- "/api/interface/data/standard/add",
- json={"name_zh": "缺少规则的标准", "tag": []},
- headers=_headers(app),
- )
- assert missing_rule.status_code == 400
- def test_governed_legacy_standard_link_is_attested_server_side(monkeypatch):
- from app import create_app
- app = create_app()
- app.config["TESTING"] = True
- _use_token_identity(monkeypatch)
- repository = PublishedAssetRepository()
- app.extensions["data_rule_repository"] = repository
- captured = {}
- monkeypatch.setattr(
- "app.api.data_interface.routes.translate_and_parse",
- lambda _value: ["published_standard"],
- )
- monkeypatch.setattr(
- "app.api.data_interface.routes.create_or_get_node",
- lambda _label, **properties: captured.update(properties) or 17,
- )
- client = app.test_client()
- rule_version_id = new_governance_uid()
- response = client.post(
- "/api/interface/data/standard/add",
- json={
- "name_zh": "已治理标准",
- "tag": [],
- "rule_version_id": rule_version_id,
- },
- headers=_headers(app),
- )
- assert response.status_code == 200
- assert repository.rule_calls == [rule_version_id]
- assert captured["rule_version_id"] == rule_version_id
- assert "rule_version_status" not in captured
- repository.published = False
- rejected = client.post(
- "/api/interface/data/standard/add",
- json={
- "name_zh": "伪造发布状态",
- "tag": [],
- "rule_version_id": new_governance_uid(),
- },
- headers=_headers(app),
- )
- assert rejected.status_code == 400
- def test_governed_dataflow_envelope_is_closed_and_uses_published_assets():
- from app.core.data_flow.dataflows import DataFlowService
- repository = PublishedAssetRepository()
- flow = valid_dataflow_spec()
- envelope = {
- "dataflow_spec": flow,
- "dataset_edges": {
- "source_table": list(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",
- },
- }
- normalized = DataFlowService.validate_governed_requirement(
- envelope, repository=repository
- )
- assert normalized["dataflow_spec"]["dataflow_uid"] == flow["dataflow_uid"]
- assert repository.dataflow_calls == [normalized["dataflow_spec"]]
- invalid = dict(envelope)
- invalid["task_list"] = []
- with pytest.raises(ValueError, match="unsupported fields"):
- DataFlowService.validate_governed_requirement(
- invalid, repository=repository
- )
- mismatched = json.loads(json.dumps(envelope))
- mismatched["dataset_edges"]["target_table"] = "bd:other:v1"
- with pytest.raises(ValueError, match="dataset edges"):
- DataFlowService.validate_governed_requirement(
- mismatched, repository=repository
- )
- def test_malformed_governed_signals_never_fall_through_to_legacy_side_effects(
- monkeypatch,
- ):
- from app.core.data_flow.dataflows import DataFlowService
- repository = PublishedAssetRepository()
- invoked = []
- monkeypatch.setattr(
- DataFlowService,
- "_save_to_pg_database",
- lambda *_args, **_kwargs: invoked.append("task"),
- )
- monkeypatch.setattr(
- DataFlowService,
- "_handle_script_relationships",
- lambda *_args, **_kwargs: invoked.append("script"),
- )
- with pytest.raises(ValueError, match="unsupported fields"):
- DataFlowService.create_dataflow(
- {
- "name_zh": "畸形治理流",
- "describe": "不能降级",
- "script_requirement": {
- "migration_metadata": {
- "status": "migrated",
- }
- },
- },
- repository=repository,
- actor_uid=new_governance_uid(),
- )
- assert invoked == []
- def test_governed_dataflow_creation_never_generates_legacy_task_or_workflow(
- monkeypatch,
- ):
- from app.core.data_flow.dataflows import DataFlowService
- repository = PublishedAssetRepository()
- flow = valid_dataflow_spec()
- data = {
- "name_zh": "客户治理生产线",
- "describe": "固定发布版本的数据生产线",
- "script_type": "python",
- "script_requirement": {
- "dataflow_spec": flow,
- "dataset_edges": {
- "source_table": list(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",
- },
- },
- }
- actor_uid = new_governance_uid()
- receipt = repository.reserve_dataflow_draft(actor_uid=actor_uid)
- receipt["dataflow_uid"] = flow["dataflow_uid"]
- data["script_type"] = "governed"
- data["draft_reservation"] = {
- key: receipt[key]
- for key in ("reservation_id", "dataflow_uid", "nonce")
- }
- monkeypatch.setattr(
- "app.core.data_flow.dataflows.translate_and_parse",
- lambda _name: pytest.fail(
- "governed create used non-deterministic translation"
- ),
- )
- created = {}
- monkeypatch.setattr(
- DataFlowService,
- "_merge_governed_dataflow",
- lambda properties: (
- 31,
- {**created, **properties, "id": 31}
- if not created.update(properties)
- else {},
- ),
- )
- monkeypatch.setattr(
- DataFlowService,
- "_save_to_pg_database",
- lambda *_args, **_kwargs: pytest.fail("task_list write was invoked"),
- )
- monkeypatch.setattr(
- DataFlowService,
- "_handle_script_relationships",
- lambda *_args, **_kwargs: pytest.fail("legacy script path was invoked"),
- )
- monkeypatch.setattr(
- DataFlowService,
- "_register_data_product",
- lambda *_args, **_kwargs: None,
- )
- result = DataFlowService.create_dataflow(
- data, repository=repository, actor_uid=actor_uid
- )
- assert result["id"] == 31
- assert created["uid"] == flow["dataflow_uid"]
- assert created["name_en"] == (
- f"dataflow_{flow['dataflow_uid'].replace('-', '')}"
- )
- assert json.loads(created["script_requirement"])["dataflow_spec"] == flow
- def test_governed_dataflow_update_preserves_identity_and_closed_envelope(
- monkeypatch,
- ):
- from app.core.data_flow.dataflows import DataFlowService
- repository = PublishedAssetRepository()
- flow = valid_dataflow_spec()
- envelope = {
- "dataflow_spec": flow,
- "dataset_edges": {
- "source_table": list(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",
- },
- }
- updated = {}
- class Result:
- def __init__(self, data=None, single=None):
- self._data = data or []
- self._single = single
- def data(self):
- return self._data
- def single(self):
- return self._single
- class Session:
- def run(self, query, params=None, **kwargs):
- values = params or kwargs
- if "RETURN n" in query and "SET " not in query:
- return Result(
- data=[
- {
- "n": {
- "uid": flow["dataflow_uid"],
- "script_type": "governed",
- "script_requirement": json.dumps(envelope),
- }
- }
- ]
- )
- if "SET " in query:
- updated.update(values)
- return Result(
- data=[
- {
- "n": {
- "uid": values["uid"],
- "script_type": values["script_type"],
- "script_requirement": values[
- "script_requirement"
- ],
- },
- "node_id": 44,
- }
- ]
- )
- return Result(single={"tags": []})
- def __enter__(self):
- return self
- def __exit__(self, *_args):
- return None
- class Driver:
- def session(self):
- return Session()
- monkeypatch.setattr(
- "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
- )
- result = DataFlowService.update_dataflow(
- 44,
- {
- "script_type": "governed",
- "script_path": "",
- "script_requirement": envelope,
- },
- repository=repository,
- )
- assert result["id"] == 44
- assert updated["uid"] == flow["dataflow_uid"]
- assert updated["script_type"] == "governed"
- assert updated["script_path"] == ""
- assert json.loads(updated["script_requirement"]) == envelope
- with pytest.raises(ValueError, match="script_type must be governed"):
- DataFlowService.update_dataflow(
- 44,
- {
- "script_type": "python",
- "script_requirement": envelope,
- },
- repository=repository,
- )
- with pytest.raises(ValueError, match="unsupported fields"):
- DataFlowService.update_dataflow(
- 44,
- {
- "script_type": "governed",
- "script_requirement": {"rule": "downgrade"},
- },
- repository=repository,
- )
- mismatched = json.loads(json.dumps(envelope))
- mismatched["dataflow_spec"]["dataflow_uid"] = new_governance_uid()
- with pytest.raises(ValueError, match="cannot replace"):
- DataFlowService.update_dataflow(
- 44,
- {"script_requirement": mismatched},
- repository=repository,
- )
|