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() assert "真实企业账号、版本、网络与 UAT 尚未提供" in page assert "{{ item.credential" not in page and "{{ credential" not in page assert "$store.state.user.userInfo.permissions" in page assert "$store.getters.roles" not in page for permission in ("connectors:read", "connectors:operate", "connectors:manage"): assert permission in page assert "isHumanDryRun(item)" in page and 'v-if="canManage"' in page assert "checkpoint_summary" in page and "cursor_summary" in page assert "仅机器凭证可操作" in page 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"})