From a07164d1aeb27e2b42a5afa41f427d098481fc08 Mon Sep 17 00:00:00 2001 From: Arek Mateusiak Date: Tue, 21 Apr 2026 19:17:33 +0200 Subject: [PATCH] refactor(lifespan): Enhance initialization flow with conditional gates - Introduced conditional gates in the `get_lifespan` function to manage the initialization of profile pictures, super users, bundles, component types, and starter projects based on their completion status. - Improved logging to provide clearer insights into which steps are being skipped or executed, enhancing the overall clarity of the initialization process. - Updated the preload logic in `preload.py` to ensure that agentic global variables and MCP server configuration are only initialized when necessary, maintaining efficient resource management in a multi-worker environment. --- src/backend/base/langflow/main.py | 95 ++++++++++++++++------------ src/backend/base/langflow/preload.py | 27 +++++--- 2 files changed, 73 insertions(+), 49 deletions(-) diff --git a/src/backend/base/langflow/main.py b/src/backend/base/langflow/main.py index c24c8de959..be5fd44c9d 100644 --- a/src/backend/base/langflow/main.py +++ b/src/backend/base/langflow/main.py @@ -151,9 +151,7 @@ def get_lifespan(*, fix_migration=False, version=None): async def lifespan(_app: FastAPI): from lfx.interface.components import get_and_cache_all_types_dict - from langflow.preload import _STATE, get_owned_temp_dirs, is_preloaded - - preloaded = is_preloaded() + from langflow.preload import _STATE, get_owned_temp_dirs configure() @@ -204,25 +202,23 @@ def get_lifespan(*, fix_migration=False, version=None): setup_llm_caching() await logger.adebug(f"LLM caching setup in {asyncio.get_event_loop().time() - current_time:.2f}s") - if not preloaded or not _STATE.profile_pictures_copied: + # Gate: Copy profile pictures + if _STATE.profile_pictures_copied: + await logger.adebug("Skipping profile-picture copy: master already completed it during preload") + else: current_time = asyncio.get_event_loop().time() await logger.adebug("Copying profile pictures") await copy_profile_pictures() await logger.adebug(f"Profile pictures copied in {asyncio.get_event_loop().time() - current_time:.2f}s") - if get_settings_service().auth_settings.AUTO_LOGIN: - current_time = asyncio.get_event_loop().time() - await logger.adebug("Initializing default super user") - await initialize_auto_login_default_superuser() - await logger.adebug( - f"Default super user initialized in {asyncio.get_event_loop().time() - current_time:.2f}s" - ) - - await logger.adebug("Initializing super user") + # Gate: Initialize default superuser (when AUTO_LOGIN is enabled) + if get_settings_service().auth_settings.AUTO_LOGIN: + current_time = asyncio.get_event_loop().time() + await logger.adebug("Initializing default super user") await initialize_auto_login_default_superuser() - await logger.adebug(f"Super user initialized in {asyncio.get_event_loop().time() - current_time:.2f}s") - elif _STATE.profile_pictures_copied: - await logger.adebug("Skipping profile-picture copy: master already completed it during preload") + await logger.adebug( + f"Default super user initialized in {asyncio.get_event_loop().time() - current_time:.2f}s" + ) if get_settings_service().settings.prometheus_enabled: try: @@ -249,13 +245,14 @@ def get_lifespan(*, fix_migration=False, version=None): telemetry_service = get_telemetry_service() - if preloaded and _STATE.starter_projects_created: + # Gate: Load bundles + if _STATE.starter_projects_created: # Inherit bundle paths and types dict from master via COW. # get_owned_temp_dirs() returns the preloaded dirs if this is # the master, or an empty list if this is a worker (workers # must NOT clean up the master's temp_dirs). temp_dirs = get_owned_temp_dirs() - await logger.adebug("Skipping bundle load, types cache and starter projects: inherited from master") + await logger.adebug("Skipping bundle load: inherited from master") else: current_time = asyncio.get_event_loop().time() await logger.adebug("Loading bundles") @@ -263,11 +260,19 @@ def get_lifespan(*, fix_migration=False, version=None): get_settings_service().settings.components_path.extend(bundles_components_paths) await logger.adebug(f"Bundles loaded in {asyncio.get_event_loop().time() - current_time:.2f}s") + # Gate: Cache component types + if _STATE.starter_projects_created: + await logger.adebug("Skipping types cache: inherited from master") + else: current_time = asyncio.get_event_loop().time() await logger.adebug("Caching types") all_types_dict = await get_and_cache_all_types_dict(get_settings_service(), telemetry_service) await logger.adebug(f"Types cached in {asyncio.get_event_loop().time() - current_time:.2f}s") + # Gate: Create/update starter projects + if _STATE.starter_projects_created: + await logger.adebug("Skipping starter projects: inherited from master") + else: # Use file-based lock to prevent multiple workers from creating duplicate starter projects # concurrently. Note that it's still possible that one worker may complete this task, release # the lock, then another worker pick it up, but the operation is idempotent so worst case it @@ -291,21 +296,24 @@ def get_lifespan(*, fix_migration=False, version=None): "Starter projects may not be created or updated." ) - # Initialize agentic global variables early (before MCP server and flows) - if get_settings_service().settings.agentic_experience and not _STATE.agentic_globals_initialized: - from langflow.api.utils.mcp.agentic_mcp import initialize_agentic_global_variables + # Gate: Initialize agentic global variables (when agentic_experience enabled) + if not get_settings_service().settings.agentic_experience: + pass # Skip when feature is disabled + elif _STATE.agentic_globals_initialized: + await logger.adebug("Skipping agentic global variables: master already completed it during preload") + else: + from langflow.api.utils.mcp.agentic_mcp import initialize_agentic_global_variables - current_time = asyncio.get_event_loop().time() - await logger.ainfo("Initializing agentic global variables...") - try: - async with session_scope() as session: - await initialize_agentic_global_variables(session) - await logger.adebug( - "Agentic global variables initialized in " - f"{asyncio.get_event_loop().time() - current_time:.2f}s" - ) - except Exception as e: # noqa: BLE001 - await logger.awarning(f"Failed to initialize agentic global variables: {e}") + current_time = asyncio.get_event_loop().time() + await logger.ainfo("Initializing agentic global variables...") + try: + async with session_scope() as session: + await initialize_agentic_global_variables(session) + await logger.adebug( + f"Agentic global variables initialized in {asyncio.get_event_loop().time() - current_time:.2f}s" + ) + except Exception as e: # noqa: BLE001 + await logger.awarning(f"Failed to initialize agentic global variables: {e}") current_time = asyncio.get_event_loop().time() await logger.adebug("Starting telemetry service") @@ -320,9 +328,12 @@ def get_lifespan(*, fix_migration=False, version=None): f"started MCP Composer service in {asyncio.get_event_loop().time() - current_time:.2f}s" ) - # Auto-configure Agentic MCP server if enabled (after variables are initialized). - # Skipped when the master already completed this during preload. - if get_settings_service().settings.agentic_experience and not _STATE.agentic_mcp_configured: + # Gate: Auto-configure agentic MCP server (when agentic_experience enabled) + if not get_settings_service().settings.agentic_experience: + pass # Skip when feature is disabled + elif _STATE.agentic_mcp_configured: + await logger.adebug("Skipping agentic MCP server config: master already completed it during preload") + else: from langflow.api.utils.mcp.agentic_mcp import auto_configure_agentic_mcp_server current_time = asyncio.get_event_loop().time() @@ -336,19 +347,21 @@ def get_lifespan(*, fix_migration=False, version=None): except Exception as e: # noqa: BLE001 await logger.awarning(f"Failed to configure agentic MCP server: {e}") + # Gate: Load flows from directory current_time = asyncio.get_event_loop().time() - if not _STATE.flows_loaded: + if _STATE.flows_loaded: + await logger.adebug("Skipping flows load: master already completed it during preload") + else: await logger.adebug("Loading flows") await load_flows_from_directory() - else: - await logger.adebug("Skipping flows load: master already completed it during preload") - # ``sync_flows_from_fs`` and the queue service MUST be started - # per-worker: they create asyncio tasks bound to this event loop. + await logger.adebug(f"Flows loaded in {asyncio.get_event_loop().time() - current_time:.2f}s") + + # Per-worker setup: sync_flows_from_fs and queue service + # (MUST be started per-worker: they create asyncio tasks bound to this event loop) sync_flows_from_fs_task = asyncio.create_task(sync_flows_from_fs()) queue_service = get_queue_service() if not queue_service.is_started(): queue_service.start() - await logger.adebug(f"Flows loaded in {asyncio.get_event_loop().time() - current_time:.2f}s") total_time = asyncio.get_event_loop().time() - start_time await logger.adebug(f"Total initialization time: {total_time:.2f}s") diff --git a/src/backend/base/langflow/preload.py b/src/backend/base/langflow/preload.py index 29a5696015..869ee9d465 100644 --- a/src/backend/base/langflow/preload.py +++ b/src/backend/base/langflow/preload.py @@ -17,6 +17,14 @@ started here. Each worker continues to set those up in its own FastAPI ``lifespan`` after fork, so workers remain fully independent and can each serve any request on their own. +Failure contract: + Fork-safety-critical steps (DB engine disposal, cache service teardown) + that fail will propagate their exception and abort preload. Best-effort + steps (profile pictures, starter projects, agentic MCP, flows) that fail + will log a warning, clear their completion flag, and allow preload to + continue so workers inherit partial progress. Workers re-run any + incomplete step during their lifespan. + Notes on CPython and copy-on-write: CPython mutates ``ob_refcnt`` on every attribute access, which triggers 4 KB page copies even for "shared" objects. Preloading @@ -152,23 +160,26 @@ async def _run_master_preload() -> None: await logger.awarning(f"[preload] starter projects init failed: {e}") if settings_service.settings.agentic_experience: - try: - from langflow.api.utils.mcp.agentic_mcp import ( - auto_configure_agentic_mcp_server, - initialize_agentic_global_variables, - ) + from langflow.api.utils.mcp.agentic_mcp import ( + auto_configure_agentic_mcp_server, + initialize_agentic_global_variables, + ) - await logger.adebug("[preload] initializing agentic global variables") + await logger.adebug("[preload] initializing agentic global variables") + try: async with session_scope() as session: await initialize_agentic_global_variables(session) _STATE.agentic_globals_initialized = True + except Exception as e: # noqa: BLE001 + await logger.awarning(f"[preload] initialize agentic global variables failed: {e}") - await logger.adebug("[preload] auto-configuring agentic MCP server") + await logger.adebug("[preload] auto-configuring agentic MCP server") + try: async with session_scope() as session: await auto_configure_agentic_mcp_server(session) _STATE.agentic_mcp_configured = True except Exception as e: # noqa: BLE001 - await logger.awarning(f"[preload] agentic MCP init failed: {e}") + await logger.awarning(f"[preload] auto-configure agentic MCP server failed: {e}") await logger.adebug("[preload] loading flows from directory") try: