Skip to content

Commit efc8bce

Browse files
Janardan S Kaviaerichare
authored andcommitted
fix(warm-registry): defer v2 host selection to settings (env-file safe); enforce exec timeout on warm sync runs
1 parent 07389ae commit efc8bce

4 files changed

Lines changed: 254 additions & 14 deletions

File tree

src/backend/base/langflow/api/router.py

Lines changed: 8 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,4 @@
11
# Router for base api
2-
import os
3-
42
from fastapi import APIRouter, Depends
53
from lfx.schema.workflow import WORKFLOW_EXECUTION_RESPONSES
64
from lfx.services.settings.feature_flags import FEATURE_FLAGS
@@ -49,8 +47,7 @@
4947
from langflow.api.v2 import registration_router as registration_router_v2
5048
from langflow.api.v2 import workflow_background_router as workflow_background_router_v2
5149
from langflow.api.v2 import workflow_public_router as workflow_public_router_v2
52-
from langflow.api.v2.warm_workflow_host import WarmWorkflowHost
53-
from langflow.api.v2.workflow_host import LangflowWorkflowHost
50+
from langflow.api.v2.host_selection import DeferredWorkflowHost
5451

5552
router_v1 = APIRouter(
5653
prefix="/v1",
@@ -155,15 +152,13 @@ def _include_agentic_router():
155152
# authenticated langflow v2 router has never carried a developer-api gate; the
156153
# default-off setting would otherwise 403 every authenticated request.
157154
# In PROD (``LANGFLOW_PROD``, execution-plane / ``--backend-only``) the warm host
158-
# serves deployed flows from the in-memory registry and skips per-flow RBAC. Read
159-
# the env directly rather than the settings service: this runs at import time,
160-
# before the settings service is guaranteed initialized. The value matches
161-
# ``settings.prod`` (env_prefix ``LANGFLOW_`` maps ``LANGFLOW_PROD`` -> ``prod``).
162-
_PROD = os.getenv("LANGFLOW_PROD", "").strip().lower() in ("1", "true", "yes", "on")
163-
# The one line that picks the runtime: warm host in PROD, DB-backed host otherwise.
164-
# Everything below is identical — only the host object swaps (the point of the seam).
165-
_workflow_host: WorkflowHost = WarmWorkflowHost() if _PROD else LangflowWorkflowHost()
166-
# Startup invariant: whichever host we picked satisfies the WorkflowHost protocol.
155+
# serves deployed flows from the in-memory registry and skips per-flow RBAC. The
156+
# warm-vs-DB choice is made per ``settings.prod`` but DEFERRED to first use: this
157+
# module is imported before ``load_dotenv(--env-file)`` runs, so reading the env
158+
# (or the settings service) here would miss ``--env-file`` values. The route
159+
# structure is identical for both hosts, so binding the deferred proxy changes
160+
# nothing structurally — only which concrete host each request resolves to.
161+
_workflow_host: WorkflowHost = DeferredWorkflowHost()
167162
assert isinstance(_workflow_host, WorkflowHost) # noqa: S101
168163
router_v2.include_router(
169164
create_workflow_router(
Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
"""Deferred selection of the v2 workflow host (warm PROD registry vs DB-backed).
2+
3+
Which host serves ``POST /api/v2/workflows`` depends on ``settings.prod``
4+
(``LANGFLOW_PROD``). The problem: ``langflow.__main__`` imports the router module
5+
(and therefore builds this router) *before* it runs ``load_dotenv(--env-file)``, so
6+
reading the env — or the settings service — at import time would miss any value
7+
supplied via ``--env-file``. That is exactly the ordering trap the extensions router
8+
documents for ``LANGFLOW_ENABLE_EXTENSION_RELOAD``.
9+
10+
``DeferredWorkflowHost`` defers the choice to first use. By the time any workflow
11+
route (or capability flag) is exercised, ``setup_app``/lifespan has initialized the
12+
settings service from the fully-loaded environment, so ``settings.prod`` is correct.
13+
The route *structure* is identical for both hosts (the shared router only reads
14+
``supports_*`` at request time, and ``auto_register_job_routes=False`` neutralizes
15+
the one mount-time read), so binding this single proxy at import changes nothing
16+
structurally — only which concrete host each request lands on.
17+
"""
18+
19+
from __future__ import annotations
20+
21+
from typing import TYPE_CHECKING, Any
22+
23+
from lfx.workflow.host import WorkflowHostBase
24+
25+
if TYPE_CHECKING:
26+
from collections.abc import AsyncIterator
27+
28+
from fastapi import BackgroundTasks, Request, Response
29+
from lfx.schema.workflow import (
30+
ParsedWorkflowRun,
31+
WorkflowExecutionResponse,
32+
WorkflowJobResponse,
33+
WorkflowStopResponse,
34+
)
35+
from lfx.workflow.host import ResolvedFlow, WorkflowAction
36+
37+
38+
class DeferredWorkflowHost(WorkflowHostBase):
39+
"""A ``WorkflowHost`` that picks the concrete host lazily from ``settings.prod``.
40+
41+
Every member delegates to the resolved host. Resolution is cached only once the
42+
settings service is initialized, so an access during module import (before
43+
``--env-file`` is loaded) never force-initializes settings nor caches a stale
44+
choice — it falls back transiently and re-resolves on the next call.
45+
"""
46+
47+
def __init__(self) -> None:
48+
self._host: WorkflowHostBase | None = None
49+
50+
def _resolve(self) -> WorkflowHostBase:
51+
if self._host is not None:
52+
return self._host
53+
from langflow.api.v2.warm_workflow_host import WarmWorkflowHost
54+
from langflow.api.v2.workflow_host import LangflowWorkflowHost
55+
from langflow.services.deps import get_settings_service, is_settings_service_initialized
56+
57+
# Do NOT force-initialize settings from a half-loaded environment at import
58+
# time (the router is built before ``load_dotenv(--env-file)``). Until the
59+
# settings service exists, fall back to the DB host WITHOUT caching so the
60+
# real choice is still made on the first post-startup call.
61+
if not is_settings_service_initialized():
62+
return LangflowWorkflowHost()
63+
host: WorkflowHostBase = WarmWorkflowHost() if get_settings_service().settings.prod else LangflowWorkflowHost()
64+
self._host = host
65+
return host
66+
67+
@property
68+
def supports_background(self) -> bool:
69+
return self._resolve().supports_background
70+
71+
@property
72+
def supports_request_overrides(self) -> bool:
73+
return self._resolve().supports_request_overrides
74+
75+
async def resolve_caller(self, request: Request) -> Any:
76+
return await self._resolve().resolve_caller(request)
77+
78+
async def get_flow(self, flow_id: str, caller: Any) -> ResolvedFlow:
79+
return await self._resolve().get_flow(flow_id, caller)
80+
81+
async def authorize(self, caller: Any, flow: ResolvedFlow, action: WorkflowAction) -> None:
82+
return await self._resolve().authorize(caller, flow, action)
83+
84+
def session(self) -> AsyncIterator[Any | None]:
85+
# Returns the concrete host's async context manager (used as ``async with``).
86+
return self._resolve().session()
87+
88+
async def run_sync(
89+
self,
90+
parsed: ParsedWorkflowRun,
91+
flow: ResolvedFlow,
92+
caller: Any,
93+
*,
94+
http_request: Request,
95+
background_tasks: BackgroundTasks,
96+
) -> WorkflowExecutionResponse:
97+
return await self._resolve().run_sync(
98+
parsed, flow, caller, http_request=http_request, background_tasks=background_tasks
99+
)
100+
101+
def stream_response(
102+
self,
103+
parsed: ParsedWorkflowRun,
104+
flow: ResolvedFlow,
105+
caller: Any,
106+
*,
107+
stream_protocol: str,
108+
http_request: Request,
109+
background_tasks: BackgroundTasks,
110+
) -> Response:
111+
return self._resolve().stream_response(
112+
parsed,
113+
flow,
114+
caller,
115+
stream_protocol=stream_protocol,
116+
http_request=http_request,
117+
background_tasks=background_tasks,
118+
)
119+
120+
async def submit_background(
121+
self,
122+
parsed: ParsedWorkflowRun,
123+
flow: ResolvedFlow,
124+
caller: Any,
125+
*,
126+
stream_protocol: str,
127+
) -> WorkflowJobResponse:
128+
return await self._resolve().submit_background(parsed, flow, caller, stream_protocol=stream_protocol)
129+
130+
async def get_job_status(
131+
self, job_id: str, caller: Any, session: Any
132+
) -> WorkflowExecutionResponse | WorkflowJobResponse:
133+
return await self._resolve().get_job_status(job_id, caller, session)
134+
135+
async def stop_job(self, job_id: str, caller: Any) -> WorkflowStopResponse:
136+
return await self._resolve().stop_job(job_id, caller)

src/backend/base/langflow/api/v2/warm_workflow_host.py

Lines changed: 40 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121

2222
from __future__ import annotations
2323

24+
import asyncio
2425
from copy import deepcopy
2526
from typing import TYPE_CHECKING, Any
2627

@@ -32,7 +33,8 @@
3233
from lfx.workflow.host import ResolvedFlow, WorkflowHostBase
3334

3435
if TYPE_CHECKING:
35-
from fastapi import Request
36+
from fastapi import BackgroundTasks, Request
37+
from lfx.schema.workflow import ParsedWorkflowRun, WorkflowExecutionResponse
3638

3739

3840
class WarmWorkflowHost(WorkflowHostBase):
@@ -113,3 +115,40 @@ async def get_flow(self, flow_id: str, caller: Any) -> ResolvedFlow: # noqa: AR
113115
# Return a run-ready Graph (not a FlowRead) as ResolvedFlow, so the base
114116
# run_sync/stream can execute it directly — no from_payload rebuild on this path.
115117
return ResolvedFlow(flow_id=flow_id, graph=graph_copy, session_id_default=flow_id)
118+
119+
async def run_sync(
120+
self,
121+
parsed: ParsedWorkflowRun,
122+
flow: ResolvedFlow,
123+
caller: Any,
124+
*,
125+
http_request: Request,
126+
background_tasks: BackgroundTasks,
127+
) -> WorkflowExecutionResponse:
128+
"""Run the graph with the langflow wall-clock ceiling enforced (408 on timeout).
129+
130+
The lean base ``run_sync`` has no timeout, so a runaway flow on the execution
131+
plane would run unbounded and hold a worker forever. Wrap it in the configured
132+
``workflow_execution_timeout`` and surface the same 408 contract the DB host
133+
uses. (Only per-flow RBAC + job/vertex-build writes stay dropped; the durable
134+
job-tracking and header-global/session mapping remain a separate adapter.)
135+
"""
136+
from langflow.services.deps import get_settings_service
137+
138+
timeout_seconds = get_settings_service().settings.workflow_execution_timeout
139+
try:
140+
return await asyncio.wait_for(
141+
super().run_sync(parsed, flow, caller, http_request=http_request, background_tasks=background_tasks),
142+
timeout=timeout_seconds,
143+
)
144+
except (TimeoutError, asyncio.TimeoutError):
145+
raise HTTPException(
146+
status_code=status.HTTP_408_REQUEST_TIMEOUT,
147+
detail={
148+
"error": "Execution timeout",
149+
"code": "EXECUTION_TIMEOUT",
150+
"message": f"Workflow execution exceeded {timeout_seconds} seconds",
151+
"flow_id": flow.flow_id,
152+
"timeout_seconds": timeout_seconds,
153+
},
154+
) from None

src/backend/tests/unit/services/warm_registry/test_warm_registry_reconcile.py

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -505,3 +505,73 @@ async def _boom_once():
505505
with contextlib.suppress(asyncio.CancelledError):
506506
await task
507507
assert calls >= 1
508+
509+
510+
# ── #2: deferred host selection (fixes import-time-vs-env-file ordering) ──────
511+
def test_deferred_host_resolves_db_when_not_prod(monkeypatch):
512+
"""settings.prod False -> DB-backed LangflowWorkflowHost."""
513+
from types import SimpleNamespace
514+
515+
from langflow.api.v2.host_selection import DeferredWorkflowHost
516+
from langflow.api.v2.workflow_host import LangflowWorkflowHost
517+
from langflow.services import deps
518+
519+
monkeypatch.setattr(deps, "is_settings_service_initialized", lambda: True)
520+
monkeypatch.setattr(deps, "get_settings_service", lambda: SimpleNamespace(settings=SimpleNamespace(prod=False)))
521+
host = DeferredWorkflowHost()
522+
assert isinstance(host._resolve(), LangflowWorkflowHost)
523+
524+
525+
def test_deferred_host_resolves_warm_when_prod(monkeypatch):
526+
"""settings.prod True -> WarmWorkflowHost, and the choice is cached."""
527+
from types import SimpleNamespace
528+
529+
from langflow.api.v2.host_selection import DeferredWorkflowHost
530+
from langflow.api.v2.warm_workflow_host import WarmWorkflowHost
531+
from langflow.services import deps
532+
533+
monkeypatch.setattr(deps, "is_settings_service_initialized", lambda: True)
534+
monkeypatch.setattr(deps, "get_settings_service", lambda: SimpleNamespace(settings=SimpleNamespace(prod=True)))
535+
host = DeferredWorkflowHost()
536+
resolved = host._resolve()
537+
assert isinstance(resolved, WarmWorkflowHost)
538+
assert host._resolve() is resolved # cached
539+
540+
541+
def test_deferred_host_does_not_cache_before_settings_init(monkeypatch):
542+
"""An access before settings are initialized (import time) must not cache a choice."""
543+
from langflow.api.v2.host_selection import DeferredWorkflowHost
544+
from langflow.api.v2.workflow_host import LangflowWorkflowHost
545+
from langflow.services import deps
546+
547+
monkeypatch.setattr(deps, "is_settings_service_initialized", lambda: False)
548+
host = DeferredWorkflowHost()
549+
resolved = host._resolve()
550+
assert isinstance(resolved, LangflowWorkflowHost) # transient DB fallback
551+
assert host._host is None # NOT cached — real choice deferred to first post-startup call
552+
553+
554+
# ── #5: warm host enforces the workflow_execution_timeout (408) on sync runs ──
555+
async def test_warm_host_run_sync_enforces_timeout(monkeypatch):
556+
"""A sync run exceeding workflow_execution_timeout -> HTTP 408 (the lean base has none)."""
557+
from types import SimpleNamespace
558+
559+
import lfx.workflow.router as lfx_router
560+
from fastapi import HTTPException
561+
from langflow.api.v2.warm_workflow_host import WarmWorkflowHost
562+
from langflow.services import deps
563+
564+
async def _slow(*_args, **_kwargs):
565+
await asyncio.sleep(1)
566+
567+
monkeypatch.setattr(lfx_router, "run_workflow_sync", _slow)
568+
monkeypatch.setattr(
569+
deps, "get_settings_service", lambda: SimpleNamespace(settings=SimpleNamespace(workflow_execution_timeout=0.01))
570+
)
571+
572+
host = WarmWorkflowHost()
573+
flow = SimpleNamespace(graph=object(), flow_id="fid")
574+
with pytest.raises(HTTPException) as exc:
575+
await host.run_sync(SimpleNamespace(), flow, None, http_request=None, background_tasks=None)
576+
assert exc.value.status_code == 408
577+
assert exc.value.detail["code"] == "EXECUTION_TIMEOUT"

0 commit comments

Comments
 (0)