diff --git a/gateway/relay/adapter.py b/gateway/relay/adapter.py index c680236deb..2e65a7c4a3 100644 --- a/gateway/relay/adapter.py +++ b/gateway/relay/adapter.py @@ -1150,6 +1150,31 @@ class RelayAdapter(BasePlatformAdapter): logger.debug("relay go_dormant failed", exc_info=True) return False + def hold_redial(self) -> bool: + """Park the transport's reconnect supervisor. False if it did not take. + + The caller suspends on the strength of this, so a swallowed failure here + would report protection that was never installed. + """ + return self._toggle_transport_redial("hold_redial") + + def release_redial(self) -> bool: + """Undo hold_redial() when the suspend did not land.""" + return self._toggle_transport_redial("release_redial") + + def _toggle_transport_redial(self, method_name: str) -> bool: + # Still never raises into the idle path, but a stub transport or a + # throwing toggle now reports False instead of passing for success. + method = getattr(self._transport, method_name, None) + if not callable(method): + return False + try: + method() + except Exception: # noqa: BLE001 - never blocks the idle path + logger.debug("relay %s failed", method_name, exc_info=True) + return False + return True + async def send_for_platform( self, logical_platform: Any, diff --git a/gateway/relay/ws_transport.py b/gateway/relay/ws_transport.py index fbf35539d7..55daa52e52 100644 --- a/gateway/relay/ws_transport.py +++ b/gateway/relay/ws_transport.py @@ -294,6 +294,12 @@ async def _await_bounded(aw: Awaitable[Any]) -> None: pass +# Ceiling on the brokered-suspend redial hold. Must outlast the client's own +# broker deadline (scale_to_zero.BROKERED_SUSPEND_TIMEOUT_S) or the supervisor +# reconnects while the stop is still in flight. +REDIAL_HOLD_MAX_S = 60.0 + + class WebSocketRelayTransport: """RelayTransport over a WebSocket connection the gateway dials to the connector.""" @@ -345,6 +351,11 @@ class WebSocketRelayTransport: # timer only advances once awake; it just needs to re-dial promptly then. self._dormant = False self._dormant_redial_s = 1.0 + # Set while a NAS-brokered suspend is in flight. See _await_redial_hold. + self._redial_held = False + self._redial_release = asyncio.Event() + # Ceiling, so a suspend that never lands cannot strand us offline. + self._redial_hold_max_s = REDIAL_HOLD_MAX_S self._ws: Any = None self._reader: Optional[asyncio.Task[None]] = None @@ -550,12 +561,17 @@ class WebSocketRelayTransport: backlog); an unexpected close re-dials immediately (the platform proxy never sees load drop, never suspends). Here the reader's fall-through still arms the supervisor on the dormant cadence; on resume the re-dial makes the - connector drain the buffered backlog. Returns the go_idle ack result; the - close happens regardless. No-op (False) when never connected. + connector drain the buffered backlog. Returns the go_idle ack result; on a + MISSED ack it returns WITHOUT closing — the caller refuses to suspend without + one, so closing would only cost a needless reconnect. No-op (False) when + never connected. """ if self._ws is None: return False acked = await self.go_idle(timeout_s=timeout_s) + if not acked: + # Nothing will suspend us, so stay connected and keep serving. + return False # Mark dormant BEFORE closing so the supervisor takes the dormant cadence. self._dormant = True try: @@ -697,6 +713,9 @@ class WebSocketRelayTransport: backoff = self._dormant_redial_s if self._dormant else self._reconnect_backoff_s while not self._closing: await asyncio.sleep(backoff) + if self._closing: + return + await self._await_redial_hold() if self._closing: return try: @@ -707,6 +726,31 @@ class WebSocketRelayTransport: logger.warning("relay ws reconnect failed: %s", exc) backoff = min(backoff * 2, self._reconnect_max_backoff_s) + def hold_redial(self) -> None: + """Park the reconnect supervisor until release_redial() or the hold cap.""" + self._redial_release.clear() + self._redial_held = True + + def release_redial(self) -> None: + """Let the supervisor re-dial again (a brokered suspend that failed).""" + self._redial_held = False + self._redial_release.set() + + async def _await_redial_hold(self) -> None: + """Block a pending re-dial while a brokered suspend is in flight: it would + clear the dormant flip. Bounded, so a lost suspend still reconnects.""" + if not self._redial_held: + return + try: + await asyncio.wait_for( + self._redial_release.wait(), timeout=self._redial_hold_max_s + ) + except asyncio.TimeoutError: + logger.info("relay: brokered suspend did not land, reconnecting") + finally: + self._redial_held = False + self._redial_release.clear() + # ── inbound frame dispatch ─────────────────────────────────────────── async def _handle_frame(self, line: str) -> None: try: diff --git a/gateway/run_shutdown.py b/gateway/run_shutdown.py index 6a16d4674a..8f0281eabe 100644 --- a/gateway/run_shutdown.py +++ b/gateway/run_shutdown.py @@ -401,7 +401,9 @@ class GatewayShutdownMixin: """Watch for idle, drive the relay dormant, then self-suspend. On sustained idle: status `draining` (NOT _running=False), relay go_dormant() (socket close, NOT disconnect()), no mark_resume_pending (suspend preserves RAM), THEN suspend via the flaps socket — Fly autostop - sees only INBOUND connections and would freeze mid-job. Off-Fly: no quiesce at all.""" + sees only INBOUND connections and would freeze mid-job. Without a flaps socket NAS brokers + the stop through the stamped GATEWAY_RELAY_SLEEP_URL; with no lever at all the watcher + abstains.""" await asyncio.sleep(min(interval, 30.0)) # let startup settle while self._running: try: @@ -413,15 +415,16 @@ class GatewayShutdownMixin: go_dormant = getattr(self._relay_adapter_for_dormancy(), "go_dormant", None) if not callable(go_dormant): continue - # Quiesce only when a suspend can follow: off-Fly the platform owns the freeze and a - # go_dormant() socket close would just re-dial (~1.4s) unflipped, dropping inbound. - from gateway.scale_to_zero import self_suspend_available - if not self_suspend_available(): + # Quiesce only when a suspend can follow: otherwise the re-dial after the socket + # close just clears the flip again. + from gateway.scale_to_zero import suspend_available + if not suspend_available(): if not self._scale_to_zero_no_suspend_logged: self._scale_to_zero_no_suspend_logged = True logger.info( - "scale-to-zero: idle, but this platform suspends on its own timer (no " - "in-machine suspend API); staying connected rather than quiescing" + "scale-to-zero: idle, but this platform offers no suspend lever (no " + "in-machine API and no brokered sleep URL); staying connected rather " + "than quiescing" ) continue logger.info( @@ -430,23 +433,44 @@ class GatewayShutdownMixin: self._scale_to_zero_idle_timeout_seconds(), ) self._scale_to_zero_status("draining", "scale-to-zero: status mark failed") + # Both levers: the 1s dormant re-dial can beat either suspend and clear the flip. + # Held BEFORE go_dormant, whose close arms it. + if not self._scale_to_zero_hold_redial(True): + # Without the hold the re-dial can clear the flip before the stop lands, so + # refuse rather than suspend unprotected. + logger.warning( + "scale-to-zero: could not hold the relay re-dial — staying awake rather " + "than suspending unprotected" + ) + self._scale_to_zero_abandon_suspend() + continue dormant_ok = True try: result = go_dormant() if asyncio.iscoroutine(result): - await result + result = await result + # The going_idle ack. Without it inbound is NOT buffered, so suspending would + # freeze a live destination: the whole bug. + if result is not True: + dormant_ok = False + logger.warning( + "scale-to-zero: connector did not ack going_idle — staying awake " + "rather than freezing a live destination" + ) except Exception: # noqa: BLE001 - dormancy is best-effort dormant_ok = False logger.debug("scale-to-zero: go_dormant failed", exc_info=True) # After a wake the drained inbound updates _last_inbound_at; give it a window so we # don't immediately re-go-dormant on the same idle reading before traffic lands. self._scale_to_zero_cooldown_until = time.time() + max(interval, 60.0) - # Self-suspend ONLY after a clean quiesce (else inbound black-holes while we sleep), and + # Suspend ONLY after an ACKED quiesce (else inbound black-holes while we sleep), and # re-check idle — inbound may have landed during the quiesce await. if not dormant_ok: + self._scale_to_zero_abandon_suspend() continue if not self._scale_to_zero_is_idle(): logger.info("scale-to-zero: inbound arrived during quiesce — skipping suspend") + self._scale_to_zero_abandon_suspend() continue await self._scale_to_zero_self_suspend() except asyncio.CancelledError: @@ -455,21 +479,102 @@ class GatewayShutdownMixin: logger.debug("scale-to-zero watcher iteration error", exc_info=True) async def _scale_to_zero_self_suspend(self) -> None: - """Suspend this Fly machine via the local flaps socket (fail-awake; off-Fly a silent no-op).""" - from gateway.scale_to_zero import self_suspend_available, suspend_self + """Suspend this machine, in-guest where possible and via NAS otherwise (fail-awake). + + Called ONLY after a clean, acked go_dormant(), with the re-dial already held. + """ + from gateway.scale_to_zero import ( + brokered_sleep_url, request_brokered_suspend, self_suspend_available, suspend_self + ) try: - if not self_suspend_available(): - logger.debug( - "scale-to-zero: flaps socket / machine identity absent — dormant without platform suspend" - ) - return - if not await asyncio.to_thread(suspend_self): + if self_suspend_available(): + accepted = await asyncio.to_thread(suspend_self) + lever = "self-suspend" + if accepted: + # flaps answers seconds BEFORE the kernel freezes, so the fence has to span + # that gap. + await self._scale_to_zero_await_freeze_gap() + self._scale_to_zero_hold_redial(False) + else: + self._scale_to_zero_abandon_suspend() + else: + # No in-guest API (Azure ACA): NAS holds the credential for the stop verb and + # brokers it for us. + url = brokered_sleep_url() + if not url: + # The watcher held on our behalf; nothing is coming to freeze the machine, so + # a held supervisor would just stay offline. + self._scale_to_zero_abandon_suspend() + logger.debug( + "scale-to-zero: no suspend lever available — dormant without platform suspend" + ) + return + # The watcher already holds the supervisor across this call. + accepted = await asyncio.to_thread(request_brokered_suspend, url) + lever = "brokered suspend" + if not accepted: + self._scale_to_zero_abandon_suspend() + if not accepted: logger.warning( - "scale-to-zero: self-suspend not accepted — machine stays " - "awake (fail-awake); will retry on the next idle window" + "scale-to-zero: %s not accepted — machine stays awake (fail-awake); will " + "retry on the next idle window", lever, ) except Exception: # noqa: BLE001 - suspend is best-effort, never crash logger.debug("scale-to-zero: self-suspend failed", exc_info=True) + self._scale_to_zero_abandon_suspend() + + async def _scale_to_zero_await_freeze_gap(self) -> None: + """Hold the re-dial fence across the flaps-2xx -> kernel-freeze gap. + + Sliced on the WALL clock rather than one ``asyncio.sleep`` because a Fly suspend stops + CLOCK_MONOTONIC while CLOCK_REALTIME keeps tracking host time. Measured on a Fly machine + (gru, 2026-09-03) across a 252.219s freeze: ``time.monotonic()`` advanced 0.501s, + ``time.time()`` advanced 252.219s. ``asyncio.sleep`` runs on ``loop.time()`` (monotonic), + so a single sleep would resume with its REMAINDER after the wake and delay the drain + re-dial by exactly that much, on every wake of every Fly agent. + + On the wall clock the deadline is already past by the time we resume, so the fence costs + nothing after a freeze while still spanning the full gap before one. That decoupling is + what lets FLY_FREEZE_GRACE_S be sized for the slowest (largest-RAM) machine. + """ + from gateway.scale_to_zero import FLY_FREEZE_GRACE_S, FLY_FREEZE_GRACE_TICK_S + deadline = time.time() + FLY_FREEZE_GRACE_S + while time.time() < deadline: + await asyncio.sleep(FLY_FREEZE_GRACE_TICK_S) + + def _scale_to_zero_abandon_suspend(self) -> None: + """Undo a quiesce we are not going to follow with a suspend. + + All three together: a released supervisor still advertising `draining` reads as + mid-shutdown until the next real inbound event, and an abort that skips the cooldown + re-runs on every tick. + """ + self._scale_to_zero_hold_redial(False) + # Same guard as _exit_external_drain: a real shutdown drain must win, so never resurrect + # a stopping gateway to `running`. + if not getattr(self, "_draining", False) and self._running: + self._scale_to_zero_status("running", "scale-to-zero: status restore failed") + # An abort before the cooldown is set would otherwise retry every tick. + self._scale_to_zero_cooldown_until = max( + self._scale_to_zero_cooldown_until, time.time() + 60.0 + ) + + def _scale_to_zero_hold_redial(self, held: bool) -> bool: + """Hold or release the relay's reconnect supervisor. Returns whether the transport actually + took it, so the caller can refuse to suspend without the protection rather than fail open.""" + try: + adapter = self._relay_adapter_for_dormancy() + if adapter is None: + return False + method = getattr(adapter, "hold_redial" if held else "release_redial", None) + if not callable(method): + return False + # Trust the adapter's answer rather than the absence of an exception: it deliberately + # never raises, so "did not throw" proves nothing. + return method() is True + except Exception: # noqa: BLE001 - never blocks the suspend it precedes + logger.debug("scale-to-zero: redial hold toggle failed", exc_info=True) + return False # External drain control: the dashboard writes/removes ``.drain_request.json`` (gateway/drain_control.py); # the watcher flips between accepting and refusing NEW turns WITHOUT exiting (reversible). diff --git a/gateway/scale_to_zero.py b/gateway/scale_to_zero.py index a668d262a9..ed8d001efd 100644 --- a/gateway/scale_to_zero.py +++ b/gateway/scale_to_zero.py @@ -15,6 +15,7 @@ import logging import os import socket import time +import urllib.parse from pathlib import Path from typing import Any, Iterable, Optional @@ -25,6 +26,13 @@ FLY_APP_NAME_ENV = "FLY_APP_NAME" # Fly-injected identity; both needed for self FLY_MACHINE_ID_ENV = "FLY_MACHINE_ID" # Local flaps (Fly Machines API) socket; POST .../suspend freezes THIS machine. FLY_API_SOCKET = "/.fly/api" + +# NAS-brokered suspend, stamped where the guest cannot suspend itself. Carries +# its own signed credential in the query string, like GATEWAY_RELAY_WAKE_URL. +SLEEP_URL_ENV = "GATEWAY_RELAY_SLEEP_URL" + +_malformed_sleep_url_logged = False + # Short is safe: real work always blocks the suspend, resume is sub-second; longer bills idle RAM. DEFAULT_IDLE_TIMEOUT_MINUTES = 2 _TRUTHY = {"1", "true", "yes", "on"} @@ -120,11 +128,110 @@ def dashboard_client_last_seen(path: Optional[os.PathLike | str] = None, *, def self_suspend_available(environ: Optional[dict] = None) -> bool: """True iff Fly machine identity is present AND the local Machines API socket exists. - Off-Fly the watcher skips the quiesce: the platform owns the freeze.""" + Off-Fly this is False; see ``suspend_available`` for whether some OTHER lever exists + before concluding the watcher must abstain.""" return bool(_env_str(environ, FLY_APP_NAME_ENV) and _env_str(environ, FLY_MACHINE_ID_ENV) and os.path.exists(FLY_API_SOCKET)) +def brokered_sleep_url(environ: Optional[dict] = None) -> Optional[str]: + """The NAS sleep endpoint to POST, or None when this backend has no broker. + + Validated here rather than at POST time: a malformed value would otherwise let + the watcher mark draining, hold the re-dial and flip the connector before + urllib rejected it, quiescing for a suspend that could never happen. + """ + env = environ if environ is not None else os.environ + url = str(env.get(SLEEP_URL_ENV, "")).strip() + if not url: + return None + parsed = urllib.parse.urlsplit(url) + if parsed.scheme != "https" or not parsed.netloc: + # Once per process: the watcher calls this every idle tick, and its own + # no-lever latch sits AFTER this, so an unlatched warning here would + # repeat for the life of a misconfigured deployment. + global _malformed_sleep_url_logged + if not _malformed_sleep_url_logged: + _malformed_sleep_url_logged = True + logger.warning( + "scale-to-zero: ignoring malformed %s (want an absolute https URL)", + SLEEP_URL_ENV, + ) + return None + return url + + +def suspend_available(environ: Optional[dict] = None) -> bool: + """Whether ANY suspend lever exists, in-guest or brokered. + + Quiescing without one is worse than not quiescing: the re-dial clears the flip. + """ + env = environ if environ is not None else os.environ + return self_suspend_available(env) or brokered_sleep_url(env) is not None + + +# Must EXCEED the broker's own hard request ceiling (NAS route maxDuration = 30s), +# or we give up while it is still working: the redial hold would be released, the +# re-dial would clear the flip, and its stop would then land on a live destination. +BROKERED_SUSPEND_TIMEOUT_S = 40.0 + +# Flaps answers a suspend BEFORE the kernel freezes (suspend_self is documented +# fire-and-forget), so the re-dial fence has to outlive the 2xx or a re-dial can +# restore a live destination in the gap. Measured on Fly (gru, shared-cpu-4x/2GB, +# 2026-09-03) across four suspends: the freeze landed 2.25s, 3.83s, 4.06s and +# 4.19s after the 2xx. The gap is a RAM snapshot, so it grows with machine size — +# 5.0s left only 0.81s of margin on the smallest instance and would be too short +# on a larger one, which is the failure this fence exists to prevent. +# +# Sized generously because the fence is NOT paid on wake: the watcher slices it on +# the wall clock (see _scale_to_zero_await_freeze_gap), so an overshoot costs only +# the reconnect delay on the rare flaps-accepted-but-never-froze path, itself +# capped by ws_transport.REDIAL_HOLD_MAX_S. +FLY_FREEZE_GRACE_S = 15.0 + +# Poll step for that wall-clock slice. Bounds how long after a resume the fence +# lingers before the drain re-dial (the deadline is already past by then). +FLY_FREEZE_GRACE_TICK_S = 0.25 + + +def request_brokered_suspend( + url: str, + *, + timeout: float = BROKERED_SUSPEND_TIMEOUT_S, + opener: Any = None, +) -> bool: + """POST the NAS sleep URL so NAS stops this machine on our behalf. + + Same contract as ``suspend_self``: never raises, True only on 2xx, fail-awake. + """ + import urllib.error + import urllib.request + + # urllib sets Content-Length itself for a bytes body. + request = urllib.request.Request(url, data=b"", method="POST") + open_url = opener or urllib.request.urlopen + try: + with open_url(request, timeout=timeout) as response: + status = int(getattr(response, "status", 0) or 0) + except urllib.error.HTTPError as exc: + # Not retried here: the watcher re-runs on its own interval. + logger.warning( + "scale-to-zero: brokered suspend rejected: %s %s", + exc.code, + exc.reason, + ) + return False + except (urllib.error.URLError, OSError, ValueError) as exc: + logger.warning("scale-to-zero: brokered suspend request failed: %s", exc) + return False + ok = 200 <= status < 300 + if ok: + logger.info("scale-to-zero: machine suspend accepted by NAS (%s)", status) + else: + logger.warning("scale-to-zero: brokered suspend returned %s", status) + return ok + + def suspend_self(environ: Optional[dict] = None, *, socket_path: str = FLY_API_SOCKET, timeout: float = 10.0) -> bool: """POST /v1/apps/{app}/machines/{id}/suspend on the local flaps socket (the socket is the diff --git a/tests/gateway/relay/test_relay_going_idle.py b/tests/gateway/relay/test_relay_going_idle.py index 069a68b867..fefb61bda4 100644 --- a/tests/gateway/relay/test_relay_going_idle.py +++ b/tests/gateway/relay/test_relay_going_idle.py @@ -227,16 +227,11 @@ async def test_go_dormant_redials_on_wake_and_drains(server): finally: await t.disconnect() - -@pytest.mark.asyncio -async def test_adapter_go_dormant_delegates_to_transport(server): - """RelayAdapter.go_dormant() drives the transport's go_dormant (going_idle + - dormant close) without the terminal teardown disconnect() does.""" - from gateway.config import PlatformConfig - from gateway.relay.adapter import RelayAdapter +def _relay_descriptor(): + """Placeholder descriptor for the adapter tests below.""" from gateway.relay.descriptor import CONTRACT_VERSION, CapabilityDescriptor - placeholder = CapabilityDescriptor( + return CapabilityDescriptor( contract_version=CONTRACT_VERSION, platform="discord", label="Relay", @@ -247,6 +242,97 @@ async def test_adapter_go_dormant_delegates_to_transport(server): markdown_dialect="plain", len_unit="chars", ) + + + +def test_adapter_reports_false_when_the_transport_hold_fails(): + """The runner suspends on the strength of this answer, so a transport that is missing the method or throws must read as 'not held' rather than pass for success on the mere absence of an exception.""" + from gateway.config import PlatformConfig + from gateway.relay.adapter import RelayAdapter + + placeholder = _relay_descriptor() + + class Throwing: + def hold_redial(self): + raise RuntimeError("transport refused the hold") + + def release_redial(self): + raise RuntimeError("transport refused the release") + + class Missing: + pass + + for transport in (Throwing(), Missing()): + adapter = RelayAdapter(PlatformConfig(), placeholder, transport=transport) + assert adapter.hold_redial() is False, type(transport).__name__ + assert adapter.release_redial() is False, type(transport).__name__ + + class Working: + def __init__(self): + self.calls = [] + + def hold_redial(self): + self.calls.append("hold") + + def release_redial(self): + self.calls.append("release") + + working = Working() + adapter = RelayAdapter(PlatformConfig(), placeholder, transport=working) + assert adapter.hold_redial() is True + assert working.calls == ["hold"] + + +@pytest.mark.asyncio +async def test_adapter_redial_hold_delegates_to_transport(server): + """The runner only holds the adapter, so this delegation is the only thing + wiring its suspend decision to the transport.""" + from gateway.config import PlatformConfig + from gateway.relay.adapter import RelayAdapter + + placeholder = _relay_descriptor() + transport = WebSocketRelayTransport( + server.url, "discord", "appShared", reconnect=True, reconnect_backoff_s=0.05 + ) + adapter = RelayAdapter(PlatformConfig(), placeholder, transport=transport) + await adapter.connect() + try: + adapter.hold_redial() + assert transport._redial_held is True + adapter.release_redial() + assert transport._redial_held is False + finally: + await adapter.disconnect() + + +@pytest.mark.asyncio +async def test_go_dormant_leaves_the_socket_open_when_the_ack_is_missed(server): + """No ack means nothing will suspend us, so closing would only cost a needless disconnect/reconnect cycle while the agent keeps serving.""" + t = WebSocketRelayTransport( + server.url, "discord", "appShared", reconnect=True, reconnect_backoff_s=0.05 + ) + await t.connect() + try: + async def no_ack(timeout_s=None): + return False + + t.go_idle = no_ack + assert await t.go_dormant() is False + assert t._dormant is False, "must not enter the dormant cadence" + assert t._ws is not None, "socket must stay open" + assert t._closing is False + finally: + await t.disconnect() + + +@pytest.mark.asyncio +async def test_adapter_go_dormant_delegates_to_transport(server): + """RelayAdapter.go_dormant() drives the transport's go_dormant (going_idle + + dormant close) without the terminal teardown disconnect() does.""" + from gateway.config import PlatformConfig + from gateway.relay.adapter import RelayAdapter + + placeholder = _relay_descriptor() transport = WebSocketRelayTransport( server.url, "discord", "appShared", reconnect=True, reconnect_backoff_s=0.05 ) diff --git a/tests/gateway/relay/test_ws_transport_hardening.py b/tests/gateway/relay/test_ws_transport_hardening.py index af923af4bc..f8b48a0030 100644 --- a/tests/gateway/relay/test_ws_transport_hardening.py +++ b/tests/gateway/relay/test_ws_transport_hardening.py @@ -371,3 +371,71 @@ async def test_send_raising_socket_returns_error_dict(): fake.reader_release.set() await t._reader + + +@pytest.mark.asyncio +async def test_redial_hold_parks_the_supervisor_until_released(monkeypatch): + """A re-dial mid-suspend clears the dormant flip, so it must be parked.""" + dials = [] + + async def _never_dials(self): + dials.append(1) + + monkeypatch.setattr( + WebSocketRelayTransport, "_dial_and_start", _never_dials, raising=True + ) + + t = WebSocketRelayTransport( + "ws://unused", "discord", "bot1", reconnect=True, reconnect_backoff_s=0.01 + ) + t._dormant_redial_s = 0.01 + t.hold_redial() + + supervisor = asyncio.create_task(t._reconnect_loop()) + try: + # Well past the cadence: unheld, this would have dialled many times. + await asyncio.sleep(0.15) + assert dials == [], "supervisor re-dialled while a brokered suspend was in flight" + + # A refused suspend releases it and the re-dial resumes. + t.release_redial() + for _ in range(100): + if dials: + break + await asyncio.sleep(0.01) + assert dials == [1] + finally: + t._closing = True + supervisor.cancel() + + +@pytest.mark.asyncio +async def test_redial_hold_expires_so_a_lost_suspend_cannot_strand_us(monkeypatch): + """Bounded: a suspend that never lands must reconnect, not hold forever.""" + dials = [] + + async def _count_dial(self): + dials.append(1) + + monkeypatch.setattr( + WebSocketRelayTransport, "_dial_and_start", _count_dial, raising=True + ) + + t = WebSocketRelayTransport( + "ws://unused", "discord", "bot1", reconnect=True, reconnect_backoff_s=0.01 + ) + t._dormant_redial_s = 0.01 + t._redial_hold_max_s = 0.05 + t.hold_redial() + + supervisor = asyncio.create_task(t._reconnect_loop()) + try: + for _ in range(100): + if dials: + break + await asyncio.sleep(0.01) + assert dials == [1] + assert t._redial_held is False + finally: + t._closing = True + supervisor.cancel() diff --git a/tests/gateway/test_scale_to_zero.py b/tests/gateway/test_scale_to_zero.py index 0f604850c4..7334633f7c 100644 --- a/tests/gateway/test_scale_to_zero.py +++ b/tests/gateway/test_scale_to_zero.py @@ -186,3 +186,106 @@ def test_self_suspend_available_needs_identity_and_socket(): assert self_suspend_available(_FLY_ENV) is False # Missing identity -> unavailable regardless of socket. assert self_suspend_available({}) is False + +# Brokered suspend: Azure's stop verb needs a credential the sandbox lacks, so +# NAS stamps a signed sleep URL and stops the machine on our POST. + +from gateway.scale_to_zero import ( # noqa: E402 - grouped with their section + SLEEP_URL_ENV, + brokered_sleep_url, + request_brokered_suspend, + suspend_available, +) + +_SLEEP_URL = "https://portal.example.com/api/agents/inst-1/sleep?t=sig" + + +def test_brokered_sleep_url_reads_the_stamp(): + assert brokered_sleep_url({SLEEP_URL_ENV: _SLEEP_URL}) == _SLEEP_URL + assert brokered_sleep_url({}) is None + assert brokered_sleep_url({SLEEP_URL_ENV: " "}) is None + + +def test_malformed_sleep_url_warns_once_not_every_idle_tick(monkeypatch, caplog): + """The watcher calls this on every tick and its own no-lever latch sits after + it, so an unlatched warning would repeat for the life of the process.""" + import gateway.scale_to_zero as sz + + monkeypatch.setattr(sz, "_malformed_sleep_url_logged", False) + monkeypatch.setenv("GATEWAY_RELAY_SLEEP_URL", "http://portal.example.com/x") + with caplog.at_level("WARNING"): + for _ in range(5): + assert sz.brokered_sleep_url() is None + assert sum("ignoring malformed" in r.message for r in caplog.records) == 1 + + +def test_brokered_sleep_url_rejects_a_malformed_stamp(monkeypatch): + """Validated before it is treated as a lever: otherwise the watcher marks draining, holds the re-dial and flips the connector, and only then does urllib reject the value, quiescing for a suspend that could never happen.""" + from gateway.scale_to_zero import brokered_sleep_url + + for bad in ("not-a-url", "http://portal.example.com/x", "/api/agents/i/sleep"): + monkeypatch.setenv("GATEWAY_RELAY_SLEEP_URL", bad) + assert brokered_sleep_url() is None, bad + monkeypatch.setenv("GATEWAY_RELAY_SLEEP_URL", "https://portal.example.com/x?t=s") + assert brokered_sleep_url() == "https://portal.example.com/x?t=s" + + +def test_suspend_available_accepts_either_lever(monkeypatch): + # No flaps socket but a broker exists, so the watcher may still quiesce. + monkeypatch.setattr( + "gateway.scale_to_zero.self_suspend_available", lambda *a, **k: False + ) + assert suspend_available({SLEEP_URL_ENV: _SLEEP_URL}) is True + assert suspend_available({}) is False + monkeypatch.setattr( + "gateway.scale_to_zero.self_suspend_available", lambda *a, **k: True + ) + assert suspend_available({}) is True + + +class _FakeResponse: + def __init__(self, status): + self.status = status + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + +def test_brokered_timeout_outlasts_the_broker_and_fits_the_redial_hold(): + """Cross-repo deadline chain.""" + from gateway.relay.ws_transport import REDIAL_HOLD_MAX_S + from gateway.scale_to_zero import BROKERED_SUSPEND_TIMEOUT_S + + NAS_ROUTE_MAX_DURATION_S = 30.0 # mirrors the sleep route's maxDuration + assert NAS_ROUTE_MAX_DURATION_S < BROKERED_SUSPEND_TIMEOUT_S + assert BROKERED_SUSPEND_TIMEOUT_S < REDIAL_HOLD_MAX_S + + +def test_request_brokered_suspend_posts_the_signed_url(): + seen = {} + + def opener(request, timeout=None): + seen["url"] = request.full_url + seen["method"] = request.get_method() + return _FakeResponse(200) + + assert request_brokered_suspend(_SLEEP_URL, opener=opener) is True + assert seen == {"url": _SLEEP_URL, "method": "POST"} + + +def test_request_brokered_suspend_fails_awake(): + # Fail-awake: any refusal leaves the machine running rather than stranding a + # frozen peer that still looks live to the connector. + import urllib.error + + def http_error(request, timeout=None): + raise urllib.error.HTTPError(_SLEEP_URL, 409, "Conflict", {}, None) + + def unreachable(request, timeout=None): + raise urllib.error.URLError("connection refused") + + for opener in (http_error, unreachable, lambda *a, **k: _FakeResponse(500)): + assert request_brokered_suspend(_SLEEP_URL, opener=opener) is False diff --git a/tests/gateway/test_scale_to_zero_watcher.py b/tests/gateway/test_scale_to_zero_watcher.py index 5ef261a921..5fa55054e9 100644 --- a/tests/gateway/test_scale_to_zero_watcher.py +++ b/tests/gateway/test_scale_to_zero_watcher.py @@ -10,6 +10,7 @@ respects the cooldown, and skips when busy — the F7/D3 + D12 behaviour. from __future__ import annotations +import contextlib import asyncio import time @@ -19,15 +20,47 @@ from gateway.run import GatewayRunner class _FakeRelayAdapter: - def __init__(self): + def __init__(self, ack=True): self.go_dormant_calls = 0 + self.redial = [] + self.ack = ack async def go_dormant(self): self.go_dormant_calls += 1 + return self.ack + + def hold_redial(self): + self.redial.append("hold") + return True + + def release_redial(self): + self.redial.append("release") return True -def _runner_with(monkeypatch, *, idle, armed_adapter=True, can_self_suspend=True): +async def _run_one_iteration(r, *, interval=0.01, settle=0.1): + """Run the watcher long enough for one iteration, then stop it cleanly.""" + task = asyncio.create_task(r._scale_to_zero_watcher(interval=interval)) + await asyncio.sleep(settle) + r._running = False + await asyncio.wait_for(task, timeout=2) + + +async def _noop_async(*a, **k): + return None + + +def _runner_with( + monkeypatch, + *, + idle, + armed_adapter=True, + can_self_suspend=True, + brokered=False, + ack=True, + idle_readings=None, + draining=False, +): """Build a GatewayRunner without booting it, stubbing just what the watcher touches. Real methods (_scale_to_zero_is_idle composition, the watcher body) run; only their dependencies are stubbed. @@ -39,17 +72,36 @@ def _runner_with(monkeypatch, *, idle, armed_adapter=True, can_self_suspend=True """ r = GatewayRunner.__new__(GatewayRunner) r._running = True + r._draining = draining r._scale_to_zero_cooldown_until = 0.0 r._scale_to_zero_no_suspend_logged = False r._last_inbound_at = time.time() r._running_agents = {} r._background_tasks = set() - adapter = _FakeRelayAdapter() if armed_adapter else None + adapter = _FakeRelayAdapter(ack=ack) if armed_adapter else None - monkeypatch.setattr(r, "_scale_to_zero_is_idle", lambda: idle, raising=False) + readings = iter(idle_readings) if idle_readings else None + monkeypatch.setattr( + r, + "_scale_to_zero_is_idle", + (lambda: next(readings, False)) if readings else (lambda: idle), + raising=False, + ) monkeypatch.setattr(r, "_relay_adapter_for_dormancy", lambda: adapter, raising=False) monkeypatch.setattr(r, "_scale_to_zero_idle_timeout_seconds", lambda: 300.0, raising=False) - monkeypatch.setattr(r, "_update_runtime_status", lambda *a, **k: None, raising=False) + r.states = [] + monkeypatch.setattr( + r, + "_update_runtime_status", + lambda *a, **k: r.states.append(a[0] if a else None), + raising=False, + ) + if brokered: + can_self_suspend = False + monkeypatch.setenv( + "GATEWAY_RELAY_SLEEP_URL", + "https://portal.example.com/api/agents/i/sleep?t=s", + ) monkeypatch.setattr( "gateway.scale_to_zero.self_suspend_available", lambda *a, **k: can_self_suspend, @@ -58,13 +110,11 @@ def _runner_with(monkeypatch, *, idle, armed_adapter=True, can_self_suspend=True @pytest.mark.asyncio -async def test_watcher_does_not_quiesce_when_the_platform_owns_the_suspend( +async def test_watcher_does_not_quiesce_when_no_suspend_lever_exists( monkeypatch, ): - """Quiescing cannot help when the platform owns the freeze, and the reconnect - that follows the socket close undoes the flip, so the destination ends up - unflipped when the freeze lands. - """ + """With no lever at all, the re-dial after the socket close just undoes the + flip, so quiescing cannot help. Stay connected instead.""" r, adapter = _runner_with(monkeypatch, idle=True, can_self_suspend=False) suspends = [] monkeypatch.setattr( @@ -74,10 +124,7 @@ async def test_watcher_does_not_quiesce_when_the_platform_owns_the_suspend( raising=False, ) - task = asyncio.create_task(r._scale_to_zero_watcher(interval=0.01)) - await asyncio.sleep(0.1) - r._running = False - await asyncio.wait_for(task, timeout=2) + await _run_one_iteration(r) assert adapter.go_dormant_calls == 0, "must not flip/close on a platform-timed suspend" assert suspends == [] @@ -87,14 +134,237 @@ async def test_watcher_does_not_quiesce_when_the_platform_owns_the_suspend( assert r._scale_to_zero_no_suspend_logged is True +@pytest.mark.asyncio +@pytest.mark.parametrize( + "lever,kwargs,redial", + [ + # Fly holds past the 2xx (flaps answers before the freeze) and releases + # once the gap closes; the gap itself is covered separately below. + ("in-guest", {"can_self_suspend": True}, ["hold", "release"]), + # The brokered stop is still in flight, so the hold stays. + ("brokered", {"brokered": True}, ["hold"]), + ], +) +async def test_watcher_quiesces_then_suspends_on_either_lever( + monkeypatch, lever, kwargs, redial +): + """Flip first, freeze second, re-dial held across it: the ordering the feature rests on.""" + r, adapter = _runner_with(monkeypatch, idle=True, **kwargs) + monkeypatch.setattr("gateway.scale_to_zero.suspend_self", lambda *a, **k: True) + monkeypatch.setattr( + "gateway.scale_to_zero.request_brokered_suspend", lambda *a, **k: True + ) + # Not what this test is about; the freeze gap has its own case below. + monkeypatch.setattr("gateway.scale_to_zero.FLY_FREEZE_GRACE_S", 0.0) + + await _run_one_iteration(r, settle=0.15) + + assert adapter.go_dormant_calls == 1, lever + assert adapter.redial == redial, lever + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "in_guest,accepted,lever,redial", + [ + # Fly holds too: flaps answers seconds BEFORE the freeze, so the fence + # spans that gap and only then releases. The gap itself has its own cases + # below; here the grace is zeroed so this stays a lever-choice test. + (True, True, "flaps", ["release"]), + # Brokered + accepted: the watcher's hold stays, the stop is still in flight. + (False, True, "brokered", []), + # Brokered + refused: nothing will freeze us, so give the supervisor back. + (False, False, "brokered", ["release"]), + ], +) +async def test_self_suspend_picks_a_lever_and_releases_only_when_nothing_will_freeze( + monkeypatch, in_guest, accepted, lever, redial +): + r, adapter = _runner_with(monkeypatch, idle=True, can_self_suspend=in_guest) + monkeypatch.setattr("gateway.scale_to_zero.FLY_FREEZE_GRACE_S", 0.0) + monkeypatch.setenv( + "GATEWAY_RELAY_SLEEP_URL", "https://portal.example.com/api/agents/i/sleep?t=s" + ) + used = [] + monkeypatch.setattr( + "gateway.scale_to_zero.suspend_self", + lambda *a, **k: used.append("flaps") or accepted, + ) + monkeypatch.setattr( + "gateway.scale_to_zero.request_brokered_suspend", + lambda *a, **k: used.append("brokered") or accepted, + ) + + await r._scale_to_zero_self_suspend() + + assert used == [lever] + assert adapter.redial == redial + + +@pytest.mark.asyncio +async def test_watcher_honours_a_false_hold_from_the_adapter(monkeypatch): + """Quiescing without the hold leaves the re-dial free to clear the flip + mid-suspend, so a False from the adapter must stop the attempt.""" + r, adapter = _runner_with(monkeypatch, idle=True, brokered=True) + adapter.hold_redial = lambda: False + suspends = [] + monkeypatch.setattr( + r, "_scale_to_zero_self_suspend", lambda: suspends.append(1) or _noop_async() + ) + + await _run_one_iteration(r) + + assert suspends == [] + assert adapter.go_dormant_calls == 0 + + + + +@pytest.mark.asyncio +async def test_hold_redial_reports_failure_when_there_is_no_adapter(monkeypatch): + """The return value gates the suspend, so an absent adapter must read as 'not held' rather than silently as success.""" + r, _ = _runner_with(monkeypatch, idle=True, armed_adapter=False) + + assert r._scale_to_zero_hold_redial(True) is False + + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "case,kwargs", + [ + # A missed ack means inbound is NOT buffered yet, so freezing would drop + # it. None counts as missed too: a partially-wired transport returning + # nothing has not acked either. + ("unacked-false", {"ack": False}), + ("unacked-none", {"ack": None}), + # Idle at the top of the tick, busy by the time the quiesce returns. + ("inbound-mid-quiesce", {"idle_readings": [True, False]}), + ], +) +async def test_watcher_abandons_cleanly_when_it_must_not_suspend( + monkeypatch, case, kwargs +): + """Every abort path leaves no trace: no suspend, hold released, running restored.""" + r, adapter = _runner_with(monkeypatch, idle=True, brokered=True, **kwargs) + suspends = [] + monkeypatch.setattr( + r, "_scale_to_zero_self_suspend", lambda: suspends.append(1) or _noop_async() + ) + + await _run_one_iteration(r) + + assert suspends == [], case + assert adapter.redial[-1] == "release", case + assert r.states[:2] == ["draining", "running"], case + +@pytest.mark.asyncio +async def test_in_guest_release_waits_for_the_freeze_gap(monkeypatch): + """flaps answers before the kernel freezes, so the fence spans that gap. A + machine that never froze must still get its supervisor back.""" + r, adapter = _runner_with(monkeypatch, idle=True, can_self_suspend=True) + monkeypatch.setattr("gateway.scale_to_zero.suspend_self", lambda *a, **k: True) + monkeypatch.setattr("gateway.scale_to_zero.FLY_FREEZE_GRACE_S", 0.0) + + await _run_one_iteration(r, settle=0.15) + + assert adapter.redial == ["hold", "release"] + + +@pytest.mark.asyncio +async def test_in_guest_fence_still_held_inside_the_freeze_gap(monkeypatch): + """Observed mid-gap: flaps has answered but the freeze has not landed, so the + supervisor must still be parked.""" + r, adapter = _runner_with(monkeypatch, idle=True, can_self_suspend=True) + monkeypatch.setattr("gateway.scale_to_zero.suspend_self", lambda *a, **k: True) + monkeypatch.setattr("gateway.scale_to_zero.FLY_FREEZE_GRACE_S", 5.0) + + task = asyncio.create_task(r._scale_to_zero_watcher(interval=0.01)) + await asyncio.sleep(0.15) + assert adapter.redial == ["hold"], "released before the freeze could land" + r._running = False + task.cancel() + with contextlib.suppress(asyncio.CancelledError): + await task + + +@pytest.mark.asyncio +async def test_in_guest_fence_releases_at_once_after_a_resume(monkeypatch): + """A Fly suspend stops CLOCK_MONOTONIC but keeps CLOCK_REALTIME tracking host + time (measured: 252.219s frozen -> monotonic +0.501s, realtime +252.219s). The + fence is therefore sliced on the wall clock, so a machine that froze mid-fence + re-dials to drain the moment it wakes instead of waiting out the remainder -- + a plain asyncio.sleep() here would cost that remainder on EVERY Fly wake.""" + r, adapter = _runner_with(monkeypatch, idle=True, can_self_suspend=True) + monkeypatch.setattr("gateway.scale_to_zero.FLY_FREEZE_GRACE_S", 30.0) + monkeypatch.setattr("gateway.scale_to_zero.FLY_FREEZE_GRACE_TICK_S", 0.01) + monkeypatch.setattr("gateway.scale_to_zero.suspend_self", lambda *a, **k: True) + + real_time = time.time + # Read 1 sets the deadline; read 2 is the first post-"resume" check, with the + # wall clock a freeze further on. asyncio.sleep would still owe ~30s here. + reads = iter([1000.0, 1000.0 + 252.219]) + monkeypatch.setattr( + "gateway.run.time.time", lambda: next(reads, 1000.0 + 252.219) + ) + + started = real_time() + await r._scale_to_zero_self_suspend() + elapsed = real_time() - started + + assert adapter.redial == ["release"] + assert elapsed < 1.0, f"fence waited out the monotonic remainder ({elapsed:.2f}s)" + + +@pytest.mark.asyncio +async def test_abort_sets_a_cooldown_so_it_does_not_retry_every_tick(monkeypatch): + r, adapter = _runner_with(monkeypatch, idle=True, brokered=True) + adapter.hold_redial = lambda: False + + await _run_one_iteration(r) + + assert r._scale_to_zero_cooldown_until > time.time() + + +@pytest.mark.asyncio +async def test_abort_never_resurrects_a_shutting_down_gateway(monkeypatch): + """A real shutdown drain must win: `running` here would clobber it.""" + r, _ = _runner_with(monkeypatch, idle=True, brokered=True, ack=False, draining=True) + + await _run_one_iteration(r) + + assert "running" not in r.states + + +@pytest.mark.asyncio +async def test_watcher_holds_redial_before_going_dormant(monkeypatch): + """The hold must precede go_dormant: its close arms the reconnect supervisor, and a re-dial would clear the flip the suspend depends on.""" + r, adapter = _runner_with(monkeypatch, idle=True, brokered=True) + order = [] + original = adapter.go_dormant + + async def recording_go_dormant(): + order.append("go_dormant") + return await original() + + adapter.go_dormant = recording_go_dormant + monkeypatch.setattr( + r, + "_scale_to_zero_hold_redial", + lambda held: bool(order.append(f"hold={held}")) or True, + ) + monkeypatch.setattr(r, "_scale_to_zero_self_suspend", _noop_async) + + await _run_one_iteration(r) + + assert order[:2] == ["hold=True", "go_dormant"] + + @pytest.mark.asyncio async def test_watcher_goes_dormant_when_idle(monkeypatch): r, adapter = _runner_with(monkeypatch, idle=True) # Run one iteration: stop after the first sleep so the loop exits cleanly. - task = asyncio.create_task(r._scale_to_zero_watcher(interval=0.01)) - await asyncio.sleep(0.1) - r._running = False - await asyncio.wait_for(task, timeout=2) + await _run_one_iteration(r) assert adapter.go_dormant_calls >= 1 # After driving dormant, a re-arm cooldown is set (0.F). assert r._scale_to_zero_cooldown_until > time.time() @@ -243,10 +513,7 @@ async def test_watcher_skips_suspend_when_dormant_fails(monkeypatch): suspend_calls.append(1) monkeypatch.setattr(r, "_scale_to_zero_self_suspend", fake_suspend, raising=False) - task = asyncio.create_task(r._scale_to_zero_watcher(interval=0.01)) - await asyncio.sleep(0.1) - r._running = False - await asyncio.wait_for(task, timeout=2) + await _run_one_iteration(r) # A failed quiesce means an UNFLIPPED relay — suspending would black-hole # inbound events. Must stay awake. assert suspend_calls == [] @@ -266,29 +533,26 @@ async def test_watcher_skips_suspend_when_inbound_lands_mid_quiesce(monkeypatch) suspend_calls.append(1) monkeypatch.setattr(r, "_scale_to_zero_self_suspend", fake_suspend, raising=False) - task = asyncio.create_task(r._scale_to_zero_watcher(interval=0.01)) - await asyncio.sleep(0.15) - r._running = False - await asyncio.wait_for(task, timeout=2) + await _run_one_iteration(r, settle=0.15) assert adapter.go_dormant_calls == 1 assert suspend_calls == [] @pytest.mark.asyncio -async def test_self_suspend_noop_off_fly(monkeypatch): - """Off-Fly (no flaps socket/identity) the helper is a silent no-op — - dormancy without platform suspend, never an error.""" - r = GatewayRunner.__new__(GatewayRunner) - monkeypatch.setattr( - "gateway.scale_to_zero.self_suspend_available", lambda *a, **k: False - ) +async def test_self_suspend_noop_with_no_lever(monkeypatch): + """Neither an in-guest API nor a brokered URL: a silent no-op, never an error.""" + r, adapter = _runner_with(monkeypatch, idle=True, can_self_suspend=False) + monkeypatch.delenv("GATEWAY_RELAY_SLEEP_URL", raising=False) called = [] monkeypatch.setattr( "gateway.scale_to_zero.suspend_self", lambda *a, **k: called.append(1) or True, ) + await r._scale_to_zero_self_suspend() + assert called == [] + assert adapter.redial == ["release"] # ── non-messaging platforms must not disarm (the api_server-key regression) ──