Example

Tool

The testing tools use mosquitto_sub/mosquitto_pub. You may also use MQTTX or test scripts for this purpose.

# Subscribe test topic
mosquitto_sub -h test.mosquitto.org -p 1884 -u rw -P readwrite -t '/cassia/test/up'

# Publish test message
mosquitto_pub -h test.mosquitto.org -p 1884 -u rw -P readwrite -t '/cassia/test/down' -m '{"hello":"world"}'

Publish

 1import asyncio
 2
 3from cassiamqtt import CassiaMQTTClient
 4
 5
 6async def main():
 7    uri = "mqtt://test.mosquitto.org:1884"
 8    username = "rw"
 9    password = "readwrite"
10    client_id = "puber"
11    topic = "/cassia/test/up"
12
13    async with CassiaMQTTClient(
14        uri,
15        username=username,
16        password=password,
17        client_id=client_id,
18    ) as client:
19        print("connect ok")
20
21        payload = "hello"
22        qos = 1
23        retain = False
24        await client.publish(topic, payload, qos=qos, retain=retain)
25
26        print("pub ok")
27
28
29asyncio.run(main())

Subscribe

 1import asyncio
 2
 3from cassiamqtt import CassiaMQTTClient
 4
 5
 6async def main():
 7    uri = "mqtt://test.mosquitto.org:1884"
 8    username = "rw"
 9    password = "readwrite"
10    client_id = "puber"
11    topic = "/cassia/test/down"
12
13    async with CassiaMQTTClient(
14        uri,
15        username=username,
16        password=password,
17        client_id=client_id,
18    ) as client:
19        print("connect ok")
20
21        qos = 1
22        await client.subscribe(topic, qos=qos)
23        print("sub ok")
24
25        async for msg in client:
26            print("recv:", msg)
27
28            for k, v in msg.items():
29                print(k, type(k), v, type(v))
30
31
32asyncio.run(main())

Reconnect

  1import asyncio
  2import json
  3import gc
  4
  5import cassiamqtt
  6
  7
  8try:
  9    from typing import Optional, Tuple
 10except ImportError:
 11    pass
 12
 13
 14class CassiaMqttWrapper:
 15    DEFAULT_TOPIC_SUB = "/cassia/test/down"
 16    DEFAULT_TOPIC_PUB = "/cassia/test/up"
 17    DEFAULT_QOS = 1
 18
 19    def __init__(
 20        self,
 21        uri: str,
 22        username: Optional[str] = None,
 23        password: Optional[str] = None,
 24        client_id: Optional[str] = None,
 25    ):
 26        self._uri = uri
 27        self._user = username
 28        self._pwd = password
 29        self._cid = client_id
 30        self._client: Optional[cassiamqtt.CassiaMQTTClient] = None
 31
 32    async def _subscribe_topic(self) -> Tuple[bool, str]:
 33        if not self._client:
 34            return False, "client not connected"
 35        return await self._client.subscribe(self.DEFAULT_TOPIC_SUB, self.DEFAULT_QOS)
 36
 37    async def _handle_messages(self) -> None:
 38        print("[MQTT] Waiting for messages...")
 39        async for msg in self._client:
 40            print("[MQTT] Received:", msg)
 41
 42    async def runner(self) -> None:
 43        while True:
 44            try:
 45                print("[MQTT] Connecting...")
 46                async with cassiamqtt.CassiaMQTTClient(
 47                    self._uri, self._user, self._pwd, self._cid
 48                ) as client:
 49                    print("[MQTT] Connected")
 50                    self._client = client
 51
 52                    ok, ret = await self._subscribe_topic()
 53                    if not ok:
 54                        print(f"[MQTT] Subscribe failed: {ret}, retry in 5s...")
 55                        self._client = None
 56                        await asyncio.sleep(5)
 57                        continue
 58
 59                    print("[MQTT] Subscribe success")
 60                    await self._handle_messages()
 61
 62            except Exception as e:
 63                print(f"[MQTT] Connection error: {e}")
 64
 65            finally:
 66                print("[MQTT] Disconnected, retrying in 5s...")
 67                self._client = None
 68                await asyncio.sleep(5)
 69
 70    async def publish(
 71        self, topic: str, payload: str, qos: int = 0, retain: bool = False
 72    ) -> Tuple[bool, str]:
 73        if not self._client:
 74            return False, "no client"
 75        return await self._client.publish(topic, payload, qos, retain)
 76
 77    async def _mem_monitor(self) -> None:
 78        while True:
 79            data = {
 80                "free": gc.mem_free(),
 81                "alloc": gc.mem_alloc(),
 82            }
 83            print(f"[GC] memory: {data}")
 84
 85            msg = json.dumps(data)
 86            await self.publish(self.DEFAULT_TOPIC_PUB, msg, qos=self.DEFAULT_QOS)
 87
 88            await asyncio.sleep(30)
 89
 90
 91async def main() -> None:
 92    uri = "mqtt://test.mosquitto.org:1884"
 93    username = "rw"
 94    password = "readwrite"
 95    client_id = "CC:1B:E0:00:00:02"
 96
 97    mqtt = CassiaMqttWrapper(
 98        uri,
 99        username=username,
100        password=password,
101        client_id=client_id,
102    )
103
104    await asyncio.gather(
105        mqtt.runner(),
106        mqtt._mem_monitor(),
107    )
108
109
110if __name__ == "__main__":
111    asyncio.run(main())