fix(gateway): release a boot claim whose adapter vanished before dispatch
The boot sweep now moves every claimed row to attempting. If the platform went fatal between the claim and the send (restart notification, flood sleep), _obligation_adapter skipped the row without releasing it, so it sat in attempting owned by this live process: the reconnect sweep only takes failed rows and the boot sweep skips live owners, so the reply waited for the next restart (and then carried a false duplicate marker). Release any undispatched claim, not only runtime ones. Found by independent review of #120450.
This commit is contained in:
@@ -478,13 +478,13 @@ class GatewayStartupMixin:
|
||||
else:
|
||||
# Startup rows preserve the historical default-adapter route.
|
||||
adapter = self.adapters.get(platform)
|
||||
# A runtime claim whose reconnect vanished before dispatch is released without spending an
|
||||
# attempt; startup claims stay 'attempting' for the next boot's marked redelivery (attempts cap +
|
||||
# stale cutoff bound retries).
|
||||
# A claim (runtime or boot) whose adapter vanished before dispatch was never sent: release it
|
||||
# without spending an attempt so the reconnect sweep, which only takes 'failed' rows, delivers it
|
||||
# instead of it waiting 'attempting' for the next restart.
|
||||
# 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"):
|
||||
if adapter is None:
|
||||
from gateway.delivery_ledger import is_flood_error
|
||||
|
||||
last_error = row.get("last_error")
|
||||
|
||||
@@ -519,6 +519,24 @@ class TestGatewayRedeliverySweep:
|
||||
|
||||
assert dl.sweep_failed_for_runtime("telegram") == []
|
||||
|
||||
@pytest.mark.parametrize("claimed_state", ["pending", "failed"])
|
||||
@pytest.mark.asyncio
|
||||
async def test_boot_claim_whose_adapter_vanishes_is_left_for_the_reconnect_sweep(self, claimed_state):
|
||||
"""A boot claim never sent (platform went away before dispatch) must not strand in 'attempting'."""
|
||||
_record()
|
||||
if claimed_state == "failed":
|
||||
dl.mark_failed("ob-1", "send_path_degraded")
|
||||
_orphan("ob-1")
|
||||
runner = self._runner(self._adapter())
|
||||
claimed = await runner._claim_pending_obligations()
|
||||
assert [r["obligation_id"] for r in claimed] == ["ob-1"]
|
||||
runner.adapters = {} # platform went fatal during the restart notification / flood sleep
|
||||
|
||||
await runner._redeliver_claimed_obligations(claimed)
|
||||
|
||||
assert _row("ob-1")["state"] == "failed" and _row("ob-1")["attempts"] == 0
|
||||
assert [r["obligation_id"] for r in dl.sweep_failed_for_runtime("slack")] == ["ob-1"]
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_runtime_failed_redelivery_clears_resume_before_send(self):
|
||||
from gateway.config import Platform
|
||||
|
||||
Reference in New Issue
Block a user