test(gateway): pin marked redelivery after an earlier boot's claim
Two invariants on the boot outbox sweep, both red on the previous ledger: - a pending row whose boot redelivery was interrupted after the platform accepted it (or that an older build claimed) is redelivered with the recovered marker, never a second plain copy; - a boot-claimed failed row is not re-claimable by the runtime reconnect sweep while the boot send is in flight (it used to stay 'failed' and owned by this process, so a reconnect could send it a second time). Also corrects the startup-claim comment in _obligation_adapter and the messaging docs bullet on mid-send recovery.
This commit is contained in:
@@ -479,7 +479,8 @@ class GatewayStartupMixin:
|
||||
# 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 keep their state (attempts cap + stale cutoff bound retries).
|
||||
# attempt; startup claims stay 'attempting' for the next boot's marked redelivery (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.
|
||||
|
||||
@@ -482,6 +482,43 @@ class TestGatewayRedeliverySweep:
|
||||
assert sent["content"].startswith(dl.RECOVERED_MARKER)
|
||||
assert sent["content"].endswith("the final answer")
|
||||
|
||||
@pytest.mark.parametrize("earlier_boot", ["killed_inside_send", "older_build_left_pending"])
|
||||
@pytest.mark.asyncio
|
||||
async def test_redelivery_after_an_earlier_boot_claim_is_marked(self, earlier_boot):
|
||||
"""Once a boot has claimed a row it may have sent it: every later copy carries the marker."""
|
||||
import asyncio
|
||||
|
||||
_record()
|
||||
_orphan("ob-1")
|
||||
if earlier_boot == "killed_inside_send":
|
||||
# The platform accepts the plain resend, then the boot dies before mark_delivered.
|
||||
first = MagicMock()
|
||||
first.send = AsyncMock(side_effect=asyncio.CancelledError)
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await self._runner(first)._redeliver_pending_obligations()
|
||||
assert first.send.call_args.kwargs["content"] == "the final answer"
|
||||
else:
|
||||
with dl._connect() as conn:
|
||||
conn.execute("UPDATE delivery_obligations SET attempts=1 WHERE obligation_id='ob-1'")
|
||||
_orphan("ob-1")
|
||||
second = self._adapter()
|
||||
|
||||
await self._runner(second)._redeliver_pending_obligations()
|
||||
|
||||
sent = second.send.call_args.kwargs["content"]
|
||||
assert sent == dl.RECOVERED_MARKER + "the final answer"
|
||||
assert _row("ob-1")["state"] == "delivered"
|
||||
|
||||
def test_boot_claimed_row_is_not_reclaimed_by_runtime_sweep_mid_send(self):
|
||||
"""A boot claim is exclusive while its send is in flight: the reconnect sweep must not resend it."""
|
||||
_record(platform="telegram")
|
||||
dl.mark_failed("ob-1", "send_path_degraded")
|
||||
_orphan("ob-1")
|
||||
|
||||
assert [r["obligation_id"] for r in dl.sweep_recoverable()] == ["ob-1"]
|
||||
|
||||
assert dl.sweep_failed_for_runtime("telegram") == []
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_runtime_failed_redelivery_clears_resume_before_send(self):
|
||||
from gateway.config import Platform
|
||||
|
||||
@@ -288,7 +288,8 @@ Semantics are honest at-least-once:
|
||||
|
||||
- A response whose send **never started** is redelivered as-is.
|
||||
- A response that was **mid-send** when the gateway died (the platform may or
|
||||
may not have received it) is redelivered with a visible
|
||||
may not have received it), including a redelivery an earlier boot was
|
||||
still sending, is redelivered with a visible
|
||||
"♻️ Recovered reply — … may be a duplicate" prefix. Ambiguity is labeled,
|
||||
never silently resent.
|
||||
- A final send refused by **flood control** (such as Telegram rate limits) is retried automatically
|
||||
|
||||
Reference in New Issue
Block a user