"""Versioned public connector SDK. Only symbols exported here and from ``app.core.connectors`` are stable. Connector packages register against the registry; the collection runtime never branches on connector names. """ from __future__ import annotations import re from abc import ABC, abstractmethod from collections.abc import Callable, Mapping from dataclasses import asdict, dataclass, field from typing import Any from jsonschema import Draft202012Validator from jsonschema.exceptions import SchemaError, ValidationError from app.core.connectors.errors import ConnectorConfigurationError SDK_VERSION = "1.0" CAPABILITIES = frozenset( { "discover", "snapshot", "incremental", "lineage", "profile", "cancel", "resume", "evidence", } ) SECRET_KEY_PATTERN = re.compile( r"(?:password|passwd|secret|token|api[_-]?key|authorization|private[_-]?key)", re.I ) SECRET_REF_PATTERN = re.compile( r"^(?:env|vault|secret):[A-Za-z0-9][A-Za-z0-9_./:-]{2,255}$" ) SEMVER_PATTERN = re.compile(r"^[0-9]+\.[0-9]+\.[0-9]+(?:-[0-9A-Za-z.-]+)?$") @dataclass(frozen=True) class ConnectorManifest: connector_id: str version: str sdk_version: str display_name: str capabilities: tuple[str, ...] config_schema: Mapping[str, Any] = field(repr=False) def __post_init__(self): if not re.fullmatch(r"[a-z][a-z0-9_-]{2,63}", self.connector_id): raise ConnectorConfigurationError("connector_id is invalid") if not SEMVER_PATTERN.fullmatch(self.version): raise ConnectorConfigurationError("connector version is invalid") if self.sdk_version != SDK_VERSION: raise ConnectorConfigurationError("SDK version is incompatible") capabilities = tuple(dict.fromkeys(self.capabilities)) if not capabilities or set(capabilities) - CAPABILITIES: raise ConnectorConfigurationError("connector capability is unsupported") object.__setattr__(self, "capabilities", capabilities) validate_schema(self.config_schema) def public_dict(self): return asdict(self) @dataclass(frozen=True) class OperationRequest: source_uid: str operation: str config: Mapping[str, Any] scope: Mapping[str, Any] = field(default_factory=dict) cursor: Mapping[str, Any] = field(default_factory=dict) checkpoint: Mapping[str, Any] = field(default_factory=dict) idempotency_key: str | None = None dry_run: bool = False principal_uid: str | None = None business_domain_uid: str | None = None environment: str | None = None process_key: str | None = None source_binding_uid: str | None = None source_binding_version: int | None = None run_key: str | None = None lease_token: str | None = None cancel_probe: Callable[[], bool] | None = field( default=None, compare=False, repr=False ) @dataclass(frozen=True) class OperationResult: records: tuple[Mapping[str, Any], ...] = () cursor: Mapping[str, Any] = field(default_factory=dict) checkpoint: Mapping[str, Any] = field(default_factory=dict) evidence: Mapping[str, Any] = field(default_factory=dict) status: str = "succeeded" @dataclass(frozen=True) class HealthResult: status: str detail: str = "" @dataclass(frozen=True) class CompatibilityResult: compatible: bool connector_version: str sdk_version: str = SDK_VERSION detail: str = "" def _walk_secrets(value, path="config"): if isinstance(value, Mapping): for key, item in value.items(): current = f"{path}.{key}" if SECRET_KEY_PATTERN.search(str(key)) and ( not str(key).endswith("_ref") or not isinstance(item, str) or not SECRET_REF_PATTERN.fullmatch(item) ): raise ConnectorConfigurationError( f"{current} must be a secret reference" ) _walk_secrets(item, current) elif isinstance(value, (list, tuple)): for index, item in enumerate(value): _walk_secrets(item, f"{path}[{index}]") def validate_schema(schema): if not isinstance(schema, Mapping) or schema.get("type") != "object": raise ConnectorConfigurationError("config schema must describe an object") if schema.get("additionalProperties") is not False: raise ConnectorConfigurationError( "config schema must reject unknown properties" ) properties = schema.get("properties") if not isinstance(properties, Mapping): raise ConnectorConfigurationError("config schema properties are required") def inspect_definition(definition): if not isinstance(definition, Mapping): return for key, child in definition.get("properties", {}).items(): if SECRET_KEY_PATTERN.search(str(key)) and not str(key).endswith("_ref"): raise ConnectorConfigurationError( "config schema cannot define plaintext secrets" ) if ( str(key).endswith("_ref") and child.get("pattern") != SECRET_REF_PATTERN.pattern ): raise ConnectorConfigurationError( "secret reference schema pattern is required" ) inspect_definition(child) inspect_definition(definition.get("items")) for keyword in ("allOf", "anyOf", "oneOf"): for child in definition.get(keyword, ()): inspect_definition(child) for child in definition.get("$defs", {}).values(): inspect_definition(child) inspect_definition(schema) try: Draft202012Validator.check_schema(dict(schema)) except SchemaError as exc: raise ConnectorConfigurationError("config JSON Schema is invalid") from exc def validate_config(schema, config): """Validate recursively with JSON Schema Draft 2020-12.""" if not isinstance(config, Mapping): raise ConnectorConfigurationError("connector config must be an object") _walk_secrets(config) try: Draft202012Validator(dict(schema)).validate(dict(config)) except ValidationError as exc: raise ConnectorConfigurationError("connector config shape is invalid") from exc return dict(config) class Connector(ABC): manifest: ConnectorManifest @abstractmethod def discover(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def snapshot(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def incremental(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def lineage(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def profile(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def cancel(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def resume(self, request: OperationRequest) -> OperationResult: ... @abstractmethod def evidence(self, request: OperationRequest) -> OperationResult: ... def health(self, config: Mapping[str, Any]) -> HealthResult: validate_config(self.manifest.config_schema, config) return HealthResult("available") def compatibility(self) -> CompatibilityResult: return CompatibilityResult(True, self.manifest.version) __all__ = [ "SDK_VERSION", "CAPABILITIES", "SECRET_REF_PATTERN", "ConnectorManifest", "OperationRequest", "OperationResult", "HealthResult", "CompatibilityResult", "Connector", "validate_config", "validate_schema", ]