fix(gateway): a released non-flood claim becomes reconnect-only

Gate review: releasing an undispatched runtime claim kept its pre-claim
error ("503 ..."), so the redelivery timer re-armed and claimed/released
the row every tier for as long as the adapter was absent - the loop the
previous commit closed for `send_path_degraded` rows only. Only a flood
row keeps its error (the platform wait must be honoured); every other
row is released as reconnect-only and re-claimed by the reconnect sweep.
`is_reconnect_only()` replaces the three spellings of that predicate.
This commit is contained in:
kshitijk4poor
2026-09-19 01:12:13 +05:30
committed by kshitij
parent 07a1d55d33
commit 877848d2f3
3 changed files with 16 additions and 5 deletions

View File

@@ -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:<seconds>`` 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:

View File

@@ -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)

View File

@@ -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