示例

测试工具

使用 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())