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:
Arek Mateusiak
2026-04-21 19:51:00 +02:00
parent a07164d1ae
commit c0e81a5401
4 changed files with 543 additions and 19 deletions

View File

@ -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

View File

@ -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

View File

@ -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"

View 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