示例
测试工具
使用 mosquitto_sub/mosquitto_pub,您也可以使用 MQTTX 或者测试脚本进行。
# 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"}'
发布
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())
订阅
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())
重连
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
标准MQTT over TLS
仅服务器(server)有证书,客户端验证服务器证书以保证通道被加密和服务器可信。客户端不需要在 TLS 层提供证书(客户端可以用用户名/密码做 MQTT 层认证,或根本不做)。
从 test.mosquitto.org 下载Server证书,设置 ca 参数为 mosquitto.org.crt 文件内容
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())
双向 TLS / mTLS(mutual TLS)
服务器要求客户端也提供证书,服务端会验证客户端证书以确认客户端身份(比单向 TLS 更强的身份验证)。
根据 Generate a TLS client certificate for test.mosquitto.org 生成测试Client证书并上传
备注
生成 CSR 时,请勿使用默认值。CSR 至少必须包含国家、组织和通用名称字段。
本部分内容和上面链接内容不一致时,请以上述链接内容为准并及时反馈
# 生成客户端证书
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())