fix(wrap): serialize shared proxy startup (#2946)

## Description
Serialize concurrent `headroom wrap` startup so separate agents can
safely share one local proxy.

## Type of Change
- [x] Bug fix

## Changes Made
- Added a per-port cross-process startup lock.
- Re-checks proxy health/configuration after waiting for the lock.
- Preserves reference-counted cleanup and Copilot subscription
isolation.
- Preserves `_ensure_proxy`'s introspectable keyword signature on the
locking wrapper.

## Testing
- 88 wrap/persistent/detach tests pass locally.
- Focused lock-boundary tests pass.
- Signature inspection exposes `learn` and the existing keyword-only
options.
- Ruff, format, compile, and diff checks pass.
- Remaining CI failures are unrelated existing shard or
external-download failures.

## Real Behavior Proof
Two wraps that start during proxy cold start now serialize: the second
waits, observes the first healthy listener, and reuses it instead of
spawning a competing listener.

## Review Readiness
- [x] I have performed a self-review.
- [x] This PR is ready for human review.

---------

Co-authored-by: Tejas Chopra <tejas@Tejass-MacBook-Pro.local>
Co-authored-by: Jerrett Davis <mxjerrett@gmail.com>
This commit is contained in:
Tejas Chopra
2026-08-12 12:50:09 -07:00
committed by GitHub
parent 039cd2431a
commit e540d64feb
3 changed files with 133 additions and 2 deletions
+85 -2
View File
@@ -33,6 +33,8 @@ import sys
import time
import urllib.parse
from collections.abc import Callable
from contextlib import contextmanager
from functools import wraps
from pathlib import Path
from typing import Any, cast
@@ -3460,7 +3462,7 @@ def _push_runtime_env(port: int, no_proxy: bool) -> None:
click.echo(f" Synced output settings to proxy: {', '.join(sorted(payload))}")
def _ensure_proxy(
def _ensure_proxy_unlocked(
port: int,
no_proxy: bool,
*,
@@ -3479,7 +3481,13 @@ def _ensure_proxy(
copilot_refresh_oauth_token: str | None = None,
copilot_api_token_expires_at: float | None = None,
) -> tuple[subprocess.Popen | None, int]:
"""Start or verify proxy. Returns (process_handle, actual_port)."""
"""Start or verify proxy. Returns (process_handle, actual_port).
The public ``_ensure_proxy`` wrapper serializes callers per port before
entering this function. Keeping the implementation separate makes the
lock boundary explicit and ensures every health/configuration check runs
under the same startup critical section.
"""
helpers = _live_wrap_module()
copilot_subscription_seed_requested = (
bool(copilot_api_token)
@@ -3846,6 +3854,81 @@ def _ensure_proxy(
return None, port
@contextmanager
def _proxy_start_lock(port: int) -> Any:
"""Serialize wrap proxy startup across processes sharing a port.
A proxy can spend tens of seconds loading optional ML components before it
binds its socket. Without this lock, two concurrent ``headroom wrap``
commands both see an unavailable health endpoint, choose the same port,
and race to spawn a listener. The lock is deliberately held through the
health/configuration checks and startup, then released once the proxy is
ready (or startup fails). Lock files are retained so an interrupted
process cannot create an inode-replacement race for another waiter.
"""
from headroom import paths as _paths
lock_path = _paths.proxy_start_lock_path(port)
try:
lock_path.parent.mkdir(parents=True, exist_ok=True)
lock_file = open(lock_path, "a+b") # noqa: SIM115
except OSError:
# Locking is a race-prevention enhancement, not a reason to make wrap
# unusable when a read-only/custom workspace cannot hold state. The
# existing port bind remains the final safety check in that degraded
# environment.
yield
return
with lock_file:
if sys.platform == "win32":
import msvcrt
# msvcrt.locking operates on bytes from the current file position.
lock_file.seek(0)
if lock_file.read(1) == b"":
lock_file.seek(0)
lock_file.write(b"0")
lock_file.flush()
lock_file.seek(0)
# LK_LOCK has implementation-dependent retry limits. A proxy may
# legitimately take longer than that to load ML components, so
# use the non-blocking primitive in a loop instead.
while True:
try:
msvcrt.locking(lock_file.fileno(), msvcrt.LK_NBLCK, 1)
break
except OSError:
time.sleep(0.05)
try:
yield
finally:
lock_file.seek(0)
msvcrt.locking(lock_file.fileno(), msvcrt.LK_UNLCK, 1)
else:
import fcntl
fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
try:
yield
finally:
fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
@wraps(_ensure_proxy_unlocked)
def _ensure_proxy(
port: int,
no_proxy: bool,
**kwargs: Any,
) -> tuple[subprocess.Popen | None, int]:
"""Start or reuse a proxy without racing another wrap on the same port."""
if no_proxy:
return _ensure_proxy_unlocked(port, no_proxy, **kwargs)
with _proxy_start_lock(port):
# Re-checking is part of the lock boundary: a concurrent wrapper may
# have finished startup while this caller was waiting for the lock.
return _ensure_proxy_unlocked(port, no_proxy, **kwargs)
def _client_marker_path(port: int) -> Path:
"""Path to this process's wrap-client marker for ``port``."""
from headroom import paths as _paths
+7
View File
@@ -342,6 +342,12 @@ def beacon_lock_path(port: int) -> Path:
return workspace_dir() / f".beacon_lock_{int(port)}"
def proxy_start_lock_path(port: int) -> Path:
"""Return the per-port lock used to serialize wrap proxy startup."""
return workspace_dir() / f".proxy_start_{int(port)}.lock"
# ---------------------------------------------------------------------------
# Per-resource helpers -- config bucket
# ---------------------------------------------------------------------------
@@ -432,6 +438,7 @@ __all__ = [
"proxy_clients_dir",
"deploy_root",
"beacon_lock_path",
"proxy_start_lock_path",
"models_config_path",
"plugin_config_dir",
"plugin_workspace_dir",
+41
View File
@@ -16,6 +16,7 @@ import errno
import json
import os
import signal
from contextlib import contextmanager
from pathlib import Path
from typing import Any
@@ -817,3 +818,43 @@ class TestFindAvailablePort:
)
with pytest.raises(RuntimeError, match="No available port found"):
wrap_mod._find_available_port(8787, max_attempts=3)
def test_ensure_proxy_serializes_startup_per_port(monkeypatch: pytest.MonkeyPatch) -> None:
"""A normal wrap must enter the per-port startup critical section."""
events: list[object] = []
@contextmanager
def fake_lock(port: int):
events.append(("lock-enter", port))
try:
yield
finally:
events.append(("lock-exit", port))
monkeypatch.setattr(wrap_mod, "_proxy_start_lock", fake_lock)
monkeypatch.setattr(
wrap_mod,
"_ensure_proxy_unlocked",
lambda port, no_proxy, **kwargs: events.append(("ensure", port, no_proxy)) or (None, port),
)
assert wrap_mod._ensure_proxy(8787, False) == (None, 8787)
assert events == [("lock-enter", 8787), ("ensure", 8787, False), ("lock-exit", 8787)]
def test_no_proxy_does_not_create_startup_lock(monkeypatch: pytest.MonkeyPatch) -> None:
"""Explicit --no-proxy reuses an existing service without taking the lock."""
entered = False
@contextmanager
def fail_lock(port: int):
nonlocal entered
entered = True
yield
monkeypatch.setattr(wrap_mod, "_proxy_start_lock", fail_lock)
monkeypatch.setattr(wrap_mod, "_ensure_proxy_unlocked", lambda *args, **kwargs: (None, 8787))
assert wrap_mod._ensure_proxy(8787, True) == (None, 8787)
assert entered is False