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