merge: refresh upstream models, notification expiry, and desktop controls
This commit is contained in:
@@ -26,9 +26,14 @@ issue #26241 for details.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any, Dict, Optional, Union
|
||||
import time
|
||||
import uuid
|
||||
from typing import Any, Callable, Dict, Optional, Union
|
||||
from urllib.parse import urlencode
|
||||
|
||||
# A 429 whose Retry-After is longer than this is reported, not waited out inside a tool call.
|
||||
MANAGED_FAL_RATE_LIMIT_RETRY_CAP_SECONDS = 30.0
|
||||
|
||||
|
||||
def import_fal_client() -> Any:
|
||||
"""Import ``fal_client`` (via ``pm`` when available) and return
|
||||
@@ -99,6 +104,61 @@ def _managed_fal_billing_error(exc: BaseException, what: str) -> Optional[str]:
|
||||
)
|
||||
|
||||
|
||||
def _managed_fal_retry_after_seconds(exc: BaseException) -> Optional[float]:
|
||||
"""Seconds the managed gateway asked us to wait after a 429: the ``Retry-After`` header,
|
||||
else the body's ``error.retryAfter``; None when the status is not 429 or neither is present."""
|
||||
response = getattr(exc, "response", None)
|
||||
if response is None or _extract_http_status(exc) != 429:
|
||||
return None
|
||||
headers = getattr(response, "headers", None)
|
||||
raw = headers.get("Retry-After") if headers is not None and hasattr(headers, "get") else None
|
||||
if raw is None:
|
||||
try:
|
||||
error = response.json().get("error")
|
||||
except Exception: # noqa: BLE001 — a non-JSON 429 body simply has no hint
|
||||
return None
|
||||
raw = error.get("retryAfter") if isinstance(error, dict) else None
|
||||
try:
|
||||
return float(raw) if raw is not None else None
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _managed_fal_rate_limit_message(what: str, name: str, retry_after: Optional[float]) -> str:
|
||||
hint = f"retry after {retry_after:g}s" if retry_after is not None else "no Retry-After given"
|
||||
return (
|
||||
f"Nous Subscription gateway rate-limited {what} '{name}' (HTTP 429; {hint}). "
|
||||
"The model is enabled — retry later instead of switching models or setting FAL_KEY."
|
||||
)
|
||||
|
||||
|
||||
def submit_managed_fal_with_rate_limit_retry(
|
||||
submit: Callable[[Dict[str, str]], Any], *, what: str, name: str,
|
||||
):
|
||||
"""Call ``submit(headers)`` with a fresh ``x-idempotency-key``; on a 429 whose Retry-After
|
||||
fits the cap, wait it out (interrupt-aware) and resubmit ONCE under a new key.
|
||||
|
||||
A second 429, or one with an unknown/too-long Retry-After, raises ValueError naming the
|
||||
rate limit — it must never fall through to the callers' "model may not be enabled" 4xx text,
|
||||
which sends agents off to switch models. Every other exception propagates untouched.
|
||||
"""
|
||||
from tools.interrupt import is_interrupted
|
||||
for attempt in (1, 2):
|
||||
try:
|
||||
return submit({"x-idempotency-key": str(uuid.uuid4())})
|
||||
except Exception as exc:
|
||||
if _extract_http_status(exc) != 429:
|
||||
raise
|
||||
retry_after = _managed_fal_retry_after_seconds(exc)
|
||||
if attempt == 2 or retry_after is None or retry_after > MANAGED_FAL_RATE_LIMIT_RETRY_CAP_SECONDS:
|
||||
raise ValueError(_managed_fal_rate_limit_message(what, name, retry_after)) from exc
|
||||
deadline = time.monotonic() + retry_after
|
||||
while time.monotonic() < deadline:
|
||||
if is_interrupted():
|
||||
raise ValueError(_managed_fal_rate_limit_message(what, name, retry_after)) from exc
|
||||
time.sleep(min(0.5, max(0.0, deadline - time.monotonic())))
|
||||
|
||||
|
||||
def _require(value: Any, what: str) -> Any:
|
||||
if value is None:
|
||||
raise RuntimeError(f"{what} is required for managed FAL gateway mode")
|
||||
|
||||
Reference in New Issue
Block a user