test_data_rule_api.py 42 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333
  1. from __future__ import annotations
  2. from datetime import UTC, datetime, timedelta
  3. import pytest
  4. from app.core.common.identifiers import new_governance_uid
  5. from app.core.data_rules.contracts import rule_spec_hash
  6. from app.core.data_rules.deployment import (
  7. OperationInProgress,
  8. OperationUnknown,
  9. TerminalConflict,
  10. )
  11. from app.core.system.tokens import decode_access_token, issue_access_token
  12. from tests.core.data_rules.test_contracts import (
  13. valid_dataflow_spec,
  14. valid_rule_spec,
  15. valid_standard_spec,
  16. )
  17. from tests.core.data_rules.test_production_line import (
  18. assertion_only_rule,
  19. published_rule,
  20. published_standard,
  21. )
  22. class FakeAuthoringAgent:
  23. def __init__(self):
  24. self.calls = []
  25. def interpret(self, **kwargs):
  26. self.calls.append(kwargs)
  27. return {
  28. "status": "ready",
  29. "source_text": kwargs["source_text"],
  30. "candidate_hash": "a" * 64,
  31. "context_hash": "b" * 64,
  32. "candidate": {
  33. "candidate_type": "rule",
  34. "rule_spec": valid_rule_spec(),
  35. },
  36. }
  37. class FakeRuleRepository:
  38. def __init__(self):
  39. self.calls = []
  40. self.catalog_items = [
  41. {
  42. "asset_type": "rule",
  43. "version_id": new_governance_uid(),
  44. "asset_uid": new_governance_uid(),
  45. "name": "手机号规范化",
  46. "version": 3,
  47. "owner": new_governance_uid(),
  48. "status": "published",
  49. "schema_compatibility": {
  50. "input": "a" * 64,
  51. "output": "b" * 64,
  52. },
  53. "impact_count": 2,
  54. "backend": "polars_batch",
  55. "latest_evidence": {
  56. "compile": {"status": "success"},
  57. "test": {"status": "success"},
  58. },
  59. }
  60. ]
  61. def create_rule_version(self, **kwargs):
  62. self.calls.append(("create_rule_version", kwargs))
  63. return {
  64. "id": new_governance_uid(),
  65. "rule_uid": kwargs["rule_spec"]["rule_uid"],
  66. "version_no": 1,
  67. "status": "draft",
  68. "spec_hash": rule_spec_hash(kwargs["rule_spec"]),
  69. "created": True,
  70. }
  71. def record_generation_run(self, **kwargs):
  72. self.calls.append(("record_generation_run", kwargs))
  73. return {
  74. "id": new_governance_uid(),
  75. "correlation_id": new_governance_uid(),
  76. "decision": kwargs["evidence"]["status"],
  77. }
  78. def resolve_validation_context(self, context):
  79. self.calls.append(("resolve_validation_context", {"context": context}))
  80. return {
  81. "input_schema_snapshot_id": new_governance_uid(),
  82. "input_schema_hash": "c" * 64,
  83. "input_fields": [{"name": "mobile", "type": "string"}],
  84. "output_schema_snapshot_id": new_governance_uid(),
  85. "output_schema_hash": "d" * 64,
  86. "output_fields": [{"name": "mobile", "type": "string"}],
  87. "input_sample_artifact_ref": "minio://trusted/input.parquet",
  88. "input_sample_artifact_digest": "e" * 64,
  89. "golden_output_artifact_ref": None,
  90. "golden_output_artifact_digest": None,
  91. }
  92. def publish_rule_version(self, **kwargs):
  93. self.calls.append(("publish_rule_version", kwargs))
  94. return {
  95. "id": kwargs["version_id"],
  96. "rule_uid": new_governance_uid(),
  97. "version_no": 1,
  98. "status": "published",
  99. "spec_hash": "a" * 64,
  100. }
  101. def create_standard_version(self, **kwargs):
  102. self.calls.append(("create_standard_version", kwargs))
  103. return {
  104. "id": new_governance_uid(),
  105. "standard_uid": kwargs["standard_spec"]["standard_uid"],
  106. "version_no": 1,
  107. "status": "validated",
  108. "spec_hash": "b" * 64,
  109. "created": True,
  110. }
  111. def publish_standard_version(self, **kwargs):
  112. self.calls.append(("publish_standard_version", kwargs))
  113. return {
  114. "id": kwargs["version_id"],
  115. "standard_uid": new_governance_uid(),
  116. "version_no": 1,
  117. "status": "published",
  118. "spec_hash": "b" * 64,
  119. }
  120. def search_published_assets(self, **kwargs):
  121. self.calls.append(("search_published_assets", kwargs))
  122. return {
  123. "items": self.catalog_items,
  124. "total": 1,
  125. "limit": kwargs["limit"],
  126. "offset": kwargs["offset"],
  127. }
  128. def get_asset_evidence(self, **kwargs):
  129. self.calls.append(("get_asset_evidence", kwargs))
  130. return {
  131. "asset_type": kwargs["asset_type"],
  132. "version_id": kwargs["version_id"],
  133. "stages": {
  134. "generation": {"status": "ready"},
  135. "logical_compile": {"status": "success"},
  136. "dry_run": {"status": "success"},
  137. "publication": {"status": "published"},
  138. "physical": [
  139. {
  140. "backend": "polars_batch",
  141. "compile": {"status": "success"},
  142. "test": {"status": "success"},
  143. }
  144. ],
  145. },
  146. }
  147. def get_published_asset(self, **kwargs):
  148. self.calls.append(("get_published_asset", kwargs))
  149. return self.catalog_items[0]
  150. class FakeReleaseService:
  151. def __init__(self):
  152. self.calls = []
  153. def release(self, **kwargs):
  154. self.calls.append(kwargs)
  155. return {
  156. "id": new_governance_uid(),
  157. "version_no": 1,
  158. "status": "released",
  159. "package_hash": "c" * 64,
  160. "package": {
  161. "package_hash": "c" * 64,
  162. "standard_version_ids": [],
  163. "rule_version_ids": [],
  164. },
  165. }
  166. class FakeDeploymentService:
  167. def __init__(self):
  168. self.calls = []
  169. self.deployment_id = new_governance_uid()
  170. self.evidence_id = new_governance_uid()
  171. def create(self, version_id, **kwargs):
  172. self.calls.append(("create", version_id, kwargs))
  173. return {"id": self.deployment_id, "status": "draft"}
  174. def deploy_disabled(self, deployment_id, actor_uid, **kwargs):
  175. self.calls.append(("deploy_disabled", deployment_id, actor_uid, kwargs))
  176. return {"id": deployment_id, "status": "disabled"}
  177. def run_canary(self, deployment_id, inputs, actor_uid, **kwargs):
  178. self.calls.append(("run_canary", deployment_id, inputs, actor_uid, kwargs))
  179. return {
  180. "id": self.evidence_id,
  181. "deployment_id": deployment_id,
  182. "status": "passed",
  183. }
  184. def activate(self, deployment_id, evidence_id, actor_uid, **kwargs):
  185. self.calls.append(("activate", deployment_id, evidence_id, actor_uid, kwargs))
  186. return {"id": deployment_id, "status": "active"}
  187. def execute_active(self, deployment_id, inputs, actor_uid, **kwargs):
  188. self.calls.append(("execute_active", deployment_id, inputs, actor_uid, kwargs))
  189. return {
  190. "deployment_id": deployment_id,
  191. "execution_id": "execution-1",
  192. "correlation_id": kwargs["correlation_id"],
  193. "status": "success",
  194. }
  195. def rollback(self, deployment_id, actor_uid, **kwargs):
  196. self.calls.append(("rollback", deployment_id, actor_uid, kwargs))
  197. return {
  198. "rolled_back": {"id": deployment_id, "status": "rolled_back"},
  199. "active": {"id": new_governance_uid(), "status": "active"},
  200. }
  201. def reconcile(self, deployment_id, action, actor_uid, **kwargs):
  202. self.calls.append(("reconcile", deployment_id, action, actor_uid, kwargs))
  203. return {
  204. "deployment_id": deployment_id,
  205. "execution_id": "execution-1",
  206. "status": "success",
  207. }
  208. def list(self, environment=None):
  209. self.calls.append(("list", environment))
  210. return [{"id": self.deployment_id, "status": "draft"}]
  211. class FakePublicationService:
  212. def __init__(self, repository):
  213. self.repository = repository
  214. self.plan_id = new_governance_uid()
  215. self.version_id = None
  216. def create_draft(self, **kwargs):
  217. self.repository.calls.append(("create_rule_version", kwargs))
  218. self.version_id = new_governance_uid()
  219. return {
  220. "id": self.version_id,
  221. "rule_uid": kwargs["rule_spec"]["rule_uid"],
  222. "version_no": 1,
  223. "status": "draft",
  224. "spec_hash": rule_spec_hash(kwargs["rule_spec"]),
  225. "generation_run_id": new_governance_uid(),
  226. "created": True,
  227. }
  228. def validate(self, version_id, actor_uid):
  229. self.repository.calls.append(
  230. (
  231. "validate_rule_version",
  232. {"version_id": version_id, "actor_uid": actor_uid},
  233. )
  234. )
  235. return {
  236. "version_id": version_id,
  237. "version_status": "draft",
  238. "plan_id": self.plan_id,
  239. "plan_status": "compiled",
  240. "plan_hash": "d" * 64,
  241. }
  242. def test(self, version_id, actor_uid, *, plan_id):
  243. self.repository.calls.append(
  244. (
  245. "test_rule_version",
  246. {
  247. "version_id": version_id,
  248. "actor_uid": actor_uid,
  249. "plan_id": plan_id,
  250. },
  251. )
  252. )
  253. return {
  254. "version_id": version_id,
  255. "version_status": "validated",
  256. "plan_id": plan_id,
  257. "plan_status": "tested",
  258. "plan_hash": "d" * 64,
  259. }
  260. def publish(self, version_id, actor_uid):
  261. self.repository.calls.append(
  262. (
  263. "publish_rule_version",
  264. {"version_id": version_id, "actor_uid": actor_uid},
  265. )
  266. )
  267. return {
  268. "id": version_id,
  269. "status": "published",
  270. "plan_id": self.plan_id,
  271. "plan_status": "published",
  272. }
  273. def evidence(self, version_id):
  274. return {"version_id": version_id, "version_status": "validated"}
  275. def catalog(self, *, query, limit):
  276. return []
  277. class FakeGraphSession:
  278. def __init__(self):
  279. self.calls = []
  280. def run(self, query, parameters):
  281. self.calls.append((query, parameters))
  282. return [
  283. {
  284. "domain_id": 9,
  285. "domain_key": "customer_raw",
  286. "revision": "v2",
  287. "field_name": "customer_id",
  288. "data_type": "string",
  289. "nullable": False,
  290. "precision": None,
  291. "scale": None,
  292. "timezone": None,
  293. }
  294. ]
  295. def __enter__(self):
  296. return self
  297. def __exit__(self, *_args):
  298. return False
  299. class FakeGraphDriver:
  300. def __init__(self):
  301. self.session = FakeGraphSession()
  302. def get_session(self):
  303. return self.session
  304. class SnapshotOnlyRepository:
  305. def __init__(self):
  306. self.snapshots = {}
  307. def find_schema_snapshot(self, *, schema_ref, schema_hash):
  308. return self.snapshots.get((schema_ref, schema_hash))
  309. def persist_schema_snapshot(self, *, snapshot):
  310. value = {"id": new_governance_uid(), **snapshot}
  311. self.snapshots[(snapshot["schema_ref"], snapshot["schema_hash"])] = value
  312. return value
  313. def _headers(app, role):
  314. token = issue_access_token(
  315. user_id=new_governance_uid(),
  316. roles=[role],
  317. secret=app.config["SECRET_KEY"],
  318. now=datetime.now(UTC),
  319. lifetime=timedelta(minutes=10),
  320. )
  321. return {"Authorization": f"Bearer {token}"}
  322. def _use_token_identity(monkeypatch):
  323. def load(token, *, secret):
  324. claims = decode_access_token(token, secret=secret)
  325. return {
  326. "id": claims["sub"],
  327. "username": "contract-test",
  328. "display_name": "Contract Test",
  329. "roles": claims["roles"],
  330. }
  331. monkeypatch.setattr(
  332. "app.core.system.auth.load_identity_from_token",
  333. load,
  334. )
  335. def test_rule_capabilities_and_validation_are_registered_and_governed(
  336. monkeypatch,
  337. ):
  338. from app import create_app
  339. app = create_app()
  340. _use_token_identity(monkeypatch)
  341. app.config["TESTING"] = True
  342. app.config["RULE_GENERATION_RECEIPT_SECRET"] = (
  343. "dedicated-test-receipt-secret-with-entropy"
  344. )
  345. client = app.test_client()
  346. response = client.get("/api/rules/capabilities", headers=_headers(app, "viewer"))
  347. assert response.status_code == 200
  348. capabilities = response.get_json()["data"]
  349. assert capabilities["natural_language_authoring"] is True
  350. assert capabilities["immutable_asset_versions"] is True
  351. assert capabilities["server_side_publishing"] is True
  352. assert capabilities["production_line_release"] is True
  353. assert capabilities["data_factory_activation"] is False
  354. spec = valid_rule_spec()
  355. response = client.post(
  356. "/api/rules/validate",
  357. json={"asset_type": "rule", "spec": spec},
  358. headers=_headers(app, "editor"),
  359. )
  360. assert response.status_code == 200
  361. result = response.get_json()["data"]
  362. assert result["spec_hash"] == rule_spec_hash(spec)
  363. assert result["normalized"]["rule_uid"] == spec["rule_uid"]
  364. forbidden = client.post(
  365. "/api/rules/validate",
  366. json={"asset_type": "rule", "spec": spec},
  367. headers=_headers(app, "viewer"),
  368. )
  369. assert forbidden.status_code == 403
  370. def test_rule_interpret_uses_configured_agent_and_preserves_surface(
  371. monkeypatch,
  372. ):
  373. from app import create_app
  374. app = create_app()
  375. _use_token_identity(monkeypatch)
  376. app.config["TESTING"] = True
  377. app.config["RULE_GENERATION_RECEIPT_SECRET"] = (
  378. "dedicated-test-receipt-secret-with-entropy"
  379. )
  380. agent = FakeAuthoringAgent()
  381. app.extensions["data_rule_authoring_agent"] = agent
  382. repository = FakeRuleRepository()
  383. app.extensions["data_rule_repository"] = repository
  384. client = app.test_client()
  385. response = client.post(
  386. "/api/rules/interpret",
  387. json={
  388. "source_text": "手机号去空格后必须为11位数字",
  389. "authoring_surface": "data_standard",
  390. "context": {
  391. "input_schema_snapshot_id": new_governance_uid(),
  392. "output_schema_snapshot_id": new_governance_uid(),
  393. "input_sample_artifact_ref": ("minio://trusted/input.parquet"),
  394. "golden_output_artifact_ref": None,
  395. },
  396. },
  397. headers=_headers(app, "editor"),
  398. )
  399. assert response.status_code == 200
  400. assert response.get_json()["data"]["status"] == "ready"
  401. assert response.get_json()["data"]["generation_run_id"]
  402. assert response.get_json()["data"]["generation_receipt"]
  403. assert agent.calls[0]["authoring_surface"] == "data_standard"
  404. assert repository.calls[0][0] == "resolve_validation_context"
  405. assert repository.calls[1][0] == "record_generation_run"
  406. def test_rule_interpret_preflights_receipt_signer_before_model_call(
  407. monkeypatch,
  408. ):
  409. from app import create_app
  410. app = create_app()
  411. _use_token_identity(monkeypatch)
  412. app.config["TESTING"] = True
  413. app.config["RULE_GENERATION_RECEIPT_SECRET"] = None
  414. agent = FakeAuthoringAgent()
  415. repository = FakeRuleRepository()
  416. app.extensions["data_rule_authoring_agent"] = agent
  417. app.extensions["data_rule_repository"] = repository
  418. client = app.test_client()
  419. response = client.post(
  420. "/api/rules/interpret",
  421. json={
  422. "source_text": "手机号去空格后必须为11位数字",
  423. "authoring_surface": "data_standard",
  424. "context": {
  425. "input_schema_snapshot_id": new_governance_uid(),
  426. "output_schema_snapshot_id": new_governance_uid(),
  427. "input_sample_artifact_ref": ("minio://trusted/input.parquet"),
  428. "golden_output_artifact_ref": None,
  429. },
  430. },
  431. headers=_headers(app, "editor"),
  432. )
  433. assert response.status_code == 503
  434. assert agent.calls == []
  435. assert repository.calls == []
  436. def test_data_standard_interpret_rejects_non_rule_candidate_before_audit(
  437. monkeypatch,
  438. ):
  439. from app import create_app
  440. class StandardAgent:
  441. def interpret(self, **kwargs):
  442. return {
  443. "status": "ready",
  444. "source_text": kwargs["source_text"],
  445. "candidate_hash": "a" * 64,
  446. "context_hash": "b" * 64,
  447. "candidate": {
  448. "candidate_type": "standard",
  449. "standard_spec": valid_standard_spec(),
  450. },
  451. }
  452. app = create_app()
  453. _use_token_identity(monkeypatch)
  454. app.config["TESTING"] = True
  455. app.config["RULE_GENERATION_RECEIPT_SECRET"] = "x" * 40
  456. repository = FakeRuleRepository()
  457. app.extensions["data_rule_repository"] = repository
  458. app.extensions["data_rule_authoring_agent"] = StandardAgent()
  459. client = app.test_client()
  460. response = client.post(
  461. "/api/rules/interpret",
  462. json={
  463. "source_text": "手机号必须为11位",
  464. "authoring_surface": "data_standard",
  465. "context": {
  466. "input_schema_snapshot_id": new_governance_uid(),
  467. "output_schema_snapshot_id": new_governance_uid(),
  468. "input_sample_artifact_ref": "minio://trusted/input.parquet",
  469. "golden_output_artifact_ref": None,
  470. },
  471. },
  472. headers=_headers(app, "editor"),
  473. )
  474. assert response.status_code == 400
  475. assert not any(call[0] == "record_generation_run" for call in repository.calls)
  476. def test_rule_interpret_and_validate_reject_unknown_fields(monkeypatch):
  477. from app import create_app
  478. app = create_app()
  479. _use_token_identity(monkeypatch)
  480. app.config["TESTING"] = True
  481. app.config["RULE_GENERATION_RECEIPT_SECRET"] = (
  482. "dedicated-test-receipt-secret-with-entropy"
  483. )
  484. repository = FakeRuleRepository()
  485. app.extensions["data_rule_repository"] = repository
  486. app.extensions["data_rule_authoring_agent"] = FakeAuthoringAgent()
  487. client = app.test_client()
  488. headers = _headers(app, "editor")
  489. interpreted = client.post(
  490. "/api/rules/interpret",
  491. json={
  492. "source_text": "手机号必须为11位数字",
  493. "authoring_surface": "data_standard",
  494. "context": {},
  495. "status": "published",
  496. },
  497. headers=headers,
  498. )
  499. validated = client.post(
  500. "/api/rules/validate",
  501. json={
  502. "asset_type": "rule",
  503. "spec": valid_rule_spec(),
  504. "evidence": {"status": "success"},
  505. },
  506. headers=headers,
  507. )
  508. assert interpreted.status_code == 400
  509. assert validated.status_code == 400
  510. assert repository.calls == []
  511. def test_published_rule_catalog_uses_canonical_closed_contract(monkeypatch):
  512. from app import create_app
  513. app = create_app()
  514. _use_token_identity(monkeypatch)
  515. app.config["TESTING"] = True
  516. repository = FakeRuleRepository()
  517. app.extensions["data_rule_repository"] = repository
  518. client = app.test_client()
  519. headers = _headers(app, "viewer")
  520. response = client.get(
  521. "/api/rules/catalog?query=mobile&asset_type=rule&limit=10&offset=20",
  522. headers=headers,
  523. )
  524. rejected = client.get(
  525. "/api/rules/catalog?query=mobile&status=published",
  526. headers=headers,
  527. )
  528. assert response.status_code == 200
  529. assert response.get_json()["data"] == {
  530. "items": repository.catalog_items,
  531. "total": 1,
  532. "limit": 10,
  533. "offset": 20,
  534. }
  535. assert repository.calls[-1] == (
  536. "search_published_assets",
  537. {
  538. "query": "mobile",
  539. "asset_type": "rule",
  540. "limit": 10,
  541. "offset": 20,
  542. },
  543. )
  544. assert rejected.status_code == 400
  545. def test_unified_catalog_defaults_to_all_assets_and_preserves_rule_alias(
  546. monkeypatch,
  547. ):
  548. from app import create_app
  549. app = create_app()
  550. _use_token_identity(monkeypatch)
  551. app.config["TESTING"] = True
  552. repository = FakeRuleRepository()
  553. app.extensions["data_rule_repository"] = repository
  554. client = app.test_client()
  555. headers = _headers(app, "viewer")
  556. canonical = client.get("/api/rules/catalog", headers=headers)
  557. alias = client.get("/api/rules/catalog/rule-versions", headers=headers)
  558. assert canonical.status_code == 200
  559. assert repository.calls[0] == (
  560. "search_published_assets",
  561. {"query": "", "asset_type": None, "limit": 50, "offset": 0},
  562. )
  563. assert alias.status_code == 200
  564. assert repository.calls[1][1]["asset_type"] == "rule"
  565. def test_catalog_passes_flow_context_and_hydrates_exact_off_page_version(
  566. monkeypatch,
  567. ):
  568. import json
  569. from app import create_app
  570. app = create_app()
  571. _use_token_identity(monkeypatch)
  572. app.config["TESTING"] = True
  573. repository = FakeRuleRepository()
  574. app.extensions["data_rule_repository"] = repository
  575. client = app.test_client()
  576. headers = _headers(app, "viewer")
  577. inputs = ["bd:customer:v7"]
  578. output = "bd:customer_clean:v3"
  579. query = f"input_schema_refs={json.dumps(inputs)}&output_schema_ref={output}"
  580. listed = client.get(f"/api/rules/catalog?{query}", headers=headers)
  581. version_id = repository.catalog_items[0]["version_id"]
  582. exact = client.get(
  583. f"/api/rules/catalog/assets/rule/{version_id}?{query}",
  584. headers=headers,
  585. )
  586. assert listed.status_code == 200
  587. assert exact.status_code == 200
  588. assert repository.calls[0][1]["input_schema_refs"] == inputs
  589. assert repository.calls[0][1]["output_schema_ref"] == output
  590. assert repository.calls[1][0] == "get_published_asset"
  591. assert repository.calls[1][1]["asset_type"] == "rule"
  592. assert repository.calls[1][1]["version_id"] == version_id
  593. assert repository.calls[1][1]["input_schema_refs"] == inputs
  594. assert repository.calls[1][1]["output_schema_ref"] == output
  595. assert repository.calls[1][1]["schema_resolver"] is not None
  596. def test_catalog_and_evidence_queries_are_closed_bounded_and_rules_read_only(
  597. monkeypatch,
  598. ):
  599. from app import create_app
  600. app = create_app()
  601. _use_token_identity(monkeypatch)
  602. app.config["TESTING"] = True
  603. repository = FakeRuleRepository()
  604. app.extensions["data_rule_repository"] = repository
  605. client = app.test_client()
  606. viewer = _headers(app, "viewer")
  607. version_id = new_governance_uid()
  608. evidence = client.get(
  609. f"/api/rules/catalog/assets/rule/{version_id}/evidence",
  610. headers=viewer,
  611. )
  612. assert evidence.status_code == 200
  613. value = evidence.get_json()["data"]
  614. assert set(value["stages"]) == {
  615. "generation",
  616. "logical_compile",
  617. "dry_run",
  618. "publication",
  619. "physical",
  620. }
  621. assert "source_text" not in str(value)
  622. assert "sample" not in str(value)
  623. assert repository.calls[-1] == (
  624. "get_asset_evidence",
  625. {"asset_type": "rule", "version_id": version_id},
  626. )
  627. assert (
  628. client.get("/api/rules/catalog?asset_type=dataflow", headers=viewer).status_code
  629. == 400
  630. )
  631. assert client.get("/api/rules/catalog?limit=101", headers=viewer).status_code == 400
  632. assert client.get("/api/rules/catalog?offset=-1", headers=viewer).status_code == 400
  633. assert (
  634. client.get(f"/api/rules/catalog/assets/rule/{version_id}/evidence").status_code
  635. == 401
  636. )
  637. def test_legacy_rule_evidence_path_uses_same_safe_repository_contract(
  638. monkeypatch,
  639. ):
  640. from app import create_app
  641. app = create_app()
  642. _use_token_identity(monkeypatch)
  643. app.config["TESTING"] = True
  644. repository = FakeRuleRepository()
  645. app.extensions["data_rule_repository"] = repository
  646. client = app.test_client()
  647. version_id = new_governance_uid()
  648. response = client.get(
  649. f"/api/rules/rule-versions/{version_id}/evidence",
  650. headers=_headers(app, "viewer"),
  651. )
  652. assert response.status_code == 200
  653. assert repository.calls[-1] == (
  654. "get_asset_evidence",
  655. {"asset_type": "rule", "version_id": version_id},
  656. )
  657. def test_production_line_resolve_preview_expands_standard_without_writing(
  658. monkeypatch,
  659. ):
  660. from app import create_app
  661. app = create_app()
  662. _use_token_identity(monkeypatch)
  663. app.config["TESTING"] = True
  664. client = app.test_client()
  665. standard_id = new_governance_uid()
  666. standard_rule_id = new_governance_uid()
  667. direct_rule_id = new_governance_uid()
  668. standard_rule = published_rule(standard_rule_id, assertion_only_rule())
  669. direct_rule = published_rule(direct_rule_id)
  670. response = client.post(
  671. "/api/rules/production-lines/resolve",
  672. json={
  673. "dataflow_spec": valid_dataflow_spec(standard_id, direct_rule_id),
  674. "standard_versions": {
  675. standard_id: published_standard(standard_id, standard_rule_id)
  676. },
  677. "rule_versions": {
  678. standard_rule_id: standard_rule,
  679. direct_rule_id: direct_rule,
  680. },
  681. "component_binding_ids": {
  682. "normalize_customer": new_governance_uid(),
  683. "customer_standard:mobile_format": new_governance_uid(),
  684. },
  685. },
  686. headers=_headers(app, "editor"),
  687. )
  688. assert response.status_code == 200
  689. result = response.get_json()["data"]
  690. assert result["preview"] is True
  691. assert result["release_ready"] is False
  692. assert result["package"]["package_hash"]
  693. assert result["package"]["standard_version_ids"] == [standard_id]
  694. def test_rule_api_rejects_invalid_or_unauthenticated_requests(monkeypatch):
  695. from app import create_app
  696. app = create_app()
  697. _use_token_identity(monkeypatch)
  698. app.config["TESTING"] = True
  699. client = app.test_client()
  700. assert client.get("/api/rules/capabilities").status_code == 401
  701. response = client.post(
  702. "/api/rules/validate",
  703. json={"asset_type": "rule", "spec": {"schema_version": "1.0"}},
  704. headers=_headers(app, "editor"),
  705. )
  706. assert response.status_code == 400
  707. assert "missing" not in str(response.get_json()).lower()
  708. def test_rule_and_standard_versions_are_created_then_published_by_separate_roles(
  709. monkeypatch,
  710. ):
  711. from app import create_app
  712. from tests.core.data_rules.test_contracts import valid_standard_spec
  713. app = create_app()
  714. _use_token_identity(monkeypatch)
  715. app.config["TESTING"] = True
  716. repository = FakeRuleRepository()
  717. app.extensions["data_rule_repository"] = repository
  718. publication = FakePublicationService(repository)
  719. app.extensions["rule_publication_service"] = publication
  720. client = app.test_client()
  721. rule_spec = valid_rule_spec()
  722. created = client.post(
  723. "/api/rules/rule-versions",
  724. json={
  725. "source_text": "手机号必须为11位数字",
  726. "rule_spec": rule_spec,
  727. "category": "standard_clause",
  728. "generation_receipt": "signed-test-receipt",
  729. },
  730. headers=_headers(app, "editor"),
  731. )
  732. assert created.status_code == 201
  733. assert created.get_json()["data"]["status"] == "draft"
  734. rule_version_id = created.get_json()["data"]["id"]
  735. compiled = client.post(
  736. f"/api/rules/rule-versions/{rule_version_id}/validate",
  737. headers=_headers(app, "editor"),
  738. )
  739. assert compiled.status_code == 200
  740. plan_id = compiled.get_json()["data"]["plan_id"]
  741. tested = client.post(
  742. f"/api/rules/rule-versions/{rule_version_id}/test",
  743. json={"plan_id": plan_id},
  744. headers=_headers(app, "editor"),
  745. )
  746. assert tested.status_code == 200
  747. assert tested.get_json()["data"]["version_status"] == "validated"
  748. forbidden = client.post(
  749. f"/api/rules/rule-versions/{rule_version_id}/publish",
  750. headers=_headers(app, "editor"),
  751. )
  752. assert forbidden.status_code == 403
  753. published = client.post(
  754. f"/api/rules/rule-versions/{rule_version_id}/publish",
  755. headers=_headers(app, "admin"),
  756. )
  757. assert published.status_code == 200
  758. assert published.get_json()["data"]["status"] == "published"
  759. standard_spec = valid_standard_spec(rule_version_id)
  760. standard = client.post(
  761. "/api/rules/standard-versions",
  762. json={
  763. "source_text": "客户手机号遵循统一格式",
  764. "standard_spec": standard_spec,
  765. },
  766. headers=_headers(app, "editor"),
  767. )
  768. assert standard.status_code == 201
  769. standard_version_id = standard.get_json()["data"]["id"]
  770. standard_published = client.post(
  771. f"/api/rules/standard-versions/{standard_version_id}/publish",
  772. headers=_headers(app, "admin"),
  773. )
  774. assert standard_published.status_code == 200
  775. assert standard_published.get_json()["data"]["status"] == "published"
  776. methods = [method for method, _kwargs in repository.calls]
  777. assert methods == [
  778. "create_rule_version",
  779. "validate_rule_version",
  780. "test_rule_version",
  781. "publish_rule_version",
  782. "create_standard_version",
  783. "publish_standard_version",
  784. ]
  785. def test_create_version_rejects_client_selected_lifecycle_status(monkeypatch):
  786. from app import create_app
  787. app = create_app()
  788. _use_token_identity(monkeypatch)
  789. app.config["TESTING"] = True
  790. repository = FakeRuleRepository()
  791. app.extensions["data_rule_repository"] = repository
  792. app.extensions["rule_publication_service"] = FakePublicationService(repository)
  793. client = app.test_client()
  794. response = client.post(
  795. "/api/rules/rule-versions",
  796. json={
  797. "source_text": "手机号必须为11位数字",
  798. "rule_spec": valid_rule_spec(),
  799. "generation_receipt": "signed-test-receipt",
  800. "status": "published",
  801. },
  802. headers=_headers(app, "editor"),
  803. )
  804. assert response.status_code == 400
  805. def test_generation_receipt_signer_requires_dedicated_secret(monkeypatch):
  806. from app import create_app
  807. from app.api.data_rules.routes import _receipt_signer
  808. monkeypatch.delenv("RULE_GENERATION_RECEIPT_SECRET", raising=False)
  809. app = create_app()
  810. app.config["RULE_GENERATION_RECEIPT_SECRET"] = None
  811. with (
  812. app.app_context(),
  813. pytest.raises(RuntimeError, match="receipt secret"),
  814. ):
  815. _receipt_signer()
  816. app.config["RULE_GENERATION_RECEIPT_SECRET"] = (
  817. "dedicated-test-receipt-secret-with-entropy"
  818. )
  819. with app.app_context():
  820. signer = _receipt_signer()
  821. assert signer is not None
  822. def test_rule_gates_reject_caller_supplied_compile_or_test_evidence(
  823. monkeypatch,
  824. ):
  825. from app import create_app
  826. app = create_app()
  827. _use_token_identity(monkeypatch)
  828. app.config["TESTING"] = True
  829. repository = FakeRuleRepository()
  830. app.extensions["rule_publication_service"] = FakePublicationService(repository)
  831. client = app.test_client()
  832. version_id = new_governance_uid()
  833. forged_compile = client.post(
  834. f"/api/rules/rule-versions/{version_id}/validate",
  835. json={"status": "success", "plan_hash": "a" * 64},
  836. headers=_headers(app, "editor"),
  837. )
  838. forged_test = client.post(
  839. f"/api/rules/rule-versions/{version_id}/test",
  840. json={
  841. "plan_id": new_governance_uid(),
  842. "evidence": {"status": "success"},
  843. },
  844. headers=_headers(app, "editor"),
  845. )
  846. assert forged_compile.status_code == 409
  847. assert forged_test.status_code == 409
  848. assert repository.calls == []
  849. def test_data_factory_lifecycle_routes_are_closed_and_independently_governed(
  850. monkeypatch,
  851. ):
  852. from app import create_app
  853. app = create_app()
  854. _use_token_identity(monkeypatch)
  855. app.config["TESTING"] = True
  856. app.config["DATA_FACTORY_ACTIVATION_ENABLED"] = True
  857. service = FakeDeploymentService()
  858. app.extensions["dataflow_deployment_service"] = service
  859. client = app.test_client()
  860. admin = _headers(app, "admin")
  861. version_id = new_governance_uid()
  862. binding_id = new_governance_uid()
  863. created = client.post(
  864. "/api/rules/deployments",
  865. json={
  866. "dataflow_version_id": version_id,
  867. "environment": "production",
  868. "binding_snapshot": {
  869. "input": {"binding_id": binding_id},
  870. "output": {"binding_id": new_governance_uid()},
  871. },
  872. "schedule_plan": {"schema_version": "1.0"},
  873. "reason": "production rollout",
  874. "idempotency_key": "create-1",
  875. },
  876. headers=admin,
  877. )
  878. assert created.status_code == 201
  879. deployment_id = created.get_json()["data"]["id"]
  880. disabled = client.post(
  881. f"/api/rules/deployments/{deployment_id}/deploy-disabled",
  882. json={"reason": "disabled first", "idempotency_key": "deploy-1"},
  883. headers=admin,
  884. )
  885. assert disabled.status_code == 200
  886. canary = client.post(
  887. f"/api/rules/deployments/{deployment_id}/canary",
  888. json={
  889. "inputs": {"biz_date": "2026-07-24"},
  890. "reason": "trial production",
  891. "idempotency_key": "canary-1",
  892. },
  893. headers=admin,
  894. )
  895. assert canary.status_code == 200
  896. evidence_id = canary.get_json()["data"]["id"]
  897. active = client.post(
  898. f"/api/rules/deployments/{deployment_id}/activate",
  899. json={
  900. "evidence_id": evidence_id,
  901. "reason": "mass production",
  902. "idempotency_key": "activate-1",
  903. },
  904. headers=admin,
  905. )
  906. assert active.status_code == 200
  907. execution_correlation = new_governance_uid()
  908. executed = client.post(
  909. f"/api/rules/deployments/{deployment_id}/execute",
  910. json={
  911. "inputs": {"biz_date": "2026-07-24"},
  912. "reason": "manual production",
  913. "idempotency_key": "execute-1",
  914. "correlation_id": execution_correlation,
  915. },
  916. headers=admin,
  917. )
  918. assert executed.status_code == 200
  919. assert executed.get_json()["data"]["correlation_id"] == (execution_correlation)
  920. reconciled = client.post(
  921. f"/api/rules/deployments/{deployment_id}/reconcile-execute",
  922. json={"idempotency_key": "execute-1"},
  923. headers=admin,
  924. )
  925. assert reconciled.status_code == 200
  926. assert reconciled.get_json()["data"]["execution_id"] == "execution-1"
  927. rolled_back = client.post(
  928. f"/api/rules/deployments/{deployment_id}/rollback",
  929. json={"reason": "restore stable", "idempotency_key": "rollback-1"},
  930. headers=admin,
  931. )
  932. assert rolled_back.status_code == 200
  933. forbidden = client.post(
  934. f"/api/rules/deployments/{deployment_id}/canary",
  935. json={"inputs": {}, "idempotency_key": "viewer"},
  936. headers=_headers(app, "viewer"),
  937. )
  938. assert forbidden.status_code == 403
  939. inline_code = client.post(
  940. f"/api/rules/deployments/{deployment_id}/canary",
  941. json={
  942. "inputs": {},
  943. "idempotency_key": "unsafe",
  944. "sql": "delete from customer",
  945. },
  946. headers=admin,
  947. )
  948. assert inline_code.status_code == 409
  949. assert [call[0] for call in service.calls] == [
  950. "create",
  951. "deploy_disabled",
  952. "run_canary",
  953. "activate",
  954. "execute_active",
  955. "reconcile",
  956. "rollback",
  957. ]
  958. def test_data_factory_api_preserves_read_write_output_binding(monkeypatch):
  959. from app import create_app
  960. app = create_app()
  961. _use_token_identity(monkeypatch)
  962. app.config["TESTING"] = True
  963. service = FakeDeploymentService()
  964. app.extensions["dataflow_deployment_service"] = service
  965. output = {
  966. "data_source_uid": new_governance_uid(),
  967. "object_kind": "table",
  968. "object_ref": "curated.clean",
  969. "schema_snapshot_id": new_governance_uid(),
  970. "schema_hash": "a" * 64,
  971. "binding_hash": "b" * 64,
  972. "access_mode": "read_write",
  973. "write_mode": "merge",
  974. "dialect": "postgresql",
  975. }
  976. response = app.test_client().post(
  977. "/api/rules/deployments",
  978. json={
  979. "dataflow_version_id": new_governance_uid(),
  980. "environment": "test",
  981. "binding_snapshot": {
  982. "input": {"binding_id": new_governance_uid()},
  983. "output": output,
  984. },
  985. "schedule_plan": {"schema_version": "1.0"},
  986. "reason": "read-write pipeline",
  987. "idempotency_key": "read-write-output",
  988. },
  989. headers=_headers(app, "admin"),
  990. )
  991. assert response.status_code == 201
  992. assert service.calls[0][2]["binding_snapshot"]["output"] == output
  993. @pytest.mark.parametrize(
  994. ("exception", "code", "disposition", "reconcile_required"),
  995. [
  996. (
  997. OperationInProgress("busy"),
  998. "operation_in_progress",
  999. "retain",
  1000. False,
  1001. ),
  1002. (OperationUnknown("unknown"), "operation_unknown", "retain", True),
  1003. (
  1004. TerminalConflict("terminal"),
  1005. "terminal_conflict",
  1006. "rotate",
  1007. False,
  1008. ),
  1009. ],
  1010. )
  1011. def test_data_factory_returns_stable_idempotency_error_contract(
  1012. monkeypatch, exception, code, disposition, reconcile_required
  1013. ):
  1014. from app import create_app
  1015. app = create_app()
  1016. _use_token_identity(monkeypatch)
  1017. app.config["TESTING"] = True
  1018. service = FakeDeploymentService()
  1019. def fail(*args, **kwargs):
  1020. del args, kwargs
  1021. raise exception
  1022. service.deploy_disabled = fail
  1023. app.extensions["dataflow_deployment_service"] = service
  1024. response = app.test_client().post(
  1025. f"/api/rules/deployments/{service.deployment_id}/deploy-disabled",
  1026. json={"idempotency_key": "stable-operation-key"},
  1027. headers=_headers(app, "admin"),
  1028. )
  1029. assert response.status_code == 409
  1030. error = response.get_json()["error"]
  1031. assert error == {
  1032. "code": code,
  1033. "idempotency_key_disposition": disposition,
  1034. "reconcile_required": reconcile_required,
  1035. }
  1036. def test_data_factory_unknown_conflict_identifies_original_reconcile_target(
  1037. monkeypatch,
  1038. ):
  1039. from app import create_app
  1040. app = create_app()
  1041. _use_token_identity(monkeypatch)
  1042. app.config["TESTING"] = True
  1043. service = FakeDeploymentService()
  1044. blocker = {
  1045. "id": new_governance_uid(),
  1046. "deployment_id": service.deployment_id,
  1047. "action": "activate",
  1048. "idempotency_key": "original-unknown-key",
  1049. "scope_unknown_count": 2,
  1050. }
  1051. def fail(*args, **kwargs):
  1052. del args, kwargs
  1053. raise OperationUnknown("reconcile original", operation=blocker)
  1054. service.deploy_disabled = fail
  1055. app.extensions["dataflow_deployment_service"] = service
  1056. response = app.test_client().post(
  1057. f"/api/rules/deployments/{service.deployment_id}/deploy-disabled",
  1058. json={"idempotency_key": "different-key"},
  1059. headers=_headers(app, "admin"),
  1060. )
  1061. assert response.status_code == 409
  1062. assert response.get_json()["error"]["blocking_operation"] == blocker
  1063. def test_create_rule_version_rejects_legacy_v1_payload_before_repository(
  1064. monkeypatch,
  1065. ):
  1066. from app import create_app
  1067. app = create_app()
  1068. _use_token_identity(monkeypatch)
  1069. app.config["TESTING"] = True
  1070. repository = FakeRuleRepository()
  1071. app.extensions["data_rule_repository"] = repository
  1072. client = app.test_client()
  1073. legacy = valid_rule_spec()
  1074. legacy["schema_version"] = "1.0"
  1075. response = client.post(
  1076. "/api/rules/rule-versions",
  1077. json={"source_text": "旧版规则不能再创建", "rule_spec": legacy},
  1078. headers=_headers(app, "editor"),
  1079. )
  1080. assert response.status_code == 400
  1081. assert repository.calls == []
  1082. def test_dataflow_release_uses_server_assets_and_release_permission(
  1083. monkeypatch,
  1084. ):
  1085. from app import create_app
  1086. app = create_app()
  1087. _use_token_identity(monkeypatch)
  1088. app.config["TESTING"] = True
  1089. service = FakeReleaseService()
  1090. app.extensions["production_line_release_service"] = service
  1091. client = app.test_client()
  1092. flow = valid_dataflow_spec()
  1093. payload = {
  1094. "source_text": "客户数据生产线",
  1095. "dataflow_spec": flow,
  1096. }
  1097. forbidden = client.post(
  1098. f"/api/rules/production-lines/{flow['dataflow_uid']}/release",
  1099. json=payload,
  1100. headers=_headers(app, "editor"),
  1101. )
  1102. assert forbidden.status_code == 403
  1103. response = client.post(
  1104. f"/api/rules/production-lines/{flow['dataflow_uid']}/release",
  1105. json=payload,
  1106. headers=_headers(app, "admin"),
  1107. )
  1108. assert response.status_code == 201
  1109. assert response.get_json()["data"]["status"] == "released"
  1110. assert service.calls[0]["dataflow_uid"] == flow["dataflow_uid"]
  1111. assert "standard_versions" not in service.calls[0]
  1112. assert "rule_versions" not in service.calls[0]
  1113. assert "component_binding_ids" not in service.calls[0]
  1114. def test_dataflow_release_rejects_client_authored_schema_hashes(monkeypatch):
  1115. from app import create_app
  1116. app = create_app()
  1117. _use_token_identity(monkeypatch)
  1118. app.config["TESTING"] = True
  1119. service = FakeReleaseService()
  1120. app.extensions["production_line_release_service"] = service
  1121. client = app.test_client()
  1122. flow = valid_dataflow_spec()
  1123. response = client.post(
  1124. f"/api/rules/production-lines/{flow['dataflow_uid']}/release",
  1125. json={
  1126. "source_text": "客户数据生产线",
  1127. "dataflow_spec": flow,
  1128. "input_schema_hashes": {"bd:customer_raw:v2": "a" * 64},
  1129. "output_schema_hash": "b" * 64,
  1130. },
  1131. headers=_headers(app, "admin"),
  1132. )
  1133. assert response.status_code == 409
  1134. assert service.calls == []
  1135. def test_default_release_service_uses_lazy_neo4j_schema_catalog(monkeypatch):
  1136. from app import create_app
  1137. from app.api.data_rules.routes import _release_service
  1138. app = create_app()
  1139. repository = SnapshotOnlyRepository()
  1140. driver = FakeGraphDriver()
  1141. monkeypatch.setattr("app.core.data_rules.schema_resolver.neo4j_driver", driver)
  1142. app.extensions["data_rule_repository"] = repository
  1143. with app.app_context():
  1144. service = _release_service()
  1145. snapshot = service.schema_resolver.resolve("bd:customer_raw:v2")
  1146. assert snapshot["source_revision"] == "neo4j:9:v2"
  1147. assert driver.session.calls[0][1] == {
  1148. "domain_key": "customer_raw",
  1149. "revision": "v2",
  1150. }