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())
TLS/SSL
Standard MQTT over TLS
Only the server holds a certificate, and the client verifies the server certificate to ensure the channel is encrypted and the server is trusted. The client does not need to provide a certificate at the TLS layer (the client may use a username/password for MQTT-level authentication, or none at all).
Download the server certificate from test.mosquitto.org and set the ca parameter to the contents of the mosquitto.org.crt file.
1import asyncio
2
3from cassiamqtt import CassiaMQTTClient
4
5
6async def main():
7 uri = "mqtts://test.mosquitto.org:8883"
8 topic = "/cassia/test/up"
9
10 # mosquitto.org.crt
11 ca = """
12-----BEGIN CERTIFICATE-----
13MIIEAzCCAuugAwIBAgIUBY1hlCGvdj4NhBXkZ/uLUZNILAwwDQYJKoZIhvcNAQEL
14BQAwgZAxCzAJBgNVBAYTAkdCMRcwFQYDVQQIDA5Vbml0ZWQgS2luZ2RvbTEOMAwG
15A1UEBwwFRGVyYnkxEjAQBgNVBAoMCU1vc3F1aXR0bzELMAkGA1UECwwCQ0ExFjAU
16BgNVBAMMDW1vc3F1aXR0by5vcmcxHzAdBgkqhkiG9w0BCQEWEHJvZ2VyQGF0Y2hv
17by5vcmcwHhcNMjAwNjA5MTEwNjM5WhcNMzAwNjA3MTEwNjM5WjCBkDELMAkGA1UE
18BhMCR0IxFzAVBgNVBAgMDlVuaXRlZCBLaW5nZG9tMQ4wDAYDVQQHDAVEZXJieTES
19MBAGA1UECgwJTW9zcXVpdHRvMQswCQYDVQQLDAJDQTEWMBQGA1UEAwwNbW9zcXVp
20dHRvLm9yZzEfMB0GCSqGSIb3DQEJARYQcm9nZXJAYXRjaG9vLm9yZzCCASIwDQYJ
21KoZIhvcNAQEBBQADggEPADCCAQoCggEBAME0HKmIzfTOwkKLT3THHe+ObdizamPg
22UZmD64Tf3zJdNeYGYn4CEXbyP6fy3tWc8S2boW6dzrH8SdFf9uo320GJA9B7U1FW
23Te3xda/Lm3JFfaHjkWw7jBwcauQZjpGINHapHRlpiCZsquAthOgxW9SgDgYlGzEA
24s06pkEFiMw+qDfLo/sxFKB6vQlFekMeCymjLCbNwPJyqyhFmPWwio/PDMruBTzPH
253cioBnrJWKXc3OjXdLGFJOfj7pP0j/dr2LH72eSvv3PQQFl90CZPFhrCUcRHSSxo
26E6yjGOdnz7f6PveLIB574kQORwt8ePn0yidrTC1ictikED3nHYhMUOUCAwEAAaNT
27MFEwHQYDVR0OBBYEFPVV6xBUFPiGKDyo5V3+Hbh4N9YSMB8GA1UdIwQYMBaAFPVV
286xBUFPiGKDyo5V3+Hbh4N9YSMA8GA1UdEwEB/wQFMAMBAf8wDQYJKoZIhvcNAQEL
29BQADggEBAGa9kS21N70ThM6/Hj9D7mbVxKLBjVWe2TPsGfbl3rEDfZ+OKRZ2j6AC
306r7jb4TZO3dzF2p6dgbrlU71Y/4K0TdzIjRj3cQ3KSm41JvUQ0hZ/c04iGDg/xWf
31+pp58nfPAYwuerruPNWmlStWAXf0UTqRtg4hQDWBuUFDJTuWuuBvEXudz74eh/wK
32sMwfu1HFvjy5Z0iMDU8PUDepjVolOCue9ashlS4EB5IECdSR2TItnAIiIwimx839
33LdUdRudafMu5T5Xma182OC0/u/xRlEm+tvKGGmfFcN0piqVl8OrSPBgIlb+1IKJE
34m/XriWr/Cq4h/JfB7NTsezVslgkBaoU=
35-----END CERTIFICATE-----
36"""
37
38 async with CassiaMQTTClient(
39 uri,
40 ca=ca,
41 ) as client:
42 print("connect ok")
43
44 payload = "hello"
45 qos = 1
46 retain = False
47 await client.publish(topic, payload, qos=qos, retain=retain)
48
49 print("pub ok")
50
51
52asyncio.run(main())
Mutual TLS (mTLS) / Two-way TLS
The server requires the client to also provide a certificate. The server verifies the client’s certificate to confirm the client’s identity, offering stronger authentication than one-way TLS.
Generate a test client certificate based on Generate a TLS client certificate for test.mosquitto.org and upload it.
Note
When generating the CSR, do not use the default values. The CSR must include at least the Country, Organization, and Common Name fields.
If the content in this section differs from the information in the above link, please follow the instructions in the link and provide feedback promptly.
# 生成客户端证书
openssl genrsa -out client.key
openssl req -out client.csr -key client.key -new
# 将client.csr文件内容提交
# 下载新生成的client.crt
1import asyncio
2
3from cassiamqtt import CassiaMQTTClient
4
5
6async def main():
7 uri = "mqtts://test.mosquitto.org:8884"
8 topic = "/cassia/test/up"
9
10 # mosquitto.org.crt
11 ca = """
12-----BEGIN CERTIFICATE-----
13MIIEAzCCAuugAwIBAgIUBY1hlCGvdj4NhBXkZ/uLUZNILAwwDQYJKoZIhvcNAQEL
14BQAwgZAxCzAJBgNVBAYTAkdCMRcwFQYDVQQIDA5Vbml0ZWQgS2luZ2RvbTEOMAwG
15A1UEBwwFRGVyYnkxEjAQBgNVBAoMCU1vc3F1aXR0bzELMAkGA1UECwwCQ0ExFjAU
16BgNVBAMMDW1vc3F1aXR0by5vcmcxHzAdBgkqhkiG9w0BCQEWEHJvZ2VyQGF0Y2hv
17by5vcmcwHhcNMjAwNjA5MTEwNjM5WhcNMzAwNjA3MTEwNjM5WjCBkDELMAkGA1UE
18BhMCR0IxFzAVBgNVBAgMDlVuaXRlZCBLaW5nZG9tMQ4wDAYDVQQHDAVEZXJieTES
19MBAGA1UECgwJTW9zcXVpdHRvMQswCQYDVQQLDAJDQTEWMBQGA1UEAwwNbW9zcXVp
20dHRvLm9yZzEfMB0GCSqGSIb3DQEJARYQcm9nZXJAYXRjaG9vLm9yZzCCASIwDQYJ
21KoZIhvcNAQEBBQADggEPADCCAQoCggEBAME0HKmIzfTOwkKLT3THHe+ObdizamPg
22UZmD64Tf3zJdNeYGYn4CEXbyP6fy3tWc8S2boW6dzrH8SdFf9uo320GJA9B7U1FW
23Te3xda/Lm3JFfaHjkWw7jBwcauQZjpGINHapHRlpiCZsquAthOgxW9SgDgYlGzEA
24s06pkEFiMw+qDfLo/sxFKB6vQlFekMeCymjLCbNwPJyqyhFmPWwio/PDMruBTzPH
253cioBnrJWKXc3OjXdLGFJOfj7pP0j/dr2LH72eSvv3PQQFl90CZPFhrCUcRHSSxo
26E6yjGOdnz7f6PveLIB574kQORwt8ePn0yidrTC1ictikED3nHYhMUOUCAwEAAaNT
27MFEwHQYDVR0OBBYEFPVV6xBUFPiGKDyo5V3+Hbh4N9YSMB8GA1UdIwQYMBaAFPVV
286xBUFPiGKDyo5V3+Hbh4N9YSMA8GA1UdEwEB/wQFMAMBAf8wDQYJKoZIhvcNAQEL
29BQADggEBAGa9kS21N70ThM6/Hj9D7mbVxKLBjVWe2TPsGfbl3rEDfZ+OKRZ2j6AC
306r7jb4TZO3dzF2p6dgbrlU71Y/4K0TdzIjRj3cQ3KSm41JvUQ0hZ/c04iGDg/xWf
31+pp58nfPAYwuerruPNWmlStWAXf0UTqRtg4hQDWBuUFDJTuWuuBvEXudz74eh/wK
32sMwfu1HFvjy5Z0iMDU8PUDepjVolOCue9ashlS4EB5IECdSR2TItnAIiIwimx839
33LdUdRudafMu5T5Xma182OC0/u/xRlEm+tvKGGmfFcN0piqVl8OrSPBgIlb+1IKJE
34m/XriWr/Cq4h/JfB7NTsezVslgkBaoU=
35-----END CERTIFICATE-----
36"""
37
38 # client.crt
39 cert = """
40-----BEGIN CERTIFICATE-----
41MIIDsDCCApigAwIBAgIBADANBgkqhkiG9w0BAQsFADCBkDELMAkGA1UEBhMCR0Ix
42FzAVBgNVBAgMDlVuaXRlZCBLaW5nZG9tMQ4wDAYDVQQHDAVEZXJieTESMBAGA1UE
43CgwJTW9zcXVpdHRvMQswCQYDVQQLDAJDQTEWMBQGA1UEAwwNbW9zcXVpdHRvLm9y
44ZzEfMB0GCSqGSIb3DQEJARYQcm9nZXJAYXRjaG9vLm9yZzAeFw0yNTA5MjkwMzIy
45MDRaFw0yNTEyMjgwMzIyMDRaMIGJMQswCQYDVQQGEwJDTjEQMA4GA1UECAwHQmVp
46amluZzEQMA4GA1UEBwwHYmVpamluZzEPMA0GA1UECgwGY2Fzc2lhMQ8wDQYDVQQL
47DAZjYXNzaWExDzANBgNVBAMMBmNhc3NpYTEjMCEGCSqGSIb3DQEJARYUa2VybmVs
48LnBpZ0BnbWFpbC5jb20wggEiMA0GCSqGSIb3DQEBAQUAA4IBDwAwggEKAoIBAQDq
49WolFeaw8e8rDcMvKYRLwh7ieExhSqj6MCsG3iIyCZmA0ZrXI4ErsmfDLoXJNRfci
5012Yo+wQ0iWafcRjO3Vc1d+NQi95pxYWWiMPX1pJaPN/si3fVRbvRwb0B4L6zs+Gq
51zk9vKQoOa+diYlS0Nfh1zIr1kykrAkmTqIUdhIPtt/Hg2uT/JltFr+fciCf+vnfz
52D2HaqNd40XPJfiQ8Glu20A3Q3VPZpI+SxLL2zlyLkz0YU2oWCkLLCMtRDaFWMid9
53/lfZD4td0XEhLfjECJO+tEoccWGDEGKiHsgN8Z+V330xlu8TWmMt4a8hlnYc/ZPv
54AI5hKXV0GmRYi2oR+ahHAgMBAAGjGjAYMAkGA1UdEwQCMAAwCwYDVR0PBAQDAgXg
55MA0GCSqGSIb3DQEBCwUAA4IBAQCgywwwYt4qxWjhACakJmCUL+E5jBB64OzsE+BH
568Eo/vpXqmst3DGPLT11JE15h2uFBGqXTRDxclloqpvvyjBPjHISKRPc3oHoPDwek
57ukMXD7q3OUbpxqv43Rq1MUPqjuYIPRJ8+bPZp8jaVw6foqMbIbtG8Lef1EjYuuj+
58AqdV5WxI/AkZYDlcmzIfsj3C8Aan39BrqjQuNC7j7rYwuCsHJMqftv8ggSH9mCBY
59a7ODn7t7IrU+rVynra8MlBaJ2wH4vRp0VGDiNCnPDkEbxS+HSMvFsLDptWCrOMOI
60NtnskeaZEqQHvIU/WQhggbxlQ+wBYzrJfAM29ociq5Yo9607
61-----END CERTIFICATE-----
62"""
63
64 # client.key
65 key = """
66-----BEGIN PRIVATE KEY-----
67MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDqWolFeaw8e8rD
68cMvKYRLwh7ieExhSqj6MCsG3iIyCZmA0ZrXI4ErsmfDLoXJNRfci12Yo+wQ0iWaf
69cRjO3Vc1d+NQi95pxYWWiMPX1pJaPN/si3fVRbvRwb0B4L6zs+Gqzk9vKQoOa+di
70YlS0Nfh1zIr1kykrAkmTqIUdhIPtt/Hg2uT/JltFr+fciCf+vnfzD2HaqNd40XPJ
71fiQ8Glu20A3Q3VPZpI+SxLL2zlyLkz0YU2oWCkLLCMtRDaFWMid9/lfZD4td0XEh
72LfjECJO+tEoccWGDEGKiHsgN8Z+V330xlu8TWmMt4a8hlnYc/ZPvAI5hKXV0GmRY
73i2oR+ahHAgMBAAECggEAYaAnkRaXpnaXSAVkD8GSpzqSsN8Zgc5D0gjlG/S5O9Uz
74/IBQ+AZfj+KtCdcOi5w60HvUpfuzi8M1SfRONlbEbpSr0DEEDSHofYYpt28+dnLn
75gV20JNcw37eKag7awneL6aCaPJ9m/gz7TePSj2PwRfpYZObR/oWsauOH2H0MjGkJ
76MvnPJNyk7gbgtYjCHFTtEOJsBAr/DfQYYwRBeuxK6kKytrDI2FH5lElgBHH82S/k
77I32/K2q2NmpRlLrSaR/V4tgLh7yA4wzKNQXK8mNUVHnF7Rv396tKVpE76dNMp/oC
78C9UPWeOG2AUsD1Yte+ABAMDk8PODeZHVKXlX8zqg2QKBgQD5GrpkcBXskT+38P5B
79ZCeB5f6DzgvU+b5+m+Jf8S7G4rAkTPCswd3JH52wWr6mMwIgAoPvoGs2e3/ndiph
80HgdLF47EdXksy4LwY5Yic/nzd9BwBgDb91ztArLk7Y6Alsg6M4LbKD5OYsj0XSF0
81jUXEc9Ea8PJmnWD9nvVvQxV5cwKBgQDw10b+UO89cQKqvEZEqcVHjqaUpGXYIb5w
82q+aIrB+wyRUNzOtVvM2TDWTp9/vIlR3BYAVLu4vY+uceeEkLsLSsfw9y+/2Vwk9t
83xdXgFPE4/d7s3ITlpOAtlLxYpbAXh6G4sYPBaT/gIvyPeolPEhv3frF74QjCjnjT
84I2jJyhzw3QKBgA46H5Ei8a2cMhZwViUn8jWyDBI9D2HvjZivkINIRBKp2cOI/Wnw
85fJlDC/+Jf0AAw8tOOXjTIaxv60Mt9YesbmA0jTvdNbmAOg8+sNFw7EKigi4TubLW
86cuE6eTsn8i6X7gGc9YlMyBoz/CQwuXttoiFxN+0g+8cuj96MWJotK6nPAoGBAIMr
873OC6V/LQ0DEJZgQTqvz0NsoSV93FUyGunlql1ITGoA7qIuqJcDW9P88mXx26CYC+
88uWOr+9jrnmE8Bhy121FvyoxHrq+YKwaQj5ICFfeCXZ4H5OHmUKrCrWpio2vNVUlw
89dWAr4LxEkeXbSVmldVHw0N21jL3aNvhX+sScrfKJAoGBAIFx8oT2tt27a6iw52Ga
904T/8qfqsyT/b6oCVsWeh2SLoL0Ai7lqboc9eUqNSovAwQf3pQKbucY3GSdm2NcTg
91i6RQf4DT2pyG+ji8LerMRwFHgFJqoU5q/6eBTYupe3OqpbAp3JlvqkUUrhHwNtrC
92X0QUZvYkcYeyI54XC8/cThTJ
93-----END PRIVATE KEY-----
94"""
95
96 async with CassiaMQTTClient(
97 uri,
98 ca=ca,
99 cert=cert,
100 key=key,
101 ) as client:
102 print("connect ok")
103
104 payload = "hello"
105 qos = 1
106 retain = False
107 await client.publish(topic, payload, qos=qos, retain=retain)
108
109 print("pub ok")
110
111
112asyncio.run(main())