fix(gateway): stop dropping the first message after an Azure agent sleeps (#99736)

* feat(gateway): let NAS broker the scale-to-zero suspend where the guest has no lever

* fix(gateway): require the going_idle ack and outlast the broker before suspending

* fix(gateway): release the redial hold on every path that abandons the suspend

* fix(gateway): hold the re-dial only for the brokered lever, and require it

* fix(gateway): make the re-dial hold contract real and fence the in-guest suspend too

* refactor(gateway): share the watcher-iteration and descriptor helpers across the sleep tests

* refactor(gateway): move the sleep tests' repeated setup into their fixtures

* refactor(gateway): one docstring line per sleep test, and parametrise the abort and lever paths

* fix(gateway): fence the freeze gap, cool down aborts, and stop clobbering a shutdown drain

* fix(gateway): slice the Fly freeze fence on the wall clock and widen it for larger machines
This commit is contained in:
Nacho Avecilla
2026-09-04 15:41:27 -03:00
committed by GitHub
parent 59d330fd15
commit a6e10e693f
8 changed files with 865 additions and 63 deletions
+25
View File
@@ -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,
+46 -2
View File
@@ -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:
+124 -19
View File
@@ -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).
+108 -1
View File
@@ -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
+94 -8
View File
@@ -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
)
@@ -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()
+103
View File
@@ -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
+297 -33
View File
@@ -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) ──