test_legacy_governance_cutover.py 16 KB

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