fix: bound skill update wait budget and lingering fetch workers
This commit is contained in:
@@ -136,7 +136,7 @@ def test_do_list_platform_env_is_ignored(three_source_env, monkeypatch):
|
||||
|
||||
|
||||
|
||||
def test_check_for_skill_updates_does_not_fall_back_across_registries():
|
||||
def test_check_for_skill_updates_does_not_fall_back_across_registries(tmp_path, monkeypatch):
|
||||
"""An entry whose source has no adapter reports `unavailable`.
|
||||
|
||||
Previously `candidate_sources ... or sources` fell back to every source, so
|
||||
@@ -168,9 +168,12 @@ def test_check_for_skill_updates_does_not_fall_back_across_registries():
|
||||
def inspect(self, identifier):
|
||||
return _ForeignBundle()
|
||||
|
||||
from tools import skills_hub as hub
|
||||
monkeypatch.setattr(hub, "SKILLS_DIR", tmp_path)
|
||||
(tmp_path / "reddit").mkdir()
|
||||
lock = _DummyLockFile([
|
||||
{"name": "reddit", "identifier": "reddit", "source": "clawhub",
|
||||
"content_hash": "hash-of-the-clawhub-copy"},
|
||||
"install_path": "reddit", "content_hash": "hash-of-the-clawhub-copy"},
|
||||
])
|
||||
|
||||
results = check_for_skill_updates(
|
||||
|
||||
112
tests/tools/test_skills_update_budget.py
Normal file
112
tests/tools/test_skills_update_budget.py
Normal file
@@ -0,0 +1,112 @@
|
||||
"""Real HTTP controls for invalid installs and bounded update-check work."""
|
||||
import threading
|
||||
import time
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
from urllib.request import urlopen
|
||||
|
||||
import pytest
|
||||
|
||||
from tools import skills_hub_install as install
|
||||
from tools.skills_hub import HubLockFile
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def upstream():
|
||||
calls = []
|
||||
release = threading.Event()
|
||||
finished = threading.Event()
|
||||
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
def log_message(self, format, *args):
|
||||
pass
|
||||
|
||||
def do_GET(self):
|
||||
calls.append(self.path)
|
||||
if self.path == "/blocked":
|
||||
release.wait(10)
|
||||
elif self.path == "/paced":
|
||||
time.sleep(0.1)
|
||||
self.send_response(200)
|
||||
self.end_headers()
|
||||
self.wfile.write(b"ok")
|
||||
|
||||
server = ThreadingHTTPServer(("127.0.0.1", 0), Handler)
|
||||
thread = threading.Thread(target=server.serve_forever, daemon=True)
|
||||
thread.start()
|
||||
|
||||
class Source:
|
||||
def source_id(self):
|
||||
return "fixture"
|
||||
|
||||
def fetch(self, identifier):
|
||||
with urlopen(f"http://127.0.0.1:{server.server_port}/{identifier}", timeout=15) as response:
|
||||
response.read()
|
||||
finished.set()
|
||||
return None
|
||||
|
||||
try:
|
||||
yield Source(), calls, release, finished
|
||||
finally:
|
||||
release.set()
|
||||
server.shutdown()
|
||||
server.server_close()
|
||||
thread.join()
|
||||
|
||||
|
||||
def record(lock, name, path, identifier="ok"):
|
||||
lock.record_install(name, "fixture", identifier, "community", "safe", "old", path, ["SKILL.md"])
|
||||
|
||||
|
||||
def test_invalid_install_never_contacts_upstream(tmp_path, monkeypatch, upstream):
|
||||
from tools import skills_hub as hub
|
||||
root = tmp_path / "skills"
|
||||
root.mkdir()
|
||||
monkeypatch.setattr(hub, "SKILLS_DIR", root)
|
||||
(root / "healthy").mkdir()
|
||||
lock = HubLockFile()
|
||||
record(lock, "invalid", "invalid", "invalid")
|
||||
record(lock, "healthy", "healthy")
|
||||
data = lock.load()
|
||||
data["installed"]["invalid"]["install_path"] = "../outside"
|
||||
lock.save(data)
|
||||
source, calls, _, _ = upstream
|
||||
rows = install.check_for_skill_updates(lock=lock, sources=[source])
|
||||
assert [row["status"] for row in rows] == ["invalid_install", "unavailable"]
|
||||
assert calls == ["/ok"]
|
||||
|
||||
|
||||
def test_update_budget_bounds_repeated_workers_and_whole_check(tmp_path, monkeypatch, upstream):
|
||||
from tools import skills_hub as hub
|
||||
root = tmp_path / "skills"
|
||||
root.mkdir()
|
||||
monkeypatch.setattr(hub, "SKILLS_DIR", root)
|
||||
monkeypatch.setattr(install, "_FETCH_TIMEOUT_SECONDS", 0.2)
|
||||
source, calls, release, finished = upstream
|
||||
lock = HubLockFile()
|
||||
for index in range(30):
|
||||
name = f"skill-{index}"
|
||||
(root / name).mkdir()
|
||||
record(lock, name, name, "blocked")
|
||||
try:
|
||||
started = time.monotonic()
|
||||
for _ in range(3):
|
||||
rows = install.check_for_skill_updates(lock=lock, sources=[source])
|
||||
assert all(row["status"] == "unavailable" for row in rows)
|
||||
assert time.monotonic() - started < 2
|
||||
assert calls == ["/blocked"]
|
||||
finally:
|
||||
release.set()
|
||||
assert finished.wait(5)
|
||||
# Wait for fetch's finally to release capacity, without assuming scheduler order.
|
||||
deadline = time.monotonic() + 5
|
||||
while any(t.name == "skills-update-fetch" for t in threading.enumerate()):
|
||||
assert time.monotonic() < deadline
|
||||
time.sleep(0.01)
|
||||
calls.clear()
|
||||
for index in range(30):
|
||||
name = f"skill-{index}"
|
||||
record(lock, name, name, "paced")
|
||||
started = time.monotonic()
|
||||
install.check_for_skill_updates(lock=lock, sources=[source])
|
||||
assert time.monotonic() - started < 2
|
||||
assert 1 <= len(calls) <= 2
|
||||
@@ -247,18 +247,23 @@ def bundle_content_hash(bundle: SkillBundle) -> str:
|
||||
_SOURCE_ID_ALIASES = {"skills.sh": "skills-sh"}
|
||||
|
||||
_FETCH_TIMEOUT_SECONDS = 30.0
|
||||
# Keep capacity occupied until fetch really exits, including after caller timeout.
|
||||
# Repeated/concurrent checks must not accumulate abandoned credential contexts.
|
||||
_FETCH_SLOT = threading.BoundedSemaphore(1)
|
||||
|
||||
|
||||
def _fetch_bundle_bounded(
|
||||
src: SkillSource, identifier: str, timeout: float = _FETCH_TIMEOUT_SECONDS,
|
||||
) -> Optional[SkillBundle]:
|
||||
"""Fetch one bundle with a hard wall-clock bound.
|
||||
"""Bound caller waiting, not execution of the synchronous source adapter.
|
||||
|
||||
A dead upstream can hang on network IO far past any useful patience, and
|
||||
every installed skill pays that cost on each update run (#104291). The
|
||||
helper thread is a daemon so an abandoned fetch cannot block CLI exit.
|
||||
One process-wide slot bounds lingering workers. While it is occupied, new
|
||||
checks fail fast rather than enqueue work or retain more request contexts.
|
||||
Python cannot cancel a running thread; the daemon may finish in background.
|
||||
"""
|
||||
deadline = time.monotonic() + timeout
|
||||
if timeout <= 0 or not _FETCH_SLOT.acquire(blocking=False):
|
||||
return None
|
||||
box: Dict[str, Tuple[float, Optional[SkillBundle]]] = {}
|
||||
|
||||
def _run() -> None:
|
||||
@@ -267,9 +272,16 @@ def _fetch_bundle_bounded(
|
||||
box["result"] = (time.monotonic(), bundle)
|
||||
except Exception:
|
||||
logger.debug("Skill update fetch failed for %s", identifier, exc_info=True)
|
||||
finally:
|
||||
_FETCH_SLOT.release()
|
||||
|
||||
worker = threading.Thread(target=copy_context().run, args=(_run,), daemon=True)
|
||||
worker.start()
|
||||
try:
|
||||
worker = threading.Thread(target=copy_context().run, args=(_run,),
|
||||
name="skills-update-fetch", daemon=True)
|
||||
worker.start()
|
||||
except Exception:
|
||||
_FETCH_SLOT.release()
|
||||
raise
|
||||
worker.join(max(0.0, deadline - time.monotonic()))
|
||||
finished, bundle = box.get("result", (float("inf"), None))
|
||||
return bundle if finished <= deadline else None
|
||||
@@ -299,6 +311,9 @@ def check_for_skill_updates(
|
||||
if sources is None:
|
||||
sources = create_source_router(auth=auth)
|
||||
|
||||
# A shared remote-wait budget, not N independent 30-second waits. Local
|
||||
# lock/path/hash work is not cancellable and is outside this wait guarantee.
|
||||
deadline = time.monotonic() + _FETCH_TIMEOUT_SECONDS
|
||||
results: List[dict] = []
|
||||
for entry in installed:
|
||||
identifier, source_name = entry.get("identifier", ""), entry.get("source", "")
|
||||
@@ -307,8 +322,9 @@ def check_for_skill_updates(
|
||||
install_dir = _resolve_lock_install_path(
|
||||
entry.get("install_path", ""), entry.get("name", "skill"))
|
||||
orphaned = not install_dir.is_dir()
|
||||
except ValueError:
|
||||
orphaned = False # unresolvable entries keep the pre-existing fetch behavior
|
||||
except (ValueError, OSError, RuntimeError):
|
||||
results.append({**row, "status": "invalid_install"})
|
||||
continue
|
||||
if orphaned:
|
||||
# The lock-file entry points at a directory that is gone: the fetched
|
||||
# bundle could never be applied, so skip the network cost entirely
|
||||
@@ -317,7 +333,7 @@ def check_for_skill_updates(
|
||||
continue
|
||||
bundle = None
|
||||
for src in filter(lambda s: _source_matches(s, source_name), sources):
|
||||
bundle = _fetch_bundle_bounded(src, identifier)
|
||||
bundle = _fetch_bundle_bounded(src, identifier, timeout=deadline - time.monotonic())
|
||||
if bundle:
|
||||
break
|
||||
if not bundle:
|
||||
|
||||
@@ -881,6 +881,10 @@ hermes skills update react --force # Overwrite a skill you've edited locally
|
||||
|
||||
This uses the stored source identifier plus the current upstream bundle content hash to detect drift.
|
||||
|
||||
Checks skip network requests for missing or non-directory installs (`orphaned`) and unsafe or unresolvable recorded paths (`invalid_install`). Missing-directory entries can be removed with `hermes skills uninstall <name>`; invalid paths require inspecting and repairing the active profile's `skills/.hub/lock.json` before retrying. No entries are removed automatically.
|
||||
|
||||
Remote fetch waiting shares a 30-second budget across the check, rather than spending 30 seconds per skill. Entries not fetched within that budget report `unavailable`, including healthy entries later in the list; checking a specific name avoids waiting behind earlier entries. This is a caller-wait limit, not cancellation or a strict whole-command deadline: local filesystem work is outside the guarantee. At most one update fetch runs in a process. A timed-out synchronous adapter may continue in its daemon thread with its request context; subsequent checks in that process report `unavailable` while it remains active, rather than accumulating more workers. Capacity returns when the adapter exits; a permanently stuck adapter requires restarting that process.
|
||||
|
||||
Skills you have edited locally (the on-disk content no longer matches the hash recorded at install time) are **skipped** by `hermes skills update` so your changes are never silently overwritten. Pass `--force` to replace them with the upstream version anyway.
|
||||
|
||||
:::tip GitHub rate limits
|
||||
|
||||
Reference in New Issue
Block a user