test_phase3_wp03_enterprise_connectors.py 57 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697
  1. from __future__ import annotations
  2. from contextlib import contextmanager
  3. from pathlib import Path
  4. from types import SimpleNamespace
  5. import pytest
  6. import yaml
  7. from jsonschema import Draft202012Validator
  8. from app.core.connectors import (
  9. SDK_VERSION,
  10. Connector,
  11. ConnectorConfigurationError,
  12. ConnectorManifest,
  13. ConnectorRegistry,
  14. ConnectorUpstreamError,
  15. EnvironmentSecretResolver,
  16. OperationRequest,
  17. OperationResult,
  18. validate_config,
  19. )
  20. from app.core.connectors.builtin import (
  21. OracleConnector,
  22. RestCatalogConnector,
  23. SafeRestTransport,
  24. SqlServerConnector,
  25. )
  26. from app.core.connectors.errors import (
  27. ConnectorAuthenticationError,
  28. ConnectorContractError,
  29. ConnectorDriverUnavailableError,
  30. ConnectorPermissionError,
  31. ConnectorRateLimitError,
  32. ConnectorTimeoutError,
  33. classify_error,
  34. )
  35. from app.core.connectors.identity import ConnectorIdentityRepository
  36. from app.core.connectors.runtime import (
  37. MAX_TOTAL_ATTEMPTS,
  38. ConnectorRuntime,
  39. InMemoryRunStore,
  40. SlidingWindowLimiter,
  41. deterministic_idempotency_key,
  42. redact_evidence,
  43. snapshot_diff,
  44. )
  45. from app.core.data_source.adapters.oracle import OracleAdapter
  46. from app.core.data_source.adapters.sqlserver import SqlServerAdapter
  47. from app.core.data_source.errors import (
  48. DataSourceConfigurationInvalid,
  49. DataSourceDriverUnavailable,
  50. )
  51. from app.core.data_source.models import DataSourceCredential, DataSourceDefinition
  52. ROOT = Path(__file__).resolve().parents[1]
  53. ALL_CAPABILITIES = (
  54. "discover",
  55. "snapshot",
  56. "incremental",
  57. "lineage",
  58. "profile",
  59. "cancel",
  60. "resume",
  61. "evidence",
  62. )
  63. class _Rows:
  64. def __init__(self, rows):
  65. self.rows = rows
  66. def mappings(self):
  67. return self
  68. def all(self):
  69. return self.rows
  70. class _Connection:
  71. def __init__(self, rows):
  72. self.rows = rows
  73. self.statement = None
  74. self.parameters = None
  75. def execute(self, statement, parameters):
  76. self.statement = str(statement)
  77. self.parameters = parameters
  78. return _Rows(self.rows)
  79. @contextmanager
  80. def _provider(_source, purpose):
  81. assert purpose == "metadata_collection"
  82. yield _Connection(
  83. [
  84. {
  85. "schema_name": "APP",
  86. "asset_name": "ORDERS",
  87. "asset_type": "TABLE",
  88. "column_name": "ID",
  89. "ordinal_position": 1,
  90. "data_type": "NUMBER",
  91. "is_nullable": "NO",
  92. "column_default": None,
  93. "column_comment": "key",
  94. }
  95. ]
  96. )
  97. def _request(operation="discover", **values):
  98. defaults = {
  99. "source_uid": "11111111-1111-4111-8111-111111111111",
  100. "operation": operation,
  101. "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  102. "scope": {},
  103. }
  104. defaults.update(values)
  105. return OperationRequest(**defaults)
  106. def test_manifest_config_and_secret_boundary_are_strict():
  107. manifest = OracleConnector.manifest
  108. assert manifest.sdk_version == SDK_VERSION
  109. assert set(manifest.capabilities) == {
  110. "discover",
  111. "snapshot",
  112. "incremental",
  113. "cancel",
  114. "resume",
  115. "evidence",
  116. }
  117. registry = ConnectorRegistry()
  118. registry.register(OracleConnector(_provider))
  119. with pytest.raises(ConnectorConfigurationError):
  120. registry.resolve("oracle", "1.0.0", "lineage")
  121. assert registry.validate("oracle", "1.0.0", {"credential_ref": "vault:team/oracle"})
  122. with pytest.raises(ConnectorConfigurationError):
  123. registry.validate("oracle", "1.0.0", {"password": "plain"})
  124. with pytest.raises(ConnectorConfigurationError):
  125. registry.validate("oracle", "9.9.9", {"credential_ref": "env:X"})
  126. with pytest.raises(ConnectorConfigurationError):
  127. ConnectorManifest(
  128. "bad",
  129. "1",
  130. SDK_VERSION,
  131. "Bad",
  132. ("discover",),
  133. {"type": "object", "additionalProperties": True, "properties": {}},
  134. )
  135. def test_draft_202012_schema_recursively_enforces_nested_arrays_bounds_and_one_of():
  136. schema = {
  137. "$schema": "https://json-schema.org/draft/2020-12/schema",
  138. "type": "object",
  139. "additionalProperties": False,
  140. "required": ["nested"],
  141. "properties": {
  142. "nested": {
  143. "type": "object",
  144. "additionalProperties": False,
  145. "required": ["items", "mode"],
  146. "properties": {
  147. "items": {
  148. "type": "array",
  149. "minItems": 1,
  150. "maxItems": 2,
  151. "items": {"type": "integer", "minimum": 1, "maximum": 5},
  152. },
  153. "mode": {"oneOf": [{"const": "a"}, {"const": "b"}]},
  154. },
  155. }
  156. },
  157. }
  158. assert validate_config(schema, {"nested": {"items": [1, 5], "mode": "a"}})
  159. for invalid in (
  160. {"nested": {"items": [], "mode": "a"}},
  161. {"nested": {"items": [6], "mode": "a"}},
  162. {"nested": {"items": [1], "mode": "c"}},
  163. {"nested": {"items": [1], "mode": "a", "extra": True}},
  164. ):
  165. with pytest.raises(ConnectorConfigurationError):
  166. validate_config(schema, invalid)
  167. with pytest.raises(ConnectorConfigurationError):
  168. ConnectorManifest(
  169. "nested-secret",
  170. "1.0.0",
  171. SDK_VERSION,
  172. "Nested Secret",
  173. ("discover",),
  174. {
  175. "type": "object",
  176. "additionalProperties": False,
  177. "properties": {
  178. "nested": {
  179. "type": "object",
  180. "properties": {"password": {"type": "string"}},
  181. }
  182. },
  183. },
  184. )
  185. @pytest.mark.parametrize("operation", ALL_CAPABILITIES)
  186. def test_all_unified_interfaces_are_callable_and_undeclared_capability_is_denied(
  187. operation,
  188. ):
  189. connector = OracleConnector(_provider)
  190. result = getattr(connector, operation)(_request(operation))
  191. assert isinstance(result, OperationResult)
  192. class DiscoverOnly(Connector):
  193. manifest = ConnectorManifest(
  194. "discover-only",
  195. "1.0.0",
  196. SDK_VERSION,
  197. "Discover",
  198. ("discover",),
  199. {"type": "object", "additionalProperties": False, "properties": {}},
  200. )
  201. def discover(self, request):
  202. return OperationResult()
  203. snapshot = incremental = lineage = profile = cancel = resume = evidence = (
  204. discover
  205. )
  206. registry = ConnectorRegistry()
  207. registry.register(DiscoverOnly())
  208. with pytest.raises(ConnectorConfigurationError):
  209. registry.resolve("discover-only", "1.0.0", "snapshot")
  210. @pytest.mark.parametrize(
  211. "connector_cls, sql_marker",
  212. [(OracleConnector, "all_tab_columns"), (SqlServerConnector, "sys.columns")],
  213. )
  214. def test_database_connectors_normalize_scope_and_use_read_only_catalog_sql(
  215. connector_cls, sql_marker
  216. ):
  217. connector = connector_cls(_provider)
  218. result = connector.discover(
  219. _request(scope={"include_schemas": ["APP"], "include_tables": ["ORDERS"]})
  220. )
  221. assert result.records[0]["asset_key"].endswith(":APP.ORDERS")
  222. assert result.records[0]["nullable"] is False
  223. assert result.evidence["query_kind"] == "read_only_metadata"
  224. assert sql_marker in connector.catalog_sql
  225. assert not any(
  226. word in connector.catalog_sql.upper()
  227. for word in ("INSERT ", "UPDATE ", "DELETE ")
  228. )
  229. def test_database_incremental_persists_bounded_snapshot_summary():
  230. result = OracleConnector(_provider).incremental(
  231. _request(
  232. "incremental",
  233. checkpoint={"snapshot_summary": {"sha256": "old", "record_count": 2}},
  234. )
  235. )
  236. diff = result.evidence["diff"]
  237. assert diff == {
  238. "changed": True,
  239. "previous_record_count": 2,
  240. "current_record_count": 1,
  241. "details_materialized": False,
  242. }
  243. assert "snapshot" not in result.checkpoint
  244. assert result.checkpoint["snapshot_summary"]["record_count"] == 1
  245. class _HttpResponse:
  246. def __init__(self, status, body, content_type="application/json"):
  247. self.status = status
  248. self.body = body
  249. self.headers = {"Content-Type": content_type}
  250. self.released = False
  251. def stream(self, **_kwargs):
  252. for offset in range(0, len(self.body), 65536):
  253. yield self.body[offset : offset + 65536]
  254. def release_conn(self):
  255. self.released = True
  256. class _HttpsPool:
  257. def __init__(self, calls, response, **kwargs):
  258. self.calls = calls
  259. self.response = response
  260. self.calls.append(("pool", kwargs))
  261. def urlopen(self, *args, **kwargs):
  262. self.calls.append(("request", {"args": args, **kwargs}))
  263. return self.response.pop(0) if isinstance(self.response, list) else self.response
  264. def close(self):
  265. self.calls.append(("close", {}))
  266. def _safe_transport(calls, response, resolver=None, secret_resolver=None):
  267. return SafeRestTransport(
  268. secret_resolver=secret_resolver or (lambda ref: "resolved-token"),
  269. resolver=resolver
  270. or (lambda *_args, **_kwargs: [(2, 1, 6, "", ("8.8.8.8", 443))]),
  271. pool_factory=lambda **kwargs: _HttpsPool(calls, response, **kwargs),
  272. )
  273. def test_reference_rest_connector_uses_same_registry_without_core_dispatcher_changes():
  274. calls, resolutions = [], []
  275. def resolver(host, port, type):
  276. return [(2, 1, 6, "", ("8.8.8.8", port))]
  277. def secret_resolver(reference):
  278. resolutions.append(reference)
  279. return "resolved-token"
  280. registry = ConnectorRegistry()
  281. page = b'{"assets":[{"key":"a","name":"Orders","namespace":"sales","type":"table"}],"next_cursor":{"page":2}}'
  282. end = b'{"assets":[]}'
  283. response = [
  284. _HttpResponse(200, page),
  285. _HttpResponse(200, end),
  286. _HttpResponse(200, page),
  287. _HttpResponse(200, end),
  288. ]
  289. registry.register(
  290. RestCatalogConnector(
  291. _safe_transport(calls, response, resolver, secret_resolver)
  292. )
  293. )
  294. result = ConnectorRuntime(registry, sleeper=lambda _: None).execute(
  295. "rest-catalog",
  296. "1.0.0",
  297. _request(
  298. config={
  299. "base_url": "https://catalog.example.test",
  300. "allowed_host": "catalog.example.test",
  301. "credential_ref": "secret:catalog/key",
  302. }
  303. ),
  304. )
  305. assert result.records[0]["name"] == "Orders"
  306. pool_call = next(value for kind, value in calls if kind == "pool")
  307. request_call = next(value for kind, value in calls if kind == "request")
  308. assert pool_call["host"] == "8.8.8.8"
  309. assert (
  310. pool_call["server_hostname"]
  311. == pool_call["assert_hostname"]
  312. == "catalog.example.test"
  313. )
  314. assert pool_call["cert_reqs"] == "CERT_REQUIRED" and pool_call["retries"] is False
  315. assert request_call["headers"]["Host"] == "catalog.example.test"
  316. assert request_call["headers"]["Authorization"] == "Bearer resolved-token"
  317. assert request_call["redirect"] is False and request_call["retries"] is False
  318. assert resolutions == ["secret:catalog/key", "secret:catalog/key"]
  319. assert "resolved-token" not in str(result.evidence)
  320. incremental = ConnectorRuntime(registry, sleeper=lambda _: None).execute(
  321. "rest-catalog",
  322. "1.0.0",
  323. _request(
  324. "incremental",
  325. config={
  326. "base_url": "https://catalog.example.test",
  327. "allowed_host": "catalog.example.test",
  328. "credential_ref": "secret:catalog/key",
  329. },
  330. checkpoint={
  331. "cursor": {"page": 1},
  332. "snapshot_summary": {"sha256": "old", "record_count": 1},
  333. },
  334. ),
  335. )
  336. incremental_request = [value for kind, value in calls if kind == "request"][-1]
  337. assert "cursor=" in incremental_request["args"][1]
  338. assert incremental.cursor["state"] == "complete"
  339. assert incremental.cursor["completion_marker"]
  340. assert incremental.evidence["diff"]["changed"] is True
  341. assert incremental.evidence["diff"]["details_materialized"] is False
  342. assert resolutions == ["secret:catalog/key"] * 4
  343. with pytest.raises(ConnectorConfigurationError):
  344. _safe_transport(
  345. calls,
  346. response,
  347. lambda *_args, **_kwargs: [(2, 1, 6, "", ("127.0.0.1", 443))],
  348. ).get_json(
  349. url="https://localhost/x",
  350. allowed_host="localhost",
  351. credential_ref="secret:test/key",
  352. )
  353. core = (ROOT / "app/core/connectors/registry.py").read_text() + (
  354. ROOT / "app/core/connectors/runtime.py"
  355. ).read_text()
  356. assert all(
  357. connector_id not in core
  358. for connector_id in ("oracle", "sqlserver", "rest-catalog")
  359. )
  360. @pytest.mark.parametrize(
  361. "status,body,error",
  362. [
  363. (302, b"{}", ConnectorUpstreamError),
  364. (200, b"x" * (2 * 1024 * 1024 + 1), Exception),
  365. (200, b"not-json", Exception),
  366. ],
  367. )
  368. def test_rest_transport_rejects_redirect_oversize_and_invalid_json(status, body, error):
  369. transport = _safe_transport([], _HttpResponse(status, body))
  370. with pytest.raises(error):
  371. transport.get_json(
  372. url="https://catalog.example.test/v1/catalog",
  373. allowed_host="catalog.example.test",
  374. credential_ref="secret:test/key",
  375. )
  376. with pytest.raises(ConnectorConfigurationError):
  377. transport.get_json(
  378. url="http://catalog.example.test/v1/catalog",
  379. allowed_host="catalog.example.test",
  380. credential_ref="secret:test/key",
  381. )
  382. with pytest.raises(ConnectorConfigurationError):
  383. transport.get_json(
  384. url="https://evil.example.test/v1/catalog",
  385. allowed_host="catalog.example.test",
  386. credential_ref="secret:test/key",
  387. )
  388. def test_rest_transport_without_secret_resolver_fails_before_dns_or_socket():
  389. calls = []
  390. transport = SafeRestTransport(
  391. resolver=lambda *_args, **_kwargs: calls.append("dns"),
  392. pool_factory=lambda **_kwargs: calls.append("pool"),
  393. )
  394. with pytest.raises(ConnectorConfigurationError):
  395. transport.get_json(
  396. url="https://catalog.example.test/v1/catalog",
  397. allowed_host="catalog.example.test",
  398. credential_ref="secret:test/key",
  399. )
  400. assert calls == []
  401. def test_default_environment_secret_resolver_is_deployable_and_fail_closed(
  402. monkeypatch,
  403. ):
  404. resolver = EnvironmentSecretResolver(
  405. {"DATAOPS_CONNECTOR_CATALOG_TOKEN": "runtime-only-value"}
  406. )
  407. assert resolver("env:DATAOPS_CONNECTOR_CATALOG_TOKEN") == "runtime-only-value"
  408. for reference in (
  409. "env:HOME",
  410. "env:DATAOPS_DATABASE_PASSWORD",
  411. "vault:catalog/token",
  412. "secret:catalog/token",
  413. ):
  414. with pytest.raises(ConnectorConfigurationError):
  415. resolver(reference)
  416. from flask import Flask
  417. from app.api.data_source import routes
  418. monkeypatch.setenv("DATAOPS_CONNECTOR_CATALOG_TOKEN", "assembled-secret")
  419. monkeypatch.setattr(
  420. "app.core.data_source.runtime.get_data_source_manager",
  421. lambda: SimpleNamespace(connect=_provider),
  422. )
  423. app = Flask(__name__)
  424. app.config["CONNECTOR_SECRET_RESOLVER"] = None
  425. with app.app_context():
  426. registry = routes._connector_registry()
  427. transport = registry.resolve("rest-catalog", "1.0.0").transport
  428. assert (
  429. transport.secret_resolver("env:DATAOPS_CONNECTOR_CATALOG_TOKEN")
  430. == "assembled-secret"
  431. )
  432. with pytest.raises(ConnectorConfigurationError):
  433. transport.secret_resolver("env:HOME")
  434. @pytest.mark.parametrize(
  435. "adapter_cls,module", [(OracleAdapter, "oracledb"), (SqlServerAdapter, "pyodbc")]
  436. )
  437. def test_optional_enterprise_drivers_fail_with_stable_safe_classification(
  438. monkeypatch, adapter_cls, module
  439. ):
  440. monkeypatch.setattr(
  441. "importlib.util.find_spec", lambda name: None if name == module else object()
  442. )
  443. definition = DataSourceDefinition(
  444. uid="u",
  445. name_en="x",
  446. name_zh="x",
  447. database_type=adapter_cls.database_type,
  448. host="db.example",
  449. port=1521,
  450. database="service",
  451. schema=None,
  452. credential_ref="u",
  453. credential_version=1,
  454. pool_size=None,
  455. max_overflow=None,
  456. tls_options={},
  457. status=True,
  458. description=None,
  459. extra_properties={},
  460. )
  461. with pytest.raises(DataSourceDriverUnavailable) as caught:
  462. adapter_cls().build_url(
  463. definition, DataSourceCredential("reader", "not-logged", {})
  464. )
  465. assert caught.value.code == "DATASOURCE_DRIVER_UNAVAILABLE"
  466. assert classify_error(caught.value).category == "configuration"
  467. @pytest.mark.parametrize(
  468. "connector_cls,module",
  469. [(OracleConnector, "oracledb"), (SqlServerConnector, "pyodbc")],
  470. )
  471. def test_health_and_compatibility_report_missing_driver_without_crash(
  472. monkeypatch, connector_cls, module
  473. ):
  474. monkeypatch.setattr(
  475. "importlib.util.find_spec",
  476. lambda name: None if name == module else object(),
  477. )
  478. connector = connector_cls(_provider)
  479. health = connector.health({"credential_ref": "env:DATAOPS_CONNECTOR_TEST"})
  480. compatibility = connector.compatibility()
  481. assert health.status == "degraded" and "driver" in health.detail
  482. assert compatibility.compatible is False
  483. assert compatibility.connector_version == "1.0.0"
  484. def test_runtime_idempotency_retry_cancel_resume_diff_and_bounded_redaction():
  485. class Flaky(Connector):
  486. manifest = OracleConnector.manifest
  487. def __init__(self):
  488. self.calls = 0
  489. def discover(self, request):
  490. self.calls += 1
  491. if self.calls == 1:
  492. raise ConnectorUpstreamError()
  493. return OperationResult(
  494. records=({"x": 1},),
  495. checkpoint={"page": 1},
  496. evidence={"password": "hidden"},
  497. )
  498. snapshot = incremental = lineage = profile = resume = evidence = discover
  499. def cancel(self, request):
  500. return OperationResult(status="cancelled")
  501. connector = Flaky()
  502. registry = ConnectorRegistry()
  503. registry.register(connector)
  504. store = InMemoryRunStore()
  505. runtime = ConnectorRuntime(registry, store=store, sleeper=lambda _: None)
  506. request = _request()
  507. first = runtime.execute("oracle", "1.0.0", request)
  508. second = runtime.execute("oracle", "1.0.0", request)
  509. assert connector.calls == 2 and second == first
  510. assert first.evidence["password"] == "[REDACTED]"
  511. assert len(deterministic_idempotency_key("oracle", "1.0.0", request)) == 64
  512. assert snapshot_diff(({"a": 1},), ({"a": 2},))["changed"] == (
  513. {"before": {"a": 1}, "after": {"a": 2}},
  514. )
  515. assert len(str(redact_evidence({"data": "x" * 40000})["data"])) == 4096
  516. key = "a" * 64
  517. store.claim(
  518. key,
  519. {
  520. "status": "running",
  521. "connector_id": "oracle",
  522. "connector_version": "1.0.0",
  523. "source_uid": request.source_uid,
  524. "config": dict(request.config),
  525. "scope": {},
  526. "cursor": {},
  527. "checkpoint": {"page": 7},
  528. "attempt_count": 1,
  529. "dry_run": False,
  530. },
  531. )
  532. assert runtime.cancel(key)["status"] == "cancelled"
  533. assert isinstance(runtime.resume(key), OperationResult)
  534. assert store.get(key)["status"] == "succeeded"
  535. assert store.get(key)["attempt_count"] == 2
  536. assert store.get(deterministic_idempotency_key("oracle", "1.0.0", request))[
  537. "checkpoint"
  538. ] == {"page": 1}
  539. def test_human_rest_dry_run_never_resolves_secret_or_opens_transport():
  540. class NoNetwork:
  541. def __init__(self):
  542. self.calls = 0
  543. def get_json(self, **_values):
  544. self.calls += 1
  545. raise AssertionError("dry-run must not use network")
  546. transport = NoNetwork()
  547. registry = ConnectorRegistry()
  548. registry.register(RestCatalogConnector(transport))
  549. result = ConnectorRuntime(registry).execute(
  550. "rest-catalog",
  551. "1.0.0",
  552. _request(
  553. dry_run=True,
  554. config={
  555. "base_url": "https://attacker.example.test",
  556. "allowed_host": "attacker.example.test",
  557. "credential_ref": "env:DATAOPS_CONNECTOR_GUESSED",
  558. },
  559. ),
  560. )
  561. assert result.status == "dry_run"
  562. assert result.evidence["network_requested"] is False
  563. assert transport.calls == 0
  564. def test_cancelled_human_dry_run_resumes_validation_only_without_network():
  565. class NoNetwork:
  566. def __init__(self):
  567. self.calls = 0
  568. def get_json(self, **_values):
  569. self.calls += 1
  570. raise AssertionError("human dry-run resume must not use transport")
  571. transport = NoNetwork()
  572. registry = ConnectorRegistry()
  573. registry.register(RestCatalogConnector(transport))
  574. store = InMemoryRunStore()
  575. key = "d" * 64
  576. config = {
  577. "base_url": "https://catalog.example.test",
  578. "allowed_host": "catalog.example.test",
  579. "credential_ref": "env:DATAOPS_CONNECTOR_TEST",
  580. }
  581. store.claim(
  582. key,
  583. {
  584. "idempotency_key": key,
  585. "status": "running",
  586. "connector_id": "rest-catalog",
  587. "connector_version": "1.0.0",
  588. "source_uid": _request().source_uid,
  589. "config": config,
  590. "scope": {},
  591. "cursor": {},
  592. "checkpoint": {},
  593. "attempt_count": 1,
  594. "attempt_lease_token": "old-lease",
  595. "dry_run": True,
  596. "principal_uid": None,
  597. },
  598. )
  599. runtime = ConnectorRuntime(registry, store=store)
  600. assert runtime.cancel(key)["cancel_requested"] is True
  601. resumed = runtime.resume(key)
  602. assert resumed.status == "dry_run"
  603. assert resumed.evidence["validation_only"] is True
  604. assert resumed.evidence["network_requested"] is False
  605. assert transport.calls == 0
  606. def test_runtime_rejects_nonterminal_operation_result_status():
  607. class InvalidStatus(Connector):
  608. manifest = OracleConnector.manifest
  609. def discover(self, _request):
  610. return OperationResult(status="running")
  611. snapshot = incremental = lineage = profile = resume = evidence = discover
  612. def cancel(self, _request):
  613. return OperationResult(status="cancelled")
  614. registry = ConnectorRegistry()
  615. registry.register(InvalidStatus())
  616. store = InMemoryRunStore()
  617. request = _request(idempotency_key="invalid-status")
  618. with pytest.raises(ConnectorConfigurationError):
  619. ConnectorRuntime(registry, store=store, max_attempts=1).execute(
  620. "oracle", "1.0.0", request
  621. )
  622. key = deterministic_idempotency_key("oracle", "1.0.0", request)
  623. assert store.get(key)["status"] == "failed"
  624. def test_cancelling_one_run_does_not_cancel_another_run_for_the_same_source():
  625. class Probe(Connector):
  626. manifest = OracleConnector.manifest
  627. def __init__(self):
  628. self.discover_calls = 0
  629. self.cancel_request = None
  630. def discover(self, _request):
  631. self.discover_calls += 1
  632. return OperationResult(records=({"asset_key": "APP.ORDERS"},))
  633. snapshot = incremental = lineage = profile = resume = evidence = discover
  634. def cancel(self, request):
  635. self.cancel_request = request
  636. return OperationResult(status="cancelled")
  637. connector = Probe()
  638. registry = ConnectorRegistry()
  639. registry.register(connector)
  640. store = InMemoryRunStore()
  641. runtime = ConnectorRuntime(registry, store=store)
  642. cancelled_key = "c" * 64
  643. store.claim(
  644. cancelled_key,
  645. {
  646. "status": "running",
  647. "connector_id": "oracle",
  648. "connector_version": "1.0.0",
  649. "source_uid": _request().source_uid,
  650. "config": dict(_request().config),
  651. "scope": {},
  652. "cursor": {},
  653. "checkpoint": {},
  654. "attempt_count": 1,
  655. "attempt_lease_token": "old-lease",
  656. "dry_run": False,
  657. },
  658. )
  659. assert runtime.cancel(cancelled_key)["status"] == "cancelled"
  660. assert connector.cancel_request.run_key == cancelled_key
  661. assert connector.cancel_request.cancel_probe() is True
  662. sibling = runtime.execute(
  663. "oracle",
  664. "1.0.0",
  665. _request(idempotency_key="sibling-run"),
  666. )
  667. assert sibling.status == "succeeded"
  668. assert connector.discover_calls == 1
  669. def test_rest_completed_cursor_restarts_full_scan_without_false_removed():
  670. class Pages:
  671. def __init__(self):
  672. self.urls = []
  673. def get_json(self, **values):
  674. self.urls.append(values["url"])
  675. return {
  676. "assets": [
  677. {
  678. "key": "orders",
  679. "name": "Orders",
  680. "namespace": "sales",
  681. "type": "table",
  682. }
  683. ],
  684. "completion_token": "server-complete-42",
  685. }
  686. pages = Pages()
  687. connector = RestCatalogConnector(pages)
  688. config = {
  689. "base_url": "https://catalog.example.test",
  690. "allowed_host": "catalog.example.test",
  691. "credential_ref": "env:DATAOPS_CONNECTOR_TEST",
  692. }
  693. first = connector.snapshot(_request("snapshot", config=config))
  694. second = connector.incremental(
  695. _request(
  696. "incremental",
  697. config=config,
  698. cursor=first.cursor,
  699. checkpoint=first.checkpoint,
  700. )
  701. )
  702. assert pages.urls == [
  703. "https://catalog.example.test/v1/catalog",
  704. "https://catalog.example.test/v1/catalog",
  705. ]
  706. assert second.cursor == {
  707. "state": "complete",
  708. "completion_marker": "server-complete-42",
  709. }
  710. assert second.evidence["diff"] == {
  711. "changed": False,
  712. "previous_record_count": 1,
  713. "current_record_count": 1,
  714. "details_materialized": False,
  715. }
  716. assert "removed" not in second.evidence["diff"]
  717. with pytest.raises(ConnectorContractError):
  718. connector.incremental(
  719. _request("incremental", config=config, cursor={"state": "invalid"})
  720. )
  721. def test_runtime_sanitizes_every_output_channel_and_rejects_oversize_records():
  722. class Unsafe(Connector):
  723. manifest = OracleConnector.manifest
  724. def discover(self, _request):
  725. return OperationResult(
  726. records=({"password": "raw", "url": "https://u:p@host/x"},),
  727. cursor={"token": "raw"},
  728. checkpoint={"authorization": "Bearer raw"},
  729. evidence={"url": "https://u:p@host/x?token=raw"},
  730. )
  731. snapshot = incremental = lineage = profile = resume = evidence = discover
  732. def cancel(self, _request):
  733. return OperationResult(status="cancelled")
  734. registry = ConnectorRegistry()
  735. registry.register(Unsafe())
  736. result = ConnectorRuntime(registry).execute("oracle", "1.0.0", _request())
  737. serialized = str(result)
  738. assert "raw" not in serialized and "u:p" not in serialized
  739. assert "[REDACTED]" in serialized
  740. class Oversize(Unsafe):
  741. def discover(self, _request):
  742. return OperationResult(records=({"value": "x" * 1_100_000},))
  743. oversized = ConnectorRegistry()
  744. oversized.register(Oversize())
  745. with pytest.raises(ConnectorConfigurationError):
  746. ConnectorRuntime(oversized, max_attempts=1).execute(
  747. "oracle", "1.0.0", _request(idempotency_key="oversize")
  748. )
  749. def test_rest_pagination_rejects_repeated_and_non_monotonic_cursors():
  750. class Repeating:
  751. def get_json(self, **_values):
  752. return {"assets": [], "next_cursor": {"page": 1}}
  753. connector = RestCatalogConnector(Repeating())
  754. with pytest.raises(ConnectorContractError):
  755. connector.discover(
  756. _request(
  757. config={
  758. "base_url": "https://catalog.example.test",
  759. "allowed_host": "catalog.example.test",
  760. "credential_ref": "env:DATAOPS_CONNECTOR_TEST",
  761. },
  762. cursor={"page": 1},
  763. )
  764. )
  765. def test_sqlserver_tls_defaults_secure_and_insecure_policy_is_development_only(
  766. monkeypatch,
  767. ):
  768. monkeypatch.setattr("importlib.util.find_spec", lambda _name: object())
  769. def definition(tls_options):
  770. return DataSourceDefinition(
  771. uid="u",
  772. name_en="sqlserver",
  773. name_zh="sqlserver",
  774. database_type="sqlserver",
  775. host="db.example",
  776. port=1433,
  777. database="catalog",
  778. credential_ref="ref",
  779. credential_version=1,
  780. tls_options=tls_options,
  781. )
  782. credential = DataSourceCredential("reader", "secret")
  783. query = SqlServerAdapter().build_url(definition({}), credential).query
  784. assert query["Encrypt"] == "yes"
  785. assert query["TrustServerCertificate"] == "no"
  786. with pytest.raises(DataSourceConfigurationInvalid):
  787. SqlServerAdapter().build_url(
  788. definition(
  789. {
  790. "Environment": "production",
  791. "Encrypt": "no",
  792. "TrustServerCertificate": "yes",
  793. }
  794. ),
  795. credential,
  796. )
  797. with pytest.raises(DataSourceConfigurationInvalid):
  798. SqlServerAdapter().build_url(
  799. definition(
  800. {
  801. "Environment": "development",
  802. "Encrypt": "no",
  803. "TrustServerCertificate": "yes",
  804. "AllowInsecureDevelopment": "yes",
  805. }
  806. ),
  807. credential,
  808. )
  809. development = SqlServerAdapter().build_url_for_environment(
  810. definition(
  811. {
  812. "Encrypt": "no",
  813. "TrustServerCertificate": "yes",
  814. }
  815. ),
  816. credential,
  817. trusted_environment="development",
  818. allow_insecure_development=True,
  819. )
  820. assert development.query["Encrypt"] == "no"
  821. def test_runtime_resume_failure_is_failed_retryable_and_preserves_checkpoint():
  822. class ResumeFlaky(Connector):
  823. manifest = OracleConnector.manifest
  824. def __init__(self):
  825. self.resume_calls = 0
  826. def resume(self, request):
  827. self.resume_calls += 1
  828. if self.resume_calls <= 2:
  829. raise ConnectorUpstreamError()
  830. return OperationResult(checkpoint={"page": 8}, evidence={"resumed": True})
  831. discover = snapshot = incremental = lineage = profile = evidence = resume
  832. def cancel(self, request):
  833. return OperationResult(status="cancelled")
  834. class RecordingLimiter:
  835. def __init__(self):
  836. self.keys = []
  837. def acquire(self, key):
  838. self.keys.append(key)
  839. connector, store, limiter = ResumeFlaky(), InMemoryRunStore(), RecordingLimiter()
  840. registry = ConnectorRegistry()
  841. registry.register(connector)
  842. key = "b" * 64
  843. original_checkpoint = {"page": 7, "snapshot": [{"key": "before"}]}
  844. store.claim(
  845. key,
  846. {
  847. "uid": "11111111-1111-4111-8111-111111111112",
  848. "status": "cancelled",
  849. "connector_id": "oracle",
  850. "connector_version": "1.0.0",
  851. "source_uid": _request().source_uid,
  852. "config": dict(_request().config),
  853. "scope": {},
  854. "cursor": {"page": 7},
  855. "checkpoint": original_checkpoint,
  856. "attempt_count": 1,
  857. "dry_run": False,
  858. },
  859. )
  860. runtime = ConnectorRuntime(
  861. registry, store=store, limiter=limiter, max_attempts=2, sleeper=lambda _: None
  862. )
  863. with pytest.raises(ConnectorUpstreamError):
  864. runtime.resume(key)
  865. failed = store.get(key)
  866. assert failed["status"] == "failed" and failed["attempt_count"] == 3
  867. assert failed["error_category"] == "upstream"
  868. assert failed["checkpoint"] == original_checkpoint
  869. result = runtime.resume(key)
  870. assert result.checkpoint == {"page": 8}
  871. assert store.get(key)["attempt_count"] == 4
  872. assert len(limiter.keys) == 2
  873. def test_runtime_never_exceeds_five_total_attempts_and_rate_limit_precedes_claim():
  874. class AlwaysFails(Connector):
  875. manifest = OracleConnector.manifest
  876. def __init__(self):
  877. self.calls = 0
  878. def discover(self, request):
  879. self.calls += 1
  880. raise ConnectorUpstreamError()
  881. snapshot = incremental = lineage = profile = resume = evidence = discover
  882. def cancel(self, request):
  883. return OperationResult(status="cancelled")
  884. connector, store = AlwaysFails(), InMemoryRunStore()
  885. registry = ConnectorRegistry()
  886. registry.register(connector)
  887. runtime = ConnectorRuntime(registry, store=store, sleeper=lambda _: None)
  888. request = _request()
  889. for expected_attempts in (3, MAX_TOTAL_ATTEMPTS):
  890. with pytest.raises(ConnectorUpstreamError):
  891. runtime.execute("oracle", "1.0.0", request)
  892. key = deterministic_idempotency_key("oracle", "1.0.0", request)
  893. assert store.get(key)["attempt_count"] == expected_attempts
  894. with pytest.raises(ConnectorConfigurationError):
  895. runtime.execute("oracle", "1.0.0", request)
  896. assert connector.calls == MAX_TOTAL_ATTEMPTS
  897. class RejectingLimiter:
  898. def acquire(self, _key):
  899. raise ConnectorRateLimitError()
  900. empty_store = InMemoryRunStore()
  901. with pytest.raises(ConnectorRateLimitError):
  902. ConnectorRuntime(
  903. registry, store=empty_store, limiter=RejectingLimiter()
  904. ).execute("oracle", "1.0.0", request)
  905. assert empty_store.list() == []
  906. def test_in_memory_attempt_cas_returns_explicit_single_owner_outcome():
  907. store = InMemoryRunStore()
  908. key = "c" * 64
  909. store.claim(key, {"status": "failed", "attempt_count": 1})
  910. winner = store.update(key, status="running", attempt_count=2)
  911. loser = store.update(key, status="running", attempt_count=2)
  912. assert winner.acquired is True and winner.record["attempt_count"] == 2
  913. assert loser.acquired is False and loser.record["attempt_count"] == 2
  914. def test_idempotency_changes_with_safe_config_and_checkpoint():
  915. base = _request(checkpoint={"page": 1})
  916. changed_checkpoint = _request(checkpoint={"page": 2})
  917. changed_config = _request(
  918. config={"credential_ref": "env:DATAOPS_CONNECTOR_OTHER"}, checkpoint={"page": 1}
  919. )
  920. keys = {
  921. deterministic_idempotency_key("oracle", "1.0.0", item)
  922. for item in (base, changed_checkpoint, changed_config)
  923. }
  924. assert len(keys) == 3
  925. def test_rate_limit_and_error_taxonomy_are_deterministic():
  926. now = [0.0]
  927. limiter = SlidingWindowLimiter(limit=2, window_seconds=10, clock=lambda: now[0])
  928. limiter.acquire("source")
  929. limiter.acquire("source")
  930. with pytest.raises(ConnectorRateLimitError):
  931. limiter.acquire("source")
  932. now[0] = 11
  933. limiter.acquire("source")
  934. assert isinstance(classify_error(TimeoutError()), ConnectorTimeoutError)
  935. assert isinstance(classify_error(PermissionError()), ConnectorPermissionError)
  936. assert isinstance(
  937. classify_error(ConnectorDriverUnavailableError()),
  938. ConnectorDriverUnavailableError,
  939. )
  940. class _Result:
  941. def __init__(self, scalar=None, mapping=None, rowcount=1):
  942. self.scalar = scalar
  943. self.mapping = mapping
  944. self.rowcount = rowcount
  945. def scalar_one_or_none(self):
  946. return self.scalar
  947. def mappings(self):
  948. return self
  949. def one_or_none(self):
  950. return self.mapping
  951. class _Session:
  952. def __init__(self, results):
  953. self.results = list(results)
  954. self.calls = []
  955. self.commits = 0
  956. self.rollbacks = 0
  957. def execute(self, statement, parameters=None):
  958. self.calls.append((str(statement), parameters or {}))
  959. return self.results.pop(0) if self.results else _Result()
  960. def commit(self):
  961. self.commits += 1
  962. def rollback(self):
  963. self.rollbacks += 1
  964. def test_machine_credential_ttl_one_time_return_and_no_plaintext_persistence():
  965. session = _Session([_Result(), _Result("active"), _Result(), _Result()])
  966. repository = ConnectorIdentityRepository(session)
  967. issued = repository.issue(
  968. "11111111-1111-4111-8111-111111111111",
  969. ttl_seconds=900,
  970. actor_uid="22222222-2222-4222-8222-222222222222",
  971. )
  972. assert issued["credential"].startswith("dopc_") and issued["returned_once"] is True
  973. persisted = [
  974. params
  975. for sql, params in session.calls
  976. if "connector_machine_credentials" in sql and "INSERT" in sql
  977. ][0]
  978. assert "credential" not in persisted and len(persisted["hash"]) == 64
  979. with pytest.raises(ConnectorConfigurationError):
  980. repository.issue("x", ttl_seconds=901, actor_uid="y")
  981. def test_machine_credential_scope_replay_revoke_and_rotation():
  982. replay_identity = {
  983. "uid": "c",
  984. "principal_uid": "p",
  985. "use_count": 1,
  986. "connector_id": "oracle",
  987. "connector_version": "1.0.0",
  988. "source_uid": "s",
  989. "business_domain_uid": "d",
  990. "environment": "staging",
  991. "allowed_operations": ["discover"],
  992. "allowed_scopes": {},
  993. "source_binding_uid": "binding",
  994. "source_binding_version": 1,
  995. "approved_config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  996. }
  997. session = _Session(
  998. [_Result(mapping=None), _Result(mapping=replay_identity), _Result(), _Result()]
  999. )
  1000. with pytest.raises(ConnectorAuthenticationError):
  1001. ConnectorIdentityRepository(session).authenticate(
  1002. "token",
  1003. connector_id="oracle",
  1004. connector_version="1.0.0",
  1005. source_uid="s",
  1006. business_domain_uid="d",
  1007. environment="staging",
  1008. operation="discover",
  1009. scope={},
  1010. )
  1011. assert any(
  1012. "credential_replay_rejected" in params.values()
  1013. for _sql, params in session.calls
  1014. )
  1015. scope_identity = {
  1016. **replay_identity,
  1017. "use_count": 0,
  1018. "allowed_scopes": {"include_schemas": ["APP"]},
  1019. }
  1020. session = _Session(
  1021. [_Result(mapping=None), _Result(mapping=scope_identity), _Result()]
  1022. )
  1023. with pytest.raises(ConnectorPermissionError):
  1024. ConnectorIdentityRepository(session).authenticate(
  1025. "token",
  1026. connector_id="oracle",
  1027. connector_version="1.0.0",
  1028. source_uid="s",
  1029. business_domain_uid="d",
  1030. environment="staging",
  1031. operation="discover",
  1032. scope={"include_schemas": ["SYS"]},
  1033. )
  1034. session = _Session([_Result("principal"), _Result()])
  1035. assert ConnectorIdentityRepository(session).revoke("credential", "actor") is True
  1036. rotation = _Session(
  1037. [
  1038. _Result("principal"),
  1039. _Result(),
  1040. _Result(),
  1041. _Result("active"),
  1042. _Result(),
  1043. _Result(),
  1044. ]
  1045. )
  1046. rotated = ConnectorIdentityRepository(rotation).rotate(
  1047. "credential", ttl_seconds=300, actor_uid="actor"
  1048. )
  1049. assert rotated["returned_once"] is True and rotated["expires_in"] == 300
  1050. def test_machine_credential_expiry_is_rejected_and_audited():
  1051. expired = {"uid": "credential", "principal_uid": "principal"}
  1052. session = _Session([_Result(mapping=expired), _Result(), _Result()])
  1053. with pytest.raises(ConnectorAuthenticationError):
  1054. ConnectorIdentityRepository(session).authenticate(
  1055. "expired",
  1056. connector_id="oracle",
  1057. connector_version="1.0.0",
  1058. source_uid="source",
  1059. business_domain_uid="domain",
  1060. environment="staging",
  1061. operation="discover",
  1062. scope={},
  1063. )
  1064. assert any(
  1065. "credential_expired_rejected" in params.values()
  1066. for _sql, params in session.calls
  1067. )
  1068. def test_machine_credential_success_is_scoped_and_audited():
  1069. identity = {
  1070. "uid": "credential",
  1071. "principal_uid": "principal",
  1072. "use_count": 0,
  1073. "connector_id": "oracle",
  1074. "connector_version": "1.0.0",
  1075. "source_uid": "source",
  1076. "business_domain_uid": "domain",
  1077. "environment": "staging",
  1078. "allowed_operations": ["discover"],
  1079. "allowed_scopes": {"include_schemas": ["APP"]},
  1080. "source_binding_uid": "binding",
  1081. "source_binding_version": 1,
  1082. "approved_config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1083. }
  1084. session = _Session([_Result(mapping=None), _Result(mapping=identity), _Result()])
  1085. result = ConnectorIdentityRepository(session).authenticate(
  1086. "token",
  1087. connector_id="oracle",
  1088. connector_version="1.0.0",
  1089. source_uid="source",
  1090. business_domain_uid="domain",
  1091. environment="staging",
  1092. operation="discover",
  1093. scope={"include_schemas": ["APP"]},
  1094. )
  1095. assert result["principal_uid"] == "principal" and session.commits == 1
  1096. assert any(
  1097. "credential_authenticated" in params.values() for _sql, params in session.calls
  1098. )
  1099. @pytest.mark.parametrize(
  1100. "override",
  1101. [
  1102. {"connector_version": "2.0.0"},
  1103. {"business_domain_uid": "other-domain"},
  1104. {"environment": "production"},
  1105. ],
  1106. )
  1107. def test_machine_credential_rejects_cross_binding(override):
  1108. identity = {
  1109. "uid": "credential",
  1110. "principal_uid": "principal",
  1111. "use_count": 0,
  1112. "connector_id": "oracle",
  1113. "connector_version": "1.0.0",
  1114. "source_uid": "source",
  1115. "business_domain_uid": "domain",
  1116. "environment": "staging",
  1117. "allowed_operations": ["discover"],
  1118. "allowed_scopes": {},
  1119. "source_binding_uid": "binding",
  1120. "source_binding_version": 1,
  1121. "approved_config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
  1122. }
  1123. session = _Session([_Result(mapping=None), _Result(mapping=identity), _Result()])
  1124. binding = {
  1125. "connector_id": "oracle",
  1126. "connector_version": "1.0.0",
  1127. "source_uid": "source",
  1128. "business_domain_uid": "domain",
  1129. "environment": "staging",
  1130. "operation": "discover",
  1131. "scope": {},
  1132. **override,
  1133. }
  1134. with pytest.raises(ConnectorPermissionError):
  1135. ConnectorIdentityRepository(session).authenticate("token", **binding)
  1136. def test_machine_scope_rejects_unknown_key_even_when_empty():
  1137. session = _Session([])
  1138. with pytest.raises(ConnectorConfigurationError):
  1139. ConnectorIdentityRepository(session).authenticate(
  1140. "token",
  1141. connector_id="oracle",
  1142. connector_version="1.0.0",
  1143. source_uid="source",
  1144. business_domain_uid="domain",
  1145. environment="staging",
  1146. operation="discover",
  1147. scope={"unknown": []},
  1148. )
  1149. assert session.calls == []
  1150. def test_api_permissions_and_machine_boundary_fail_closed_without_credential():
  1151. from app import create_app
  1152. from app.core.system.permissions import permission_for_request
  1153. assert permission_for_request("/api/datasource/connectors/manifests", "GET") == (
  1154. "connectors:read",
  1155. )
  1156. assert permission_for_request("/api/datasource/connectors/runs", "POST") == (
  1157. "connectors:operate",
  1158. )
  1159. assert permission_for_request("/api/datasource/connectors/principals", "POST") == (
  1160. "connectors:manage",
  1161. )
  1162. assert permission_for_request(
  1163. "/api/datasource/connectors/machine/runs", "POST"
  1164. ) == ("public",)
  1165. app = create_app()
  1166. app.config.update(TESTING=True)
  1167. response = app.test_client().post(
  1168. "/api/datasource/connectors/machine/runs", json={}
  1169. )
  1170. assert (
  1171. response.status_code == 401
  1172. and "credential" not in response.get_data(as_text=True).lower()
  1173. )
  1174. def test_migration_permissions_frontend_docs_and_release_parity_contract():
  1175. migration = (
  1176. ROOT / "migrations/versions/20260802_472_enterprise_connectors.py"
  1177. ).read_text()
  1178. for table in (
  1179. "connector_manifests",
  1180. "connector_principals",
  1181. "connector_machine_credentials",
  1182. "connector_runs",
  1183. "connector_run_attempts",
  1184. "connector_checkpoints",
  1185. "connector_evidence",
  1186. "connector_graph_edges",
  1187. "connector_audit_events",
  1188. ):
  1189. assert f"CREATE TABLE public.{table}" in migration
  1190. for marker in (
  1191. "INTERVAL '15 minutes'",
  1192. "idempotency_key CHAR(64) NOT NULL UNIQUE",
  1193. "use_count",
  1194. "replayed",
  1195. "DROP TABLE IF EXISTS",
  1196. ):
  1197. assert marker in migration
  1198. hardening = (
  1199. ROOT / "migrations/versions/20260802_473_connector_runtime_hardening.py"
  1200. ).read_text()
  1201. assert 'down_revision = "20260802_472"' in hardening
  1202. for marker in (
  1203. "connector_version VARCHAR(40)",
  1204. "connector_rate_limits",
  1205. "safe_config JSONB",
  1206. "business_domain_uid",
  1207. "process_key",
  1208. "ck_connector_machine_run_binding",
  1209. ):
  1210. assert marker in hardening
  1211. guardrails = (
  1212. ROOT / "migrations/versions/20260802_474_connector_version_guardrails.py"
  1213. ).read_text()
  1214. assert 'down_revision = "20260802_473"' in guardrails
  1215. assert "multiple active manifest versions" in guardrails
  1216. assert "multiple principal versions" in guardrails
  1217. security_bindings = (
  1218. ROOT / "migrations/versions/20260802_475_connector_security_bindings.py"
  1219. ).read_text()
  1220. assert 'down_revision = "20260802_474"' in security_bindings
  1221. for marker in (
  1222. "connector_source_bindings",
  1223. "request_hash",
  1224. "client_hint_hash",
  1225. "attempt_lease_token",
  1226. "cancel_requested",
  1227. "explicitly remove source bindings",
  1228. ):
  1229. assert marker in security_bindings
  1230. binding_enforcement = (
  1231. ROOT / "migrations/versions/20260802_476_connector_binding_enforcement.py"
  1232. ).read_text()
  1233. assert 'down_revision = "20260802_475"' in binding_enforcement
  1234. assert "bind every active enterprise principal first" in binding_enforcement
  1235. assert (
  1236. "process_key = db.Column(db.String(300))"
  1237. in (ROOT / "app/models/connectors.py").read_text()
  1238. )
  1239. permissions = (ROOT / "app/core/system/permissions.py").read_text()
  1240. for permission in ("connectors:read", "connectors:operate", "connectors:manage"):
  1241. assert permission in permissions
  1242. api = (ROOT / "app/api/data_source/routes.py").read_text()
  1243. for endpoint in (
  1244. "/connectors/manifests",
  1245. "/connectors/config/validate",
  1246. "/compatibility",
  1247. "/connectors/runs",
  1248. "/cancel",
  1249. "/resume",
  1250. "/connectors/principals",
  1251. "/credentials",
  1252. ):
  1253. assert endpoint in api
  1254. assert "DATASOURCE_GRAPH_NOT_IMPLEMENTED" not in api
  1255. page = (
  1256. ROOT / "frontend/src/views/dataGovernance/development/enterpriseConnectors.vue"
  1257. ).read_text()
  1258. connector_operations = (
  1259. ROOT / "frontend/src/components/connectors/ConnectorOperations.vue"
  1260. ).read_text()
  1261. assert "真实企业账号、版本、网络与 UAT 尚未提供" in page
  1262. assert "{{ item.credential" not in page and "{{ credential" not in page
  1263. assert "$store.state.user.userInfo.permissions" in connector_operations
  1264. assert "$store.getters.roles" not in connector_operations
  1265. for permission in ("connectors:read", "connectors:operate", "connectors:manage"):
  1266. assert permission in connector_operations
  1267. assert "isHumanDryRun(item)" in connector_operations and 'v-if="canManage"' in connector_operations
  1268. assert "checkpoint_summary" in connector_operations and "cursor_summary" in connector_operations
  1269. assert "仅机器凭证可操作" in connector_operations
  1270. assert "human connector runs must use dry_run=true" in api
  1271. connector_files = [
  1272. path.relative_to(ROOT).as_posix()
  1273. for path in (ROOT / "app/core/connectors").rglob("*.py")
  1274. ]
  1275. for relative in (
  1276. *connector_files,
  1277. "app/core/data_source/adapters/__init__.py",
  1278. "app/core/data_source/adapters/oracle.py",
  1279. "app/core/data_source/adapters/sqlserver.py",
  1280. "app/core/data_source/errors.py",
  1281. "app/core/data_source/service.py",
  1282. "app/models/connectors.py",
  1283. "app/models/__init__.py",
  1284. "app/api/data_source/routes.py",
  1285. "app/core/system/permissions.py",
  1286. ):
  1287. assert (ROOT / relative).read_bytes() == (
  1288. ROOT / "deployment" / relative
  1289. ).read_bytes()
  1290. assert (
  1291. ROOT / "migrations/versions/20260802_472_enterprise_connectors.py"
  1292. ).read_bytes() == (
  1293. ROOT / "deployment/migrations/versions/20260802_472_enterprise_connectors.py"
  1294. ).read_bytes()
  1295. def test_openapi_connector_contract_is_strict_and_matches_behavior():
  1296. contract = yaml.safe_load((ROOT / "docs/architecture/OPENAPI.yaml").read_text())
  1297. paths = contract["paths"]
  1298. machine = paths["/api/datasource/connectors/machine/runs"]["post"]
  1299. assert machine["requestBody"]["required"] is True
  1300. assert machine["requestBody"]["content"]["application/json"]["schema"][
  1301. "$ref"
  1302. ].endswith("MachineConnectorRunRequest")
  1303. header = next(
  1304. item
  1305. for item in machine["parameters"]
  1306. if item["name"] == "X-Connector-Credential"
  1307. )
  1308. assert header["in"] == "header" and header["required"] is True
  1309. assert set(machine["responses"]) == {"201", "default"}
  1310. assert machine["responses"]["201"]["content"]["application/json"]["schema"][
  1311. "$ref"
  1312. ].endswith("ConnectorRunResultEnvelope")
  1313. assert machine["x-required-permission"] == "machine-credential"
  1314. for action in ("cancel", "resume"):
  1315. operation = paths[
  1316. f"/api/datasource/connectors/runs/{{idempotency_key}}/{action}"
  1317. ]["post"]
  1318. assert "requestBody" not in operation
  1319. human = paths["/api/datasource/connectors/runs"]["post"]
  1320. assert (
  1321. human["responses"].get("201")
  1322. and human["x-required-permission"] == "connectors:operate"
  1323. )
  1324. schemas = contract["components"]["schemas"]
  1325. assert schemas["ConnectorRunRequest"]["properties"]["dry_run"] == {"const": True}
  1326. assert schemas["ConnectorScope"]["additionalProperties"] is False
  1327. assert schemas["ConnectorRunResultEnvelope"]["additionalProperties"] is False
  1328. assert "config" not in schemas["MachineConnectorRunRequest"]["properties"]
  1329. assert schemas["ConnectorOperationResult"]["properties"]["status"]["enum"] == [
  1330. "succeeded",
  1331. "dry_run",
  1332. ]
  1333. assert "allOf" not in schemas["ConnectorHealthEnvelope"]
  1334. categories = schemas["ConnectorErrorEnvelope"]["properties"]["error"]["properties"][
  1335. "category"
  1336. ]["enum"]
  1337. assert {
  1338. "configuration",
  1339. "authentication",
  1340. "permission",
  1341. "conflict",
  1342. "rate_limit",
  1343. "timeout",
  1344. "upstream",
  1345. "contract",
  1346. "cancelled",
  1347. } == set(categories)
  1348. for action in ("cancel", "resume"):
  1349. machine_action = paths[
  1350. f"/api/datasource/connectors/machine/runs/{{idempotency_key}}/{action}"
  1351. ]["post"]
  1352. assert "requestBody" not in machine_action
  1353. assert machine_action["x-required-permission"] == "machine-credential"
  1354. credential_action = paths[
  1355. "/api/datasource/connectors/credentials/{credential_uid}/{action}"
  1356. ]["post"]
  1357. assert credential_action["requestBody"]["required"] is False
  1358. credential_issue = paths[
  1359. "/api/datasource/connectors/principals/{principal_uid}/credentials"
  1360. ]["post"]
  1361. assert credential_issue["requestBody"]["required"] is False
  1362. binding_create = paths["/api/datasource/connectors/source-bindings"]["post"]
  1363. assert set(binding_create["responses"]) == {"201", "default"}
  1364. assert binding_create["responses"]["201"]["content"]["application/json"][
  1365. "schema"
  1366. ]["$ref"].endswith("ConnectorSourceBindingEnvelope")
  1367. action_parameter = next(
  1368. item for item in credential_action["parameters"] if item["name"] == "action"
  1369. )
  1370. assert action_parameter["schema"]["enum"] == ["revoke", "rotate"]
  1371. assert paths["/api/datasource/graph"]["post"]["x-required-permission"] == (
  1372. "connectors:read"
  1373. )
  1374. for envelope in (
  1375. "ConnectorRunResultEnvelope",
  1376. "ConnectorErrorEnvelope",
  1377. "ConnectorHealthEnvelope",
  1378. ):
  1379. properties = schemas[envelope]["properties"]
  1380. assert "success" not in properties and "timestamp" not in properties
  1381. assert {"code", "message", "data"}.issubset(properties)
  1382. assert (
  1383. ROOT / "migrations/versions/20260802_473_connector_runtime_hardening.py"
  1384. ).read_bytes() == (
  1385. ROOT
  1386. / "deployment/migrations/versions/20260802_473_connector_runtime_hardening.py"
  1387. ).read_bytes()
  1388. assert (
  1389. ROOT / "migrations/versions/20260802_474_connector_version_guardrails.py"
  1390. ).read_bytes() == (
  1391. ROOT
  1392. / "deployment/migrations/versions/20260802_474_connector_version_guardrails.py"
  1393. ).read_bytes()
  1394. assert (
  1395. ROOT / "migrations/versions/20260802_475_connector_security_bindings.py"
  1396. ).read_bytes() == (
  1397. ROOT
  1398. / "deployment/migrations/versions/20260802_475_connector_security_bindings.py"
  1399. ).read_bytes()
  1400. assert (
  1401. ROOT / "migrations/versions/20260802_476_connector_binding_enforcement.py"
  1402. ).read_bytes() == (
  1403. ROOT
  1404. / "deployment/migrations/versions/20260802_476_connector_binding_enforcement.py"
  1405. ).read_bytes()
  1406. def test_flask_machine_responses_validate_against_openapi_201_and_401(
  1407. monkeypatch,
  1408. ):
  1409. from app import create_app
  1410. from app.core.connectors.identity import ConnectorIdentityRepository
  1411. from app.core.connectors.repository import ConnectorRepository
  1412. from app.core.connectors.runtime import ConnectorRuntime
  1413. contract = yaml.safe_load((ROOT / "docs/architecture/OPENAPI.yaml").read_text())
  1414. def expand(schema):
  1415. if isinstance(schema, dict) and "$ref" in schema:
  1416. target = contract
  1417. for part in schema["$ref"].removeprefix("#/").split("/"):
  1418. target = target[part]
  1419. return expand(target)
  1420. if isinstance(schema, dict):
  1421. return {key: expand(value) for key, value in schema.items()}
  1422. if isinstance(schema, list):
  1423. return [expand(value) for value in schema]
  1424. return schema
  1425. monkeypatch.setattr(
  1426. ConnectorIdentityRepository,
  1427. "authenticate",
  1428. lambda _self, _token, **_binding: {
  1429. "principal_uid": "principal",
  1430. "source_binding_uid": "binding",
  1431. "source_binding_version": 1,
  1432. "approved_config": {
  1433. "credential_ref": "env:DATAOPS_CONNECTOR_TEST"
  1434. },
  1435. },
  1436. )
  1437. monkeypatch.setattr(
  1438. ConnectorRepository,
  1439. "register_manifest",
  1440. lambda _self, _manifest, _actor: None,
  1441. )
  1442. monkeypatch.setattr(
  1443. ConnectorRuntime,
  1444. "execute",
  1445. lambda _self, _connector, _version, _request: OperationResult(
  1446. records=({"asset_key": "source:APP.ORDERS"},),
  1447. checkpoint={"page": 1},
  1448. cursor={"page": 2},
  1449. evidence={"safe": True},
  1450. ),
  1451. )
  1452. app = create_app()
  1453. app.config.update(TESTING=True)
  1454. app.extensions["connector_registry"] = SimpleNamespace(
  1455. resolve=lambda *_args, **_kwargs: SimpleNamespace(
  1456. manifest=OracleConnector.manifest
  1457. )
  1458. )
  1459. client = app.test_client()
  1460. path = "/api/datasource/connectors/machine/runs"
  1461. unauthorized = client.post(path, json={})
  1462. assert unauthorized.status_code == 401
  1463. error_schema = expand(
  1464. contract["paths"][path]["post"]["responses"]["default"]["content"][
  1465. "application/json"
  1466. ]["schema"]
  1467. )
  1468. Draft202012Validator(error_schema).validate(unauthorized.get_json())
  1469. assert set(unauthorized.get_json()) == {"code", "message", "data", "error"}
  1470. assert unauthorized.get_json()["error"] == {
  1471. "code": "CONNECTOR_ERROR",
  1472. "category": "authentication",
  1473. "retryable": False,
  1474. }
  1475. created = client.post(
  1476. path,
  1477. headers={"X-Connector-Credential": "one-time-token"},
  1478. json={
  1479. "connector_id": "oracle",
  1480. "version": "1.0.0",
  1481. "source_uid": "11111111-1111-4111-8111-111111111111",
  1482. "business_domain_uid": "22222222-2222-4222-8222-222222222222",
  1483. "environment": "staging",
  1484. "process_key": "catalog-sync",
  1485. "operation": "discover",
  1486. "scope": {},
  1487. },
  1488. )
  1489. assert created.status_code == 201
  1490. success_schema = expand(
  1491. contract["paths"][path]["post"]["responses"]["201"]["content"][
  1492. "application/json"
  1493. ]["schema"]
  1494. )
  1495. Draft202012Validator(success_schema).validate(created.get_json())
  1496. assert set(created.get_json()) == {"code", "message", "data"}
  1497. assert created.get_json()["code"] == created.status_code == 201
  1498. invalid = client.post(
  1499. path,
  1500. headers={"X-Connector-Credential": "one-time-token"},
  1501. json={
  1502. "connector_id": "oracle",
  1503. "version": "1.0.0",
  1504. "source_uid": "11111111-1111-4111-8111-111111111111",
  1505. "business_domain_uid": "22222222-2222-4222-8222-222222222222",
  1506. "environment": "staging",
  1507. "process_key": "catalog-sync",
  1508. "operation": "cancel",
  1509. "scope": {},
  1510. },
  1511. )
  1512. assert invalid.status_code == invalid.get_json()["code"] == 400
  1513. Draft202012Validator(error_schema).validate(invalid.get_json())
  1514. assert invalid.get_json()["error"]["category"] == "configuration"
  1515. def test_human_run_requires_explicit_dry_run():
  1516. from app.api.data_source.routes import (
  1517. _require_collection_operation,
  1518. _require_human_dry_run,
  1519. )
  1520. _require_human_dry_run({"dry_run": True})
  1521. with pytest.raises(ConnectorConfigurationError):
  1522. _require_human_dry_run({"dry_run": False})
  1523. with pytest.raises(ConnectorConfigurationError):
  1524. _require_human_dry_run({})
  1525. _require_collection_operation({"operation": "incremental"})
  1526. with pytest.raises(ConnectorConfigurationError):
  1527. _require_collection_operation({"operation": "cancel"})