From 2ba1e08ff5d0a942b769887908e2766da7dd698c Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sun, 27 Sep 2026 19:45:24 +0300 Subject: [PATCH 1/2] fix: reject a legacy timer envelope whose body is not hex --- faststream_redis_timers/envelope.py | 5 +++-- tests/test_envelope.py | 35 +++++++++++++++++++++++++++-- 2 files changed, 36 insertions(+), 4 deletions(-) diff --git a/faststream_redis_timers/envelope.py b/faststream_redis_timers/envelope.py index 41cd14d..3ac5563 100644 --- a/faststream_redis_timers/envelope.py +++ b/faststream_redis_timers/envelope.py @@ -22,8 +22,9 @@ def parse(cls, data: bytes) -> tuple[bytes, dict[str, Any]]: if isinstance(d, dict) and "b" in d: try: body = bytes.fromhex(d["b"]) - except (TypeError, ValueError): - body = b"" + except (TypeError, ValueError) as e: + msg = "legacy timer envelope body is not hex" + raise ValueError(msg) from e headers: dict[str, Any] = {} if d.get("ct"): headers["content-type"] = d["ct"] diff --git a/tests/test_envelope.py b/tests/test_envelope.py index 3a2730b..c773834 100644 --- a/tests/test_envelope.py +++ b/tests/test_envelope.py @@ -4,6 +4,7 @@ import time import uuid +import pytest from faststream import Context from redis.asyncio import Redis @@ -108,8 +109,10 @@ def test_non_json_payload_starting_with_brace_parses_as_raw_body() -> None: assert TimerMessageFormat.parse(b"{not json") == (b"{not json", {}) -def test_legacy_envelope_with_non_hex_body_parses_as_empty_body() -> None: - assert TimerMessageFormat.parse(b'{"b": "zz"}') == (b"", {}) +@pytest.mark.parametrize("envelope", [b'{"b": "zz"}', b'{"b": 12}', b'{"b": null}']) +def test_legacy_envelope_with_non_hex_body_is_rejected(envelope: bytes) -> None: + with pytest.raises(ValueError, match="legacy timer envelope"): + TimerMessageFormat.parse(envelope) async def test_works_with_decode_responses_true() -> None: @@ -174,3 +177,31 @@ async def handler(body: str) -> None: await asyncio.wait_for(event.wait(), timeout=5.0) assert seen == ["hello"] + + +async def test_legacy_envelope_with_non_hex_body_never_reaches_handler(redis_client: Redis) -> None: + suffix = uuid.uuid4().hex + broker = TimersBroker( + redis_client, + timeline_key=f"corrupt_tl_{suffix}", + payloads_key=f"corrupt_pl_{suffix}", + ) + timeline_key = f"corrupt_tl_{suffix}:topic" + payloads_key = f"corrupt_pl_{suffix}:topic" + + due_at = time.time() - 1 + await redis_client.zadd(timeline_key, {"old-timer": due_at}) + await redis_client.hset(payloads_key, "old-timer", b'{"b": "zz"}') + + seen: list[bytes] = [] + + @broker.subscriber("topic") + async def handler(body: bytes) -> None: # pragma: no cover - never invoked; the parser rejects the timer + seen.append(body) + + async with broker: + await asyncio.sleep(0.3) + + assert seen == [] + assert (await redis_client.zscore(timeline_key, "old-timer") or 0) > due_at + assert await redis_client.hexists(payloads_key, "old-timer") From 9314f51434dfd6d01bb6308f8b9bbd02906a0056 Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sun, 27 Sep 2026 20:10:03 +0300 Subject: [PATCH 2/2] fix: remove a timer whose payload cannot be parsed instead of retrying it forever --- faststream_redis_timers/parser/parser.py | 16 ++++++++++++++-- tests/test_envelope.py | 21 ++++++++++++++------- 2 files changed, 28 insertions(+), 9 deletions(-) diff --git a/faststream_redis_timers/parser/parser.py b/faststream_redis_timers/parser/parser.py index 07f8ccd..0a6064e 100644 --- a/faststream_redis_timers/parser/parser.py +++ b/faststream_redis_timers/parser/parser.py @@ -1,3 +1,4 @@ +import logging import typing from functools import partial @@ -17,9 +18,20 @@ def __init__(self, config: "TimersSubscriberConfig") -> None: self._config = config async def parse_message(self, msg: "TimerMessage") -> TimerStreamMessage: - body, headers = TimerMessageFormat.parse(msg["data"]) timer_id = msg["timer_id"] - store = self._config._outer_config.store # noqa: SLF001 + outer_config = self._config._outer_config # noqa: SLF001 + store = outer_config.store + try: + body, headers = TimerMessageFormat.parse(msg["data"]) + except ValueError as e: + await store.remove(self._config.full_topic, timer_id) + outer_config.logger.log( + f"Timer {timer_id!r} on {self._config.full_topic!r} removed: " + f"its {len(msg['data'])}-byte payload cannot be parsed", + logging.ERROR, + exc_info=e, + ) + raise return TimerStreamMessage( raw_message=msg, body=body, diff --git a/tests/test_envelope.py b/tests/test_envelope.py index c773834..e9accd1 100644 --- a/tests/test_envelope.py +++ b/tests/test_envelope.py @@ -1,5 +1,6 @@ import asyncio import json +import logging import os import time import uuid @@ -179,18 +180,20 @@ async def handler(body: str) -> None: assert seen == ["hello"] -async def test_legacy_envelope_with_non_hex_body_never_reaches_handler(redis_client: Redis) -> None: +async def test_legacy_envelope_with_non_hex_body_is_removed_and_logged_once( + redis_client: Redis, caplog: pytest.LogCaptureFixture +) -> None: suffix = uuid.uuid4().hex broker = TimersBroker( redis_client, timeline_key=f"corrupt_tl_{suffix}", payloads_key=f"corrupt_pl_{suffix}", + logger=logging.getLogger(f"corrupt-{suffix}"), ) timeline_key = f"corrupt_tl_{suffix}:topic" payloads_key = f"corrupt_pl_{suffix}:topic" - due_at = time.time() - 1 - await redis_client.zadd(timeline_key, {"old-timer": due_at}) + await redis_client.zadd(timeline_key, {"old-timer": time.time() - 1}) await redis_client.hset(payloads_key, "old-timer", b'{"b": "zz"}') seen: list[bytes] = [] @@ -199,9 +202,13 @@ async def test_legacy_envelope_with_non_hex_body_never_reaches_handler(redis_cli async def handler(body: bytes) -> None: # pragma: no cover - never invoked; the parser rejects the timer seen.append(body) - async with broker: - await asyncio.sleep(0.3) + with caplog.at_level(logging.ERROR, logger=f"corrupt-{suffix}"): + async with broker: + await asyncio.sleep(0.3) assert seen == [] - assert (await redis_client.zscore(timeline_key, "old-timer") or 0) > due_at - assert await redis_client.hexists(payloads_key, "old-timer") + assert await redis_client.zscore(timeline_key, "old-timer") is None + assert not await redis_client.hexists(payloads_key, "old-timer") + removals = [record for record in caplog.records if "removed" in record.getMessage()] + assert len(removals) == 1 + assert "'old-timer'" in removals[0].getMessage()