|
1 | 1 | import abc |
2 | 2 | from _typeshed import Incomplete |
3 | | -from collections.abc import Callable |
| 3 | +from collections.abc import Callable, Mapping, Sequence |
4 | 4 | from logging import Logger |
| 5 | +from typing import Final, Generic, Literal, TypeVar |
5 | 6 | from typing_extensions import Self |
6 | 7 |
|
7 | | -from ..adapters.utils.nbio_interface import AbstractIOServices, AbstractStreamProtocol |
8 | | -from ..connection import Connection, Parameters |
| 8 | +from pika.adapters.utils.connection_workflow import AbstractAMQPConnectionWorkflow, AMQPConnectorException |
| 9 | +from pika.adapters.utils.nbio_interface import AbstractIOServices, AbstractStreamProtocol, AbstractStreamTransport |
| 10 | +from pika.callback import CallbackManager |
| 11 | +from pika.channel import Channel |
| 12 | +from pika.connection import Connection, Parameters |
| 13 | +from pika.frame import Method |
| 14 | +from pika.spec import Connection as SpecConnection |
9 | 15 |
|
10 | 16 | LOGGER: Logger |
11 | 17 |
|
12 | | -class BaseConnection(Connection, metaclass=abc.ABCMeta): |
| 18 | +_IOLoop = TypeVar("_IOLoop") |
| 19 | + |
| 20 | +class BaseConnection(Connection, Generic[_IOLoop], metaclass=abc.ABCMeta): |
13 | 21 | def __init__( |
14 | 22 | self, |
15 | 23 | parameters: Parameters | None, |
16 | | - on_open_callback: Callable[[Self], object] | None, |
17 | | - on_open_error_callback: Callable[[Self, BaseException], object] | None, |
18 | | - on_close_callback: Callable[[Self, BaseException], object] | None, |
| 24 | + on_open_callback: Callable[[Connection], object] | None, |
| 25 | + on_open_error_callback: Callable[[Connection, BaseException], object] | None, |
| 26 | + on_close_callback: Callable[[Connection, BaseException], object] | None, |
19 | 27 | nbio: AbstractIOServices, |
20 | 28 | internal_connection_workflow: bool = True, |
21 | 29 | ) -> None: ... |
22 | 30 | @classmethod |
23 | 31 | @abc.abstractmethod |
24 | | - def create_connection(cls, connection_configs, on_done, custom_ioloop=None, workflow=None): ... |
| 32 | + def create_connection( |
| 33 | + cls, |
| 34 | + connection_configs: Sequence[Parameters], |
| 35 | + on_done: Callable[[Connection | AMQPConnectorException], object], |
| 36 | + custom_ioloop: _IOLoop | None = None, |
| 37 | + workflow: AbstractAMQPConnectionWorkflow | None = None, |
| 38 | + ) -> AbstractAMQPConnectionWorkflow: ... |
| 39 | + @property |
| 40 | + def ioloop(self) -> _IOLoop: ... |
| 41 | + |
| 42 | +class _StreamingProtocolShim(AbstractStreamProtocol, Generic[_IOLoop]): |
| 43 | + conn: BaseConnection[_IOLoop] |
| 44 | + def __init__(self, conn: BaseConnection[_IOLoop]) -> None: ... |
| 45 | + # These are defined as None, but on initialization are set as callable attributes |
| 46 | + def connection_made(self, transport: AbstractStreamTransport) -> None: ... |
| 47 | + def connection_lost(self, error: BaseException | None) -> None: ... |
| 48 | + def eof_received(self) -> bool: ... |
| 49 | + def data_received(self, data: bytes) -> None: ... |
| 50 | + |
| 51 | + # Next attributes are accessed via getattr() from connection.Connection class: |
| 52 | + ON_CONNECTION_CLOSED: Final = "_on_connection_closed" |
| 53 | + ON_CONNECTION_ERROR: Final = "_on_connection_error" |
| 54 | + ON_CONNECTION_OPEN_OK: Final = "_on_connection_open_ok" |
| 55 | + CONNECTION_CLOSED: Final = 0 |
| 56 | + CONNECTION_INIT: Final = 1 |
| 57 | + CONNECTION_PROTOCOL: Final = 2 |
| 58 | + CONNECTION_START: Final = 3 |
| 59 | + CONNECTION_TUNE: Final = 4 |
| 60 | + CONNECTION_OPEN: Final = 5 |
| 61 | + CONNECTION_CLOSING: Final = 6 |
| 62 | + connection_state: Literal[0, 1, 2, 3, 4, 5, 6] # one of the constants above |
| 63 | + params: Parameters |
| 64 | + callbacks: CallbackManager |
| 65 | + server_capabilities: Mapping[str, bool] | None |
| 66 | + server_properties: Mapping[str, Incomplete] | None |
| 67 | + known_hosts: str | None |
| 68 | + def add_on_close_callback(self, callback: Callable[[Self, BaseException], object]) -> None: ... |
| 69 | + def add_on_connection_blocked_callback(self, callback: Callable[[Self, Method[SpecConnection.Blocked]], object]) -> None: ... |
| 70 | + def add_on_connection_unblocked_callback( |
| 71 | + self, callback: Callable[[Self, Method[SpecConnection.Unblocked]], object] |
| 72 | + ) -> None: ... |
| 73 | + def add_on_open_callback(self, callback: Callable[[Self], object]) -> None: ... |
| 74 | + def add_on_open_error_callback( |
| 75 | + self, callback: Callable[[Self, BaseException], object], remove_default: bool = True |
| 76 | + ) -> None: ... |
| 77 | + def channel( |
| 78 | + self, channel_number: int | None = None, on_open_callback: Callable[[Channel], object] | None = None |
| 79 | + ) -> Channel: ... |
| 80 | + def update_secret( |
| 81 | + self, |
| 82 | + new_secret: str | bytes, |
| 83 | + reason: str | bytes, |
| 84 | + callback: Callable[[Method[SpecConnection.UpdateSecretOk]], object] | None = None, |
| 85 | + ) -> None: ... |
| 86 | + def close(self, reply_code: int = 200, reply_text: str = "Normal shutdown") -> None: ... |
25 | 87 | @property |
26 | | - def ioloop(self): ... |
| 88 | + def is_closed(self) -> bool: ... |
| 89 | + @property |
| 90 | + def is_closing(self) -> bool: ... |
| 91 | + @property |
| 92 | + def is_open(self) -> bool: ... |
| 93 | + @property |
| 94 | + def basic_nack(self) -> bool: ... |
| 95 | + @property |
| 96 | + def consumer_cancel_notify(self) -> bool: ... |
| 97 | + @property |
| 98 | + def exchange_exchange_bindings(self) -> bool: ... |
| 99 | + @property |
| 100 | + def publisher_confirms(self) -> bool: ... |
27 | 101 |
|
28 | | -class _StreamingProtocolShim(AbstractStreamProtocol): |
29 | | - connection_made: Incomplete |
30 | | - connection_lost: Incomplete |
31 | | - eof_received: Incomplete |
32 | | - data_received: Incomplete |
33 | | - conn: Incomplete |
34 | | - def __init__(self, conn) -> None: ... |
35 | | - def __getattr__(self, attr: str): ... |
| 102 | + # Next attributes are accessed via getattr() from BaseConnection class: |
| 103 | + @classmethod |
| 104 | + def create_connection( |
| 105 | + cls, |
| 106 | + connection_configs: Sequence[Parameters], |
| 107 | + on_done: Callable[[Connection | AMQPConnectorException], object], |
| 108 | + custom_ioloop: _IOLoop | None = None, |
| 109 | + workflow: AbstractAMQPConnectionWorkflow | None = None, |
| 110 | + ) -> AbstractAMQPConnectionWorkflow: ... |
| 111 | + @property |
| 112 | + def ioloop(self) -> _IOLoop: ... |
0 commit comments