| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697 |
- from __future__ import annotations
- from contextlib import contextmanager
- from pathlib import Path
- from types import SimpleNamespace
- import pytest
- import yaml
- from jsonschema import Draft202012Validator
- from app.core.connectors import (
- SDK_VERSION,
- Connector,
- ConnectorConfigurationError,
- ConnectorManifest,
- ConnectorRegistry,
- ConnectorUpstreamError,
- EnvironmentSecretResolver,
- OperationRequest,
- OperationResult,
- validate_config,
- )
- from app.core.connectors.builtin import (
- OracleConnector,
- RestCatalogConnector,
- SafeRestTransport,
- SqlServerConnector,
- )
- from app.core.connectors.errors import (
- ConnectorAuthenticationError,
- ConnectorContractError,
- ConnectorDriverUnavailableError,
- ConnectorPermissionError,
- ConnectorRateLimitError,
- ConnectorTimeoutError,
- classify_error,
- )
- from app.core.connectors.identity import ConnectorIdentityRepository
- from app.core.connectors.runtime import (
- MAX_TOTAL_ATTEMPTS,
- ConnectorRuntime,
- InMemoryRunStore,
- SlidingWindowLimiter,
- deterministic_idempotency_key,
- redact_evidence,
- snapshot_diff,
- )
- from app.core.data_source.adapters.oracle import OracleAdapter
- from app.core.data_source.adapters.sqlserver import SqlServerAdapter
- from app.core.data_source.errors import (
- DataSourceConfigurationInvalid,
- DataSourceDriverUnavailable,
- )
- from app.core.data_source.models import DataSourceCredential, DataSourceDefinition
- ROOT = Path(__file__).resolve().parents[1]
- ALL_CAPABILITIES = (
- "discover",
- "snapshot",
- "incremental",
- "lineage",
- "profile",
- "cancel",
- "resume",
- "evidence",
- )
- class _Rows:
- def __init__(self, rows):
- self.rows = rows
- def mappings(self):
- return self
- def all(self):
- return self.rows
- class _Connection:
- def __init__(self, rows):
- self.rows = rows
- self.statement = None
- self.parameters = None
- def execute(self, statement, parameters):
- self.statement = str(statement)
- self.parameters = parameters
- return _Rows(self.rows)
- @contextmanager
- def _provider(_source, purpose):
- assert purpose == "metadata_collection"
- yield _Connection(
- [
- {
- "schema_name": "APP",
- "asset_name": "ORDERS",
- "asset_type": "TABLE",
- "column_name": "ID",
- "ordinal_position": 1,
- "data_type": "NUMBER",
- "is_nullable": "NO",
- "column_default": None,
- "column_comment": "key",
- }
- ]
- )
- def _request(operation="discover", **values):
- defaults = {
- "source_uid": "11111111-1111-4111-8111-111111111111",
- "operation": operation,
- "config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- "scope": {},
- }
- defaults.update(values)
- return OperationRequest(**defaults)
- def test_manifest_config_and_secret_boundary_are_strict():
- manifest = OracleConnector.manifest
- assert manifest.sdk_version == SDK_VERSION
- assert set(manifest.capabilities) == {
- "discover",
- "snapshot",
- "incremental",
- "cancel",
- "resume",
- "evidence",
- }
- registry = ConnectorRegistry()
- registry.register(OracleConnector(_provider))
- with pytest.raises(ConnectorConfigurationError):
- registry.resolve("oracle", "1.0.0", "lineage")
- assert registry.validate("oracle", "1.0.0", {"credential_ref": "vault:team/oracle"})
- with pytest.raises(ConnectorConfigurationError):
- registry.validate("oracle", "1.0.0", {"password": "plain"})
- with pytest.raises(ConnectorConfigurationError):
- registry.validate("oracle", "9.9.9", {"credential_ref": "env:X"})
- with pytest.raises(ConnectorConfigurationError):
- ConnectorManifest(
- "bad",
- "1",
- SDK_VERSION,
- "Bad",
- ("discover",),
- {"type": "object", "additionalProperties": True, "properties": {}},
- )
- def test_draft_202012_schema_recursively_enforces_nested_arrays_bounds_and_one_of():
- schema = {
- "$schema": "https://json-schema.org/draft/2020-12/schema",
- "type": "object",
- "additionalProperties": False,
- "required": ["nested"],
- "properties": {
- "nested": {
- "type": "object",
- "additionalProperties": False,
- "required": ["items", "mode"],
- "properties": {
- "items": {
- "type": "array",
- "minItems": 1,
- "maxItems": 2,
- "items": {"type": "integer", "minimum": 1, "maximum": 5},
- },
- "mode": {"oneOf": [{"const": "a"}, {"const": "b"}]},
- },
- }
- },
- }
- assert validate_config(schema, {"nested": {"items": [1, 5], "mode": "a"}})
- for invalid in (
- {"nested": {"items": [], "mode": "a"}},
- {"nested": {"items": [6], "mode": "a"}},
- {"nested": {"items": [1], "mode": "c"}},
- {"nested": {"items": [1], "mode": "a", "extra": True}},
- ):
- with pytest.raises(ConnectorConfigurationError):
- validate_config(schema, invalid)
- with pytest.raises(ConnectorConfigurationError):
- ConnectorManifest(
- "nested-secret",
- "1.0.0",
- SDK_VERSION,
- "Nested Secret",
- ("discover",),
- {
- "type": "object",
- "additionalProperties": False,
- "properties": {
- "nested": {
- "type": "object",
- "properties": {"password": {"type": "string"}},
- }
- },
- },
- )
- @pytest.mark.parametrize("operation", ALL_CAPABILITIES)
- def test_all_unified_interfaces_are_callable_and_undeclared_capability_is_denied(
- operation,
- ):
- connector = OracleConnector(_provider)
- result = getattr(connector, operation)(_request(operation))
- assert isinstance(result, OperationResult)
- class DiscoverOnly(Connector):
- manifest = ConnectorManifest(
- "discover-only",
- "1.0.0",
- SDK_VERSION,
- "Discover",
- ("discover",),
- {"type": "object", "additionalProperties": False, "properties": {}},
- )
- def discover(self, request):
- return OperationResult()
- snapshot = incremental = lineage = profile = cancel = resume = evidence = (
- discover
- )
- registry = ConnectorRegistry()
- registry.register(DiscoverOnly())
- with pytest.raises(ConnectorConfigurationError):
- registry.resolve("discover-only", "1.0.0", "snapshot")
- @pytest.mark.parametrize(
- "connector_cls, sql_marker",
- [(OracleConnector, "all_tab_columns"), (SqlServerConnector, "sys.columns")],
- )
- def test_database_connectors_normalize_scope_and_use_read_only_catalog_sql(
- connector_cls, sql_marker
- ):
- connector = connector_cls(_provider)
- result = connector.discover(
- _request(scope={"include_schemas": ["APP"], "include_tables": ["ORDERS"]})
- )
- assert result.records[0]["asset_key"].endswith(":APP.ORDERS")
- assert result.records[0]["nullable"] is False
- assert result.evidence["query_kind"] == "read_only_metadata"
- assert sql_marker in connector.catalog_sql
- assert not any(
- word in connector.catalog_sql.upper()
- for word in ("INSERT ", "UPDATE ", "DELETE ")
- )
- def test_database_incremental_persists_bounded_snapshot_summary():
- result = OracleConnector(_provider).incremental(
- _request(
- "incremental",
- checkpoint={"snapshot_summary": {"sha256": "old", "record_count": 2}},
- )
- )
- diff = result.evidence["diff"]
- assert diff == {
- "changed": True,
- "previous_record_count": 2,
- "current_record_count": 1,
- "details_materialized": False,
- }
- assert "snapshot" not in result.checkpoint
- assert result.checkpoint["snapshot_summary"]["record_count"] == 1
- class _HttpResponse:
- def __init__(self, status, body, content_type="application/json"):
- self.status = status
- self.body = body
- self.headers = {"Content-Type": content_type}
- self.released = False
- def stream(self, **_kwargs):
- for offset in range(0, len(self.body), 65536):
- yield self.body[offset : offset + 65536]
- def release_conn(self):
- self.released = True
- class _HttpsPool:
- def __init__(self, calls, response, **kwargs):
- self.calls = calls
- self.response = response
- self.calls.append(("pool", kwargs))
- def urlopen(self, *args, **kwargs):
- self.calls.append(("request", {"args": args, **kwargs}))
- return self.response.pop(0) if isinstance(self.response, list) else self.response
- def close(self):
- self.calls.append(("close", {}))
- def _safe_transport(calls, response, resolver=None, secret_resolver=None):
- return SafeRestTransport(
- secret_resolver=secret_resolver or (lambda ref: "resolved-token"),
- resolver=resolver
- or (lambda *_args, **_kwargs: [(2, 1, 6, "", ("8.8.8.8", 443))]),
- pool_factory=lambda **kwargs: _HttpsPool(calls, response, **kwargs),
- )
- def test_reference_rest_connector_uses_same_registry_without_core_dispatcher_changes():
- calls, resolutions = [], []
- def resolver(host, port, type):
- return [(2, 1, 6, "", ("8.8.8.8", port))]
- def secret_resolver(reference):
- resolutions.append(reference)
- return "resolved-token"
- registry = ConnectorRegistry()
- page = b'{"assets":[{"key":"a","name":"Orders","namespace":"sales","type":"table"}],"next_cursor":{"page":2}}'
- end = b'{"assets":[]}'
- response = [
- _HttpResponse(200, page),
- _HttpResponse(200, end),
- _HttpResponse(200, page),
- _HttpResponse(200, end),
- ]
- registry.register(
- RestCatalogConnector(
- _safe_transport(calls, response, resolver, secret_resolver)
- )
- )
- result = ConnectorRuntime(registry, sleeper=lambda _: None).execute(
- "rest-catalog",
- "1.0.0",
- _request(
- config={
- "base_url": "https://catalog.example.test",
- "allowed_host": "catalog.example.test",
- "credential_ref": "secret:catalog/key",
- }
- ),
- )
- assert result.records[0]["name"] == "Orders"
- pool_call = next(value for kind, value in calls if kind == "pool")
- request_call = next(value for kind, value in calls if kind == "request")
- assert pool_call["host"] == "8.8.8.8"
- assert (
- pool_call["server_hostname"]
- == pool_call["assert_hostname"]
- == "catalog.example.test"
- )
- assert pool_call["cert_reqs"] == "CERT_REQUIRED" and pool_call["retries"] is False
- assert request_call["headers"]["Host"] == "catalog.example.test"
- assert request_call["headers"]["Authorization"] == "Bearer resolved-token"
- assert request_call["redirect"] is False and request_call["retries"] is False
- assert resolutions == ["secret:catalog/key", "secret:catalog/key"]
- assert "resolved-token" not in str(result.evidence)
- incremental = ConnectorRuntime(registry, sleeper=lambda _: None).execute(
- "rest-catalog",
- "1.0.0",
- _request(
- "incremental",
- config={
- "base_url": "https://catalog.example.test",
- "allowed_host": "catalog.example.test",
- "credential_ref": "secret:catalog/key",
- },
- checkpoint={
- "cursor": {"page": 1},
- "snapshot_summary": {"sha256": "old", "record_count": 1},
- },
- ),
- )
- incremental_request = [value for kind, value in calls if kind == "request"][-1]
- assert "cursor=" in incremental_request["args"][1]
- assert incremental.cursor["state"] == "complete"
- assert incremental.cursor["completion_marker"]
- assert incremental.evidence["diff"]["changed"] is True
- assert incremental.evidence["diff"]["details_materialized"] is False
- assert resolutions == ["secret:catalog/key"] * 4
- with pytest.raises(ConnectorConfigurationError):
- _safe_transport(
- calls,
- response,
- lambda *_args, **_kwargs: [(2, 1, 6, "", ("127.0.0.1", 443))],
- ).get_json(
- url="https://localhost/x",
- allowed_host="localhost",
- credential_ref="secret:test/key",
- )
- core = (ROOT / "app/core/connectors/registry.py").read_text() + (
- ROOT / "app/core/connectors/runtime.py"
- ).read_text()
- assert all(
- connector_id not in core
- for connector_id in ("oracle", "sqlserver", "rest-catalog")
- )
- @pytest.mark.parametrize(
- "status,body,error",
- [
- (302, b"{}", ConnectorUpstreamError),
- (200, b"x" * (2 * 1024 * 1024 + 1), Exception),
- (200, b"not-json", Exception),
- ],
- )
- def test_rest_transport_rejects_redirect_oversize_and_invalid_json(status, body, error):
- transport = _safe_transport([], _HttpResponse(status, body))
- with pytest.raises(error):
- transport.get_json(
- url="https://catalog.example.test/v1/catalog",
- allowed_host="catalog.example.test",
- credential_ref="secret:test/key",
- )
- with pytest.raises(ConnectorConfigurationError):
- transport.get_json(
- url="http://catalog.example.test/v1/catalog",
- allowed_host="catalog.example.test",
- credential_ref="secret:test/key",
- )
- with pytest.raises(ConnectorConfigurationError):
- transport.get_json(
- url="https://evil.example.test/v1/catalog",
- allowed_host="catalog.example.test",
- credential_ref="secret:test/key",
- )
- def test_rest_transport_without_secret_resolver_fails_before_dns_or_socket():
- calls = []
- transport = SafeRestTransport(
- resolver=lambda *_args, **_kwargs: calls.append("dns"),
- pool_factory=lambda **_kwargs: calls.append("pool"),
- )
- with pytest.raises(ConnectorConfigurationError):
- transport.get_json(
- url="https://catalog.example.test/v1/catalog",
- allowed_host="catalog.example.test",
- credential_ref="secret:test/key",
- )
- assert calls == []
- def test_default_environment_secret_resolver_is_deployable_and_fail_closed(
- monkeypatch,
- ):
- resolver = EnvironmentSecretResolver(
- {"DATAOPS_CONNECTOR_CATALOG_TOKEN": "runtime-only-value"}
- )
- assert resolver("env:DATAOPS_CONNECTOR_CATALOG_TOKEN") == "runtime-only-value"
- for reference in (
- "env:HOME",
- "env:DATAOPS_DATABASE_PASSWORD",
- "vault:catalog/token",
- "secret:catalog/token",
- ):
- with pytest.raises(ConnectorConfigurationError):
- resolver(reference)
- from flask import Flask
- from app.api.data_source import routes
- monkeypatch.setenv("DATAOPS_CONNECTOR_CATALOG_TOKEN", "assembled-secret")
- monkeypatch.setattr(
- "app.core.data_source.runtime.get_data_source_manager",
- lambda: SimpleNamespace(connect=_provider),
- )
- app = Flask(__name__)
- app.config["CONNECTOR_SECRET_RESOLVER"] = None
- with app.app_context():
- registry = routes._connector_registry()
- transport = registry.resolve("rest-catalog", "1.0.0").transport
- assert (
- transport.secret_resolver("env:DATAOPS_CONNECTOR_CATALOG_TOKEN")
- == "assembled-secret"
- )
- with pytest.raises(ConnectorConfigurationError):
- transport.secret_resolver("env:HOME")
- @pytest.mark.parametrize(
- "adapter_cls,module", [(OracleAdapter, "oracledb"), (SqlServerAdapter, "pyodbc")]
- )
- def test_optional_enterprise_drivers_fail_with_stable_safe_classification(
- monkeypatch, adapter_cls, module
- ):
- monkeypatch.setattr(
- "importlib.util.find_spec", lambda name: None if name == module else object()
- )
- definition = DataSourceDefinition(
- uid="u",
- name_en="x",
- name_zh="x",
- database_type=adapter_cls.database_type,
- host="db.example",
- port=1521,
- database="service",
- schema=None,
- credential_ref="u",
- credential_version=1,
- pool_size=None,
- max_overflow=None,
- tls_options={},
- status=True,
- description=None,
- extra_properties={},
- )
- with pytest.raises(DataSourceDriverUnavailable) as caught:
- adapter_cls().build_url(
- definition, DataSourceCredential("reader", "not-logged", {})
- )
- assert caught.value.code == "DATASOURCE_DRIVER_UNAVAILABLE"
- assert classify_error(caught.value).category == "configuration"
- @pytest.mark.parametrize(
- "connector_cls,module",
- [(OracleConnector, "oracledb"), (SqlServerConnector, "pyodbc")],
- )
- def test_health_and_compatibility_report_missing_driver_without_crash(
- monkeypatch, connector_cls, module
- ):
- monkeypatch.setattr(
- "importlib.util.find_spec",
- lambda name: None if name == module else object(),
- )
- connector = connector_cls(_provider)
- health = connector.health({"credential_ref": "env:DATAOPS_CONNECTOR_TEST"})
- compatibility = connector.compatibility()
- assert health.status == "degraded" and "driver" in health.detail
- assert compatibility.compatible is False
- assert compatibility.connector_version == "1.0.0"
- def test_runtime_idempotency_retry_cancel_resume_diff_and_bounded_redaction():
- class Flaky(Connector):
- manifest = OracleConnector.manifest
- def __init__(self):
- self.calls = 0
- def discover(self, request):
- self.calls += 1
- if self.calls == 1:
- raise ConnectorUpstreamError()
- return OperationResult(
- records=({"x": 1},),
- checkpoint={"page": 1},
- evidence={"password": "hidden"},
- )
- snapshot = incremental = lineage = profile = resume = evidence = discover
- def cancel(self, request):
- return OperationResult(status="cancelled")
- connector = Flaky()
- registry = ConnectorRegistry()
- registry.register(connector)
- store = InMemoryRunStore()
- runtime = ConnectorRuntime(registry, store=store, sleeper=lambda _: None)
- request = _request()
- first = runtime.execute("oracle", "1.0.0", request)
- second = runtime.execute("oracle", "1.0.0", request)
- assert connector.calls == 2 and second == first
- assert first.evidence["password"] == "[REDACTED]"
- assert len(deterministic_idempotency_key("oracle", "1.0.0", request)) == 64
- assert snapshot_diff(({"a": 1},), ({"a": 2},))["changed"] == (
- {"before": {"a": 1}, "after": {"a": 2}},
- )
- assert len(str(redact_evidence({"data": "x" * 40000})["data"])) == 4096
- key = "a" * 64
- store.claim(
- key,
- {
- "status": "running",
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": request.source_uid,
- "config": dict(request.config),
- "scope": {},
- "cursor": {},
- "checkpoint": {"page": 7},
- "attempt_count": 1,
- "dry_run": False,
- },
- )
- assert runtime.cancel(key)["status"] == "cancelled"
- assert isinstance(runtime.resume(key), OperationResult)
- assert store.get(key)["status"] == "succeeded"
- assert store.get(key)["attempt_count"] == 2
- assert store.get(deterministic_idempotency_key("oracle", "1.0.0", request))[
- "checkpoint"
- ] == {"page": 1}
- def test_human_rest_dry_run_never_resolves_secret_or_opens_transport():
- class NoNetwork:
- def __init__(self):
- self.calls = 0
- def get_json(self, **_values):
- self.calls += 1
- raise AssertionError("dry-run must not use network")
- transport = NoNetwork()
- registry = ConnectorRegistry()
- registry.register(RestCatalogConnector(transport))
- result = ConnectorRuntime(registry).execute(
- "rest-catalog",
- "1.0.0",
- _request(
- dry_run=True,
- config={
- "base_url": "https://attacker.example.test",
- "allowed_host": "attacker.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_GUESSED",
- },
- ),
- )
- assert result.status == "dry_run"
- assert result.evidence["network_requested"] is False
- assert transport.calls == 0
- def test_cancelled_human_dry_run_resumes_validation_only_without_network():
- class NoNetwork:
- def __init__(self):
- self.calls = 0
- def get_json(self, **_values):
- self.calls += 1
- raise AssertionError("human dry-run resume must not use transport")
- transport = NoNetwork()
- registry = ConnectorRegistry()
- registry.register(RestCatalogConnector(transport))
- store = InMemoryRunStore()
- key = "d" * 64
- config = {
- "base_url": "https://catalog.example.test",
- "allowed_host": "catalog.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_TEST",
- }
- store.claim(
- key,
- {
- "idempotency_key": key,
- "status": "running",
- "connector_id": "rest-catalog",
- "connector_version": "1.0.0",
- "source_uid": _request().source_uid,
- "config": config,
- "scope": {},
- "cursor": {},
- "checkpoint": {},
- "attempt_count": 1,
- "attempt_lease_token": "old-lease",
- "dry_run": True,
- "principal_uid": None,
- },
- )
- runtime = ConnectorRuntime(registry, store=store)
- assert runtime.cancel(key)["cancel_requested"] is True
- resumed = runtime.resume(key)
- assert resumed.status == "dry_run"
- assert resumed.evidence["validation_only"] is True
- assert resumed.evidence["network_requested"] is False
- assert transport.calls == 0
- def test_runtime_rejects_nonterminal_operation_result_status():
- class InvalidStatus(Connector):
- manifest = OracleConnector.manifest
- def discover(self, _request):
- return OperationResult(status="running")
- snapshot = incremental = lineage = profile = resume = evidence = discover
- def cancel(self, _request):
- return OperationResult(status="cancelled")
- registry = ConnectorRegistry()
- registry.register(InvalidStatus())
- store = InMemoryRunStore()
- request = _request(idempotency_key="invalid-status")
- with pytest.raises(ConnectorConfigurationError):
- ConnectorRuntime(registry, store=store, max_attempts=1).execute(
- "oracle", "1.0.0", request
- )
- key = deterministic_idempotency_key("oracle", "1.0.0", request)
- assert store.get(key)["status"] == "failed"
- def test_cancelling_one_run_does_not_cancel_another_run_for_the_same_source():
- class Probe(Connector):
- manifest = OracleConnector.manifest
- def __init__(self):
- self.discover_calls = 0
- self.cancel_request = None
- def discover(self, _request):
- self.discover_calls += 1
- return OperationResult(records=({"asset_key": "APP.ORDERS"},))
- snapshot = incremental = lineage = profile = resume = evidence = discover
- def cancel(self, request):
- self.cancel_request = request
- return OperationResult(status="cancelled")
- connector = Probe()
- registry = ConnectorRegistry()
- registry.register(connector)
- store = InMemoryRunStore()
- runtime = ConnectorRuntime(registry, store=store)
- cancelled_key = "c" * 64
- store.claim(
- cancelled_key,
- {
- "status": "running",
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": _request().source_uid,
- "config": dict(_request().config),
- "scope": {},
- "cursor": {},
- "checkpoint": {},
- "attempt_count": 1,
- "attempt_lease_token": "old-lease",
- "dry_run": False,
- },
- )
- assert runtime.cancel(cancelled_key)["status"] == "cancelled"
- assert connector.cancel_request.run_key == cancelled_key
- assert connector.cancel_request.cancel_probe() is True
- sibling = runtime.execute(
- "oracle",
- "1.0.0",
- _request(idempotency_key="sibling-run"),
- )
- assert sibling.status == "succeeded"
- assert connector.discover_calls == 1
- def test_rest_completed_cursor_restarts_full_scan_without_false_removed():
- class Pages:
- def __init__(self):
- self.urls = []
- def get_json(self, **values):
- self.urls.append(values["url"])
- return {
- "assets": [
- {
- "key": "orders",
- "name": "Orders",
- "namespace": "sales",
- "type": "table",
- }
- ],
- "completion_token": "server-complete-42",
- }
- pages = Pages()
- connector = RestCatalogConnector(pages)
- config = {
- "base_url": "https://catalog.example.test",
- "allowed_host": "catalog.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_TEST",
- }
- first = connector.snapshot(_request("snapshot", config=config))
- second = connector.incremental(
- _request(
- "incremental",
- config=config,
- cursor=first.cursor,
- checkpoint=first.checkpoint,
- )
- )
- assert pages.urls == [
- "https://catalog.example.test/v1/catalog",
- "https://catalog.example.test/v1/catalog",
- ]
- assert second.cursor == {
- "state": "complete",
- "completion_marker": "server-complete-42",
- }
- assert second.evidence["diff"] == {
- "changed": False,
- "previous_record_count": 1,
- "current_record_count": 1,
- "details_materialized": False,
- }
- assert "removed" not in second.evidence["diff"]
- with pytest.raises(ConnectorContractError):
- connector.incremental(
- _request("incremental", config=config, cursor={"state": "invalid"})
- )
- def test_runtime_sanitizes_every_output_channel_and_rejects_oversize_records():
- class Unsafe(Connector):
- manifest = OracleConnector.manifest
- def discover(self, _request):
- return OperationResult(
- records=({"password": "raw", "url": "https://u:p@host/x"},),
- cursor={"token": "raw"},
- checkpoint={"authorization": "Bearer raw"},
- evidence={"url": "https://u:p@host/x?token=raw"},
- )
- snapshot = incremental = lineage = profile = resume = evidence = discover
- def cancel(self, _request):
- return OperationResult(status="cancelled")
- registry = ConnectorRegistry()
- registry.register(Unsafe())
- result = ConnectorRuntime(registry).execute("oracle", "1.0.0", _request())
- serialized = str(result)
- assert "raw" not in serialized and "u:p" not in serialized
- assert "[REDACTED]" in serialized
- class Oversize(Unsafe):
- def discover(self, _request):
- return OperationResult(records=({"value": "x" * 1_100_000},))
- oversized = ConnectorRegistry()
- oversized.register(Oversize())
- with pytest.raises(ConnectorConfigurationError):
- ConnectorRuntime(oversized, max_attempts=1).execute(
- "oracle", "1.0.0", _request(idempotency_key="oversize")
- )
- def test_rest_pagination_rejects_repeated_and_non_monotonic_cursors():
- class Repeating:
- def get_json(self, **_values):
- return {"assets": [], "next_cursor": {"page": 1}}
- connector = RestCatalogConnector(Repeating())
- with pytest.raises(ConnectorContractError):
- connector.discover(
- _request(
- config={
- "base_url": "https://catalog.example.test",
- "allowed_host": "catalog.example.test",
- "credential_ref": "env:DATAOPS_CONNECTOR_TEST",
- },
- cursor={"page": 1},
- )
- )
- def test_sqlserver_tls_defaults_secure_and_insecure_policy_is_development_only(
- monkeypatch,
- ):
- monkeypatch.setattr("importlib.util.find_spec", lambda _name: object())
- def definition(tls_options):
- return DataSourceDefinition(
- uid="u",
- name_en="sqlserver",
- name_zh="sqlserver",
- database_type="sqlserver",
- host="db.example",
- port=1433,
- database="catalog",
- credential_ref="ref",
- credential_version=1,
- tls_options=tls_options,
- )
- credential = DataSourceCredential("reader", "secret")
- query = SqlServerAdapter().build_url(definition({}), credential).query
- assert query["Encrypt"] == "yes"
- assert query["TrustServerCertificate"] == "no"
- with pytest.raises(DataSourceConfigurationInvalid):
- SqlServerAdapter().build_url(
- definition(
- {
- "Environment": "production",
- "Encrypt": "no",
- "TrustServerCertificate": "yes",
- }
- ),
- credential,
- )
- with pytest.raises(DataSourceConfigurationInvalid):
- SqlServerAdapter().build_url(
- definition(
- {
- "Environment": "development",
- "Encrypt": "no",
- "TrustServerCertificate": "yes",
- "AllowInsecureDevelopment": "yes",
- }
- ),
- credential,
- )
- development = SqlServerAdapter().build_url_for_environment(
- definition(
- {
- "Encrypt": "no",
- "TrustServerCertificate": "yes",
- }
- ),
- credential,
- trusted_environment="development",
- allow_insecure_development=True,
- )
- assert development.query["Encrypt"] == "no"
- def test_runtime_resume_failure_is_failed_retryable_and_preserves_checkpoint():
- class ResumeFlaky(Connector):
- manifest = OracleConnector.manifest
- def __init__(self):
- self.resume_calls = 0
- def resume(self, request):
- self.resume_calls += 1
- if self.resume_calls <= 2:
- raise ConnectorUpstreamError()
- return OperationResult(checkpoint={"page": 8}, evidence={"resumed": True})
- discover = snapshot = incremental = lineage = profile = evidence = resume
- def cancel(self, request):
- return OperationResult(status="cancelled")
- class RecordingLimiter:
- def __init__(self):
- self.keys = []
- def acquire(self, key):
- self.keys.append(key)
- connector, store, limiter = ResumeFlaky(), InMemoryRunStore(), RecordingLimiter()
- registry = ConnectorRegistry()
- registry.register(connector)
- key = "b" * 64
- original_checkpoint = {"page": 7, "snapshot": [{"key": "before"}]}
- store.claim(
- key,
- {
- "uid": "11111111-1111-4111-8111-111111111112",
- "status": "cancelled",
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": _request().source_uid,
- "config": dict(_request().config),
- "scope": {},
- "cursor": {"page": 7},
- "checkpoint": original_checkpoint,
- "attempt_count": 1,
- "dry_run": False,
- },
- )
- runtime = ConnectorRuntime(
- registry, store=store, limiter=limiter, max_attempts=2, sleeper=lambda _: None
- )
- with pytest.raises(ConnectorUpstreamError):
- runtime.resume(key)
- failed = store.get(key)
- assert failed["status"] == "failed" and failed["attempt_count"] == 3
- assert failed["error_category"] == "upstream"
- assert failed["checkpoint"] == original_checkpoint
- result = runtime.resume(key)
- assert result.checkpoint == {"page": 8}
- assert store.get(key)["attempt_count"] == 4
- assert len(limiter.keys) == 2
- def test_runtime_never_exceeds_five_total_attempts_and_rate_limit_precedes_claim():
- class AlwaysFails(Connector):
- manifest = OracleConnector.manifest
- def __init__(self):
- self.calls = 0
- def discover(self, request):
- self.calls += 1
- raise ConnectorUpstreamError()
- snapshot = incremental = lineage = profile = resume = evidence = discover
- def cancel(self, request):
- return OperationResult(status="cancelled")
- connector, store = AlwaysFails(), InMemoryRunStore()
- registry = ConnectorRegistry()
- registry.register(connector)
- runtime = ConnectorRuntime(registry, store=store, sleeper=lambda _: None)
- request = _request()
- for expected_attempts in (3, MAX_TOTAL_ATTEMPTS):
- with pytest.raises(ConnectorUpstreamError):
- runtime.execute("oracle", "1.0.0", request)
- key = deterministic_idempotency_key("oracle", "1.0.0", request)
- assert store.get(key)["attempt_count"] == expected_attempts
- with pytest.raises(ConnectorConfigurationError):
- runtime.execute("oracle", "1.0.0", request)
- assert connector.calls == MAX_TOTAL_ATTEMPTS
- class RejectingLimiter:
- def acquire(self, _key):
- raise ConnectorRateLimitError()
- empty_store = InMemoryRunStore()
- with pytest.raises(ConnectorRateLimitError):
- ConnectorRuntime(
- registry, store=empty_store, limiter=RejectingLimiter()
- ).execute("oracle", "1.0.0", request)
- assert empty_store.list() == []
- def test_in_memory_attempt_cas_returns_explicit_single_owner_outcome():
- store = InMemoryRunStore()
- key = "c" * 64
- store.claim(key, {"status": "failed", "attempt_count": 1})
- winner = store.update(key, status="running", attempt_count=2)
- loser = store.update(key, status="running", attempt_count=2)
- assert winner.acquired is True and winner.record["attempt_count"] == 2
- assert loser.acquired is False and loser.record["attempt_count"] == 2
- def test_idempotency_changes_with_safe_config_and_checkpoint():
- base = _request(checkpoint={"page": 1})
- changed_checkpoint = _request(checkpoint={"page": 2})
- changed_config = _request(
- config={"credential_ref": "env:DATAOPS_CONNECTOR_OTHER"}, checkpoint={"page": 1}
- )
- keys = {
- deterministic_idempotency_key("oracle", "1.0.0", item)
- for item in (base, changed_checkpoint, changed_config)
- }
- assert len(keys) == 3
- def test_rate_limit_and_error_taxonomy_are_deterministic():
- now = [0.0]
- limiter = SlidingWindowLimiter(limit=2, window_seconds=10, clock=lambda: now[0])
- limiter.acquire("source")
- limiter.acquire("source")
- with pytest.raises(ConnectorRateLimitError):
- limiter.acquire("source")
- now[0] = 11
- limiter.acquire("source")
- assert isinstance(classify_error(TimeoutError()), ConnectorTimeoutError)
- assert isinstance(classify_error(PermissionError()), ConnectorPermissionError)
- assert isinstance(
- classify_error(ConnectorDriverUnavailableError()),
- ConnectorDriverUnavailableError,
- )
- class _Result:
- def __init__(self, scalar=None, mapping=None, rowcount=1):
- self.scalar = scalar
- self.mapping = mapping
- self.rowcount = rowcount
- def scalar_one_or_none(self):
- return self.scalar
- def mappings(self):
- return self
- def one_or_none(self):
- return self.mapping
- class _Session:
- def __init__(self, results):
- self.results = list(results)
- self.calls = []
- self.commits = 0
- self.rollbacks = 0
- def execute(self, statement, parameters=None):
- self.calls.append((str(statement), parameters or {}))
- return self.results.pop(0) if self.results else _Result()
- def commit(self):
- self.commits += 1
- def rollback(self):
- self.rollbacks += 1
- def test_machine_credential_ttl_one_time_return_and_no_plaintext_persistence():
- session = _Session([_Result(), _Result("active"), _Result(), _Result()])
- repository = ConnectorIdentityRepository(session)
- issued = repository.issue(
- "11111111-1111-4111-8111-111111111111",
- ttl_seconds=900,
- actor_uid="22222222-2222-4222-8222-222222222222",
- )
- assert issued["credential"].startswith("dopc_") and issued["returned_once"] is True
- persisted = [
- params
- for sql, params in session.calls
- if "connector_machine_credentials" in sql and "INSERT" in sql
- ][0]
- assert "credential" not in persisted and len(persisted["hash"]) == 64
- with pytest.raises(ConnectorConfigurationError):
- repository.issue("x", ttl_seconds=901, actor_uid="y")
- def test_machine_credential_scope_replay_revoke_and_rotation():
- replay_identity = {
- "uid": "c",
- "principal_uid": "p",
- "use_count": 1,
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": "s",
- "business_domain_uid": "d",
- "environment": "staging",
- "allowed_operations": ["discover"],
- "allowed_scopes": {},
- "source_binding_uid": "binding",
- "source_binding_version": 1,
- "approved_config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- }
- session = _Session(
- [_Result(mapping=None), _Result(mapping=replay_identity), _Result(), _Result()]
- )
- with pytest.raises(ConnectorAuthenticationError):
- ConnectorIdentityRepository(session).authenticate(
- "token",
- connector_id="oracle",
- connector_version="1.0.0",
- source_uid="s",
- business_domain_uid="d",
- environment="staging",
- operation="discover",
- scope={},
- )
- assert any(
- "credential_replay_rejected" in params.values()
- for _sql, params in session.calls
- )
- scope_identity = {
- **replay_identity,
- "use_count": 0,
- "allowed_scopes": {"include_schemas": ["APP"]},
- }
- session = _Session(
- [_Result(mapping=None), _Result(mapping=scope_identity), _Result()]
- )
- with pytest.raises(ConnectorPermissionError):
- ConnectorIdentityRepository(session).authenticate(
- "token",
- connector_id="oracle",
- connector_version="1.0.0",
- source_uid="s",
- business_domain_uid="d",
- environment="staging",
- operation="discover",
- scope={"include_schemas": ["SYS"]},
- )
- session = _Session([_Result("principal"), _Result()])
- assert ConnectorIdentityRepository(session).revoke("credential", "actor") is True
- rotation = _Session(
- [
- _Result("principal"),
- _Result(),
- _Result(),
- _Result("active"),
- _Result(),
- _Result(),
- ]
- )
- rotated = ConnectorIdentityRepository(rotation).rotate(
- "credential", ttl_seconds=300, actor_uid="actor"
- )
- assert rotated["returned_once"] is True and rotated["expires_in"] == 300
- def test_machine_credential_expiry_is_rejected_and_audited():
- expired = {"uid": "credential", "principal_uid": "principal"}
- session = _Session([_Result(mapping=expired), _Result(), _Result()])
- with pytest.raises(ConnectorAuthenticationError):
- ConnectorIdentityRepository(session).authenticate(
- "expired",
- connector_id="oracle",
- connector_version="1.0.0",
- source_uid="source",
- business_domain_uid="domain",
- environment="staging",
- operation="discover",
- scope={},
- )
- assert any(
- "credential_expired_rejected" in params.values()
- for _sql, params in session.calls
- )
- def test_machine_credential_success_is_scoped_and_audited():
- identity = {
- "uid": "credential",
- "principal_uid": "principal",
- "use_count": 0,
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": "source",
- "business_domain_uid": "domain",
- "environment": "staging",
- "allowed_operations": ["discover"],
- "allowed_scopes": {"include_schemas": ["APP"]},
- "source_binding_uid": "binding",
- "source_binding_version": 1,
- "approved_config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- }
- session = _Session([_Result(mapping=None), _Result(mapping=identity), _Result()])
- result = ConnectorIdentityRepository(session).authenticate(
- "token",
- connector_id="oracle",
- connector_version="1.0.0",
- source_uid="source",
- business_domain_uid="domain",
- environment="staging",
- operation="discover",
- scope={"include_schemas": ["APP"]},
- )
- assert result["principal_uid"] == "principal" and session.commits == 1
- assert any(
- "credential_authenticated" in params.values() for _sql, params in session.calls
- )
- @pytest.mark.parametrize(
- "override",
- [
- {"connector_version": "2.0.0"},
- {"business_domain_uid": "other-domain"},
- {"environment": "production"},
- ],
- )
- def test_machine_credential_rejects_cross_binding(override):
- identity = {
- "uid": "credential",
- "principal_uid": "principal",
- "use_count": 0,
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": "source",
- "business_domain_uid": "domain",
- "environment": "staging",
- "allowed_operations": ["discover"],
- "allowed_scopes": {},
- "source_binding_uid": "binding",
- "source_binding_version": 1,
- "approved_config": {"credential_ref": "env:DATAOPS_CONNECTOR_TEST"},
- }
- session = _Session([_Result(mapping=None), _Result(mapping=identity), _Result()])
- binding = {
- "connector_id": "oracle",
- "connector_version": "1.0.0",
- "source_uid": "source",
- "business_domain_uid": "domain",
- "environment": "staging",
- "operation": "discover",
- "scope": {},
- **override,
- }
- with pytest.raises(ConnectorPermissionError):
- ConnectorIdentityRepository(session).authenticate("token", **binding)
- def test_machine_scope_rejects_unknown_key_even_when_empty():
- session = _Session([])
- with pytest.raises(ConnectorConfigurationError):
- ConnectorIdentityRepository(session).authenticate(
- "token",
- connector_id="oracle",
- connector_version="1.0.0",
- source_uid="source",
- business_domain_uid="domain",
- environment="staging",
- operation="discover",
- scope={"unknown": []},
- )
- assert session.calls == []
- def test_api_permissions_and_machine_boundary_fail_closed_without_credential():
- from app import create_app
- from app.core.system.permissions import permission_for_request
- assert permission_for_request("/api/datasource/connectors/manifests", "GET") == (
- "connectors:read",
- )
- assert permission_for_request("/api/datasource/connectors/runs", "POST") == (
- "connectors:operate",
- )
- assert permission_for_request("/api/datasource/connectors/principals", "POST") == (
- "connectors:manage",
- )
- assert permission_for_request(
- "/api/datasource/connectors/machine/runs", "POST"
- ) == ("public",)
- app = create_app()
- app.config.update(TESTING=True)
- response = app.test_client().post(
- "/api/datasource/connectors/machine/runs", json={}
- )
- assert (
- response.status_code == 401
- and "credential" not in response.get_data(as_text=True).lower()
- )
- def test_migration_permissions_frontend_docs_and_release_parity_contract():
- migration = (
- ROOT / "migrations/versions/20260802_472_enterprise_connectors.py"
- ).read_text()
- for table in (
- "connector_manifests",
- "connector_principals",
- "connector_machine_credentials",
- "connector_runs",
- "connector_run_attempts",
- "connector_checkpoints",
- "connector_evidence",
- "connector_graph_edges",
- "connector_audit_events",
- ):
- assert f"CREATE TABLE public.{table}" in migration
- for marker in (
- "INTERVAL '15 minutes'",
- "idempotency_key CHAR(64) NOT NULL UNIQUE",
- "use_count",
- "replayed",
- "DROP TABLE IF EXISTS",
- ):
- assert marker in migration
- hardening = (
- ROOT / "migrations/versions/20260802_473_connector_runtime_hardening.py"
- ).read_text()
- assert 'down_revision = "20260802_472"' in hardening
- for marker in (
- "connector_version VARCHAR(40)",
- "connector_rate_limits",
- "safe_config JSONB",
- "business_domain_uid",
- "process_key",
- "ck_connector_machine_run_binding",
- ):
- assert marker in hardening
- guardrails = (
- ROOT / "migrations/versions/20260802_474_connector_version_guardrails.py"
- ).read_text()
- assert 'down_revision = "20260802_473"' in guardrails
- assert "multiple active manifest versions" in guardrails
- assert "multiple principal versions" in guardrails
- security_bindings = (
- ROOT / "migrations/versions/20260802_475_connector_security_bindings.py"
- ).read_text()
- assert 'down_revision = "20260802_474"' in security_bindings
- for marker in (
- "connector_source_bindings",
- "request_hash",
- "client_hint_hash",
- "attempt_lease_token",
- "cancel_requested",
- "explicitly remove source bindings",
- ):
- assert marker in security_bindings
- binding_enforcement = (
- ROOT / "migrations/versions/20260802_476_connector_binding_enforcement.py"
- ).read_text()
- assert 'down_revision = "20260802_475"' in binding_enforcement
- assert "bind every active enterprise principal first" in binding_enforcement
- assert (
- "process_key = db.Column(db.String(300))"
- in (ROOT / "app/models/connectors.py").read_text()
- )
- permissions = (ROOT / "app/core/system/permissions.py").read_text()
- for permission in ("connectors:read", "connectors:operate", "connectors:manage"):
- assert permission in permissions
- api = (ROOT / "app/api/data_source/routes.py").read_text()
- for endpoint in (
- "/connectors/manifests",
- "/connectors/config/validate",
- "/compatibility",
- "/connectors/runs",
- "/cancel",
- "/resume",
- "/connectors/principals",
- "/credentials",
- ):
- assert endpoint in api
- assert "DATASOURCE_GRAPH_NOT_IMPLEMENTED" not in api
- page = (
- ROOT / "frontend/src/views/dataGovernance/development/enterpriseConnectors.vue"
- ).read_text()
- connector_operations = (
- ROOT / "frontend/src/components/connectors/ConnectorOperations.vue"
- ).read_text()
- assert "真实企业账号、版本、网络与 UAT 尚未提供" in page
- assert "{{ item.credential" not in page and "{{ credential" not in page
- assert "$store.state.user.userInfo.permissions" in connector_operations
- assert "$store.getters.roles" not in connector_operations
- for permission in ("connectors:read", "connectors:operate", "connectors:manage"):
- assert permission in connector_operations
- assert "isHumanDryRun(item)" in connector_operations and 'v-if="canManage"' in connector_operations
- assert "checkpoint_summary" in connector_operations and "cursor_summary" in connector_operations
- assert "仅机器凭证可操作" in connector_operations
- assert "human connector runs must use dry_run=true" in api
- connector_files = [
- path.relative_to(ROOT).as_posix()
- for path in (ROOT / "app/core/connectors").rglob("*.py")
- ]
- for relative in (
- *connector_files,
- "app/core/data_source/adapters/__init__.py",
- "app/core/data_source/adapters/oracle.py",
- "app/core/data_source/adapters/sqlserver.py",
- "app/core/data_source/errors.py",
- "app/core/data_source/service.py",
- "app/models/connectors.py",
- "app/models/__init__.py",
- "app/api/data_source/routes.py",
- "app/core/system/permissions.py",
- ):
- assert (ROOT / relative).read_bytes() == (
- ROOT / "deployment" / relative
- ).read_bytes()
- assert (
- ROOT / "migrations/versions/20260802_472_enterprise_connectors.py"
- ).read_bytes() == (
- ROOT / "deployment/migrations/versions/20260802_472_enterprise_connectors.py"
- ).read_bytes()
- def test_openapi_connector_contract_is_strict_and_matches_behavior():
- contract = yaml.safe_load((ROOT / "docs/architecture/OPENAPI.yaml").read_text())
- paths = contract["paths"]
- machine = paths["/api/datasource/connectors/machine/runs"]["post"]
- assert machine["requestBody"]["required"] is True
- assert machine["requestBody"]["content"]["application/json"]["schema"][
- "$ref"
- ].endswith("MachineConnectorRunRequest")
- header = next(
- item
- for item in machine["parameters"]
- if item["name"] == "X-Connector-Credential"
- )
- assert header["in"] == "header" and header["required"] is True
- assert set(machine["responses"]) == {"201", "default"}
- assert machine["responses"]["201"]["content"]["application/json"]["schema"][
- "$ref"
- ].endswith("ConnectorRunResultEnvelope")
- assert machine["x-required-permission"] == "machine-credential"
- for action in ("cancel", "resume"):
- operation = paths[
- f"/api/datasource/connectors/runs/{{idempotency_key}}/{action}"
- ]["post"]
- assert "requestBody" not in operation
- human = paths["/api/datasource/connectors/runs"]["post"]
- assert (
- human["responses"].get("201")
- and human["x-required-permission"] == "connectors:operate"
- )
- schemas = contract["components"]["schemas"]
- assert schemas["ConnectorRunRequest"]["properties"]["dry_run"] == {"const": True}
- assert schemas["ConnectorScope"]["additionalProperties"] is False
- assert schemas["ConnectorRunResultEnvelope"]["additionalProperties"] is False
- assert "config" not in schemas["MachineConnectorRunRequest"]["properties"]
- assert schemas["ConnectorOperationResult"]["properties"]["status"]["enum"] == [
- "succeeded",
- "dry_run",
- ]
- assert "allOf" not in schemas["ConnectorHealthEnvelope"]
- categories = schemas["ConnectorErrorEnvelope"]["properties"]["error"]["properties"][
- "category"
- ]["enum"]
- assert {
- "configuration",
- "authentication",
- "permission",
- "conflict",
- "rate_limit",
- "timeout",
- "upstream",
- "contract",
- "cancelled",
- } == set(categories)
- for action in ("cancel", "resume"):
- machine_action = paths[
- f"/api/datasource/connectors/machine/runs/{{idempotency_key}}/{action}"
- ]["post"]
- assert "requestBody" not in machine_action
- assert machine_action["x-required-permission"] == "machine-credential"
- credential_action = paths[
- "/api/datasource/connectors/credentials/{credential_uid}/{action}"
- ]["post"]
- assert credential_action["requestBody"]["required"] is False
- credential_issue = paths[
- "/api/datasource/connectors/principals/{principal_uid}/credentials"
- ]["post"]
- assert credential_issue["requestBody"]["required"] is False
- binding_create = paths["/api/datasource/connectors/source-bindings"]["post"]
- assert set(binding_create["responses"]) == {"201", "default"}
- assert binding_create["responses"]["201"]["content"]["application/json"][
- "schema"
- ]["$ref"].endswith("ConnectorSourceBindingEnvelope")
- action_parameter = next(
- item for item in credential_action["parameters"] if item["name"] == "action"
- )
- assert action_parameter["schema"]["enum"] == ["revoke", "rotate"]
- assert paths["/api/datasource/graph"]["post"]["x-required-permission"] == (
- "connectors:read"
- )
- for envelope in (
- "ConnectorRunResultEnvelope",
- "ConnectorErrorEnvelope",
- "ConnectorHealthEnvelope",
- ):
- properties = schemas[envelope]["properties"]
- assert "success" not in properties and "timestamp" not in properties
- assert {"code", "message", "data"}.issubset(properties)
- assert (
- ROOT / "migrations/versions/20260802_473_connector_runtime_hardening.py"
- ).read_bytes() == (
- ROOT
- / "deployment/migrations/versions/20260802_473_connector_runtime_hardening.py"
- ).read_bytes()
- assert (
- ROOT / "migrations/versions/20260802_474_connector_version_guardrails.py"
- ).read_bytes() == (
- ROOT
- / "deployment/migrations/versions/20260802_474_connector_version_guardrails.py"
- ).read_bytes()
- assert (
- ROOT / "migrations/versions/20260802_475_connector_security_bindings.py"
- ).read_bytes() == (
- ROOT
- / "deployment/migrations/versions/20260802_475_connector_security_bindings.py"
- ).read_bytes()
- assert (
- ROOT / "migrations/versions/20260802_476_connector_binding_enforcement.py"
- ).read_bytes() == (
- ROOT
- / "deployment/migrations/versions/20260802_476_connector_binding_enforcement.py"
- ).read_bytes()
- def test_flask_machine_responses_validate_against_openapi_201_and_401(
- monkeypatch,
- ):
- from app import create_app
- from app.core.connectors.identity import ConnectorIdentityRepository
- from app.core.connectors.repository import ConnectorRepository
- from app.core.connectors.runtime import ConnectorRuntime
- contract = yaml.safe_load((ROOT / "docs/architecture/OPENAPI.yaml").read_text())
- def expand(schema):
- if isinstance(schema, dict) and "$ref" in schema:
- target = contract
- for part in schema["$ref"].removeprefix("#/").split("/"):
- target = target[part]
- return expand(target)
- if isinstance(schema, dict):
- return {key: expand(value) for key, value in schema.items()}
- if isinstance(schema, list):
- return [expand(value) for value in schema]
- return schema
- monkeypatch.setattr(
- ConnectorIdentityRepository,
- "authenticate",
- lambda _self, _token, **_binding: {
- "principal_uid": "principal",
- "source_binding_uid": "binding",
- "source_binding_version": 1,
- "approved_config": {
- "credential_ref": "env:DATAOPS_CONNECTOR_TEST"
- },
- },
- )
- monkeypatch.setattr(
- ConnectorRepository,
- "register_manifest",
- lambda _self, _manifest, _actor: None,
- )
- monkeypatch.setattr(
- ConnectorRuntime,
- "execute",
- lambda _self, _connector, _version, _request: OperationResult(
- records=({"asset_key": "source:APP.ORDERS"},),
- checkpoint={"page": 1},
- cursor={"page": 2},
- evidence={"safe": True},
- ),
- )
- app = create_app()
- app.config.update(TESTING=True)
- app.extensions["connector_registry"] = SimpleNamespace(
- resolve=lambda *_args, **_kwargs: SimpleNamespace(
- manifest=OracleConnector.manifest
- )
- )
- client = app.test_client()
- path = "/api/datasource/connectors/machine/runs"
- unauthorized = client.post(path, json={})
- assert unauthorized.status_code == 401
- error_schema = expand(
- contract["paths"][path]["post"]["responses"]["default"]["content"][
- "application/json"
- ]["schema"]
- )
- Draft202012Validator(error_schema).validate(unauthorized.get_json())
- assert set(unauthorized.get_json()) == {"code", "message", "data", "error"}
- assert unauthorized.get_json()["error"] == {
- "code": "CONNECTOR_ERROR",
- "category": "authentication",
- "retryable": False,
- }
- created = client.post(
- path,
- headers={"X-Connector-Credential": "one-time-token"},
- json={
- "connector_id": "oracle",
- "version": "1.0.0",
- "source_uid": "11111111-1111-4111-8111-111111111111",
- "business_domain_uid": "22222222-2222-4222-8222-222222222222",
- "environment": "staging",
- "process_key": "catalog-sync",
- "operation": "discover",
- "scope": {},
- },
- )
- assert created.status_code == 201
- success_schema = expand(
- contract["paths"][path]["post"]["responses"]["201"]["content"][
- "application/json"
- ]["schema"]
- )
- Draft202012Validator(success_schema).validate(created.get_json())
- assert set(created.get_json()) == {"code", "message", "data"}
- assert created.get_json()["code"] == created.status_code == 201
- invalid = client.post(
- path,
- headers={"X-Connector-Credential": "one-time-token"},
- json={
- "connector_id": "oracle",
- "version": "1.0.0",
- "source_uid": "11111111-1111-4111-8111-111111111111",
- "business_domain_uid": "22222222-2222-4222-8222-222222222222",
- "environment": "staging",
- "process_key": "catalog-sync",
- "operation": "cancel",
- "scope": {},
- },
- )
- assert invalid.status_code == invalid.get_json()["code"] == 400
- Draft202012Validator(error_schema).validate(invalid.get_json())
- assert invalid.get_json()["error"]["category"] == "configuration"
- def test_human_run_requires_explicit_dry_run():
- from app.api.data_source.routes import (
- _require_collection_operation,
- _require_human_dry_run,
- )
- _require_human_dry_run({"dry_run": True})
- with pytest.raises(ConnectorConfigurationError):
- _require_human_dry_run({"dry_run": False})
- with pytest.raises(ConnectorConfigurationError):
- _require_human_dry_run({})
- _require_collection_operation({"operation": "incremental"})
- with pytest.raises(ConnectorConfigurationError):
- _require_collection_operation({"operation": "cancel"})
|