From b8fe970493fd5fb2e1fc71dfccc79f90b76058fa Mon Sep 17 00:00:00 2001 From: Eric Hare Date: Tue, 21 Apr 2026 13:06:13 -0700 Subject: [PATCH] fix(mcp): close path traversal + cross-user disclosure (PVR0754098) (#12818) * fix(mcp): close path traversal + cross-user disclosure in MCP endpoint (PVR0754098) Path-containment on the storage service was only enforced in save_file; get_file/get_file_stream/delete_file/get_file_size all resolved user-supplied names directly, letting an authenticated user read arbitrary files via the MCP resources/read handler. resources/list and tools/list additionally returned every user's flows regardless of ownership. - storage/local.py: extract path-containment into shared _validated_path() and call it from every read/write/delete entry point (langflow + lfx). - mcp_utils.handle_read_resource: reject filenames containing ../, /, \ at the handler layer; require a current user and verify flow ownership (or self-owned user bucket) before dispatching to storage. - mcp_utils.handle_read_resource: accept optional project_id so project servers can't read resources outside the project's flows. - mcp_projects: pass project_id into handle_read_resource. - mcp_utils.handle_list_resources / handle_list_tools: scope flow queries to the authenticated user on the global server. Regression tests cover traversal rejection on every read path, cross-user flow access denial, project-scope enforcement, and unauthenticated list queries returning empty. * fix(mcp): drop user-bucket files from project-scoped resources/list Project-scoped handle_list_resources still appended every UserFile owned by the caller, so a project MCP client could enumerate user-level files unrelated to the project. User files have no project association, so skip them entirely when project_id is set. Adds a regression test proving that project-scoped resources/list returns only files from flows in the project. (cherry picked from commit f0fd436fe9829192ee550e6cb46961a01dd37032) --- .../base/langflow/api/v1/mcp_projects.py | 2 +- src/backend/base/langflow/api/v1/mcp_utils.py | 92 +++++++-- .../base/langflow/services/storage/local.py | 111 ++++++----- .../tests/unit/api/v1/test_mcp_utils.py | 174 +++++++++++++++++- .../storage/test_local_storage_service.py | 56 ++++++ src/lfx/src/lfx/services/storage/local.py | 58 +++++- 6 files changed, 419 insertions(+), 74 deletions(-) diff --git a/src/backend/base/langflow/api/v1/mcp_projects.py b/src/backend/base/langflow/api/v1/mcp_projects.py index 196158caf6..100a45e7b5 100644 --- a/src/backend/base/langflow/api/v1/mcp_projects.py +++ b/src/backend/base/langflow/api/v1/mcp_projects.py @@ -1249,7 +1249,7 @@ class ProjectMCPServer: @self.server.read_resource() async def handle_read_project_resource(uri: str) -> bytes: """Handle resource read requests for this specific project.""" - return await handle_read_resource(uri=uri) + return await handle_read_resource(uri=uri, project_id=self.project_id) @self.server.call_tool() @handle_mcp_errors diff --git a/src/backend/base/langflow/api/v1/mcp_utils.py b/src/backend/base/langflow/api/v1/mcp_utils.py index 8b6700d01e..5407e1fbb0 100644 --- a/src/backend/base/langflow/api/v1/mcp_utils.py +++ b/src/backend/base/langflow/api/v1/mcp_utils.py @@ -11,7 +11,7 @@ from functools import wraps from pathlib import Path from typing import Any, ParamSpec, TypeVar from urllib.parse import quote, unquote, urlparse -from uuid import uuid4 +from uuid import UUID, uuid4 from lfx.base.mcp.constants import MAX_MCP_TOOL_NAME_LENGTH from lfx.base.mcp.util import get_flow_snake_case, get_unique_name, sanitize_mcp_name @@ -100,9 +100,20 @@ async def handle_list_resources(project_id=None): msg = f"Error getting current user: {e!s}" await logger.aexception(msg) current_user = None + + # SECURITY: The current_user context is required to scope resources. + # Without it we cannot safely list files from any flow because the + # global server previously leaked every user's flow URIs (PVR0754098). + if current_user is None: + await logger.awarning("handle_list_resources called without a current user; returning empty list") + return resources + async with session_scope() as session: - # Build query based on whether project_id is provided - flows_query = select(Flow).where(Flow.folder_id == project_id) if project_id else select(Flow) + # SECURITY: Always scope to the calling user to prevent cross-user enumeration. + if project_id: + flows_query = select(Flow).where(Flow.folder_id == project_id, Flow.user_id == current_user.id) + else: + flows_query = select(Flow).where(Flow.user_id == current_user.id) flows = (await session.exec(flows_query)).all() @@ -132,8 +143,13 @@ async def handle_list_resources(project_id=None): # So the above query for flow files is not enough. # So we list all user files for the current user. # This is not good. We need to fix this for 1.8.0. + # + # SECURITY (PVR0754098): user-level files have no project association, + # so they must not be exposed through a project-scoped MCP server — + # doing so would let a project client enumerate files unrelated to + # the project. Only include them on the global (project_id is None) server. ################################################### - if current_user: + if project_id is None: user_files_stmt = select(UserFile).where(UserFile.user_id == current_user.id) user_files = (await session.exec(user_files_stmt)).all() for user_file in user_files: @@ -158,8 +174,15 @@ async def handle_list_resources(project_id=None): return resources -async def handle_read_resource(uri: str) -> bytes: - """Handle resource read requests.""" +async def handle_read_resource(uri: str, project_id: UUID | str | None = None) -> bytes: + """Handle resource read requests. + + Args: + uri: The resource URI; last two path segments are the namespace (flow_id or user_id) + and filename. + project_id: When invoked from a project-scoped server, restricts the lookup so a + caller cannot read resources that live outside the project. + """ try: # Parse the URI properly parsed_uri = urlparse(str(uri)) @@ -174,15 +197,50 @@ async def handle_read_resource(uri: str) -> bytes: msg = f"Invalid URI format: {uri}" raise ValueError(msg) - flow_id = path_parts[-2] + namespace_id = path_parts[-2] filename = unquote(path_parts[-1]) # URL decode the filename + # SECURITY (defense-in-depth): reject obvious traversal attempts before any + # service call. The storage service validates as well, but failing fast here + # keeps error logs from the storage layer off the hot path and closes the gap + # between the MCP decode step and the storage layer for future refactors. + if not filename or ".." in filename or "/" in filename or "\\" in filename: + await logger.awarning(f"Rejected MCP resource read with invalid filename: {filename!r}") + msg = "Invalid filename" + raise ValueError(msg) + + # SECURITY: authorise the caller before reading. The storage layer alone is + # not enough because the filesystem doesn't know about Langflow users, and + # previously any authenticated user could request any flow_id. + try: + current_user = current_user_ctx.get() + except LookupError as exc: + msg = "Authenticated user context is required to read MCP resources" + raise ValueError(msg) from exc + + async with session_scope() as session: + flow_query = select(Flow).where(Flow.id == namespace_id, Flow.user_id == current_user.id) + if project_id is not None: + flow_query = flow_query.where(Flow.folder_id == project_id) + flow = (await session.exec(flow_query)).first() + + if flow is None: + # The namespace segment may refer to the user's own bucket (user-level + # files uploaded via /api/v2/files) rather than a flow id. + if str(current_user.id) != str(namespace_id): + msg = "Resource not found or access denied" + raise ValueError(msg) + # User-level access is never in-scope for a project-scoped server. + if project_id is not None: + msg = "Resource not found or access denied" + raise ValueError(msg) + storage_service = get_storage_service() # Read the file content - content = await storage_service.get_file(flow_id=flow_id, file_name=filename) + content = await storage_service.get_file(flow_id=namespace_id, file_name=filename) if not content: - msg = f"File {filename} not found in flow {flow_id}" + msg = f"File {filename} not found in flow {namespace_id}" raise ValueError(msg) # Ensure content is base64 encoded @@ -335,6 +393,13 @@ async def handle_list_tools(project_id=None, *, mcp_enabled_only=False): """ tools = [] try: + # SECURITY: tools returned from the global server previously included every + # user's flows (PVR0754098). Always scope to the authenticated caller. + try: + current_user = current_user_ctx.get() + except LookupError: + current_user = None + async with session_scope() as session: # Build query based on parameters if project_id: @@ -342,9 +407,14 @@ async def handle_list_tools(project_id=None, *, mcp_enabled_only=False): flows_query = select(Flow).where(Flow.folder_id == project_id, Flow.is_component == False) # noqa: E712 if mcp_enabled_only: flows_query = flows_query.where(Flow.mcp_enabled == True) # noqa: E712 + elif current_user is not None: + # Global server: scope to the calling user only. + flows_query = select(Flow).where(Flow.user_id == current_user.id) else: - # Get all flows - flows_query = select(Flow) + await logger.awarning( + "handle_list_tools called without a current user and no project_id; returning empty list" + ) + return tools flows = (await session.exec(flows_query)).all() diff --git a/src/backend/base/langflow/services/storage/local.py b/src/backend/base/langflow/services/storage/local.py index d9e33e179c..ed8410e2bf 100644 --- a/src/backend/base/langflow/services/storage/local.py +++ b/src/backend/base/langflow/services/storage/local.py @@ -13,6 +13,8 @@ from langflow.services.storage.service import StorageService if TYPE_CHECKING: from collections.abc import AsyncIterator + import anyio + from langflow.services.session.service import SessionService from langflow.services.settings.service import SettingsService @@ -38,6 +40,48 @@ class LocalStorageService(StorageService): super().__init__(session_service, settings_service) # Base class already sets self.data_dir as anyio.Path from settings_service.settings.config_dir + async def _validated_path(self, flow_id: str, file_name: str) -> anyio.Path: + """Return a path inside the flow directory or raise if traversal is attempted. + + Centralizes the path-containment check used across every file operation so + a single bug cannot re-introduce the traversal issue on a subset of methods. + """ + if not isinstance(file_name, str) or "/" in file_name or "\\" in file_name or ".." in file_name: + await logger.aerror( + f"Invalid file_name contains path separators or traversal sequences: flow_id='{flow_id}'" + ) + msg = "Invalid file name: contains path separators" + raise ValueError(msg) + + folder_path = self.data_dir / flow_id + file_path = folder_path / file_name + + try: + resolved_folder = await folder_path.resolve() + resolved_file = await file_path.resolve() + if not resolved_file.is_relative_to(resolved_folder): + await logger.aerror( + f"Path traversal attempt detected for flow_id='{flow_id}'. " + "File path would escape flow directory boundary." + ) + msg = "Invalid file path: path traversal detected" + raise ValueError(msg) + except ValueError: + raise + except AttributeError: + resolved_folder_str = str(await folder_path.resolve()) + resolved_file_str = str(await file_path.resolve()) + if not resolved_file_str.startswith(resolved_folder_str): + await logger.aerror(f"Path traversal attempt detected for flow_id='{flow_id}' (fallback check)") + msg = "Invalid file path: path traversal detected" + raise ValueError(msg) from None + except Exception as e: + await logger.aerror(f"Error validating file path for flow_id='{flow_id}': {type(e).__name__}") + msg = "Invalid file path" + raise ValueError(msg) from e + + return file_path + def resolve_component_path(self, logical_path: str) -> str: """Convert logical path to absolute filesystem path for local storage. @@ -112,59 +156,11 @@ class LocalStorageService(StorageService): PermissionError: If there is no permission to write the file. ValueError: If path traversal is detected (security violation). """ - # SECURITY FIX: Defense-in-depth path containment check to prevent directory traversal - # Validate BEFORE creating any directories to prevent race conditions - # This ensures the resolved absolute path stays within the intended flow directory - # even if the API layer validation is bypassed or has bugs. + # SECURITY: Validate BEFORE creating any directories to prevent race conditions + # and path traversal. Shared with all other read/write/delete entry points so + # the storage boundary cannot be bypassed by any single call site. + file_path = await self._validated_path(flow_id, file_name) folder_path = self.data_dir / flow_id - file_path = folder_path / file_name - - try: - # Resolve paths to their absolute canonical forms - # Note: resolve() works even if the path doesn't exist yet - resolved_folder_path = await folder_path.resolve() - resolved_file_path = await file_path.resolve() - - # Check if the resolved file path is actually within the flow directory - # Using is_relative_to() ensures no path traversal can escape the boundary - # Example attack: file_name="../../etc/passwd" would resolve outside folder_path - if not resolved_file_path.is_relative_to(resolved_folder_path): - # SECURITY: Don't log the actual paths to avoid information disclosure - await logger.aerror( - f"Path traversal attempt detected for flow_id='{flow_id}'. " - "File path would escape flow directory boundary." - ) - msg = "Invalid file path: path traversal detected" - raise ValueError(msg) - - # Additional check: ensure file_name doesn't contain path separators - # This catches cases where Path() might normalize away dangerous sequences - if "/" in file_name or "\\" in file_name or ".." in file_name: - await logger.aerror( - f"Invalid file_name contains path separators or traversal sequences: flow_id='{flow_id}'" - ) - msg = "Invalid file name: contains path separators" - raise ValueError(msg) - - except ValueError: - # Re-raise ValueError (our security checks) - raise - except AttributeError: - # is_relative_to() not available (Python < 3.9), fall back to string comparison - # This should never happen in production but provides compatibility - resolved_folder_str = str(await folder_path.resolve()) - resolved_file_str = str(await file_path.resolve()) - if not resolved_file_str.startswith(resolved_folder_str): - await logger.aerror(f"Path traversal attempt detected for flow_id='{flow_id}' (fallback check)") - msg = "Invalid file path: path traversal detected" - raise ValueError(msg) from None - except Exception as e: - # Log unexpected errors during path resolution but don't expose details to user - await logger.aerror(f"Error validating file path for flow_id='{flow_id}': {type(e).__name__}") - msg = "Invalid file path" - raise ValueError(msg) from e - - # Only create directory after validation passes await folder_path.mkdir(parents=True, exist_ok=True) try: @@ -189,8 +185,9 @@ class LocalStorageService(StorageService): Raises: FileNotFoundError: If the file does not exist. + ValueError: If path traversal is detected (security violation). """ - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if not await file_path.exists(): await logger.awarning(f"File {file_name} not found in flow {flow_id}.") msg = f"File {file_name} not found in flow {flow_id}" @@ -204,7 +201,7 @@ class LocalStorageService(StorageService): async def get_file_stream(self, flow_id: str, file_name: str, chunk_size: int = 8192) -> AsyncIterator[bytes]: """Retrieve a file from storage as a stream.""" - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if not await file_path.exists(): await logger.awarning(f"File {file_name} not found in flow {flow_id}.") msg = f"File {file_name} not found in flow {flow_id}" @@ -252,8 +249,9 @@ class LocalStorageService(StorageService): Raises: FileNotFoundError: If the file does not exist. + ValueError: If path traversal is detected (security violation). """ - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if await file_path.exists(): await file_path.unlink() await logger.ainfo(f"File {file_name} deleted successfully from flow {flow_id}.") @@ -272,8 +270,9 @@ class LocalStorageService(StorageService): Raises: FileNotFoundError: If the file does not exist. + ValueError: If path traversal is detected (security violation). """ - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if not await file_path.exists(): await logger.awarning(f"File {file_name} not found in flow {flow_id}.") msg = f"File {file_name} not found in flow {flow_id}" diff --git a/src/backend/tests/unit/api/v1/test_mcp_utils.py b/src/backend/tests/unit/api/v1/test_mcp_utils.py index 0dba6f8728..e548eaf377 100644 --- a/src/backend/tests/unit/api/v1/test_mcp_utils.py +++ b/src/backend/tests/unit/api/v1/test_mcp_utils.py @@ -12,6 +12,9 @@ class FakeResult: def all(self): return list(self._rows) + def first(self): + return self._rows[0] if self._rows else None + class FakeSession: def __init__(self, flows, user_files): @@ -40,12 +43,20 @@ class FakeSessionContext: class FakeStorageService: - def __init__(self, files_by_flow): + def __init__(self, files_by_flow, file_bytes: dict[str, bytes] | None = None): self._files_by_flow = files_by_flow + self._file_bytes = file_bytes or {} async def list_files(self, flow_id: str): return self._files_by_flow.get(flow_id, []) + async def get_file(self, flow_id: str, file_name: str) -> bytes: + key = f"{flow_id}/{file_name}" + if key not in self._file_bytes: + msg = f"File {file_name} not found in flow {flow_id}" + raise FileNotFoundError(msg) + return self._file_bytes[key] + @pytest.mark.asyncio async def test_handle_list_resources_includes_flow_and_user_files(monkeypatch): @@ -119,6 +130,165 @@ async def test_handle_list_tools_skips_blocked_custom_flows(monkeypatch): monkeypatch.setattr(component_cache, "type_to_current_hash", {"ChatInput": "known-hash"}) monkeypatch.setattr(component_cache, "all_types_dict", None) - tools = await mcp_utils.handle_list_tools() + token = mcp_utils.current_user_ctx.set(SimpleNamespace(id="user-1")) + try: + tools = await mcp_utils.handle_list_tools() + finally: + mcp_utils.current_user_ctx.reset(token) assert tools == [] + + +# ============================================================================ +# PVR0754098 regression tests — MCP path traversal and cross-user disclosure. +# ============================================================================ + + +@pytest.mark.asyncio +async def test_handle_list_resources_requires_current_user(monkeypatch): + """Without a user context, the global server must not enumerate any flows.""" + flows = [SimpleNamespace(id="flow-attacker-saw", name="Leaked Flow")] + fake_session = FakeSession(flows=flows, user_files=[]) + storage_service = FakeStorageService({"flow-attacker-saw": ["leaked.txt"]}) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + monkeypatch.setattr(mcp_utils, "get_storage_service", lambda: storage_service) + monkeypatch.setattr( + mcp_utils, + "get_settings_service", + lambda: SimpleNamespace(settings=SimpleNamespace(host="localhost", port=4000)), + ) + + # Intentionally no user context set. + resources = await mcp_utils.handle_list_resources() + assert resources == [] + + +@pytest.mark.asyncio +async def test_handle_read_resource_rejects_path_traversal(monkeypatch): + """Encoded ../ sequences must be rejected before reaching storage.""" + own_flow = SimpleNamespace(id="flow-own", user_id="user-bob", folder_id=None) + fake_session = FakeSession(flows=[own_flow], user_files=[]) + storage_service = FakeStorageService({}, {"flow-own/legit.txt": b"ok"}) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + monkeypatch.setattr(mcp_utils, "get_storage_service", lambda: storage_service) + + token = mcp_utils.current_user_ctx.set(SimpleNamespace(id="user-bob")) + try: + uri = "http://host/api/v1/files/download/flow-own/..%2F..%2F..%2Fetc%2Fpasswd" + with pytest.raises(ValueError, match="Invalid filename"): + await mcp_utils.handle_read_resource(uri) + finally: + mcp_utils.current_user_ctx.reset(token) + + +@pytest.mark.asyncio +async def test_handle_read_resource_denies_other_users_flow(monkeypatch): + """Bob must not be able to read a flow owned by Alice, even with a valid filename.""" + # Session returns nothing for the ownership query — i.e. the flow does not belong to bob. + fake_session = FakeSession(flows=[], user_files=[]) + storage_service = FakeStorageService({}, {"flow-alice/secret.txt": b"alice-secret"}) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + monkeypatch.setattr(mcp_utils, "get_storage_service", lambda: storage_service) + + token = mcp_utils.current_user_ctx.set(SimpleNamespace(id="user-bob")) + try: + uri = "http://host/api/v1/files/download/flow-alice/secret.txt" + with pytest.raises(ValueError, match="access denied"): + await mcp_utils.handle_read_resource(uri) + finally: + mcp_utils.current_user_ctx.reset(token) + + +@pytest.mark.asyncio +async def test_handle_read_resource_allows_user_level_bucket(monkeypatch): + """A user reading from their own user-level bucket (not a flow id) is allowed.""" + # No flow match returned; namespace equals current_user.id so access is allowed. + fake_session = FakeSession(flows=[], user_files=[]) + storage_service = FakeStorageService({}, {"user-bob/my-file.txt": b"mine"}) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + monkeypatch.setattr(mcp_utils, "get_storage_service", lambda: storage_service) + + token = mcp_utils.current_user_ctx.set(SimpleNamespace(id="user-bob")) + try: + uri = "http://host/api/v1/files/download/user-bob/my-file.txt" + result = await mcp_utils.handle_read_resource(uri) + finally: + mcp_utils.current_user_ctx.reset(token) + + # handle_read_resource returns base64-encoded bytes. + import base64 + + assert base64.b64decode(result) == b"mine" + + +@pytest.mark.asyncio +async def test_handle_read_resource_denies_user_bucket_under_project_scope(monkeypatch): + """User-level bucket access must not leak through a project-scoped server.""" + fake_session = FakeSession(flows=[], user_files=[]) + storage_service = FakeStorageService({}, {"user-bob/my-file.txt": b"mine"}) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + monkeypatch.setattr(mcp_utils, "get_storage_service", lambda: storage_service) + + token = mcp_utils.current_user_ctx.set(SimpleNamespace(id="user-bob")) + try: + uri = "http://host/api/v1/files/download/user-bob/my-file.txt" + with pytest.raises(ValueError, match="access denied"): + await mcp_utils.handle_read_resource(uri, project_id="project-xyz") + finally: + mcp_utils.current_user_ctx.reset(token) + + +@pytest.mark.asyncio +async def test_handle_list_resources_project_scoped_excludes_user_bucket(monkeypatch): + """A project-scoped resources/list must not leak user-bucket files unrelated to the project.""" + user_id = "user-bob" + project_flow = SimpleNamespace(id="flow-in-project", name="Project Flow") + + # If the implementation incorrectly includes user-bucket files, this one would show up. + user_files = [ + SimpleNamespace( + name="unrelated.pdf", + path=f"{user_id}/unrelated.pdf", + provider="File Manager", + ) + ] + + fake_session = FakeSession(flows=[project_flow], user_files=user_files) + storage_service = FakeStorageService({"flow-in-project": ["project-doc.txt"]}) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + monkeypatch.setattr(mcp_utils, "get_storage_service", lambda: storage_service) + monkeypatch.setattr( + mcp_utils, + "get_settings_service", + lambda: SimpleNamespace(settings=SimpleNamespace(host="localhost", port=4000)), + ) + + token = mcp_utils.current_user_ctx.set(SimpleNamespace(id=user_id)) + try: + resources = await mcp_utils.handle_list_resources(project_id="project-xyz") + finally: + mcp_utils.current_user_ctx.reset(token) + + uris = {str(resource.uri) for resource in resources} + assert "http://localhost:4000/api/v1/files/download/flow-in-project/project-doc.txt" in uris + # User-bucket file must not leak through a project-scoped server. + assert not any("unrelated.pdf" in uri for uri in uris) + + +@pytest.mark.asyncio +async def test_handle_list_tools_requires_current_user_on_global_server(monkeypatch): + """Global list_tools must refuse to enumerate flows without a user context.""" + flows = [SimpleNamespace(id="flow-leak", user_id="someone-else", name="Leaked", description=None, data={})] + fake_session = FakeSession(flows=flows, user_files=[]) + + monkeypatch.setattr(mcp_utils, "session_scope", lambda: FakeSessionContext(fake_session)) + + # No user context set — must return empty. + tools = await mcp_utils.handle_list_tools() + assert tools == [] diff --git a/src/backend/tests/unit/services/storage/test_local_storage_service.py b/src/backend/tests/unit/services/storage/test_local_storage_service.py index bd4a4fafa4..099d9a9d61 100644 --- a/src/backend/tests/unit/services/storage/test_local_storage_service.py +++ b/src/backend/tests/unit/services/storage/test_local_storage_service.py @@ -320,6 +320,62 @@ class TestLocalStorageServiceTeardown: # Local storage teardown is a no-op, so just verify it doesn't raise +@pytest.mark.asyncio +class TestLocalStorageServicePathTraversal: + """Regression tests for PVR0754098 — path traversal on read/delete/size endpoints.""" + + @pytest.mark.parametrize( + "malicious_filename", + [ + "../../../etc/passwd", + "../other_flow/file.txt", + "..\\..\\etc\\passwd", + "..", + "sub/dir.txt", + ], + ) + async def test_get_file_rejects_traversal(self, local_storage_service, malicious_filename): + with pytest.raises(ValueError, match="Invalid file"): + await local_storage_service.get_file("flow_a", malicious_filename) + + @pytest.mark.parametrize( + "malicious_filename", + [ + "../../../etc/passwd", + "../other_flow/file.txt", + "..", + ], + ) + async def test_delete_file_rejects_traversal(self, local_storage_service, malicious_filename): + with pytest.raises(ValueError, match="Invalid file"): + await local_storage_service.delete_file("flow_a", malicious_filename) + + @pytest.mark.parametrize( + "malicious_filename", + [ + "../../../etc/passwd", + "../other_flow/file.txt", + "..", + ], + ) + async def test_get_file_size_rejects_traversal(self, local_storage_service, malicious_filename): + with pytest.raises(ValueError, match="Invalid file"): + await local_storage_service.get_file_size("flow_a", malicious_filename) + + async def test_get_file_stream_rejects_traversal(self, local_storage_service): + with pytest.raises(ValueError, match="Invalid file"): + # AsyncIterator functions don't raise until first iteration. + async for _ in local_storage_service.get_file_stream("flow_a", "../escape.txt"): + pass + + async def test_cross_flow_read_is_blocked(self, local_storage_service): + """Writing a file to flow_b then attempting to read it from flow_a must fail.""" + await local_storage_service.save_file("flow_b", "secret.txt", b"hidden") + + with pytest.raises(ValueError, match="Invalid file"): + await local_storage_service.get_file("flow_a", "../flow_b/secret.txt") + + @pytest.mark.asyncio class TestLocalStorageServiceEdgeCases: """Test edge cases and error conditions.""" diff --git a/src/lfx/src/lfx/services/storage/local.py b/src/lfx/src/lfx/services/storage/local.py index a59cb185aa..810d649cbf 100644 --- a/src/lfx/src/lfx/services/storage/local.py +++ b/src/lfx/src/lfx/services/storage/local.py @@ -11,6 +11,7 @@ from lfx.services.base import Service from lfx.services.storage.service import StorageService if TYPE_CHECKING: + import anyio from langflow.services.session.service import SessionService from lfx.services.settings.service import SettingsService @@ -37,6 +38,48 @@ class LocalStorageService(StorageService, Service): super().__init__(session_service, settings_service) # Base class already sets self.data_dir as anyio.Path from settings_service.settings.config_dir + async def _validated_path(self, flow_id: str, file_name: str) -> anyio.Path: + """Return a path inside the flow directory or raise if traversal is attempted. + + Centralizes the path-containment check used across every file operation so + a single bug cannot re-introduce the traversal issue on a subset of methods. + """ + if not isinstance(file_name, str) or "/" in file_name or "\\" in file_name or ".." in file_name: + await logger.aerror( + f"Invalid file_name contains path separators or traversal sequences: flow_id='{flow_id}'" + ) + msg = "Invalid file name: contains path separators" + raise ValueError(msg) + + folder_path = self.data_dir / flow_id + file_path = folder_path / file_name + + try: + resolved_folder = await folder_path.resolve() + resolved_file = await file_path.resolve() + if not resolved_file.is_relative_to(resolved_folder): + await logger.aerror( + f"Path traversal attempt detected for flow_id='{flow_id}'. " + "File path would escape flow directory boundary." + ) + msg = "Invalid file path: path traversal detected" + raise ValueError(msg) + except ValueError: + raise + except AttributeError: + resolved_folder_str = str(await folder_path.resolve()) + resolved_file_str = str(await file_path.resolve()) + if not resolved_file_str.startswith(resolved_folder_str): + await logger.aerror(f"Path traversal attempt detected for flow_id='{flow_id}' (fallback check)") + msg = "Invalid file path: path traversal detected" + raise ValueError(msg) from None + except Exception as e: + await logger.aerror(f"Error validating file path for flow_id='{flow_id}': {type(e).__name__}") + msg = "Invalid file path" + raise ValueError(msg) from e + + return file_path + def resolve_component_path(self, logical_path: str) -> str: """Convert logical path to absolute filesystem path for local storage. @@ -112,10 +155,14 @@ class LocalStorageService(StorageService, Service): FileNotFoundError: If the specified flow does not exist. IsADirectoryError: If the file name is a directory. PermissionError: If there is no permission to write the file. + ValueError: If path traversal is detected (security violation). """ + # SECURITY: Validate BEFORE creating any directories to prevent race conditions + # and path traversal. Shared with all read/write/delete entry points so the + # storage boundary cannot be bypassed by any single call site. + file_path = await self._validated_path(flow_id, file_name) folder_path = self.data_dir / flow_id await folder_path.mkdir(parents=True, exist_ok=True) - file_path = folder_path / file_name try: mode = "ab" if append else "wb" @@ -139,8 +186,9 @@ class LocalStorageService(StorageService, Service): Raises: FileNotFoundError: If the file does not exist. + ValueError: If path traversal is detected (security violation). """ - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if not await file_path.exists(): await logger.awarning(f"File {file_name} not found in flow {flow_id}.") msg = f"File {file_name} not found in flow {flow_id}" @@ -187,8 +235,9 @@ class LocalStorageService(StorageService, Service): Raises: FileNotFoundError: If the file does not exist. + ValueError: If path traversal is detected (security violation). """ - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if await file_path.exists(): await file_path.unlink() await logger.ainfo(f"File {file_name} deleted successfully from flow {flow_id}.") @@ -207,8 +256,9 @@ class LocalStorageService(StorageService, Service): Raises: FileNotFoundError: If the file does not exist. + ValueError: If path traversal is detected (security violation). """ - file_path = self.data_dir / flow_id / file_name + file_path = await self._validated_path(flow_id, file_name) if not await file_path.exists(): await logger.awarning(f"File {file_name} not found in flow {flow_id}.") msg = f"File {file_name} not found in flow {flow_id}"