Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions faststream_redis_timers/envelope.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
Expand Down
16 changes: 14 additions & 2 deletions faststream_redis_timers/parser/parser.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import logging
import typing
from functools import partial

Expand All @@ -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,
Expand Down
42 changes: 40 additions & 2 deletions tests/test_envelope.py
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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()
Loading