test_legacy_governance_cutover.py 15 KB

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