test_legacy_governance_cutover.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408
  1. from __future__ import annotations
  2. import json
  3. import uuid
  4. from datetime import UTC, datetime, timedelta
  5. import pytest
  6. from app.core.common.identifiers import new_governance_uid
  7. from app.core.system.tokens import decode_access_token, issue_access_token
  8. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  9. class PublishedAssetRepository:
  10. def __init__(self, *, published=True):
  11. self.published = published
  12. self.rule_calls = []
  13. self.dataflow_calls = []
  14. def require_published_rule_version(self, rule_version_id):
  15. self.rule_calls.append(rule_version_id)
  16. if not self.published:
  17. raise ValueError("rule version is not published")
  18. return {"id": rule_version_id, "status": "published"}
  19. def load_published_assets(self, dataflow_spec):
  20. self.dataflow_calls.append(dataflow_spec)
  21. if not self.published:
  22. raise ValueError("dataflow references unpublished assets")
  23. return {}, {}
  24. def _headers(app, role="editor"):
  25. token = issue_access_token(
  26. user_id=new_governance_uid(),
  27. roles=[role],
  28. secret=app.config["SECRET_KEY"],
  29. now=datetime.now(UTC),
  30. lifetime=timedelta(minutes=10),
  31. )
  32. return {"Authorization": f"Bearer {token}"}
  33. def _use_token_identity(monkeypatch):
  34. def load(token, *, secret):
  35. claims = decode_access_token(token, secret=secret)
  36. return {
  37. "id": claims["sub"],
  38. "username": "cutover-test",
  39. "display_name": "Cutover Test",
  40. "roles": claims["roles"],
  41. }
  42. monkeypatch.setattr("app.core.system.auth.load_identity_from_token", load)
  43. def test_production_line_draft_identity_is_server_owned_closed_and_governed(
  44. monkeypatch,
  45. ):
  46. from app import create_app
  47. app = create_app()
  48. app.config["TESTING"] = True
  49. _use_token_identity(monkeypatch)
  50. client = app.test_client()
  51. response = client.post(
  52. "/api/rules/production-lines/draft-identity",
  53. json={},
  54. headers=_headers(app),
  55. )
  56. assert response.status_code == 201
  57. value = response.get_json()["data"]["dataflow_uid"]
  58. parsed = uuid.UUID(value)
  59. assert parsed.version == 7
  60. assert parsed.variant == uuid.RFC_4122
  61. assert (
  62. client.post(
  63. "/api/rules/production-lines/draft-identity",
  64. json={"dataflow_uid": value},
  65. headers=_headers(app),
  66. ).status_code
  67. == 400
  68. )
  69. assert (
  70. client.post(
  71. "/api/rules/production-lines/draft-identity",
  72. json={},
  73. headers=_headers(app, "viewer"),
  74. ).status_code
  75. == 403
  76. )
  77. def test_legacy_standard_code_generation_is_closed_and_code_cannot_be_written(
  78. monkeypatch,
  79. ):
  80. from app import create_app
  81. app = create_app()
  82. app.config["TESTING"] = True
  83. _use_token_identity(monkeypatch)
  84. client = app.test_client()
  85. monkeypatch.setattr(
  86. "app.api.data_interface.routes.create_or_get_node",
  87. lambda *_args, **_kwargs: pytest.fail("legacy code was persisted"),
  88. )
  89. generated = client.post(
  90. "/api/interface/data/standard/code",
  91. json={"input": [], "describe": "生成代码", "output": []},
  92. headers=_headers(app),
  93. )
  94. assert generated.status_code == 410
  95. assert generated.get_json()["data"]["semantics"] == "read_only_migration"
  96. added = client.post(
  97. "/api/interface/data/standard/add",
  98. json={
  99. "name_zh": "旧代码标准",
  100. "tag": [],
  101. "code": "print('must not persist')",
  102. },
  103. headers=_headers(app),
  104. )
  105. assert added.status_code == 400
  106. updated = client.post(
  107. "/api/interface/data/standard/update",
  108. json={
  109. "name_zh": "旧代码标准",
  110. "tag": [],
  111. "code": "print('must not persist')",
  112. },
  113. headers=_headers(app),
  114. )
  115. assert updated.status_code == 400
  116. def test_governed_legacy_standard_link_is_attested_server_side(monkeypatch):
  117. from app import create_app
  118. app = create_app()
  119. app.config["TESTING"] = True
  120. _use_token_identity(monkeypatch)
  121. repository = PublishedAssetRepository()
  122. app.extensions["data_rule_repository"] = repository
  123. captured = {}
  124. monkeypatch.setattr(
  125. "app.api.data_interface.routes.translate_and_parse",
  126. lambda _value: ["published_standard"],
  127. )
  128. monkeypatch.setattr(
  129. "app.api.data_interface.routes.create_or_get_node",
  130. lambda _label, **properties: captured.update(properties) or 17,
  131. )
  132. client = app.test_client()
  133. rule_version_id = new_governance_uid()
  134. response = client.post(
  135. "/api/interface/data/standard/add",
  136. json={
  137. "name_zh": "已治理标准",
  138. "tag": [],
  139. "rule_version_id": rule_version_id,
  140. "rule_version_status": "draft",
  141. },
  142. headers=_headers(app),
  143. )
  144. assert response.status_code == 200
  145. assert repository.rule_calls == [rule_version_id]
  146. assert captured["rule_version_id"] == rule_version_id
  147. assert "rule_version_status" not in captured
  148. repository.published = False
  149. rejected = client.post(
  150. "/api/interface/data/standard/add",
  151. json={
  152. "name_zh": "伪造发布状态",
  153. "tag": [],
  154. "rule_version_id": new_governance_uid(),
  155. "rule_version_status": "published",
  156. },
  157. headers=_headers(app),
  158. )
  159. assert rejected.status_code == 400
  160. def test_governed_dataflow_envelope_is_closed_and_uses_published_assets():
  161. from app.core.data_flow.dataflows import DataFlowService
  162. repository = PublishedAssetRepository()
  163. flow = valid_dataflow_spec()
  164. envelope = {
  165. "dataflow_spec": flow,
  166. "dataset_edges": {
  167. "source_table": list(flow["input_schema_refs"]),
  168. "target_table": flow["output_schema_ref"],
  169. },
  170. "migration_metadata": {
  171. "status": "migrated",
  172. "legacy_fields_present": False,
  173. "preserved_for_read_only": True,
  174. "governed_semantics": "dataflow_spec",
  175. },
  176. }
  177. normalized = DataFlowService.validate_governed_requirement(
  178. envelope, repository=repository
  179. )
  180. assert normalized["dataflow_spec"]["dataflow_uid"] == flow["dataflow_uid"]
  181. assert repository.dataflow_calls == [normalized["dataflow_spec"]]
  182. invalid = dict(envelope)
  183. invalid["task_list"] = []
  184. with pytest.raises(ValueError, match="unsupported fields"):
  185. DataFlowService.validate_governed_requirement(
  186. invalid, repository=repository
  187. )
  188. mismatched = json.loads(json.dumps(envelope))
  189. mismatched["dataset_edges"]["target_table"] = "bd:other:v1"
  190. with pytest.raises(ValueError, match="dataset edges"):
  191. DataFlowService.validate_governed_requirement(
  192. mismatched, repository=repository
  193. )
  194. def test_governed_dataflow_creation_never_generates_legacy_task_or_workflow(
  195. monkeypatch,
  196. ):
  197. from app.core.data_flow.dataflows import DataFlowService
  198. repository = PublishedAssetRepository()
  199. flow = valid_dataflow_spec()
  200. data = {
  201. "name_zh": "客户治理生产线",
  202. "describe": "固定发布版本的数据生产线",
  203. "script_type": "python",
  204. "script_requirement": {
  205. "dataflow_spec": flow,
  206. "dataset_edges": {
  207. "source_table": list(flow["input_schema_refs"]),
  208. "target_table": flow["output_schema_ref"],
  209. },
  210. "migration_metadata": {
  211. "status": "migrated",
  212. "legacy_fields_present": False,
  213. "preserved_for_read_only": True,
  214. "governed_semantics": "dataflow_spec",
  215. },
  216. },
  217. }
  218. created = {}
  219. class Result:
  220. def single(self):
  221. return None
  222. class Session:
  223. def run(self, *_args, **_kwargs):
  224. return Result()
  225. def __enter__(self):
  226. return self
  227. def __exit__(self, *_args):
  228. return None
  229. class Driver:
  230. def session(self):
  231. return Session()
  232. monkeypatch.setattr(
  233. "app.core.data_flow.dataflows.translate_and_parse",
  234. lambda _name: ["customer_governed_line"],
  235. )
  236. monkeypatch.setattr(
  237. "app.core.data_flow.dataflows.get_node", lambda *_args, **_kwargs: None
  238. )
  239. monkeypatch.setattr(
  240. "app.core.data_flow.dataflows.create_or_get_node",
  241. lambda _label, **properties: created.update(properties) or 31,
  242. )
  243. monkeypatch.setattr(
  244. "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
  245. )
  246. monkeypatch.setattr(
  247. DataFlowService,
  248. "_save_to_pg_database",
  249. lambda *_args, **_kwargs: pytest.fail("task_list write was invoked"),
  250. )
  251. monkeypatch.setattr(
  252. DataFlowService,
  253. "_handle_script_relationships",
  254. lambda *_args, **_kwargs: pytest.fail("legacy script path was invoked"),
  255. )
  256. monkeypatch.setattr(
  257. DataFlowService,
  258. "_register_data_product",
  259. lambda *_args, **_kwargs: None,
  260. )
  261. result = DataFlowService.create_dataflow(data, repository=repository)
  262. assert result["id"] == 31
  263. assert created["uid"] == flow["dataflow_uid"]
  264. assert json.loads(created["script_requirement"])["dataflow_spec"] == flow
  265. def test_governed_dataflow_update_preserves_identity_and_closed_envelope(
  266. monkeypatch,
  267. ):
  268. from app.core.data_flow.dataflows import DataFlowService
  269. repository = PublishedAssetRepository()
  270. flow = valid_dataflow_spec()
  271. envelope = {
  272. "dataflow_spec": flow,
  273. "dataset_edges": {
  274. "source_table": list(flow["input_schema_refs"]),
  275. "target_table": flow["output_schema_ref"],
  276. },
  277. "migration_metadata": {
  278. "status": "migrated",
  279. "legacy_fields_present": False,
  280. "preserved_for_read_only": True,
  281. "governed_semantics": "dataflow_spec",
  282. },
  283. }
  284. updated = {}
  285. class Result:
  286. def __init__(self, data=None, single=None):
  287. self._data = data or []
  288. self._single = single
  289. def data(self):
  290. return self._data
  291. def single(self):
  292. return self._single
  293. class Session:
  294. def run(self, query, params=None, **kwargs):
  295. values = params or kwargs
  296. if "RETURN n" in query and "SET " not in query:
  297. return Result(data=[{"n": {"uid": flow["dataflow_uid"]}}])
  298. if "SET " in query:
  299. updated.update(values)
  300. return Result(
  301. data=[
  302. {
  303. "n": {
  304. "uid": values["uid"],
  305. "script_type": values["script_type"],
  306. "script_requirement": values[
  307. "script_requirement"
  308. ],
  309. },
  310. "node_id": 44,
  311. }
  312. ]
  313. )
  314. return Result(single={"tags": []})
  315. def __enter__(self):
  316. return self
  317. def __exit__(self, *_args):
  318. return None
  319. class Driver:
  320. def session(self):
  321. return Session()
  322. monkeypatch.setattr(
  323. "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
  324. )
  325. result = DataFlowService.update_dataflow(
  326. 44,
  327. {
  328. "script_type": "python",
  329. "script_path": "/tmp/forged.py",
  330. "script_requirement": envelope,
  331. },
  332. repository=repository,
  333. )
  334. assert result["id"] == 44
  335. assert updated["uid"] == flow["dataflow_uid"]
  336. assert updated["script_type"] == "governed"
  337. assert updated["script_path"] == ""
  338. assert json.loads(updated["script_requirement"]) == envelope
  339. mismatched = json.loads(json.dumps(envelope))
  340. mismatched["dataflow_spec"]["dataflow_uid"] = new_governance_uid()
  341. with pytest.raises(ValueError, match="cannot replace"):
  342. DataFlowService.update_dataflow(
  343. 44,
  344. {"script_requirement": mismatched},
  345. repository=repository,
  346. )