Skip to content

hololinked.server.mqtt.MQTTPublisher

Bases: BaseProtocolServer

MQTT Publisher.

All events and observable properties defined on the Thing will be published to MQTT topics with topic name "{thing id}/{event name}".

For setting up an MQTT broker if one does not exist, see infrastructure project.

Source code in repo/hololinked/hololinked/server/mqtt/server.py
class MQTTPublisher(BaseProtocolServer):
    """
    MQTT Publisher.

    All events and observable properties defined on the Thing will be published to MQTT topics
    with topic name "{thing id}/{event name}".

    For setting up an MQTT broker if one does not exist,
    see [infrastructure project](https://github.com/hololinked-dev/daq-system-infrastructure).
    """

    hostname = String(default="localhost")
    """The MQTT broker hostname"""

    ssl_context = ClassSelector(class_=ssl.SSLContext, allow_None=True, default=None)
    """The SSL context to use for secure connections, or None for no SSL"""

    config = ClassSelector(class_=RuntimeConfig, default=None, allow_None=True)  # type: RuntimeConfig
    """Runtime configuration for the MQTT publisher"""

    def __init__(
        self,
        hostname: str,
        port: int,
        username: str,
        password: str,
        qos: int = 1,
        things: Optional[list[CoreThing]] = None,
        config: Optional[dict] = None,
        **kwargs,
    ):
        """
        Initialize the MQTT publisher.

        Parameters
        ----------
        hostname: str
            The MQTT broker hostname
        port: int
            The MQTT broker port
        username: str
            The MQTT broker username
        password: str
            The MQTT broker password
        qos: int
            The (global) MQTT QoS level to use for publishing messages
        things: list[Thing]
            The `Thing`s that need to publish their events/properties to MQTT broker
        config: dict, optional
            Additional runtime configuration for the MQTT publisher, see `RuntimeConfig` object under
            `hololinked.server.mqtt.config`
        kwargs: dict
            Additional keyword arguments
        """
        default_config: dict[str, Any] = dict(
            topic_publisher=kwargs.get("topic_publisher", TopicPublisher),
            thing_description_publisher=kwargs.get("thing_description_publisher", ThingDescriptionPublisher),
            thing_description_service=kwargs.get("thing_description_service", ThingDescriptionService),
            qos=qos,
        )
        default_config.update(config or dict())

        super().__init__(config=RuntimeConfig(**default_config))

        endpoint = f"{hostname}{f':{port}' if port else ''}"

        self.hostname = hostname
        self.port = port
        self.username = username
        self.password = password
        self.publishers = dict()  # type: dict[str, TopicPublisher]
        self.logger = kwargs.get("logger", structlog.get_logger()).bind(component="mqtt-publisher", hostname=endpoint)
        self.ssl_context = kwargs.get("ssl_context", None)
        self.id = endpoint
        self.add_things(*(things or []))

    @classmethod
    def from_params(cls, id: str, params: str | int | dict | list[str] | None) -> Self:  # noqa: D102
        # doc already provided in the base class
        if isinstance(params, str):
            params = dict(hostname=params)
        elif not isinstance(params, dict):
            raise ValueError("MQTT parameters must be supplied as a dictionary or the broker hostname as a string.")
        return cls(**params)

    def welcome_lines(self) -> list[str]:  # noqa: D102
        # doc already provided in the base class
        lines = ["", "📡 MQTT:", f" • Broker:   {self.hostname}:{self.port}"]
        for thing in self.things.values():
            lines.append(f"   ➜ Topic tree: {thing.id}/thing-description")
        return lines

    async def start(self):
        """
        Sets up the MQTT client and starts publishing events from the `Thing`s.

        All events are dispatched to their own async tasks. This method returns and
        creates side-effects only & does not block. Use the `run()` method instead for a blocking call.
        """
        await self.setup()
        loop = get_current_async_loop()
        for thing in self.things.values():
            loop.create_task(self.start_publishers(thing))

    async def start_publishers(self, thing: CoreThing) -> None:
        """
        Start the publishers for a given `Thing`.

        Raises
        ------
        ValueError
            if the `Thing` is not bound to an event loop
        """
        loop = get_current_async_loop()
        if not thing.eventloop:
            raise ValueError(f"Thing {thing.id} is not associated with any event loop")
        TD = thing.get_thing_model(ignore_errors=True).json()

        for event_name in TD.get("events", {}).keys():
            event_affordance = EventAffordance.from_TD(event_name, TD)
            topic_publisher = self.config.topic_publisher(
                client=self.client,
                resource=event_affordance,
                logger=self.logger,
                config=self.config,
                thing=thing,
            )
            self.publishers[topic_publisher.topic] = topic_publisher
            loop.create_task(topic_publisher.publish())
            self.logger.info(f"MQTT will publish events for {event_name} of thing {thing.id}")
        for prop_name in TD.get("properties", {}).keys():
            property_affordance = PropertyAffordance.from_TD(prop_name, TD)
            if not property_affordance.observable:
                continue
            topic_publisher = self.config.topic_publisher(
                client=self.client,
                resource=property_affordance,
                logger=self.logger,
                config=self.config,
                thing=thing,
            )
            self.publishers[topic_publisher.topic] = topic_publisher
            loop.create_task(topic_publisher.publish())
            self.logger.info(f"MQTT will publish observable property changes for {prop_name} of thing {thing.id}")
        # TD publisher
        td_publisher = self.config.thing_description_publisher(
            client=self.client,
            logger=self.logger,
            thing=thing,
            config=self.config,
        )
        self.publishers[td_publisher.topic] = td_publisher
        loop.create_task(td_publisher.publish())

    async def setup(self) -> None:
        """Setup MQTT publishers per `Thing` post connection to broker."""
        self.client = aiomqtt.Client(
            hostname=self.hostname,
            port=self.port,
            username=self.username,
            password=self.password,
            tls_context=self.ssl_context,
        )
        try:
            await self.client.__aenter__()
            endpoint = f"{self.hostname}{f':{self.port}' if self.port else ''}"
            self.logger.info(f"Connected to MQTT broker at {endpoint}")
        except aiomqtt.MqttReentrantError:
            pass

    def stop(self):
        """Stop publishing, the client is not closed automatically."""
        for publisher in self.publishers.values():
            publisher.stop()

Attributes

id instance-attribute

id = endpoint

hostname class-attribute instance-attribute

hostname = hostname

The MQTT broker hostname

port instance-attribute

port = port

username instance-attribute

username = username

password instance-attribute

password = password

ssl_context class-attribute instance-attribute

ssl_context = get('ssl_context', None)

The SSL context to use for secure connections, or None for no SSL

config class-attribute instance-attribute

config = ClassSelector(class_=RuntimeConfig, default=None, allow_None=True)

Runtime configuration for the MQTT publisher

logger instance-attribute

logger = bind(component='mqtt-publisher', hostname=endpoint)

things instance-attribute

things: MutableMapping[str, Thing] = TypeConstrainedDict({}, key_type=str, item_type=Thing)

Every served Thing, by id. Sub-things are not served.

publishers instance-attribute

publishers = dict()

Functions

__init__

__init__(hostname: str, port: int, username: str, password: str, qos: int = 1, things: Optional[list[Thing]] = None, config: Optional[dict] = None, **kwargs)

Initialize the MQTT publisher.

Parameters:

Name Type Description Default

hostname

str

The MQTT broker hostname

required

port

int

The MQTT broker port

required

username

str

The MQTT broker username

required

password

str

The MQTT broker password

required

qos

int

The (global) MQTT QoS level to use for publishing messages

1

things

Optional[list[Thing]]

The Things that need to publish their events/properties to MQTT broker

None

config

Optional[dict]

Additional runtime configuration for the MQTT publisher, see RuntimeConfig object under hololinked.server.mqtt.config

None

kwargs

Additional keyword arguments

{}
Source code in repo/hololinked/hololinked/server/mqtt/server.py
def __init__(
    self,
    hostname: str,
    port: int,
    username: str,
    password: str,
    qos: int = 1,
    things: Optional[list[CoreThing]] = None,
    config: Optional[dict] = None,
    **kwargs,
):
    """
    Initialize the MQTT publisher.

    Parameters
    ----------
    hostname: str
        The MQTT broker hostname
    port: int
        The MQTT broker port
    username: str
        The MQTT broker username
    password: str
        The MQTT broker password
    qos: int
        The (global) MQTT QoS level to use for publishing messages
    things: list[Thing]
        The `Thing`s that need to publish their events/properties to MQTT broker
    config: dict, optional
        Additional runtime configuration for the MQTT publisher, see `RuntimeConfig` object under
        `hololinked.server.mqtt.config`
    kwargs: dict
        Additional keyword arguments
    """
    default_config: dict[str, Any] = dict(
        topic_publisher=kwargs.get("topic_publisher", TopicPublisher),
        thing_description_publisher=kwargs.get("thing_description_publisher", ThingDescriptionPublisher),
        thing_description_service=kwargs.get("thing_description_service", ThingDescriptionService),
        qos=qos,
    )
    default_config.update(config or dict())

    super().__init__(config=RuntimeConfig(**default_config))

    endpoint = f"{hostname}{f':{port}' if port else ''}"

    self.hostname = hostname
    self.port = port
    self.username = username
    self.password = password
    self.publishers = dict()  # type: dict[str, TopicPublisher]
    self.logger = kwargs.get("logger", structlog.get_logger()).bind(component="mqtt-publisher", hostname=endpoint)
    self.ssl_context = kwargs.get("ssl_context", None)
    self.id = endpoint
    self.add_things(*(things or []))

from_params classmethod

from_params(id: str, params: str | int | dict | list[str] | None) -> Self
Source code in repo/hololinked/hololinked/server/mqtt/server.py
@classmethod
def from_params(cls, id: str, params: str | int | dict | list[str] | None) -> Self:  # noqa: D102
    # doc already provided in the base class
    if isinstance(params, str):
        params = dict(hostname=params)
    elif not isinstance(params, dict):
        raise ValueError("MQTT parameters must be supplied as a dictionary or the broker hostname as a string.")
    return cls(**params)

add_thing

add_thing(thing: Thing) -> None

Adds a thing to the things being served.

Sub-things are not served - see EventLoop.add_thing, which does not register them either.

Source code in repo/hololinked/hololinked/core/interfaces/protocol_server.py
def add_thing(self, thing: Thing) -> None:
    """
    Adds a thing to the things being served.

    Sub-things are not served - see `EventLoop.add_thing`, which does not register them
    either.
    """
    self.things[thing.id] = thing

add_things

add_things(*things: Thing) -> None

Adds multiple things to be served.

Source code in repo/hololinked/hololinked/core/interfaces/protocol_server.py
def add_things(self, *things: Thing) -> None:
    """Adds multiple things to be served."""
    for thing in things:
        self.add_thing(thing)

add_property

add_property(*args, **kwargs) -> None

Add a property to be served.

Raises:

Type Description
NotImplementedError

if the protocol does not support this operation

Source code in repo/hololinked/hololinked/core/interfaces/protocol_server.py
def add_property(self, *args, **kwargs) -> None:
    """
    Add a property to be served.

    Raises
    ------
    NotImplementedError
        if the protocol does not support this operation
    """
    raise NotImplementedError("Not implemented for this protocol")

add_action

add_action(*args, **kwargs) -> None

Add an action to be served.

Raises:

Type Description
NotImplementedError

if the protocol does not support this operation

Source code in repo/hololinked/hololinked/core/interfaces/protocol_server.py
def add_action(self, *args, **kwargs) -> None:
    """
    Add an action to be served.

    Raises
    ------
    NotImplementedError
        if the protocol does not support this operation
    """
    raise NotImplementedError("Not implemented for this protocol")

add_event

add_event(*args, **kwargs) -> None

Add an event to be served.

Raises:

Type Description
NotImplementedError

if the protocol does not support this operation

Source code in repo/hololinked/hololinked/core/interfaces/protocol_server.py
def add_event(self, *args, **kwargs) -> None:
    """
    Add an event to be served.

    Raises
    ------
    NotImplementedError
        if the protocol does not support this operation
    """
    raise NotImplementedError("Not implemented for this protocol")

setup async

setup() -> None

Setup MQTT publishers per Thing post connection to broker.

Source code in repo/hololinked/hololinked/server/mqtt/server.py
async def setup(self) -> None:
    """Setup MQTT publishers per `Thing` post connection to broker."""
    self.client = aiomqtt.Client(
        hostname=self.hostname,
        port=self.port,
        username=self.username,
        password=self.password,
        tls_context=self.ssl_context,
    )
    try:
        await self.client.__aenter__()
        endpoint = f"{self.hostname}{f':{self.port}' if self.port else ''}"
        self.logger.info(f"Connected to MQTT broker at {endpoint}")
    except aiomqtt.MqttReentrantError:
        pass

start async

start()

Sets up the MQTT client and starts publishing events from the Things.

All events are dispatched to their own async tasks. This method returns and creates side-effects only & does not block. Use the run() method instead for a blocking call.

Source code in repo/hololinked/hololinked/server/mqtt/server.py
async def start(self):
    """
    Sets up the MQTT client and starts publishing events from the `Thing`s.

    All events are dispatched to their own async tasks. This method returns and
    creates side-effects only & does not block. Use the `run()` method instead for a blocking call.
    """
    await self.setup()
    loop = get_current_async_loop()
    for thing in self.things.values():
        loop.create_task(self.start_publishers(thing))

start_publishers async

start_publishers(thing: Thing) -> None

Start the publishers for a given Thing.

Raises:

Type Description
ValueError

if the Thing is not bound to an event loop

Source code in repo/hololinked/hololinked/server/mqtt/server.py
async def start_publishers(self, thing: CoreThing) -> None:
    """
    Start the publishers for a given `Thing`.

    Raises
    ------
    ValueError
        if the `Thing` is not bound to an event loop
    """
    loop = get_current_async_loop()
    if not thing.eventloop:
        raise ValueError(f"Thing {thing.id} is not associated with any event loop")
    TD = thing.get_thing_model(ignore_errors=True).json()

    for event_name in TD.get("events", {}).keys():
        event_affordance = EventAffordance.from_TD(event_name, TD)
        topic_publisher = self.config.topic_publisher(
            client=self.client,
            resource=event_affordance,
            logger=self.logger,
            config=self.config,
            thing=thing,
        )
        self.publishers[topic_publisher.topic] = topic_publisher
        loop.create_task(topic_publisher.publish())
        self.logger.info(f"MQTT will publish events for {event_name} of thing {thing.id}")
    for prop_name in TD.get("properties", {}).keys():
        property_affordance = PropertyAffordance.from_TD(prop_name, TD)
        if not property_affordance.observable:
            continue
        topic_publisher = self.config.topic_publisher(
            client=self.client,
            resource=property_affordance,
            logger=self.logger,
            config=self.config,
            thing=thing,
        )
        self.publishers[topic_publisher.topic] = topic_publisher
        loop.create_task(topic_publisher.publish())
        self.logger.info(f"MQTT will publish observable property changes for {prop_name} of thing {thing.id}")
    # TD publisher
    td_publisher = self.config.thing_description_publisher(
        client=self.client,
        logger=self.logger,
        thing=thing,
        config=self.config,
    )
    self.publishers[td_publisher.topic] = td_publisher
    loop.create_task(td_publisher.publish())

run

run(forked: bool = False, print_welcome_message: bool = True) -> None

Run the server and serve your things.

Use this method if this is the only running protocol. Blocks.

Parameters:

Name Type Description Default

forked

bool

whether to run in a forked thread

False

print_welcome_message

bool

whether to print a welcome message on startup, like the ports and access points

True
Source code in repo/hololinked/hololinked/core/interfaces/protocol_server.py
@forkable
def run(self, forked: bool = False, print_welcome_message: bool = True) -> None:
    """
    Run the server and serve your things.

    Use this method if this is the only running protocol. Blocks.

    Parameters
    ----------
    forked: bool, default False
        whether to run in a forked thread
    print_welcome_message: bool, default True
        whether to print a welcome message on startup, like the ports and access points
    """
    from hololinked.server import run

    run(self, print_welcome_message=print_welcome_message)

stop

stop()

Stop publishing, the client is not closed automatically.

Source code in repo/hololinked/hololinked/server/mqtt/server.py
def stop(self):
    """Stop publishing, the client is not closed automatically."""
    for publisher in self.publishers.values():
        publisher.stop()

welcome_lines

welcome_lines() -> list[str]
Source code in repo/hololinked/hololinked/server/mqtt/server.py
def welcome_lines(self) -> list[str]:  # noqa: D102
    # doc already provided in the base class
    lines = ["", "📡 MQTT:", f" • Broker:   {self.hostname}:{self.port}"]
    for thing in self.things.values():
        lines.append(f"   ➜ Topic tree: {thing.id}/thing-description")
    return lines