Skip to content

ConcurrentMixin

faststream.broker.subscriber.mixins.ConcurrentMixin #

ConcurrentMixin(*args, max_workers, **kwargs)

Bases: TasksMixin, Generic[MsgType]

Source code in faststream/broker/subscriber/mixins.py
def __init__(
    self,
    *args: Any,
    max_workers: int,
    **kwargs: Any,
) -> None:
    self.max_workers = max_workers

    self.send_stream, self.receive_stream = anyio.create_memory_object_stream(
        max_buffer_size=max_workers
    )
    self.limiter = anyio.Semaphore(max_workers)

    super().__init__(*args, **kwargs)

title_ instance-attribute #

title_ = title_

description_ instance-attribute #

description_ = description_

include_in_schema instance-attribute #

include_in_schema = include_in_schema

name property #

name

Returns the name of the API operation.

description property #

description

Returns the description of the API operation.

calls instance-attribute #

calls = []

running instance-attribute #

running = False

call_name property #

call_name

Returns the name of the handler call.

lock instance-attribute #

extra_watcher_options instance-attribute #

extra_watcher_options = {}

extra_context instance-attribute #

extra_context = {}

graceful_timeout instance-attribute #

graceful_timeout = None

tasks instance-attribute #

tasks = []

send_stream instance-attribute #

send_stream

receive_stream instance-attribute #

receive_stream

max_workers instance-attribute #

max_workers = max_workers

limiter instance-attribute #

limiter = Semaphore(max_workers)

setup #

setup(*, logger, producer, graceful_timeout, extra_context, broker_parser, broker_decoder, apply_types, is_validate, _get_dependant, _call_decorators)
Source code in faststream/broker/subscriber/usecase.py
@override
def setup(  # type: ignore[override]
    self,
    *,
    logger: Optional["LoggerProto"],
    producer: Optional["ProducerProto"],
    graceful_timeout: Optional[float],
    extra_context: "AnyDict",
    # broker options
    broker_parser: Optional["CustomCallable"],
    broker_decoder: Optional["CustomCallable"],
    # dependant args
    apply_types: bool,
    is_validate: bool,
    _get_dependant: Optional[Callable[..., Any]],
    _call_decorators: Iterable["Decorator"],
) -> None:
    self.lock = MultiLock()

    self._producer = producer
    self.graceful_timeout = graceful_timeout
    self.extra_context = extra_context

    self.watcher = get_watcher_context(logger, self._no_ack, self._retry)

    for call in self.calls:
        if parser := call.item_parser or broker_parser:
            async_parser = resolve_custom_func(to_async(parser), self._parser)
        else:
            async_parser = self._parser

        if decoder := call.item_decoder or broker_decoder:
            async_decoder = resolve_custom_func(to_async(decoder), self._decoder)
        else:
            async_decoder = self._decoder

        self._parser = async_parser
        self._decoder = async_decoder

        call.setup(
            parser=async_parser,
            decoder=async_decoder,
            apply_types=apply_types,
            is_validate=is_validate,
            _get_dependant=_get_dependant,
            _call_decorators=(*self._call_decorators, *_call_decorators),
            broker_dependencies=self._broker_dependencies,
        )

        call.handler.refresh(with_mock=False)

add_prefix abstractmethod #

add_prefix(prefix)
Source code in faststream/broker/proto.py
@abstractmethod
def add_prefix(self, prefix: str) -> None: ...

schema #

schema()

Returns the schema of the API operation as a dictionary of channel names and channel objects.

Source code in faststream/asyncapi/abc.py
def schema(self) -> Dict[str, Channel]:
    """Returns the schema of the API operation as a dictionary of channel names and channel objects."""
    if self.include_in_schema:
        return self.get_schema()
    else:
        return {}

add_middleware #

add_middleware(middleware)
Source code in faststream/broker/subscriber/usecase.py
def add_middleware(self, middleware: "BrokerMiddleware[MsgType]") -> None:
    self._broker_middlewares = (*self._broker_middlewares, middleware)

get_log_context #

get_log_context(message)

Generate log context.

Source code in faststream/broker/subscriber/usecase.py
def get_log_context(
    self,
    message: Optional["StreamMessage[MsgType]"],
) -> Dict[str, str]:
    """Generate log context."""
    return {
        "message_id": getattr(message, "message_id", ""),
    }

start abstractmethod async #

start()

Start the handler.

Source code in faststream/broker/subscriber/usecase.py
@abstractmethod
async def start(self) -> None:
    """Start the handler."""
    self.running = True

close async #

close()

Clean up handler subscription, cancel consume task in graceful mode.

Source code in faststream/broker/subscriber/mixins.py
async def close(self) -> None:
    """Clean up handler subscription, cancel consume task in graceful mode."""
    await super().close()

    for task in self.tasks:
        if not task.done():
            task.cancel()

    self.tasks = []

consume async #

consume(msg)

Consume a message asynchronously.

Source code in faststream/broker/subscriber/usecase.py
async def consume(self, msg: MsgType) -> Any:
    """Consume a message asynchronously."""
    if not self.running:
        return None

    try:
        return await self.process_message(msg)

    except StopConsume:
        # Stop handler at StopConsume exception
        await self.close()

    except SystemExit:
        # Stop handler at `exit()` call
        await self.close()

        if app := context.get("app"):
            app.exit()

    except Exception:  # nosec B110
        # All other exceptions were logged by CriticalLogMiddleware
        pass

process_message async #

process_message(msg)

Execute all message processing stages.

Source code in faststream/broker/subscriber/usecase.py
async def process_message(self, msg: MsgType) -> "Response":
    """Execute all message processing stages."""
    async with AsyncExitStack() as stack:
        stack.enter_context(self.lock)

        # Enter context before middlewares
        for k, v in self.extra_context.items():
            stack.enter_context(context.scope(k, v))

        stack.enter_context(context.scope("handler_", self))

        # enter all middlewares
        middlewares: List[BaseMiddleware] = []
        for base_m in self._broker_middlewares:
            middleware = base_m(msg)
            middlewares.append(middleware)
            await middleware.__aenter__()

        cache: Dict[Any, Any] = {}
        parsing_error: Optional[Exception] = None
        for h in self.calls:
            try:
                message = await h.is_suitable(msg, cache)
            except Exception as e:
                parsing_error = e
                break

            if message is not None:
                # Acknowledgement scope
                # TODO: move it to scope enter at `retry` option deprecation
                await stack.enter_async_context(
                    self.watcher(
                        message,
                        **self.extra_watcher_options,
                    )
                )

                stack.enter_context(
                    context.scope("log_context", self.get_log_context(message))
                )
                stack.enter_context(context.scope("message", message))

                # Middlewares should be exited before scope release
                for m in middlewares:
                    stack.push_async_exit(m.__aexit__)

                result_msg = ensure_response(
                    await h.call(
                        message=message,
                        # consumer middlewares
                        _extra_middlewares=(
                            m.consume_scope for m in middlewares[::-1]
                        ),
                    )
                )

                if not result_msg.correlation_id:
                    result_msg.correlation_id = message.correlation_id

                for p in chain(
                    self.__get_response_publisher(message),
                    h.handler._publishers,
                ):
                    await p.publish(
                        result_msg.body,
                        **result_msg.as_publish_kwargs(),
                        # publisher middlewares
                        _extra_middlewares=[
                            m.publish_scope for m in middlewares[::-1]
                        ],
                    )

                # Return data for tests
                return result_msg

        # Suitable handler was not found or
        # parsing/decoding exception occurred
        for m in middlewares:
            stack.push_async_exit(m.__aexit__)

        if parsing_error:
            raise parsing_error

        else:
            raise SubscriberNotFound(f"There is no suitable handler for {msg=}")

    # An error was raised and processed by some middleware
    return ensure_response(None)

get_one abstractmethod async #

get_one(*, timeout=5.0)
Source code in faststream/broker/subscriber/proto.py
@abstractmethod
async def get_one(
    self, *, timeout: float = 5.0
) -> "Optional[StreamMessage[MsgType]]": ...

add_call #

add_call(*, filter_, parser_, decoder_, middlewares_, dependencies_)
Source code in faststream/broker/subscriber/usecase.py
def add_call(
    self,
    *,
    filter_: "Filter[Any]",
    parser_: Optional["CustomCallable"],
    decoder_: Optional["CustomCallable"],
    middlewares_: Sequence["SubscriberMiddleware[Any]"],
    dependencies_: Iterable["Depends"],
) -> Self:
    self._call_options = _CallOptions(
        filter=filter_,
        parser=parser_,
        decoder=decoder_,
        middlewares=middlewares_,
        dependencies=dependencies_,
    )
    return self

get_name abstractmethod #

get_name()

Name property fallback.

Source code in faststream/asyncapi/abc.py
@abstractmethod
def get_name(self) -> str:
    """Name property fallback."""
    raise NotImplementedError()

get_description #

get_description()

Returns the description of the handler.

Source code in faststream/broker/subscriber/usecase.py
def get_description(self) -> Optional[str]:
    """Returns the description of the handler."""
    if not self.calls:  # pragma: no cover
        return None

    else:
        return self.calls[0].description

get_schema abstractmethod #

get_schema()

Generate AsyncAPI schema.

Source code in faststream/asyncapi/abc.py
@abstractmethod
def get_schema(self) -> Dict[str, Channel]:
    """Generate AsyncAPI schema."""
    raise NotImplementedError()

get_payloads #

get_payloads()

Get the payloads of the handler.

Source code in faststream/broker/subscriber/usecase.py
def get_payloads(self) -> List[Tuple["AnyDict", str]]:
    """Get the payloads of the handler."""
    payloads: List[Tuple[AnyDict, str]] = []

    for h in self.calls:
        if h.dependant is None:
            raise SetupError("You should setup `Handler` at first.")

        body = parse_handler_params(
            h.dependant,
            prefix=f"{self.title_ or self.call_name}:Message",
        )

        payloads.append((body, to_camelcase(h.call_name)))

    if not self.calls:
        payloads.append(
            (
                {
                    "title": f"{self.title_ or self.call_name}:Message:Payload",
                },
                to_camelcase(self.call_name),
            )
        )

    return payloads

add_task #

add_task(coro)
Source code in faststream/broker/subscriber/mixins.py
def add_task(self, coro: Coroutine[Any, Any, Any]) -> None:
    self.tasks.append(asyncio.create_task(coro))

start_consume_task #

start_consume_task()
Source code in faststream/broker/subscriber/mixins.py
def start_consume_task(self) -> None:
    self.add_task(self._serve_consume_queue())