mirror of
https://github.com/langflow-ai/langflow.git
synced 2026-07-24 16:12:19 +08:00
fix: resolve double-call of initialize_auto_login_default_superuser and add preload tests
- Fix double-call issue: setup_superuser() now handles AUTO_LOGIN completely with file lock - Add comprehensive unit tests for preload.py covering failure-fallback contract - Simplify code by doing superuser initialization in initialize_services() (called early in both preload and worker startup) - File lock protects multi-worker race conditions when preload is disabled - Tests verify critical step failures propagate, best-effort steps continue on failure Made-with: Cursor
This commit is contained in:
@ -1234,6 +1234,12 @@ async def create_or_update_starter_projects(all_types_dict: dict) -> None:
|
||||
|
||||
|
||||
async def initialize_auto_login_default_superuser() -> None:
|
||||
"""Initialize the default superuser for AUTO_LOGIN mode.
|
||||
|
||||
Note: In production, this is called indirectly via setup_superuser() during
|
||||
initialize_services(), which includes file lock protection for multi-worker
|
||||
environments. This standalone function is kept for testing and CLI usage.
|
||||
"""
|
||||
settings_service = get_settings_service()
|
||||
if not settings_service.auth_settings.AUTO_LOGIN:
|
||||
return
|
||||
|
||||
@ -33,7 +33,6 @@ from langflow.api.v1.mcp_projects import init_mcp_servers
|
||||
from langflow.initial_setup.setup import (
|
||||
copy_profile_pictures,
|
||||
create_or_update_starter_projects,
|
||||
initialize_auto_login_default_superuser,
|
||||
load_bundles_from_urls,
|
||||
load_flows_from_directory,
|
||||
sync_flows_from_fs,
|
||||
@ -211,15 +210,6 @@ def get_lifespan(*, fix_migration=False, version=None):
|
||||
await copy_profile_pictures()
|
||||
await logger.adebug(f"Profile pictures copied in {asyncio.get_event_loop().time() - current_time:.2f}s")
|
||||
|
||||
# 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"Default super user initialized in {asyncio.get_event_loop().time() - current_time:.2f}s"
|
||||
)
|
||||
|
||||
if get_settings_service().settings.prometheus_enabled:
|
||||
try:
|
||||
from prometheus_client import start_http_server
|
||||
|
||||
@ -2,10 +2,11 @@ from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
from importlib import import_module
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from lfx.log.logger import logger
|
||||
from lfx.services.settings.constants import DEFAULT_SUPERUSER, DEFAULT_SUPERUSER_PASSWORD
|
||||
from lfx.services.settings.constants import DEFAULT_SUPERUSER
|
||||
from lfx.services.settings.feature_flags import FEATURE_FLAGS
|
||||
from sqlalchemy import delete
|
||||
from sqlalchemy import exc as sqlalchemy_exc
|
||||
@ -71,16 +72,52 @@ async def get_or_create_super_user(session: AsyncSession, username, password, is
|
||||
|
||||
async def setup_superuser(settings_service: SettingsService, session: AsyncSession) -> None:
|
||||
if settings_service.auth_settings.AUTO_LOGIN:
|
||||
await logger.adebug("AUTO_LOGIN is set to True. Creating default superuser.")
|
||||
await logger.adebug("AUTO_LOGIN is set to True. Creating default superuser with full initialization.")
|
||||
# Use file lock to prevent race conditions in multi-worker environments
|
||||
from tempfile import gettempdir
|
||||
|
||||
from filelock import FileLock
|
||||
from lfx.services.settings.constants import DEFAULT_SUPERUSER, DEFAULT_SUPERUSER_PASSWORD
|
||||
|
||||
username = DEFAULT_SUPERUSER
|
||||
password = DEFAULT_SUPERUSER_PASSWORD.get_secret_value()
|
||||
else:
|
||||
# Remove the default superuser if it exists
|
||||
await teardown_superuser(settings_service, session)
|
||||
# If AUTO_LOGIN is disabled, attempt to use configured credentials
|
||||
# or fall back to default credentials if none are provided.
|
||||
username = settings_service.auth_settings.SUPERUSER or DEFAULT_SUPERUSER
|
||||
password = (settings_service.auth_settings.SUPERUSER_PASSWORD or DEFAULT_SUPERUSER_PASSWORD).get_secret_value()
|
||||
|
||||
if not username or not password:
|
||||
msg = "SUPERUSER and SUPERUSER_PASSWORD must be set in the settings if AUTO_LOGIN is true."
|
||||
raise ValueError(msg)
|
||||
|
||||
# Use file lock similar to starter projects
|
||||
lock_file = Path(gettempdir()) / "langflow_auto_login_superuser.lock"
|
||||
lock = FileLock(lock_file, timeout=5)
|
||||
|
||||
try:
|
||||
with lock:
|
||||
# Create user and initialize all related resources
|
||||
super_user = await get_or_create_super_user(session, username, password, is_default=True)
|
||||
if super_user: # Only initialize if user was created
|
||||
from langflow.initial_setup.setup import get_or_create_default_folder
|
||||
from langflow.services.deps import get_variable_service
|
||||
|
||||
await get_variable_service().initialize_user_variables(super_user.id, session)
|
||||
|
||||
# Initialize agentic variables if enabled
|
||||
if settings_service.settings.agentic_experience:
|
||||
from langflow.api.utils.mcp.agentic_mcp import initialize_agentic_user_variables
|
||||
|
||||
await initialize_agentic_user_variables(super_user.id, session)
|
||||
|
||||
_ = await get_or_create_default_folder(session, super_user.id)
|
||||
await logger.adebug("Auto-login superuser initialized successfully")
|
||||
except TimeoutError:
|
||||
# Another worker is handling it - all operations are idempotent
|
||||
await logger.adebug("Another worker is initializing auto-login superuser, skipping")
|
||||
return
|
||||
# Remove the default superuser if it exists
|
||||
await teardown_superuser(settings_service, session)
|
||||
# If AUTO_LOGIN is disabled, attempt to use configured credentials
|
||||
# or fall back to default credentials if none are provided.
|
||||
username = settings_service.auth_settings.SUPERUSER or DEFAULT_SUPERUSER
|
||||
password = (settings_service.auth_settings.SUPERUSER_PASSWORD or DEFAULT_SUPERUSER_PASSWORD).get_secret_value()
|
||||
|
||||
if not username or not password:
|
||||
msg = "Username and password must be set"
|
||||
|
||||
491
src/backend/tests/unit/base/test_preload.py
Normal file
491
src/backend/tests/unit/base/test_preload.py
Normal file
@ -0,0 +1,491 @@
|
||||
"""Unit tests for langflow.preload module.
|
||||
|
||||
These tests verify the failure-fallback contract and state management
|
||||
of the Gunicorn master preload functionality.
|
||||
"""
|
||||
|
||||
import os
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
import pytest
|
||||
from langflow.preload import (
|
||||
_STATE,
|
||||
_run_master_preload,
|
||||
get_owned_temp_dirs,
|
||||
get_preloaded_temp_dirs,
|
||||
is_master,
|
||||
is_preloaded,
|
||||
preload_master,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def reset_preload_state():
|
||||
"""Reset the preload state before and after each test."""
|
||||
original_state = {
|
||||
"preloaded": _STATE.preloaded,
|
||||
"master_pid": _STATE.master_pid,
|
||||
"temp_dirs": _STATE.temp_dirs.copy(),
|
||||
"bundles_components_paths": _STATE.bundles_components_paths.copy(),
|
||||
"profile_pictures_copied": _STATE.profile_pictures_copied,
|
||||
"starter_projects_created": _STATE.starter_projects_created,
|
||||
"agentic_globals_initialized": _STATE.agentic_globals_initialized,
|
||||
"agentic_mcp_configured": _STATE.agentic_mcp_configured,
|
||||
"flows_loaded": _STATE.flows_loaded,
|
||||
}
|
||||
|
||||
# Reset to defaults
|
||||
_STATE.preloaded = False
|
||||
_STATE.master_pid = None
|
||||
_STATE.temp_dirs = []
|
||||
_STATE.bundles_components_paths = []
|
||||
_STATE.profile_pictures_copied = False
|
||||
_STATE.starter_projects_created = False
|
||||
_STATE.agentic_globals_initialized = False
|
||||
_STATE.agentic_mcp_configured = False
|
||||
_STATE.flows_loaded = False
|
||||
|
||||
yield
|
||||
|
||||
# Restore original state
|
||||
_STATE.preloaded = original_state["preloaded"]
|
||||
_STATE.master_pid = original_state["master_pid"]
|
||||
_STATE.temp_dirs = original_state["temp_dirs"]
|
||||
_STATE.bundles_components_paths = original_state["bundles_components_paths"]
|
||||
_STATE.profile_pictures_copied = original_state["profile_pictures_copied"]
|
||||
_STATE.starter_projects_created = original_state["starter_projects_created"]
|
||||
_STATE.agentic_globals_initialized = original_state["agentic_globals_initialized"]
|
||||
_STATE.agentic_mcp_configured = original_state["agentic_mcp_configured"]
|
||||
_STATE.flows_loaded = original_state["flows_loaded"]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Helper function tests
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_is_preloaded_false_by_default():
|
||||
"""is_preloaded() should return False when preload hasn't run."""
|
||||
assert is_preloaded() is False
|
||||
|
||||
|
||||
def test_is_preloaded_true_after_preload():
|
||||
"""is_preloaded() should return True after successful preload."""
|
||||
_STATE.preloaded = True
|
||||
assert is_preloaded() is True
|
||||
|
||||
|
||||
def test_is_master_false_by_default():
|
||||
"""is_master() should return False when master_pid is not set."""
|
||||
assert is_master() is False
|
||||
|
||||
|
||||
def test_is_master_true_in_master_process():
|
||||
"""is_master() should return True when called in the master process."""
|
||||
_STATE.master_pid = os.getpid()
|
||||
assert is_master() is True
|
||||
|
||||
|
||||
def test_is_master_false_in_worker_process():
|
||||
"""is_master() should return False when called in a forked worker."""
|
||||
_STATE.master_pid = os.getpid() + 1 # Simulate different PID
|
||||
assert is_master() is False
|
||||
|
||||
|
||||
def test_get_preloaded_temp_dirs_empty_by_default():
|
||||
"""get_preloaded_temp_dirs() should return empty list before preload."""
|
||||
assert get_preloaded_temp_dirs() == []
|
||||
|
||||
|
||||
def test_get_preloaded_temp_dirs_returns_list():
|
||||
"""get_preloaded_temp_dirs() should return the temp_dirs list."""
|
||||
fake_temp_dir = MagicMock()
|
||||
_STATE.temp_dirs = [fake_temp_dir]
|
||||
result = get_preloaded_temp_dirs()
|
||||
assert result == [fake_temp_dir]
|
||||
|
||||
|
||||
def test_get_owned_temp_dirs_master_owns_dirs():
|
||||
"""get_owned_temp_dirs() should return temp_dirs when called in master."""
|
||||
fake_temp_dir = MagicMock()
|
||||
_STATE.preloaded = True
|
||||
_STATE.master_pid = os.getpid()
|
||||
_STATE.temp_dirs = [fake_temp_dir]
|
||||
|
||||
result = get_owned_temp_dirs()
|
||||
assert result == [fake_temp_dir]
|
||||
|
||||
|
||||
def test_get_owned_temp_dirs_worker_owns_nothing():
|
||||
"""get_owned_temp_dirs() should return empty list when called in worker."""
|
||||
fake_temp_dir = MagicMock()
|
||||
_STATE.preloaded = True
|
||||
_STATE.master_pid = os.getpid() + 1 # Simulate worker
|
||||
_STATE.temp_dirs = [fake_temp_dir]
|
||||
|
||||
result = get_owned_temp_dirs()
|
||||
assert result == []
|
||||
|
||||
|
||||
def test_get_owned_temp_dirs_not_preloaded():
|
||||
"""get_owned_temp_dirs() should return empty list when not preloaded."""
|
||||
fake_temp_dir = MagicMock()
|
||||
_STATE.preloaded = False
|
||||
_STATE.temp_dirs = [fake_temp_dir]
|
||||
|
||||
result = get_owned_temp_dirs()
|
||||
assert result == []
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# preload_master() function tests
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@patch("langflow.preload.asyncio.run")
|
||||
def test_preload_master_idempotent(mock_asyncio_run):
|
||||
"""preload_master() should be idempotent (no-op on subsequent calls)."""
|
||||
_STATE.preloaded = True
|
||||
|
||||
preload_master()
|
||||
|
||||
# Should not call asyncio.run if already preloaded
|
||||
mock_asyncio_run.assert_not_called()
|
||||
|
||||
|
||||
@patch("langflow.preload.asyncio.run")
|
||||
def test_preload_master_sets_state_on_success(mock_asyncio_run):
|
||||
"""preload_master() should set preloaded flag on success."""
|
||||
mock_asyncio_run.return_value = None # Simulate successful completion
|
||||
|
||||
preload_master()
|
||||
|
||||
assert _STATE.preloaded is True
|
||||
assert _STATE.master_pid == os.getpid()
|
||||
mock_asyncio_run.assert_called_once()
|
||||
|
||||
|
||||
@patch("langflow.preload.logger")
|
||||
@patch("langflow.preload.asyncio.run")
|
||||
def test_preload_master_no_flag_on_failure(mock_asyncio_run, mock_logger):
|
||||
"""preload_master() should NOT set preloaded flag if preload fails.
|
||||
|
||||
This verifies the failure-fallback contract: if preload fails,
|
||||
workers will fall back to running full lifespan initialization.
|
||||
"""
|
||||
mock_asyncio_run.side_effect = RuntimeError("Preload failed")
|
||||
|
||||
preload_master()
|
||||
|
||||
assert _STATE.preloaded is False
|
||||
mock_logger.exception.assert_called_once()
|
||||
|
||||
|
||||
@patch("langflow.preload.gc")
|
||||
@patch("langflow.preload.asyncio.run")
|
||||
def test_preload_master_calls_gc_freeze(mock_asyncio_run, mock_gc):
|
||||
"""preload_master() should call gc.freeze() after successful preload."""
|
||||
mock_asyncio_run.return_value = None
|
||||
|
||||
preload_master()
|
||||
|
||||
mock_gc.collect.assert_called_once()
|
||||
mock_gc.freeze.assert_called_once()
|
||||
|
||||
|
||||
@patch("langflow.preload.gc")
|
||||
@patch("langflow.preload.logger")
|
||||
@patch("langflow.preload.asyncio.run")
|
||||
def test_preload_master_continues_if_gc_freeze_fails(mock_asyncio_run, mock_logger, mock_gc):
|
||||
"""preload_master() should continue (not abort) if gc.freeze() fails.
|
||||
|
||||
gc.freeze() failure is not critical, so preload should still succeed.
|
||||
"""
|
||||
mock_asyncio_run.return_value = None
|
||||
mock_gc.freeze.side_effect = RuntimeError("gc.freeze() failed")
|
||||
|
||||
preload_master()
|
||||
|
||||
assert _STATE.preloaded is True # Should still succeed
|
||||
mock_logger.exception.assert_called()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# _run_master_preload() function tests - Failure-Fallback Contract
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@patch("langflow.preload.get_db_service")
|
||||
@patch("langflow.preload.get_settings_service")
|
||||
@patch("langflow.preload.initialize_services")
|
||||
async def test_run_master_preload_disposes_db_engine(
|
||||
mock_initialize_services,
|
||||
mock_get_settings_service,
|
||||
mock_get_db_service,
|
||||
):
|
||||
"""_run_master_preload() must dispose DB engine before returning.
|
||||
|
||||
This is CRITICAL for fork-safety. Failure to dispose the engine
|
||||
should propagate and abort preload.
|
||||
"""
|
||||
mock_engine = AsyncMock()
|
||||
mock_db_service = MagicMock()
|
||||
mock_db_service.engine = mock_engine
|
||||
mock_get_db_service.return_value = mock_db_service
|
||||
|
||||
mock_settings_service = MagicMock()
|
||||
mock_settings_service.settings.agentic_experience = False
|
||||
mock_get_settings_service.return_value = mock_settings_service
|
||||
|
||||
mock_initialize_services.return_value = None
|
||||
|
||||
# Mock all the setup functions to avoid actual operations
|
||||
with (
|
||||
patch("langflow.preload.copy_profile_pictures", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_bundles_with_error_handling", new_callable=AsyncMock) as mock_load_bundles,
|
||||
patch("langflow.preload.get_and_cache_all_types_dict", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_flows_from_directory", new_callable=AsyncMock),
|
||||
patch("langflow.preload.get_service"),
|
||||
):
|
||||
mock_load_bundles.return_value = ([], [])
|
||||
await _run_master_preload()
|
||||
|
||||
# Verify DB engine was disposed
|
||||
mock_engine.dispose.assert_called_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@patch("langflow.preload.get_db_service")
|
||||
@patch("langflow.preload.get_settings_service")
|
||||
@patch("langflow.preload.initialize_services")
|
||||
async def test_run_master_preload_critical_step_failure_propagates(
|
||||
mock_initialize_services,
|
||||
mock_get_settings_service,
|
||||
mock_get_db_service,
|
||||
):
|
||||
"""_run_master_preload() should propagate exceptions from critical steps.
|
||||
|
||||
If DB engine disposal fails, the exception should propagate and
|
||||
abort preload, forcing workers to fall back to full initialization.
|
||||
"""
|
||||
mock_engine = AsyncMock()
|
||||
mock_engine.dispose.side_effect = RuntimeError("Failed to dispose DB engine")
|
||||
mock_db_service = MagicMock()
|
||||
mock_db_service.engine = mock_engine
|
||||
mock_get_db_service.return_value = mock_db_service
|
||||
|
||||
mock_settings_service = MagicMock()
|
||||
mock_settings_service.settings.agentic_experience = False
|
||||
mock_get_settings_service.return_value = mock_settings_service
|
||||
|
||||
mock_initialize_services.return_value = None
|
||||
|
||||
with (
|
||||
patch("langflow.preload.copy_profile_pictures", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_bundles_with_error_handling", new_callable=AsyncMock) as mock_load_bundles,
|
||||
patch("langflow.preload.get_and_cache_all_types_dict", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_flows_from_directory", new_callable=AsyncMock),
|
||||
patch("langflow.preload.get_service"),
|
||||
):
|
||||
mock_load_bundles.return_value = ([], [])
|
||||
with pytest.raises(RuntimeError, match="Failed to dispose DB engine"):
|
||||
await _run_master_preload()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@patch("langflow.preload.get_db_service")
|
||||
@patch("langflow.preload.get_settings_service")
|
||||
@patch("langflow.preload.initialize_services")
|
||||
@patch("langflow.preload.logger")
|
||||
async def test_run_master_preload_best_effort_step_failure_continues(
|
||||
mock_logger,
|
||||
mock_initialize_services,
|
||||
mock_get_settings_service,
|
||||
mock_get_db_service,
|
||||
):
|
||||
"""_run_master_preload() should continue if best-effort steps fail.
|
||||
|
||||
Best-effort steps (profile pictures, starter projects, etc.) should
|
||||
log a warning and clear their completion flag, but not abort preload.
|
||||
Workers will re-run incomplete steps during their lifespan.
|
||||
"""
|
||||
mock_engine = AsyncMock()
|
||||
mock_db_service = MagicMock()
|
||||
mock_db_service.engine = mock_engine
|
||||
mock_get_db_service.return_value = mock_db_service
|
||||
|
||||
mock_settings_service = MagicMock()
|
||||
mock_settings_service.settings.agentic_experience = False
|
||||
mock_get_settings_service.return_value = mock_settings_service
|
||||
|
||||
mock_initialize_services.return_value = None
|
||||
|
||||
# Mock copy_profile_pictures to fail
|
||||
with (
|
||||
patch("langflow.preload.copy_profile_pictures", new_callable=AsyncMock) as mock_copy_pics,
|
||||
patch("langflow.preload.load_bundles_with_error_handling", new_callable=AsyncMock) as mock_load_bundles,
|
||||
patch("langflow.preload.get_and_cache_all_types_dict", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_flows_from_directory", new_callable=AsyncMock),
|
||||
patch("langflow.preload.get_service"),
|
||||
):
|
||||
mock_copy_pics.side_effect = RuntimeError("Failed to copy profile pictures")
|
||||
mock_load_bundles.return_value = ([], [])
|
||||
# Should NOT raise exception
|
||||
await _run_master_preload()
|
||||
|
||||
# Verify warning was logged
|
||||
assert any("copy_profile_pictures failed" in str(call) for call in mock_logger.awarning.call_args_list)
|
||||
|
||||
# Verify completion flag was NOT set (workers will re-run this step)
|
||||
assert _STATE.profile_pictures_copied is False
|
||||
|
||||
# Verify DB engine was still disposed (critical step)
|
||||
mock_engine.dispose.assert_called_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@patch("langflow.preload.get_db_service")
|
||||
@patch("langflow.preload.get_settings_service")
|
||||
@patch("langflow.preload.initialize_services")
|
||||
@patch("langflow.preload.component_cache")
|
||||
async def test_run_master_preload_sets_completion_flags_on_success(
|
||||
mock_component_cache,
|
||||
mock_initialize_services,
|
||||
mock_get_settings_service,
|
||||
mock_get_db_service,
|
||||
):
|
||||
"""_run_master_preload() should set completion flags for successful steps.
|
||||
|
||||
This verifies that workers inherit the completion state and can
|
||||
skip redundant initialization for steps that completed during preload.
|
||||
"""
|
||||
mock_engine = AsyncMock()
|
||||
mock_db_service = MagicMock()
|
||||
mock_db_service.engine = mock_engine
|
||||
mock_get_db_service.return_value = mock_db_service
|
||||
|
||||
mock_settings_service = MagicMock()
|
||||
mock_settings_service.settings.agentic_experience = False
|
||||
mock_settings_service.settings.components_path = []
|
||||
mock_get_settings_service.return_value = mock_settings_service
|
||||
|
||||
mock_initialize_services.return_value = None
|
||||
|
||||
# Mock component_cache.all_types_dict
|
||||
mock_component_cache.all_types_dict = {"fake": "types_dict"}
|
||||
|
||||
with (
|
||||
patch("langflow.preload.copy_profile_pictures", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_bundles_with_error_handling", new_callable=AsyncMock) as mock_load_bundles,
|
||||
patch("langflow.preload.get_and_cache_all_types_dict", new_callable=AsyncMock),
|
||||
patch("langflow.preload.create_or_update_starter_projects", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_flows_from_directory", new_callable=AsyncMock),
|
||||
patch("langflow.preload.get_service"),
|
||||
):
|
||||
mock_load_bundles.return_value = ([], [])
|
||||
await _run_master_preload()
|
||||
|
||||
# Verify completion flags were set
|
||||
assert _STATE.profile_pictures_copied is True
|
||||
assert _STATE.starter_projects_created is True
|
||||
assert _STATE.flows_loaded is True
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@patch("langflow.preload.get_db_service")
|
||||
@patch("langflow.preload.get_settings_service")
|
||||
@patch("langflow.preload.initialize_services")
|
||||
async def test_run_master_preload_closes_cache_service_socket(
|
||||
mock_initialize_services,
|
||||
mock_get_settings_service,
|
||||
mock_get_db_service,
|
||||
):
|
||||
"""_run_master_preload() should teardown external cache service before fork.
|
||||
|
||||
This is CRITICAL for fork-safety with Redis or other external cache services.
|
||||
"""
|
||||
mock_engine = AsyncMock()
|
||||
mock_db_service = MagicMock()
|
||||
mock_db_service.engine = mock_engine
|
||||
mock_get_db_service.return_value = mock_db_service
|
||||
|
||||
mock_settings_service = MagicMock()
|
||||
mock_settings_service.settings.agentic_experience = False
|
||||
mock_get_settings_service.return_value = mock_settings_service
|
||||
|
||||
mock_initialize_services.return_value = None
|
||||
|
||||
# Mock cache service with teardown method
|
||||
mock_cache_service = MagicMock()
|
||||
mock_teardown = AsyncMock()
|
||||
mock_cache_service.teardown = mock_teardown
|
||||
|
||||
with (
|
||||
patch("langflow.preload.copy_profile_pictures", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_bundles_with_error_handling", new_callable=AsyncMock) as mock_load_bundles,
|
||||
patch("langflow.preload.get_and_cache_all_types_dict", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_flows_from_directory", new_callable=AsyncMock),
|
||||
patch("langflow.preload.get_service") as mock_get_service,
|
||||
patch("langflow.preload.ExternalAsyncBaseCacheService") as mock_base_class,
|
||||
):
|
||||
mock_load_bundles.return_value = ([], [])
|
||||
mock_get_service.return_value = mock_cache_service
|
||||
# Make isinstance check pass
|
||||
mock_base_class.__instancecheck__ = lambda *_args: True
|
||||
|
||||
await _run_master_preload()
|
||||
|
||||
# Verify teardown was called
|
||||
mock_teardown.assert_called_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@patch("langflow.preload.get_db_service")
|
||||
@patch("langflow.preload.get_settings_service")
|
||||
@patch("langflow.preload.initialize_services")
|
||||
@patch("langflow.preload.component_cache")
|
||||
async def test_run_master_preload_agentic_experience_enabled(
|
||||
mock_component_cache,
|
||||
mock_initialize_services,
|
||||
mock_get_settings_service,
|
||||
mock_get_db_service,
|
||||
):
|
||||
"""_run_master_preload() should initialize agentic features when enabled."""
|
||||
mock_engine = AsyncMock()
|
||||
mock_db_service = MagicMock()
|
||||
mock_db_service.engine = mock_engine
|
||||
mock_get_db_service.return_value = mock_db_service
|
||||
|
||||
mock_settings_service = MagicMock()
|
||||
mock_settings_service.settings.agentic_experience = True # Enable agentic
|
||||
mock_settings_service.settings.components_path = []
|
||||
mock_get_settings_service.return_value = mock_settings_service
|
||||
|
||||
mock_initialize_services.return_value = None
|
||||
mock_component_cache.all_types_dict = {"fake": "types_dict"}
|
||||
|
||||
with (
|
||||
patch("langflow.preload.copy_profile_pictures", new_callable=AsyncMock),
|
||||
patch("langflow.preload.load_bundles_with_error_handling", new_callable=AsyncMock) as mock_load_bundles,
|
||||
patch("langflow.preload.get_and_cache_all_types_dict", new_callable=AsyncMock),
|
||||
patch("langflow.preload.create_or_update_starter_projects", new_callable=AsyncMock),
|
||||
patch("langflow.preload.initialize_agentic_global_variables", new_callable=AsyncMock) as mock_init_agentic,
|
||||
patch("langflow.preload.auto_configure_agentic_mcp_server", new_callable=AsyncMock) as mock_auto_config,
|
||||
patch("langflow.preload.load_flows_from_directory", new_callable=AsyncMock),
|
||||
patch("langflow.preload.session_scope") as mock_session_scope,
|
||||
patch("langflow.preload.get_service"),
|
||||
):
|
||||
mock_load_bundles.return_value = ([], [])
|
||||
mock_session = AsyncMock()
|
||||
mock_session_scope.return_value.__aenter__.return_value = mock_session
|
||||
|
||||
await _run_master_preload()
|
||||
|
||||
# Verify agentic functions were called
|
||||
mock_init_agentic.assert_called_once()
|
||||
mock_auto_config.assert_called_once()
|
||||
|
||||
# Verify completion flags
|
||||
assert _STATE.agentic_globals_initialized is True
|
||||
assert _STATE.agentic_mcp_configured is True
|
||||
Reference in New Issue
Block a user