test_phase3_wp03_connectors_postgres.py 64 KB

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