Skip to content

StreamVisitor

faststream.redis.testing.StreamVisitor #

Bases: Visitor

visit #

visit(
    *,
    sub: LogicSubscriber,
    channel: str | None = None,
    list: str | None = None,
    stream: str | None = None,
) -> str | None
Source code in faststream/redis/testing.py
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
@override
def visit(
    self,
    *,
    sub: "LogicSubscriber",
    channel: str | None = None,
    list: str | None = None,
    stream: str | None = None,
) -> str | None:
    if stream is None or not isinstance(sub, _StreamHandlerMixin):
        return None

    if stream == sub.stream_sub.name:
        return stream

    return None

get_message #

get_message(
    channel: str, body: Any, sub: _StreamHandlerMixin
) -> Any
Source code in faststream/redis/testing.py
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
@override
def get_message(  # type: ignore[override]
    self,
    channel: str,
    body: Any,
    sub: "_StreamHandlerMixin",
) -> Any:
    message: BatchStreamMessage | DefaultStreamMessage
    if sub.stream_sub.batch:
        message = BatchStreamMessage(
            type="bstream",
            channel=channel,
            data=[{bDATA_KEY: body}],
            message_ids=[],
        )
    else:
        message = DefaultStreamMessage(
            type="stream",
            channel=channel,
            data={bDATA_KEY: body},
            message_ids=[],
        )

    if sub.stream_sub.claim_min_idle_time is not None:
        # The in-memory broker has no claiming: every delivery is a new
        # message, so expose the new-message metadata values
        message["idle_times"] = [0]
        message["delivery_counts"] = [0]

    return message