fix(tui-gateway): register the readiness done-callback off the lock and give joiners a short grace
Two defects in the single-flight (#65151) hung the desktop at the readiness seam, timing out onboarding-first-chat and lineage-rotation in Desktop core E2E: - add_done_callback was registered while holding _readiness_lock, and the callback acquires the same non-reentrant lock. A probe settling in the window before registration runs the callback inline on the RPC worker, self-deadlocking it: every later setup.status / setup.runtime_check blocks on the held lock and the gateway stops answering readiness at all ('Gateway checking' forever). Register the callback outside the lock. - An overlapping poll answered the retryable error AT ONCE. The desktop fires setup.status + setup.runtime_check from independent consumers at the same seam (boot, the post-assignment setup.ready broadcast); both legs erroring read as unknown/fallback, which the onboarding gate treated as not-ready and the blocking overlay never closed. Joiners now wait a 0.5s grace for a fast (config-read) probe and read its shared result; a probe still running after the grace answers the retryable error so overlapping polls cannot starve the shared pool (#65151 invariant, kept by test_overlapping_runtime_checks...).
This commit is contained in:
@@ -192,3 +192,69 @@ def test_probe_outliving_its_budget_answers_retryable_unknown_then_reprobes(monk
|
||||
assert second["result"]["ok"] is True
|
||||
assert second["result"]["provider"] == "custom"
|
||||
assert len(calls) == 2
|
||||
|
||||
|
||||
def test_fast_probe_cannot_deadlock_the_share_lock(monkeypatch):
|
||||
"""A probe that settles before ``add_done_callback`` is registered must not
|
||||
self-deadlock: the done callback acquires the same non-reentrant lock the
|
||||
owner still held when registering it (a probe settling in that window runs
|
||||
the callback inline on the RPC worker thread, freezing every later
|
||||
readiness call and hanging the desktop at "Gateway checking")."""
|
||||
def instant_resolve(requested=None, **kwargs):
|
||||
return {"provider": "custom", "api_key": "no-key-required", "source": "config"}
|
||||
|
||||
_patch_fast_probe_env(monkeypatch, instant_resolve)
|
||||
transport = _RecordingTransport()
|
||||
|
||||
# Real dispatch through the shared pool; the probe resolves immediately, so
|
||||
# the future is routinely already done when the callback is registered.
|
||||
for index in range(20):
|
||||
_dispatch(transport, f"fast-{index}", "setup.runtime_check")
|
||||
|
||||
for index in range(20):
|
||||
response = transport.wait_for(f"fast-{index}", timeout=5)
|
||||
# An overlapping poll may legitimately answer the retryable error while
|
||||
# the probe is in flight; it must NEVER hang the worker that owns it.
|
||||
assert "result" in response or response["error"]["code"] == server._READINESS_IN_PROGRESS_ERR
|
||||
|
||||
# The lock must be free: the inflight entry was cleared, and a further
|
||||
# readiness call neither blocks nor answers the retryable error.
|
||||
deadline = time.monotonic() + 2
|
||||
while server._readiness_inflight and time.monotonic() < deadline:
|
||||
time.sleep(0.01)
|
||||
assert not server._readiness_inflight
|
||||
_dispatch(transport, "after", "setup.runtime_check")
|
||||
after = transport.wait_for("after", timeout=5)
|
||||
assert after["result"]["ok"] is True
|
||||
|
||||
|
||||
def test_overlapping_poll_joins_a_fast_probe_instead_of_erroring(monkeypatch):
|
||||
"""Two consumers poll at the same seam (boot, the post-assignment
|
||||
``setup.ready`` broadcast): the joiner must read the shared result of a
|
||||
fast probe, not answer the retryable error at once — the onboarding gate
|
||||
treats an unknown pair (setup.status + setup.runtime_check both errored)
|
||||
as not-ready and the blocking overlay never closes (the E2E hang)."""
|
||||
release = threading.Event()
|
||||
|
||||
def brief_resolve(requested=None, **kwargs):
|
||||
release.wait(timeout=10)
|
||||
return {"provider": "custom", "api_key": "no-key-required", "source": "config"}
|
||||
|
||||
_patch_fast_probe_env(monkeypatch, brief_resolve)
|
||||
transport = _RecordingTransport()
|
||||
|
||||
_dispatch(transport, "owner", "setup.runtime_check")
|
||||
|
||||
# Give the owner's probe a head start, then fire the overlapping poll.
|
||||
time.sleep(0.05)
|
||||
_dispatch(transport, "joiner", "setup.runtime_check")
|
||||
time.sleep(0.05)
|
||||
|
||||
release.set()
|
||||
|
||||
owner = transport.wait_for("owner", timeout=3)
|
||||
joiner = transport.wait_for("joiner", timeout=3)
|
||||
assert owner["result"]["ok"] is True
|
||||
# Inside the join grace the joiner reads the shared result — the same
|
||||
# authoritative answer, not the retryable-in-progress error.
|
||||
assert joiner["result"]["ok"] is True
|
||||
|
||||
@@ -27,10 +27,17 @@ _profile_scoped = _registry.profile_scoped
|
||||
# ``(kind, profile, requested provider)``:
|
||||
#
|
||||
# * the first caller submits the probe and waits a bounded budget;
|
||||
# * an overlapping poll for a still-running probe answers a retryable error
|
||||
# immediately — a JSON-RPC error, never a fabricated ``ok`` (the result
|
||||
# contract requires the real shape, and the desktop already treats an errored
|
||||
# runtime_check as unknown, keeping setup.status authoritative);
|
||||
# * an overlapping poll for a still-running probe waits a short join grace for
|
||||
# it — the Desktop fires setup.status + setup.runtime_check from independent
|
||||
# consumers at the same seam (boot, the post-assignment ``setup.ready``
|
||||
# broadcast), and answering the retryable error AT ONCE made both legs of
|
||||
# one consumer transiently unknown, which the onboarding gate read as
|
||||
# not-ready and the overlay never closed. Fast (config-read) probes settle
|
||||
# well inside the grace, so an overlapping poll reads the shared result; a
|
||||
# probe still running after it answers a retryable error instead —
|
||||
# a JSON-RPC error, never a fabricated ``ok`` (the result contract requires
|
||||
# the real shape, and the desktop already treats an errored runtime_check as
|
||||
# unknown, keeping setup.status authoritative);
|
||||
# * a probe that outlives the budget answers the same retryable error while it
|
||||
# keeps running in the background; the in-flight entry is cleared when it
|
||||
# settles, so the next poll starts a fresh probe and never reads a stale one.
|
||||
@@ -41,6 +48,14 @@ atexit.register(lambda: _readiness_pool.shutdown(wait=False, cancel_futures=True
|
||||
_readiness_lock = threading.Lock()
|
||||
_readiness_inflight: dict[tuple, concurrent.futures.Future] = {}
|
||||
_READINESS_SHARE_WAIT_SECONDS = 4.0
|
||||
# Join grace for an overlapping poll on a still-running probe: fast (config-read)
|
||||
# probes settle in milliseconds, so an overlapping consumer at the same seam
|
||||
# (statusbar + onboarding, both polling setup.status / setup.runtime_check) reads
|
||||
# the shared result instead of an unknown-readiness error pair — which the
|
||||
# onboarding gate read as not-ready and the blocking overlay never closed.
|
||||
# Kept well under a second: a joiner occupies its shared RPC worker for at most
|
||||
# the grace, so overlapping polls still cannot starve the pool (#65151).
|
||||
_READINESS_JOIN_GRACE_SECONDS = 0.5
|
||||
# setup.status's probe legitimately blocks on the boot bootstrap's record
|
||||
# (free_tier_bootstrap.SETUP_READY_WAIT_SECONDS = 8s): its join budget must
|
||||
# cover that wait or every boot poll would answer the retryable error.
|
||||
@@ -301,28 +316,36 @@ def _readiness_cleared(key):
|
||||
def _readiness_share(rid, key, run_probe, wait_seconds):
|
||||
"""Run ``run_probe`` single-flighted under ``key`` on the dedicated readiness pool.
|
||||
|
||||
The first caller submits the probe and waits up to ``wait_seconds``; a caller that
|
||||
finds a still-running probe answers the retryable error immediately (its shared RPC
|
||||
worker is freed at once — the probe keeps running for the first caller), and one that
|
||||
finds it settled reads the shared result. A probe that outlives the budget answers
|
||||
the same retryable error while it continues in the background."""
|
||||
The first caller submits the probe and waits up to ``wait_seconds``; a caller that finds a
|
||||
still-running probe waits the short join grace for it first — the Desktop fires
|
||||
setup.status + setup.runtime_check from independent consumers at the same seam (boot, the
|
||||
post-assignment ``setup.ready`` broadcast), and answering the retryable error at once made
|
||||
both legs of one consumer transiently unknown, which its onboarding gate read as not-ready.
|
||||
Fast (config-read) probes settle well inside the grace, so an overlapping poll reads the
|
||||
shared result; one that outlives the grace answers a retryable error, freeing its shared RPC
|
||||
worker (the probe keeps running for the first caller). A probe that outlives ``wait_seconds``
|
||||
answers the same retryable error while it continues in the background."""
|
||||
with _readiness_lock:
|
||||
future = _readiness_inflight.get(key)
|
||||
owner = future is None
|
||||
if owner:
|
||||
future = _readiness_pool.submit(run_probe)
|
||||
_readiness_inflight[key] = future
|
||||
future.add_done_callback(_readiness_cleared(key))
|
||||
if not owner and not future.done():
|
||||
return _err(rid, _READINESS_IN_PROGRESS_ERR,
|
||||
"readiness check still in progress; retrying next tick")
|
||||
if owner:
|
||||
# Registered OUTSIDE the lock: a probe that already settled runs the
|
||||
# callback inline on this thread, and the callback acquires the same
|
||||
# non-reentrant lock — under the lock that self-deadlocks the RPC
|
||||
# worker and every later readiness call blocks on it (the gateway
|
||||
# hangs answering setup.status / setup.runtime_check at all).
|
||||
future.add_done_callback(_readiness_cleared(key))
|
||||
try:
|
||||
return _ok(rid, future.result(timeout=wait_seconds))
|
||||
return _ok(rid, future.result(timeout=wait_seconds if owner else
|
||||
min(wait_seconds, _READINESS_JOIN_GRACE_SECONDS)))
|
||||
except concurrent.futures.TimeoutError:
|
||||
logger.warning("readiness probe %s exceeded %.1fs; it continues in the background",
|
||||
logger.warning("readiness probe %s exceeded its budget (%.1fs); it continues in the background",
|
||||
key, wait_seconds)
|
||||
return _err(rid, _READINESS_IN_PROGRESS_ERR,
|
||||
f"readiness check timed out after {wait_seconds:.0f}s; retrying next tick")
|
||||
"readiness check still in progress; retrying next tick")
|
||||
|
||||
|
||||
def _readiness_check(rid, params, probe, *, probe_key, wait_seconds):
|
||||
|
||||
Reference in New Issue
Block a user