diff --git a/src/backend/base/langflow/services/__init__.py b/src/backend/base/langflow/services/__init__.py index 06aa0020aa..248ca173b1 100644 --- a/src/backend/base/langflow/services/__init__.py +++ b/src/backend/base/langflow/services/__init__.py @@ -1,3 +1,4 @@ +from .manager import get_service_manager from .schema import ServiceType -__all__ = ["ServiceType"] +__all__ = ["ServiceType", "get_service_manager"] diff --git a/src/backend/base/langflow/services/enhanced_manager.py b/src/backend/base/langflow/services/enhanced_manager.py deleted file mode 100644 index d431c069fa..0000000000 --- a/src/backend/base/langflow/services/enhanced_manager.py +++ /dev/null @@ -1,72 +0,0 @@ -"""Enhanced ServiceManager that extends lfx's ServiceManager with langflow features.""" - -from __future__ import annotations - -import importlib -import inspect -from typing import TYPE_CHECKING - -from lfx.log.logger import logger -from lfx.services.manager import NoFactoryRegisteredError -from lfx.services.manager import ServiceManager as BaseServiceManager -from lfx.utils.concurrency import KeyedMemoryLockManager - -if TYPE_CHECKING: - from langflow.services.base import Service - from langflow.services.factory import ServiceFactory - from langflow.services.schema import ServiceType - - -__all__ = ["NoFactoryRegisteredError", "ServiceManager"] - - -class ServiceManager(BaseServiceManager): - """Enhanced ServiceManager with langflow factory system and dependency injection.""" - - def __init__(self) -> None: - super().__init__() - self.register_factories() - self.keyed_lock = KeyedMemoryLockManager() - - def register_factories(self, factories: list[ServiceFactory] | None = None) -> None: - """Register all available service factories.""" - for factory in factories or self.get_factories(): - try: - self.register_factory(factory) - except Exception: # noqa: BLE001 - logger.exception(f"Error initializing {factory}") - - def get(self, service_name: ServiceType, default: ServiceFactory | None = None) -> Service: - """Get (or create) a service by its name with keyed locking.""" - with self.keyed_lock.lock(service_name): - return super().get(service_name, default) - - @classmethod - def get_factories(cls) -> list[ServiceFactory]: - """Auto-discover and return all service factories.""" - from langflow.services.factory import ServiceFactory - from langflow.services.schema import ServiceType - - service_names = [ServiceType(service_type).value.replace("_service", "") for service_type in ServiceType] - base_module = "langflow.services" - factories = [] - - for name in service_names: - try: - # Special handling for services that are in lfx module - base_module = "lfx.services" if name in ["settings", "mcp_composer"] else "langflow.services" - module_name = f"{base_module}.{name}.factory" - module = importlib.import_module(module_name) - - # Find all classes in the module that are subclasses of ServiceFactory - for _, obj in inspect.getmembers(module, inspect.isclass): - if issubclass(obj, ServiceFactory) and obj is not ServiceFactory: - factories.append(obj()) - break - - except Exception as exc: - logger.exception(exc) - msg = f"Could not initialize services. Please check your settings. Error in {name}." - raise RuntimeError(msg) from exc - - return factories diff --git a/src/backend/base/langflow/services/manager.py b/src/backend/base/langflow/services/manager.py index 3c95ccc2ce..e3db3085a4 100644 --- a/src/backend/base/langflow/services/manager.py +++ b/src/backend/base/langflow/services/manager.py @@ -1,19 +1,15 @@ -"""Langflow ServiceManager that extends lfx's ServiceManager with enhanced features. - -This maintains backward compatibility while using lfx as the foundation. -""" +"""Langflow ServiceManager - re-exports from lfx for backwards compatibility.""" from __future__ import annotations -# Import the enhanced manager that extends lfx -from langflow.services.enhanced_manager import NoFactoryRegisteredError, ServiceManager +# Re-export everything from lfx +from lfx.services.manager import NoFactoryRegisteredError, ServiceManager, get_service_manager -__all__ = ["NoFactoryRegisteredError", "ServiceManager"] +__all__ = ["NoFactoryRegisteredError", "ServiceManager", "get_service_manager"] def initialize_settings_service() -> None: """Initialize the settings manager.""" - from lfx.services.manager import get_service_manager from lfx.services.settings import factory as settings_factory get_service_manager().register_factory(settings_factory.SettingsServiceFactory()) @@ -21,13 +17,10 @@ def initialize_settings_service() -> None: def initialize_session_service() -> None: """Initialize the session manager.""" - from lfx.services.manager import get_service_manager - from langflow.services.cache import factory as cache_factory from langflow.services.session import factory as session_service_factory initialize_settings_service() get_service_manager().register_factory(cache_factory.CacheServiceFactory()) - get_service_manager().register_factory(session_service_factory.SessionServiceFactory()) diff --git a/src/lfx/pyproject.toml b/src/lfx/pyproject.toml index 6308200033..76878bb0c6 100644 --- a/src/lfx/pyproject.toml +++ b/src/lfx/pyproject.toml @@ -39,6 +39,7 @@ dependencies = [ "loguru>=0.7.3,<1.0.0", "langchain~=0.3.23", "validators>=0.34.0,<1.0.0", + "filelock>=3.20.0", ] [project.scripts] diff --git a/src/lfx/src/lfx/services/manager.py b/src/lfx/src/lfx/services/manager.py index 5057f75cfb..91e8e19821 100644 --- a/src/lfx/src/lfx/services/manager.py +++ b/src/lfx/src/lfx/services/manager.py @@ -14,6 +14,7 @@ from typing import TYPE_CHECKING from lfx.log.logger import logger from lfx.services.schema import ServiceType +from lfx.utils.concurrency import KeyedMemoryLockManager if TYPE_CHECKING: from lfx.services.base import Service @@ -31,6 +32,7 @@ class ServiceManager: self.services: dict[str, Service] = {} self.factories: dict[str, ServiceFactory] = {} self._lock = threading.RLock() + self.keyed_lock = KeyedMemoryLockManager() self.factory_registered = False from lfx.services.settings.factory import SettingsServiceFactory @@ -65,7 +67,7 @@ class ServiceManager: def get(self, service_name: ServiceType, default: ServiceFactory | None = None) -> Service: """Get (or create) a service by its name.""" - with self._lock: + with self.keyed_lock.lock(service_name): if service_name not in self.services: self._create_service(service_name, default) return self.services[service_name] @@ -161,12 +163,23 @@ class ServiceManager: return factories -# Global service manager instance -_service_manager = None +# Global variables for lazy initialization +_service_manager: ServiceManager | None = None +_service_manager_lock = threading.Lock() -def get_service_manager(): +def get_service_manager() -> ServiceManager: + """Get or create the service manager instance using lazy initialization. + + This function ensures thread-safe lazy initialization of the service manager, + preventing automatic service creation during module import. + + Returns: + ServiceManager: The singleton service manager instance. + """ global _service_manager # noqa: PLW0603 if _service_manager is None: - _service_manager = ServiceManager() + with _service_manager_lock: + if _service_manager is None: + _service_manager = ServiceManager() return _service_manager diff --git a/uv.lock b/uv.lock index a31be2af5f..c34c098f10 100644 --- a/uv.lock +++ b/uv.lock @@ -6200,6 +6200,7 @@ dependencies = [ { name = "docstring-parser" }, { name = "emoji" }, { name = "fastapi" }, + { name = "filelock" }, { name = "httpx", extra = ["http2"] }, { name = "json-repair" }, { name = "langchain" }, @@ -6247,6 +6248,7 @@ requires-dist = [ { name = "docstring-parser", specifier = ">=0.16,<1.0.0" }, { name = "emoji", specifier = ">=2.14.1,<3.0.0" }, { name = "fastapi", specifier = ">=0.115.13,<1.0.0" }, + { name = "filelock", specifier = ">=3.20.0" }, { name = "httpx", extras = ["http2"], specifier = ">=0.24.0,<1.0.0" }, { name = "json-repair", specifier = ">=0.30.3,<1.0.0" }, { name = "langchain", specifier = "~=0.3.23" },