registry.py 1.7 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152
  1. """Thread-safe, deny-by-default connector registry."""
  2. from __future__ import annotations
  3. import threading
  4. from app.core.connectors.errors import ConnectorConfigurationError
  5. from app.core.connectors.sdk import Connector, validate_config
  6. class ConnectorRegistry:
  7. def __init__(self):
  8. self._lock = threading.RLock()
  9. self._connectors = {}
  10. def register(self, connector: Connector):
  11. manifest = connector.manifest
  12. key = (manifest.connector_id, manifest.version)
  13. with self._lock:
  14. if key in self._connectors:
  15. raise ConnectorConfigurationError(
  16. "connector version is already registered"
  17. )
  18. self._connectors[key] = connector
  19. return connector
  20. def resolve(self, connector_id, version, capability=None):
  21. with self._lock:
  22. connector = self._connectors.get((str(connector_id), str(version)))
  23. if connector is None:
  24. raise ConnectorConfigurationError("connector or version is not registered")
  25. if capability and capability not in connector.manifest.capabilities:
  26. raise ConnectorConfigurationError("connector capability is not declared")
  27. return connector
  28. def validate(self, connector_id, version, config):
  29. connector = self.resolve(connector_id, version)
  30. return validate_config(connector.manifest.config_schema, config)
  31. def manifests(self):
  32. with self._lock:
  33. items = tuple(self._connectors.values())
  34. return [
  35. item.manifest.public_dict()
  36. for item in sorted(
  37. items,
  38. key=lambda item: (item.manifest.connector_id, item.manifest.version),
  39. )
  40. ]
  41. __all__ = ["ConnectorRegistry"]