一、代码总体功能

这是一个运行在嵌入式 Linux 设备上的 MQTT 客户端,主要功能:

  1. 从配置文件 /etc/wechat.cfg 读取连接参数(Broker URI、ClientID、用户名、密码、产品 ID、设备名称)。

  2. 初始化 RPC 客户端(用于调用板载硬件接口,如读取温湿度传感器、控制 LED)。

  3. 连接 MQTT 服务器(Broker)。

  4. 订阅下行主题(接收云端或手机端下发的控制命令,例如开灯/关灯)。

  5. 循环读取温湿度传感器,构造 JSON 消息,发布到上行主题(上报设备数据)。

  6. 在回调函数中解析收到的 JSON 指令,控制 LED。

📋 核心结构体一览 (以同步客户端 MQTTClient.h 为例)

我把这些常用的结构体整理了出来,这样理解起来会更清晰:

结构体名称 主要用途 在代码示例中的角色
MQTTClient 客户端实例的句柄 (Handle)。它本身是一个不透明的指针,代表了你的 MQTT 客户端对象。 在 main() 函数中声明的 MQTTClient client;。后续的 MQTTClient_create() 等函数都通过它来操作这个特定的客户端。
MQTTClient_message 代表一个具体的MQTT消息。它不包含主题(topic),只包含消息的具体内容,如payload(数据负载)、payloadlen(长度)、qos(服务质量)等。 出现在 msgarrvd 回调函数中。当设备收到消息时,库会填充好这个结构体并传给你的回调函数。你也用它来构造要发布的消息。
MQTTClient_connectOptions 连接选项。用于配置客户端连接到Broker(MQTT服务器)时的各项参数,如keepAliveInterval(心跳间隔)、cleansession(会话清理)、用户名/密码等。 在 main() 函数中声明,并设置好参数后传给 MQTTClient_connect() 发起连接。
MQTTClient_willOptions 遗嘱(Last Will and Testament, LWT)选项。它让你能提前指示Broker,当你的客户端异常断开时,自动代你发布一条预先定义好的消息到某个主题。 如果你的设备需要有“上线/离线”通知功能,就会用到这个结构体。它通常作为 MQTTClient_connectOptions 的一个子结构体被使用。
MQTTClient_SSLOptions TLS/SSL安全连接选项。用于配置如何建立安全的MQTT连接,例如指定CA证书、客户端证书和私钥等。 当你连接的Broker要求进行SSL/TLS加密通信时,你需要配置这个结构体,并将其关联到 MQTTClient_connectOptions 中。
MQTTClient_createOptions 客户端创建选项。这是一个更高级的配置结构,主要用于在创建客户端时设定一些不太常用的选项。例如,你可以通过它指定要使用的MQTT协议版本(如5.0)。 这是 MQTTClient_create() 函数的高级版本 MQTTClient_createWithOptions() 的参数。
MQTTClient_persistence 消息持久化接口。它不是一个简单的数据容器,而是一个包含函数指针的结构体。它允许你自定义消息持久化的方式,比如将离线期间收到的消息存入文件、数据库或内存中。 MQTTClient_create() 函数会用到它,你通过传入 MQTTCLIENT_PERSISTENCE_NONE 或 MQTTCLIENT_PERSISTENCE_DEFAULT 来选择默认的持久化行为。
MQTTProperties / MQTTProperty MQTT 5.0 协议属性。这两个结构体用于支持MQTT 5.0协议中引入的属性(Properties)机制。MQTTProperties 是一个属性的集合。 如果你的项目需要使用MQTT 5.0的高级特性(如用户属性、请求-响应模式等),你就需要操作这些结构体。

二、逐段详解

1. 头文件和宏定义

c

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "MQTTClient.h"
#include "cJSON.h"
#include "rpc_client.h"
#include "cfg.h"

#define ADDRESS     "tcp://broker.emqx.io:1883"   // 备用 Broker 地址(实际被配置文件覆盖)
#define CLIENTID    "mqttx_6e973856"              // 备用 ClientID(实际被配置文件覆盖)
#define TOPIC_SUBSCRIBE     "/iot/down"           // 备用订阅主题(实际被字符串覆盖)
#define TOPIC_PUBLISH       "/iot/up"             // 备用发布主题(实际被字符串覆盖)
#define QOS         1                             // 服务质量等级:0/1/2,1表示至少一次送达
#define TIMEOUT     10000L                        // 等待发布完成的超时时间,单位毫秒

#define USER_NAME "100ask"      // 备用用户名(实际被配置文件覆盖)
#define PASSWORD  "100ask"      // 备用密码(实际被配置文件覆盖)
  • MQTTClient.h:Paho MQTT C 客户端库头文件。

  • cJSON.h:轻量级 JSON 解析库。

  • rpc_client.h:自定义 RPC 客户端,用于与底层硬件交互(如 rpc_dht11_readrpc_led_control)。

  • cfg.h:声明 read_cfg 函数。

注意:代码后面会用 strcpy 将 pub_topic/sub_topic 硬覆盖为这些宏,但实际上动态构造的主题(包含 ProductID/DeviceName)被覆盖掉了,这可能是调试遗留代码。我们分析时仍按功能逻辑理解。


2. 全局变量和回调函数

c

volatile MQTTClient_deliveryToken deliveredtoken;
  • deliveredtoken:标记最近一次消息的发送 token,用于确认消息是否被 Broker 接收。volatile 防止编译器优化。token 用于唯一标识某次发布操作,以便在 delivered 回调中确认是哪一条消息被 Broker 成功接收。delivery为服务器成功接收后,返回给客户端参数的回调函数

2.1 delivered 回调 —— 消息送达确认
void delivered(void *context, MQTTClient_deliveryToken dt)
{
    printf("Message with token value %d delivery confirmed\n", dt);
    deliveredtoken = dt;
}  //confirmed(已确认)

delivered 回调 —— 消息被服务器确认

  • 触发时机:当客户端发布的消息(QoS ≥ 1)被服务器成功接收并确认(收到 PUBACK)后,库调用此回调。

  • 数据流向:客户端 → 服务器(的确认反馈)。

  • 典型用途:确认发布操作完成,可用于日志、释放资源、触发下一批发送。

  • 是谁调用的:MQTT 客户端库在收到 PUBACK 后,调用用户注册的 delivered

  • QoS(Quality of Service,服务质量)0:至多一次 1:至少成功一次 2:只成功一次

  • 当发布的消息(QoS≥1)成功被 Broker 确认接收后,库会调用此函数。

  • dt:该消息的唯一 token(由 MQTTClient_publishMessage 返回)。int 类型

  • 作用:让用户知道发送已完成,一般用于日志或触发下一步操作。

  • Message with token value 10 delivery confirmed
    如10,为第10次发布

2.2 msgarrvd 回调 —— 收到消息
int msgarrvd(void *context, char *topicName, int topicLen, MQTTClient_message *message)
  • 当订阅的主题上有新消息到达时,库自动调用此函数。

  • msgarrvd 回调 —— 消息到达客户端

  • 触发时机:当客户端从服务器(Broker)收到一条订阅的消息时,MQTT 客户端库会自动调用此回调。 MQTT 客户端库在收到一条订阅的消息时自动调用的函数。

  • 数据流向:服务器 → 客户端。

  • 典型用途:解析收到的 JSON 指令,执行对应的动作(如控制 LED、上报 ACK)。

  • 是谁调用的:MQTT 客户端库在内部收到 PUBLISH 报文后,调用用户注册的 msgarrvd

  • 参数详解

    • context:用户传入的上下文指针(本例为 NULL)。

    • topicName:消息所属的主题字符串。

    • topicLen:主题长度(如果为 0,则需自己 strlen)。

    • message:指向消息结构体的指针,包含 payload(数据)、payloadlenqosretained 等。//如发送,用于实际调动开发板上的硬件指令
      {
          "method": "control",
          "clientToken": "mqttx_6e973856",
          "params": {
              "LED1": 1
          }
      } //可调动下面代码的LED指令,由客户端发出
       

  • 函数内部:

    1. 打印主题和消息内容。

    2. 用 cJSON 解析 message->payload

    3. 查找 "params" 对象,再查找 "LED1"

    4. 调用 rpc_led_control(1, led1->valueint) 控制 LED(valueint 为 0 或 1)。

    5. 释放 cJSON 根对象,以及 MQTT 库分配的消息内存。

  • 返回值:1 表示成功(非 0 即可)。

关键点message->payload 是二进制数据,可能不以 '\0' 结尾,但本例中发送的是合法 JSON 字符串,通常发送方会加结束符,所以 cJSON_Parse 直接转换风险较低。更安全的是用 cJSON_ParseWithLength

2.3 connlost 回调 —— 连接断开
void connlost(void *context, char *cause)
  • 当 MQTT 连接意外断开(网络故障、Broker 重启等)时调用。

  • cause:描述断开原因的字符串。

  • 本例只打印信息,没有尝试重连(主循环中虽然无限重连,但连接成功后才会进入发布循环。若运行中失去连接,只会打印,程序不会自动重连——这是一个潜在缺陷)。


3. main 函数逐步解析

3.1 变量定义
MQTTClient client;                      // MQTT 客户端句柄
MQTTClient_connectOptions conn_opts;    // 连接选项结构体
int rc;                                 // 返回码,用于检查库函数是否成功

char URI[1000], clientid[1000], username[1000], password[1000], ProductID[1000], DeviceName[1000];
char address[1000];                     // 构建的 Broker 地址(形如 "tcp://xxx:1883")
char pub_topic[1000], sub_topic[1000];  // 动态构造的发布和订阅主题
3.2 步骤1:读取配置文件
if (0 != read_cfg(URI, clientid, username, password, ProductID, DeviceName))
{
    printf("read cfg err\n");
    return -1;
}
printf("%s\n%s\n%s\n%s\n%s\n%s\n", URI, clientid, username, password, ProductID, DeviceName);
  • read_cfg 解析 /etc/wechat.cfg(JSON 格式),将各个字段填入传入的字符数组。

  • 如果失败(文件不存在、解析错误等),程序退出。

  • 打印出来供调试。

3.3 构造连接地址和主题(有逻辑问题)
sprintf(address, "tcp://%s:1883", URI);  //假设配置文件中 URI 仅包含 IP 或域名,不含协议和端口
/*此为微信小程序本程序中不需要*/
sprintf(pub_topic, "$thing/up/property/%s/%s", ProductID, DeviceName); 
sprintf(sub_topic, "$thing/down/property/%s/%s", ProductID, DeviceName);
/******************************************************************/
strcpy(pub_topic, TOPIC_PUBLISH);               // 硬覆盖为 "/iot/up"发布(主题)
strcpy(sub_topic, TOPIC_SUBSCRIBE);             // 硬覆盖为 "/iot/down"订阅(主题)
  • 正常逻辑:根据配置的 URI(例如 "broker.emqx.io")加上 tcp:// 和端口 1883 构成完整地址。

  • $thing/up/property/{ProductID}/{DeviceName} 是腾讯云 IoT 等平台的标准格式,但随后被 strcpy 覆盖成简单字符串。最终发布/订阅主题变成了 /iot/up 和 /iot/down。根据实际运行时要根据需求选择使用哪一种。

  • 教学解读中,我们忽略覆盖,认为最终使用的是 /iot/up 和 /iot/down(因为后面 MQTTClient_subscribe 和 MQTTClient_publishMessage 用的就是 sub_topicpub_topic)。

3.4 步骤2:初始化 RPC 客户端
if (-1 == RPC_Client_Init())
{
    printf("RPC_Client_Init err\n");
    return -1;
}
  • 连接到本地的 RPC 服务(可能是一个后台进程,提供读写传感器、控制 GPIO 的能力)。

  • 失败则退出。

3.5 步骤3:创建 MQTT 客户端
if ((rc = MQTTClient_create(&client, address, clientid, MQTTCLIENT_PERSISTENCE_NONE, NULL)) != MQTTCLIENT_SUCCESS)
  • MQTTClient_create 初始化一个客户端实例。

    • &client:输出参数,保存句柄。

    • address:Broker 地址字符串(如 "tcp://192.168.1.100:1883")。

    • clientid:客户端标识符,必须唯一(Broker 用此区分不同设备)。

    • MQTTCLIENT_PERSISTENCE_NONE:不使用持久化(消息不存磁盘)。

    • NULL:持久化目录(不用)。

  • 返回值 rcMQTTCLIENT_SUCCESS(0)表示成功。

3.6 步骤3.1:设置回调函数
MQTTClient_setCallbacks(client, NULL, connlost, msgarrvd, delivered)
  • 参数:

    • client:客户端句柄。

    • NULL:上下文指针,会传给各个回调(本例不需要)。

    • connlost:连接丢失回调。

    • msgarrvd:消息到达回调。

msgarrvd 回调

你看到的代码中,除了 printf 打印消息内容,还做了:

  • 用 cJSON_Parse 解析 message->payload

  • 提取 params 对象和 LED1 字段

  • 调用 rpc_led_control 控制 LED

    这些不是库内部实现的,而是(开发者)自己写的业务逻辑。库只负责在收到消息时调用这个回调,并把消息数据通过参数传进来。

  • delivered:消息已发送回调。

    目前只打印了一条确认信息,没有其他动作。你也可以在里面添加自己的逻辑(例如记录日志、释放资源、触发下一次发送等)。

  • 必须在连接之前设置,回调才能生效。

3.7 步骤3.2:配置连接参数
conn_opts.keepAliveInterval = 100;    // 保活周期,单位秒。100秒内无任何报文则发送 PINGREQ
conn_opts.cleansession = 1;           // 1 = 清除之前的会话(订阅等),0 = 恢复之前会话
conn_opts.username = username;        // 用户名指针
conn_opts.password = password;        // 密码指针
  • conn_opts 必须先初始化为 MQTTClient_connectOptions_initializer,再修改所需字段。

  • 如果 Broker 不需要认证,可以设为 NULL

3.8 步骤4:连接服务器(带重试)
while (1)
{
    if ((rc = MQTTClient_connect(client, &conn_opts)) != MQTTCLIENT_SUCCESS)
    {
        printf("Failed to connect, return code %d\n", rc);
        // 不退出,继续重试
        sleep(2);
    }
    else
    {
        break;
    }
}
  • MQTTClient_connect 发起 TCP 连接和 MQTT 握手。

  • 失败则打印并循环(没有延时,会疯狂重试,最好加 sleep(1))。

  • 成功后跳出循环。

3.9 步骤5:订阅主题

c

if ((rc = MQTTClient_subscribe(client, sub_topic, QOS)) != MQTTCLIENT_SUCCESS)
{
    printf("Failed to subscribe, return code %d\n", rc);
    rc = EXIT_FAILURE;
}
else
{
    // 进入发布循环
}
  • MQTTClient_subscribe 发送 SUBSCRIBE 报文。

  • sub_topic:要订阅的主题(如 /iot/down)。

  • QOS:订阅的 QoS 等级(本例为 1,Broker 会保证至少一次送达)。

  • 订阅成功后,后续发布到该主题的消息会触发 msgarrvd

逻辑原因

1. 确保能收到下行消息

  • 设备需要接收云端或手机端发来的控制命令(如开关 LED)。

  • 这些命令可能随时到达,设备必须在连接成功后尽快注册自己的兴趣(即订阅主题),否则一旦有命令在设备订阅之前到达,设备就会错过。

  • 先订阅可以避免消息丢失(尤其是 QoS ≥ 1 时,服务器会为离线客户端保留消息的前提是会话被持久化,但你的代码中 cleansession=1,会话是全新的,先订阅才能保证后续消息到达)。

2. 发布是周期性的,订阅是一次性的

  • 发布(上报传感器数据)是在主循环中每隔 5 秒执行一次,可以随时开始。

  • 订阅通常只需要在连接成功后执行一次,之后一直有效。

  • 先完成订阅,再进入发布循环,符合逻辑顺序。

3. 大多数物联网场景的典型模式

  • 设备上电 → 连接 Broker → 订阅下行主题 → 进入主循环(发布上行数据 + 处理下行命令)。

  • 如果不先订阅,可能在发布几次之后才能收到命令,造成响应延迟。


技术上的可行性

  • MQTT 协议本身不强制订阅必须在发布之前。你可以先发布后订阅,甚至只发布不订阅。

  • 但如果设备需要接收消息,就必须在对方发送消息之前完成订阅。由于你无法预测控制命令何时下发,最保险的做法就是连接成功后立即订阅

3.10 发布循环内的操作
char humi, temp;
while (0 != rpc_dht11_read(&humi, &temp));   // 阻塞直到读取成功
  • rpc_dht11_read 通过 RPC 读取温湿度,返回 0 表示成功。

  • humi 和 temp 是 char 类型,湿度一般为 0~100,温度可能负值?实际应使用 int,但这里示例用了 char,可能有截断风险。

然后构造 JSON:

c

sprintf(buf, "{\"method\":\"report\",\"clientToken\":\"mqttx_6e973856\",\"timestamp\":1628646783,\"params\":{\"temp_value\":%d,\"humi_value\":%d}}", temp, humi);
  • 使用 sprintf 将 JSON 写入 buf。注意时间戳是固定的,实际应该动态获取。

填充发布消息结构体:

c

pubmsg.payload = buf;                // 指向 JSON 字符串
pubmsg.payloadlen = (int)strlen(buf);
pubmsg.qos = QOS;                    // 1
pubmsg.retained = 0;                 // 不保留消息,即不保存最新值给后续订阅者

发布:

if ((rc = MQTTClient_publishMessage(client, pub_topic, &pubmsg, &token)) != MQTTCLIENT_SUCCESS)
{
    printf("Failed to publish message, return code %d\n", rc);
    continue;
}
  • MQTTClient_publishMessage:将消息发送到 Broker。

  • pub_topic:主题(/iot/up)。

  • &token:输出参数,代表本次发布的传输令牌。

  • 返回值非 0 表示发布请求失败(网络断开、无效参数等)。

等待发布完成(仅对 QoS≥1 有意义):

c

rc = MQTTClient_waitForCompletion(client, token, TIMEOUT);
printf("Message with delivery token %d delivered\n", token);
  • MQTTClient_waitForCompletion 阻塞直到该 token 对应的消息被确认(收到 PUBACK)或超时。

  • TIMEOUT 为 10000 毫秒。

  • 对于 QoS 0,该函数会立即返回成功。

最后 sleep(5),每 5 秒上报一次。

3.11 清理(正常不会执行到)

c

MQTTClient_unsubscribe(client, sub_topic);
MQTTClient_disconnect(client, 10000);
MQTTClient_destroy(&client);
  • 由于发布循环是 while(1),这些代码不会执行。仅当订阅失败或循环被 break 时才会走到。实际产品中应加入退出机制(如检测到 'Q' 按键)。


三、参数配合和整体框架

各组件如何配合?

  1. MQTTClient_create + MQTTClient_setCallbacks
    → 创建客户端并注册回调,为后续收发消息做准备。

  2. MQTTClient_connect
    → 建立网络连接和 MQTT 会话,Broker 会验证 clientid/username/password

  3. MQTTClient_subscribe
    → 告诉 Broker:“请把主题 sub_topic 的消息转发给我”。
    之后 Broker 收到匹配 sub_topic 的消息时,就会推送给该客户端。

  4. msgarrvd 回调
    → 当 Broker 推送消息到达时,库自动调用它。在回调内解析 JSON 并执行动作(控制 LED)。

  5. MQTTClient_publishMessage
    → 发布消息到 pub_topic,Broker 会转发给所有订阅了该主题的客户端。

  6. MQTTClient_waitForCompletion
    → 确保(QoS≥1)消息已正确送达,否则可重试。

回调函数与主循环的关系

  • 主循环负责主动发布传感器数据。

  • 回调函数msgarrvd 由库内部线程(或主循环中的 MQTTClient_yield,本例未显式调用)触发。Paho 库通常会在 MQTTClient_connect 和 publishMessage 等函数内部自行处理网络 I/O 并分发回调。所以不需要额外的 MQTTClient_receive 调用。


四、整体框架总结(可复用的模板)

以下是编写类似 MQTT 客户端程序的通用步骤:

text

1. 读取配置(Broker地址、ClientID、用户名密码、业务主题参数)
2. 初始化硬件接口(RPC、GPIO、传感器驱动)
3. 创建 MQTT 客户端(MQTTClient_create)
4. 设置回调(connlost, msgarrvd, delivered)
5. 配置连接选项(keepAlive, cleansession, 用户名密码)
        cleansession会话状态的处理方式
标志值 会话行为 订阅是否持久 离线期间的消息是否接收
cleansession = 1 每次连接都是全新会话,之前的订阅、未确认的消息等全部丢弃 每次连接后必须重新订阅 不接收离线期间的消息
cleansession = 0 Broker 记住会话,包括订阅、未完成的 QoS 1/2 消息 订阅持续有效,断线重连后无需重新订阅 设备离线时到达的消息,Broker 会保留(取决于 QoS),重连后会收到
6. 循环连接直到成功(MQTTClient_connect)
7. 订阅下行主题(MQTTClient_subscribe)
8. 进入主循环:
   - 读取传感器数据
   - 构造 JSON 消息
   - 发布到上行主题(MQTTClient_publishMessage)
   - 等待送达确认(可选)
   - 延时或等待下次周期
9. (如果退出)取消订阅、断开连接、销毁客户端

五、编写类似程序时需要注意的要点

  1. 缓冲区大小:用宏定义统一长度,避免硬编码 1000。

  2. JSON 解析:使用 cJSON_ParseWithLength 传入 payloadlen,防止缺 '\0'

  3. 连接重试:重试循环内应加入延迟(sleep(1)),避免 CPU 彪高。

  4. 断开重连机制:在 connlost 回调中设置标志,主循环检测后重新调用 MQTTClient_connect

  5. 资源释放:在 msgarrvd 中必须调用 MQTTClient_freeMessage 和 MQTTClient_freeTopicName 防止内存泄漏。

  6. 主题设计:区分上行(设备→云)和下行(云→设备),保持一致性。

  7. 线程安全:如果主线程和多线程回调共享数据,需加锁。本例只有单线程回调(Paho 默认在调用 MQTTClient_connect/publish 的线程中处理 I/O),相对安全。

希望这个详尽的分析能帮助你完全理解代码,并自己编写出健壮的 MQTT 设备端程序。

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "MQTTClient.h"
#include "cJSON.h"
#include "rpc_client.h"
#include "cfg.h"

#define ADDRESS     "tcp://broker.emqx.io:1883"
#define CLIENTID    "mqttx_6e973856"
#define TOPIC_SUBSCRIBE     "/iot/down"
#define TOPIC_PUBLISH       "/iot/up"
#define QOS         1
#define TIMEOUT     10000L

#define USER_NAME "100ask"
#define PASSWORD  "100ask"

volatile MQTTClient_deliveryToken deliveredtoken;

//送达至服务器成功
void delivered(void *context, MQTTClient_deliveryToken dt)
{
    printf("Message with token value %d delivery confirmed\n", dt);
    deliveredtoken = dt;
}

//传给客户端成功 //由 MQTT 客户端库在收到消息时自动调用
int msgarrvd(void *context, char *topicName, int topicLen, MQTTClient_message *message)
{
    printf("Message arrived\n");
    printf("     topic: %s\n", topicName);
    printf("   message: %.*s\n", message->payloadlen, (char*)message->payload);

    /* 根据接收到的消息控制LED */
	/* 消息格式: {"method":"control","clientToken":"v2177557380FdRVD::DU2S%vwG7","params":{"LED1":0}} */
	cJSON *root = cJSON_Parse((char*)message->payload);  //有效数据
	cJSON *ptTemp = cJSON_GetObjectItem(root, "params"); 
	if (ptTemp)
	{
		cJSON *led1 = cJSON_GetObjectItem(ptTemp, "LED1");
		if (led1)
		{
			rpc_led_control(1,led1->valueint);
		}
	}
	cJSON_Delete(root);

    MQTTClient_freeMessage(&message);
    MQTTClient_free(topicName);
    return 1;
}
/*{
 *   "method": "control",
 *   "clientToken": "mqttx_6e973856",
 *   "params": {
 *     "LED1": 1
 *  }
 *}
*/ //下行设备(客户端)手机,可发送此命令控制LED

void connlost(void *context, char *cause)
{
    printf("\nConnection lost\n");
    if (cause)
    	printf("     cause: %s\n", cause);
}

int main(int argc, char* argv[])
{
    MQTTClient client;
    MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer;
    int rc;

	char URI[1000];
	char clientid[1000]; 
	char username[1000]; 
	char password[1000]; 
	char ProductID[1000]; 
	char DeviceName[1000];

	char address[1000];
    //发布主题+接收主题
	char pub_topic[1000]; /* $thing/up/property/{ProductID}/{DeviceName} */
	char sub_topic[1000]; /* $thing/down/property/{ProductID}/{DeviceName} */

    /* 1. 读配置文件 /etc/wechat.cfg得到URI,CLIENTID,USERNAME,PASSWD,ProductID,DeviceName */
	if ( 0 != read_cfg(URI, clientid, username, password, ProductID, DeviceName))
	{
		printf("read cfg err\n");
		return -1;
	}

    printf("%s\n%s\n%s\n%s\n%s\n%s\n",URI,clientid, username, password, ProductID, DeviceName);

    //strcpy(URI, ADDRESS);
    //strcpy(clientid, CLIENTID);
    //strcpy(username, "100ask");
    //strcpy(password, "100ask");



    //sprintf(address, "%s", URI);       
	sprintf(address, "tcp://%s:1883", URI);       //产品ID + 设备名称
                            //属性 
	//sprintf(pub_topic, "$thing/up/property/%s/%s", ProductID, DeviceName);
	//sprintf(sub_topic, "$thing/down/property/%s/%s", ProductID, DeviceName);

    strcpy(pub_topic, TOPIC_PUBLISH);
    strcpy(sub_topic, TOPIC_SUBSCRIBE);                            

    /* INIT RPC: connect RPC Server */    
    if (-1 == RPC_Client_Init()) // 2. 启动RPC客户端
	{
		printf("RPC_Client_Init err\n");
		return -1;
	}
	
    // 3. 创建MQTT客户端
    if ((rc = MQTTClient_create(&client, address, clientid,MQTTCLIENT_PERSISTENCE_NONE, NULL)) != MQTTCLIENT_SUCCESS)
    {
        printf("Failed to create client, return code %d\n", rc);
        rc = EXIT_FAILURE;
        goto exit;
    }

    // 3.1设置回调函数,打印客户端在连接过程中的错误
    if ((rc = MQTTClient_setCallbacks(client, NULL, connlost, msgarrvd, delivered)) != MQTTCLIENT_SUCCESS)
    {
        printf("Failed to set callbacks, return code %d\n", rc);
        rc = EXIT_FAILURE;
        goto destroy_exit;
    }

    // 3.2配置客户端连接参数
    conn_opts.keepAliveInterval = 100;
    conn_opts.cleansession = 1;
    conn_opts.username = username;
    conn_opts.password = password;

    while (1)
    {
                // 4.客户端连接服务器
	    if ((rc = MQTTClient_connect(client, &conn_opts)) != MQTTCLIENT_SUCCESS)
	    {
		    printf("Failed to connect, return code %d\n", rc);
		    rc = EXIT_FAILURE;
            sleep(1);
		//goto destroy_exit;
	    }
        else
        {
            break;
        }
    }

    printf("Subscribing to topic %s\nfor client %s using QoS%d\n\n"
           "Press Q<Enter> to quit\n\n", TOPIC_SUBSCRIBE, CLIENTID, QOS);

        //5.1订阅消息
    if ((rc = MQTTClient_subscribe(client, sub_topic, QOS)) != MQTTCLIENT_SUCCESS)
    {
    	printf("Failed to subscribe, return code %d\n", rc);
    	rc = EXIT_FAILURE;
    }
    else
    {
    	int ch;
        int cnt = 0;
        //初始化订阅服务
        MQTTClient_message pubmsg = MQTTClient_message_initializer;
        char buf[1000];
        MQTTClient_deliveryToken token; //送达
        
    	while (1)
    	{
            /* rpc_dht11_read */
			char humi, temp;
			
			while (0 !=rpc_dht11_read(&humi, &temp));

            sprintf(buf, "\
				{ \
					\"method\":\"report\",\
					\"clientToken\":\"mqttx_6e973856\",\
					\"timestamp\":1628646783,\
					\"params\":{\
						\"temp_value\":%d,\
					    \"humi_value\":%d\
					}\
				}", temp, humi);
				
            pubmsg.payload = buf;
            pubmsg.payloadlen = (int)strlen(buf);
            pubmsg.qos = QOS;
            pubmsg.retained = 0;

            /* 发布消息 */
                            //6.发布消息至服务器
            if ((rc = MQTTClient_publishMessage(client, pub_topic, &pubmsg, &token)) != MQTTCLIENT_SUCCESS)
            {
                 printf("Failed to publish message, return code %d\n", rc);
                 continue;
            }
                            //等待发布完成
            rc = MQTTClient_waitForCompletion(client, token, TIMEOUT);
            printf("Message with delivery token %d delivered\n", token);                        

			sleep(5);
    	} 
            //取消订阅
        if ((rc = MQTTClient_unsubscribe(client, sub_topic)) != MQTTCLIENT_SUCCESS)
        {
        	printf("Failed to unsubscribe, return code %d\n", rc);
        	rc = EXIT_FAILURE;
        }
    }

    if ((rc = MQTTClient_disconnect(client, 10000)) != MQTTCLIENT_SUCCESS)
    {
    	printf("Failed to disconnect, return code %d\n", rc);
    	rc = EXIT_FAILURE;
    }
destroy_exit:
    MQTTClient_destroy(&client);
exit:
    return rc;
}




#include <sys/types.h>
#include <sys/stat.h>
#include <fcntl.h>
#include <stdio.h>
#include <unistd.h>
#include "cfg.h"
#include "cJSON.h"

#define CFG_STR_LEN 256

int read_cfg(char *URI,char *clientid,char *username,char *password,char *ProductID,char *DeviceName)
{
    int ret;
    cJSON *ptTemp;
    char buf[1024];


    int fd = open(CFG_FILE,O_RDONLY);
    if(fd < 0)
    {
        perror("open(CFG_FILE) ");
        return -1;
    }

    ret = read(fd,buf,sizeof(buf));
    if(ret < 0)
    {
        close(fd);
        return -1;
    }    
    else 
        buf[ret] = '\0';
    
    cJSON *root = cJSON_Parse(buf);

    ptTemp = cJSON_GetObjectItem(root,"URI");
    if(ptTemp)
        snprintf(URI, CFG_STR_LEN,"%s",ptTemp->valuestring);

    ptTemp = cJSON_GetObjectItem(root,"clientid");
    if(ptTemp)
        snprintf(clientid, CFG_STR_LEN,"%s",ptTemp->valuestring); 

    ptTemp = cJSON_GetObjectItem(root,"username");
    if(ptTemp)
        snprintf(username, CFG_STR_LEN,"%s",ptTemp->valuestring);  

    ptTemp = cJSON_GetObjectItem(root,"password");
    if(ptTemp)
        snprintf(password, CFG_STR_LEN,"%s",ptTemp->valuestring); 

    ptTemp = cJSON_GetObjectItem(root,"ProductID");
        if(ptTemp)
            snprintf(ProductID, CFG_STR_LEN,"%s",ptTemp->valuestring); 

    ptTemp = cJSON_GetObjectItem(root,"DeviceName");
        if(ptTemp)
            snprintf(DeviceName, CFG_STR_LEN,"%s",ptTemp->valuestring); 


    return 0;
}

Logo

智能硬件社区聚焦AI智能硬件技术生态,汇聚嵌入式AI、物联网硬件开发者,打造交流分享平台,同步全国赛事资讯、开展 OPC 核心人才招募,助力技术落地与开发者成长。

更多推荐