Skip to content

概念先知道

  • MQTTS:走 TLS 加密通道的 MQTT(MQTT over TLS),默认端口 8883,防止消息被窃听/篡改。
  • TLS 握手:客户端与服务器先交换证书、协商密钥,确认身份并建立加密通道,之后 MQTT 报文在通道内传输。
  • 证书验证:例程用 mbedTLS,MBEDTLS_SSL_VERIFY_REQUIRED 要求服务器证书必须由内置 CA 签发(例程用 ssl/ca_1.crt 签发的测试证书)。
  • SNI:TLS 里告诉服务器“我要访问哪个域名”的字段,例程填 server.local,与测试证书 CN 对应。

例程功能简介

本页对应博流官方 SDK 的 wifi_mqtt 例程中的 MQTTS 发布端examples/wifi/sta/wifi_mqtt/wifi_mqtts_pub.c):

  • mqtts_pub <ip> <port> 命令启动;
  • tcp_client_connect 建立 TCP 连接(SO_REUSEADDR 可复用端口);
  • ssl_client_connect 完成 mbedTLS 初始化、CA 解析、握手与证书验证;
  • MQTTC_PAL_CONNTION_TYPE_TLS 类型创建 MQTT 客户端,mqtt_connect 连接 Broker;
  • 每 3 秒 mqtt_publish + mqtt_sync 发布 {"hello mqtt !"} 到主题 mqtt_test
  • 同族例程:同目录 mqtts_sub.c 提供订阅端命令 mqtts_sub;明文版见 MQTT 页(mqtt_pub / mqtt_sub)。

操作步骤

1
进入例程目录

在终端进入 SDK 的 MQTT 例程目录(前提:环境已按快速开始(Linux)Windows搭好):

cd examples/wifi/sta/wifi_mqtt
2
编译工程

统一填写 bl616(MQTTS 依赖 mbedTLS,例程 defconfig 已开启 CONFIG_MBEDTLS_V2):

make CHIP=bl616 BOARD=bl616dk
3
烧录固件

按住 BOOT 键短按 EN/RST 进入下载模式后烧录:

make flash CHIP=bl616 COMX=/dev/ttyUSB0
4
准备 TLS 版 Broker

例程默认使用 ssl/ 目录的测试证书(CA 的 CN 是 server.local),按 README_ssl.md 先在电脑上启动带 TLS 的 mosquitto:

mosquitto -c ssl/mos.conf
mosquitto_sub -h server.local -p 8883 --cafile ssl/ca_1.crt -t mqtt_test
5
连接 Wi-Fi 并发布消息

串口助手波特率 2000000,连接路由器后执行 mqtts_pub <服务器IP> <端口>(例程默认 8883):

wifi_sta_connect Your_SSID 12345678
mqtts_pub 192.168.1.143 8883
6
运行验证

日志依次出现 Performing the SSL/TLS handshakeVerifying peer X.509 certificatessl connect ok,随后每 3 秒发布一条消息;电脑订阅端同步收到即成功。

代码执行流程

MQTTS 发布端从启动到发布消息的完整流程如下:

例程调用的 API 介绍

mbedtls_x509_crt_parse(&ca_cert, ca_pem, ca_len)

解析 CA 证书(PEM 格式),用于后续验证服务器证书链。

参数

  • ca_certmbedtls_x509_crt 对象
  • ca_pem / ca_len:证书内容与长度(例程来自 certs.h

返回值0 成功;负值错误码

mbedtls_ssl_config_defaults / mbedtls_ssl_conf_authmode(...)

配置 TLS:客户端角色 + 流式传输 + 默认预设;MBEDTLS_SSL_VERIFY_REQUIRED 强制校验服务器证书。

参数

  • confmbedtls_ssl_config
  • MBEDTLS_SSL_IS_CLIENT:客户端模式
  • MBEDTLS_SSL_VERIFY_REQUIRED:必须验证

返回值0 成功;负值错误码

mbedtls_ssl_handshake(&ssl)

执行 TLS 握手。返回 MBEDTLS_ERR_SSL_WANT_READ/WRITE 时需继续调用直到完成。

参数

  • sslmbedtls_ssl_context

返回值0 握手完成;WANT_READ/WRITE 需重试;其他为错误

mqtt_init / mqtt_connect / mqtt_publish / mqtt_sync(...)

MQTT 客户端接口(与明文版相同),区别在 handle.type = MQTTC_PAL_CONNTION_TYPE_TLS,socket 句柄换成 TLS 通道。

参数

  • clientstruct mqtt_client
  • handlecustom_socket_handle,TLS 类型 + fd + ssl 上下文

返回值MQTT_OK 成功

完整代码

以下为 wifi_mqtts_pub.c 完整源码,与官方示例(examples/wifi/sta/wifi_mqtt)一致,默认折叠,点击展开:

📜 点击展开 wifi_mqtt/main.c 完整代码
c
/****************************************************************************
 *
 * Licensed to the Apache Software Foundation (ASF) under one or more
 * contributor license agreements.  See the NOTICE file distributed with
 * this work for additional information regarding copyright ownership.  The
 * ASF licenses this file to you under the Apache License, Version 2.0 (the
 * "License"); you may not use this file except in compliance with the
 * License.  You may obtain a copy of the License at
 *
 *   http://www.apache.org/licenses/LICENSE-2.0
 *
 * Unless required by applicable law or agreed to in writing, software
 * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
 * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.  See the
 * License for the specific language governing permissions and limitations
 * under the License.
 *
 ****************************************************************************/

/****************************************************************************
 * Included Files
 ****************************************************************************/

#include "FreeRTOS.h"
#include "task.h"
#include "timers.h"

#include <lwip/tcpip.h>
#include <lwip/sockets.h>
#include <lwip/netdb.h>

#include "wifi_mgmr_ext.h"

#include "bflb_irq.h"
#include "bflb_uart.h"

#include "rfparam_adapter.h"

#include "board.h"
#include "shell.h"

#ifndef BL602
#include "fhost_api.h"
#include "wifi_mgmr.h"
#endif
#include "async_event.h"
#include "mm.h"

#define DBG_TAG "MAIN"
#include "log.h"

struct bflb_device_s *gpio;

/****************************************************************************
 * Pre-processor Definitions
 ****************************************************************************/

#define WIFI_STACK_SIZE  (1536)
#define TASK_PRIORITY_FW (16)

/****************************************************************************
 * Private Types
 ****************************************************************************/

/****************************************************************************
 * Private Data
 ****************************************************************************/

static struct bflb_device_s *uart0;

static TaskHandle_t wifi_fw_task __attribute__((unused));


extern void shell_init_with_task(struct bflb_device_s *shell);
extern void wifi_event_handler(async_input_event_t ev, void *priv);
#ifdef BL602
extern void wifi_task_create(void);
extern int fhost_init(void);
extern int wifi_mgmr_task_start(void);
#endif

/****************************************************************************
 * Private Function Prototypes
 ****************************************************************************/

/****************************************************************************
 * Functions
 ****************************************************************************/

void wifi_start_firmware_task(void *param)
{
    LOG_I("Starting wifi ...\r\n");

    async_register_event_filter(EV_WIFI, wifi_event_handler, NULL);


    wifi_task_create();

    LOG_I("Starting fhost ...\r\n");
    fhost_init();

    vTaskDelete(NULL);
}

volatile uint32_t wifi_state = 0;
void wifi_event_handler(async_input_event_t ev, void *priv)
{
    uint32_t code = ev->code;

    switch (code) {
        case CODE_WIFI_ON_INIT_DONE: {
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_INIT_DONE\r\n", __func__);
            wifi_mgmr_task_start();
        } break;
        case CODE_WIFI_ON_MGMR_DONE: {
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_MGMR_DONE\r\n", __func__);
        } break;
        case CODE_WIFI_ON_SCAN_DONE: {
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_SCAN_DONE\r\n", __func__);
            wifi_mgmr_sta_scanlist();
        } break;
        case CODE_WIFI_ON_CONNECTED: {
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_CONNECTED\r\n", __func__);
            void mm_sec_keydump();
            mm_sec_keydump();
        } break;
        case CODE_WIFI_ON_GOT_IP: {
            wifi_state = 1;
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_GOT_IP\r\n", __func__);
            LOG_I("[SYS] Memory left is %d Bytes\r\n", kfree_size(0));
        } break;
        case CODE_WIFI_ON_DISCONNECT: {
            wifi_state = 0;
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_DISCONNECT\r\n", __func__);
        } break;
        case CODE_WIFI_ON_AP_STARTED: {
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_AP_STARTED\r\n", __func__);
        } break;
        case CODE_WIFI_ON_AP_STOPPED: {
            LOG_I("[APP] [EVT] %s, CODE_WIFI_ON_AP_STOPPED\r\n", __func__);
        } break;
        case CODE_WIFI_ON_AP_STA_ADD: {
            LOG_I("[APP] [EVT] [AP] [ADD] %lld\r\n", xTaskGetTickCount());
        } break;
        case CODE_WIFI_ON_AP_STA_DEL: {
            LOG_I("[APP] [EVT] [AP] [DEL] %lld\r\n", xTaskGetTickCount());
        } break;
        default: {
            LOG_I("[APP] [EVT] Unknown code %u \r\n", code);
        }
    }
}

int main(void)
{
    board_init();

    uart0 = bflb_device_get_by_name("uart0");
    shell_init_with_task(uart0);

    if (0 != rfparam_init(0, NULL, 0)) {
        LOG_I("PHY RF init failed!\r\n");
        return 0;
    }

    LOG_I("PHY RF init success!\r\n");

    tcpip_init(NULL, NULL);
    xTaskCreate(wifi_start_firmware_task, "wifi init", 1024, NULL, 10, NULL);

    vTaskStartScheduler();

    while (1) {
    }
}
📜 点击展开 wifi_mqtts_pub.c 完整代码
c
#include "FreeRTOS_POSIX.h"
#include <unistd.h>
#include <stdlib.h>
#include <stdio.h>
#include <sys/socket.h>
#include <sys/types.h>
#include <lwip/errno.h>
#include <netdb.h>
#if !defined(MBEDTLS_CONFIG_FILE)
#include "mbedtls/config.h"
#else
#include MBEDTLS_CONFIG_FILE
#endif

#include "mbedtls/debug.h"
#include "mbedtls/ssl.h"
#include "mbedtls/x509_crt.h"
#include "mbedtls/net_sockets.h"

#include "utils_getopt.h"

#include "mqtt.h"
#include "shell.h"
#include "certs.h"
#include "bflb_sec_trng.h"

static uint8_t sendbuf[2048]; /* sendbuf should be large enough to hold multiple whole mqtt messages */
static uint8_t recvbuf[1024]; /* recvbuf should be large enough any whole mqtt message expected to be received */

static shell_sig_func_ptr abort_exec;
static int test_sockfd;
static const char* addr;


typedef struct {
    int ssl_inited;
    mbedtls_net_context net;
    mbedtls_x509_crt ca_cert;
    mbedtls_x509_crt owncert;
    mbedtls_ssl_config conf;
    mbedtls_ssl_context ssl;
    mbedtls_pk_context pkey;
} ssl_param_t;

typedef struct {
    char *ca_cert;
    int ca_cert_len;
    char *own_cert;
    int own_cert_len;
    char *private_cert;
    int private_cert_len;

    char **alpn;
    int alpn_num;

    char *psk;
    int psk_len;
    char *pskhint;
    int pskhint_len;

    char *sni;
} ssl_conn_param_t;

static int ssl_random(void *prng, unsigned char *output, size_t output_len)
{
    (void)prng;
    bflb_trng_readlen(output, output_len);
    return 0;
}


static int tcp_client_connect(const char *ip, const char *port_str)
{
    int fd;
    int res;

    struct sockaddr_in addr;
    ip4_addr_t remote_ip;
    int port = atoi(port_str);

    ip4addr_aton(ip, &remote_ip);

    printf("tcp client connect %s:%d\r\n", ip4addr_ntoa(&remote_ip), port);

    {
        if ( (fd =  socket(AF_INET, SOCK_STREAM, 0))  < 0) {
            printf("socket create failed\r\n");
            return -2;
        }

        memset(&addr, 0, sizeof(addr));
        addr.sin_family = AF_INET;
        addr.sin_len = sizeof(addr);
        addr.sin_port = htons(port);
        //addr.sin_addr.s_addr = ((struct in_addr *) hostinfo->h_addr)->s_addr;
        addr.sin_addr.s_addr = ip4_addr_get_u32(&remote_ip);
    }

    printf("tcp_client_connect fd:%d\r\n", fd);

    int on= 1;
    res = setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &on, sizeof(on) );
    if (res != 0) {
        printf("setsockopt failed, res:%d\r\n", res);
    }

    res = connect(fd, (struct sockaddr *)&addr, sizeof(addr));
    if (res < 0) {
        printf("connect failed, res:%d\r\n", res);
        close(fd);
    }

    return fd;
}

static int ssl_client_connect(const char *ip, const char *port, ssl_param_t *ssl_param)
{
    int fd, ret;
    ssl_conn_param_t param = {0};
    //int mbedtls_platform_set_printf(int (*printf_func)(const char *, ...));
    //mbedtls_platform_set_printf(printf);

    memset(ssl_param, 0, sizeof(*ssl_param));
    memset(&param, 0, sizeof(param));

    param.ca_cert = ca_cert;
    param.ca_cert_len = sizeof(ca_cert);

    param.own_cert = own_cert;
    param.own_cert_len = sizeof(own_cert);

    param.private_cert = private_cert;
    param.private_cert_len = sizeof(private_cert);

    fd = tcp_client_connect(ip, port);
    if (fd < 0) {
        printf("tcp_client_connect fd:%d\r\n", fd);
        return -1;
    }

    param.sni = "server.local";

    ssl_param->ssl_inited = 1;
    /*
     * Initialize the connection
     */
    ssl_param->net.fd = fd;

    mbedtls_ssl_config_init(&ssl_param->conf);
    mbedtls_ssl_init(&ssl_param->ssl);
    mbedtls_x509_crt_init(&ssl_param->ca_cert);
    mbedtls_x509_crt_init(&ssl_param->owncert);
    mbedtls_pk_init(&ssl_param->pkey);

    ret = mbedtls_x509_crt_parse(&ssl_param->ca_cert, (unsigned char *)param.ca_cert, (size_t)param.ca_cert_len);
    if (ret < 0) {
        printf("[MBEDTLS] ssl connect: root parse failed- 0x%x\r\n", -ret);
        goto err;
    }
    // ret = mbedtls_x509_crt_parse(&ssl_param->owncert, (unsigned char *)param.own_cert, (size_t)param.own_cert_len);
    // if (ret < 0) {
    //  printf("[MBEDTLS] ssl connect: x509 parse failed- 0x%x\r\n", -ret);
    //  goto err;
    // }
    // ret = mbedtls_pk_parse_key(&ssl_param->pkey, (unsigned char *)param.private_cert, param.private_cert_len, NULL, 0);
    // if (ret != 0) {
    //  printf("[MBEDTLS] ssl connect: x509 parse failed- 0x%x\r\n", -ret);
    //     goto err;
    // }
    ret = mbedtls_ssl_config_defaults(&ssl_param->conf, MBEDTLS_SSL_IS_CLIENT, MBEDTLS_SSL_TRANSPORT_STREAM, MBEDTLS_SSL_PRESET_DEFAULT);
    if (ret != 0) {
        printf("[MBEDTLS] ssl connect: x509 config failed- 0x%x\r\n", -ret);
        goto err;
    }
    //mbedtls_ssl_conf_authmode(&ssl_param->conf, MBEDTLS_SSL_VERIFY_NONE);
    mbedtls_ssl_conf_authmode(&ssl_param->conf, MBEDTLS_SSL_VERIFY_REQUIRED);
    mbedtls_ssl_conf_ca_chain(&ssl_param->conf, &ssl_param->ca_cert, NULL);
    // mbedtls_ssl_conf_own_cert(&ssl_param->conf, &ssl_param.owncert, &ssl_param.pkey);
    mbedtls_ssl_conf_rng(&ssl_param->conf, ssl_random, NULL);
    //mbedtls_ssl_conf_dbg(&ssl_param->conf, my_debug, NULL);
    //mbedtls_ssl_conf_alpn_protocols(&ssl_param->conf, (const char **)alpn_str);
    //mbedtls_ssl_conf_psk(&ssl_param->conf, (const unsigned char *)"psk123456", 10, (const unsigned char *)"identity", 9);
    if((ret = mbedtls_ssl_setup(&ssl_param->ssl, &ssl_param->conf)) != 0) {
        printf("mbedtls ssl_setup fail \r\n", -ret);
        goto err;
    }

    if (param.sni) {
        mbedtls_ssl_set_hostname(&ssl_param->ssl, param.sni);
    }
    mbedtls_ssl_set_bio(&ssl_param->ssl, &ssl_param->net, mbedtls_net_send, mbedtls_net_recv, NULL);

    /*
     * handshake
     */
    printf("[MBEDTLS] Performing the SSL/TLS handshake ... \r\n");
    while ((ret = mbedtls_ssl_handshake(&ssl_param->ssl)) != 0) {
        if ((ret != MBEDTLS_ERR_SSL_WANT_READ) && (ret != MBEDTLS_ERR_SSL_WANT_WRITE)) {
            printf("[MBEDTLS] ssl connect: mbedtls_ssl_handshake returned -0x%x\r\n", -ret);
            goto err;
        }
    }

    /*
     * verify the server certificate
     */
    printf("[MBEDTLS] ...... Verifying peer X.509 certificate ... \r\n");
    ret = mbedtls_ssl_get_verify_result(&ssl_param->ssl);
    if (ret != 0) {
        printf("[MBEDTLS] ssl connect: verify result not confirmed - %d\r\n", -ret);
        goto err;
    }

    printf("[MBEDTLS] ssl connect ok\r\n");

        /* make non-blocking */
    if (fd != -1) {
        int iMode = 1;
        ioctlsocket(fd, FIONBIO, &iMode);
    }


        return fd;

err:
        mbedtls_ssl_close_notify(&ssl_param->ssl);
    mbedtls_x509_crt_free(&ssl_param->ca_cert);
    mbedtls_ssl_free(&ssl_param->ssl);
    mbedtls_ssl_config_free(&ssl_param->conf);

    close(ssl_param->net.fd);
        return -1;
}

/**
 * @brief The function will be called whenever a PUBLISH message is received.
 */
static void publish_callback_1(void** unused, struct mqtt_response_publish *published);

static ssl_param_t g_ssl_param = {0};
static void test_close(int sig)
{
    if (test_sockfd)
    {
        close(test_sockfd);
    }

    printf("mqtt_pub stop publish to %s\r\n", addr);

    abort_exec(sig);
}

static int example_mqtt(int argc, const char *argv[])
{
    const char* port = NULL;
    const char* topic;

    int ret = 0;
    // int argc = 0;

    abort_exec = shell_signal(1, test_close);

    /* get address (argv[1] if present) */
    if (argc > 1) {
        addr = argv[1];
    }

    /* get port number (argv[2] if present) */
    if (argc > 2) {
        port = argv[2];
    }

    /* get the topic name to publish */
        topic = "mqtt_test";

    /* open the non-blocking TCP socket (connecting to the broker) */
        test_sockfd = ssl_client_connect(addr, port, &g_ssl_param);
        struct custom_socket_handle handle;
        handle.type = MQTTC_PAL_CONNTION_TYPE_TLS;
        handle.ctx.fd = test_sockfd;
        handle.ctx.ssl_ctx = &g_ssl_param.ssl;

    if (test_sockfd < 0) {
        printf("Failed to open socket: %d\r\n", test_sockfd);
        test_close(SHELL_SIGINT);
    }

    /* setup a client */
    struct mqtt_client client;

    mqtt_init(&client, &handle, sendbuf, sizeof(sendbuf), recvbuf, sizeof(recvbuf), publish_callback_1);
    /* Create an anonymous session */
    const char* client_id = NULL;
    /* Ensure we have a clean session */
    uint8_t connect_flags = MQTT_CONNECT_CLEAN_SESSION;
    /* Send connection request to the broker. */
    ret = mqtt_connect(&client, client_id, NULL, NULL, 0, NULL, NULL, connect_flags, 400);

    if (ret != MQTT_OK)
    {
        printf("fail \r\n");
    }
    /* check that we don't have any errors */
    if (client.error != MQTT_OK) {
        printf("error: %s\r\n", mqtt_error_str(client.error));
        test_close(SHELL_SIGINT);
    }

    printf("%s is ready to begin publishing hello.\r\n", argv[0]);

    while(1) {
        /* print a message */
        char application_message[256] = {"{\"hello mqtt !\"}\r\n"};
        printf("%s published : \"%s\"\r\n", argv[0], application_message);

        mqtt_publish(&client, topic, application_message, strlen(application_message) + 1, MQTT_PUBLISH_QOS_0);
        mqtt_sync(&client);

        /* check for errors */
        if (client.error != MQTT_OK) {
            printf("error: %s\r\n", mqtt_error_str(client.error));
            test_close(SHELL_SIGINT);
        }
        vTaskDelay(3000);
    }

    /* disconnect */
    /* exit */
    test_close(SHELL_SIGINT);

    return 0;
}

static void publish_callback_1(void** unused, struct mqtt_response_publish *published)
{
    /* not used in this example */
}

#ifdef CONFIG_SHELL
#include <shell.h>

extern uint32_t wifi_state;
static int check_wifi_state(void)
{
    if (wifi_state == 1)
    {
        return 0;
    } else {
        return 1;
    }
}

static int cmd_mqtts_publisher(int argc, const char **argv)
{
    uint32_t ret = 0;

    ret = check_wifi_state();
    if (ret != 0) {
        printf("your wifi not connected!\r\n");
        return 0;
    }

    // xTaskCreate(example_mqtt,(char*)"test_mqtt", 8192, argv, 10, NULL);
    example_mqtt(argc, argv);

    return 0;
}

SHELL_CMD_EXPORT_ALIAS(cmd_mqtts_publisher, mqtts_pub, mqtts publisher);
#endif
📜 点击展开 wifi_mqtt/wifi_mqtts_sub.c 完整代码
c

/**
 * @file
 * A simple program that subscribes to a topic.
 */
#include "FreeRTOS_POSIX.h"
#include <unistd.h>
#include <stdlib.h>
#include <stdio.h>
#include <sys/socket.h>
#include <sys/types.h>
#include <lwip/errno.h>
#include <netdb.h>
#if !defined(MBEDTLS_CONFIG_FILE)
#include "mbedtls/config.h"
#else
#include MBEDTLS_CONFIG_FILE
#endif

#include "mbedtls/debug.h"
#include "mbedtls/ssl.h"
#include "mbedtls/x509_crt.h"
#include "mbedtls/net_sockets.h"

#include "utils_getopt.h"

#include "mqtt.h"
#include "shell.h"
#include "certs.h"
#include "bflb_sec_trng.h"

static uint8_t sendbuf[2048]; /* sendbuf should be large enough to hold multiple whole mqtt messages */
static uint8_t recvbuf[1024]; /* recvbuf should be large enough any whole mqtt message expected to be received */

static shell_sig_func_ptr abort_exec;
static TaskHandle_t client_daemon;
static int test_sockfd;
static const char* addr;

typedef struct {
    int ssl_inited;
    mbedtls_net_context net;
    mbedtls_x509_crt ca_cert;
    mbedtls_x509_crt owncert;
    mbedtls_ssl_config conf;
    mbedtls_ssl_context ssl;
    mbedtls_pk_context pkey;
} ssl_param_t;

typedef struct {
    char *ca_cert;
    int ca_cert_len;
    char *own_cert;
    int own_cert_len;
    char *private_cert;
    int private_cert_len;

    char **alpn;
    int alpn_num;

    char *psk;
    int psk_len;
    char *pskhint;
    int pskhint_len;

    char *sni;
} ssl_conn_param_t;

static int ssl_random(void *prng, unsigned char *output, size_t output_len)
{
    (void)prng;
    bflb_trng_readlen(output, output_len);
    return 0;
}

static int tcp_client_connect(const char *ip, const char *port_str)
{
    int fd;
    int res;

    struct sockaddr_in addr;
    ip4_addr_t remote_ip;
    int port = atoi(port_str);

    ip4addr_aton(ip, &remote_ip);

    printf("tcp client connect %s:%d\r\n", ip4addr_ntoa(&remote_ip), port);

    {
        if ( (fd =  socket(AF_INET, SOCK_STREAM, 0))  < 0) {
            printf("socket create failed\r\n");
            return -2;
        }

        memset(&addr, 0, sizeof(addr));
        addr.sin_family = AF_INET;
        addr.sin_len = sizeof(addr);
        addr.sin_port = htons(port);
        //addr.sin_addr.s_addr = ((struct in_addr *) hostinfo->h_addr)->s_addr;
        addr.sin_addr.s_addr = ip4_addr_get_u32(&remote_ip);
    }

    printf("tcp_client_connect fd:%d\r\n", fd);

    int on= 1;
    res = setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &on, sizeof(on) );
    if (res != 0) {
        printf("setsockopt failed, res:%d\r\n", res);
    }

    res = connect(fd, (struct sockaddr *)&addr, sizeof(addr));
    if (res < 0) {
        printf("connect failed, res:%d\r\n", res);
        close(fd);
    }

    return fd;
}

static int ssl_client_connect(const char *ip, const char *port, ssl_param_t *ssl_param)
{
    int fd, ret;
    ssl_conn_param_t param = {0};
    //int mbedtls_platform_set_printf(int (*printf_func)(const char *, ...));
    //mbedtls_platform_set_printf(printf);

    memset(ssl_param, 0, sizeof(*ssl_param));
    memset(&param, 0, sizeof(param));

    param.ca_cert = ca_cert;
    param.ca_cert_len = sizeof(ca_cert);

    param.own_cert = own_cert;
    param.own_cert_len = sizeof(own_cert);

    param.private_cert = private_cert;
    param.private_cert_len = sizeof(private_cert);

    fd = tcp_client_connect(ip, port);
    if (fd < 0) {
        printf("tcp_client_connect fd:%d\r\n", fd);
        return -1;
    }

    param.sni = "server.local";

    ssl_param->ssl_inited = 1;
    /*
     * Initialize the connection
     */
    ssl_param->net.fd = fd;

    mbedtls_ssl_config_init(&ssl_param->conf);
    mbedtls_ssl_init(&ssl_param->ssl);
    mbedtls_x509_crt_init(&ssl_param->ca_cert);
    mbedtls_x509_crt_init(&ssl_param->owncert);
    mbedtls_pk_init(&ssl_param->pkey);

    ret = mbedtls_x509_crt_parse(&ssl_param->ca_cert, (unsigned char *)param.ca_cert, (size_t)param.ca_cert_len);
    if (ret < 0) {
        printf("[MBEDTLS] ssl connect: root parse failed- 0x%x\r\n", -ret);
        goto err;
    }
    // ret = mbedtls_x509_crt_parse(&ssl_param->owncert, (unsigned char *)param.own_cert, (size_t)param.own_cert_len);
    // if (ret < 0) {
    //  printf("[MBEDTLS] ssl connect: x509 parse failed- 0x%x\r\n", -ret);
    //  goto err;
    // }
    // ret = mbedtls_pk_parse_key(&ssl_param->pkey, (unsigned char *)param.private_cert, param.private_cert_len, NULL, 0);
    // if (ret != 0) {
    //  printf("[MBEDTLS] ssl connect: x509 parse failed- 0x%x\r\n", -ret);
    //     goto err;
    // }
    ret = mbedtls_ssl_config_defaults(&ssl_param->conf, MBEDTLS_SSL_IS_CLIENT, MBEDTLS_SSL_TRANSPORT_STREAM, MBEDTLS_SSL_PRESET_DEFAULT);
    if (ret != 0) {
        printf("[MBEDTLS] ssl connect: x509 config failed- 0x%x\r\n", -ret);
        goto err;
    }
    //mbedtls_ssl_conf_authmode(&ssl_param->conf, MBEDTLS_SSL_VERIFY_NONE);
    mbedtls_ssl_conf_authmode(&ssl_param->conf, MBEDTLS_SSL_VERIFY_REQUIRED);
    mbedtls_ssl_conf_ca_chain(&ssl_param->conf, &ssl_param->ca_cert, NULL);
    // mbedtls_ssl_conf_own_cert(&ssl_param->conf, &ssl_param.owncert, &ssl_param.pkey);
    mbedtls_ssl_conf_rng(&ssl_param->conf, ssl_random, NULL);
    //mbedtls_ssl_conf_dbg(&ssl_param->conf, my_debug, NULL);
    //mbedtls_ssl_conf_alpn_protocols(&ssl_param->conf, (const char **)alpn_str);
    //mbedtls_ssl_conf_psk(&ssl_param->conf, (const unsigned char *)"psk123456", 10, (const unsigned char *)"identity", 9);
    if((ret = mbedtls_ssl_setup(&ssl_param->ssl, &ssl_param->conf)) != 0) {
        printf("mbedtls ssl_setup fail \r\n", -ret);
        goto err;
    }

    if (param.sni) {
        mbedtls_ssl_set_hostname(&ssl_param->ssl, param.sni);
    }
    mbedtls_ssl_set_bio(&ssl_param->ssl, &ssl_param->net, mbedtls_net_send, mbedtls_net_recv, NULL);

    /*
     * handshake
     */
    printf("[MBEDTLS] Performing the SSL/TLS handshake ... \r\n");
    while ((ret = mbedtls_ssl_handshake(&ssl_param->ssl)) != 0) {
        if ((ret != MBEDTLS_ERR_SSL_WANT_READ) && (ret != MBEDTLS_ERR_SSL_WANT_WRITE)) {
            printf("[MBEDTLS] ssl connect: mbedtls_ssl_handshake returned -0x%x\r\n", -ret);
            goto err;
        }
    }

    /*
     * verify the server certificate
     */
    printf("[MBEDTLS] ...... Verifying peer X.509 certificate ... \r\n");
    ret = mbedtls_ssl_get_verify_result(&ssl_param->ssl);
    if (ret != 0) {
        printf("[MBEDTLS] ssl connect: verify result not confirmed - %d\r\n", -ret);
        goto err;
    }

    printf("[MBEDTLS] ssl connect ok\r\n");

        /* make non-blocking */
    if (fd != -1) {
        int iMode = 1;
        ioctlsocket(fd, FIONBIO, &iMode);
    }


        return fd;

err:
        mbedtls_ssl_close_notify(&ssl_param->ssl);
    mbedtls_x509_crt_free(&ssl_param->ca_cert);
    mbedtls_ssl_free(&ssl_param->ssl);
    mbedtls_ssl_config_free(&ssl_param->conf);

    close(ssl_param->net.fd);
        return -1;
}

/**
 * @brief The function will be called whenever a PUBLISH message is received.
 */
static void publish_callback_1(void** unused, struct mqtt_response_publish *published);

/**
 * @brief The client's refresher. This function triggers back-end routines to
 *        handle ingress/egress traffic to the broker.
 *
 * @note All this function needs to do is call \ref __mqtt_recv and
 *       \ref __mqtt_send every so often. I've picked 100 ms meaning that
 *       client ingress/egress traffic will be handled every 100 ms.
 */
static void client_refresher(void* client);

/**
 * @brief Safelty closes the \p sockfd and cancels the \p client_daemon before \c exit.
 */
static ssl_param_t g_ssl_param = {0};
static void test_close(int sig)
{
    if (test_sockfd)
    {
        close(test_sockfd);
    }
    printf("mqtt_sub disconnecting from %s\r\n", addr);

    abort_exec(sig);

    vTaskDelete(client_daemon);

}

static int example_mqtt(int argc, const char *argv[])
{
    const char* port = NULL;
    const char* topic;

    int ret = 0;
    // int argc = 0;

    abort_exec = shell_signal(1, test_close);

    /* get address (argv[1] if present) */
    if (argc > 1) {
        addr = argv[1];
    }

    /* get port number (argv[2] if present) */
    if (argc > 2) {
        port = argv[2];
    }

    /* get the topic name to publish */
        topic = "mqtt_test";

    /* open the non-blocking TCP socket (connecting to the broker) */
    test_sockfd = ssl_client_connect(addr, port, &g_ssl_param);
    struct custom_socket_handle handle;
    handle.type = MQTTC_PAL_CONNTION_TYPE_TLS;
    handle.ctx.fd = test_sockfd;
    handle.ctx.ssl_ctx = &g_ssl_param.ssl;

    if (test_sockfd < 0) {
        printf("Failed to open socket: %d\r\n", test_sockfd);
        test_close(SHELL_SIGINT);
    }

    /* setup a client */
    struct mqtt_client client;

    mqtt_init(&client, &handle, sendbuf, sizeof(sendbuf), recvbuf, sizeof(recvbuf), publish_callback_1);
    /* Create an anonymous session */
    const char* client_id = NULL;
    /* Ensure we have a clean session */
    uint8_t connect_flags = MQTT_CONNECT_CLEAN_SESSION;
    /* Send connection request to the broker. */
    ret = mqtt_connect(&client, client_id, NULL, NULL, 0, NULL, NULL, connect_flags, 400);

    if (ret != MQTT_OK)
    {
        printf("fail \r\n");
    }
    /* check that we don't have any errors */
    if (client.error != MQTT_OK) {
        printf("error: %s\r\n", mqtt_error_str(client.error));
        test_close(SHELL_SIGINT);
    }

    /* start a thread to refresh the client (handle egress and ingree client traffic) */
    xTaskCreate(client_refresher, (char*)"client_ref", 1024,  &client, 10, &client_daemon);

    /* subscribe */
    mqtt_subscribe(&client, topic, 0);

    printf("%s listening for '%s' messages.\r\n", argv[0], topic);

    /* block wait CTRL-C exit */
    while(1) {
        vTaskDelay(100);
    }

    /* disconnect */
    /* exit */
    test_close(SHELL_SIGINT);

    return 0;
}

static void publish_callback_1(void** unused, struct mqtt_response_publish *published)
{
    /* note that published->topic_name is NOT null-terminated (here we'll change it to a c-string) */
    char* topic_name = (char*) malloc(published->topic_name_size + 1);
    char* topic_msg = (char*) malloc(published->application_message_size + 1);
    if (topic_name && topic_msg) {
        memcpy(topic_name, published->topic_name, published->topic_name_size);
        topic_name[published->topic_name_size] = '\0';

        memcpy(topic_msg, published->application_message, published->application_message_size);
        topic_msg[published->application_message_size] = '\0';

        printf("Received publish('%s'): %s\r\n", topic_name, topic_msg);
    } else {
        printf("No memory to receive published msg\r\n");
    }

    if(topic_name) {
        free(topic_name);
    }
    if(topic_msg) {
        free(topic_msg);
    }
}

static void client_refresher(void* client)
{
    while(1)
    {
        mqtt_sync((struct mqtt_client*) client);
        vTaskDelay(100);
    }

}

#ifdef CONFIG_SHELL
#include <shell.h>

extern uint32_t wifi_state;
static int check_wifi_state(void)
{
    if (wifi_state == 1)
    {
        return 0;
    } else {
        return 1;
    }
}

static int cmd_mqtts_subscribe(int argc, const char **argv)
{
    uint32_t ret = 0;

    ret = check_wifi_state();
    if (ret != 0) {
        printf("your wifi not connected!\r\n");
        return 0;
    }

    // xTaskCreate(example_mqtt,(char*)"test_mqtt", 8192, argv, 10, NULL);
    example_mqtt(argc, argv);

    return 0;
}

SHELL_CMD_EXPORT_ALIAS(cmd_mqtts_subscribe, mqtts_sub, mqtts subscribe);
#endif

FAQ

TLS 握手失败 / 证书验证不通过

按 README_ssl.md 在电脑 /etc/hosts 里把 server.local 指向本机,并确保 mosquitto 的 TLS 配置(ssl/mos.conf)使用的证书由 ssl/ca_1.crt 签发;模组内置的 CA 是例程自带的测试 CA,连其他公网 Broker 需要换 CA。

想连公共 MQTTS 服务器(如 test.mosquitto.org:8883)

需要把该服务器的 CA 证书写进 certs.h 并重新编译,同时把 SNI 改成对应域名;证书内容超长时注意 sendbuf/内存占用。这也是例程默认使用本地测试证书的原因。

遇到问题?

如有其他问题,请到统一的提问与讨论区:Ai-Thinker Discussions

Released under the MIT License. Build Time 2026-09-11 14:52:23