test_legacy_governance_cutover.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503
  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. self.consumed = set()
  15. def require_published_rule_version(self, rule_version_id):
  16. self.rule_calls.append(rule_version_id)
  17. if not self.published:
  18. raise ValueError("rule version is not published")
  19. return {"id": rule_version_id, "status": "published"}
  20. def load_published_assets(self, dataflow_spec):
  21. self.dataflow_calls.append(dataflow_spec)
  22. if not self.published:
  23. raise ValueError("dataflow references unpublished assets")
  24. return {}, {}
  25. def reserve_dataflow_draft(self, *, actor_uid):
  26. return {
  27. "reservation_id": new_governance_uid(),
  28. "dataflow_uid": new_governance_uid(),
  29. "nonce": "single-use-test-nonce",
  30. "expires_at": "2099-01-01T00:00:00+00:00",
  31. }
  32. def consume_dataflow_draft(self, receipt, *, actor_uid):
  33. key = (receipt["reservation_id"], actor_uid)
  34. if key in self.consumed:
  35. raise ValueError("already used")
  36. self.consumed.add(key)
  37. return receipt["dataflow_uid"]
  38. def _headers(app, role="editor"):
  39. token = issue_access_token(
  40. user_id=new_governance_uid(),
  41. roles=[role],
  42. secret=app.config["SECRET_KEY"],
  43. now=datetime.now(UTC),
  44. lifetime=timedelta(minutes=10),
  45. )
  46. return {"Authorization": f"Bearer {token}"}
  47. def _use_token_identity(monkeypatch):
  48. def load(token, *, secret):
  49. claims = decode_access_token(token, secret=secret)
  50. return {
  51. "id": claims["sub"],
  52. "username": "cutover-test",
  53. "display_name": "Cutover Test",
  54. "roles": claims["roles"],
  55. }
  56. monkeypatch.setattr("app.core.system.auth.load_identity_from_token", load)
  57. def test_production_line_draft_identity_is_server_owned_closed_and_governed(
  58. monkeypatch,
  59. ):
  60. from app import create_app
  61. app = create_app()
  62. app.config["TESTING"] = True
  63. _use_token_identity(monkeypatch)
  64. app.extensions["data_rule_repository"] = PublishedAssetRepository()
  65. client = app.test_client()
  66. response = client.post(
  67. "/api/rules/production-lines/draft-identity",
  68. json={},
  69. headers=_headers(app),
  70. )
  71. assert response.status_code == 201
  72. receipt = response.get_json()["data"]
  73. value = receipt["dataflow_uid"]
  74. parsed = uuid.UUID(value)
  75. assert parsed.version == 7
  76. assert parsed.variant == uuid.RFC_4122
  77. assert (
  78. client.post(
  79. "/api/rules/production-lines/draft-identity",
  80. json={"dataflow_uid": value},
  81. headers=_headers(app),
  82. ).status_code
  83. == 400
  84. )
  85. assert (
  86. client.post(
  87. "/api/rules/production-lines/draft-identity",
  88. json={},
  89. headers=_headers(app, "viewer"),
  90. ).status_code
  91. == 403
  92. )
  93. def test_legacy_standard_code_generation_is_closed_and_code_cannot_be_written(
  94. monkeypatch,
  95. ):
  96. from app import create_app
  97. app = create_app()
  98. app.config["TESTING"] = True
  99. _use_token_identity(monkeypatch)
  100. client = app.test_client()
  101. monkeypatch.setattr(
  102. "app.api.data_interface.routes.create_or_get_node",
  103. lambda *_args, **_kwargs: pytest.fail("legacy code was persisted"),
  104. )
  105. generated = client.post(
  106. "/api/interface/data/standard/code",
  107. json={"input": [], "describe": "生成代码", "output": []},
  108. headers=_headers(app),
  109. )
  110. assert generated.status_code == 410
  111. assert generated.get_json()["data"]["semantics"] == "read_only_migration"
  112. added = client.post(
  113. "/api/interface/data/standard/add",
  114. json={
  115. "name_zh": "旧代码标准",
  116. "tag": [],
  117. "code": "print('must not persist')",
  118. },
  119. headers=_headers(app),
  120. )
  121. assert added.status_code == 400
  122. updated = client.post(
  123. "/api/interface/data/standard/update",
  124. json={
  125. "name_zh": "旧代码标准",
  126. "tag": [],
  127. "code": "print('must not persist')",
  128. },
  129. headers=_headers(app),
  130. )
  131. assert updated.status_code == 400
  132. missing_rule = client.post(
  133. "/api/interface/data/standard/add",
  134. json={"name_zh": "缺少规则的标准", "tag": []},
  135. headers=_headers(app),
  136. )
  137. assert missing_rule.status_code == 400
  138. def test_governed_legacy_standard_link_is_attested_server_side(monkeypatch):
  139. from app import create_app
  140. app = create_app()
  141. app.config["TESTING"] = True
  142. _use_token_identity(monkeypatch)
  143. repository = PublishedAssetRepository()
  144. app.extensions["data_rule_repository"] = repository
  145. captured = {}
  146. monkeypatch.setattr(
  147. "app.api.data_interface.routes.translate_and_parse",
  148. lambda _value: ["published_standard"],
  149. )
  150. monkeypatch.setattr(
  151. "app.api.data_interface.routes.create_or_get_node",
  152. lambda _label, **properties: captured.update(properties) or 17,
  153. )
  154. client = app.test_client()
  155. rule_version_id = new_governance_uid()
  156. response = client.post(
  157. "/api/interface/data/standard/add",
  158. json={
  159. "name_zh": "已治理标准",
  160. "tag": [],
  161. "rule_version_id": rule_version_id,
  162. },
  163. headers=_headers(app),
  164. )
  165. assert response.status_code == 200
  166. assert repository.rule_calls == [rule_version_id]
  167. assert captured["rule_version_id"] == rule_version_id
  168. assert "rule_version_status" not in captured
  169. repository.published = False
  170. rejected = client.post(
  171. "/api/interface/data/standard/add",
  172. json={
  173. "name_zh": "伪造发布状态",
  174. "tag": [],
  175. "rule_version_id": new_governance_uid(),
  176. },
  177. headers=_headers(app),
  178. )
  179. assert rejected.status_code == 400
  180. def test_governed_dataflow_envelope_is_closed_and_uses_published_assets():
  181. from app.core.data_flow.dataflows import DataFlowService
  182. repository = PublishedAssetRepository()
  183. flow = valid_dataflow_spec()
  184. envelope = {
  185. "dataflow_spec": flow,
  186. "dataset_edges": {
  187. "source_table": list(flow["input_schema_refs"]),
  188. "target_table": flow["output_schema_ref"],
  189. },
  190. "migration_metadata": {
  191. "status": "migrated",
  192. "legacy_fields_present": False,
  193. "preserved_for_read_only": True,
  194. "governed_semantics": "dataflow_spec",
  195. },
  196. }
  197. normalized = DataFlowService.validate_governed_requirement(
  198. envelope, repository=repository
  199. )
  200. assert normalized["dataflow_spec"]["dataflow_uid"] == flow["dataflow_uid"]
  201. assert repository.dataflow_calls == [normalized["dataflow_spec"]]
  202. invalid = dict(envelope)
  203. invalid["task_list"] = []
  204. with pytest.raises(ValueError, match="unsupported fields"):
  205. DataFlowService.validate_governed_requirement(
  206. invalid, repository=repository
  207. )
  208. mismatched = json.loads(json.dumps(envelope))
  209. mismatched["dataset_edges"]["target_table"] = "bd:other:v1"
  210. with pytest.raises(ValueError, match="dataset edges"):
  211. DataFlowService.validate_governed_requirement(
  212. mismatched, repository=repository
  213. )
  214. def test_malformed_governed_signals_never_fall_through_to_legacy_side_effects(
  215. monkeypatch,
  216. ):
  217. from app.core.data_flow.dataflows import DataFlowService
  218. repository = PublishedAssetRepository()
  219. invoked = []
  220. monkeypatch.setattr(
  221. DataFlowService,
  222. "_save_to_pg_database",
  223. lambda *_args, **_kwargs: invoked.append("task"),
  224. )
  225. monkeypatch.setattr(
  226. DataFlowService,
  227. "_handle_script_relationships",
  228. lambda *_args, **_kwargs: invoked.append("script"),
  229. )
  230. with pytest.raises(ValueError, match="unsupported fields"):
  231. DataFlowService.create_dataflow(
  232. {
  233. "name_zh": "畸形治理流",
  234. "describe": "不能降级",
  235. "script_requirement": {
  236. "migration_metadata": {
  237. "status": "migrated",
  238. }
  239. },
  240. },
  241. repository=repository,
  242. actor_uid=new_governance_uid(),
  243. )
  244. assert invoked == []
  245. def test_governed_dataflow_creation_never_generates_legacy_task_or_workflow(
  246. monkeypatch,
  247. ):
  248. from app.core.data_flow.dataflows import DataFlowService
  249. repository = PublishedAssetRepository()
  250. flow = valid_dataflow_spec()
  251. data = {
  252. "name_zh": "客户治理生产线",
  253. "describe": "固定发布版本的数据生产线",
  254. "script_type": "python",
  255. "script_requirement": {
  256. "dataflow_spec": flow,
  257. "dataset_edges": {
  258. "source_table": list(flow["input_schema_refs"]),
  259. "target_table": flow["output_schema_ref"],
  260. },
  261. "migration_metadata": {
  262. "status": "migrated",
  263. "legacy_fields_present": False,
  264. "preserved_for_read_only": True,
  265. "governed_semantics": "dataflow_spec",
  266. },
  267. },
  268. }
  269. actor_uid = new_governance_uid()
  270. receipt = repository.reserve_dataflow_draft(actor_uid=actor_uid)
  271. receipt["dataflow_uid"] = flow["dataflow_uid"]
  272. data["script_type"] = "governed"
  273. data["draft_reservation"] = {
  274. key: receipt[key]
  275. for key in ("reservation_id", "dataflow_uid", "nonce")
  276. }
  277. created = {}
  278. class Result:
  279. def single(self):
  280. return None
  281. class Session:
  282. def run(self, *_args, **_kwargs):
  283. return Result()
  284. def __enter__(self):
  285. return self
  286. def __exit__(self, *_args):
  287. return None
  288. class Driver:
  289. def session(self):
  290. return Session()
  291. monkeypatch.setattr(
  292. "app.core.data_flow.dataflows.translate_and_parse",
  293. lambda _name: ["customer_governed_line"],
  294. )
  295. monkeypatch.setattr(
  296. "app.core.data_flow.dataflows.get_node", lambda *_args, **_kwargs: None
  297. )
  298. monkeypatch.setattr(
  299. "app.core.data_flow.dataflows.create_or_get_node",
  300. lambda _label, **properties: created.update(properties) or 31,
  301. )
  302. monkeypatch.setattr(
  303. "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
  304. )
  305. monkeypatch.setattr(
  306. DataFlowService,
  307. "_save_to_pg_database",
  308. lambda *_args, **_kwargs: pytest.fail("task_list write was invoked"),
  309. )
  310. monkeypatch.setattr(
  311. DataFlowService,
  312. "_handle_script_relationships",
  313. lambda *_args, **_kwargs: pytest.fail("legacy script path was invoked"),
  314. )
  315. monkeypatch.setattr(
  316. DataFlowService,
  317. "_register_data_product",
  318. lambda *_args, **_kwargs: None,
  319. )
  320. result = DataFlowService.create_dataflow(
  321. data, repository=repository, actor_uid=actor_uid
  322. )
  323. assert result["id"] == 31
  324. assert created["uid"] == flow["dataflow_uid"]
  325. assert json.loads(created["script_requirement"])["dataflow_spec"] == flow
  326. def test_governed_dataflow_update_preserves_identity_and_closed_envelope(
  327. monkeypatch,
  328. ):
  329. from app.core.data_flow.dataflows import DataFlowService
  330. repository = PublishedAssetRepository()
  331. flow = valid_dataflow_spec()
  332. envelope = {
  333. "dataflow_spec": flow,
  334. "dataset_edges": {
  335. "source_table": list(flow["input_schema_refs"]),
  336. "target_table": flow["output_schema_ref"],
  337. },
  338. "migration_metadata": {
  339. "status": "migrated",
  340. "legacy_fields_present": False,
  341. "preserved_for_read_only": True,
  342. "governed_semantics": "dataflow_spec",
  343. },
  344. }
  345. updated = {}
  346. class Result:
  347. def __init__(self, data=None, single=None):
  348. self._data = data or []
  349. self._single = single
  350. def data(self):
  351. return self._data
  352. def single(self):
  353. return self._single
  354. class Session:
  355. def run(self, query, params=None, **kwargs):
  356. values = params or kwargs
  357. if "RETURN n" in query and "SET " not in query:
  358. return Result(
  359. data=[
  360. {
  361. "n": {
  362. "uid": flow["dataflow_uid"],
  363. "script_type": "governed",
  364. "script_requirement": json.dumps(envelope),
  365. }
  366. }
  367. ]
  368. )
  369. if "SET " in query:
  370. updated.update(values)
  371. return Result(
  372. data=[
  373. {
  374. "n": {
  375. "uid": values["uid"],
  376. "script_type": values["script_type"],
  377. "script_requirement": values[
  378. "script_requirement"
  379. ],
  380. },
  381. "node_id": 44,
  382. }
  383. ]
  384. )
  385. return Result(single={"tags": []})
  386. def __enter__(self):
  387. return self
  388. def __exit__(self, *_args):
  389. return None
  390. class Driver:
  391. def session(self):
  392. return Session()
  393. monkeypatch.setattr(
  394. "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
  395. )
  396. result = DataFlowService.update_dataflow(
  397. 44,
  398. {
  399. "script_type": "governed",
  400. "script_path": "",
  401. "script_requirement": envelope,
  402. },
  403. repository=repository,
  404. )
  405. assert result["id"] == 44
  406. assert updated["uid"] == flow["dataflow_uid"]
  407. assert updated["script_type"] == "governed"
  408. assert updated["script_path"] == ""
  409. assert json.loads(updated["script_requirement"]) == envelope
  410. with pytest.raises(ValueError, match="script_type must be governed"):
  411. DataFlowService.update_dataflow(
  412. 44,
  413. {
  414. "script_type": "python",
  415. "script_requirement": envelope,
  416. },
  417. repository=repository,
  418. )
  419. with pytest.raises(ValueError, match="unsupported fields"):
  420. DataFlowService.update_dataflow(
  421. 44,
  422. {
  423. "script_type": "governed",
  424. "script_requirement": {"rule": "downgrade"},
  425. },
  426. repository=repository,
  427. )
  428. mismatched = json.loads(json.dumps(envelope))
  429. mismatched["dataflow_spec"]["dataflow_uid"] = new_governance_uid()
  430. with pytest.raises(ValueError, match="cannot replace"):
  431. DataFlowService.update_dataflow(
  432. 44,
  433. {"script_requirement": mismatched},
  434. repository=repository,
  435. )