test_data_research_docker_api.py 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282
  1. """Opt-in end-to-end acceptance for the containerized data-research API."""
  2. from __future__ import annotations
  3. import os
  4. import uuid
  5. import pytest
  6. import requests
  7. from minio import Minio
  8. from sqlalchemy import create_engine, text
  9. pytestmark = pytest.mark.skipif(
  10. os.getenv("RUN_DATA_RESEARCH_DOCKER_API") != "1",
  11. reason="set RUN_DATA_RESEARCH_DOCKER_API=1 for Docker API acceptance",
  12. )
  13. def _data(response, expected_status):
  14. assert response.status_code == expected_status, response.text
  15. payload = response.json()
  16. assert payload["code"] == 200, payload
  17. return payload["data"]
  18. def test_file_ingestion_and_dynamic_ontology_publish_through_live_api():
  19. base_url = os.getenv("TEST_BACKEND_URL", "http://127.0.0.1:15500/api")
  20. database_url = os.environ["TEST_DATABASE_URL"]
  21. engine = create_engine(database_url, pool_pre_ping=True)
  22. minio = Minio(
  23. os.getenv("TEST_MINIO_HOST", "127.0.0.1:19000"),
  24. access_key="dataops-test",
  25. secret_key="dataops-test-password",
  26. secure=False,
  27. )
  28. source_uid = str(uuid.uuid4())
  29. owner_domain_uid = str(uuid.uuid4())
  30. contributor_domain_uid = str(uuid.uuid4())
  31. ontology_uid = None
  32. artifact_uid = None
  33. artifact_object_key = None
  34. login = requests.post(
  35. f"{base_url}/system/auth/login",
  36. json={"username": "admin", "password": os.environ["TEST_ADMIN_PASSWORD"]},
  37. timeout=20,
  38. )
  39. token = _data(login, 200)["token"]
  40. headers = {"Authorization": f"Bearer {token}"}
  41. try:
  42. with engine.begin() as connection:
  43. connection.execute(
  44. text(
  45. "INSERT INTO public.ingestion_sources "
  46. "(uid, source_type, name, created_by) "
  47. "VALUES (CAST(:uid AS uuid), 'file', :name, 'docker-api-test')"
  48. ),
  49. {"uid": source_uid, "name": f"docker-api-{source_uid}"},
  50. )
  51. csv_content = b"code,name\nCUSTOMER_ID,Customer identifier\n"
  52. uploaded = requests.post(
  53. f"{base_url}/development/v1/sources/files",
  54. headers=headers,
  55. data={"source_uid": source_uid, "parser_version": "csv-v1"},
  56. files={"file": ("elements.csv", csv_content, "text/csv")},
  57. timeout=20,
  58. )
  59. artifact = _data(uploaded, 201)
  60. artifact_uid = artifact["uid"]
  61. duplicate = requests.post(
  62. f"{base_url}/development/v1/sources/files",
  63. headers=headers,
  64. data={"source_uid": source_uid, "parser_version": "csv-v1"},
  65. files={"file": ("elements.csv", csv_content, "text/csv")},
  66. timeout=20,
  67. )
  68. assert _data(duplicate, 200)["uid"] == artifact_uid
  69. job_payload = {
  70. "source_uid": source_uid,
  71. "artifact_uid": artifact_uid,
  72. "job_type": "file_extract",
  73. "parser_version": "csv-v1",
  74. "parameters": {"delimiter": ","},
  75. }
  76. created_job = _data(
  77. requests.post(
  78. f"{base_url}/development/v1/ingestion-jobs",
  79. headers=headers,
  80. json=job_payload,
  81. timeout=20,
  82. ),
  83. 201,
  84. )
  85. duplicate_job = _data(
  86. requests.post(
  87. f"{base_url}/development/v1/ingestion-jobs",
  88. headers=headers,
  89. json=job_payload,
  90. timeout=20,
  91. ),
  92. 200,
  93. )
  94. assert duplicate_job["uid"] == created_job["uid"]
  95. cancelled = _data(
  96. requests.post(
  97. f"{base_url}/development/v1/ingestion-jobs/"
  98. f"{created_job['uid']}/cancel",
  99. headers=headers,
  100. timeout=20,
  101. ),
  102. 200,
  103. )
  104. assert cancelled["status"] == "cancelled"
  105. ontology = _data(
  106. requests.post(
  107. f"{base_url}/development/v1/ontologies",
  108. headers=headers,
  109. json={
  110. "code": f"DOCKER_{uuid.uuid4().hex[:12].upper()}",
  111. "name": "Docker API customer ontology",
  112. "domain_links": [
  113. {"domain_uid": owner_domain_uid, "role": "owner"},
  114. {
  115. "domain_uid": contributor_domain_uid,
  116. "role": "contributor",
  117. },
  118. ],
  119. },
  120. timeout=20,
  121. ),
  122. 201,
  123. )
  124. ontology_uid = ontology["uid"]
  125. graph = {
  126. "classes": [{"uid": "customer", "name": "Customer"}],
  127. "properties": [
  128. {
  129. "uid": "customer-id",
  130. "name": "customerId",
  131. "class_uid": "customer",
  132. "required": True,
  133. "cardinality": "1",
  134. "data_element_uid": "customer-id-element",
  135. }
  136. ],
  137. "relations": [],
  138. "constraints": [],
  139. "domain_links": [
  140. {"domain_uid": owner_domain_uid, "role": "owner"},
  141. {
  142. "domain_uid": contributor_domain_uid,
  143. "role": "contributor",
  144. },
  145. ],
  146. "element_mappings": [
  147. {
  148. "property_uid": "customer-id",
  149. "data_element_uid": "customer-id-element",
  150. }
  151. ],
  152. }
  153. draft = _data(
  154. requests.patch(
  155. f"{base_url}/development/v1/ontologies/{ontology_uid}/graph",
  156. headers={**headers, "If-Match": '"0"'},
  157. json=graph,
  158. timeout=20,
  159. ),
  160. 200,
  161. )
  162. assert draft["version"] == 1
  163. assert _data(
  164. requests.post(
  165. f"{base_url}/development/v1/ontologies/{ontology_uid}/validate",
  166. headers=headers,
  167. timeout=20,
  168. ),
  169. 200,
  170. ) == []
  171. published = _data(
  172. requests.post(
  173. f"{base_url}/development/v1/ontologies/{ontology_uid}/publish",
  174. headers={**headers, "Idempotency-Key": f"docker-{ontology_uid}"},
  175. timeout=20,
  176. ),
  177. 200,
  178. )
  179. assert published["status"] == "published"
  180. exported = requests.get(
  181. f"{base_url}/development/v1/ontologies/{ontology_uid}/export?format=json",
  182. headers=headers,
  183. timeout=20,
  184. )
  185. assert exported.status_code == 200
  186. assert exported.json()["graph_document"]["classes"][0]["uid"] == "customer"
  187. with engine.connect() as connection:
  188. artifact_object_key = connection.execute(
  189. text(
  190. "SELECT storage_ref FROM public.source_artifacts "
  191. "WHERE uid = CAST(:uid AS uuid)"
  192. ),
  193. {"uid": artifact_uid},
  194. ).scalar_one().removeprefix("minio://dataops-bucket/")
  195. outbox_payload = connection.execute(
  196. text(
  197. "SELECT payload FROM public.outbox_events "
  198. "WHERE aggregate_id = :uid "
  199. "AND event_type = 'ontology.version_published'"
  200. ),
  201. {"uid": ontology_uid},
  202. ).scalar_one()
  203. assert outbox_payload["ontology_uid"] == ontology_uid
  204. finally:
  205. with engine.begin() as connection:
  206. if ontology_uid:
  207. connection.execute(
  208. text(
  209. "DELETE FROM public.outbox_consumptions WHERE event_id IN "
  210. "(SELECT event_id FROM public.outbox_events "
  211. "WHERE aggregate_id = :uid)"
  212. ),
  213. {"uid": ontology_uid},
  214. )
  215. connection.execute(
  216. text("DELETE FROM public.outbox_events WHERE aggregate_id = :uid"),
  217. {"uid": ontology_uid},
  218. )
  219. connection.execute(
  220. text(
  221. "UPDATE public.ontologies SET active_version_uid = NULL "
  222. "WHERE uid = CAST(:uid AS uuid)"
  223. ),
  224. {"uid": ontology_uid},
  225. )
  226. for table in (
  227. "ontology_publish_runs",
  228. "ontology_change_sets",
  229. "ontology_domain_links",
  230. "ontology_versions",
  231. "ontologies",
  232. ):
  233. connection.execute(
  234. text(
  235. f"DELETE FROM public.{table} "
  236. "WHERE ontology_uid = CAST(:uid AS uuid)"
  237. if table != "ontologies"
  238. else "DELETE FROM public.ontologies "
  239. "WHERE uid = CAST(:uid AS uuid)"
  240. ),
  241. {"uid": ontology_uid},
  242. )
  243. connection.execute(
  244. text(
  245. "DELETE FROM public.ingestion_jobs "
  246. "WHERE source_uid = CAST(:uid AS uuid)"
  247. ),
  248. {"uid": source_uid},
  249. )
  250. connection.execute(
  251. text(
  252. "DELETE FROM public.source_artifacts "
  253. "WHERE source_uid = CAST(:uid AS uuid)"
  254. ),
  255. {"uid": source_uid},
  256. )
  257. connection.execute(
  258. text(
  259. "DELETE FROM public.ingestion_sources "
  260. "WHERE uid = CAST(:uid AS uuid)"
  261. ),
  262. {"uid": source_uid},
  263. )
  264. if artifact_object_key:
  265. minio.remove_object("dataops-bucket", artifact_object_key)
  266. engine.dispose()