fix(pm): isolate worker control stdin and share bootstrap cache

This commit is contained in:
ethernet
2026-09-12 18:41:04 -04:00
parent 3d0b6f9a0c
commit 062b7138e7
7 changed files with 144 additions and 17 deletions

View File

@@ -57,6 +57,8 @@ def _request(operation, arguments, *, callbacks=None, pause_event=None, project_
"lockfile": str(paths.lockfile_path())},
}
worker = Path(__file__).with_name("worker.py").resolve()
# Bootstrap precedes dispatch and must share the operation's selected cache.
cache = Path(arguments["cache"]) if arguments.get("cache") is not None else None
environment = runtime_environment()
state_sync = operation in ("sync_venv", "build_environment", "lock_project",
"ensure_environment", "ensure_python_tool", "check_project_lock",
@@ -67,7 +69,7 @@ def _request(operation, arguments, *, callbacks=None, pause_event=None, project_
# A ready PM still decides no-op/refusal under its install lock. A cold
# PM is itself a missing prerequisite, not permission to bootstrap tools.
try:
command = runtime_command(worker, bootstrap=False)
command = runtime_command(worker, bootstrap=False, cache=cache)
except InstallError as exc:
token = receipt.begin("sync")
try:
@@ -78,9 +80,9 @@ def _request(operation, arguments, *, callbacks=None, pause_event=None, project_
raise
environment["HERMES_DISABLE_LAZY_INSTALLS"] = "1"
elif operation == "venv_is_current":
command = runtime_command(worker, bootstrap=False)
command = runtime_command(worker, bootstrap=False, cache=cache)
else:
command = runtime_command(worker)
command = runtime_command(worker, cache=cache)
callback_error = None
stopped = threading.Event()
write_lock = threading.Lock()

View File

@@ -100,7 +100,8 @@ def _validate(python: Path, env: dict[str, str]) -> str:
def prepare_runtime(uv: Path, python: Path, root: Path, *, offline: bool = False,
project: Path | None = None, bootstrap: bool = True) -> Path:
project: Path | None = None, bootstrap: bool = True,
cache: Path | None = None) -> Path:
"""Publish a locked PM environment without resolving the application.
Generations are immutable after publication. Failed preparation leaves the
@@ -132,7 +133,7 @@ def prepare_runtime(uv: Path, python: Path, root: Path, *, offline: bool = False
environment = root / generation
try:
print("Preparing the isolated PM runtime…", file=sys.stderr, flush=True)
executable = stage_runtime(uv, python, environment, project=project, offline=offline)
executable = stage_runtime(uv, python, environment, project=project, offline=offline, cache=cache)
_write(environment / "pm-runtime.json", {"inputs": identity})
_write(selected, {"inputs": identity, "generation": generation.as_posix()})
return executable
@@ -142,7 +143,7 @@ def prepare_runtime(uv: Path, python: Path, root: Path, *, offline: bool = False
def runtime_python(*, bootstrap: bool = True) -> Path:
def runtime_python(*, bootstrap: bool = True, cache: Path | None = None) -> Path:
"""Resolve PM without selecting, repairing, or importing the app environment."""
if is_runtime():
return Path(sys.executable)
@@ -179,14 +180,15 @@ def runtime_python(*, bootstrap: bool = True) -> Path:
raise InstallError("pm-runtime", "pinned uv and Python are unavailable")
uv, python = tools
return prepare_runtime(uv, python, install_state_dir(project) / "pm-runtime",
bootstrap=bootstrap)
bootstrap=bootstrap, cache=cache)
def runtime_command(script: Path, args: tuple[str, ...] | list[str] = (), *, bootstrap: bool = True) -> list[str]:
def runtime_command(script: Path, args: tuple[str, ...] | list[str] = (), *,
bootstrap: bool = True, cache: Path | None = None) -> list[str]:
"""One launch contract for mutable venvs and resident signed payloads."""
resident = _resident_runtime()
if resident is None:
python = runtime_python() if bootstrap else runtime_python(bootstrap=False)
python = runtime_python(bootstrap=bootstrap, cache=cache)
return [str(python), "-I", "-B", str(script), *args]
python, site = resident
launcher = (

View File

@@ -18,12 +18,12 @@ def _members(value):
return [Path(path) for path in value["paths"]]
def _read_controls(messages, pause):
def _read_controls(messages, pause, fd):
# Raw reads avoid a daemon thread holding sys.stdin's buffered lock at exit.
pending = b""
request_id = None
try:
while block := os.read(0, 65536):
while block := os.read(fd, 65536):
pending += block
while b"\n" in pending:
line, pending = pending.split(b"\n", 1)
@@ -39,6 +39,7 @@ def _read_controls(messages, pause):
except (OSError, ValueError, KeyError) as exc:
messages.put(exc)
finally:
os.close(fd)
pause.set()
messages.put(None)
@@ -51,9 +52,15 @@ def main():
wire = os.fdopen(os.dup(sys.stdout.fileno()), "w", encoding="utf-8", buffering=1)
os.dup2(sys.stderr.fileno(), sys.stdout.fileno())
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
# A pending read on inherited control stdin can block child Python startup
# on Windows. Keep the protocol private and give every ordinary child EOF.
controls = os.dup(0)
os.set_inheritable(controls, False)
with open(os.devnull, "rb") as null:
os.dup2(null.fileno(), 0)
messages = queue.Queue()
pause = threading.Event()
threading.Thread(target=_read_controls, args=(messages, pause), daemon=True).start()
threading.Thread(target=_read_controls, args=(messages, pause, controls), daemon=True).start()
def receive():
message = messages.get()

View File

@@ -53,6 +53,53 @@ print(json.dumps({{"prefix": sys.prefix, "yaml": importlib.util.find_spec("ruame
assert checked.returncode == 0, checked.stdout + checked.stderr
def test_cold_worker_bootstrap_reuses_the_requests_cache(tmp_path, monkeypatch):
import pm
from hermes_constants import get_default_hermes_root
from pm import client, runtime
from pm.runtime_stage import stage_runtime
uv = shutil.which("uv")
assert uv, "the bootstrap cache contract requires real uv"
tools = Path(uv), Path(sys.executable)
cache = tmp_path / "shared-cache"
monkeypatch.setenv("HOME", str(tmp_path / "home"))
monkeypatch.setenv("USERPROFILE", str(tmp_path / "home"))
monkeypatch.setenv("HERMES_HOME", str(tmp_path / "home/.hermes"))
monkeypatch.setenv("HERMES_RUNTIME_DIR", str(tmp_path / "tools"))
monkeypatch.setattr("pm.paths.repo_root", lambda: tmp_path / "project")
monkeypatch.setattr("pm._uv._toolchain", lambda **kwargs: tools)
monkeypatch.setattr(client, "is_runtime", lambda: False)
stage_runtime(*tools, tmp_path / "warmup", cache=cache)
shutil.rmtree(tmp_path / "warmup")
def offline_stage(*args, **kwargs):
# A fresh manager must use the populated explicit cache, not download
# its dependencies again under the bundle's isolated HOME.
kwargs["offline"] = True
return stage_runtime(*args, **kwargs)
monkeypatch.setattr("pm.runtime_stage.stage_runtime", offline_stage)
worker = Path(client.__file__).with_name("worker.py")
script = (
"import runpy, sys; from pathlib import Path; "
f"sys.path.insert(0, {str(worker.parent.parent)!r}); import pm._uv; "
f"pm._uv._toolchain = lambda **kwargs: (Path({uv!r}), Path({sys.executable!r})); "
f"runpy.run_path({str(worker)!r}, run_name='__main__')"
)
def command(*args, **kwargs):
prepared = runtime.runtime_command(*args, **kwargs)
return [*prepared[:3], "-c", script]
monkeypatch.setattr(client, "runtime_command", command)
before = dict(os.environ)
pm.prune_cache(cache)
assert cache.is_dir()
assert not (get_default_hermes_root() / "cache/uv").exists(), "bootstrap created an unshared private cache"
assert dict(os.environ) == before
def test_sealed_worker_command_uses_only_its_recorded_site(tmp_path, monkeypatch):
from pm import paths
from pm.runtime import runtime_command

View File

@@ -69,7 +69,7 @@ def test_refused_or_already_paused_install_does_not_acquire_runtime(client, monk
import threading
from pm.downloader import DownloadPaused
monkeypatch.setattr(client, "runtime_command", lambda path: pytest.fail("refusal acquired PM runtime"))
monkeypatch.setattr(client, "runtime_command", lambda path, **kwargs: pytest.fail("refusal acquired PM runtime"))
monkeypatch.setenv("HERMES_DISABLE_LAZY_INSTALLS", "1")
with pytest.raises(InstallError, match="lazy installs are disabled"):
client.ensure("node")
@@ -125,7 +125,7 @@ def test_currency_probe_preserves_union_and_candidate_inputs(client, tmp_path, m
root_args = {"project_root": repo} if route != "worker" else {}
acquisitions = []
def ready_runtime(*, bootstrap):
def ready_runtime(*, bootstrap, cache):
assert bootstrap is False, "currency probe attempted to bootstrap PM"
acquisitions.append(bootstrap)
return isolated_python
@@ -461,7 +461,7 @@ def _patch_worker_apply(client, monkeypatch, isolated_python, body):
"Venv.apply = apply\n"
f"runpy.run_path({str(worker)!r}, run_name='__main__')\n"
)
monkeypatch.setattr(client, "runtime_command", lambda path: [str(isolated_python), "-I", "-B", "-c", script])
monkeypatch.setattr(client, "runtime_command", lambda path, **kwargs: [str(isolated_python), "-I", "-B", "-c", script])
def test_resolution_conflict_survives_worker_and_receipt(client, tmp_path, monkeypatch, isolated_python, capfd):
@@ -523,7 +523,7 @@ def test_invalid_arguments_keep_the_engine_exception_type(client):
def test_worker_death_reports_transport_failure(client, monkeypatch, isolated_python):
monkeypatch.setattr(client, "runtime_command", lambda path: [str(isolated_python), "-I", "-c", "import os; os._exit(7)"])
monkeypatch.setattr(client, "runtime_command", lambda path, **kwargs: [str(isolated_python), "-I", "-c", "import os; os._exit(7)"])
with pytest.raises(InstallError, match="worker.*result"):
client.ensure("node", explicit=True)

View File

@@ -40,7 +40,7 @@ def worker_python(tmp_path_factory):
def test_registered_package_installs_archive_in_real_worker(tmp_path, monkeypatch, worker_python, dl_server, operation):
from pm import paths
monkeypatch.setattr("pm.runtime.runtime_python", lambda: worker_python)
monkeypatch.setattr("pm.runtime.runtime_python", lambda **kwargs: worker_python)
monkeypatch.setenv("HERMES_RUNTIME_DIR", str(tmp_path / "store"))
monkeypatch.setattr(paths, "lockfile_path", lambda: tmp_path / "lock.json")
source = tmp_path / "package.py"

View File

@@ -0,0 +1,69 @@
"""Worker control traffic must never become a subprocess's standard input."""
from __future__ import annotations
import json
import os
from pathlib import Path
import subprocess
import textwrap
import pytest
from tests.pm._fixtures import isolated_python as isolated_python
@pytest.mark.parametrize("streaming", [False, True])
def test_worker_children_get_eof_while_control_pipe_stays_open(tmp_path, isolated_python, streaming):
root = Path(__file__).resolve().parents[2]
# Inject an operation, not a process mock: main owns the real protocol reader
# and the Python engine launches an actual child with inherited standard IO.
script = textwrap.dedent(f"""
import os, sys
from pathlib import Path
sys.path.insert(0, {str(root)!r})
from pm import build_operations, worker
from pm.environment import PythonEnvironment
def probe(cache, ci=False):
environment = PythonEnvironment(
uv=Path(sys.executable), python=Path(sys.executable),
destination=Path(cache) / 'unused', cache=Path(cache),
env=dict(os.environ), output=sys.stderr if {streaming!r} else None,
)
result = environment._run(
['-I', '-c', "import sys; assert sys.stdin.read() == ''; print('EOF_OK')"],
cwd=Path(cache), timeout=10,
)
if result.returncode:
raise RuntimeError(result.stderr)
assert 'EOF_OK' in result.stdout + result.stderr
return 'child completed'
build_operations.prune_cache = probe
worker.main()
""")
request = {
"id": "stdio-probe", "operation": "prune_cache",
"arguments": {"cache": str(tmp_path)}, "callbacks": [], "packages": [],
"context": {"repo": str(root), "lockfile": str(root / "pm/lock.json")},
}
env = {**os.environ, "HERMES_HOME": str(tmp_path / "home"),
"HERMES_RUNTIME_DIR": str(tmp_path / "tools")}
with (tmp_path / "diagnostics.log").open("w+") as diagnostics:
with subprocess.Popen([str(isolated_python), "-I", "-c", script],
stdin=subprocess.PIPE, stdout=subprocess.PIPE,
stderr=diagnostics, text=True, env=env) as process:
assert process.stdin is not None and process.stdout is not None
try:
process.stdin.write(json.dumps(request) + "\n")
process.stdin.flush()
# communicate() would close stdin and erase the bug's precondition.
process.wait(timeout=30)
response = json.loads(process.stdout.read())
diagnostics.seek(0)
assert response.get("result") == "child completed", (response, diagnostics.read())
assert process.returncode == 0
finally:
if process.poll() is None:
process.kill()
process.wait(timeout=5)