test_phase3_wp03_connectors_postgres.py 62 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635
  1. from __future__ import annotations
  2. import hashlib
  3. import os
  4. import subprocess
  5. import uuid
  6. from concurrent.futures import ThreadPoolExecutor
  7. from contextlib import contextmanager
  8. from pathlib import Path
  9. from threading import Barrier, Event, Lock
  10. import pytest
  11. from sqlalchemy import create_engine, inspect, text
  12. from sqlalchemy.exc import DBAPIError, IntegrityError
  13. from sqlalchemy.orm import Session
  14. from app.core.connectors.builtin.oracle import OracleConnector
  15. from app.core.connectors.builtin.rest_catalog import RestCatalogConnector
  16. from app.core.connectors.builtin.sqlserver import SqlServerConnector
  17. from app.core.connectors.errors import (
  18. ConnectorAuthenticationError,
  19. ConnectorConfigurationError,
  20. ConnectorRateLimitError,
  21. ConnectorUpstreamError,
  22. )
  23. from app.core.connectors.identity import ConnectorIdentityRepository
  24. from app.core.connectors.registry import ConnectorRegistry
  25. from app.core.connectors.repository import ConnectorRepository
  26. from app.core.connectors.runtime import (
  27. ConnectorRuntime,
  28. deterministic_idempotency_key,
  29. )
  30. from app.core.connectors.sdk import (
  31. SDK_VERSION,
  32. SECRET_REF_PATTERN,
  33. Connector,
  34. ConnectorManifest,
  35. OperationRequest,
  36. OperationResult,
  37. )
  38. ROOT = Path(__file__).resolve().parents[2]
  39. pytestmark = pytest.mark.integration
  40. def _uid():
  41. return str(uuid.uuid4())
  42. def _alembic(url, command, target):
  43. subprocess.run(
  44. [
  45. str(ROOT / ".venv/bin/alembic"),
  46. "-c",
  47. str(ROOT / "alembic.ini"),
  48. command,
  49. target,
  50. ],
  51. cwd=ROOT,
  52. env={**os.environ, "SQLALCHEMY_DATABASE_URI": url},
  53. check=True,
  54. capture_output=True,
  55. text=True,
  56. )
  57. def _clear_connector_test_data(engine):
  58. tables = set(inspect(engine).get_table_names(schema="public"))
  59. with engine.begin() as connection:
  60. for table in (
  61. "connector_audit_events",
  62. "connector_graph_edges",
  63. "connector_evidence",
  64. "connector_checkpoints",
  65. "connector_run_attempts",
  66. "connector_runs",
  67. "connector_machine_credentials",
  68. "connector_principals",
  69. "connector_source_bindings",
  70. "connector_manifests",
  71. "connector_rate_limits",
  72. ):
  73. if table in tables:
  74. connection.execute(text(f"DELETE FROM public.{table}"))
  75. def test_connector_flask_postgres_security_bindings_and_machine_actions(
  76. monkeypatch,
  77. ):
  78. url = os.environ.get("TEST_DATABASE_URL")
  79. if not url:
  80. pytest.skip("TEST_DATABASE_URL is required")
  81. engine = create_engine(url, pool_pre_ping=True)
  82. _clear_connector_test_data(engine)
  83. engine.dispose()
  84. _alembic(url, "upgrade", "head")
  85. monkeypatch.setenv("DATABASE_URL", url)
  86. from app import create_app
  87. actor = _uid()
  88. monkeypatch.setattr(
  89. "app.core.system.permissions.authenticate_request",
  90. lambda: {"id": actor, "sub": actor, "roles": ["admin"]},
  91. )
  92. class SpyTransport:
  93. def __init__(self):
  94. self.calls = []
  95. def get_json(self, **values):
  96. self.calls.append(values)
  97. return {
  98. "assets": [
  99. {
  100. "key": "approved-asset",
  101. "name": "Approved",
  102. "namespace": "security",
  103. "type": "table",
  104. }
  105. ]
  106. }
  107. spy = SpyTransport()
  108. registry = ConnectorRegistry()
  109. registry.register(RestCatalogConnector(spy))
  110. app = create_app()
  111. app.config.update(TESTING=True)
  112. app.extensions["connector_registry"] = registry
  113. client = app.test_client()
  114. source, domain = _uid(), _uid()
  115. binding_response = client.post(
  116. "/api/datasource/connectors/source-bindings",
  117. json={
  118. "connector_id": "rest-catalog",
  119. "version": "1.0.0",
  120. "source_uid": source,
  121. "business_domain_uid": domain,
  122. "environment": "staging",
  123. "approved_config": {
  124. "base_url": "https://approved.example.test",
  125. "allowed_host": "approved.example.test",
  126. "credential_ref": "env:DATAOPS_CONNECTOR_APPROVED",
  127. },
  128. },
  129. )
  130. assert binding_response.status_code == 201
  131. binding = binding_response.get_json()["data"]
  132. principal_response = client.post(
  133. "/api/datasource/connectors/principals",
  134. json={
  135. "connector_id": "rest-catalog",
  136. "version": "1.0.0",
  137. "source_uid": source,
  138. "business_domain_uid": domain,
  139. "environment": "staging",
  140. "operations": ["discover", "cancel", "resume"],
  141. "scopes": {},
  142. "source_binding_uid": binding["binding_uid"],
  143. "source_binding_version": binding["binding_version"],
  144. },
  145. )
  146. assert principal_response.status_code == 201
  147. principal = principal_response.get_json()["data"]["principal_uid"]
  148. rebound_response = client.post(
  149. "/api/datasource/connectors/source-bindings",
  150. json={
  151. "binding_uid": binding["binding_uid"],
  152. "connector_id": "rest-catalog",
  153. "version": "1.0.0",
  154. "source_uid": source,
  155. "business_domain_uid": domain,
  156. "environment": "staging",
  157. "approved_config": {
  158. "base_url": "https://approved.example.test",
  159. "allowed_host": "approved.example.test",
  160. "credential_ref": "env:DATAOPS_CONNECTOR_APPROVED_V2",
  161. },
  162. },
  163. )
  164. assert rebound_response.status_code == 201
  165. rebound = rebound_response.get_json()["data"]
  166. assert rebound["binding_version"] == 2
  167. assert rebound["rebound_principals"] == 1
  168. with create_engine(url).connect() as connection:
  169. versions = connection.execute(
  170. text("""
  171. SELECT p.source_binding_version,
  172. ARRAY_AGG(b.status ORDER BY b.binding_version)
  173. FROM public.connector_principals p
  174. JOIN public.connector_source_bindings b
  175. ON b.uid=p.source_binding_uid
  176. WHERE p.uid=CAST(:principal AS uuid)
  177. GROUP BY p.source_binding_version
  178. """),
  179. {"principal": principal},
  180. ).one()
  181. assert versions == (2, ["revoked", "approved"])
  182. def issue(principal_uid):
  183. response = client.post(
  184. f"/api/datasource/connectors/principals/{principal_uid}/credentials",
  185. json={"ttl_seconds": 300},
  186. )
  187. assert response.status_code == 201
  188. return response.get_json()["data"]["credential"]
  189. hint = uuid.uuid4().hex + uuid.uuid4().hex
  190. token = issue(principal)
  191. empty_config = client.post(
  192. "/api/datasource/connectors/machine/runs",
  193. headers={"X-Connector-Credential": token},
  194. json={
  195. "connector_id": "rest-catalog",
  196. "version": "1.0.0",
  197. "source_uid": source,
  198. "business_domain_uid": domain,
  199. "environment": "staging",
  200. "process_key": "security-probe",
  201. "operation": "discover",
  202. "config": {},
  203. "scope": {},
  204. "idempotency_key": hint,
  205. },
  206. )
  207. assert empty_config.status_code == 400 and spy.calls == []
  208. attack = client.post(
  209. "/api/datasource/connectors/machine/runs",
  210. headers={"X-Connector-Credential": token},
  211. json={
  212. "connector_id": "rest-catalog",
  213. "version": "1.0.0",
  214. "source_uid": source,
  215. "business_domain_uid": domain,
  216. "environment": "staging",
  217. "process_key": "security-probe",
  218. "operation": "discover",
  219. "config": {
  220. "base_url": "https://attacker.example.test",
  221. "allowed_host": "attacker.example.test",
  222. "credential_ref": "env:DATAOPS_CONNECTOR_APPROVED",
  223. },
  224. "scope": {},
  225. "idempotency_key": hint,
  226. },
  227. )
  228. assert attack.status_code == 400 and spy.calls == []
  229. created = client.post(
  230. "/api/datasource/connectors/machine/runs",
  231. headers={"X-Connector-Credential": token},
  232. json={
  233. "connector_id": "rest-catalog",
  234. "version": "1.0.0",
  235. "source_uid": source,
  236. "business_domain_uid": domain,
  237. "environment": "staging",
  238. "process_key": "security-probe",
  239. "operation": "discover",
  240. "scope": {},
  241. "idempotency_key": hint,
  242. },
  243. )
  244. assert created.status_code == 201
  245. assert len(spy.calls) == 1
  246. assert spy.calls[0]["url"].startswith("https://approved.example.test/")
  247. assert spy.calls[0]["allowed_host"] == "approved.example.test"
  248. assert spy.calls[0]["credential_ref"] == "env:DATAOPS_CONNECTOR_APPROVED_V2"
  249. assert client.post(
  250. f"/api/datasource/connectors/runs/{hint}/cancel"
  251. ).status_code == 400
  252. assert client.post(
  253. f"/api/datasource/connectors/runs/{hint}/resume"
  254. ).status_code == 400
  255. second_source, second_domain = _uid(), _uid()
  256. second_binding = client.post(
  257. "/api/datasource/connectors/source-bindings",
  258. json={
  259. "connector_id": "rest-catalog",
  260. "version": "1.0.0",
  261. "source_uid": second_source,
  262. "business_domain_uid": second_domain,
  263. "environment": "production",
  264. "approved_config": {
  265. "base_url": "https://second.example.test",
  266. "allowed_host": "second.example.test",
  267. "credential_ref": "env:DATAOPS_CONNECTOR_SECOND",
  268. },
  269. },
  270. ).get_json()["data"]
  271. second_principal = client.post(
  272. "/api/datasource/connectors/principals",
  273. json={
  274. "connector_id": "rest-catalog",
  275. "version": "1.0.0",
  276. "source_uid": second_source,
  277. "business_domain_uid": second_domain,
  278. "environment": "production",
  279. "operations": ["discover", "cancel", "resume"],
  280. "scopes": {},
  281. "source_binding_uid": second_binding["binding_uid"],
  282. "source_binding_version": second_binding["binding_version"],
  283. },
  284. ).get_json()["data"]["principal_uid"]
  285. cross = client.post(
  286. f"/api/datasource/connectors/machine/runs/{hint}/cancel",
  287. headers={"X-Connector-Credential": issue(second_principal)},
  288. )
  289. assert cross.status_code == 403
  290. replay = client.post(
  291. f"/api/datasource/connectors/machine/runs/{hint}/cancel",
  292. headers={"X-Connector-Credential": token},
  293. )
  294. assert replay.status_code == 401
  295. with create_engine(url).begin() as connection:
  296. connection.execute(
  297. text("""
  298. UPDATE public.connector_runs
  299. SET status='running',cancel_requested=FALSE
  300. WHERE client_hint_hash=:hint
  301. """),
  302. {"hint": hashlib.sha256(hint.encode()).hexdigest()},
  303. )
  304. cancelled = client.post(
  305. f"/api/datasource/connectors/machine/runs/{hint}/cancel",
  306. headers={"X-Connector-Credential": issue(principal)},
  307. )
  308. assert cancelled.status_code == 200
  309. resumed = client.post(
  310. f"/api/datasource/connectors/machine/runs/{hint}/resume",
  311. headers={"X-Connector-Credential": issue(principal)},
  312. )
  313. assert resumed.status_code == 200
  314. conflict = client.post(
  315. "/api/datasource/connectors/machine/runs",
  316. headers={"X-Connector-Credential": issue(second_principal)},
  317. json={
  318. "connector_id": "rest-catalog",
  319. "version": "1.0.0",
  320. "source_uid": second_source,
  321. "business_domain_uid": second_domain,
  322. "environment": "production",
  323. "process_key": "different-process",
  324. "operation": "discover",
  325. "scope": {},
  326. "idempotency_key": hint,
  327. },
  328. )
  329. assert conflict.status_code == 409
  330. listed = client.get("/api/datasource/connectors/source-bindings")
  331. assert listed.status_code == 200
  332. serialized = str(listed.get_json())
  333. assert "credential_ref" not in serialized
  334. assert "DATAOPS_CONNECTOR_APPROVED" not in serialized
  335. with create_engine(url).connect() as connection:
  336. persisted = connection.execute(
  337. text("""
  338. SELECT COALESCE(string_agg(payload::text,''),'')
  339. FROM public.connector_evidence
  340. WHERE run_uid IN (
  341. SELECT uid FROM public.connector_runs WHERE principal_uid=CAST(:principal AS uuid)
  342. )
  343. """),
  344. {"principal": principal},
  345. ).scalar_one()
  346. assert "DATAOPS_CONNECTOR_APPROVED" not in persisted
  347. revoked = client.post(
  348. f"/api/datasource/connectors/source-bindings/{binding['binding_uid']}/revoke"
  349. )
  350. assert revoked.status_code == 200
  351. assert revoked.get_json()["data"] == {
  352. "revoked": True,
  353. "principals_deactivated": 1,
  354. }
  355. with create_engine(url).connect() as connection:
  356. statuses = connection.execute(
  357. text("""
  358. SELECT p.status,COALESCE(MAX(c.status),'none')
  359. FROM public.connector_principals p
  360. LEFT JOIN public.connector_machine_credentials c
  361. ON c.principal_uid=p.uid
  362. WHERE p.uid=CAST(:principal AS uuid)
  363. GROUP BY p.status
  364. """),
  365. {"principal": principal},
  366. ).one()
  367. assert statuses[0] == "revoked"
  368. def test_oracle_sqlserver_bindings_authentication_and_trusted_environment(monkeypatch):
  369. url = os.environ.get("TEST_DATABASE_URL")
  370. if not url:
  371. pytest.skip("TEST_DATABASE_URL is required")
  372. engine = create_engine(url, pool_pre_ping=True)
  373. _clear_connector_test_data(engine)
  374. _alembic(url, "upgrade", "head")
  375. monkeypatch.setenv("DATABASE_URL", url)
  376. from app import create_app
  377. actor = _uid()
  378. monkeypatch.setattr(
  379. "app.core.system.permissions.authenticate_request",
  380. lambda: {"id": actor, "sub": actor, "roles": ["admin"]},
  381. )
  382. provider_calls = []
  383. class Rows:
  384. def mappings(self):
  385. return self
  386. def all(self):
  387. return [
  388. {
  389. "schema_name": "APP",
  390. "asset_name": "ORDERS",
  391. "asset_type": "TABLE",
  392. "column_name": "ID",
  393. "ordinal_position": 1,
  394. "data_type": "NUMBER",
  395. "is_nullable": "NO",
  396. "column_default": None,
  397. "column_comment": None,
  398. }
  399. ]
  400. class Connection:
  401. def execute(self, _statement, _parameters):
  402. return Rows()
  403. @contextmanager
  404. def provider(
  405. source_uid,
  406. purpose,
  407. *,
  408. environment="production",
  409. allow_insecure_development=False,
  410. ):
  411. provider_calls.append(
  412. {
  413. "source_uid": source_uid,
  414. "purpose": purpose,
  415. "environment": environment,
  416. "allow_insecure_development": allow_insecure_development,
  417. }
  418. )
  419. yield Connection()
  420. class TrackingOracle(OracleConnector):
  421. seen_configs = []
  422. def discover(self, request):
  423. self.seen_configs.append(dict(request.config))
  424. return super().discover(request)
  425. class TrackingSqlServer(SqlServerConnector):
  426. seen_configs = []
  427. def discover(self, request):
  428. self.seen_configs.append(dict(request.config))
  429. return super().discover(request)
  430. registry = ConnectorRegistry()
  431. oracle = TrackingOracle(provider)
  432. sqlserver = TrackingSqlServer(provider)
  433. registry.register(oracle)
  434. registry.register(sqlserver)
  435. app = create_app()
  436. app.config.update(TESTING=True)
  437. app.extensions["connector_registry"] = registry
  438. client = app.test_client()
  439. rejected = client.post(
  440. "/api/datasource/connectors/principals",
  441. json={
  442. "connector_id": "oracle",
  443. "version": "1.0.0",
  444. "source_uid": _uid(),
  445. "business_domain_uid": _uid(),
  446. "environment": "staging",
  447. "operations": ["discover"],
  448. "scopes": {},
  449. },
  450. )
  451. assert rejected.status_code == 400
  452. for connector_id, environment, approved_config in (
  453. (
  454. "oracle",
  455. "staging",
  456. {"credential_ref": "env:DATAOPS_CONNECTOR_ORACLE"},
  457. ),
  458. (
  459. "sqlserver",
  460. "production",
  461. {
  462. "credential_ref": "env:DATAOPS_CONNECTOR_SQLSERVER",
  463. "allow_insecure_development": True,
  464. },
  465. ),
  466. ):
  467. source, domain = _uid(), _uid()
  468. binding_response = client.post(
  469. "/api/datasource/connectors/source-bindings",
  470. json={
  471. "connector_id": connector_id,
  472. "version": "1.0.0",
  473. "source_uid": source,
  474. "business_domain_uid": domain,
  475. "environment": environment,
  476. "approved_config": approved_config,
  477. },
  478. )
  479. assert binding_response.status_code == 201
  480. binding = binding_response.get_json()["data"]
  481. principal_response = client.post(
  482. "/api/datasource/connectors/principals",
  483. json={
  484. "connector_id": connector_id,
  485. "version": "1.0.0",
  486. "source_uid": source,
  487. "business_domain_uid": domain,
  488. "environment": environment,
  489. "operations": ["discover"],
  490. "scopes": {},
  491. "source_binding_uid": binding["binding_uid"],
  492. "source_binding_version": binding["binding_version"],
  493. },
  494. )
  495. assert principal_response.status_code == 201
  496. principal = principal_response.get_json()["data"]["principal_uid"]
  497. credential_response = client.post(
  498. f"/api/datasource/connectors/principals/{principal}/credentials"
  499. )
  500. assert credential_response.status_code == 201
  501. token = credential_response.get_json()["data"]["credential"]
  502. payload = {
  503. "connector_id": connector_id,
  504. "version": "1.0.0",
  505. "source_uid": source,
  506. "business_domain_uid": domain,
  507. "environment": environment,
  508. "process_key": f"{connector_id}-catalog",
  509. "operation": "discover",
  510. "scope": {},
  511. }
  512. if connector_id == "sqlserver":
  513. forged = client.post(
  514. "/api/datasource/connectors/machine/runs",
  515. headers={"X-Connector-Credential": token},
  516. json={**payload, "environment": "development"},
  517. )
  518. assert forged.status_code == 403
  519. created = client.post(
  520. "/api/datasource/connectors/machine/runs",
  521. headers={"X-Connector-Credential": token},
  522. json=payload,
  523. )
  524. assert created.status_code == 201
  525. assert created.get_json()["data"]["records"][0]["name"] == "ORDERS"
  526. assert oracle.seen_configs == [
  527. {"credential_ref": "env:DATAOPS_CONNECTOR_ORACLE"}
  528. ]
  529. assert sqlserver.seen_configs == [
  530. {
  531. "credential_ref": "env:DATAOPS_CONNECTOR_SQLSERVER",
  532. "allow_insecure_development": True,
  533. }
  534. ]
  535. assert provider_calls[0]["environment"] == "production"
  536. assert provider_calls[0]["allow_insecure_development"] is False
  537. assert provider_calls[1]["environment"] == "production"
  538. assert provider_calls[1]["allow_insecure_development"] is True
  539. with engine.connect() as connection:
  540. configs = connection.execute(
  541. text("""
  542. SELECT connector_id,safe_config
  543. FROM public.connector_runs
  544. WHERE connector_id IN ('oracle','sqlserver')
  545. ORDER BY connector_id
  546. """)
  547. ).all()
  548. assert configs == [
  549. ("oracle", {"credential_ref": "env:DATAOPS_CONNECTOR_ORACLE"}),
  550. (
  551. "sqlserver",
  552. {
  553. "credential_ref": "env:DATAOPS_CONNECTOR_SQLSERVER",
  554. "allow_insecure_development": True,
  555. },
  556. ),
  557. ]
  558. engine.dispose()
  559. def test_connector_postgres_attempt_lease_rejects_late_terminal_writes():
  560. url = os.environ.get("TEST_DATABASE_URL")
  561. if not url:
  562. pytest.skip("TEST_DATABASE_URL is required")
  563. engine = create_engine(url, pool_pre_ping=True)
  564. _clear_connector_test_data(engine)
  565. _alembic(url, "upgrade", "head")
  566. actor, source = _uid(), _uid()
  567. connector_id = f"lease-test-{uuid.uuid4().hex[:8]}"
  568. with engine.begin() as connection:
  569. connection.execute(
  570. text("""
  571. INSERT INTO public.connector_manifests
  572. (uid,connector_id,connector_version,sdk_version,display_name,
  573. capabilities,config_schema,status,created_by)
  574. VALUES(CAST(:uid AS uuid),:connector,'1.0.0','1.0','lease-test',
  575. '["discover","cancel","resume"]'::jsonb,
  576. CAST(:schema AS jsonb),
  577. 'active',CAST(:actor AS uuid))
  578. """),
  579. {
  580. "uid": _uid(),
  581. "connector": connector_id,
  582. "actor": actor,
  583. "schema": '{"type":"object","additionalProperties":false,"properties":{}}',
  584. },
  585. )
  586. key = uuid.uuid4().hex + uuid.uuid4().hex
  587. first = Session(engine)
  588. second = Session(engine)
  589. try:
  590. repository = ConnectorRepository(first)
  591. repository.claim(
  592. key,
  593. {
  594. "idempotency_key": key,
  595. "request_hash": key,
  596. "client_hint_hash": None,
  597. "connector_id": connector_id,
  598. "connector_version": "1.0.0",
  599. "source_uid": source,
  600. "operation": "discover",
  601. "config": {},
  602. "scope": {},
  603. "checkpoint": {},
  604. "cursor": {},
  605. "dry_run": True,
  606. "actor_uid": actor,
  607. },
  608. )
  609. lease_one, lease_two = _uid(), _uid()
  610. assert repository.update(
  611. key, status="running", attempt_count=1, lease_token=lease_one
  612. ).acquired
  613. assert repository.cancel(key)["status"] == "cancelled"
  614. with Session(engine) as observer:
  615. observed = ConnectorRepository(observer).get(key)
  616. assert observed["status"] == "cancelled"
  617. assert observed["cancel_requested"] is True
  618. assert ConnectorRepository(observer).is_cancel_requested(key) is True
  619. resumed = ConnectorRepository(second).update(
  620. key, status="resumable", cancel_requested=False
  621. )
  622. assert resumed.acquired
  623. assert ConnectorRepository(second).update(
  624. key, status="running", attempt_count=2, lease_token=lease_two
  625. ).acquired
  626. stale_success = ConnectorRepository(first).update(
  627. key,
  628. status="succeeded",
  629. expected_attempt=1,
  630. lease_token=lease_one,
  631. )
  632. stale_failure = ConnectorRepository(first).update(
  633. key,
  634. status="failed",
  635. error_category="upstream",
  636. error_code="late",
  637. expected_attempt=1,
  638. lease_token=lease_one,
  639. )
  640. assert stale_success.acquired is False
  641. assert stale_failure.acquired is False
  642. winner = ConnectorRepository(second).update(
  643. key,
  644. status="succeeded",
  645. expected_attempt=2,
  646. lease_token=lease_two,
  647. )
  648. assert winner.acquired
  649. final = ConnectorRepository(second).get(key)
  650. assert final["status"] == "succeeded"
  651. assert final["attempt_count"] == 2
  652. attempts = second.execute(
  653. text("""
  654. SELECT attempt_number,status,lease_token::text
  655. FROM public.connector_run_attempts
  656. WHERE run_uid=CAST(:run AS uuid)
  657. ORDER BY attempt_number
  658. """),
  659. {"run": final["uid"]},
  660. ).all()
  661. assert attempts == [
  662. (1, "cancelled", lease_one),
  663. (2, "succeeded", lease_two),
  664. ]
  665. finally:
  666. first.close()
  667. second.close()
  668. engine.dispose()
  669. def test_connector_postgres_repository_identity_graph_cas_and_473_rollback():
  670. url = os.environ.get("TEST_DATABASE_URL")
  671. if not url:
  672. pytest.skip("TEST_DATABASE_URL is required")
  673. engine = create_engine(url, pool_pre_ping=True)
  674. _clear_connector_test_data(engine)
  675. actor, source, domain = _uid(), _uid(), _uid()
  676. connector_id, version = f"pg-test-{uuid.uuid4().hex[:8]}", "1.0.0"
  677. _alembic(url, "upgrade", "head")
  678. try:
  679. tables = inspect(engine).get_table_names(schema="public")
  680. assert {"connector_runs", "connector_rate_limits"}.issubset(tables)
  681. binding_uid = _uid()
  682. with engine.begin() as connection:
  683. connection.execute(
  684. text("""
  685. INSERT INTO public.connector_manifests
  686. (uid,connector_id,connector_version,sdk_version,display_name,capabilities,config_schema,status,created_by)
  687. VALUES(CAST(:uid AS uuid),:connector,:version,'1.0','test','["discover","cancel","resume"]'::jsonb,
  688. '{"type":"object"}'::jsonb,'active',CAST(:actor AS uuid))
  689. """),
  690. {
  691. "uid": _uid(),
  692. "connector": connector_id,
  693. "version": version,
  694. "actor": actor,
  695. },
  696. )
  697. connection.execute(
  698. text("""
  699. INSERT INTO public.connector_source_bindings
  700. (uid,binding_version,connector_id,connector_version,source_uid,
  701. business_domain_uid,environment,approved_config,status,approved_by)
  702. VALUES(CAST(:uid AS uuid),1,:connector,:version,CAST(:source AS uuid),
  703. CAST(:domain AS uuid),'staging',
  704. '{"credential_ref":"env:DATAOPS_CONNECTOR_TEST"}'::jsonb,
  705. 'approved',CAST(:actor AS uuid))
  706. """),
  707. {
  708. "uid": binding_uid,
  709. "connector": connector_id,
  710. "version": version,
  711. "source": source,
  712. "domain": domain,
  713. "actor": actor,
  714. },
  715. )
  716. idem = uuid.uuid4().hex + uuid.uuid4().hex
  717. def claim(_index):
  718. with engine.begin() as connection:
  719. return connection.execute(
  720. text("""
  721. INSERT INTO public.connector_runs
  722. (uid,idempotency_key,connector_id,connector_version,source_uid,operation,status,attempt_count,
  723. checkpoint,cursor,dry_run,actor_uid,safe_config,scope,request_hash)
  724. VALUES(CAST(:uid AS uuid),:key,:connector,:version,CAST(:source AS uuid),'discover','running',0,
  725. '{}'::jsonb,'{}'::jsonb,TRUE,CAST(:actor AS uuid),'{}'::jsonb,'{}'::jsonb,:key)
  726. ON CONFLICT(idempotency_key) DO NOTHING RETURNING uid
  727. """),
  728. {
  729. "uid": _uid(),
  730. "key": idem,
  731. "connector": connector_id,
  732. "version": version,
  733. "source": source,
  734. "actor": actor,
  735. },
  736. ).scalar_one_or_none()
  737. with ThreadPoolExecutor(max_workers=4) as executor:
  738. assert sum(item is not None for item in executor.map(claim, range(4))) == 1
  739. principal = _uid()
  740. with engine.begin() as connection:
  741. connection.execute(
  742. text("""
  743. INSERT INTO public.connector_principals
  744. (uid,connector_id,connector_version,source_uid,business_domain_uid,environment,
  745. allowed_operations,allowed_scopes,status,created_by,
  746. source_binding_uid,source_binding_version)
  747. VALUES(CAST(:uid AS uuid),:connector,:version,CAST(:source AS uuid),CAST(:domain AS uuid),
  748. 'staging',ARRAY['discover'],'{}'::jsonb,'active',CAST(:actor AS uuid),
  749. CAST(:binding AS uuid),1)
  750. """),
  751. {
  752. "uid": principal,
  753. "connector": connector_id,
  754. "version": version,
  755. "source": source,
  756. "domain": domain,
  757. "actor": actor,
  758. "binding": binding_uid,
  759. },
  760. )
  761. with pytest.raises(IntegrityError), engine.begin() as connection:
  762. connection.execute(
  763. text("""
  764. INSERT INTO public.connector_machine_credentials(uid,principal_uid,token_hash,status,issued_by,expires_at)
  765. VALUES(CAST(:uid AS uuid),CAST(:principal AS uuid),:hash,'active',CAST(:actor AS uuid),
  766. CURRENT_TIMESTAMP+INTERVAL '901 seconds')
  767. """),
  768. {
  769. "uid": _uid(),
  770. "principal": principal,
  771. "hash": "a" * 64,
  772. "actor": actor,
  773. },
  774. )
  775. session = Session(engine)
  776. identity = ConnectorIdentityRepository(session)
  777. issued = identity.issue(principal, ttl_seconds=300, actor_uid=actor)
  778. authenticated = identity.authenticate(
  779. issued["credential"],
  780. connector_id=connector_id,
  781. connector_version=version,
  782. source_uid=source,
  783. business_domain_uid=domain,
  784. environment="staging",
  785. operation="discover",
  786. scope={},
  787. )
  788. assert authenticated["principal_uid"] == principal
  789. with pytest.raises(ConnectorAuthenticationError):
  790. identity.authenticate(
  791. issued["credential"],
  792. connector_id=connector_id,
  793. connector_version=version,
  794. source_uid=source,
  795. business_domain_uid=domain,
  796. environment="staging",
  797. operation="discover",
  798. scope={},
  799. )
  800. second = identity.issue(principal, ttl_seconds=300, actor_uid=actor)
  801. rotated = identity.rotate(
  802. second["credential_uid"], ttl_seconds=300, actor_uid=actor
  803. )
  804. assert identity.revoke(rotated["credential_uid"], actor) is True
  805. limiter_key = f"{connector_id}:{source}"
  806. for _ in range(30):
  807. with Session(engine) as limiter_session:
  808. ConnectorRepository(limiter_session).acquire_rate_limit(limiter_key)
  809. with Session(engine) as limiter_session, pytest.raises(ConnectorRateLimitError):
  810. ConnectorRepository(limiter_session).acquire_rate_limit(limiter_key)
  811. run_key = uuid.uuid4().hex + uuid.uuid4().hex
  812. repository = ConnectorRepository(session, principal)
  813. repository.claim(
  814. run_key,
  815. {
  816. "connector_id": connector_id,
  817. "connector_version": version,
  818. "source_uid": source,
  819. "operation": "discover",
  820. "checkpoint": {},
  821. "cursor": {},
  822. "dry_run": False,
  823. "principal_uid": principal,
  824. "business_domain_uid": domain,
  825. "environment": "staging",
  826. "process_key": "catalog-sync",
  827. "config": {"credential_ref": "secret:test/key"},
  828. "scope": {},
  829. },
  830. )
  831. running_graph = repository.graph(run_uid=repository.get(run_key)["uid"])
  832. assert running_graph["summary"]["edge_count"] == 3
  833. assert {node["type"] for node in running_graph["nodes"]} == {
  834. "source",
  835. "process",
  836. "business_domain",
  837. "run",
  838. }
  839. assert {edge["run_status"] for edge in running_graph["edges"]} == {"running"}
  840. run_lease = _uid()
  841. repository.update(
  842. run_key, status="running", attempt_count=1, lease_token=run_lease
  843. )
  844. repository.update(
  845. run_key,
  846. status="succeeded",
  847. result=OperationResult(
  848. records=({"asset_key": f"{source}:APP.ORDERS"},),
  849. cursor={"page": 2},
  850. checkpoint={"cursor": {"page": 2}, "snapshot": [{"key": "orders"}]},
  851. evidence={"safe": True},
  852. ),
  853. checkpoint={"cursor": {"page": 2}},
  854. cursor={"page": 2},
  855. error_category=None,
  856. expected_attempt=1,
  857. lease_token=run_lease,
  858. )
  859. graph = repository.graph(source_uid=source, business_domain_uid=domain)
  860. assert {node["type"] for node in graph["nodes"]} == {
  861. "source",
  862. "asset",
  863. "process",
  864. "business_domain",
  865. "run",
  866. }
  867. assert graph["summary"]["edge_count"] == 5
  868. process_graph = repository.graph(process_key="catalog-sync")
  869. assert process_graph["summary"]["edge_count"] >= 1
  870. assert {"process", "run"}.issubset(
  871. {node["type"] for node in process_graph["nodes"]}
  872. )
  873. failed_key = uuid.uuid4().hex + uuid.uuid4().hex
  874. repository.claim(
  875. failed_key,
  876. {
  877. "connector_id": connector_id,
  878. "connector_version": version,
  879. "source_uid": source,
  880. "operation": "discover",
  881. "checkpoint": {"page": 4},
  882. "cursor": {"page": 4},
  883. "dry_run": False,
  884. "principal_uid": principal,
  885. "business_domain_uid": domain,
  886. "environment": "staging",
  887. "process_key": "failed-catalog-sync",
  888. "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  889. "scope": {},
  890. },
  891. )
  892. failed_lease = _uid()
  893. repository.update(
  894. failed_key,
  895. status="running",
  896. attempt_count=1,
  897. lease_token=failed_lease,
  898. )
  899. failed = repository.update(
  900. failed_key,
  901. status="failed",
  902. error_category="upstream",
  903. error_code="ConnectorUpstreamError",
  904. expected_attempt=1,
  905. lease_token=failed_lease,
  906. ).record
  907. assert failed["checkpoint"] == {"page": 4}
  908. failed_graph = repository.graph(run_uid=failed["uid"])
  909. assert failed_graph["summary"]["edge_count"] == 3
  910. assert {edge["run_status"] for edge in failed_graph["edges"]} == {"failed"}
  911. cancelled_key = uuid.uuid4().hex + uuid.uuid4().hex
  912. repository.claim(
  913. cancelled_key,
  914. {
  915. "connector_id": connector_id,
  916. "connector_version": version,
  917. "source_uid": source,
  918. "operation": "discover",
  919. "checkpoint": {"page": 1},
  920. "cursor": {"page": 1},
  921. "dry_run": False,
  922. "principal_uid": principal,
  923. "business_domain_uid": domain,
  924. "environment": "staging",
  925. "process_key": "catalog-sync",
  926. "config": {"credential_ref": "secret:test/key"},
  927. "scope": {},
  928. },
  929. )
  930. cancelled = Event()
  931. def cancel_worker():
  932. with Session(engine) as worker_session:
  933. result = ConnectorRepository(worker_session, principal).cancel(
  934. cancelled_key
  935. )
  936. cancelled.set()
  937. return result
  938. def completion_worker():
  939. assert cancelled.wait(timeout=5)
  940. with Session(engine) as worker_session:
  941. return ConnectorRepository(worker_session, principal).update(
  942. cancelled_key,
  943. status="succeeded",
  944. result=OperationResult(),
  945. checkpoint={},
  946. cursor={},
  947. error_category=None,
  948. )
  949. with ThreadPoolExecutor(max_workers=2) as executor:
  950. cancel_future = executor.submit(cancel_worker)
  951. complete_future = executor.submit(completion_worker)
  952. assert cancel_future.result()["status"] == "cancelled"
  953. completion_outcome = complete_future.result()
  954. assert completion_outcome.acquired is False
  955. after = completion_outcome.record
  956. assert after["status"] == "cancelled"
  957. cancelled_graph = repository.graph(run_uid=after["uid"])
  958. assert cancelled_graph["summary"]["edge_count"] == 3
  959. assert {edge["run_status"] for edge in cancelled_graph["edges"]} == {
  960. "cancelled"
  961. }
  962. with pytest.raises(IntegrityError), engine.begin() as connection:
  963. connection.execute(
  964. text("""
  965. INSERT INTO public.connector_runs(uid,idempotency_key,connector_id,connector_version,source_uid,
  966. operation,status,attempt_count,checkpoint,cursor,dry_run,actor_uid,safe_config,scope)
  967. VALUES(CAST(:uid AS uuid),:key,:connector,:version,CAST(:source AS uuid),'discover','running',0,
  968. '{}'::jsonb,'{}'::jsonb,FALSE,CAST(:actor AS uuid),'{}'::jsonb,'{}'::jsonb)
  969. """),
  970. {
  971. "uid": _uid(),
  972. "key": uuid.uuid4().hex + uuid.uuid4().hex,
  973. "connector": connector_id,
  974. "version": version,
  975. "source": source,
  976. "actor": actor,
  977. },
  978. )
  979. class RuntimeProbeConnector(Connector):
  980. manifest = ConnectorManifest(
  981. connector_id=connector_id,
  982. version=version,
  983. sdk_version=SDK_VERSION,
  984. display_name="PostgreSQL Runtime Probe",
  985. capabilities=(
  986. "discover",
  987. "snapshot",
  988. "incremental",
  989. "lineage",
  990. "profile",
  991. "cancel",
  992. "resume",
  993. "evidence",
  994. ),
  995. config_schema={
  996. "type": "object",
  997. "additionalProperties": False,
  998. "required": ["credential_ref"],
  999. "properties": {
  1000. "credential_ref": {
  1001. "type": "string",
  1002. "pattern": SECRET_REF_PATTERN.pattern,
  1003. }
  1004. },
  1005. },
  1006. )
  1007. def __init__(self):
  1008. self.discover_calls = 0
  1009. self.resume_calls = 0
  1010. self.discover_succeeds = False
  1011. self.resume_failures = 2
  1012. self.call_lock = Lock()
  1013. def discover(self, request):
  1014. with self.call_lock:
  1015. self.discover_calls += 1
  1016. if self.discover_succeeds:
  1017. return OperationResult(evidence={"concurrent": True})
  1018. raise ConnectorUpstreamError()
  1019. snapshot = incremental = lineage = profile = evidence = discover
  1020. def cancel(self, request):
  1021. return OperationResult(status="cancelled")
  1022. def resume(self, request):
  1023. with self.call_lock:
  1024. self.resume_calls += 1
  1025. call = self.resume_calls
  1026. if call <= self.resume_failures:
  1027. raise ConnectorUpstreamError()
  1028. return OperationResult(
  1029. checkpoint={"page": 8},
  1030. cursor={"page": 9},
  1031. evidence={"resumed": True},
  1032. )
  1033. probe = RuntimeProbeConnector()
  1034. registry = ConnectorRegistry()
  1035. registry.register(probe)
  1036. capped_source = _uid()
  1037. capped_request = OperationRequest(
  1038. source_uid=capped_source,
  1039. operation="discover",
  1040. config={"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1041. principal_uid=principal,
  1042. business_domain_uid=domain,
  1043. environment="staging",
  1044. process_key="attempt-cap",
  1045. )
  1046. capped_runtime = ConnectorRuntime(
  1047. registry,
  1048. store=ConnectorRepository(session, principal),
  1049. sleeper=lambda _seconds: None,
  1050. )
  1051. with pytest.raises(ConnectorUpstreamError):
  1052. capped_runtime.execute(connector_id, version, capped_request)
  1053. with pytest.raises(ConnectorUpstreamError):
  1054. capped_runtime.execute(connector_id, version, capped_request)
  1055. with pytest.raises(ConnectorConfigurationError):
  1056. capped_runtime.execute(connector_id, version, capped_request)
  1057. capped_key = capped_request.idempotency_key or deterministic_idempotency_key(
  1058. connector_id, version, capped_request
  1059. )
  1060. capped = ConnectorRepository(session).get(capped_key)
  1061. assert capped["attempt_count"] == 5 and probe.discover_calls == 5
  1062. attempts = session.execute(
  1063. text("""
  1064. SELECT COUNT(*),MAX(attempt_number)
  1065. FROM public.connector_run_attempts
  1066. WHERE run_uid=CAST(:run AS uuid)
  1067. """),
  1068. {"run": capped["uid"]},
  1069. ).one()
  1070. assert attempts == (5, 5)
  1071. rejected_source = _uid()
  1072. rejected_key = f"{connector_id}:{rejected_source}"
  1073. for _ in range(30):
  1074. ConnectorRepository(session).acquire_rate_limit(rejected_key)
  1075. rejected_request = OperationRequest(
  1076. source_uid=rejected_source,
  1077. operation="discover",
  1078. config={"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1079. dry_run=True,
  1080. )
  1081. rejected_idempotency = deterministic_idempotency_key(
  1082. connector_id, version, rejected_request
  1083. )
  1084. with pytest.raises(ConnectorRateLimitError):
  1085. ConnectorRuntime(
  1086. registry, store=ConnectorRepository(session), sleeper=lambda _: None
  1087. ).execute(connector_id, version, rejected_request)
  1088. assert ConnectorRepository(session).get(rejected_idempotency) is None
  1089. resume_source = _uid()
  1090. resume_key = uuid.uuid4().hex + uuid.uuid4().hex
  1091. resume_repository = ConnectorRepository(session, principal)
  1092. resume_repository.claim(
  1093. resume_key,
  1094. {
  1095. "connector_id": connector_id,
  1096. "connector_version": version,
  1097. "source_uid": resume_source,
  1098. "operation": "discover",
  1099. "checkpoint": {"page": 7},
  1100. "cursor": {"page": 7},
  1101. "dry_run": False,
  1102. "principal_uid": principal,
  1103. "business_domain_uid": domain,
  1104. "environment": "staging",
  1105. "process_key": "resume-state-machine",
  1106. "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1107. "scope": {},
  1108. },
  1109. )
  1110. resume_repository.update(
  1111. resume_key, status="running", attempt_count=1, lease_token=_uid()
  1112. )
  1113. resume_repository.cancel(resume_key)
  1114. resume_runtime = ConnectorRuntime(
  1115. registry,
  1116. store=resume_repository,
  1117. max_attempts=2,
  1118. sleeper=lambda _seconds: None,
  1119. )
  1120. with pytest.raises(ConnectorUpstreamError):
  1121. resume_runtime.resume(resume_key)
  1122. resume_failed = resume_repository.get(resume_key)
  1123. assert resume_failed["status"] == "failed"
  1124. assert resume_failed["attempt_count"] == 3
  1125. assert resume_failed["checkpoint"] == {"page": 7}
  1126. resumed_result = resume_runtime.resume(resume_key)
  1127. assert resumed_result.checkpoint == {"page": 8}
  1128. resumed = resume_repository.get(resume_key)
  1129. assert resumed["status"] == "succeeded" and resumed["attempt_count"] == 4
  1130. persisted = session.execute(
  1131. text("""
  1132. SELECT
  1133. (SELECT COUNT(*) FROM public.connector_run_attempts WHERE run_uid=CAST(:run AS uuid)),
  1134. (SELECT COUNT(*) FROM public.connector_evidence WHERE run_uid=CAST(:run AS uuid)),
  1135. (SELECT COUNT(*) FROM public.connector_checkpoints WHERE run_uid=CAST(:run AS uuid))
  1136. """),
  1137. {"run": resumed["uid"]},
  1138. ).one()
  1139. assert persisted == (4, 1, 1)
  1140. probe.discover_succeeds = True
  1141. probe.discover_calls = 0
  1142. retry_hint = uuid.uuid4().hex + uuid.uuid4().hex
  1143. retry_source = _uid()
  1144. retry_request = OperationRequest(
  1145. source_uid=retry_source,
  1146. operation="discover",
  1147. config={"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1148. idempotency_key=retry_hint,
  1149. principal_uid=principal,
  1150. business_domain_uid=domain,
  1151. environment="staging",
  1152. process_key="concurrent-retry",
  1153. )
  1154. retry_key = deterministic_idempotency_key(
  1155. connector_id, version, retry_request
  1156. )
  1157. retry_seed = ConnectorRepository(session, principal)
  1158. retry_seed.claim(
  1159. retry_key,
  1160. {
  1161. "connector_id": connector_id,
  1162. "connector_version": version,
  1163. "source_uid": retry_source,
  1164. "operation": "discover",
  1165. "checkpoint": {"page": 1},
  1166. "cursor": {"page": 1},
  1167. "dry_run": False,
  1168. "principal_uid": principal,
  1169. "business_domain_uid": domain,
  1170. "environment": "staging",
  1171. "process_key": "concurrent-retry",
  1172. "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1173. "scope": {},
  1174. },
  1175. )
  1176. retry_lease = _uid()
  1177. retry_seed.update(
  1178. retry_key, status="running", attempt_count=1, lease_token=retry_lease
  1179. )
  1180. retry_seed.update(
  1181. retry_key,
  1182. status="failed",
  1183. error_category="upstream",
  1184. error_code="ConnectorUpstreamError",
  1185. expected_attempt=1,
  1186. lease_token=retry_lease,
  1187. )
  1188. retry_gate = Barrier(2)
  1189. class RetryBarrierRepository(ConnectorRepository):
  1190. def claim(self, key, record):
  1191. claimed, created = super().claim(key, record)
  1192. if not created:
  1193. retry_gate.wait(timeout=5)
  1194. return claimed, created
  1195. def competing_retry():
  1196. with Session(engine) as worker_session:
  1197. worker_session.execute(text("SET statement_timeout='5s'"))
  1198. worker_session.commit()
  1199. try:
  1200. result = ConnectorRuntime(
  1201. registry,
  1202. store=RetryBarrierRepository(worker_session, principal),
  1203. max_attempts=1,
  1204. sleeper=lambda _seconds: None,
  1205. ).execute(connector_id, version, retry_request)
  1206. return "success", result.status
  1207. except Exception as error:
  1208. return "error", type(error).__name__
  1209. with ThreadPoolExecutor(max_workers=2) as executor:
  1210. retry_futures = [executor.submit(competing_retry) for _index in range(2)]
  1211. retry_results = [future.result(timeout=10) for future in retry_futures]
  1212. assert sorted(item[0] for item in retry_results) == ["error", "success"]
  1213. assert [item[1] for item in retry_results if item[0] == "error"] == [
  1214. "ConnectorConfigurationError"
  1215. ]
  1216. assert probe.discover_calls == 1
  1217. retried = ConnectorRepository(session).get(retry_key)
  1218. assert retried["status"] == "succeeded" and retried["attempt_count"] == 2
  1219. retry_attempts = session.execute(
  1220. text("""
  1221. SELECT ARRAY_AGG(attempt_number ORDER BY attempt_number)
  1222. FROM public.connector_run_attempts
  1223. WHERE run_uid=CAST(:run AS uuid)
  1224. """),
  1225. {"run": retried["uid"]},
  1226. ).scalar_one()
  1227. assert retry_attempts == [1, 2]
  1228. probe.resume_failures = 0
  1229. probe.resume_calls = 0
  1230. competing_resume_key = uuid.uuid4().hex + uuid.uuid4().hex
  1231. competing_resume_source = _uid()
  1232. resume_seed = ConnectorRepository(session, principal)
  1233. resume_seed.claim(
  1234. competing_resume_key,
  1235. {
  1236. "connector_id": connector_id,
  1237. "connector_version": version,
  1238. "source_uid": competing_resume_source,
  1239. "operation": "discover",
  1240. "checkpoint": {"page": 5},
  1241. "cursor": {"page": 5},
  1242. "dry_run": False,
  1243. "principal_uid": principal,
  1244. "business_domain_uid": domain,
  1245. "environment": "staging",
  1246. "process_key": "concurrent-resume",
  1247. "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1248. "scope": {},
  1249. },
  1250. )
  1251. resume_seed.update(
  1252. competing_resume_key,
  1253. status="running",
  1254. attempt_count=1,
  1255. lease_token=_uid(),
  1256. )
  1257. resume_seed.cancel(competing_resume_key)
  1258. resume_gate = Barrier(2)
  1259. class ResumeBarrierRepository(ConnectorRepository):
  1260. def __init__(self, worker_session, actor_uid):
  1261. super().__init__(worker_session, actor_uid)
  1262. self.waited = False
  1263. def get(self, key):
  1264. record = super().get(key)
  1265. if not self.waited and record and record.get("status") == "cancelled":
  1266. self.waited = True
  1267. resume_gate.wait(timeout=5)
  1268. return record
  1269. def competing_resume():
  1270. with Session(engine) as worker_session:
  1271. worker_session.execute(text("SET statement_timeout='5s'"))
  1272. worker_session.commit()
  1273. try:
  1274. result = ConnectorRuntime(
  1275. registry,
  1276. store=ResumeBarrierRepository(worker_session, principal),
  1277. max_attempts=1,
  1278. sleeper=lambda _seconds: None,
  1279. ).resume(competing_resume_key)
  1280. return "success", result.status
  1281. except Exception as error:
  1282. return "error", type(error).__name__
  1283. with ThreadPoolExecutor(max_workers=2) as executor:
  1284. resume_futures = [executor.submit(competing_resume) for _index in range(2)]
  1285. resume_results = [future.result(timeout=10) for future in resume_futures]
  1286. assert sorted(item[0] for item in resume_results) == ["error", "success"]
  1287. assert [item[1] for item in resume_results if item[0] == "error"] == [
  1288. "ConnectorConfigurationError"
  1289. ]
  1290. assert probe.resume_calls == 1
  1291. concurrently_resumed = ConnectorRepository(session).get(competing_resume_key)
  1292. assert concurrently_resumed["status"] == "succeeded"
  1293. assert concurrently_resumed["attempt_count"] == 2
  1294. resume_attempts = session.execute(
  1295. text("""
  1296. SELECT ARRAY_AGG(attempt_number ORDER BY attempt_number)
  1297. FROM public.connector_run_attempts
  1298. WHERE run_uid=CAST(:run AS uuid)
  1299. """),
  1300. {"run": concurrently_resumed["uid"]},
  1301. ).scalar_one()
  1302. assert resume_attempts == [1, 2]
  1303. session.close()
  1304. with engine.begin() as connection:
  1305. for table in (
  1306. "connector_audit_events",
  1307. "connector_graph_edges",
  1308. "connector_evidence",
  1309. "connector_checkpoints",
  1310. "connector_run_attempts",
  1311. "connector_runs",
  1312. "connector_machine_credentials",
  1313. "connector_principals",
  1314. "connector_source_bindings",
  1315. "connector_manifests",
  1316. "connector_rate_limits",
  1317. ):
  1318. connection.execute(text(f"DELETE FROM public.{table}"))
  1319. _alembic(url, "downgrade", "20260802_472")
  1320. ambiguous_connector = f"ambiguous-{uuid.uuid4().hex[:8]}"
  1321. ambiguous_source, ambiguous_domain = _uid(), _uid()
  1322. with engine.begin() as connection:
  1323. for ambiguous_version in ("1.0.0", "2.0.0"):
  1324. connection.execute(
  1325. text("""
  1326. INSERT INTO public.connector_manifests
  1327. (uid,connector_id,connector_version,sdk_version,display_name,
  1328. capabilities,config_schema,status,created_by)
  1329. VALUES(CAST(:uid AS uuid),:connector,:version,'1.0','ambiguous',
  1330. '["discover"]'::jsonb,'{"type":"object"}'::jsonb,
  1331. 'active',CAST(:actor AS uuid))
  1332. """),
  1333. {
  1334. "uid": _uid(),
  1335. "connector": ambiguous_connector,
  1336. "version": ambiguous_version,
  1337. "actor": actor,
  1338. },
  1339. )
  1340. connection.execute(
  1341. text("""
  1342. INSERT INTO public.connector_principals
  1343. (uid,connector_id,source_uid,business_domain_uid,environment,
  1344. allowed_operations,allowed_scopes,status,created_by)
  1345. VALUES(CAST(:uid AS uuid),:connector,CAST(:source AS uuid),
  1346. CAST(:domain AS uuid),'staging',ARRAY['discover'],
  1347. '{}'::jsonb,'active',CAST(:actor AS uuid))
  1348. """),
  1349. {
  1350. "uid": _uid(),
  1351. "connector": ambiguous_connector,
  1352. "source": ambiguous_source,
  1353. "domain": ambiguous_domain,
  1354. "actor": actor,
  1355. },
  1356. )
  1357. with pytest.raises(subprocess.CalledProcessError) as ambiguous_upgrade:
  1358. _alembic(url, "upgrade", "20260802_473")
  1359. assert "rejected before backfill" in ambiguous_upgrade.value.stderr
  1360. with engine.connect() as connection:
  1361. assert connection.execute(
  1362. text("SELECT version_num FROM alembic_version")
  1363. ).scalar_one() == "20260802_472"
  1364. assert "connector_version" not in {
  1365. item["name"]
  1366. for item in inspect(engine).get_columns(
  1367. "connector_principals", schema="public"
  1368. )
  1369. }
  1370. single_connector = f"single-{uuid.uuid4().hex[:8]}"
  1371. single_principal, single_source, single_domain = _uid(), _uid(), _uid()
  1372. with engine.begin() as connection:
  1373. connection.execute(
  1374. text("DELETE FROM public.connector_principals")
  1375. )
  1376. connection.execute(text("DELETE FROM public.connector_manifests"))
  1377. connection.execute(
  1378. text("""
  1379. INSERT INTO public.connector_manifests
  1380. (uid,connector_id,connector_version,sdk_version,display_name,
  1381. capabilities,config_schema,status,created_by)
  1382. VALUES(CAST(:uid AS uuid),:connector,'1.0.0','1.0','single',
  1383. '["discover"]'::jsonb,'{"type":"object"}'::jsonb,
  1384. 'active',CAST(:actor AS uuid))
  1385. """),
  1386. {"uid": _uid(), "connector": single_connector, "actor": actor},
  1387. )
  1388. connection.execute(
  1389. text("""
  1390. INSERT INTO public.connector_principals
  1391. (uid,connector_id,source_uid,business_domain_uid,environment,
  1392. allowed_operations,allowed_scopes,status,created_by)
  1393. VALUES(CAST(:uid AS uuid),:connector,CAST(:source AS uuid),
  1394. CAST(:domain AS uuid),'staging',ARRAY['discover'],
  1395. '{}'::jsonb,'active',CAST(:actor AS uuid))
  1396. """),
  1397. {
  1398. "uid": single_principal,
  1399. "connector": single_connector,
  1400. "source": single_source,
  1401. "domain": single_domain,
  1402. "actor": actor,
  1403. },
  1404. )
  1405. _alembic(url, "upgrade", "20260802_475")
  1406. with pytest.raises(subprocess.CalledProcessError) as unbound_upgrade:
  1407. _alembic(url, "upgrade", "20260802_476")
  1408. assert "bind every active enterprise principal first" in (
  1409. unbound_upgrade.value.stderr
  1410. )
  1411. binding_uid = _uid()
  1412. revoked_principal = _uid()
  1413. revoked_source, revoked_domain, revoked_binding = _uid(), _uid(), _uid()
  1414. with engine.begin() as connection:
  1415. connection.execute(
  1416. text("""
  1417. INSERT INTO public.connector_source_bindings
  1418. (uid,binding_version,connector_id,connector_version,source_uid,
  1419. business_domain_uid,environment,approved_config,status,approved_by)
  1420. VALUES(CAST(:uid AS uuid),1,:connector,'1.0.0',CAST(:source AS uuid),
  1421. CAST(:domain AS uuid),'staging','{}'::jsonb,'approved',CAST(:actor AS uuid))
  1422. """),
  1423. {
  1424. "uid": binding_uid,
  1425. "connector": single_connector,
  1426. "source": single_source,
  1427. "domain": single_domain,
  1428. "actor": actor,
  1429. },
  1430. )
  1431. connection.execute(
  1432. text("""
  1433. UPDATE public.connector_principals
  1434. SET source_binding_uid=CAST(:binding AS uuid),source_binding_version=1
  1435. WHERE uid=CAST(:principal AS uuid)
  1436. """),
  1437. {"binding": binding_uid, "principal": single_principal},
  1438. )
  1439. connection.execute(
  1440. text("""
  1441. INSERT INTO public.connector_principals
  1442. (uid,connector_id,connector_version,source_uid,business_domain_uid,
  1443. environment,allowed_operations,allowed_scopes,status,created_by)
  1444. VALUES(CAST(:uid AS uuid),:connector,'1.0.0',CAST(:source AS uuid),
  1445. CAST(:domain AS uuid),'staging',ARRAY['discover'],
  1446. '{}'::jsonb,'revoked',CAST(:actor AS uuid))
  1447. """),
  1448. {
  1449. "uid": revoked_principal,
  1450. "connector": single_connector,
  1451. "source": revoked_source,
  1452. "domain": revoked_domain,
  1453. "actor": actor,
  1454. },
  1455. )
  1456. _alembic(url, "upgrade", "head")
  1457. with pytest.raises(DBAPIError), engine.begin() as connection:
  1458. connection.execute(
  1459. text("""
  1460. UPDATE public.connector_principals SET status='active'
  1461. WHERE uid=CAST(:principal AS uuid)
  1462. """),
  1463. {"principal": revoked_principal},
  1464. )
  1465. with engine.begin() as connection:
  1466. connection.execute(
  1467. text("""
  1468. INSERT INTO public.connector_source_bindings
  1469. (uid,binding_version,connector_id,connector_version,source_uid,
  1470. business_domain_uid,environment,approved_config,status,approved_by)
  1471. VALUES(CAST(:uid AS uuid),1,:connector,'1.0.0',CAST(:source AS uuid),
  1472. CAST(:domain AS uuid),'staging','{}'::jsonb,'approved',CAST(:actor AS uuid))
  1473. """),
  1474. {
  1475. "uid": revoked_binding,
  1476. "connector": single_connector,
  1477. "source": revoked_source,
  1478. "domain": revoked_domain,
  1479. "actor": actor,
  1480. },
  1481. )
  1482. connection.execute(
  1483. text("""
  1484. UPDATE public.connector_principals
  1485. SET source_binding_uid=CAST(:binding AS uuid),source_binding_version=1
  1486. WHERE uid=CAST(:principal AS uuid)
  1487. """),
  1488. {"binding": revoked_binding, "principal": revoked_principal},
  1489. )
  1490. connection.execute(
  1491. text("""
  1492. UPDATE public.connector_principals SET status='active'
  1493. WHERE uid=CAST(:principal AS uuid)
  1494. """),
  1495. {"principal": revoked_principal},
  1496. )
  1497. assert connection.execute(
  1498. text("""
  1499. SELECT status FROM public.connector_principals
  1500. WHERE uid=CAST(:principal AS uuid)
  1501. """),
  1502. {"principal": revoked_principal},
  1503. ).scalar_one() == "active"
  1504. with pytest.raises(DBAPIError), engine.begin() as connection:
  1505. connection.execute(
  1506. text("""
  1507. INSERT INTO public.connector_manifests
  1508. (uid,connector_id,connector_version,sdk_version,display_name,
  1509. capabilities,config_schema,status,created_by)
  1510. VALUES(CAST(:uid AS uuid),:connector,'2.0.0','1.0','ambiguous',
  1511. '["discover"]'::jsonb,'{"type":"object"}'::jsonb,
  1512. 'active',CAST(:actor AS uuid))
  1513. """),
  1514. {"uid": _uid(), "connector": single_connector, "actor": actor},
  1515. )
  1516. with pytest.raises(subprocess.CalledProcessError) as guarded_downgrade:
  1517. _alembic(url, "downgrade", "20260802_474")
  1518. assert "explicitly remove source bindings" in guarded_downgrade.value.stderr
  1519. _alembic(url, "downgrade", "20260802_475")
  1520. with engine.begin() as connection:
  1521. connection.execute(
  1522. text("""
  1523. UPDATE public.connector_principals
  1524. SET source_binding_uid=NULL,source_binding_version=NULL
  1525. WHERE uid IN (CAST(:principal AS uuid),CAST(:revoked AS uuid))
  1526. """),
  1527. {"principal": single_principal, "revoked": revoked_principal},
  1528. )
  1529. connection.execute(
  1530. text(
  1531. "DELETE FROM public.connector_source_bindings WHERE uid=CAST(:uid AS uuid)"
  1532. ),
  1533. {"uid": binding_uid},
  1534. )
  1535. connection.execute(
  1536. text(
  1537. "DELETE FROM public.connector_source_bindings WHERE uid=CAST(:uid AS uuid)"
  1538. ),
  1539. {"uid": revoked_binding},
  1540. )
  1541. _alembic(url, "downgrade", "20260802_474")
  1542. assert "connector_source_bindings" not in inspect(engine).get_table_names(
  1543. schema="public"
  1544. )
  1545. _clear_connector_test_data(engine)
  1546. _alembic(url, "upgrade", "head")
  1547. finally:
  1548. _clear_connector_test_data(engine)
  1549. _alembic(url, "upgrade", "head")
  1550. engine.dispose()