test_development_api.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498
  1. from __future__ import annotations
  2. from dataclasses import replace
  3. from datetime import datetime
  4. import pytest
  5. class FakeDevelopmentService:
  6. def __init__(self, record):
  7. self.record = record
  8. self.actions = []
  9. def create_job(self, payload, actor_uid):
  10. self.actions.append(("create", payload, actor_uid))
  11. return replace(self.record, actor_uid=actor_uid), True
  12. def list_jobs(self, filters=None):
  13. self.actions.append(("list", filters or {}))
  14. return [self.record]
  15. def get_job(self, uid):
  16. self.actions.append(("get", uid))
  17. return self.record
  18. def retry(self, uid):
  19. self.actions.append(("retry", uid))
  20. return replace(self.record, status="queued", last_error=None)
  21. def cancel(self, uid):
  22. self.actions.append(("cancel", uid))
  23. return replace(self.record, status="cancelled")
  24. class FakeSourceRegistrar:
  25. def __init__(self):
  26. self.actions = []
  27. def ensure(self, source_uid, actor_uid):
  28. self.actions.append((source_uid, actor_uid))
  29. return object(), True
  30. class FakeCatalogExecutor:
  31. def __init__(self, record):
  32. self.record = record
  33. self.actions = []
  34. def execute(self, uid):
  35. self.actions.append(uid)
  36. return replace(
  37. self.record,
  38. status="awaiting_review",
  39. attempt_count=1,
  40. statistics={
  41. "snapshot_uid": "snapshot-1",
  42. "asset_count": 2,
  43. "field_count": 8,
  44. "evidence_count": 8,
  45. },
  46. )
  47. class FakeCatalogSnapshotRepository:
  48. def list(self, job_uid):
  49. from app.core.data_research.catalog.models import CatalogSnapshotRecord
  50. return [
  51. CatalogSnapshotRecord(
  52. uid="snapshot-1",
  53. job_uid=job_uid,
  54. source_uid="00000000-0000-0000-0000-000000000001",
  55. attempt=1,
  56. database_type="postgresql",
  57. content_hash="b" * 64,
  58. snapshot={
  59. "data_source_uid": "00000000-0000-0000-0000-000000000001",
  60. "database_type": "postgresql",
  61. "assets": [],
  62. },
  63. evidence_count=8,
  64. )
  65. ]
  66. class FakeEvidenceService:
  67. def list_for_job(self, job_uid):
  68. return [
  69. {
  70. "uid": "evidence-1",
  71. "job_uid": job_uid,
  72. "locator": {
  73. "kind": "database.column",
  74. "schema": "asset",
  75. "table": "equipment",
  76. "column": "equipment_code",
  77. },
  78. "excerpt": "password=[redacted]",
  79. "confidence": 1.0,
  80. }
  81. ]
  82. class FakeDeviceAssetService:
  83. def __init__(self):
  84. from app.core.data_research.device_assets import (
  85. DeviceAssetDetail,
  86. DeviceAssetImportItem,
  87. DeviceAssetImportResult,
  88. DeviceAssetMappingRecord,
  89. DeviceAssetRecord,
  90. DeviceAssetVersionRecord,
  91. )
  92. now = datetime(2026, 7, 29, 11, 0)
  93. self.asset = DeviceAssetRecord(
  94. uid="00000000-0000-7000-8000-000000000301",
  95. asset_type="device",
  96. name="一号循环泵",
  97. status="active",
  98. current_version=2,
  99. content_hash="c" * 64,
  100. location="动力车间",
  101. organization="设备动力部",
  102. responsible_person="张工",
  103. attributes={"model": "P-100"},
  104. created_by="editor-1",
  105. updated_by="editor-1",
  106. created_at=now,
  107. updated_at=now,
  108. )
  109. self.mapping = DeviceAssetMappingRecord(
  110. uid="00000000-0000-7000-8000-000000000302",
  111. asset_uid=self.asset.uid,
  112. source_uid="00000000-0000-0000-0000-000000000101",
  113. source_entity="asset.equipment",
  114. asset_type="device",
  115. source_code="EQ-001",
  116. source_updated_at=now,
  117. first_seen_at=now,
  118. last_seen_at=now,
  119. )
  120. self.version = DeviceAssetVersionRecord(
  121. uid="00000000-0000-7000-8000-000000000303",
  122. asset_uid=self.asset.uid,
  123. version=2,
  124. content_hash=self.asset.content_hash,
  125. snapshot={
  126. "asset_type": "device",
  127. "source_code": "EQ-001",
  128. "name": "一号循环泵",
  129. "status": "active",
  130. "location": "动力车间",
  131. "organization": "设备动力部",
  132. "responsible_person": "张工",
  133. "attributes": {"model": "P-100"},
  134. },
  135. source_mapping_uid=self.mapping.uid,
  136. actor_uid="editor-1",
  137. created_at=now,
  138. )
  139. self.detail = DeviceAssetDetail(
  140. asset=self.asset,
  141. mappings=(self.mapping,),
  142. )
  143. self.import_result = DeviceAssetImportResult(
  144. items=(
  145. DeviceAssetImportItem(
  146. action="updated",
  147. asset=self.asset,
  148. mapping=self.mapping,
  149. ),
  150. ),
  151. created_count=0,
  152. updated_count=1,
  153. unchanged_count=0,
  154. )
  155. self.actions = []
  156. def import_records(self, payload, actor_uid):
  157. self.actions.append(("import", payload, actor_uid))
  158. return self.import_result
  159. def search(self, filters, *, page, page_size):
  160. self.actions.append(("search", filters, page, page_size))
  161. return [self.detail], 1
  162. def get(self, asset_uid):
  163. self.actions.append(("get", asset_uid))
  164. return self.detail
  165. def versions(self, asset_uid):
  166. self.actions.append(("versions", asset_uid))
  167. return [self.version]
  168. @pytest.fixture()
  169. def development_client(monkeypatch):
  170. from flask import request
  171. from app import create_app
  172. from app.api.data_development import routes
  173. from app.core.data_research.models import IngestionJobRecord
  174. from app.core.system import permissions
  175. record = IngestionJobRecord(
  176. uid="00000000-0000-0000-0000-000000000010",
  177. source_uid="00000000-0000-0000-0000-000000000001",
  178. artifact_uid="00000000-0000-0000-0000-000000000002",
  179. job_type="file_extract",
  180. parser_version="csv-v1",
  181. idempotency_key="a" * 64,
  182. actor_uid="editor-1",
  183. parameters={"schema": "public"},
  184. )
  185. service = FakeDevelopmentService(record)
  186. registrar = FakeSourceRegistrar()
  187. executor = FakeCatalogExecutor(record)
  188. asset_service = FakeDeviceAssetService()
  189. def identity():
  190. header = request.headers.get("Authorization", "")
  191. if header == "Bearer viewer":
  192. return {"id": "viewer-1", "roles": ["viewer"]}
  193. if header == "Bearer editor":
  194. return {"id": "editor-1", "roles": ["editor"]}
  195. if header == "Bearer other-editor":
  196. return {"id": "editor-2", "roles": ["editor"]}
  197. if header == "Bearer admin":
  198. return {"id": "admin-1", "roles": ["admin"]}
  199. return None
  200. monkeypatch.setattr(permissions, "authenticate_request", identity)
  201. monkeypatch.setattr(routes, "get_ingestion_service", lambda: service)
  202. monkeypatch.setattr(
  203. routes,
  204. "get_database_source_registration_service",
  205. lambda: registrar,
  206. )
  207. monkeypatch.setattr(
  208. routes,
  209. "get_catalog_ingestion_executor",
  210. lambda: executor,
  211. )
  212. monkeypatch.setattr(
  213. routes,
  214. "get_catalog_snapshot_repository",
  215. lambda: FakeCatalogSnapshotRepository(),
  216. )
  217. monkeypatch.setattr(
  218. routes,
  219. "get_evidence_service",
  220. lambda: FakeEvidenceService(),
  221. )
  222. monkeypatch.setattr(
  223. routes,
  224. "get_device_asset_service",
  225. lambda: asset_service,
  226. raising=False,
  227. )
  228. app = create_app()
  229. app.config.update(TESTING=True)
  230. service.registrar = registrar
  231. service.executor = executor
  232. service.asset_service = asset_service
  233. return app.test_client(), service
  234. def job_payload():
  235. return {
  236. "source_uid": "00000000-0000-0000-0000-000000000001",
  237. "artifact_uid": "00000000-0000-0000-0000-000000000002",
  238. "job_type": "file_extract",
  239. "parser_version": "csv-v1",
  240. "parameters": {"schema": "public"},
  241. "password": "must-not-return",
  242. }
  243. def test_ingestion_api_requires_authentication_and_run_permission(development_client):
  244. client, _service = development_client
  245. assert client.post("/api/development/v1/ingestion-jobs", json=job_payload()).status_code == 401
  246. assert client.post(
  247. "/api/development/v1/ingestion-jobs",
  248. json=job_payload(),
  249. headers={"Authorization": "Bearer viewer"},
  250. ).status_code == 403
  251. def test_editor_creates_and_lists_secret_free_jobs(development_client):
  252. client, service = development_client
  253. created = client.post(
  254. "/api/development/v1/ingestion-jobs",
  255. json=job_payload(),
  256. headers={"Authorization": "Bearer editor"},
  257. )
  258. listed = client.get(
  259. "/api/development/v1/ingestion-jobs?status=created",
  260. headers={"Authorization": "Bearer editor"},
  261. )
  262. assert created.status_code == 201
  263. assert created.get_json()["data"]["uid"].endswith("0010")
  264. assert "must-not-return" not in created.get_data(as_text=True)
  265. assert listed.status_code == 200
  266. assert listed.get_json()["data"]["total"] == 1
  267. assert service.actions[-1] == ("list", {"status": "created"})
  268. def test_retry_is_admin_only_and_cancel_is_owner_or_admin(development_client):
  269. client, _service = development_client
  270. uid = "00000000-0000-0000-0000-000000000010"
  271. editor_retry = client.post(
  272. f"/api/development/v1/ingestion-jobs/{uid}/retry",
  273. headers={"Authorization": "Bearer editor"},
  274. )
  275. admin_retry = client.post(
  276. f"/api/development/v1/ingestion-jobs/{uid}/retry",
  277. headers={"Authorization": "Bearer admin"},
  278. )
  279. other_cancel = client.post(
  280. f"/api/development/v1/ingestion-jobs/{uid}/cancel",
  281. headers={"Authorization": "Bearer other-editor"},
  282. )
  283. owner_cancel = client.post(
  284. f"/api/development/v1/ingestion-jobs/{uid}/cancel",
  285. headers={"Authorization": "Bearer editor"},
  286. )
  287. assert editor_retry.status_code == 403
  288. assert admin_retry.status_code == 200
  289. assert other_cancel.status_code == 403
  290. assert owner_cancel.status_code == 200
  291. def test_catalog_job_registers_source_then_executes_with_report(
  292. development_client,
  293. ):
  294. client, service = development_client
  295. payload = {
  296. **job_payload(),
  297. "artifact_uid": None,
  298. "job_type": "catalog_collect",
  299. "parser_version": "catalog-v1",
  300. }
  301. created = client.post(
  302. "/api/development/v1/ingestion-jobs",
  303. json=payload,
  304. headers={"Authorization": "Bearer editor"},
  305. )
  306. executed = client.post(
  307. "/api/development/v1/ingestion-jobs/"
  308. "00000000-0000-0000-0000-000000000010/execute",
  309. headers={"Authorization": "Bearer editor"},
  310. )
  311. assert created.status_code == 201
  312. assert service.registrar.actions == [
  313. ("00000000-0000-0000-0000-000000000001", "editor-1")
  314. ]
  315. assert executed.status_code == 200
  316. assert executed.get_json()["data"]["status"] == "awaiting_review"
  317. assert executed.get_json()["data"]["attempt_count"] == 1
  318. assert service.executor.actions == [
  319. "00000000-0000-0000-0000-000000000010"
  320. ]
  321. def test_catalog_snapshots_and_evidence_are_readable_without_secrets(
  322. development_client,
  323. ):
  324. client, _service = development_client
  325. uid = "00000000-0000-0000-0000-000000000010"
  326. snapshots = client.get(
  327. f"/api/development/v1/ingestion-jobs/{uid}/catalog-snapshots",
  328. headers={"Authorization": "Bearer viewer"},
  329. )
  330. evidence = client.get(
  331. f"/api/development/v1/ingestion-jobs/{uid}/evidence",
  332. headers={"Authorization": "Bearer viewer"},
  333. )
  334. assert snapshots.status_code == 200
  335. assert snapshots.get_json()["data"]["records"][0]["attempt"] == 1
  336. assert evidence.status_code == 200
  337. text = evidence.get_data(as_text=True)
  338. assert "clear-secret" not in text
  339. assert "equipment_code" in text
  340. def test_device_asset_catalog_is_readable_and_import_requires_edit_permission(
  341. development_client,
  342. ):
  343. client, service = development_client
  344. endpoint = "/api/development/v1/device-assets"
  345. payload = {
  346. "source_uid": "00000000-0000-0000-0000-000000000101",
  347. "source_entity": "asset.equipment",
  348. "records": [
  349. {
  350. "asset_type": "device",
  351. "source_code": "EQ-001",
  352. "name": "一号循环泵",
  353. "attributes": {"model": "P-100"},
  354. }
  355. ],
  356. }
  357. viewed = client.get(
  358. f"{endpoint}?keyword=循环泵&asset_type=device&status=active"
  359. "&source_uid=00000000-0000-0000-0000-000000000101"
  360. "&page=2&page_size=10",
  361. headers={"Authorization": "Bearer viewer"},
  362. )
  363. denied = client.post(
  364. f"{endpoint}/import",
  365. json=payload,
  366. headers={"Authorization": "Bearer viewer"},
  367. )
  368. imported = client.post(
  369. f"{endpoint}/import",
  370. json=payload,
  371. headers={"Authorization": "Bearer editor"},
  372. )
  373. assert viewed.status_code == 200
  374. data = viewed.get_json()["data"]
  375. assert data["page"] == 2
  376. assert data["page_size"] == 10
  377. assert data["total"] == 1
  378. assert data["records"][0]["uid"].endswith("0301")
  379. assert data["records"][0]["source_mappings"][0]["source_code"] == "EQ-001"
  380. assert service.asset_service.actions[0] == (
  381. "search",
  382. {
  383. "keyword": "循环泵",
  384. "asset_type": "device",
  385. "status": "active",
  386. "source_uid": "00000000-0000-0000-0000-000000000101",
  387. },
  388. 2,
  389. 10,
  390. )
  391. assert denied.status_code == 403
  392. assert imported.status_code == 200
  393. assert imported.get_json()["data"]["updated_count"] == 1
  394. assert service.asset_service.actions[-1][2] == "editor-1"
  395. def test_device_asset_detail_and_versions_return_traceability_without_secrets(
  396. development_client,
  397. ):
  398. client, _service = development_client
  399. uid = "00000000-0000-7000-8000-000000000301"
  400. detail = client.get(
  401. f"/api/development/v1/device-assets/{uid}",
  402. headers={"Authorization": "Bearer viewer"},
  403. )
  404. versions = client.get(
  405. f"/api/development/v1/device-assets/{uid}/versions",
  406. headers={"Authorization": "Bearer viewer"},
  407. )
  408. assert detail.status_code == 200
  409. assert versions.status_code == 200
  410. assert detail.get_json()["data"]["current_version"] == 2
  411. assert detail.get_json()["data"]["source_mappings"][0][
  412. "source_entity"
  413. ] == "asset.equipment"
  414. assert versions.get_json()["data"]["records"][0]["version"] == 2
  415. text = detail.get_data(as_text=True) + versions.get_data(as_text=True)
  416. assert "password" not in text.lower()
  417. assert "credential" not in text.lower()
  418. def test_device_asset_catalog_rejects_unbounded_or_malformed_pagination(
  419. development_client,
  420. ):
  421. client, _service = development_client
  422. endpoint = "/api/development/v1/device-assets"
  423. malformed = client.get(
  424. f"{endpoint}?page=not-a-number",
  425. headers={"Authorization": "Bearer viewer"},
  426. )
  427. oversized = client.get(
  428. f"{endpoint}?page_size=101",
  429. headers={"Authorization": "Bearer viewer"},
  430. )
  431. assert malformed.status_code == 422
  432. assert oversized.status_code == 422
  433. assert malformed.get_json()["error"]["code"] == "DEVICE_ASSET_INVALID"