Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions stubs/pika/@tests/stubtest_allowlist.txt
Original file line number Diff line number Diff line change
Expand Up @@ -12,3 +12,29 @@ pika.connection.ConnectionParameters.T

# Arguments have a sentinel default, which is not reflected in the stubs.
pika.connection.ConnectionParameters.__init__

# These are defined as None, but on initialization are set as callable attributes.
pika.adapters.base_connection._StreamingProtocolShim.connection_made
pika.adapters.base_connection._StreamingProtocolShim.connection_lost
pika.adapters.base_connection._StreamingProtocolShim.eof_received
pika.adapters.base_connection._StreamingProtocolShim.data_received

# The following methods are not defined directly on this class;
# they are resolved via __getattr__.
pika.adapters.base_connection._StreamingProtocolShim.add_on_close_callback
pika.adapters.base_connection._StreamingProtocolShim.add_on_connection_blocked_callback
pika.adapters.base_connection._StreamingProtocolShim.add_on_connection_unblocked_callback
pika.adapters.base_connection._StreamingProtocolShim.add_on_open_callback
pika.adapters.base_connection._StreamingProtocolShim.add_on_open_error_callback
pika.adapters.base_connection._StreamingProtocolShim.channel
pika.adapters.base_connection._StreamingProtocolShim.update_secret
pika.adapters.base_connection._StreamingProtocolShim.close
pika.adapters.base_connection._StreamingProtocolShim.is_closed
pika.adapters.base_connection._StreamingProtocolShim.is_closing
pika.adapters.base_connection._StreamingProtocolShim.is_open
pika.adapters.base_connection._StreamingProtocolShim.basic_nack
pika.adapters.base_connection._StreamingProtocolShim.consumer_cancel_notify
pika.adapters.base_connection._StreamingProtocolShim.exchange_exchange_bindings
pika.adapters.base_connection._StreamingProtocolShim.publisher_confirms
pika.adapters.base_connection._StreamingProtocolShim.create_connection
pika.adapters.base_connection._StreamingProtocolShim.ioloop
1 change: 1 addition & 0 deletions stubs/pika/METADATA.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ extra-description = """\
The `types-pika` package contains alternate, more complete type stubs, that \
are maintained outside of typeshed.\
"""
optional-dependencies = ["types-gevent"]

[tool.stubtest]
stubtest-dependencies = ["gevent", "tornado", "twisted"]
44 changes: 26 additions & 18 deletions stubs/pika/pika/adapters/asyncio_connection.pyi
Original file line number Diff line number Diff line change
@@ -1,33 +1,37 @@
import asyncio
from _typeshed import Incomplete
from asyncio import AbstractEventLoop, Future, Handle
from collections.abc import Callable, Sequence
from logging import Logger
from typing_extensions import Self

from ..connection import Connection, Parameters
from .base_connection import BaseConnection
from .utils import io_services_utils
from .utils.connection_workflow import AbstractAMQPConnectionWorkflow, AMQPConnectorException
from .utils.nbio_interface import AbstractFileDescriptorServices, AbstractIOReference, AbstractIOServices, AbstractTimerReference
from pika.adapters.base_connection import BaseConnection
from pika.adapters.utils import io_services_utils
from pika.adapters.utils.connection_workflow import AbstractAMQPConnectionWorkflow, AMQPConnectorException
from pika.adapters.utils.nbio_interface import (
AbstractFileDescriptorServices,
AbstractIOReference,
AbstractIOServices,
AbstractTimerReference,
)
from pika.connection import Connection, Parameters

LOGGER: Logger

class AsyncioConnection(BaseConnection):
class AsyncioConnection(BaseConnection[asyncio.AbstractEventLoop]):
def __init__(
self,
parameters: Parameters | None = None,
on_open_callback: Callable[[Self], object] | None = None,
on_open_error_callback: Callable[[Self, BaseException], object] | None = None,
on_close_callback: Callable[[Self, BaseException], object] | None = None,
custom_ioloop: AbstractEventLoop | None = None,
on_open_callback: Callable[[Connection], object] | None = None,
on_open_error_callback: Callable[[Connection, BaseException], object] | None = None,
on_close_callback: Callable[[Connection, BaseException], object] | None = None,
custom_ioloop: asyncio.AbstractEventLoop | AbstractIOServices | None = None,
internal_connection_workflow: bool = True,
) -> None: ...
@classmethod
def create_connection(
cls,
connection_configs: Sequence[Parameters],
on_done: Callable[[Connection | AMQPConnectorException], object],
custom_ioloop: AbstractEventLoop | None = None,
custom_ioloop: asyncio.AbstractEventLoop | None = None,
workflow: AbstractAMQPConnectionWorkflow | None = None,
) -> AbstractAMQPConnectionWorkflow: ...

Expand All @@ -37,8 +41,8 @@ class _AsyncioIOServicesAdapter(
AbstractIOServices,
AbstractFileDescriptorServices,
):
def __init__(self, loop: AbstractEventLoop | None = None) -> None: ...
def get_native_ioloop(self) -> AbstractEventLoop: ...
def __init__(self, loop: asyncio.AbstractEventLoop | None = None) -> None: ...
def get_native_ioloop(self) -> asyncio.AbstractEventLoop: ...
def close(self) -> None: ...
def run(self) -> None: ...
def stop(self) -> None: ...
Expand All @@ -48,7 +52,7 @@ class _AsyncioIOServicesAdapter(
self,
host: str | bytes | None,
port: str | bytes | int | None,
on_done: Callable[[BaseConnection | BaseException], object], # type: ignore[override]
on_done: Callable[[BaseConnection[asyncio.AbstractEventLoop] | BaseException], object], # type: ignore[override]
family: int = 0,
socktype: int = 0,
proto: int = 0,
Expand All @@ -60,9 +64,13 @@ class _AsyncioIOServicesAdapter(
def remove_writer(self, fd: int) -> bool: ...

class _TimerHandle(AbstractTimerReference):
def __init__(self, handle: Handle) -> None: ...
def __init__(self, handle: asyncio.Handle) -> None: ...
def cancel(self) -> None: ...

class _AsyncioIOReference(AbstractIOReference):
def __init__(self, future: Future[Incomplete], on_done: Callable[[BaseConnection | BaseException], object]) -> None: ...
def __init__(
self,
future: asyncio.Future[Incomplete],
on_done: Callable[[BaseConnection[asyncio.AbstractEventLoop] | BaseException], object],
) -> None: ...
def cancel(self) -> bool: ...
111 changes: 94 additions & 17 deletions stubs/pika/pika/adapters/base_connection.pyi
Original file line number Diff line number Diff line change
@@ -1,35 +1,112 @@
import abc
from _typeshed import Incomplete
from collections.abc import Callable
from collections.abc import Callable, Mapping, Sequence
from logging import Logger
from typing import Final, Generic, Literal, TypeVar
from typing_extensions import Self

from ..adapters.utils.nbio_interface import AbstractIOServices, AbstractStreamProtocol
from ..connection import Connection, Parameters
from pika.adapters.utils.connection_workflow import AbstractAMQPConnectionWorkflow, AMQPConnectorException
from pika.adapters.utils.nbio_interface import AbstractIOServices, AbstractStreamProtocol, AbstractStreamTransport
from pika.callback import CallbackManager
from pika.channel import Channel
from pika.connection import Connection, Parameters
from pika.frame import Method
from pika.spec import Connection as SpecConnection

LOGGER: Logger

class BaseConnection(Connection, metaclass=abc.ABCMeta):
_IOLoop = TypeVar("_IOLoop")

class BaseConnection(Connection, Generic[_IOLoop], metaclass=abc.ABCMeta):
def __init__(
self,
parameters: Parameters | None,
on_open_callback: Callable[[Self], object] | None,
on_open_error_callback: Callable[[Self, BaseException], object] | None,
on_close_callback: Callable[[Self, BaseException], object] | None,
on_open_callback: Callable[[Connection], object] | None,
on_open_error_callback: Callable[[Connection, BaseException], object] | None,
on_close_callback: Callable[[Connection, BaseException], object] | None,
nbio: AbstractIOServices,
internal_connection_workflow: bool = True,
) -> None: ...
@classmethod
@abc.abstractmethod
def create_connection(cls, connection_configs, on_done, custom_ioloop=None, workflow=None): ...
def create_connection(
cls,
connection_configs: Sequence[Parameters],
on_done: Callable[[Connection | AMQPConnectorException], object],
custom_ioloop: _IOLoop | None = None,
workflow: AbstractAMQPConnectionWorkflow | None = None,
) -> AbstractAMQPConnectionWorkflow: ...
@property
def ioloop(self) -> _IOLoop: ...

class _StreamingProtocolShim(AbstractStreamProtocol, Generic[_IOLoop]):
conn: BaseConnection[_IOLoop]
def __init__(self, conn: BaseConnection[_IOLoop]) -> None: ...
# These are defined as None, but on initialization are set as callable attributes
def connection_made(self, transport: AbstractStreamTransport) -> None: ...
def connection_lost(self, error: BaseException | None) -> None: ...
def eof_received(self) -> bool: ...
def data_received(self, data: bytes) -> None: ...

# Next attributes are accessed via getattr() from connection.Connection class:
ON_CONNECTION_CLOSED: Final = "_on_connection_closed"
ON_CONNECTION_ERROR: Final = "_on_connection_error"
ON_CONNECTION_OPEN_OK: Final = "_on_connection_open_ok"
CONNECTION_CLOSED: Final = 0
CONNECTION_INIT: Final = 1
CONNECTION_PROTOCOL: Final = 2
CONNECTION_START: Final = 3
CONNECTION_TUNE: Final = 4
CONNECTION_OPEN: Final = 5
CONNECTION_CLOSING: Final = 6
connection_state: Literal[0, 1, 2, 3, 4, 5, 6] # one of the constants above
params: Parameters
callbacks: CallbackManager
server_capabilities: Mapping[str, bool] | None
server_properties: Mapping[str, Incomplete] | None
known_hosts: str | None
def add_on_close_callback(self, callback: Callable[[Self, BaseException], object]) -> None: ...
def add_on_connection_blocked_callback(self, callback: Callable[[Self, Method[SpecConnection.Blocked]], object]) -> None: ...
def add_on_connection_unblocked_callback(
self, callback: Callable[[Self, Method[SpecConnection.Unblocked]], object]
) -> None: ...
def add_on_open_callback(self, callback: Callable[[Self], object]) -> None: ...
def add_on_open_error_callback(
self, callback: Callable[[Self, BaseException], object], remove_default: bool = True
) -> None: ...
def channel(
self, channel_number: int | None = None, on_open_callback: Callable[[Channel], object] | None = None
) -> Channel: ...
def update_secret(
self,
new_secret: str | bytes,
reason: str | bytes,
callback: Callable[[Method[SpecConnection.UpdateSecretOk]], object] | None = None,
) -> None: ...
def close(self, reply_code: int = 200, reply_text: str = "Normal shutdown") -> None: ...
@property
def ioloop(self): ...
def is_closed(self) -> bool: ...
@property
def is_closing(self) -> bool: ...
@property
def is_open(self) -> bool: ...
@property
def basic_nack(self) -> bool: ...
@property
def consumer_cancel_notify(self) -> bool: ...
@property
def exchange_exchange_bindings(self) -> bool: ...
@property
def publisher_confirms(self) -> bool: ...

class _StreamingProtocolShim(AbstractStreamProtocol):
connection_made: Incomplete
connection_lost: Incomplete
eof_received: Incomplete
data_received: Incomplete
conn: Incomplete
def __init__(self, conn) -> None: ...
def __getattr__(self, attr: str): ...
# Next attributes are accessed via getattr() from BaseConnection class:
@classmethod
def create_connection(
cls,
connection_configs: Sequence[Parameters],
on_done: Callable[[Connection | AMQPConnectorException], object],
custom_ioloop: _IOLoop | None = None,
workflow: AbstractAMQPConnectionWorkflow | None = None,
) -> AbstractAMQPConnectionWorkflow: ...
@property
def ioloop(self) -> _IOLoop: ...
Loading