diff --git a/modules/skills_manager.py b/modules/skills_manager.py index af4b453a..ae839ad7 100644 --- a/modules/skills_manager.py +++ b/modules/skills_manager.py @@ -8,6 +8,7 @@ from __future__ import annotations import re import shutil import threading +import time from pathlib import Path from typing import Dict, List, Optional, Sequence @@ -32,6 +33,52 @@ def _get_sync_lock(skills_dir: Path) -> threading.Lock: return lock +def wait_skill_file_ready( + skill_file: str | Path, + skills_dir: str | Path, + *, + max_wait_seconds: float = 2.5, + poll_interval: float = 0.2, +) -> bool: + """Wait until a workspace skill file exists / 等待工作区 skill 文件就绪。 + + 工作区 skills 目录可能被并发的全量同步(rmtree + 重建)短暂清空, + 读取方在窗口期会看到「文件不存在」。这里先等当前这一轮同步结束 + (复用同步锁做有界等待),再以固定间隔轮询到预算耗尽,避免把瞬时 + 缺失误判为硬错误。 + + 返回 True 表示文件已存在;超过预算仍不存在返回 False。 + """ + target = Path(skill_file) + if target.is_file(): + return True + try: + lock: Optional[threading.Lock] = _get_sync_lock(Path(skills_dir)) + except Exception: + lock = None + deadline = time.monotonic() + max(0.0, max_wait_seconds) + while True: + if lock is not None: + try: + remaining = deadline - time.monotonic() + if remaining <= 0: + break + # 锁空闲时会立即拿到并释放,开销可忽略; + # 锁被同步占用则有界等到本轮重建完成。 + acquired = lock.acquire(timeout=remaining) + if acquired: + lock.release() + except Exception: + pass + if target.is_file(): + return True + remaining = deadline - time.monotonic() + if remaining <= 0: + break + time.sleep(min(poll_interval, remaining)) + return target.is_file() + + def ensure_agent_skills_dir(base_dir: Optional[str] = None) -> Path: """Ensure the global skills directory exists / 确保全局技能目录存在。""" root = Path(base_dir or AGENT_SKILLS_DIR).expanduser().resolve() diff --git a/server/tasks/skills.py b/server/tasks/skills.py index 5eaa9cd6..b50c227e 100644 --- a/server/tasks/skills.py +++ b/server/tasks/skills.py @@ -12,6 +12,7 @@ from server.auth_helpers import api_login_required, get_current_username from server.context import get_user_resources from server.utils_common import debug_log from config import WORKSPACE_SKILLS_DIRNAME +from modules.skills_manager import wait_skill_file_ready SKILL_FRONTMATTER_RE = re.compile(r"^---\s*\n(?P.*?)\n---\s*\n?", re.S) @@ -60,8 +61,11 @@ def _resolve_workspace_skill_path(workspace, raw_path: str) -> Path: except ValueError as exc: debug_log(f"[SkillsAPI] resolve skill path outside skills_dir: target={target!r} skills_dir={skills_dir!r}") raise ValueError("skill 路径必须位于当前工作区 .astrion/skills/ 内") from exc - if not target.is_file(): - debug_log(f"[SkillsAPI] resolve skill file not found: {target!r}") + # 工作区 skills 目录可能被并发的全量同步(rmtree+重建)短暂清空, + # 直接 is_file() 一次失败会把瞬时窗口误判为「skill 文件不存在」拒绝任务, + # 这里改为有界等待重试,等过同步窗口再下结论。 + if not wait_skill_file_ready(target, skills_dir): + debug_log(f"[SkillsAPI] resolve skill file not found (after wait): {target!r}") raise ValueError("skill 文件不存在") return target diff --git a/test/test_skills_manager.py b/test/test_skills_manager.py index 7a8d4a23..2ba8f739 100644 --- a/test/test_skills_manager.py +++ b/test/test_skills_manager.py @@ -2,6 +2,8 @@ from __future__ import annotations import os import tempfile +import threading +import time import unittest from pathlib import Path @@ -9,11 +11,13 @@ from pathlib import Path os.environ["TERMINAL_SANDBOX_MODE"] = "web" from modules.skills_manager import ( + _get_sync_lock, archive_skill_directory, get_skills_catalog, infer_private_skills_dir, sync_workspace_skills, validate_skill_directory, + wait_skill_file_ready, ) @@ -97,5 +101,73 @@ class SkillsManagerTest(unittest.TestCase): self.assertTrue((project / ".astrion" / "skills" / "private-skill" / "SKILL.md").exists()) +class WaitSkillFileReadyTest(unittest.TestCase): + """读取方等待原语:覆盖并发全量同步(rmtree+重建)的瞬时窗口。""" + + def test_existing_file_returns_true_immediately(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + skills_dir = root / ".astrion" / "skills" + target = skills_dir / "demo" / "SKILL.md" + target.parent.mkdir(parents=True) + target.write_text("x", encoding="utf-8") + start = time.monotonic() + self.assertTrue(wait_skill_file_ready(target, skills_dir, max_wait_seconds=0.5)) + self.assertLess(time.monotonic() - start, 0.2) + + def test_file_appearing_during_wait_returns_true(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + skills_dir = root / ".astrion" / "skills" + skills_dir.mkdir(parents=True) + target = skills_dir / "demo" / "SKILL.md" + + def create_later(): + time.sleep(0.3) + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text("x", encoding="utf-8") + + threading.Thread(target=create_later, daemon=True).start() + self.assertTrue( + wait_skill_file_ready(target, skills_dir, max_wait_seconds=2.0, poll_interval=0.05) + ) + + def test_in_flight_sync_lock_is_awaited(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + skills_dir = root / ".astrion" / "skills" + skills_dir.mkdir(parents=True) + target = skills_dir / "demo" / "SKILL.md" + lock = _get_sync_lock(skills_dir) + + def fake_sync(): + with lock: + time.sleep(0.3) + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text("x", encoding="utf-8") + + threading.Thread(target=fake_sync, daemon=True).start() + # 等假同步先拿到锁,模拟读取方撞上重建窗口 + time.sleep(0.05) + self.assertFalse(target.is_file()) + self.assertTrue( + wait_skill_file_ready(target, skills_dir, max_wait_seconds=2.0, poll_interval=0.05) + ) + + def test_missing_file_returns_false_within_budget(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + skills_dir = root / ".astrion" / "skills" + skills_dir.mkdir(parents=True) + target = skills_dir / "nope" / "SKILL.md" + start = time.monotonic() + self.assertFalse( + wait_skill_file_ready(target, skills_dir, max_wait_seconds=0.4, poll_interval=0.05) + ) + elapsed = time.monotonic() - start + self.assertGreaterEqual(elapsed, 0.4) + self.assertLess(elapsed, 1.5) + + if __name__ == "__main__": unittest.main()