sdk.py 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233
  1. """Versioned public connector SDK.
  2. Only symbols exported here and from ``app.core.connectors`` are stable. Connector
  3. packages register against the registry; the collection runtime never branches on
  4. connector names.
  5. """
  6. from __future__ import annotations
  7. import re
  8. from abc import ABC, abstractmethod
  9. from collections.abc import Callable, Mapping
  10. from dataclasses import asdict, dataclass, field
  11. from typing import Any
  12. from jsonschema import Draft202012Validator
  13. from jsonschema.exceptions import SchemaError, ValidationError
  14. from app.core.connectors.errors import ConnectorConfigurationError
  15. SDK_VERSION = "1.0"
  16. CAPABILITIES = frozenset(
  17. {
  18. "discover",
  19. "snapshot",
  20. "incremental",
  21. "lineage",
  22. "profile",
  23. "cancel",
  24. "resume",
  25. "evidence",
  26. }
  27. )
  28. SECRET_KEY_PATTERN = re.compile(
  29. r"(?:password|passwd|secret|token|api[_-]?key|authorization|private[_-]?key)", re.I
  30. )
  31. SECRET_REF_PATTERN = re.compile(
  32. r"^(?:env|vault|secret):[A-Za-z0-9][A-Za-z0-9_./:-]{2,255}$"
  33. )
  34. SEMVER_PATTERN = re.compile(r"^[0-9]+\.[0-9]+\.[0-9]+(?:-[0-9A-Za-z.-]+)?$")
  35. @dataclass(frozen=True)
  36. class ConnectorManifest:
  37. connector_id: str
  38. version: str
  39. sdk_version: str
  40. display_name: str
  41. capabilities: tuple[str, ...]
  42. config_schema: Mapping[str, Any] = field(repr=False)
  43. def __post_init__(self):
  44. if not re.fullmatch(r"[a-z][a-z0-9_-]{2,63}", self.connector_id):
  45. raise ConnectorConfigurationError("connector_id is invalid")
  46. if not SEMVER_PATTERN.fullmatch(self.version):
  47. raise ConnectorConfigurationError("connector version is invalid")
  48. if self.sdk_version != SDK_VERSION:
  49. raise ConnectorConfigurationError("SDK version is incompatible")
  50. capabilities = tuple(dict.fromkeys(self.capabilities))
  51. if not capabilities or set(capabilities) - CAPABILITIES:
  52. raise ConnectorConfigurationError("connector capability is unsupported")
  53. object.__setattr__(self, "capabilities", capabilities)
  54. validate_schema(self.config_schema)
  55. def public_dict(self):
  56. return asdict(self)
  57. @dataclass(frozen=True)
  58. class OperationRequest:
  59. source_uid: str
  60. operation: str
  61. config: Mapping[str, Any]
  62. scope: Mapping[str, Any] = field(default_factory=dict)
  63. cursor: Mapping[str, Any] = field(default_factory=dict)
  64. checkpoint: Mapping[str, Any] = field(default_factory=dict)
  65. idempotency_key: str | None = None
  66. dry_run: bool = False
  67. principal_uid: str | None = None
  68. business_domain_uid: str | None = None
  69. environment: str | None = None
  70. process_key: str | None = None
  71. source_binding_uid: str | None = None
  72. source_binding_version: int | None = None
  73. run_key: str | None = None
  74. lease_token: str | None = None
  75. cancel_probe: Callable[[], bool] | None = field(
  76. default=None, compare=False, repr=False
  77. )
  78. @dataclass(frozen=True)
  79. class OperationResult:
  80. records: tuple[Mapping[str, Any], ...] = ()
  81. cursor: Mapping[str, Any] = field(default_factory=dict)
  82. checkpoint: Mapping[str, Any] = field(default_factory=dict)
  83. evidence: Mapping[str, Any] = field(default_factory=dict)
  84. status: str = "succeeded"
  85. @dataclass(frozen=True)
  86. class HealthResult:
  87. status: str
  88. detail: str = ""
  89. @dataclass(frozen=True)
  90. class CompatibilityResult:
  91. compatible: bool
  92. connector_version: str
  93. sdk_version: str = SDK_VERSION
  94. detail: str = ""
  95. def _walk_secrets(value, path="config"):
  96. if isinstance(value, Mapping):
  97. for key, item in value.items():
  98. current = f"{path}.{key}"
  99. if SECRET_KEY_PATTERN.search(str(key)) and (
  100. not str(key).endswith("_ref")
  101. or not isinstance(item, str)
  102. or not SECRET_REF_PATTERN.fullmatch(item)
  103. ):
  104. raise ConnectorConfigurationError(
  105. f"{current} must be a secret reference"
  106. )
  107. _walk_secrets(item, current)
  108. elif isinstance(value, (list, tuple)):
  109. for index, item in enumerate(value):
  110. _walk_secrets(item, f"{path}[{index}]")
  111. def validate_schema(schema):
  112. if not isinstance(schema, Mapping) or schema.get("type") != "object":
  113. raise ConnectorConfigurationError("config schema must describe an object")
  114. if schema.get("additionalProperties") is not False:
  115. raise ConnectorConfigurationError(
  116. "config schema must reject unknown properties"
  117. )
  118. properties = schema.get("properties")
  119. if not isinstance(properties, Mapping):
  120. raise ConnectorConfigurationError("config schema properties are required")
  121. def inspect_definition(definition):
  122. if not isinstance(definition, Mapping):
  123. return
  124. for key, child in definition.get("properties", {}).items():
  125. if SECRET_KEY_PATTERN.search(str(key)) and not str(key).endswith("_ref"):
  126. raise ConnectorConfigurationError(
  127. "config schema cannot define plaintext secrets"
  128. )
  129. if (
  130. str(key).endswith("_ref")
  131. and child.get("pattern") != SECRET_REF_PATTERN.pattern
  132. ):
  133. raise ConnectorConfigurationError(
  134. "secret reference schema pattern is required"
  135. )
  136. inspect_definition(child)
  137. inspect_definition(definition.get("items"))
  138. for keyword in ("allOf", "anyOf", "oneOf"):
  139. for child in definition.get(keyword, ()):
  140. inspect_definition(child)
  141. for child in definition.get("$defs", {}).values():
  142. inspect_definition(child)
  143. inspect_definition(schema)
  144. try:
  145. Draft202012Validator.check_schema(dict(schema))
  146. except SchemaError as exc:
  147. raise ConnectorConfigurationError("config JSON Schema is invalid") from exc
  148. def validate_config(schema, config):
  149. """Validate recursively with JSON Schema Draft 2020-12."""
  150. if not isinstance(config, Mapping):
  151. raise ConnectorConfigurationError("connector config must be an object")
  152. _walk_secrets(config)
  153. try:
  154. Draft202012Validator(dict(schema)).validate(dict(config))
  155. except ValidationError as exc:
  156. raise ConnectorConfigurationError("connector config shape is invalid") from exc
  157. return dict(config)
  158. class Connector(ABC):
  159. manifest: ConnectorManifest
  160. @abstractmethod
  161. def discover(self, request: OperationRequest) -> OperationResult: ...
  162. @abstractmethod
  163. def snapshot(self, request: OperationRequest) -> OperationResult: ...
  164. @abstractmethod
  165. def incremental(self, request: OperationRequest) -> OperationResult: ...
  166. @abstractmethod
  167. def lineage(self, request: OperationRequest) -> OperationResult: ...
  168. @abstractmethod
  169. def profile(self, request: OperationRequest) -> OperationResult: ...
  170. @abstractmethod
  171. def cancel(self, request: OperationRequest) -> OperationResult: ...
  172. @abstractmethod
  173. def resume(self, request: OperationRequest) -> OperationResult: ...
  174. @abstractmethod
  175. def evidence(self, request: OperationRequest) -> OperationResult: ...
  176. def health(self, config: Mapping[str, Any]) -> HealthResult:
  177. validate_config(self.manifest.config_schema, config)
  178. return HealthResult("available")
  179. def compatibility(self) -> CompatibilityResult:
  180. return CompatibilityResult(True, self.manifest.version)
  181. __all__ = [
  182. "SDK_VERSION",
  183. "CAPABILITIES",
  184. "SECRET_REF_PATTERN",
  185. "ConnectorManifest",
  186. "OperationRequest",
  187. "OperationResult",
  188. "HealthResult",
  189. "CompatibilityResult",
  190. "Connector",
  191. "validate_config",
  192. "validate_schema",
  193. ]