| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233 |
- """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",
- ]
|