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/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 3a2730b..e9accd1 100644 --- a/tests/test_envelope.py +++ b/tests/test_envelope.py @@ -1,9 +1,11 @@ import asyncio import json +import logging import os import time import uuid +import pytest from faststream import Context from redis.asyncio import Redis @@ -108,8 +110,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 +178,37 @@ 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_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" + + 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] = [] + + @broker.subscriber("topic") + async def handler(body: bytes) -> None: # pragma: no cover - never invoked; the parser rejects the timer + seen.append(body) + + 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") 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()