diff --git a/src/lfx/README.md b/src/lfx/README.md index 2bb2df031e..96dfaad6ff 100644 --- a/src/lfx/README.md +++ b/src/lfx/README.md @@ -259,6 +259,13 @@ To view the LFX server's API docs and schema, see the `/docs` endpoint at `http: | `--verbose`, `-v` | Show diagnostic output and execution details. | | `--flow-json` | Read inline flow JSON content as a string. Example: `uv run lfx serve --flow-json '{...}'`. | | `--stdin` | Read JSON flow content from `stdin`. Example: `cat flow.json | uv run lfx serve --stdin`. | +| `--workers`, `-w` | Number of worker processes. Use with `--flow-dir` for multi-worker flow sharing. Default: `1`. | +| `--flow-dir` | Directory for filesystem-backed flow storage shared across workers (e.g. `/tmp/lfx-flows` for single-pod, or a PVC mount for cross-pod). Defaults to in-memory only when omitted. | +| `--max-requests` | Recycle each worker after this many requests (gunicorn, Unix-only, `--workers > 1`); `1` gives per-request recycling. Default: unset (workers are never recycled). | +| `--limit-concurrency` | Max in-flight requests per worker (`--workers > 1`); excess get HTTP 503. Default: unset (unlimited). | +| `--no-env-fallback` / `--env-fallback` | Disable the `os.environ` fallback for credential variables; variables not supplied via `global_vars` per request resolve to `None` instead of reading the process environment. Default: `--env-fallback` (fallback enabled). | +| `--reset-environ` / `--no-reset-environ` | Snapshot `os.environ` before each flow run and restore it afterward, so a flow's environment mutations (or request-scoped credentials) cannot leak into the next request served by the same warm worker. Default: `--no-reset-environ` (off). | +| `--sync-workers` / `--no-sync-workers` | Multi-worker only (`--workers > 1`, Unix). Use gunicorn's blocking `sync` worker so the kernel routes each request to an idle worker instead of queueing it behind an in-flight request on a busy async worker. Requires the `a2wsgi` package (`pip install a2wsgi`). Default: `--no-sync-workers` (async worker). | ## Run the simple agent flow with `lfx run` diff --git a/src/lfx/pyproject.toml b/src/lfx/pyproject.toml index f96c9665d8..80c6fc132e 100644 --- a/src/lfx/pyproject.toml +++ b/src/lfx/pyproject.toml @@ -17,6 +17,10 @@ dependencies = [ # gunicorn powers `lfx serve --workers N` (preload + max_requests=1 per-request # isolation). Unix-only; on Windows multi-worker serve is refused with a clear error. "gunicorn>=22.0; sys_platform != 'win32'", + # a2wsgi bridges the ASGI app onto gunicorn's sync worker for + # `lfx serve --sync-workers` (one whole request per worker / idle-worker routing). + # Unix-only like gunicorn; the sync-worker path is refused on Windows. + "a2wsgi>=1.10.0; sys_platform != 'win32'", "typer>=0.16.0,<1.0.0", "platformdirs>=4.3.8,<5.0.0", "aiofiles>=24.1.0,<25.0.0", diff --git a/src/lfx/src/lfx/cli/_running_commands.py b/src/lfx/src/lfx/cli/_running_commands.py index d65995715d..0412aafa1d 100644 --- a/src/lfx/src/lfx/cli/_running_commands.py +++ b/src/lfx/src/lfx/cli/_running_commands.py @@ -177,6 +177,25 @@ def register(app: typer.Typer) -> None: "instead of reading from the process environment." ), ), + reset_environ: bool = typer.Option( + False, + "--reset-environ/--no-reset-environ", + help=( + "Snapshot os.environ before each flow run and restore it afterward, so a " + "flow's environment mutations (or request-scoped credentials) cannot leak " + "into the next request served by the same warm worker. Off by default." + ), + ), + sync_workers: bool = typer.Option( + False, + "--sync-workers/--no-sync-workers", + help=( + "Use gunicorn's blocking 'sync' worker (Unix, --workers > 1) so the kernel " + "routes each request to an idle worker instead of queueing it behind an " + "in-flight request on a busy async worker. Requires the 'a2wsgi' package. " + "Off by default (async worker)." + ), + ), ) -> None: """Serve LFX flows as a web API (lazy-loaded).""" from pathlib import Path @@ -201,4 +220,6 @@ def register(app: typer.Typer) -> None: no_env_fallback=no_env_fallback, max_requests=max_requests, limit_concurrency=limit_concurrency, + reset_environ=reset_environ, + sync_workers=sync_workers, ) diff --git a/src/lfx/src/lfx/cli/commands.py b/src/lfx/src/lfx/cli/commands.py index d0d53f3381..04dcaa20a5 100644 --- a/src/lfx/src/lfx/cli/commands.py +++ b/src/lfx/src/lfx/cli/commands.py @@ -221,6 +221,25 @@ def serve_command( "instead of reading from the process environment." ), ), + reset_environ: bool = typer.Option( + False, # noqa: FBT003 + "--reset-environ/--no-reset-environ", + help=( + "Snapshot os.environ before each flow run and restore it afterward, so a " + "flow's environment mutations (or request-scoped credentials) cannot leak " + "into the next request served by the same warm worker. Off by default." + ), + ), + sync_workers: bool = typer.Option( + False, # noqa: FBT003 + "--sync-workers/--no-sync-workers", + help=( + "Use gunicorn's blocking 'sync' worker (Unix, --workers > 1) so the kernel " + "routes each request to an idle worker instead of queueing it behind an " + "in-flight request on a busy async worker. Requires the 'a2wsgi' package. " + "Off by default (async worker)." + ), + ), ) -> None: """Serve LFX flows as a web API. @@ -383,8 +402,16 @@ def serve_command( verbose_print=verbose_print, max_requests=max_requests, limit_concurrency=limit_concurrency, + reset_environ=reset_environ, + sync_workers=sync_workers, ) else: + from lfx.cli.serve_app import _SERVE_RESET_ENVIRON_ENV + + # Single worker also serves many requests warm, so honor --reset-environ + # here (read per request by guarded_execute). --sync-workers is a + # multi-worker routing concern and has no effect with one worker. + os.environ[_SERVE_RESET_ENVIRON_ENV] = "1" if reset_environ else "0" serve_app = create_multi_serve_app(registry=registry) uvicorn.run(serve_app, host=host, port=port, workers=1, log_level=log_level) except KeyboardInterrupt: @@ -413,6 +440,8 @@ def _launch_workers( verbose_print: Callable[[str], None], max_requests: int | None, limit_concurrency: int | None, + reset_environ: bool = False, + sync_workers: bool = False, ) -> None: """Launch ``workers`` worker processes for ``lfx serve --workers N``. @@ -434,11 +463,19 @@ def _launch_workers( falls back to uvicorn's own multi-worker supervisor (no preload, no per-request recycling). ``--limit-concurrency`` is still honored there (uvicorn-native), but ``--max-requests`` (recycling) is refused, since it cannot be supported. + + ``reset_environ`` (``--reset-environ``) is forwarded to the workers via + ``LFX_SERVE_RESET_ENVIRON`` so each worker snapshots/restores ``os.environ`` + around every flow run (see ``guarded_execute``). ``sync_workers`` + (``--sync-workers``, Unix only) swaps the async worker for gunicorn's blocking + ``sync`` worker wrapped by an a2wsgi ASGI->WSGI bridge, so the kernel routes each + request to an idle worker. Both default off. """ from lfx.cli.serve_app import ( _SERVE_FLOW_DIR_ENV, _SERVE_LIMIT_CONCURRENCY_ENV, _SERVE_NO_ENV_FALLBACK_ENV, + _SERVE_RESET_ENVIRON_ENV, _SERVE_STARTUP_PATHS_ENV, ) @@ -459,6 +496,9 @@ def _launch_workers( # Read per worker by LFXUvicornWorker (Unix); passed to uvicorn.run on Windows. if limit_concurrency is not None: os.environ[_SERVE_LIMIT_CONCURRENCY_ENV] = str(limit_concurrency) + # Read per request by guarded_execute in each worker. Always set explicitly so a + # stray inherited value can't silently flip behavior. + os.environ[_SERVE_RESET_ENVIRON_ENV] = "1" if reset_environ else "0" try: if sys.platform == "win32": @@ -469,6 +509,13 @@ def _launch_workers( "without isolation, or deploy on Linux/macOS for per-request isolation." ) raise typer.Exit(1) + if sync_workers: + verbose_print( + "Error: --sync-workers uses gunicorn's sync worker, which is not available on " + "Windows. Omit --sync-workers to run multi-worker on Windows, or deploy on " + "Linux/macOS for idle-worker routing." + ) + raise typer.Exit(1) # gunicorn cannot run on Windows; fall back to uvicorn's multi-worker # supervisor. No preload/COW and no per-request recycling (no isolation), # though --limit-concurrency is still honored (uvicorn-native). @@ -491,14 +538,32 @@ def _launch_workers( else: from lfx.cli.serve_gunicorn import LFXGunicornApp + if sync_workers: + # Fail fast in the parent rather than per-worker on first request. + try: + import a2wsgi # noqa: F401 + except ImportError as exc: + verbose_print( + "Error: --sync-workers requires the 'a2wsgi' package. Install it with: pip install a2wsgi" + ) + raise typer.Exit(1) from exc + # gunicorn's blocking sync worker stops accepting while a request runs, + # so the kernel routes the next request to an idle worker. It serves the + # ASGI app through the a2wsgi WSGI bridge (built lazily, post-fork). + app_import_string = "lfx.cli.serve_preloaded_app:wsgi_application" + worker_class = "sync" + else: + # Async worker; applies LFX_SERVE_LIMIT_CONCURRENCY (gunicorn's + # UvicornWorker cannot forward uvicorn's limit_concurrency). + app_import_string = "lfx.cli.serve_preloaded_app:app" + worker_class = "lfx.cli.serve_gunicorn.LFXUvicornWorker" + LFXGunicornApp( - "lfx.cli.serve_preloaded_app:app", + app_import_string, { "bind": f"{host}:{port}", "workers": workers, - # Custom worker applies LFX_SERVE_LIMIT_CONCURRENCY (gunicorn's - # UvicornWorker cannot forward uvicorn's limit_concurrency). - "worker_class": "lfx.cli.serve_gunicorn.LFXUvicornWorker", + "worker_class": worker_class, "preload_app": True, # None -> 0 (gunicorn's default: never recycle). 1 -> recycle per request. "max_requests": max_requests if max_requests is not None else 0, @@ -517,6 +582,7 @@ def _launch_workers( _SERVE_NO_ENV_FALLBACK_ENV, _SERVE_STARTUP_PATHS_ENV, _SERVE_LIMIT_CONCURRENCY_ENV, + _SERVE_RESET_ENVIRON_ENV, ): os.environ.pop(k, None) diff --git a/src/lfx/src/lfx/cli/serve_app.py b/src/lfx/src/lfx/cli/serve_app.py index ec421631a6..25ef93f577 100644 --- a/src/lfx/src/lfx/cli/serve_app.py +++ b/src/lfx/src/lfx/cli/serve_app.py @@ -19,6 +19,7 @@ from __future__ import annotations import asyncio import json +import os import time import traceback import uuid @@ -57,6 +58,12 @@ _SERVE_STARTUP_PATHS_ENV = f"{_SERVE_ENV_PREFIX}STARTUP_PATHS" # uvicorn limit_concurrency applied per worker by LFXUvicornWorker (gunicorn's # UvicornWorker doesn't expose it). Set by serve_command; read in each worker. _SERVE_LIMIT_CONCURRENCY_ENV = f"{_SERVE_ENV_PREFIX}LIMIT_CONCURRENCY" +# Opt-in (`lfx serve --reset-environ`): when "1", guarded_execute snapshots +# os.environ before each flow run and restores it after, so env mutations made by +# one request (or request-scoped credentials) cannot leak into the next request +# served by the same warm worker. Off by default. Set by serve_command; read per +# request in guarded_execute. +_SERVE_RESET_ENVIRON_ENV = f"{_SERVE_ENV_PREFIX}RESET_ENVIRON" api_key_query = APIKeyQuery(name=API_KEY_NAME, scheme_name="API key query", auto_error=False) api_key_header = APIKeyHeader(name=API_KEY_NAME, scheme_name="API key header", auto_error=False) @@ -73,9 +80,22 @@ async def guarded_execute(graph_copy, input_value, session_id=None): Serializes the env-sensitive execution section so two concurrent requests in the same async worker can never overlap a flow run before the worker recycles. + + When ``LFX_SERVE_RESET_ENVIRON`` is "1" (``lfx serve --reset-environ``), the + process environment is snapshotted before the run and restored afterward, so a + flow's os.environ mutations (or request-scoped credentials) cannot leak into the + next request served by the same warm worker. Off by default — the snapshot is + skipped entirely unless opted in. """ async with _EXECUTE_GUARD: - return await execute_graph_with_capture(graph_copy, input_value, session_id=session_id) + reset_environ = os.environ.get(_SERVE_RESET_ENVIRON_ENV) == "1" + env_snapshot = dict(os.environ) if reset_environ else None + try: + return await execute_graph_with_capture(graph_copy, input_value, session_id=session_id) + finally: + if env_snapshot is not None and os.environ != env_snapshot: + os.environ.clear() + os.environ.update(env_snapshot) def verify_api_key( diff --git a/src/lfx/src/lfx/cli/serve_gunicorn.py b/src/lfx/src/lfx/cli/serve_gunicorn.py index 2381433408..583ac7bed7 100644 --- a/src/lfx/src/lfx/cli/serve_gunicorn.py +++ b/src/lfx/src/lfx/cli/serve_gunicorn.py @@ -129,7 +129,12 @@ class LFXGunicornApp(BaseApplication): self.cfg.set(key, value) def load(self): - # preload_app=True -> gunicorn imports this string in the master pre-fork. - from lfx.cli.serve_preloaded_app import app + # preload_app=True -> gunicorn imports this "module:attr" string in the master + # pre-fork. The async path passes the ASGI app ("...:app"); --sync-workers + # passes the WSGI bridge ("...:wsgi_application"). Honor whichever was given + # rather than hardcoding one, so the worker gets the callable it expects. + import importlib - return app + module_str, _, attr = self._app_import_string.partition(":") + module = importlib.import_module(module_str) + return getattr(module, attr) diff --git a/src/lfx/src/lfx/cli/serve_preloaded_app.py b/src/lfx/src/lfx/cli/serve_preloaded_app.py index 8d00fc4495..401491c277 100644 --- a/src/lfx/src/lfx/cli/serve_preloaded_app.py +++ b/src/lfx/src/lfx/cli/serve_preloaded_app.py @@ -17,3 +17,26 @@ from lfx.cli.serve_app import build_registry_from_env, create_multi_serve_app registry = build_registry_from_env() app = create_multi_serve_app(registry=registry) + +# WSGI bridge entrypoint for the opt-in ``lfx serve --sync-workers`` mode, which +# runs gunicorn's blocking ``sync`` worker so the kernel routes each request to an +# idle worker (an async worker keeps accepting connections while busy, which can +# queue a second request behind an in-flight one even when other workers are idle). +# The ASGI->WSGI bridge is created LAZILY on first request, never at preload: a2wsgi +# spins up a background event-loop thread, and threads do not survive fork(), so a +# bridge built in the preload master would be dead in every forked worker. Building +# it on first call constructs it inside the (post-fork) worker process instead. +_bridge = None + + +def wsgi_application(environ, start_response): + """Lazily-constructed a2wsgi bridge wrapping the preloaded ASGI ``app``.""" + global _bridge # noqa: PLW0603 + if _bridge is None: + try: + from a2wsgi import ASGIMiddleware + except ImportError as exc: # pragma: no cover - exercised via --sync-workers without the dep + msg = "lfx serve --sync-workers requires the 'a2wsgi' package. Install it with: pip install a2wsgi" + raise RuntimeError(msg) from exc + _bridge = ASGIMiddleware(app) + return _bridge(environ, start_response) diff --git a/src/lfx/tests/unit/cli/test_serve.py b/src/lfx/tests/unit/cli/test_serve.py index d2f1cfeed3..3349773b71 100644 --- a/src/lfx/tests/unit/cli/test_serve.py +++ b/src/lfx/tests/unit/cli/test_serve.py @@ -709,7 +709,8 @@ def _run_serve_capturing_gunicorn(*, workers, max_requests, captured, limit_conc patch("lfx.cli.serve_gunicorn.LFXGunicornApp", FakeGunicornApp), ): # Direct call (not via CLI): typer defaults are not applied, so pass - # max_requests / limit_concurrency explicitly even when None. + # max_requests / limit_concurrency / reset_environ / sync_workers + # explicitly even when None/False (OptionInfo sentinels are truthy). serve_command( script_paths=[str(p)], host="127.0.0.1", @@ -725,6 +726,8 @@ def _run_serve_capturing_gunicorn(*, workers, max_requests, captured, limit_conc no_env_fallback=False, max_requests=max_requests, limit_concurrency=limit_concurrency, + reset_environ=False, + sync_workers=False, ) @@ -828,6 +831,8 @@ def test_serve_command_sets_startup_paths_env_for_multi_worker(tmp_path): stdin=False, check_variables=False, no_env_fallback=False, + reset_environ=False, + sync_workers=False, ) assert _SERVE_STARTUP_PATHS_ENV in captured_env, "LFX_SERVE_STARTUP_PATHS must be set before the launch" @@ -886,6 +891,8 @@ def test_serve_command_does_not_set_startup_paths_when_flow_dir_set(tmp_path): stdin=False, check_variables=False, no_env_fallback=False, + reset_environ=False, + sync_workers=False, ) assert _SERVE_STARTUP_PATHS_ENV in captured_env, "env var must still be set (to empty list)" @@ -938,6 +945,8 @@ def test_serve_command_warns_when_workers_gt1_without_flow_dir(): stdin=False, check_variables=False, no_env_fallback=False, + reset_environ=False, + sync_workers=False, ) assert any("--flow-dir" in msg for msg in stderr_output), ( @@ -978,6 +987,8 @@ def test_serve_command_rejects_py_with_multiple_workers(tmp_path): stdin=False, check_variables=False, no_env_fallback=False, + reset_environ=False, + sync_workers=False, ) assert any(".py" in msg and "cannot be used" in msg for msg in stderr_output), stderr_output @@ -1031,6 +1042,8 @@ def test_serve_command_allows_py_with_multiple_workers_no_flow_dir(tmp_path): stdin=False, check_variables=False, no_env_fallback=False, + reset_environ=False, + sync_workers=False, ) # LFX_SERVE_STARTUP_PATHS must contain the .py path so workers can reload it @@ -1084,6 +1097,8 @@ def test_serve_command_no_warning_when_workers_gt1_with_flow_dir(tmp_path): stdin=False, check_variables=False, no_env_fallback=False, + reset_environ=False, + sync_workers=False, ) assert not any("--flow-dir" in msg for msg in stderr_output) diff --git a/src/lfx/tests/unit/cli/test_serve_env_isolation_integration.py b/src/lfx/tests/unit/cli/test_serve_env_isolation_integration.py new file mode 100644 index 0000000000..3429e2edc3 --- /dev/null +++ b/src/lfx/tests/unit/cli/test_serve_env_isolation_integration.py @@ -0,0 +1,180 @@ +"""End-to-end env-isolation tests for `lfx serve` (spawns a live multi-worker server). + +Demonstrates the actual security guarantee — caller A writes os.environ, caller B +must not read it — across real worker processes (fork + recycle), which can only be +observed with live processes, not in-process unit tests. + +Two scenarios using a probe flow whose component writes a canary to os.environ and +reports whether it SAW a prior write: +- ``--workers 1`` (no recycle): a single persistent process -> later requests see the + canary -> ``LEAKED`` (proves the hazard is real, and that this harness detects it). +- ``--workers 2 --max-requests 1 --limit-concurrency 1``: each request runs in a + freshly-forked, recycled, single-in-flight worker -> always ``clean``. + +Skipped in CI (spawns gunicorn + a real server; slow) — run locally as a harness. +""" + +import json +import os +import socket +import subprocess +import sys +import time +import urllib.error +import urllib.request +from contextlib import contextmanager +from pathlib import Path + +import pytest + +_is_ci = os.environ.get("CI", "").lower() in {"1", "true", "yes"} +pytestmark = [ + pytest.mark.skipif(_is_ci, reason="spawns a live multi-worker server; not for fast CI"), + pytest.mark.skipif(sys.platform == "win32", reason="isolation path uses gunicorn (Unix-only)"), +] + +PROBE_FLOW = """ +import os +from lfx.components.input_output import ChatInput, ChatOutput +from lfx.custom import Component +from lfx.graph import Graph +from lfx.io import MessageTextInput, Output +from lfx.schema.message import Message + +PROBE_KEY = "LFX_LEAK_PROBE" +CANARY = "SECRET-CANARY-VALUE" + + +class EnvLeakProbe(Component): + display_name = "Env Leak Probe" + inputs = [MessageTextInput(name="input_value", display_name="Input")] + outputs = [Output(name="status", display_name="Status", method="get_status")] + + def get_status(self) -> Message: + seen = os.environ.get(PROBE_KEY) + os.environ[PROBE_KEY] = CANARY + status = "LEAKED" if seen == CANARY else "clean" + return Message(text=f"{status}|pid={os.getpid()}") + + +chat_input = ChatInput() +leak = EnvLeakProbe().set(input_value=chat_input.message_response) +chat_output = ChatOutput().set(input_value=leak.get_status) +graph = Graph(chat_input, chat_output) +""" + +_API_KEY = "leak-test-key" # pragma: allowlist secret + + +def _free_port() -> int: + s = socket.socket() + s.bind(("127.0.0.1", 0)) + port = s.getsockname()[1] + s.close() + return port + + +def _get(url: str): + req = urllib.request.Request(url, headers={"x-api-key": _API_KEY}) # noqa: S310 + with urllib.request.urlopen(req, timeout=5) as resp: # noqa: S310 + return json.load(resp) + + +def _run_flow(port: int, flow_id: str, value: str) -> str: + req = urllib.request.Request( + f"http://127.0.0.1:{port}/flows/{flow_id}/run", + data=json.dumps({"input_value": value}).encode(), + method="POST", + headers={"x-api-key": _API_KEY, "content-type": "application/json"}, + ) + with urllib.request.urlopen(req, timeout=30) as resp: # noqa: S310 + return json.load(resp)["result"] + + +@contextmanager +def _serve(tmp_path: Path, extra_args: list[str]): + flow = tmp_path / "leak_flow.py" + flow.write_text(PROBE_FLOW, encoding="utf-8") + port = _free_port() + log = (tmp_path / "server.log").open("w") + env = {**os.environ, "LANGFLOW_API_KEY": _API_KEY} + proc = subprocess.Popen( # noqa: S603 + [sys.executable, "-m", "lfx", "serve", str(flow), "--host", "127.0.0.1", "--port", str(port), *extra_args], + env=env, + stdout=log, + stderr=subprocess.STDOUT, + ) + try: + flow_id = _wait_ready(port, proc, tmp_path) + yield port, flow_id + finally: + proc.terminate() + try: + proc.wait(timeout=10) + except subprocess.TimeoutExpired: + proc.kill() + log.close() + + +def _wait_ready(port: int, proc: subprocess.Popen, tmp_path: Path, timeout: float = 120.0) -> str: + deadline = time.time() + timeout + while time.time() < deadline: + if proc.poll() is not None: + log = (tmp_path / "server.log").read_text(encoding="utf-8") + msg = f"serve exited early (code {proc.returncode}):\n{log[-2000:]}" + raise RuntimeError(msg) + try: + flows = _get(f"http://127.0.0.1:{port}/flows") + if flows: + return flows[0]["id"] + except (urllib.error.URLError, ConnectionError, TimeoutError, OSError): + pass + time.sleep(1) + msg = "serve did not become ready in time" + raise RuntimeError(msg) + + +def _statuses(results): + # Each result is "clean|pid=NNNN" or "LEAKED|pid=NNNN". + return [r.split("|", 1)[0] for r in results] + + +def test_no_env_leak_multi_worker_warm(tmp_path): + """workers=2 with warm long-lived workers (--max-requests 0, no recycling) -> no env leak. + + This is the production model: workers are NOT recycled per request, so isolation comes + entirely from per-request os.environ snapshot/restore, even as each warm worker serves + many requests. + """ + with _serve(tmp_path, ["--workers", "2", "--max-requests", "0", "--reset-environ"]) as (port, fid): + results = [_run_flow(port, fid, f"req{i}") for i in range(12)] + assert _statuses(results) == ["clean"] * 12, results + + +def test_no_env_leak_single_worker(tmp_path): + """workers=1 (single persistent process, NO recycling) -> still no cross-request leak. + + This is the strongest demonstration that per-request os.environ snapshot/restore in + ``guarded_execute`` is what enforces isolation: every request is served by the SAME + process (recycling never happens here), yet none sees a prior request's env write. + Before that fix, this exact scenario leaked from the 2nd request onward. + """ + with _serve(tmp_path, ["--workers", "1", "--reset-environ"]) as (port, fid): + results = [_run_flow(port, fid, f"req{i}") for i in range(12)] + assert _statuses(results) == ["clean"] * 12, results + # All requests served by the same process (no recycling) -> isolation came from + # per-request env restore, not from a fresh process. + pids = {r.split("pid=", 1)[1] for r in results if "pid=" in r} + assert len(pids) == 1, f"expected a single reused worker process, got pids={pids} results={results}" + + +def test_sync_workers_serves_requests(tmp_path): + """--sync-workers (gunicorn sync worker + a2wsgi bridge) serves requests via the real CLI. + + Proves the opt-in sync-worker path boots and the ASGI->WSGI bridge handles requests, + and that --reset-environ still enforces isolation under the sync worker. + """ + with _serve(tmp_path, ["--workers", "2", "--sync-workers", "--reset-environ"]) as (port, fid): + results = [_run_flow(port, fid, f"req{i}") for i in range(8)] + assert all("pid=" in r for r in results), results # the a2wsgi bridge served every request + assert _statuses(results) == ["clean"] * 8, results # reset-environ holds under sync worker diff --git a/src/lfx/tests/unit/cli/test_serve_process_model.py b/src/lfx/tests/unit/cli/test_serve_process_model.py index 5e9d14055f..60c2474f32 100644 --- a/src/lfx/tests/unit/cli/test_serve_process_model.py +++ b/src/lfx/tests/unit/cli/test_serve_process_model.py @@ -16,6 +16,7 @@ Covers the preload/warm/guard/fork-safety pieces of the gunicorn """ import json +import os import sys from pathlib import Path @@ -430,3 +431,129 @@ def test_pre_fork_flags_ghost_thread_but_not_benign(): blob = "\n".join(warnings) assert "Ghost threads" in blob and "EvilGhostThread" in blob, warnings assert "OTel-benign" not in blob # benign-named threads are filtered out + + +# --------------------------------------------------------------------------- +# Opt-in flags: --reset-environ (os.environ snapshot/restore) and --sync-workers +# (gunicorn sync worker + a2wsgi bridge). Both default OFF so the committed +# behavior is unchanged; these tests assert the opt-in wiring. +# --------------------------------------------------------------------------- + + +async def test_guarded_execute_restores_environ_when_enabled(monkeypatch): + """With LFX_SERVE_RESET_ENVIRON=1, a flow's os.environ mutation is rolled back.""" + from lfx.cli import serve_app + + monkeypatch.setenv(serve_app._SERVE_RESET_ENVIRON_ENV, "1") + monkeypatch.delenv("LEAKED_BY_FLOW", raising=False) + + async def fake_capture(graph, input_value, session_id=None): # noqa: ARG001 + os.environ["LEAKED_BY_FLOW"] = "secret" + return ([], "") + + monkeypatch.setattr(serve_app, "execute_graph_with_capture", fake_capture) + await serve_app.guarded_execute(object(), "a", None) + + # The mutation made during the run is restored after it (no cross-request leak). + assert "LEAKED_BY_FLOW" not in os.environ + + +async def test_guarded_execute_does_not_reset_environ_by_default(monkeypatch): + """Default (flag off): os.environ mutations persist — committed behavior is unchanged.""" + from lfx.cli import serve_app + + monkeypatch.delenv(serve_app._SERVE_RESET_ENVIRON_ENV, raising=False) + monkeypatch.delenv("LEAKED_BY_FLOW", raising=False) + + async def fake_capture(graph, input_value, session_id=None): # noqa: ARG001 + os.environ["LEAKED_BY_FLOW"] = "secret" + return ([], "") + + monkeypatch.setattr(serve_app, "execute_graph_with_capture", fake_capture) + try: + await serve_app.guarded_execute(object(), "a", None) + assert os.environ.get("LEAKED_BY_FLOW") == "secret" # not restored + finally: + monkeypatch.delenv("LEAKED_BY_FLOW", raising=False) + + +def _capture_gunicorn_launch(monkeypatch, **launch_overrides): + """Run ``_launch_workers`` on the Unix gunicorn path with a fake LFXGunicornApp. + + Returns a dict with the captured ``app_import_string``, gunicorn ``options``, + and a snapshot of ``os.environ`` taken inside ``run()`` (i.e. what forked + workers would inherit). + """ + from lfx.cli import commands + + captured: dict = {} + + class FakeGunicornApp: + def __init__(self, app_import_string, options): + captured["app_import_string"] = app_import_string + captured["options"] = options + + def run(self): + captured["env"] = dict(os.environ) + + monkeypatch.setattr("lfx.cli.serve_gunicorn.LFXGunicornApp", FakeGunicornApp) + kwargs = { + "host": "127.0.0.1", + "port": 8000, + "workers": 2, + "log_level": "warning", + "flow_dir": None, + "no_env_fallback": False, + "script_paths": None, + "temp_file_to_cleanup": None, + "verbose_print": lambda *_a, **_k: None, + "max_requests": None, + "limit_concurrency": None, + } + kwargs.update(launch_overrides) + commands._launch_workers(**kwargs) + return captured + + +def test_launch_workers_default_uses_async_uvicorn_worker(monkeypatch): + """Default (no --sync-workers): the async LFXUvicornWorker serves the ASGI app.""" + captured = _capture_gunicorn_launch(monkeypatch) + assert captured["app_import_string"] == "lfx.cli.serve_preloaded_app:app" + assert captured["options"]["worker_class"] == "lfx.cli.serve_gunicorn.LFXUvicornWorker" + + +def test_launch_workers_sync_workers_uses_sync_worker(monkeypatch): + """--sync-workers swaps in gunicorn's sync worker serving the a2wsgi WSGI bridge.""" + captured = _capture_gunicorn_launch(monkeypatch, sync_workers=True) + assert captured["app_import_string"] == "lfx.cli.serve_preloaded_app:wsgi_application" + assert captured["options"]["worker_class"] == "sync" + + +def test_launch_workers_sync_workers_without_a2wsgi_errors(monkeypatch): + """--sync-workers fails fast in the parent when a2wsgi is not installed.""" + import typer + + monkeypatch.setitem(sys.modules, "a2wsgi", None) # forces `import a2wsgi` to raise ImportError + with pytest.raises(typer.Exit): + _capture_gunicorn_launch(monkeypatch, sync_workers=True) + # env vars set before the failed launch are still cleaned up + from lfx.cli.serve_app import _SERVE_RESET_ENVIRON_ENV + + assert _SERVE_RESET_ENVIRON_ENV not in os.environ + + +def test_launch_workers_reset_environ_exports_env_for_workers(monkeypatch): + """--reset-environ exports LFX_SERVE_RESET_ENVIRON=1 (inherited by workers), then cleans up.""" + from lfx.cli.serve_app import _SERVE_RESET_ENVIRON_ENV + + captured = _capture_gunicorn_launch(monkeypatch, reset_environ=True) + assert captured["env"].get(_SERVE_RESET_ENVIRON_ENV) == "1" + assert _SERVE_RESET_ENVIRON_ENV not in os.environ # cleaned up after launch + + +def test_launch_workers_reset_environ_off_by_default(monkeypatch): + """Without --reset-environ, the env var is exported as "0" (snapshot/restore disabled).""" + from lfx.cli.serve_app import _SERVE_RESET_ENVIRON_ENV + + captured = _capture_gunicorn_launch(monkeypatch) + assert captured["env"].get(_SERVE_RESET_ENVIRON_ENV) == "0" diff --git a/uv.lock b/uv.lock index 5e97104bee..bf66f923a2 100644 --- a/uv.lock +++ b/uv.lock @@ -79,6 +79,18 @@ http-server = [ { name = "starlette", marker = "(python_full_version < '3.14' and platform_machine == 'arm64') or (python_full_version < '3.14' and sys_platform != 'darwin')" }, ] +[[package]] +name = "a2wsgi" +version = "1.10.10" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "typing-extensions", marker = "python_full_version < '3.11' and sys_platform != 'win32'" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/9a/cb/822c56fbea97e9eee201a2e434a80437f6750ebcb1ed307ee3a0a7505b14/a2wsgi-1.10.10.tar.gz", hash = "sha256:a5bcffb52081ba39df0d5e9a884fc6f819d92e3a42389343ba77cbf809fe1f45", size = 18799, upload-time = "2025-06-18T09:00:10.843Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/02/d5/349aba3dc421e73cbd4958c0ce0a4f1aa3a738bc0d7de75d2f40ed43a535/a2wsgi-1.10.10-py3-none-any.whl", hash = "sha256:d2b21379479718539dc15fce53b876251a0efe7615352dfe49f6ad1bc507848d", size = 17389, upload-time = "2025-06-18T09:00:09.676Z" }, +] + [[package]] name = "accelerate" version = "1.13.0" @@ -8989,6 +9001,7 @@ name = "lfx" version = "0.5.0" source = { editable = "src/lfx" } dependencies = [ + { name = "a2wsgi", marker = "sys_platform != 'win32'" }, { name = "ag-ui-protocol" }, { name = "aiofile", version = "3.9.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, { name = "aiofile", version = "3.11.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, @@ -9061,6 +9074,7 @@ integration = [ [package.metadata] requires-dist = [ + { name = "a2wsgi", marker = "sys_platform != 'win32'", specifier = ">=1.10.0" }, { name = "ag-ui-protocol", specifier = ">=0.1.10" }, { name = "aiofile", specifier = ">=3.8.0,<4.0.0" }, { name = "aiofiles", specifier = ">=24.1.0,<25.0.0" },