mirror of
https://github.com/langflow-ai/langflow.git
synced 2026-07-24 11:47:11 +08:00
refactor(service_manager): implement lazy initalization of service manager (#8828)
* refactor: implement lazy initialization for ServiceManager with thread safety - Replaced direct instantiation of ServiceManager with a lazy initialization approach using a global variable and threading lock. - Updated the public API to expose `get_service_manager` for retrieving the singleton instance. - Ensured thread-safe access to the ServiceManager instance to prevent issues during module import. * refactor: update service manager imports to use get_service_manager - Replaced direct imports of service_manager with get_service_manager in multiple files to ensure consistent access to the singleton instance. - This change enhances code clarity and maintains the lazy initialization approach for the ServiceManager. * refactor: remove deprecated Enhanced ServiceManager implementation * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * refactor: implement thread-safe lazy initialization for ServiceManager * feat: add filelock dependency in lfx * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * [autofix.ci] apply automated fixes (attempt 3/3) * update component index * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * [autofix.ci] apply automated fixes (attempt 3/3) --------- Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
This commit is contained in:
committed by
GitHub
parent
2d4469d9a6
commit
5bcede4dec
@ -1,3 +1,4 @@
|
||||
from .manager import get_service_manager
|
||||
from .schema import ServiceType
|
||||
|
||||
__all__ = ["ServiceType"]
|
||||
__all__ = ["ServiceType", "get_service_manager"]
|
||||
|
||||
@ -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
|
||||
@ -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())
|
||||
|
||||
@ -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]
|
||||
|
||||
@ -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
|
||||
|
||||
2
uv.lock
generated
2
uv.lock
generated
@ -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" },
|
||||
|
||||
Reference in New Issue
Block a user