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
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
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