diff --git a/gateway/delivery_ledger.py b/gateway/delivery_ledger.py index 0bda393f3d..255ca60d09 100644 --- a/gateway/delivery_ledger.py +++ b/gateway/delivery_ledger.py @@ -93,6 +93,11 @@ def _raw_flood_wait(text: str) -> Optional[float]: return None +def is_reconnect_only(error: Any) -> bool: + """True for a row that only an adapter reconnect may retry (no timer, no backoff).""" + return str(error or "").strip().lower() in _RUNTIME_RETRYABLE_ERRORS + + def is_flood_error(error: Any) -> bool: """True for a flood refusal: the adapters' fail-closed ``flood_control:`` result, or a row still carrying the platform's own flood wording (see ``_RAW_FLOOD_RE``).""" @@ -153,7 +158,7 @@ def retry_not_before(updated_at: Any, last_error: Any, attempts: Any) -> Optiona if is_flood_error(last_error): return flood_not_before(updated_at, last_error) text = str(last_error or "").strip().lower() - if text in _RUNTIME_RETRYABLE_ERRORS: + if is_reconnect_only(text): return _failed_stamp(updated_at) if classify_dead_error(text): return None @@ -485,7 +490,7 @@ def pending_retries(now: Optional[float] = None) -> List[Dict[str, Any]]: for platform, adapter_profile, updated_at, last_error, attempts, created_at in rows: # Reconnect-only rows (a claim released because the adapter was gone) are re-claimed by the # reconnect sweep; a timer would claim and release them every tick until the adapter is back. - if str(last_error or "").strip().lower() in _RUNTIME_RETRYABLE_ERRORS: + if is_reconnect_only(last_error): continue due = retry_not_before(updated_at, last_error, attempts) if due is None or attempts >= MAX_ATTEMPTS or (now - created_at) > STALE_AFTER_SECONDS: diff --git a/gateway/platforms/base.py b/gateway/platforms/base.py index 935ff4878a..83b8a03985 100644 --- a/gateway/platforms/base.py +++ b/gateway/platforms/base.py @@ -4009,13 +4009,13 @@ class BasePlatformAdapter(ABC): backoff has passed instead of waiting for the next restart (#91653).""" try: from gateway.dead_targets import classify_dead_error - from gateway.delivery_ledger import mark_delivered, mark_failed + from gateway.delivery_ledger import is_reconnect_only, mark_delivered, mark_failed if getattr(result, "success", False): await asyncio.to_thread(mark_delivered, obligation_id) return error = str(getattr(result, "error", "") or "") await asyncio.to_thread(mark_failed, obligation_id, error) - if error == "send_path_degraded": + if is_reconnect_only(error): redeliver = getattr( self.gateway_runner, "_redeliver_failed_obligations_for_platform", None) live = self._final_delivery_adapter(event.source) diff --git a/gateway/run_startup.py b/gateway/run_startup.py index ca06c85794..4686582bb4 100644 --- a/gateway/run_startup.py +++ b/gateway/run_startup.py @@ -449,10 +449,16 @@ class GatewayStartupMixin: adapter = self.adapters.get(platform) # A runtime claim whose reconnect vanished before dispatch is released without spending an # attempt; startup claims keep their state (attempts cap + stale cutoff bound retries). + # Only a flood row keeps its error (the platform's wait must be honoured); any other row + # becomes reconnect-only, or the redelivery timer would claim and release it until the + # adapter is back. if adapter is None and row.get("runtime_recovery"): + from gateway.delivery_ledger import is_flood_error + + last_error = row.get("last_error") await self._release_runtime_claim_quiet( row["obligation_id"], "failed to release undispatched runtime obligation %s", - error=row.get("last_error") or "send_path_degraded", + error=last_error if is_flood_error(last_error) else "send_path_degraded", ) return adapter