From 9b62dc6f98d7d1ee79a31d40b18bda683db9e7be Mon Sep 17 00:00:00 2001 From: "ark-hand[bot]" <315378070+ark-hand[bot]@users.noreply.github.com> Date: Thu, 3 Sep 2026 02:53:09 +0000 Subject: [PATCH] fix(selfHosted): align worker workspace layout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 简述 统一 self-hosted worker 的工作目录语义并对齐 CMA,同时修正内部隐藏目录命名,并隔离不同 session 的工具结果恢复账本。 ## 修改前 - Run 模式会把配置的 Workdir 当作根目录,并自动追加 session ID。 - HandleItem 模式直接使用配置的 Workdir,两种入口行为不一致。 - 工具结果账本使用 `.ma_self_host_worker/tool_ledger`。 - Workdir 扁平化后,如果账本仍直接放在共享 Workdir 下,旧 session 的 pending tool result 可能被恢复到新 session。 ## 修改后 - Run 和 HandleItem 都直接使用配置的 Workdir,不再自动创建 session ID 子目录。 - Skill 在两种模式下都安装到 `/skills/`,与 CMA 一致。 - EnvironmentWorker 的工具结果账本改为 `/.ma_self_hosted_worker/tool_ledger/`,避免跨 session 恢复和误投递。 - 非法路径字符的 session ID 会映射为稳定哈希目录,避免路径逃逸。 - 保留 `FileToolResultStore(workdir)` 的原有调用方式;EnvironmentWorker 传入 session ID 启用隔离。 - 删除不再需要的 session workdir 计算逻辑,并增加目录、安全与跨 session 隔离回归测试。 ## 兼容边界 - 多个并发 worker 需要工具文件隔离时,应由调用方分别配置独立 Workdir。 - 不自动迁移旧 `.ma_self_host_worker` 账本:旧共享账本无法可靠判定记录所属 session,自动搬迁存在误投递风险。 ## 验证 - Python 全量测试:54 passed - `ruff check src/` 通过 - `ruff format --check src/` 通过 - `test/run.sh --sdk` 真实 STG 三语言 smoke:Go、Python 通过;Java 首轮遇到模型限流,单 case 重跑通过 ## 二次 Review 修复 - 共享 Workdir 下,tool-result ledger 按 session 隔离,避免旧 session 的结果被新 session 恢复或误投递。 - worker 结束时只清理本轮成功安装的 skill 目录,保留 `skills/` 根目录及非 SDK 管理内容,避免 skill 跨 session 泄漏;该生命周期与 CMA 的 Cleanup 行为一致。 ## 补充验证 - `.venv/bin/pytest -q`(54 passed)、ruff check/format - STG 真实 SkillHub + 私有 skill 下载通过,且 Python worker 均真实执行 Python/Java 版本命令。 See merge request: !93 Sync-Source-Commit: 890523c9625536ca93332db652054be572e50a10 Hand-Written-Reason: No Ark-APIs provenance marker; treated as a hand-written source commit. Release-Version: 0.4.0 --- src/arkruntime/selfhosted/envinit.py | 18 +++++++- .../selfhosted/tool_result_store.py | 18 ++++++-- src/arkruntime/selfhosted/worker.py | 41 ++++++++----------- tests/selfhosted/test_envinit.py | 7 ++++ tests/selfhosted/test_tool_result_store.py | 19 +++++++++ tests/selfhosted/test_worker.py | 11 ++--- 6 files changed, 80 insertions(+), 34 deletions(-) diff --git a/src/arkruntime/selfhosted/envinit.py b/src/arkruntime/selfhosted/envinit.py index 303af05..e56c893 100644 --- a/src/arkruntime/selfhosted/envinit.py +++ b/src/arkruntime/selfhosted/envinit.py @@ -14,7 +14,7 @@ from contextlib import suppress from dataclasses import dataclass from pathlib import Path -from typing import Any, Optional +from typing import Any, List, Optional from .types import Session, SkillRef @@ -38,6 +38,7 @@ class Initializer: def __init__(self, api: Any, options: InitializerOptions) -> None: self.api = api self.options = options + self._installed_skill_dirs: List[Path] = [] if not self.options.skills_dir: self.options.skills_dir = str(Path(self.options.workdir) / "skills") @@ -63,6 +64,20 @@ def setup(self, session: Session) -> None: ) continue + def cleanup(self) -> None: + """Remove only the skill directories installed by this initializer.""" + errors = [] + for path in self._installed_skill_dirs: + try: + shutil.rmtree(path) + except FileNotFoundError: + pass + except OSError as exc: + errors.append(f"{path}: {exc}") + self._installed_skill_dirs.clear() + if errors: + raise OSError("remove installed skills: " + "; ".join(errors)) + def install_skill(self, session_id: str, skill: SkillRef) -> None: Path(self.options.workdir).mkdir(parents=True, exist_ok=True) Path(self.options.skills_dir).mkdir(parents=True, exist_ok=True) @@ -93,6 +108,7 @@ def install_skill(self, session_id: str, skill: SkillRef) -> None: source = _install_source_dir(Path(tmp)) target = Path(self.options.skills_dir) / name backup = _replace_skill_dir(Path(source), target) + self._installed_skill_dirs.append(target) if backup is not None: try: shutil.rmtree(backup) diff --git a/src/arkruntime/selfhosted/tool_result_store.py b/src/arkruntime/selfhosted/tool_result_store.py index 0d11b00..40a38bf 100644 --- a/src/arkruntime/selfhosted/tool_result_store.py +++ b/src/arkruntime/selfhosted/tool_result_store.py @@ -6,10 +6,11 @@ import hashlib import json import os +import re import tempfile from dataclasses import dataclass from pathlib import Path -from typing import Dict, Tuple +from typing import Dict, Optional, Tuple from .types import ContentBlock, Event, new_user_custom_tool_result_event, new_user_tool_result_event, utc_now_iso @@ -27,10 +28,14 @@ class ToolCallStoreDecision: class FileToolResultStore: """File-backed ledger that avoids re-running side-effectful tool calls.""" - def __init__(self, workdir: str) -> None: + def __init__(self, workdir: str, session_id: Optional[str] = None) -> None: if not workdir: raise ValueError("workdir must not be empty") - self.dir = Path(workdir) / ".ma_self_host_worker" / "tool_ledger" + self.dir = Path(workdir) / ".ma_self_hosted_worker" / "tool_ledger" + if session_id is not None: + if not session_id: + raise ValueError("session id must not be empty") + self.dir /= _session_ledger_name(session_id) self.dir.mkdir(parents=True, exist_ok=True, mode=0o700) def recover(self) -> Tuple[Dict[str, Event], Dict[str, bool]]: @@ -144,6 +149,13 @@ def _unknown_tool_execution_result(call_id: str, event: Event) -> Event: return new_user_tool_result_event(call_id, content, True, event.session_thread_id) +def _session_ledger_name(session_id: str) -> str: + if re.fullmatch(r"[A-Za-z0-9._-]+", session_id) and session_id not in (".", ".."): + return session_id + digest = hashlib.sha256(session_id.encode("utf-8")).hexdigest() + return f"session-{digest}" + + def _sync_directory(path: Path) -> None: try: fd = os.open(str(path), os.O_RDONLY) diff --git a/src/arkruntime/selfhosted/worker.py b/src/arkruntime/selfhosted/worker.py index a549846..c052b82 100644 --- a/src/arkruntime/selfhosted/worker.py +++ b/src/arkruntime/selfhosted/worker.py @@ -3,11 +3,9 @@ from __future__ import annotations -import hashlib import logging import os import random -import re import socket import threading import time @@ -239,7 +237,7 @@ def run(self) -> None: raise poller.error return try: - self._handle_item(_claimed_work_from_item(item), use_workdir_as_session=False) + self._handle_item(_claimed_work_from_item(item)) except (IdleTimeout, SessionTerminated): pass except Exception as exc: # noqa: BLE001 - continue polling after a bad work item. @@ -250,11 +248,11 @@ def run(self) -> None: def handle_item(self, options: HandleItemOptions) -> None: work = self._claimed_work_from_options(options) try: - self._handle_item(work, use_workdir_as_session=True) + self._handle_item(work) except (IdleTimeout, SessionTerminated): return - def _handle_item(self, work: _ClaimedWork, *, use_workdir_as_session: bool) -> None: + def _handle_item(self, work: _ClaimedWork) -> None: if not work.environment_id: work.environment_id = self.options.environment_id or os.environ.get("MA_ENVIRONMENT_ID", "") heartbeat_stop = threading.Event() @@ -262,8 +260,9 @@ def _handle_item(self, work: _ClaimedWork, *, use_workdir_as_session: bool) -> N heartbeat_done = threading.Event() heartbeat_cause = {"value": ""} heartbeat = None + initializer = None try: - workdir = self._workdir_for(work.session_id, use_workdir_as_session) + workdir = self._workdir() heartbeat_thread = threading.Thread( target=self._heartbeat_loop, args=(work, heartbeat_stop, heartbeat_done, heartbeat_cause), @@ -278,14 +277,15 @@ def _handle_item(self, work: _ClaimedWork, *, use_workdir_as_session: bool) -> N raise ValueError("session response is empty") if not session.id: session.id = work.session_id - Initializer( + initializer = Initializer( self.api, InitializerOptions(workdir=workdir, logger=self.options.logger), - ).setup(session) + ) + initializer.setup(session) if work_stop.is_set(): return tool_context = self._tool_context(workdir, work_stop) - store = FileToolResultStore(workdir) + store = FileToolResultStore(workdir, work.session_id) runner = SessionToolRunner( self.api, work.session_id, @@ -303,6 +303,11 @@ def _handle_item(self, work: _ClaimedWork, *, use_workdir_as_session: bool) -> N ) runner.run() finally: + if initializer is not None: + try: + initializer.cleanup() + except OSError as exc: + self.options.logger.warning("cleanup session skills failed: %s", exc) heartbeat_stop.set() if heartbeat is not None: heartbeat_done.wait(timeout=DEFAULT_HEARTBEAT_SECONDS + 1) @@ -400,14 +405,11 @@ def _tool_context(self, workdir: str, cancel_event: Any) -> ToolContext: cancel_event=cancel_event, ) - def _workdir_for(self, session_id: str, use_workdir_as_session: bool) -> str: + def _workdir(self) -> str: + """Return the shared worker workdir used for tool cwd and installed skills.""" root = str(Path(self.options.workdir or ".").resolve()) - if use_workdir_as_session: - Path(root).mkdir(parents=True, exist_ok=True) - return root - workdir = str(Path(root) / _session_workdir_name(session_id)) - Path(workdir).mkdir(parents=True, exist_ok=True) - return workdir + Path(root).mkdir(parents=True, exist_ok=True) + return root def _claimed_work_from_options(self, options: HandleItemOptions) -> _ClaimedWork: work_id = options.work_id or os.environ.get("MA_WORK_ID", "") @@ -464,13 +466,6 @@ def _should_stop_item(heartbeat_cause: str) -> bool: } -def _session_workdir_name(session_id: str) -> str: - if re.fullmatch(r"[A-Za-z0-9._-]+", session_id) and session_id not in (".", ".."): - return session_id - digest = hashlib.sha256(session_id.encode("utf-8")).hexdigest() - return f"session-{digest}" - - class _CombinedStopEvent: def __init__(self, *events: threading.Event) -> None: self._events = events diff --git a/tests/selfhosted/test_envinit.py b/tests/selfhosted/test_envinit.py index 18d820d..6b50ef8 100644 --- a/tests/selfhosted/test_envinit.py +++ b/tests/selfhosted/test_envinit.py @@ -92,6 +92,13 @@ def test_setup_installs_skill_under_resolved_metadata_name(tmp_path) -> None: assert (tmp_path / "skills" / "canonical-skill-name" / "SKILL.md").read_text() == "hello" assert not (tmp_path / "skills" / "skill-1").exists() + retained = tmp_path / "skills" / "retained" + retained.mkdir() + initializer.cleanup() + + assert not (tmp_path / "skills" / "canonical-skill-name").exists() + assert retained.is_dir() + def test_install_closes_skill_body_when_archive_copy_fails(tmp_path) -> None: body = _CloseTrackingBody(b"archive-too-large") diff --git a/tests/selfhosted/test_tool_result_store.py b/tests/selfhosted/test_tool_result_store.py index 5bf592d..7d5e91f 100644 --- a/tests/selfhosted/test_tool_result_store.py +++ b/tests/selfhosted/test_tool_result_store.py @@ -6,6 +6,7 @@ def test_recovery_uses_persisted_call_id(tmp_path) -> None: store = FileToolResultStore(str(tmp_path)) + assert store.dir == tmp_path / ".ma_self_hosted_worker" / "tool_ledger" store.begin("call-1", Event(id="event-1", type="agent.tool_use", name="bash")) pending, _ = store.recover() @@ -21,3 +22,21 @@ def test_recovery_removes_stale_temporary_records(tmp_path) -> None: store.recover() assert not stale.exists() + + +def test_session_store_isolates_sessions(tmp_path) -> None: + first = FileToolResultStore(str(tmp_path), "session-a") + second = FileToolResultStore(str(tmp_path), "session-b") + assert first.dir == tmp_path / ".ma_self_hosted_worker" / "tool_ledger" / "session-a" + + first.begin("call-1", Event(id="event-1", type="agent.tool_use", name="bash")) + + assert second.recover() == ({}, {}) + + +def test_session_store_sanitizes_session_id(tmp_path) -> None: + store = FileToolResultStore(str(tmp_path), "../../outside") + base = tmp_path / ".ma_self_hosted_worker" / "tool_ledger" + + assert store.dir.parent == base + assert store.dir.name.startswith("session-") diff --git a/tests/selfhosted/test_worker.py b/tests/selfhosted/test_worker.py index 81798bb..cce337e 100644 --- a/tests/selfhosted/test_worker.py +++ b/tests/selfhosted/test_worker.py @@ -185,22 +185,19 @@ def test_worker_options_preserve_legacy_positional_order() -> None: custom_tools = {"custom": object()} logger = logging.getLogger("legacy-positional-worker") - options = EnvironmentWorkerOptions( - "env-1", "worker-1", ".", False, None, None, 60, custom_tools, logger - ) + options = EnvironmentWorkerOptions("env-1", "worker-1", ".", False, None, None, 60, custom_tools, logger) assert options.custom_tools is custom_tools assert options.logger is logger assert options.tool_timeout_seconds is None -def test_session_id_cannot_escape_worker_root(tmp_path) -> None: +def test_worker_uses_configured_workdir(tmp_path) -> None: worker = EnvironmentWorker(object(), EnvironmentWorkerOptions(workdir=str(tmp_path))) - workdir = worker._workdir_for("../../outside", use_workdir_as_session=False) + workdir = worker._workdir() - assert str(tmp_path.resolve()) in workdir - assert ".." not in workdir + assert workdir == str(tmp_path.resolve()) @pytest.mark.parametrize("status_code", [408, 409, 412, 429])