Skip to content

FakeProducer

faststream.redis.testing.FakeProducer #

FakeProducer(
    broker: RedisBroker,
    brokers: Sequence[RedisBroker],
    config: ParserConfig,
    pel: PEL | None = None,
)

Bases: RedisFastProducer

Source code in faststream/redis/testing.py
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
def __init__(
    self,
    broker: RedisBroker,
    brokers: Sequence[RedisBroker],
    config: ParserConfig,
    pel: PEL | None = None,
) -> None:
    self.broker = broker
    self.brokers = brokers
    self._fake_config = config

    default = RedisPubSubParser(config)

    self._parser = ParserComposition(
        broker._parser,
        default.parse_message,
    )
    self._decoder = ParserComposition(
        broker._decoder,
        default.decode_message,
    )
    self.codec = broker.config.broker_codec or DefaultCodec()
    self.pel = pel or PEL()

broker instance-attribute #

broker = broker

brokers instance-attribute #

brokers = brokers

codec instance-attribute #

codec = broker.config.broker_codec or DefaultCodec()

pel instance-attribute #

pel = pel or PEL()

subscribers property #

subscribers: Iterable[LogicSubscriber]

serializer instance-attribute #

serializer = serializer

publish async #

publish(cmd: RedisPublishCommand) -> int | bytes
Source code in faststream/redis/testing.py
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
@override
async def publish(self, cmd: "RedisPublishCommand") -> int | bytes:
    body = await build_message(
        message=cmd.body,
        reply_to=cmd.reply_to,
        correlation_id=cmd.correlation_id or self.broker.config.id_generator(),
        headers=cmd.headers,
        message_format=cmd.message_format,
        serializer=self.broker.config.fd_config._serializer,
        codec=self.codec,
    )
    destination = _make_destination_kwargs(cmd)
    visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor())
    session_id = uuid.uuid4()
    for visitor, visited_ch, handler in self._find_handlers(
        destination=destination,
        visitors=visitors,
        cmd=cmd,
        session_id=session_id,
    ):
        msg = visitor.get_message(
            visited_ch,
            body,
            handler,
        )
        self._put_pel(msg=msg, cmd=cmd, handler=handler, session_id=session_id)
        await self._execute_handler(msg, handler, session_id=session_id)

    return 0

request async #

Source code in faststream/redis/testing.py
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
@override
async def request(self, cmd: "RedisPublishCommand") -> "PubSubMessage":
    body = await build_message(
        message=cmd.body,
        correlation_id=cmd.correlation_id or self.broker.config.id_generator(),
        headers=cmd.headers,
        message_format=cmd.message_format,
        serializer=self.broker.config.fd_config._serializer,
        codec=self.codec,
    )

    destination = _make_destination_kwargs(cmd)
    visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor())
    session_id = uuid.uuid4()

    for visitor, visited_ch, handler in self._find_handlers(
        destination=destination,
        visitors=visitors,
        cmd=cmd,
        session_id=session_id,
    ):
        msg = visitor.get_message(
            visited_ch,
            body,
            handler,
        )
        self._put_pel(msg=msg, cmd=cmd, handler=handler, session_id=session_id)
        with anyio.fail_after(cmd.timeout):
            return await self._execute_handler(msg, handler, session_id=session_id)

    raise SubscriberNotFound

publish_batch async #

publish_batch(cmd: RedisPublishCommand) -> int
Source code in faststream/redis/testing.py
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
@override
async def publish_batch(self, cmd: "RedisPublishCommand") -> int:
    data_to_send = [
        await build_message(
            m,
            correlation_id=cmd.correlation_id or self.broker.config.id_generator(),
            headers=cmd.headers,
            message_format=cmd.message_format,
            serializer=self.broker.config.fd_config._serializer,
            codec=self.codec,
        )
        for m in cmd.batch_bodies
    ]
    session_id = uuid.uuid4()

    for visitor, visited_ch, handler in self._find_handlers(
        {"list": cmd.destination},
        (ListVisitor(),),
        cmd=cmd,
        session_id=session_id,
    ):
        casted_handler = cast("_ListHandlerMixin", handler)

        if casted_handler.list_sub.batch:
            msg = visitor.get_message(
                visited_ch,
                data_to_send,
                casted_handler,
            )
            self._put_pel(msg=msg, cmd=cmd, handler=handler, session_id=session_id)
            await self._execute_handler(msg, handler, session_id=session_id)

    return 0

connect #

connect(
    serializer: Optional[SerializerProto] = None,
    codec: Optional[CodecProto] = None,
) -> None
Source code in faststream/redis/publisher/producer.py
83
84
85
86
87
88
89
90
def connect(
    self,
    serializer: Optional["SerializerProto"] = None,
    codec: Optional["CodecProto"] = None,
) -> None:
    self.serializer = serializer
    if codec is not None:
        self.codec = codec