diff --git a/tests/tui_gateway/test_readiness_singleflight.py b/tests/tui_gateway/test_readiness_singleflight.py index 76d38d563e..6dee859c30 100644 --- a/tests/tui_gateway/test_readiness_singleflight.py +++ b/tests/tui_gateway/test_readiness_singleflight.py @@ -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 diff --git a/tui_gateway/methods_config.py b/tui_gateway/methods_config.py index 9605fa27a3..a7597e0354 100644 --- a/tui_gateway/methods_config.py +++ b/tui_gateway/methods_config.py @@ -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):