diff --git a/backend/AGENTS.md b/backend/AGENTS.md index 20f63e0a4..95b9a1cda 100644 --- a/backend/AGENTS.md +++ b/backend/AGENTS.md @@ -347,8 +347,8 @@ Proxied through nginx: `/api/langgraph/*` → Gateway LangGraph-compatible runti **Provider Pattern**: `SandboxProvider` with `acquire`, `acquire_async`, `get`, `release` lifecycle. Async agent/tool paths call async sandbox lifecycle hooks so Docker sandbox creation, discovery, cross-process locking, readiness polling, and release stay off the event loop. **Environment policy** (`sandbox/env_policy.py`): `execute_command` no longer inherits the full `os.environ`. `build_sandbox_env()` scrubs secret-looking names (`*KEY*`/`*SECRET*`/`*TOKEN*`/`*PASSWORD*`/`*CREDENTIAL*`) from the inherited environment before layering injected request secrets on top, so platform credentials (e.g. `OPENAI_API_KEY`) never leak into skill subprocesses. Benign vars (`PATH`, `HOME`, `LANG`, `VIRTUAL_ENV`, ...) are preserved. **Implementations**: -- `LocalSandboxProvider` - Local filesystem execution. `acquire(thread_id)` returns a per-thread `LocalSandbox` (id `local:{thread_id}`) whose `path_mappings` resolve `/mnt/user-data/{workspace,uploads,outputs}` and `/mnt/acp-workspace` to that thread's host directories, so the public `Sandbox` API honours the `/mnt/user-data` contract uniformly with AIO. `acquire()` / `acquire(None)` keeps the legacy generic singleton (id `local`) for callers without a thread context. Per-thread sandboxes are held in an LRU cache (default 256 entries) guarded by a `threading.Lock`. -- `AioSandboxProvider` (`packages/harness/deerflow/community/`) - Docker-based isolation. Active-cache and warm-pool entries are checked with the backend during acquire/reuse; definitively dead containers are dropped from all in-process maps so the thread can discover or create a fresh sandbox instead of reusing a stale client. Backend health-check failures are treated as unknown, not dead; local discovery likewise treats an unverifiable container as not adoptable and falls through to create rather than failing acquire. `get()` remains an in-memory lookup for event-loop-safe tool paths. +- `LocalSandboxProvider` - Local filesystem execution. `acquire(thread_id)` returns a per-thread `LocalSandbox` (id `local:{thread_id}`) whose `path_mappings` resolve `/mnt/user-data/{workspace,uploads,outputs}` and `/mnt/acp-workspace` to that thread's host directories, so the public `Sandbox` API honours the `/mnt/user-data` contract uniformly with AIO. `acquire()` / `acquire(None)` keeps the legacy generic singleton (id `local`) for callers without a thread context. Per-thread sandboxes are held in an LRU cache (default 256 entries) guarded by a `threading.Lock`. Legacy global-custom mounts are gated by the same user-scoped skill discovery rule used for prompt/list visibility; providers must not infer visibility from raw directory presence alone. +- `AioSandboxProvider` (`packages/harness/deerflow/community/`) - Docker-based isolation. Active-cache and warm-pool entries are checked with the backend during acquire/reuse; definitively dead containers are dropped from all in-process maps so the thread can discover or create a fresh sandbox instead of reusing a stale client. Backend health-check failures are treated as unknown, not dead; local discovery likewise treats an unverifiable container as not adoptable and falls through to create rather than failing acquire. `get()` remains an in-memory lookup for event-loop-safe tool paths. Legacy global-custom mounts follow the same shared visibility helper as local and remote providers. - `BoxliteProvider` (`packages/harness/deerflow/community/boxlite/`) - BoxLite micro-VM isolation. The `boxlite` runtime is optional (`deerflow-harness[boxlite]`) and lazy-imported only when this provider is selected. The provider owns one private asyncio event loop on a daemon thread because BoxLite handles are loop-affine; sync `Sandbox` calls marshal onto that loop with `run_coroutine_threadsafe`. Boxes are named deterministically from `user_id:thread_id`, released into an in-process warm pool after each agent turn, and reclaimed only by the same user/thread. Warm-pool health checks use a short explicit timeout and forward that timeout through both BoxLite `exec(timeout=...)` and the private-loop `.result(timeout)` bridge so a hung VM cannot pin the per-thread acquire lock indefinitely. `sandbox.replicas` caps active + warm VMs per gateway process; if capacity is exhausted, only warm-pool VMs are evicted. `sandbox.idle_timeout` stops idle warm VMs after the configured seconds. `reset()` is intentionally a lightweight registry clear for `reset_sandbox_provider()` and does not close boxes, stop the idle reaper, or close the private loop; full teardown remains `shutdown()`. diff --git a/backend/packages/harness/deerflow/community/aio_sandbox/aio_sandbox_provider.py b/backend/packages/harness/deerflow/community/aio_sandbox/aio_sandbox_provider.py index 8874cd47e..663d1cbd3 100644 --- a/backend/packages/harness/deerflow/community/aio_sandbox/aio_sandbox_provider.py +++ b/backend/packages/harness/deerflow/community/aio_sandbox/aio_sandbox_provider.py @@ -36,10 +36,11 @@ from deerflow.community.warm_pool_lifecycle import ( IDLE_CHECK_INTERVAL as _SHARED_IDLE_CHECK_INTERVAL, ) from deerflow.config import get_app_config -from deerflow.config.paths import VIRTUAL_PATH_PREFIX, get_paths +from deerflow.config.paths import VIRTUAL_PATH_PREFIX, get_paths, join_host_path from deerflow.runtime.user_context import get_effective_user_id from deerflow.sandbox.sandbox import Sandbox from deerflow.sandbox.sandbox_provider import SandboxProvider +from deerflow.skills.storage import user_should_see_legacy_skills from .aio_sandbox import AioSandbox from .backend import SandboxBackend, wait_for_sandbox_ready, wait_for_sandbox_ready_async @@ -308,10 +309,10 @@ class AioSandboxProvider(WarmPoolLifecycleMixin[SandboxInfo], SandboxProvider): mounts.extend(self._get_thread_mounts(thread_id, user_id=user_id)) logger.info(f"Adding thread mounts for thread {thread_id}: {mounts}") - skills_mount = self._get_skills_mount() - if skills_mount: - mounts.append(skills_mount) - logger.info(f"Adding skills mount: {skills_mount}") + skills_mounts = self._get_skills_mounts(user_id=user_id) + if skills_mounts: + mounts.extend(skills_mounts) + logger.info(f"Adding skills mounts: {skills_mounts}") return mounts @@ -337,24 +338,76 @@ class AioSandboxProvider(WarmPoolLifecycleMixin[SandboxInfo], SandboxProvider): ] @staticmethod - def _get_skills_mount() -> tuple[str, str, bool] | None: - """Get the skills directory mount configuration. + def _get_skills_mounts(*, user_id: str | None = None) -> list[tuple[str, str, bool]]: + """Get skills directory mount configurations for three-way skills layout. - Mount source uses DEER_FLOW_HOST_SKILLS_PATH when running inside Docker (DooD) - so the host Docker daemon can resolve the path. + Mirrors ``LocalSandboxProvider._build_thread_path_mappings`` for AIO + sandboxes: public, per-user custom, and legacy (pre-migration + global-custom) skills are mounted to separate container subdirectories so + that ``Skill.get_container_path()`` category-aware paths resolve + correctly inside the sandbox. + + Mount sources use ``DEER_FLOW_HOST_SKILLS_PATH`` and + ``DEER_FLOW_HOST_BASE_DIR`` when running inside Docker (DooD) so the + host Docker daemon can resolve the paths. """ + mounts: list[tuple[str, str, bool]] = [] try: config = get_app_config() skills_path = config.skills.get_skills_path() container_path = config.skills.container_path - if skills_path.exists(): - # When running inside Docker with DooD, use host-side skills path. - host_skills = os.environ.get("DEER_FLOW_HOST_SKILLS_PATH") or str(skills_path) - return (host_skills, container_path, True) # Read-only for security + # When running inside Docker with DooD, use host-side skills path. + host_skills_root = os.environ.get("DEER_FLOW_HOST_SKILLS_PATH") or str(skills_path) + + # 1. Public skills: global, read-only — static, shared by all threads + public_skills_path = skills_path / "public" + if public_skills_path.exists(): + mounts.append( + ( + join_host_path(host_skills_root, "public"), + f"{container_path}/public", + True, + ) + ) + + # 2. Per-user custom skills: read-only, per-thread/per-user + effective_user_id = AioSandboxProvider._effective_acquire_user_id(user_id) + paths = get_paths() + user_custom_path = paths.user_custom_skills_dir(effective_user_id) + user_custom_path.mkdir(parents=True, exist_ok=True) + + host_user_custom = join_host_path( + str(paths.host_base_dir), + "users", + effective_user_id, + "skills", + "custom", + ) + mounts.append( + ( + host_user_custom, + f"{container_path}/custom", + True, + ) + ) + + # 3. Legacy (pre-migration global-custom) skills: only mount for + # users who have no per-user custom skills yet, mirroring + # ``UserScopedSkillStorage._iter_skill_files`` visibility rule. + legacy_skills_path = skills_path / "custom" + if user_should_see_legacy_skills(effective_user_id, host_path=str(skills_path)) and legacy_skills_path.exists(): + mounts.append( + ( + join_host_path(host_skills_root, "custom"), + f"{container_path}/legacy", + True, + ) + ) except Exception as e: - logger.warning(f"Could not setup skills mount: {e}") - return None + logger.warning("Could not setup skills mounts: %s", e) + + return mounts # ── Idle timeout management ────────────────────────────────────────── diff --git a/backend/packages/harness/deerflow/community/aio_sandbox/remote_backend.py b/backend/packages/harness/deerflow/community/aio_sandbox/remote_backend.py index ee9848d48..9c448c4a5 100644 --- a/backend/packages/harness/deerflow/community/aio_sandbox/remote_backend.py +++ b/backend/packages/harness/deerflow/community/aio_sandbox/remote_backend.py @@ -22,6 +22,7 @@ import logging import requests from deerflow.runtime.user_context import get_effective_user_id +from deerflow.skills.storage import user_should_see_legacy_skills from .backend import SandboxBackend from .sandbox_info import SandboxInfo @@ -145,6 +146,7 @@ class RemoteSandboxBackend(SandboxBackend): """POST /api/sandboxes → create Pod + Service.""" del extra_mounts effective_user_id = user_id or get_effective_user_id() + include_legacy_skills = user_should_see_legacy_skills(effective_user_id) try: resp = requests.post( f"{self._provisioner_url}/api/sandboxes", @@ -152,6 +154,7 @@ class RemoteSandboxBackend(SandboxBackend): "sandbox_id": sandbox_id, "thread_id": thread_id, "user_id": effective_user_id, + "include_legacy_skills": include_legacy_skills, }, timeout=30, ) diff --git a/backend/packages/harness/deerflow/sandbox/local/local_sandbox_provider.py b/backend/packages/harness/deerflow/sandbox/local/local_sandbox_provider.py index 192fc3271..c02f2cb32 100644 --- a/backend/packages/harness/deerflow/sandbox/local/local_sandbox_provider.py +++ b/backend/packages/harness/deerflow/sandbox/local/local_sandbox_provider.py @@ -6,6 +6,7 @@ from pathlib import Path from deerflow.sandbox.local.local_sandbox import LocalSandbox, PathMapping from deerflow.sandbox.sandbox import Sandbox from deerflow.sandbox.sandbox_provider import SandboxProvider +from deerflow.skills.storage import user_should_see_legacy_skills logger = logging.getLogger(__name__) @@ -310,8 +311,7 @@ class LocalSandboxProvider(SandboxProvider): skills_container_path = config.skills.container_path user_custom_path = paths.user_custom_skills_dir(effective_user_id) legacy_skills_path = config.skills.get_skills_path() / "custom" - user_has_no_custom_skills = not any(p.is_dir() and not p.name.startswith(".") for p in user_custom_path.iterdir()) if user_custom_path.exists() else True - if user_has_no_custom_skills and legacy_skills_path.exists() and any((legacy_skills_path / d / "SKILL.md").exists() for d in legacy_skills_path.iterdir() if d.is_dir() and not d.name.startswith(".")): + if user_should_see_legacy_skills(effective_user_id, host_path=str(config.skills.get_skills_path())) and legacy_skills_path.exists(): mappings.append( PathMapping( container_path=f"{skills_container_path}/legacy", diff --git a/backend/packages/harness/deerflow/skills/storage/__init__.py b/backend/packages/harness/deerflow/skills/storage/__init__.py index f91a1ace0..1d8940140 100644 --- a/backend/packages/harness/deerflow/skills/storage/__init__.py +++ b/backend/packages/harness/deerflow/skills/storage/__init__.py @@ -12,6 +12,7 @@ from collections import OrderedDict from deerflow.skills.storage.local_skill_storage import LocalSkillStorage from deerflow.skills.storage.skill_storage import SkillStorage from deerflow.skills.storage.user_scoped_skill_storage import UserScopedSkillStorage +from deerflow.skills.types import SkillCategory logger = logging.getLogger(__name__) @@ -142,6 +143,22 @@ def get_or_new_user_skill_storage(user_id: str, **kwargs) -> SkillStorage: return cached +def user_should_see_legacy_skills(user_id: str, **kwargs) -> bool: + """Return whether discovery exposes any LEGACY skills for this user. + + Sandbox mounts must not be more permissive than skill discovery. This + helper centralizes that contract so local, AIO, and remote providers all + follow the same visibility rule. + """ + if kwargs: + from deerflow.config.paths import make_safe_user_id + + storage = UserScopedSkillStorage(make_safe_user_id(user_id), **kwargs) + else: + storage = get_or_new_user_skill_storage(user_id) + return any((skill.category.value if hasattr(skill.category, "value") else skill.category) == SkillCategory.LEGACY.value for skill in storage.load_skills(enabled_only=False)) + + def reset_skill_storage() -> None: """Clear all cached storage instances (used in tests and hot-reload scenarios).""" global _default_skill_storage, _default_skill_storage_config @@ -180,6 +197,7 @@ __all__ = [ "UserScopedSkillStorage", "get_or_new_skill_storage", "get_or_new_user_skill_storage", + "user_should_see_legacy_skills", "reset_skill_storage", "reset_user_skill_storage", ] diff --git a/backend/tests/test_aio_sandbox_provider.py b/backend/tests/test_aio_sandbox_provider.py index e0349f0a0..df1745873 100644 --- a/backend/tests/test_aio_sandbox_provider.py +++ b/backend/tests/test_aio_sandbox_provider.py @@ -373,6 +373,7 @@ def test_remote_backend_create_forwards_effective_user_id(monkeypatch): return _Response() monkeypatch.setattr(remote_mod.requests, "post", _post) + monkeypatch.setattr(remote_mod, "user_should_see_legacy_skills", lambda user_id: True) try: backend.create("thread-42", "sandbox-42") @@ -384,6 +385,7 @@ def test_remote_backend_create_forwards_effective_user_id(monkeypatch): "sandbox_id": "sandbox-42", "thread_id": "thread-42", "user_id": "user-7", + "include_legacy_skills": True, } @@ -406,10 +408,12 @@ def test_remote_backend_create_prefers_explicit_user_id(monkeypatch): monkeypatch.setattr(remote_mod.requests, "post", _post) monkeypatch.setattr(remote_mod, "get_effective_user_id", lambda: "default") + monkeypatch.setattr(remote_mod, "user_should_see_legacy_skills", lambda user_id: False) backend.create("thread-42", "sandbox-42", user_id="ou-user") assert posted["json"]["user_id"] == "ou-user" + assert posted["json"]["include_legacy_skills"] is False # ── Sandbox client teardown (#2872) ────────────────────────────────────────── diff --git a/backend/tests/test_provisioner_pvc_volumes.py b/backend/tests/test_provisioner_pvc_volumes.py index d5b66a2c7..2d0668de3 100644 --- a/backend/tests/test_provisioner_pvc_volumes.py +++ b/backend/tests/test_provisioner_pvc_volumes.py @@ -1,40 +1,84 @@ -"""Regression tests for provisioner PVC volume support.""" +"""Regression tests for provisioner three-way skills + PVC volume support.""" # ── _build_volumes ───────────────────────────────────────────────────── class TestBuildVolumes: - """Tests for _build_volumes: PVC vs hostPath selection.""" + """Tests for _build_volumes: hostPath three-way vs PVC fallback.""" - def test_default_uses_hostpath_for_skills(self, provisioner_module): - """When SKILLS_PVC_NAME is empty, skills volume should use hostPath.""" + # ── hostPath mode (default) ──────────────────────────────────────── + + def test_hostpath_without_legacy_returns_three_volumes(self, provisioner_module): + """hostPath mode omits legacy volume unless the backend requests it.""" provisioner_module.SKILLS_PVC_NAME = "" - volumes = provisioner_module._build_volumes("thread-1") - skills_vol = volumes[0] - assert skills_vol.host_path is not None - assert skills_vol.host_path.path == provisioner_module.SKILLS_HOST_PATH - assert skills_vol.host_path.type == "Directory" - assert skills_vol.persistent_volume_claim is None - - def test_default_uses_hostpath_for_userdata(self, provisioner_module): - """When USERDATA_PVC_NAME is empty, user-data volume should use hostPath.""" provisioner_module.USERDATA_PVC_NAME = "" volumes = provisioner_module._build_volumes("thread-1") - userdata_vol = volumes[1] - assert userdata_vol.host_path is not None - assert userdata_vol.persistent_volume_claim is None + assert len(volumes) == 3 + + def test_hostpath_skills_public_volume(self, provisioner_module): + """First skills volume mounts public/ subdirectory.""" + provisioner_module.SKILLS_PVC_NAME = "" + volumes = provisioner_module._build_volumes("thread-1") + pub = volumes[0] + assert pub.name == "skills-public" + assert pub.host_path is not None + assert pub.host_path.path.endswith("/public") + assert pub.host_path.type == "Directory" + assert pub.persistent_volume_claim is None + + def test_hostpath_skills_custom_volume(self, provisioner_module): + """Second skills volume mounts per-user custom directory.""" + provisioner_module.SKILLS_PVC_NAME = "" + volumes = provisioner_module._build_volumes("thread-1", user_id="user-7") + custom = volumes[1] + assert custom.name == "skills-custom" + assert custom.host_path is not None + assert "users/user-7/skills/custom" in custom.host_path.path + assert custom.host_path.type == "DirectoryOrCreate" + + def test_hostpath_skills_legacy_volume(self, provisioner_module): + """Legacy global-custom directory is mounted only when requested.""" + provisioner_module.SKILLS_PVC_NAME = "" + volumes = provisioner_module._build_volumes( + "thread-1", + include_legacy_skills=True, + ) + legacy = volumes[2] + assert legacy.name == "skills-legacy" + assert legacy.host_path is not None + assert legacy.host_path.path.endswith("/custom") + assert legacy.host_path.type == "Directory" + + def test_hostpath_without_legacy_has_no_legacy_volume(self, provisioner_module): + """Fresh installs should not require a missing global legacy directory.""" + provisioner_module.SKILLS_PVC_NAME = "" + volumes = provisioner_module._build_volumes("thread-1") + assert [volume.name for volume in volumes] == [ + "skills-public", + "skills-custom", + "user-data", + ] def test_hostpath_userdata_includes_thread_id(self, provisioner_module): """hostPath user-data path should include thread_id.""" provisioner_module.USERDATA_PVC_NAME = "" volumes = provisioner_module._build_volumes("my-thread-42") - userdata_vol = volumes[1] + userdata_vol = volumes[-1] path = userdata_vol.host_path.path assert "my-thread-42" in path assert path.endswith("user-data") assert userdata_vol.host_path.type == "DirectoryOrCreate" + # ── PVC mode (single-volume fallback) ────────────────────────────── + + def test_pvc_returns_two_volumes(self, provisioner_module): + """PVC mode falls back to 1 skills volume + 1 user-data volume.""" + provisioner_module.SKILLS_PVC_NAME = "my-skills-pvc" + provisioner_module.USERDATA_PVC_NAME = "" + volumes = provisioner_module._build_volumes("thread-1") + assert len(volumes) == 2 + def test_skills_pvc_overrides_hostpath(self, provisioner_module): """When SKILLS_PVC_NAME is set, skills volume should use PVC.""" provisioner_module.SKILLS_PVC_NAME = "my-skills-pvc" @@ -49,7 +93,7 @@ class TestBuildVolumes: """When USERDATA_PVC_NAME is set, user-data volume should use PVC.""" provisioner_module.USERDATA_PVC_NAME = "my-userdata-pvc" volumes = provisioner_module._build_volumes("thread-1") - userdata_vol = volumes[1] + userdata_vol = volumes[-1] assert userdata_vol.persistent_volume_claim is not None assert userdata_vol.persistent_volume_claim.claim_name == "my-userdata-pvc" assert userdata_vol.host_path is None @@ -60,78 +104,112 @@ class TestBuildVolumes: provisioner_module.USERDATA_PVC_NAME = "userdata-pvc" volumes = provisioner_module._build_volumes("thread-1") assert volumes[0].persistent_volume_claim is not None - assert volumes[1].persistent_volume_claim is not None + assert volumes[-1].persistent_volume_claim is not None - def test_returns_two_volumes(self, provisioner_module): - """Should always return exactly two volumes.""" - provisioner_module.SKILLS_PVC_NAME = "" - provisioner_module.USERDATA_PVC_NAME = "" - assert len(provisioner_module._build_volumes("t")) == 2 - - provisioner_module.SKILLS_PVC_NAME = "a" - provisioner_module.USERDATA_PVC_NAME = "b" - assert len(provisioner_module._build_volumes("t")) == 2 - - def test_volume_names_are_stable(self, provisioner_module): - """Volume names must stay 'skills' and 'user-data'.""" + def test_pvc_volume_names_are_stable(self, provisioner_module): + """PVC mode volume names must stay 'skills' and 'user-data'.""" + provisioner_module.SKILLS_PVC_NAME = "x" volumes = provisioner_module._build_volumes("thread-1") assert volumes[0].name == "skills" - assert volumes[1].name == "user-data" + assert volumes[-1].name == "user-data" # ── _build_volume_mounts ─────────────────────────────────────────────── class TestBuildVolumeMounts: - """Tests for _build_volume_mounts: mount paths and subPath behavior.""" + """Tests for _build_volume_mounts: three-way mount paths and subPath.""" - def test_default_no_subpath(self, provisioner_module): + # ── hostPath mode ────────────────────────────────────────────────── + + def test_hostpath_without_legacy_returns_three_mounts(self, provisioner_module): + """hostPath mode omits legacy mount unless the backend requests it.""" + provisioner_module.SKILLS_PVC_NAME = "" + provisioner_module.USERDATA_PVC_NAME = "" + mounts = provisioner_module._build_volume_mounts("thread-1") + assert len(mounts) == 3 + + def test_hostpath_skills_public_mount(self, provisioner_module): + """Public skills mount at /mnt/skills/public, read-only.""" + provisioner_module.SKILLS_PVC_NAME = "" + mounts = provisioner_module._build_volume_mounts("thread-1") + assert mounts[0].name == "skills-public" + assert mounts[0].mount_path == "/mnt/skills/public" + assert mounts[0].read_only is True + + def test_hostpath_skills_custom_mount(self, provisioner_module): + """Per-user custom skills mount at /mnt/skills/custom, read-only.""" + provisioner_module.SKILLS_PVC_NAME = "" + mounts = provisioner_module._build_volume_mounts("thread-1") + assert mounts[1].name == "skills-custom" + assert mounts[1].mount_path == "/mnt/skills/custom" + assert mounts[1].read_only is True + + def test_hostpath_skills_legacy_mount(self, provisioner_module): + """Legacy skills mount at /mnt/skills/legacy, read-only.""" + provisioner_module.SKILLS_PVC_NAME = "" + mounts = provisioner_module._build_volume_mounts( + "thread-1", + include_legacy_skills=True, + ) + assert mounts[2].name == "skills-legacy" + assert mounts[2].mount_path == "/mnt/skills/legacy" + assert mounts[2].read_only is True + + def test_hostpath_without_legacy_has_no_legacy_mount(self, provisioner_module): + """Users with custom skills should not see hidden legacy content in the sandbox.""" + provisioner_module.SKILLS_PVC_NAME = "" + mounts = provisioner_module._build_volume_mounts("thread-1") + assert [mount.name for mount in mounts] == [ + "skills-public", + "skills-custom", + "user-data", + ] + + def test_hostpath_userdata_read_write(self, provisioner_module): + """User-data mount should always be read-write.""" + provisioner_module.SKILLS_PVC_NAME = "" + mounts = provisioner_module._build_volume_mounts("thread-1") + userdata = mounts[-1] + assert userdata.name == "user-data" + assert userdata.mount_path == "/mnt/user-data" + assert userdata.read_only is False + + # ── PVC mode ─────────────────────────────────────────────────────── + + def test_pvc_returns_two_mounts(self, provisioner_module): + """PVC mode falls back to 1 skills mount + 1 user-data mount.""" + provisioner_module.SKILLS_PVC_NAME = "x" + mounts = provisioner_module._build_volume_mounts("thread-1") + assert len(mounts) == 2 + + def test_pvc_skills_mount_is_single_root(self, provisioner_module): + """PVC mode skills mount is at /mnt/skills.""" + provisioner_module.SKILLS_PVC_NAME = "x" + mounts = provisioner_module._build_volume_mounts("thread-1") + assert mounts[0].mount_path == "/mnt/skills" + + def test_pvc_no_subpath_on_userdata(self, provisioner_module): """hostPath mode should not set sub_path on user-data mount.""" provisioner_module.USERDATA_PVC_NAME = "" mounts = provisioner_module._build_volume_mounts("thread-1") - userdata_mount = mounts[1] + userdata_mount = mounts[-1] assert userdata_mount.sub_path is None def test_pvc_sets_user_scoped_subpath(self, provisioner_module): """PVC mode should include user_id in the user-data subPath.""" provisioner_module.USERDATA_PVC_NAME = "my-pvc" mounts = provisioner_module._build_volume_mounts("thread-42", user_id="user-7") - userdata_mount = mounts[1] + userdata_mount = mounts[-1] assert userdata_mount.sub_path == "deer-flow/users/user-7/threads/thread-42/user-data" def test_pvc_defaults_to_default_user_subpath(self, provisioner_module): """Older callers should still land under a stable default user namespace.""" provisioner_module.USERDATA_PVC_NAME = "my-pvc" mounts = provisioner_module._build_volume_mounts("thread-42") - userdata_mount = mounts[1] + userdata_mount = mounts[-1] assert userdata_mount.sub_path == "deer-flow/users/default/threads/thread-42/user-data" - def test_skills_mount_read_only(self, provisioner_module): - """Skills mount should always be read-only.""" - mounts = provisioner_module._build_volume_mounts("thread-1") - assert mounts[0].read_only is True - - def test_userdata_mount_read_write(self, provisioner_module): - """User-data mount should always be read-write.""" - mounts = provisioner_module._build_volume_mounts("thread-1") - assert mounts[1].read_only is False - - def test_mount_paths_are_stable(self, provisioner_module): - """Mount paths must stay /mnt/skills and /mnt/user-data.""" - mounts = provisioner_module._build_volume_mounts("thread-1") - assert mounts[0].mount_path == "/mnt/skills" - assert mounts[1].mount_path == "/mnt/user-data" - - def test_mount_names_match_volumes(self, provisioner_module): - """Mount names should match the volume names.""" - mounts = provisioner_module._build_volume_mounts("thread-1") - assert mounts[0].name == "skills" - assert mounts[1].name == "user-data" - - def test_returns_two_mounts(self, provisioner_module): - """Should always return exactly two mounts.""" - assert len(provisioner_module._build_volume_mounts("t")) == 2 - # ── _build_pod integration ───────────────────────────────────────────── @@ -139,16 +217,52 @@ class TestBuildVolumeMounts: class TestBuildPodVolumes: """Integration: _build_pod should wire volumes and mounts correctly.""" - def test_pod_spec_has_volumes(self, provisioner_module): - """Pod spec should contain exactly 2 volumes.""" + def test_pod_hostpath_without_legacy_has_three_volumes(self, provisioner_module): + """hostPath Pod spec should omit legacy volume by default.""" provisioner_module.SKILLS_PVC_NAME = "" provisioner_module.USERDATA_PVC_NAME = "" pod = provisioner_module._build_pod("sandbox-1", "thread-1") + assert len(pod.spec.volumes) == 3 + + def test_pod_hostpath_without_legacy_has_three_mounts(self, provisioner_module): + """hostPath container should omit legacy mount by default.""" + provisioner_module.SKILLS_PVC_NAME = "" + provisioner_module.USERDATA_PVC_NAME = "" + pod = provisioner_module._build_pod("sandbox-1", "thread-1") + assert len(pod.spec.containers[0].volume_mounts) == 3 + + def test_pod_hostpath_with_legacy_has_four_volumes(self, provisioner_module): + """Legacy volume should be present when the backend requests it.""" + provisioner_module.SKILLS_PVC_NAME = "" + provisioner_module.USERDATA_PVC_NAME = "" + pod = provisioner_module._build_pod( + "sandbox-1", + "thread-1", + include_legacy_skills=True, + ) + assert len(pod.spec.volumes) == 4 + + def test_pod_hostpath_with_legacy_has_four_mounts(self, provisioner_module): + """Legacy mount should be present when the backend requests it.""" + provisioner_module.SKILLS_PVC_NAME = "" + provisioner_module.USERDATA_PVC_NAME = "" + pod = provisioner_module._build_pod( + "sandbox-1", + "thread-1", + include_legacy_skills=True, + ) + assert len(pod.spec.containers[0].volume_mounts) == 4 + + def test_pod_pvc_has_two_volumes(self, provisioner_module): + """PVC Pod spec should contain exactly 2 volumes.""" + provisioner_module.SKILLS_PVC_NAME = "skills-pvc" + provisioner_module.USERDATA_PVC_NAME = "" + pod = provisioner_module._build_pod("sandbox-1", "thread-1") assert len(pod.spec.volumes) == 2 - def test_pod_spec_has_volume_mounts(self, provisioner_module): - """Container should have exactly 2 volume mounts.""" - provisioner_module.SKILLS_PVC_NAME = "" + def test_pod_pvc_has_two_mounts(self, provisioner_module): + """PVC container should have exactly 2 volume mounts.""" + provisioner_module.SKILLS_PVC_NAME = "skills-pvc" provisioner_module.USERDATA_PVC_NAME = "" pod = provisioner_module._build_pod("sandbox-1", "thread-1") assert len(pod.spec.containers[0].volume_mounts) == 2 @@ -159,6 +273,20 @@ class TestBuildPodVolumes: provisioner_module.USERDATA_PVC_NAME = "userdata-pvc" pod = provisioner_module._build_pod("sandbox-1", "thread-1", user_id="user-7") assert pod.spec.volumes[0].persistent_volume_claim is not None - assert pod.spec.volumes[1].persistent_volume_claim is not None - userdata_mount = pod.spec.containers[0].volume_mounts[1] + assert pod.spec.volumes[-1].persistent_volume_claim is not None + userdata_mount = pod.spec.containers[0].volume_mounts[-1] assert userdata_mount.sub_path == "deer-flow/users/user-7/threads/thread-1/user-data" + + def test_pod_three_way_skills_mount_paths(self, provisioner_module): + """Ensure public/custom/legacy mount paths are correct.""" + provisioner_module.SKILLS_PVC_NAME = "" + provisioner_module.USERDATA_PVC_NAME = "" + pod = provisioner_module._build_pod( + "sandbox-1", + "thread-1", + include_legacy_skills=True, + ) + mount_paths = {m.name: m.mount_path for m in pod.spec.containers[0].volume_mounts} + assert mount_paths["skills-public"] == "/mnt/skills/public" + assert mount_paths["skills-custom"] == "/mnt/skills/custom" + assert mount_paths["skills-legacy"] == "/mnt/skills/legacy" diff --git a/backend/tests/test_provisioner_request_threading.py b/backend/tests/test_provisioner_request_threading.py index 24ea81b3a..a834d085d 100644 --- a/backend/tests/test_provisioner_request_threading.py +++ b/backend/tests/test_provisioner_request_threading.py @@ -27,6 +27,7 @@ class _RecordingCoreV1: self.ready_after_service_reads = ready_after_service_reads or {} self.service_read_counts: dict[str, int] = {} self.created_pods: list[str] = [] + self.created_pod_specs: dict[str, object] = {} self.created_services: list[str] = [] def _record_k8s_call(self) -> None: @@ -58,6 +59,7 @@ class _RecordingCoreV1: self._record_k8s_call() sandbox_id = pod.metadata.labels["sandbox-id"] self.created_pods.append(sandbox_id) + self.created_pod_specs[sandbox_id] = pod def create_namespaced_service(self, _namespace: str, service) -> None: self._record_k8s_call() @@ -157,3 +159,46 @@ async def test_sandbox_business_routes_run_k8s_client_off_event_loop_thread( if expected_created_sandbox is not None: assert fake_core_v1.created_pods == [expected_created_sandbox] assert fake_core_v1.created_services == [expected_created_sandbox] + + +@pytest.mark.parametrize( + ("include_legacy_skills", "expected_mount_names"), + [ + ( + False, + ["skills-public", "skills-custom", "user-data"], + ), + ( + True, + ["skills-public", "skills-custom", "skills-legacy", "user-data"], + ), + ], + ids=["without-legacy", "with-legacy"], +) +def test_create_sandbox_route_builds_expected_skills_mount_layout( + include_legacy_skills: bool, + expected_mount_names: list[str], + monkeypatch: pytest.MonkeyPatch, + provisioner_module, +) -> None: + fake_core_v1 = _RecordingCoreV1( + event_loop_thread_id=-1, + ready_after_service_reads={"sandbox-layout": 1}, + ) + monkeypatch.setattr(provisioner_module, "core_v1", fake_core_v1) + + response = provisioner_module.create_sandbox( + provisioner_module.CreateSandboxRequest( + sandbox_id="sandbox-layout", + thread_id="thread-1", + user_id="user-1", + include_legacy_skills=include_legacy_skills, + ) + ) + + assert response.status == "Running" + pod = fake_core_v1.created_pod_specs["sandbox-layout"] + volume_names = [volume.name for volume in pod.spec.volumes] + mount_names = [mount.name for mount in pod.spec.containers[0].volume_mounts] + assert volume_names == expected_mount_names + assert mount_names == expected_mount_names diff --git a/backend/tests/test_remote_sandbox_backend.py b/backend/tests/test_remote_sandbox_backend.py index 6550cae10..f6abe4e08 100644 --- a/backend/tests/test_remote_sandbox_backend.py +++ b/backend/tests/test_remote_sandbox_backend.py @@ -3,8 +3,11 @@ from __future__ import annotations import pytest import requests +import deerflow.skills.storage as storage_mod +from deerflow.community.aio_sandbox import remote_backend as remote_backend_mod from deerflow.community.aio_sandbox.remote_backend import RemoteSandboxBackend from deerflow.community.aio_sandbox.sandbox_info import SandboxInfo +from deerflow.skills.types import SkillCategory class _StubResponse: @@ -123,6 +126,25 @@ def test_provisioner_list_skips_non_dict_sandbox_entries(monkeypatch): assert infos[0].sandbox_url == "http://k3s:31001" +@pytest.mark.parametrize( + ("categories", "expected"), + [ + ([SkillCategory.LEGACY], True), + (["legacy"], True), + ([SkillCategory.CUSTOM], False), + ], +) +def test_user_should_see_legacy_skills_follows_storage_visibility_rule(monkeypatch, categories, expected): + class _Storage: + def load_skills(self, *, enabled_only: bool = False): + assert enabled_only is False + return [type("SkillStub", (), {"category": category})() for category in categories] + + monkeypatch.setattr(storage_mod, "get_or_new_user_skill_storage", lambda user_id: _Storage()) + + assert storage_mod.user_should_see_legacy_skills("user-1") is expected + + @pytest.mark.parametrize("expected_user_id", [None, "owner-1"]) def test_create_delegates_to_provisioner_create(monkeypatch, expected_user_id): backend = RemoteSandboxBackend("http://provisioner:8002") @@ -148,6 +170,7 @@ def test_create_delegates_to_provisioner_create(monkeypatch, expected_user_id): def test_provisioner_create_returns_sandbox_info(monkeypatch): backend = RemoteSandboxBackend("http://provisioner:8002") + monkeypatch.setattr(remote_backend_mod, "user_should_see_legacy_skills", lambda user_id: True) def mock_post(url: str, json: dict, timeout: int): assert url == "http://provisioner:8002/api/sandboxes" @@ -155,6 +178,7 @@ def test_provisioner_create_returns_sandbox_info(monkeypatch): "sandbox_id": "abc123", "thread_id": "thread-1", "user_id": "test-user-autouse", + "include_legacy_skills": True, } assert timeout == 30 return _StubResponse(payload={"sandbox_id": "abc123", "sandbox_url": "http://k3s:31001"}) @@ -168,6 +192,7 @@ def test_provisioner_create_returns_sandbox_info(monkeypatch): def test_provisioner_create_accepts_anonymous_thread_id(monkeypatch): backend = RemoteSandboxBackend("http://provisioner:8002") + monkeypatch.setattr(remote_backend_mod, "user_should_see_legacy_skills", lambda user_id: False) def mock_post(url: str, json: dict, timeout: int): assert url == "http://provisioner:8002/api/sandboxes" @@ -175,6 +200,7 @@ def test_provisioner_create_accepts_anonymous_thread_id(monkeypatch): "sandbox_id": "anon123", "thread_id": None, "user_id": "test-user-autouse", + "include_legacy_skills": False, } assert timeout == 30 return _StubResponse(payload={"sandbox_id": "anon123", "sandbox_url": "http://k3s:31002"}) @@ -188,6 +214,7 @@ def test_provisioner_create_accepts_anonymous_thread_id(monkeypatch): def test_provisioner_create_raises_runtime_error_on_request_exception(monkeypatch): backend = RemoteSandboxBackend("http://provisioner:8002") + monkeypatch.setattr(remote_backend_mod, "user_should_see_legacy_skills", lambda user_id: False) def mock_post(url: str, json: dict, timeout: int): raise requests.RequestException("boom") diff --git a/backend/tests/test_three_way_skills_mount_e2e.py b/backend/tests/test_three_way_skills_mount_e2e.py new file mode 100644 index 000000000..bb166740a --- /dev/null +++ b/backend/tests/test_three_way_skills_mount_e2e.py @@ -0,0 +1,336 @@ +"""End-to-end tests for three-way skills mount across sandbox providers. + +Verifies that (a) public, (b) per-user custom, and (c) legacy global-custom +skills all resolve to correct container paths that the sandbox providers +actually mount — covering ``LocalSandboxProvider`` and +``AioSandboxProvider`` (DooD / local-backend path). + +Includes a full-pipeline test that exercises the actual path the model +takes: ``UserScopedSkillStorage`` category assignment → ``Skill.get_container_file_path()`` → ``sandbox.read_file()``. +""" + +import importlib +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import patch + +import pytest + +from deerflow.config.paths import Paths +from deerflow.sandbox.local.local_sandbox import PathMapping +from deerflow.sandbox.local.local_sandbox_provider import LocalSandboxProvider +from deerflow.skills.types import SKILL_MD_FILE, Skill, SkillCategory + +_AIO_MODULE = "deerflow.community.aio_sandbox.aio_sandbox_provider" +_AIO_GET_CONFIG = f"{_AIO_MODULE}.get_app_config" + + +def _write_skill(base: Path, name: str, description: str = "test skill") -> Path: + skill_dir = base / name + skill_dir.mkdir(parents=True, exist_ok=True) + skill_md = skill_dir / SKILL_MD_FILE + skill_md.write_text( + f"---\nname: {name}\ndescription: {description}\n---\n\n# {name}\n", + encoding="utf-8", + ) + return skill_md + + +def _build_config(skills_root: Path): + from deerflow.config.sandbox_config import SandboxConfig + + return SimpleNamespace( + skills=SimpleNamespace( + container_path="/mnt/skills", + get_skills_path=lambda sk=skills_root: sk, + use="deerflow.skills.storage.local_skill_storage:LocalSkillStorage", + ), + sandbox=SandboxConfig( + use="deerflow.sandbox.local:LocalSandboxProvider", + mounts=[], + ), + ) + + +def _local_mounts(provider: LocalSandboxProvider, thread_id: str, user_id: str) -> dict[str, PathMapping]: + mappings = list(provider._path_mappings) + provider._build_thread_path_mappings(thread_id, user_id=user_id) + return {m.container_path: m for m in mappings} + + +@pytest.fixture +def skills_fs(tmp_path: Path) -> dict: + root = tmp_path / "skills" + pub = root / "public" + legacy = root / "custom" + users_dir = tmp_path / "users" + user_custom = users_dir / "user-1" / "skills" / "custom" + + return { + "root": root, + "public": pub, + "legacy_global": legacy, + "user_custom": user_custom, + "users_dir": users_dir, + "pub_skill": _write_skill(pub, "pub-skill", "public skill"), + "legacy_skill": _write_skill(legacy, "leg-skill", "legacy skill"), + "user_skill": _write_skill(user_custom, "usr-skill", "user custom skill"), + } + + +@pytest.fixture +def aio_mod(): + return importlib.import_module(_AIO_MODULE) + + +class TestThreeWayMountEndToEnd: + # ── LocalSandboxProvider: mount structure ────────────────────────── + + def test_local_public_skill_mounted(self, skills_fs): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + idx = _local_mounts(provider, "thread-1", user_id="user-1") + assert "/mnt/skills/public" in idx + assert idx["/mnt/skills/public"].read_only is True + + def test_local_per_user_custom_skill_mounted(self, skills_fs): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + idx = _local_mounts(provider, "thread-1", user_id="user-1") + assert "/mnt/skills/custom" in idx + assert str(skills_fs["user_custom"]) in idx["/mnt/skills/custom"].local_path + + def test_local_legacy_mounted_for_user_without_custom(self, skills_fs): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + idx = _local_mounts(provider, "thread-1", user_id="noob") + assert "/mnt/skills/legacy" in idx + assert str(skills_fs["legacy_global"]) in idx["/mnt/skills/legacy"].local_path + + def test_local_legacy_not_mounted_when_user_has_custom(self, skills_fs): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + idx = _local_mounts(provider, "thread-1", user_id="user-1") + assert "/mnt/skills/legacy" not in idx + + def test_local_legacy_still_mounted_when_user_has_only_non_skill_subdir(self, skills_fs): + (skills_fs["users_dir"] / "ghost" / "skills" / "custom" / "dangling-dir").mkdir(parents=True, exist_ok=True) + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + idx = _local_mounts(provider, "thread-1", user_id="ghost") + assert "/mnt/skills/legacy" in idx + + # ── LocalSandboxProvider: read_file on container paths ───────────── + + def test_local_read_file_resolves_public_and_custom(self, skills_fs): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + sid = provider.acquire("thread-1", user_id="user-1") + sandbox = provider.get(sid) + assert "pub-skill" in sandbox.read_file("/mnt/skills/public/pub-skill/SKILL.md") + assert "usr-skill" in sandbox.read_file("/mnt/skills/custom/usr-skill/SKILL.md") + + def test_local_read_file_resolves_legacy_skill(self, skills_fs): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + sid = provider.acquire("thread-1", user_id="noob") + sandbox = provider.get(sid) + assert "leg-skill" in sandbox.read_file("/mnt/skills/legacy/leg-skill/SKILL.md") + + # ── Full pipeline: registry → container path → sandbox read ──────── + + def test_registry_to_sandbox_full_pipeline(self, skills_fs): + """Model's exact path: storage category → get_container_file_path → sandbox.read_file.""" + from deerflow.skills.storage.user_scoped_skill_storage import UserScopedSkillStorage + + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + + with patch("deerflow.config.get_app_config", return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + provider = LocalSandboxProvider() + sid_user = provider.acquire("t1", user_id="user-1") + sid_noob = provider.acquire("t2", user_id="noob") + sandbox_user = provider.get(sid_user) + sandbox_noob = provider.get(sid_noob) + + # user-1 storage: sees public + custom, no legacy + with patch("deerflow.config.paths.get_paths", return_value=paths): + storage = UserScopedSkillStorage(user_id="user-1", host_path=str(skills_fs["root"])) + skills = list(storage._iter_skill_files()) + by_name = {sf.parent.name: (cat, sf) for cat, _root, sf in skills} + + # public + assert "pub-skill" in by_name + cat, _ = by_name["pub-skill"] + assert cat == SkillCategory.PUBLIC + s = Skill(name="pub-skill", description="p", license=None, skill_dir=skills_fs["public"] / "pub-skill", skill_file=skills_fs["pub_skill"], relative_path=Path("pub-skill"), category=cat) + cp = s.get_container_file_path("/mnt/skills") + assert cp == "/mnt/skills/public/pub-skill/SKILL.md" + assert "pub-skill" in sandbox_user.read_file(cp) + + # custom + assert "usr-skill" in by_name + cat, _ = by_name["usr-skill"] + assert cat == SkillCategory.CUSTOM + s = Skill(name="usr-skill", description="u", license=None, skill_dir=skills_fs["user_custom"] / "usr-skill", skill_file=skills_fs["user_skill"], relative_path=Path("usr-skill"), category=cat) + cp = s.get_container_file_path("/mnt/skills") + assert cp == "/mnt/skills/custom/usr-skill/SKILL.md" + assert "usr-skill" in sandbox_user.read_file(cp) + + # noob storage: sees public + legacy (no per-user custom) + with patch("deerflow.config.paths.get_paths", return_value=paths): + storage = UserScopedSkillStorage(user_id="noob", host_path=str(skills_fs["root"])) + skills = list(storage._iter_skill_files()) + by_name = {sf.parent.name: (cat, sf) for cat, _root, sf in skills} + + assert "leg-skill" in by_name + cat, _ = by_name["leg-skill"] + assert cat == SkillCategory.LEGACY + s = Skill(name="leg-skill", description="l", license=None, skill_dir=skills_fs["legacy_global"] / "leg-skill", skill_file=skills_fs["legacy_skill"], relative_path=Path("leg-skill"), category=cat) + cp = s.get_container_file_path("/mnt/skills") + assert cp == "/mnt/skills/legacy/leg-skill/SKILL.md" + assert "leg-skill" in sandbox_noob.read_file(cp) + + # ── AioSandboxProvider ────────────────────────────────────────────── + + def test_aio_public_skill_mount(self, skills_fs, aio_mod): + cfg = _build_config(skills_fs["root"]) + with patch(_AIO_GET_CONFIG, return_value=cfg): + mounts = aio_mod.AioSandboxProvider._get_skills_mounts(user_id="user-1") + idx = {m[1]: m for m in mounts} + assert "/mnt/skills/public" in idx + + def test_aio_per_user_custom_skill_mount(self, skills_fs, aio_mod, monkeypatch): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + monkeypatch.setattr(aio_mod, "get_paths", lambda: paths) + with patch(_AIO_GET_CONFIG, return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + mounts = aio_mod.AioSandboxProvider._get_skills_mounts(user_id="user-1") + idx = {m[1]: m for m in mounts} + assert "/mnt/skills/custom" in idx + host, _, _ = idx["/mnt/skills/custom"] + assert "users/user-1/skills/custom" in host.replace("\\", "/") + + def test_aio_legacy_mounted_for_user_without_custom(self, skills_fs, aio_mod, monkeypatch): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + monkeypatch.setattr(aio_mod, "get_paths", lambda: paths) + with patch(_AIO_GET_CONFIG, return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + mounts = aio_mod.AioSandboxProvider._get_skills_mounts(user_id="noob") + idx = {m[1]: m for m in mounts} + assert "/mnt/skills/legacy" in idx + + def test_aio_legacy_not_mounted_when_user_has_custom(self, skills_fs, aio_mod, monkeypatch): + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + monkeypatch.setattr(aio_mod, "get_paths", lambda: paths) + with patch(_AIO_GET_CONFIG, return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + mounts = aio_mod.AioSandboxProvider._get_skills_mounts(user_id="user-1") + idx = {m[1]: m for m in mounts} + assert "/mnt/skills/legacy" not in idx + + def test_aio_legacy_still_mounted_when_user_has_only_non_skill_subdir(self, skills_fs, aio_mod, monkeypatch): + (skills_fs["users_dir"] / "ghost" / "skills" / "custom" / "dangling-dir").mkdir(parents=True, exist_ok=True) + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + monkeypatch.setattr(aio_mod, "get_paths", lambda: paths) + with patch(_AIO_GET_CONFIG, return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + mounts = aio_mod.AioSandboxProvider._get_skills_mounts(user_id="ghost") + idx = {m[1]: m for m in mounts} + assert "/mnt/skills/legacy" in idx + + # ── AIO → Docker --mount translation ─────────────────────────────── + + def test_aio_extra_mounts_translate_to_docker_bind_mounts(self, skills_fs, aio_mod, monkeypatch): + """extra_mounts → _format_container_mount → correct Docker --mount args.""" + from deerflow.community.aio_sandbox.local_backend import _format_container_mount + + cfg = _build_config(skills_fs["root"]) + paths = Paths(base_dir=skills_fs["users_dir"].parent) + monkeypatch.setattr(aio_mod, "get_paths", lambda: paths) + + with patch(_AIO_GET_CONFIG, return_value=cfg), patch("deerflow.config.paths.get_paths", return_value=paths): + extra = aio_mod.AioSandboxProvider._get_extra_mounts( + aio_mod.AioSandboxProvider.__new__(aio_mod.AioSandboxProvider), + "thread-1", + user_id="noob", + ) + + # extra includes thread mounts + skills mounts + docker_args: list[str] = [] + mount_entries: dict[str, str] = {} + for host, container, ro in extra: + args = _format_container_mount("docker", host, container, ro) + docker_args.extend(args) + if args[0] == "--mount": + mount_entries[container] = args[1] + + assert "--mount" in docker_args + # Skills mounts must be present + assert "/mnt/skills/public" in mount_entries + assert "dst=/mnt/skills/public" in mount_entries["/mnt/skills/public"] + assert "readonly" in mount_entries["/mnt/skills/public"] + + assert "/mnt/skills/custom" in mount_entries + assert "dst=/mnt/skills/custom" in mount_entries["/mnt/skills/custom"] + assert "users/noob/skills/custom" in mount_entries["/mnt/skills/custom"] + + # noob has no per-user custom → legacy is mounted + assert "/mnt/skills/legacy" in mount_entries + assert "dst=/mnt/skills/legacy" in mount_entries["/mnt/skills/legacy"] + + # ── Path alignment ────────────────────────────────────────────────── + + def test_skill_container_paths_match_expected_mounts(self, skills_fs): + cr = "/mnt/skills" + assert ( + Skill( + name="p", + description="", + license=None, + skill_dir=skills_fs["public"] / "pub-skill", + skill_file=skills_fs["pub_skill"], + relative_path=Path("pub-skill"), + category=SkillCategory.PUBLIC, + ).get_container_path(cr) + == "/mnt/skills/public/pub-skill" + ) + + assert ( + Skill( + name="u", + description="", + license=None, + skill_dir=skills_fs["user_custom"] / "usr-skill", + skill_file=skills_fs["user_skill"], + relative_path=Path("usr-skill"), + category=SkillCategory.CUSTOM, + ).get_container_path(cr) + == "/mnt/skills/custom/usr-skill" + ) + + assert ( + Skill( + name="l", + description="", + license=None, + skill_dir=skills_fs["legacy_global"] / "leg-skill", + skill_file=skills_fs["legacy_skill"], + relative_path=Path("leg-skill"), + category=SkillCategory.LEGACY, + ).get_container_path(cr) + == "/mnt/skills/legacy/leg-skill" + ) diff --git a/docker/docker-compose-dev.yaml b/docker/docker-compose-dev.yaml index 19ffc2e4c..69c185baa 100644 --- a/docker/docker-compose-dev.yaml +++ b/docker/docker-compose-dev.yaml @@ -52,6 +52,8 @@ services: # export DEER_FLOW_ROOT=/absolute/path/to/deer-flow - SKILLS_HOST_PATH=${DEER_FLOW_ROOT}/skills - THREADS_HOST_PATH=${DEER_FLOW_ROOT}/backend/.deer-flow/threads + # Per-user data base directory for user-scoped skill mounts + - DEER_FLOW_HOST_BASE_DIR=${DEER_FLOW_ROOT}/backend/.deer-flow # Production: use PVC instead of hostPath to avoid data loss on node failure. # When set, hostPath vars above are ignored for the corresponding volume. # USERDATA_PVC_NAME uses subPath (deer-flow/users/{user_id}/threads/{thread_id}/user-data) automatically. diff --git a/docker/docker-compose.yaml b/docker/docker-compose.yaml index 30e7880b5..32b07f7a7 100644 --- a/docker/docker-compose.yaml +++ b/docker/docker-compose.yaml @@ -158,6 +158,7 @@ services: - SANDBOX_IMAGE=enterprise-public-cn-beijing.cr.volces.com/vefaas-public/all-in-one-sandbox:latest - SKILLS_HOST_PATH=${DEER_FLOW_REPO_ROOT}/skills - THREADS_HOST_PATH=${DEER_FLOW_HOME}/threads + - DEER_FLOW_HOST_BASE_DIR=${DEER_FLOW_HOME} - KUBECONFIG_PATH=/root/.kube/config - NODE_HOST=host.docker.internal - K8S_API_SERVER=https://host.docker.internal:26443 diff --git a/docker/provisioner/app.py b/docker/provisioner/app.py index 8015ba5ea..77f559c5f 100644 --- a/docker/provisioner/app.py +++ b/docker/provisioner/app.py @@ -60,6 +60,7 @@ SANDBOX_IMAGE = os.environ.get( ) SKILLS_HOST_PATH = os.environ.get("SKILLS_HOST_PATH", "/skills") THREADS_HOST_PATH = os.environ.get("THREADS_HOST_PATH", "/.deer-flow/threads") +DEER_FLOW_HOST_BASE_DIR = os.environ.get("DEER_FLOW_HOST_BASE_DIR", "/.deer-flow") SKILLS_PVC_NAME = os.environ.get("SKILLS_PVC_NAME", "") USERDATA_PVC_NAME = os.environ.get("USERDATA_PVC_NAME", "") SANDBOX_CONTAINER_PORT_RAW = os.environ.get("SANDBOX_CONTAINER_PORT", "8080") @@ -227,6 +228,7 @@ class CreateSandboxRequest(BaseModel): sandbox_id: str thread_id: str = Field(pattern=SAFE_THREAD_ID_PATTERN) user_id: str = Field(default=DEFAULT_USER_ID, pattern=SAFE_USER_ID_PATTERN) + include_legacy_skills: bool = False class SandboxResponse(BaseModel): @@ -250,26 +252,79 @@ def _sandbox_url(node_port: int) -> str: """Build the sandbox URL using the configured NODE_HOST.""" return f"http://{NODE_HOST}:{node_port}" +def _build_volumes( + thread_id: str, + user_id: str = DEFAULT_USER_ID, + *, + include_legacy_skills: bool = False, +) -> list[k8s_client.V1Volume]: + """Build volume list: PVC when configured, otherwise hostPath. + + Skills are split into public, per-user custom, and legacy (global-custom) + volumes so that ``/mnt/skills/{public,custom,legacy}/`` paths resolve + correctly inside the sandbox — matching the hostPath layout produced by + ``LocalSandboxProvider`` and ``AioSandboxProvider``. + """ + volumes: list[k8s_client.V1Volume] = [] + + # ── Skills volumes ──────────────────────────────────────────────── -def _build_volumes(thread_id: str) -> list[k8s_client.V1Volume]: - """Build volume list: PVC when configured, otherwise hostPath.""" if SKILLS_PVC_NAME: - skills_vol = k8s_client.V1Volume( - name="skills", - persistent_volume_claim=k8s_client.V1PersistentVolumeClaimVolumeSource( - claim_name=SKILLS_PVC_NAME, - read_only=True, - ), + # PVC mode: three-way subPath not yet supported; fall back to + # single-volume mount for backward compatibility. + logger.warning( + "SKILLS_PVC_NAME is set — three-way skills layout is not " + "supported in PVC mode yet; falling back to single /mnt/skills mount" + ) + volumes.append( + k8s_client.V1Volume( + name="skills", + persistent_volume_claim=k8s_client.V1PersistentVolumeClaimVolumeSource( + claim_name=SKILLS_PVC_NAME, + read_only=True, + ), + ) ) else: - skills_vol = k8s_client.V1Volume( - name="skills", - host_path=k8s_client.V1HostPathVolumeSource( - path=SKILLS_HOST_PATH, - type="Directory", - ), + # hostPath mode: three-way layout + public_path = join_host_path(SKILLS_HOST_PATH, "public") + volumes.append( + k8s_client.V1Volume( + name="skills-public", + host_path=k8s_client.V1HostPathVolumeSource( + path=public_path, + type="Directory", + ), + ) ) + user_custom_path = join_host_path( + DEER_FLOW_HOST_BASE_DIR, "users", user_id, "skills", "custom", + ) + volumes.append( + k8s_client.V1Volume( + name="skills-custom", + host_path=k8s_client.V1HostPathVolumeSource( + path=user_custom_path, + type="DirectoryOrCreate", + ), + ) + ) + + if include_legacy_skills: + legacy_path = join_host_path(SKILLS_HOST_PATH, "custom") + volumes.append( + k8s_client.V1Volume( + name="skills-legacy", + host_path=k8s_client.V1HostPathVolumeSource( + path=legacy_path, + type="Directory", + ), + ) + ) + + # ── User-data volume ────────────────────────────────────────────── + if USERDATA_PVC_NAME: userdata_vol = k8s_client.V1Volume( name="user-data", @@ -286,13 +341,56 @@ def _build_volumes(thread_id: str) -> list[k8s_client.V1Volume]: ), ) - return [skills_vol, userdata_vol] + volumes.append(userdata_vol) + return volumes def _build_volume_mounts( - thread_id: str, user_id: str = DEFAULT_USER_ID + thread_id: str, + user_id: str = DEFAULT_USER_ID, + *, + include_legacy_skills: bool = False, ) -> list[k8s_client.V1VolumeMount]: - """Build volume mount list, using subPath for PVC user-data.""" + """Build volume mount list, mirroring three-way skills layout. + + Skills are mounted to ``/mnt/skills/{public,custom,legacy}/`` so that + category-aware ``Skill.get_container_path()`` paths resolve correctly. + PVC mode falls back to a single ``/mnt/skills`` mount. + """ + mounts: list[k8s_client.V1VolumeMount] = [] + + if SKILLS_PVC_NAME: + mounts.append( + k8s_client.V1VolumeMount( + name="skills", + mount_path="/mnt/skills", + read_only=True, + ) + ) + else: + mounts.extend( + [ + k8s_client.V1VolumeMount( + name="skills-public", + mount_path="/mnt/skills/public", + read_only=True, + ), + k8s_client.V1VolumeMount( + name="skills-custom", + mount_path="/mnt/skills/custom", + read_only=True, + ), + ] + ) + if include_legacy_skills: + mounts.append( + k8s_client.V1VolumeMount( + name="skills-legacy", + mount_path="/mnt/skills/legacy", + read_only=True, + ) + ) + userdata_mount = k8s_client.V1VolumeMount( name="user-data", mount_path="/mnt/user-data", @@ -302,19 +400,17 @@ def _build_volume_mounts( userdata_mount.sub_path = ( f"deer-flow/users/{user_id}/threads/{thread_id}/user-data" ) + mounts.append(userdata_mount) - return [ - k8s_client.V1VolumeMount( - name="skills", - mount_path="/mnt/skills", - read_only=True, - ), - userdata_mount, - ] + return mounts def _build_pod( - sandbox_id: str, thread_id: str, user_id: str = DEFAULT_USER_ID + sandbox_id: str, + thread_id: str, + user_id: str = DEFAULT_USER_ID, + *, + include_legacy_skills: bool = False, ) -> k8s_client.V1Pod: """Construct a Pod manifest for a single sandbox.""" return k8s_client.V1Pod( @@ -373,14 +469,22 @@ def _build_pod( "ephemeral-storage": "500Mi", }, ), - volume_mounts=_build_volume_mounts(thread_id, user_id=user_id), + volume_mounts=_build_volume_mounts( + thread_id, + user_id=user_id, + include_legacy_skills=include_legacy_skills, + ), security_context=k8s_client.V1SecurityContext( privileged=False, allow_privilege_escalation=True, ), ) ], - volumes=_build_volumes(thread_id), + volumes=_build_volumes( + thread_id, + user_id=user_id, + include_legacy_skills=include_legacy_skills, + ), restart_policy="Always", ), ) @@ -457,12 +561,14 @@ def create_sandbox(req: CreateSandboxRequest): sandbox_id = req.sandbox_id thread_id = req.thread_id user_id = req.user_id + include_legacy_skills = req.include_legacy_skills logger.info( - "Received request to create sandbox '%s' for thread '%s' user '%s'", + "Received request to create sandbox '%s' for thread '%s' user '%s' include_legacy_skills=%s", sandbox_id, thread_id, user_id, + include_legacy_skills, ) # ── Fast path: sandbox already exists ──────────────────────────── @@ -477,7 +583,13 @@ def create_sandbox(req: CreateSandboxRequest): # ── Create Pod ─────────────────────────────────────────────────── try: core_v1.create_namespaced_pod( - K8S_NAMESPACE, _build_pod(sandbox_id, thread_id, user_id=user_id) + K8S_NAMESPACE, + _build_pod( + sandbox_id, + thread_id, + user_id=user_id, + include_legacy_skills=include_legacy_skills, + ), ) logger.info(f"Created Pod {_pod_name(sandbox_id)}") except ApiException as exc: