| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152 |
- """Thread-safe, deny-by-default connector registry."""
- from __future__ import annotations
- import threading
- from app.core.connectors.errors import ConnectorConfigurationError
- from app.core.connectors.sdk import Connector, validate_config
- class ConnectorRegistry:
- def __init__(self):
- self._lock = threading.RLock()
- self._connectors = {}
- def register(self, connector: Connector):
- manifest = connector.manifest
- key = (manifest.connector_id, manifest.version)
- with self._lock:
- if key in self._connectors:
- raise ConnectorConfigurationError(
- "connector version is already registered"
- )
- self._connectors[key] = connector
- return connector
- def resolve(self, connector_id, version, capability=None):
- with self._lock:
- connector = self._connectors.get((str(connector_id), str(version)))
- if connector is None:
- raise ConnectorConfigurationError("connector or version is not registered")
- if capability and capability not in connector.manifest.capabilities:
- raise ConnectorConfigurationError("connector capability is not declared")
- return connector
- def validate(self, connector_id, version, config):
- connector = self.resolve(connector_id, version)
- return validate_config(connector.manifest.config_schema, config)
- def manifests(self):
- with self._lock:
- items = tuple(self._connectors.values())
- return [
- item.manifest.public_dict()
- for item in sorted(
- items,
- key=lambda item: (item.manifest.connector_id, item.manifest.version),
- )
- ]
- __all__ = ["ConnectorRegistry"]
|