diff --git a/backend/packages/harness/deerflow/config/paths.py b/backend/packages/harness/deerflow/config/paths.py index c0683904..0c90756c 100644 --- a/backend/packages/harness/deerflow/config/paths.py +++ b/backend/packages/harness/deerflow/config/paths.py @@ -10,6 +10,7 @@ VIRTUAL_PATH_PREFIX = "/mnt/user-data" _SAFE_THREAD_ID_RE = re.compile(r"^[A-Za-z0-9_\-]+$") _SAFE_USER_ID_RE = re.compile(r"^[A-Za-z0-9_\-]+$") +_SAFE_WORKSPACE_ID_RE = re.compile(r"^[A-Za-z0-9_\-]+$") def _default_local_base_dir() -> Path: @@ -31,6 +32,13 @@ def _validate_user_id(user_id: str) -> str: return user_id +def _validate_workspace_id(workspace_id: str) -> str: + """Validate a workspace ID before using it in filesystem paths.""" + if not _SAFE_WORKSPACE_ID_RE.match(workspace_id): + raise ValueError(f"Invalid workspace_id {workspace_id!r}: only alphanumeric characters, hyphens, and underscores are allowed.") + return workspace_id + + def _join_host_path(base: str, *parts: str) -> str: """Join host filesystem path segments while preserving native style. @@ -148,116 +156,129 @@ class Paths: """Legacy per-agent memory file: `{base_dir}/agents/{name}/memory.json`.""" return self.agent_dir(name) / "memory.json" - def user_dir(self, user_id: str) -> Path: - """Directory for a specific user: `{base_dir}/users/{user_id}/`.""" + def workspace_dir(self, workspace_id: str) -> Path: + """Directory for a specific workspace: `{base_dir}/workspaces/{workspace_id}/`. + + PR6 introduces this as the top-level isolation dimension. Per-user + state and per-thread state both live underneath their workspace so + a user with access to two workspaces never sees state bleed between + them on the filesystem. + """ + return self.base_dir / "workspaces" / _validate_workspace_id(workspace_id) + + def user_dir(self, user_id: str, *, workspace_id: str | None = None) -> Path: + """Directory for a specific user. + + When ``workspace_id`` is provided (PR6+): + ``{base_dir}/workspaces/{wid}/users/{user_id}/`` + + Otherwise (legacy layout): + ``{base_dir}/users/{user_id}/`` + """ + if workspace_id is not None: + return self.workspace_dir(workspace_id) / "users" / _validate_user_id(user_id) return self.base_dir / "users" / _validate_user_id(user_id) - def user_memory_file(self, user_id: str) -> Path: - """Per-user memory file: `{base_dir}/users/{user_id}/memory.json`.""" - return self.user_dir(user_id) / "memory.json" + def user_memory_file(self, user_id: str, *, workspace_id: str | None = None) -> Path: + """Per-user memory file under the active workspace (legacy without).""" + return self.user_dir(user_id, workspace_id=workspace_id) / "memory.json" - def user_agents_dir(self, user_id: str) -> Path: - """Per-user root for that user's custom agents: `{base_dir}/users/{user_id}/agents/`.""" - return self.user_dir(user_id) / "agents" + def user_agents_dir(self, user_id: str, *, workspace_id: str | None = None) -> Path: + """Per-user root for custom agents under the active workspace.""" + return self.user_dir(user_id, workspace_id=workspace_id) / "agents" - def user_agent_dir(self, user_id: str, agent_name: str) -> Path: - """Per-user per-agent directory: `{base_dir}/users/{user_id}/agents/{name}/`.""" - return self.user_agents_dir(user_id) / agent_name.lower() + def user_agent_dir(self, user_id: str, agent_name: str, *, workspace_id: str | None = None) -> Path: + """Per-user per-agent directory under the active workspace.""" + return self.user_agents_dir(user_id, workspace_id=workspace_id) / agent_name.lower() - def user_agent_memory_file(self, user_id: str, agent_name: str) -> Path: - """Per-user per-agent memory: `{base_dir}/users/{user_id}/agents/{name}/memory.json`.""" - return self.user_agent_dir(user_id, agent_name) / "memory.json" + def user_agent_memory_file(self, user_id: str, agent_name: str, *, workspace_id: str | None = None) -> Path: + """Per-user per-agent memory file under the active workspace.""" + return self.user_agent_dir(user_id, agent_name, workspace_id=workspace_id) / "memory.json" - def thread_dir(self, thread_id: str, *, user_id: str | None = None) -> Path: + def thread_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> Path: """ Host path for a thread's data. - When *user_id* is provided: - `{base_dir}/users/{user_id}/threads/{thread_id}/` - Otherwise (legacy layout): - `{base_dir}/threads/{thread_id}/` + Precedence — workspace beats user, both beat legacy: - This directory contains a `user-data/` subdirectory that is mounted - as `/mnt/user-data/` inside the sandbox. + * ``workspace_id`` given (PR6+): + ``{base_dir}/workspaces/{wid}/threads/{thread_id}/`` + * ``user_id`` only (legacy after user-isolation migration): + ``{base_dir}/users/{user_id}/threads/{thread_id}/`` + * neither (very legacy, pre-isolation): + ``{base_dir}/threads/{thread_id}/`` + + The contained ``user-data/`` subdirectory is mounted as + ``/mnt/user-data/`` inside the sandbox regardless of which form is + chosen — only the host-side parent differs. Raises: - ValueError: If `thread_id` or `user_id` contains unsafe characters (path - separators or `..`) that could cause directory traversal. + ValueError: If any of the supplied ids contains unsafe characters. """ + if workspace_id is not None: + return self.workspace_dir(workspace_id) / "threads" / _validate_thread_id(thread_id) if user_id is not None: - return self.user_dir(user_id) / "threads" / _validate_thread_id(thread_id) + return self.base_dir / "users" / _validate_user_id(user_id) / "threads" / _validate_thread_id(thread_id) return self.base_dir / "threads" / _validate_thread_id(thread_id) - def sandbox_work_dir(self, thread_id: str, *, user_id: str | None = None) -> Path: + def sandbox_work_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> Path: """ Host path for the agent's workspace directory. - Host: `{base_dir}/threads/{thread_id}/user-data/workspace/` - Sandbox: `/mnt/user-data/workspace/` + Host: ``{thread_dir}/user-data/workspace/`` + Sandbox: ``/mnt/user-data/workspace/`` """ - return self.thread_dir(thread_id, user_id=user_id) / "user-data" / "workspace" + return self.thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id) / "user-data" / "workspace" - def sandbox_uploads_dir(self, thread_id: str, *, user_id: str | None = None) -> Path: - """ - Host path for user-uploaded files. - Host: `{base_dir}/threads/{thread_id}/user-data/uploads/` - Sandbox: `/mnt/user-data/uploads/` - """ - return self.thread_dir(thread_id, user_id=user_id) / "user-data" / "uploads" + def sandbox_uploads_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> Path: + """Host path for user-uploaded files; sandbox: ``/mnt/user-data/uploads/``.""" + return self.thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id) / "user-data" / "uploads" - def sandbox_outputs_dir(self, thread_id: str, *, user_id: str | None = None) -> Path: - """ - Host path for agent-generated artifacts. - Host: `{base_dir}/threads/{thread_id}/user-data/outputs/` - Sandbox: `/mnt/user-data/outputs/` - """ - return self.thread_dir(thread_id, user_id=user_id) / "user-data" / "outputs" + def sandbox_outputs_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> Path: + """Host path for agent-generated artifacts; sandbox: ``/mnt/user-data/outputs/``.""" + return self.thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id) / "user-data" / "outputs" - def acp_workspace_dir(self, thread_id: str, *, user_id: str | None = None) -> Path: + def acp_workspace_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> Path: """ - Host path for the ACP workspace of a specific thread. - Host: `{base_dir}/threads/{thread_id}/acp-workspace/` - Sandbox: `/mnt/acp-workspace/` + Host path for the ACP workspace of a specific thread; sandbox: ``/mnt/acp-workspace/``. Each thread gets its own isolated ACP workspace so that concurrent sessions cannot read each other's ACP agent outputs. """ - return self.thread_dir(thread_id, user_id=user_id) / "acp-workspace" + return self.thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id) / "acp-workspace" - def sandbox_user_data_dir(self, thread_id: str, *, user_id: str | None = None) -> Path: - """ - Host path for the user-data root. - Host: `{base_dir}/threads/{thread_id}/user-data/` - Sandbox: `/mnt/user-data/` - """ - return self.thread_dir(thread_id, user_id=user_id) / "user-data" + def sandbox_user_data_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> Path: + """Host path for the user-data root; sandbox: ``/mnt/user-data/``.""" + return self.thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id) / "user-data" - def host_thread_dir(self, thread_id: str, *, user_id: str | None = None) -> str: + def host_thread_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> str: """Host path for a thread directory, preserving Windows path syntax.""" + if workspace_id is not None: + return _join_host_path(self._host_base_dir_str(), "workspaces", _validate_workspace_id(workspace_id), "threads", _validate_thread_id(thread_id)) if user_id is not None: return _join_host_path(self._host_base_dir_str(), "users", _validate_user_id(user_id), "threads", _validate_thread_id(thread_id)) return _join_host_path(self._host_base_dir_str(), "threads", _validate_thread_id(thread_id)) - def host_sandbox_user_data_dir(self, thread_id: str, *, user_id: str | None = None) -> str: + def host_sandbox_user_data_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> str: """Host path for a thread's user-data root.""" - return _join_host_path(self.host_thread_dir(thread_id, user_id=user_id), "user-data") + return _join_host_path(self.host_thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id), "user-data") - def host_sandbox_work_dir(self, thread_id: str, *, user_id: str | None = None) -> str: + def host_sandbox_work_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> str: """Host path for the workspace mount source.""" - return _join_host_path(self.host_sandbox_user_data_dir(thread_id, user_id=user_id), "workspace") + return _join_host_path(self.host_sandbox_user_data_dir(thread_id, workspace_id=workspace_id, user_id=user_id), "workspace") - def host_sandbox_uploads_dir(self, thread_id: str, *, user_id: str | None = None) -> str: + def host_sandbox_uploads_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> str: """Host path for the uploads mount source.""" - return _join_host_path(self.host_sandbox_user_data_dir(thread_id, user_id=user_id), "uploads") + return _join_host_path(self.host_sandbox_user_data_dir(thread_id, workspace_id=workspace_id, user_id=user_id), "uploads") - def host_sandbox_outputs_dir(self, thread_id: str, *, user_id: str | None = None) -> str: + def host_sandbox_outputs_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> str: """Host path for the outputs mount source.""" - return _join_host_path(self.host_sandbox_user_data_dir(thread_id, user_id=user_id), "outputs") + return _join_host_path(self.host_sandbox_user_data_dir(thread_id, workspace_id=workspace_id, user_id=user_id), "outputs") - def host_acp_workspace_dir(self, thread_id: str, *, user_id: str | None = None) -> str: + def host_acp_workspace_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> str: """Host path for the ACP workspace mount source.""" - return _join_host_path(self.host_thread_dir(thread_id, user_id=user_id), "acp-workspace") + return _join_host_path(self.host_thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id), "acp-workspace") - def ensure_thread_dirs(self, thread_id: str, *, user_id: str | None = None) -> None: + def ensure_thread_dirs(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> None: """Create all standard sandbox directories for a thread. Directories are created with mode 0o777 so that sandbox containers @@ -271,24 +292,28 @@ class Paths: ACP agent invocation. """ for d in [ - self.sandbox_work_dir(thread_id, user_id=user_id), - self.sandbox_uploads_dir(thread_id, user_id=user_id), - self.sandbox_outputs_dir(thread_id, user_id=user_id), - self.acp_workspace_dir(thread_id, user_id=user_id), + self.sandbox_work_dir(thread_id, workspace_id=workspace_id, user_id=user_id), + self.sandbox_uploads_dir(thread_id, workspace_id=workspace_id, user_id=user_id), + self.sandbox_outputs_dir(thread_id, workspace_id=workspace_id, user_id=user_id), + self.acp_workspace_dir(thread_id, workspace_id=workspace_id, user_id=user_id), ]: d.mkdir(parents=True, exist_ok=True) d.chmod(0o777) - def delete_thread_dir(self, thread_id: str, *, user_id: str | None = None) -> None: - """Delete all persisted data for a thread. - - The operation is idempotent: missing thread directories are ignored. - """ - thread_dir = self.thread_dir(thread_id, user_id=user_id) + def delete_thread_dir(self, thread_id: str, *, workspace_id: str | None = None, user_id: str | None = None) -> None: + """Delete all persisted data for a thread. Idempotent.""" + thread_dir = self.thread_dir(thread_id, workspace_id=workspace_id, user_id=user_id) if thread_dir.exists(): shutil.rmtree(thread_dir) - def resolve_virtual_path(self, thread_id: str, virtual_path: str, *, user_id: str | None = None) -> Path: + def resolve_virtual_path( + self, + thread_id: str, + virtual_path: str, + *, + workspace_id: str | None = None, + user_id: str | None = None, + ) -> Path: """Resolve a sandbox virtual path to the actual host filesystem path. Args: @@ -296,7 +321,8 @@ class Paths: virtual_path: Virtual path as seen inside the sandbox, e.g. ``/mnt/user-data/outputs/report.pdf``. Leading slashes are stripped before matching. - user_id: Optional user ID for user-scoped path resolution. + workspace_id: Optional workspace ID for workspace-scoped resolution. + user_id: Optional user ID for legacy user-scoped resolution. Returns: The resolved absolute host filesystem path. @@ -314,7 +340,7 @@ class Paths: raise ValueError(f"Path must start with /{prefix}") relative = stripped[len(prefix) :].lstrip("/") - base = self.sandbox_user_data_dir(thread_id, user_id=user_id).resolve() + base = self.sandbox_user_data_dir(thread_id, workspace_id=workspace_id, user_id=user_id).resolve() actual = (base / relative).resolve() try: diff --git a/backend/tests/test_paths_workspace.py b/backend/tests/test_paths_workspace.py new file mode 100644 index 00000000..2ce02a9c --- /dev/null +++ b/backend/tests/test_paths_workspace.py @@ -0,0 +1,115 @@ +"""PR6 T6.9 / T6.10 — workspace-scoped path resolution. + +The Paths class learns a new top-level dimension for multi-tenant +filesystems: ``{base_dir}/workspaces/{wid}/...``. Precedence: + +- ``workspace_id`` given → new shape ``workspaces/{wid}/threads/{tid}/...`` +- only ``user_id`` given → legacy shape ``users/{uid}/threads/{tid}/...`` +- neither → very-legacy shape ``threads/{tid}/...`` + +Per-user filesystem state (memory.json, custom agents) lives under the +workspace too: ``workspaces/{wid}/users/{uid}/memory.json`` etc. — so a +user's memory cannot be reused across workspaces by mistake. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from deerflow.config.paths import Paths + + +@pytest.fixture +def paths(tmp_path: Path) -> Paths: + return Paths(tmp_path) + + +class TestValidateWorkspaceId: + def test_valid_workspace_id(self, paths: Paths): + d = paths.workspace_dir("ws-abc-123") + assert d == paths.base_dir / "workspaces" / "ws-abc-123" + + def test_rejects_path_traversal(self, paths: Paths): + with pytest.raises(ValueError, match="Invalid workspace_id"): + paths.workspace_dir("../escape") + + def test_rejects_slash(self, paths: Paths): + with pytest.raises(ValueError, match="Invalid workspace_id"): + paths.workspace_dir("ws/bar") + + def test_rejects_empty(self, paths: Paths): + with pytest.raises(ValueError, match="Invalid workspace_id"): + paths.workspace_dir("") + + +class TestWorkspaceScopedThreadDir: + def test_workspace_takes_precedence_over_user(self, paths: Paths): + """When both are given, workspace wins — user_id is recorded in the row, not the filesystem.""" + expected = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" + assert paths.thread_dir("t1", workspace_id="ws-alpha", user_id="alice") == expected + + def test_workspace_only(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" + assert paths.thread_dir("t1", workspace_id="ws-alpha") == expected + + def test_user_only_still_legacy(self, paths: Paths): + """Legacy callers keep the user_id-only shape until migration runs.""" + expected = paths.base_dir / "users" / "alice" / "threads" / "t1" + assert paths.thread_dir("t1", user_id="alice") == expected + + def test_no_ids_very_legacy(self, paths: Paths): + expected = paths.base_dir / "threads" / "t1" + assert paths.thread_dir("t1") == expected + + +class TestEnsureThreadDirsWorkspace: + def test_creates_workspace_layout(self, paths: Paths): + paths.ensure_thread_dirs("t1", workspace_id="ws-alpha") + root = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" + for sub in ("user-data/workspace", "user-data/uploads", "user-data/outputs", "acp-workspace"): + assert (root / sub).is_dir(), f"missing {sub}" + + +class TestSandboxDirsWorkspace: + def test_sandbox_work_dir(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" / "user-data" / "workspace" + assert paths.sandbox_work_dir("t1", workspace_id="ws-alpha") == expected + + def test_sandbox_uploads_dir(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" / "user-data" / "uploads" + assert paths.sandbox_uploads_dir("t1", workspace_id="ws-alpha") == expected + + def test_sandbox_outputs_dir(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" / "user-data" / "outputs" + assert paths.sandbox_outputs_dir("t1", workspace_id="ws-alpha") == expected + + +class TestUserMemoryUnderWorkspace: + def test_user_memory_file_under_workspace(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "users" / "alice" / "memory.json" + assert paths.user_memory_file("alice", workspace_id="ws-alpha") == expected + + def test_user_memory_file_legacy_without_workspace(self, paths: Paths): + expected = paths.base_dir / "users" / "alice" / "memory.json" + assert paths.user_memory_file("alice") == expected + + def test_user_agents_dir_under_workspace(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "users" / "alice" / "agents" + assert paths.user_agents_dir("alice", workspace_id="ws-alpha") == expected + + def test_user_agent_memory_file_under_workspace(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "users" / "alice" / "agents" / "code-reviewer" / "memory.json" + assert paths.user_agent_memory_file("alice", "code-reviewer", workspace_id="ws-alpha") == expected + + +class TestVirtualPathResolutionWorkspace: + def test_resolve_virtual_path_workspace_scope(self, paths: Paths): + expected = paths.base_dir / "workspaces" / "ws-alpha" / "threads" / "t1" / "user-data" / "outputs" / "x.json" + actual = paths.resolve_virtual_path("t1", "/mnt/user-data/outputs/x.json", workspace_id="ws-alpha") + assert actual == expected.resolve() + + def test_resolve_virtual_path_rejects_traversal_under_workspace(self, paths: Paths): + with pytest.raises(ValueError, match="path traversal"): + paths.resolve_virtual_path("t1", "/mnt/user-data/../../etc/passwd", workspace_id="ws-alpha")