Broker API reference¤
The amqtt.broker.Broker class provides a complete MQTT 3.1.1 broker implementation. This class allows Python developers to embed a MQTT broker in their own applications.
Usage example¤
The following example shows how to start a broker using the default configuration:
import asyncio
from asyncio import CancelledError
import logging
from amqtt.broker import Broker
"""
This sample shows how to run a broker
"""
formatter = "[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s"
logging.basicConfig(level=logging.INFO, format=formatter)
async def run_server() -> None:
broker = Broker()
try:
await broker.start()
while True:
await asyncio.sleep(1)
except CancelledError:
await broker.shutdown()
def __main__():
try:
asyncio.run(run_server())
except KeyboardInterrupt:
print("Server exiting...")
if __name__ == "__main__":
__main__()
This will start the broker and let it run until it is shutdown by ^c.
Reference¤
Broker API¤
The amqtt.broker module provides the following key methods in the Broker class:
start(): Starts the broker and begins servingshutdown(): Gracefully shuts down the broker
Broker
¤
Broker(
config: BrokerConfig | dict[str, Any] | None = None,
loop: AbstractEventLoop | None = None,
plugin_namespace: str | None = None,
)
MQTT 3.1.1 compliant broker implementation.
Parameters:
-
(config¤BrokerConfig | dict[str, Any] | None, default:None) –BrokerConfigor dictionary of equivalent structure options (see broker configuration). -
(loop¤AbstractEventLoop | None, default:None) –asyncio loop. defaults to
asyncio.new_event_loop(). -
(plugin_namespace¤str | None, default:None) –plugin namespace to use when loading plugin entry_points. defaults to
amqtt.broker.plugins.
Raises:
-
BrokerError–problem with broker configuration
-
PluginImportError–if importing a plugin from configuration
-
PluginInitError–if initialization plugin fails
Methods:
-
external_connected–Engage the broker in handling the data stream to/from an established connection.
-
shutdown–Stop broker instance.
-
start–Start the broker to serve with the given configuration.
external_connected
async
¤
external_connected(
reader: ReaderAdapter,
writer: WriterAdapter,
listener_name: str,
) -> None
Engage the broker in handling the data stream to/from an established connection.
start
async
¤
start() -> None
Start the broker to serve with the given configuration.
Start method opens network sockets and will start listening for incoming connections.
Broker Context API¤
Broker plugins receive BrokerContext as their runtime context. It exposes
public broker state and helper methods for plugin implementations.
BrokerContext
¤
BrokerContext(broker: Broker)
Bases: BaseContext
Broker runtime context passed to broker plugins.
Broker plugins use this object to inspect broker state, publish internal messages, retain messages, and manage sessions or subscriptions without reaching into the broker's private attributes.
Methods:
-
add_subscription–Create a topic subscription for the given
client_id. -
broadcast_message–Send message to all client sessions subscribing to
topic. -
get_session–Return the session associated with
client_id, if it exists. -
retain_message–Retain a message on behalf of the broker.
Attributes:
-
retained_messages(dict[str, RetainedApplicationMessage]) –Retained messages keyed by topic name.
-
sessions(Generator[Session]) –All known broker sessions.
-
subscriptions(dict[str, list[tuple[Session, int]]]) –Active subscriptions keyed by topic filter.
retained_messages
property
¤
retained_messages: dict[str, RetainedApplicationMessage]
Retained messages keyed by topic name.
subscriptions
property
¤
subscriptions: dict[str, list[tuple[Session, int]]]
Active subscriptions keyed by topic filter.
add_subscription
async
¤
add_subscription(
client_id: str, topic: str | None, qos: int | None
) -> None
Create a topic subscription for the given client_id.
If a client session doesn't exist for client_id, create a disconnected session.
If topic and qos are both None, only create the client session.
broadcast_message
async
¤
broadcast_message(
topic: str, data: bytes, qos: int | None = None
) -> None
Send message to all client sessions subscribing to topic.
get_session
¤
get_session(client_id: str) -> Session | None
Return the session associated with client_id, if it exists.
retain_message
async
¤
retain_message(
topic_name: str,
data: bytes | bytearray,
qos: int | None = None,
) -> None
Retain a message on behalf of the broker.