引言

在物联网系统中,远程控制是最核心的业务能力之一。从智能家居的"开灯/关灯",到工业设备的"启停电机",再到农业灌溉的"开阀/关阀",所有这些场景都依赖一个安全可靠的远程指令通道。

我们沧州虎王科技技术团队在为某智慧农业园区设计灌溉控制系统时,遇到了一系列远程控制的工程挑战:4G网络不稳定导致的指令丢失、指令被篡改引发的安全风险、高并发控制请求下的系统过载、离线设备指令的缓存与重发策略。这些问题如果解决不好,轻则控制指令延迟,重则引发安全事故。

本文将从架构设计到代码实现,系统性地讲解如何基于ESP32构建一套安全、可靠、高效的远程指令系统。内容涵盖指令协议设计、传输安全保障、离线指令处理、并发控制和权限验证等工程实践。


一、远程控制通道整体架构

1.1 架构设计

远程控制系统的核心挑战在于:在不可靠的网络环境下实现可靠的控制指令传输。我们的架构采用三层设计:指令生成层、指令传输层和指令执行层。

┌──────────────────────────────────────────────────────────────┐
│           ESP32 物联网远程控制通道架构                         │
├──────────────────────────────────────────────────────────────┤
│                                                              │
│  ┌─── 云端: 指令生成层 ───────────────────────────────┐     │
│  │                                                      │   │
│  │  ┌──────────┐  ┌──────────┐  ┌──────────┐         │   │
│  │  │ Web控制台 │  │ 移动APP  │  │ API调用  │         │   │
│  │  └─────┬─────┘  └─────┬────┘  └─────┬────┘         │   │
│  │        │              │              │               │   │
│  │  ┌─────▼──────────────▼──────────────▼─────┐        │   │
│  │  │           指令调度服务                    │        │   │
│  │  │  • 权限验证  • 指令加密  • 优先级排序    │        │   │
│  │  │  • 离线缓存  • 重发策略  • 结果回调      │        │   │
│  │  └────────────────────┬──────────────────────┘        │   │
│  │                       │                                │   │
│  └───────────────────────┼────────────────────────────────┘   │
│                          │ MQTT/TLS                          │
│  ┌───────────────────────┼────────────────────────────────┐   │
│  │                       ▼                                │   │
│  │  ┌─── 边缘: 指令传输层 (ESP32) ───────────────┐       │   │
│  │  │                                              │       │   │
│  │  │  ┌──────────────────────────────────┐      │       │   │
│  │  │  │     MQTT客户端 + TLS加密通道      │      │       │   │
│  │  │  └──────────────┬───────────────────┘      │       │   │
│  │  │                 │                           │       │   │
│  │  │  ┌──────────────▼───────────────────┐      │       │   │
│  │  │  │     指令解析与验证模块             │      │       │   │
│  │  │  │  • 签名校验 • 时序检查 • 去重    │      │       │   │
│  │  │  └──────────────┬───────────────────┘      │       │   │
│  │  │                 │                           │       │   │
│  │  │  ┌──────────────▼───────────────────┐      │       │   │
│  │  │  │     指令执行队列                  │      │       │   │
│  │  │  │  • 优先级调度 • 并发控制          │      │       │   │
│  │  │  └──────────────┬───────────────────┘      │       │   │
│  │  └─────────────────┼──────────────────────────┘       │   │
│  │                    │                                    │   │
│  └────────────────────┼────────────────────────────────────┘   │
│                       │                                        │
│  ┌────────────────────▼── 设备: 指令执行层 ───────────────┐  │
│  │                                                        │  │
│  │  ┌─────────┐ ┌─────────┐ ┌─────────┐ ┌─────────┐    │  │
│  │  │继电器   │ │电机控制 │ │阀门控制 │ │传感器   │    │  │
│  │  │GPIO输出 │ │PWM输出  │ │I2C/SPI  │ │ADC读取  │    │  │
│  │  └─────────┘ └─────────┘ └─────────┘ └─────────┘    │  │
│  │                                                        │  │
│  └────────────────────────────────────────────────────────┘  │
│                                                              │
└──────────────────────────────────────────────────────────────┘

1.2 指令协议设计

#include "cJSON.h"
#include "mbedtls/base64.h"
#include "mbedtls/sha256.h"
#include "mbedtls/pk.h"

/* 指令类型定义 */
typedef enum {
    CMD_TYPE_ACTUATOR = 0x01,    /* 执行器控制 */
    CMD_TYPE_CONFIG   = 0x02,    /* 配置修改 */
    CMD_TYPE_QUERY    = 0x03,    /* 数据查询 */
    CMD_TYPE_OTA      = 0x04,    /* 固件升级 */
    CMD_TYPE_DIAG     = 0x05,    /* 诊断命令 */
    CMD_TYPE_SYSTEM   = 0x06,    /* 系统命令(重启等) */
} cmd_type_t;

/* 指令优先级 */
typedef enum {
    CMD_PRIORITY_LOW    = 0,
    CMD_PRIORITY_NORMAL  = 1,
    CMD_PRIORITY_HIGH    = 2,
    CMD_PRIORITY_URGENT = 3,    /* 紧急:插队执行 */
} cmd_priority_t;

/* 指令结构体 */
typedef struct {
    uint32_t  cmd_id;           /* 指令唯一ID */
    uint8_t   cmd_type;         /* 指令类型 */
    uint8_t   priority;         /* 优先级 */
    uint8_t   target_device;    /* 目标设备ID */
    uint8_t   reserved;
    uint64_t  timestamp;        /* 服务端时间戳(ms) */
    uint64_t  expire_time;      /* 过期时间(ms),0=不过期 */
    uint16_t  payload_len;      /* 负载长度 */
    uint8_t   signature[32];   /* SHA-256签名 */
    uint8_t   payload[];        /* 柔性数组,指令负载 */
} __attribute__((packed)) remote_cmd_t;

/* 指令响应结构 */
typedef struct {
    uint32_t  cmd_id;           /* 对应的指令ID */
    uint8_t   result_code;      /* 执行结果码 */
    uint8_t   reserved[3];
    uint64_t  execute_time;     /* 执行时间戳 */
    uint16_t  data_len;         /* 返回数据长度 */
    uint8_t   data[];           /* 返回数据 */
} __attribute__((packed)) cmd_response_t;

/* 结果码定义 */
#define CMD_RESULT_SUCCESS         0x00
#define CMD_RESULT_FAIL            0x01
#define CMD_RESULT_TIMEOUT         0x02
#define CMD_RESULT_INVALID         0x03
#define CMD_RESULT_UNAUTHORIZED    0x04
#define CMD_RESULT_EXPIRED         0x05
#define CMD_RESULT_BUSY            0x06
#define CMD_RESULT_NOT_SUPPORTED  0x07

二、MQTT安全传输通道

2.1 MQTT TLS连接配置

#include "mqtt_client.h"
#include "esp_certs.h"

/* MQTT配置 */
#define MQTT_BROKER_URI   "mqtts://broker.example.com:8883"
#define MQTT_CMD_TOPIC    "device/cmd/+/+"    /* 订阅指令主题 */
#define MQTT_RESP_TOPIC   "device/resp/"      /* 发布响应主题 */
#define MQTT_HEARTBEAT_TOPIC "device/heartbeat/"

static esp_mqtt_client_handle_t g_mqtt_client = NULL;

/* 设备证书 (实际项目中应存储在NVS或eFuse中) */
extern const uint8_t ca_cert_pem_start[] asm("_binary_ca_cert_pem_start");
extern const uint8_t ca_cert_pem_end[]   asm("_binary_ca_cert_pem_end");
extern const uint8_t client_cert_start[] asm("_binary_client_cert_pem_start");
extern const uint8_t client_cert_end[]   asm("_binary_client_cert_pem_end");
extern const uint8_t client_key_start[]  asm("_binary_client_key_pem_start");
extern const uint8_t client_key_end[]    asm("_binary_client_key_pem_end");

/**
 * @brief 初始化MQTT安全连接
 */
esp_err_t init_mqtt_secure_channel(const char *device_id)
{
    char cmd_topic[64];
    char resp_topic[64];
    char hb_topic[64];

    snprintf(cmd_topic, sizeof(cmd_topic), "device/cmd/%s/+", device_id);
    snprintf(resp_topic, sizeof(resp_topic), "device/resp/%s", device_id);
    snprintf(hb_topic, sizeof(hb_topic), "device/heartbeat/%s", device_id);

    /* MQTT配置:TLS双向认证 */
    esp_mqtt_client_config_t mqtt_cfg = {
        .uri = MQTT_BROKER_URI,
        .port = 8883,

        /* TLS证书配置 */
        .cert_pem = (const char *)ca_cert_pem_start,
        .client_cert = (const char *)client_cert_start,
        .client_key = (const char *)client_key_start,

        /* 客户端ID */
        .client_id = device_id,

        /* 认证信息 */
        .username = device_id,
        .password = get_device_token(),  /* 从NVS读取设备Token */

        /* Keep Alive */
        .keepalive = 60,
        .disable_keepalive = false,

        /* 重连配置 */
        .reconnect_timeout_ms = 5000,
        .network_timeout_ms = 10000,

        /* 最后遗嘱消息 */
        .lwt_topic = hb_topic,
        .lwt_msg = "{\"status\":\"offline\"}",
        .lwt_qos = 1,
        .lwt_retain = true,

        /* 缓冲区配置 */
        .buffer_size = 2048,
        .buffer_out_size = 2048,

        /* QoS配置 */
        .disable_clean_session = false,
    };

    g_mqtt_client = esp_mqtt_client_init(&mqtt_cfg);
    if (g_mqtt_client == NULL) {
        ESP_LOGE(TAG, "MQTT客户端初始化失败");
        return ESP_FAIL;
    }

    /* 注册事件回调 */
    esp_mqtt_client_register_event(g_mqtt_client, ESP_MQTT_EVENT_ANY,
                                    mqtt_event_handler, NULL);

    ESP_LOGI(TAG, "MQTT安全通道初始化完成, broker=%s", MQTT_BROKER_URI);
    return ESP_OK;
}

2.2 MQTT事件处理

/* MQTT事件处理回调 */
static void mqtt_event_handler(void *handler_args, esp_event_base_t base,
                                int32_t event_id, void *event_data)
{
    esp_mqtt_event_handle_t event = event_data;

    switch ((esp_mqtt_event_id_t)event_id) {
    case MQTT_EVENT_CONNECTED:
        ESP_LOGI(TAG, "MQTT已连接到Broker");
        g_mqtt_connected = true;

        /* 订阅指令主题 */
        esp_mqtt_client_subscribe(g_mqtt_client, "device/cmd/+/+", 1);
        /* 订阅广播主题 */
        esp_mqtt_client_subscribe(g_mqtt_client, "device/broadcast", 0);

        /* 发布在线状态 */
        publish_online_status();
        break;

    case MQTT_EVENT_DISCONNECTED:
        ESP_LOGW(TAG, "MQTT连接断开");
        g_mqtt_connected = false;
        /* 断线时切换到离线指令缓存模式 */
        enable_offline_cache_mode();
        break;

    case MQTT_EVENT_DATA:
        ESP_LOGD(TAG, "MQTT收到消息: topic=%.*s, len=%d",
                 event->topic_len, event->topic, event->data_len);
        /* 处理收到的指令 */
        handle_incoming_mqtt_message(event->topic, event->topic_len,
                                      event->data, event->data_len);
        break;

    case MQTT_EVENT_ERROR:
        ESP_LOGE(TAG, "MQTT错误: type=%d", event->error_handle->error_type);
        if (event->error_handle->error_type == MQTT_ERROR_TYPE_ESP_TLS) {
            ESP_LOGE(TAG, "  TLS错误: %s",
                     esp_err_to_name(event->error_handle->esp_tls_last_esp_err));
        }
        break;

    default:
        ESP_LOGD(TAG, "MQTT事件: %d", (int)event_id);
        break;
    }
}

/**
 * @brief 处理收到的MQTT消息
 */
static void handle_incoming_mqtt_message(const char *topic, int topic_len,
                                          const char *data, int data_len)
{
    /* 解析主题,提取指令类型 */
    char topic_buf[128];
    if (topic_len >= sizeof(topic_buf)) {
        ESP_LOGW(TAG, "主题过长");
        return;
    }
    memcpy(topic_buf, topic, topic_len);
    topic_buf[topic_len] = '\0';

    /* 主题格式: device/cmd/{device_id}/{cmd_category} */
    char *tokens[5];
    int token_count = 0;
    char *tok = strtok(topic_buf, "/");
    while (tok && token_count < 5) {
        tokens[token_count++] = tok;
        tok = strtok(NULL, "/");
    }

    if (token_count < 3) {
        ESP_LOGW(TAG, "无效指令主题: %s", topic_buf);
        return;
    }

    const char *cmd_category = tokens[2];

    /* 根据类别处理 */
    if (strcmp(cmd_category, "actuator") == 0) {
        process_actuator_command(data, data_len);
    } else if (strcmp(cmd_category, "config") == 0) {
        process_config_command(data, data_len);
    } else if (strcmp(cmd_category, "query") == 0) {
        process_query_command(data, data_len);
    } else if (strcmp(cmd_category, "system") == 0) {
        process_system_command(data, data_len);
    } else {
        ESP_LOGW(TAG, "未知指令类别: %s", cmd_category);
    }
}

三、指令安全验证

3.1 指令签名验证

#include "mbedtls/md.h"
#include "mbedtls/pk.h"

/* 设备密钥 (从安全存储中加载) */
static uint8_t g_device_secret[32];

/**
 * @brief 加载设备密钥
 */
esp_err_t load_device_secret(void)
{
    nvs_handle_t handle;
    esp_err_t ret = nvs_open("security", NVS_READONLY, &handle);
    if (ret != ESP_OK) {
        ESP_LOGE(TAG, "无法打开安全NVS: %s", esp_err_to_name(ret));
        return ret;
    }

    size_t secret_len = sizeof(g_device_secret);
    ret = nvs_get_blob(handle, "device_secret", g_device_secret, &secret_len);
    nvs_close(handle);

    if (ret != ESP_OK || secret_len != 32) {
        ESP_LOGE(TAG, "设备密钥加载失败");
        return ESP_FAIL;
    }

    ESP_LOGI(TAG, "设备密钥加载成功");
    return ESP_OK;
}

/**
 * @brief 验证指令签名 (HMAC-SHA256)
 * @param cmd 指令指针
 * @return true=验证通过
 */
bool verify_command_signature(const remote_cmd_t *cmd, size_t total_len)
{
    /* 构造待验签数据: cmd_id + type + priority + target + timestamp + payload */
    uint8_t sign_data[256];
    size_t sign_data_len = 0;

    /* 拷贝头部字段(不含签名) */
    memcpy(sign_data, &cmd->cmd_id, sizeof(cmd->cmd_id));
    sign_data_len += sizeof(cmd->cmd_id);
    memcpy(sign_data + sign_data_len, &cmd->cmd_type, 1);
    sign_data_len += 1;
    memcpy(sign_data + sign_data_len, &cmd->priority, 1);
    sign_data_len += 1;
    memcpy(sign_data + sign_data_len, &cmd->target_device, 1);
    sign_data_len += 1;
    memcpy(sign_data + sign_data_len, &cmd->timestamp, sizeof(cmd->timestamp));
    sign_data_len += sizeof(cmd->timestamp);

    /* 拷贝负载 */
    size_t payload_offset = offsetof(remote_cmd_t, payload);
    if (total_len > payload_offset) {
        size_t payload_len = total_len - payload_offset;
        if (sign_data_len + payload_len > sizeof(sign_data)) {
            ESP_LOGE(TAG, "验签数据过长");
            return false;
        }
        memcpy(sign_data + sign_data_len, cmd->payload, payload_len);
        sign_data_len += payload_len;
    }

    /* 计算HMAC-SHA256 */
    uint8_t calculated_hmac[32];
    const mbedtls_md_info_t *md_info = mbedtls_md_info_from_type(MBEDTLS_MD_SHA256);

    int ret = mbedtls_md_hmac(md_info,
                               g_device_secret, sizeof(g_device_secret),
                               sign_data, sign_data_len,
                               calculated_hmac);
    if (ret != 0) {
        ESP_LOGE(TAG, "HMAC计算失败: -0x%04X", -ret);
        return false;
    }

    /* 比较签名 (使用恒定时间比较防止时序攻击) */
    int diff = 0;
    for (int i = 0; i < 32; i++) {
        diff |= calculated_hmac[i] ^ cmd->signature[i];
    }

    if (diff != 0) {
        ESP_LOGW(TAG, "指令签名验证失败, cmd_id=%lu",
                 (unsigned long)cmd->cmd_id);
        return false;
    }

    ESP_LOGD(TAG, "指令签名验证通过, cmd_id=%lu",
             (unsigned long)cmd->cmd_id);
    return true;
}

3.2 时序与重放保护

/* 已处理指令ID缓存 (防重放攻击) */
#define PROCESSED_CMD_CACHE_SIZE 32

typedef struct {
    uint32_t cmd_ids[PROCESSED_CMD_CACHE_SIZE];
    uint8_t  write_idx;
} cmd_id_cache_t;

static cmd_id_cache_t g_cmd_cache;
static portMUX_TYPE g_cmd_cache_lock = portMUX_INITIALIZER_UNLOCKED;

/* 最近时间戳偏差容忍度 */
#define MAX_TIMESTAMP_SKEW_MS  300000  /* 5分钟 */

/**
 * @brief 检查指令时序有效性
 */
bool verify_command_timing(const remote_cmd_t *cmd)
{
    uint64_t current_time = get_synced_time_ms(); /* NTP同步时间 */

    /* 1. 检查时间戳偏差 */
    int64_t skew = (int64_t)current_time - (int64_t)cmd->timestamp;
    if (skew > MAX_TIMESTAMP_SKEW_MS) {
        ESP_LOGW(TAG, "指令时间戳过期: skew=%lld ms", (long long)skew);
        return false;
    }
    if (skew < -MAX_TIMESTAMP_SKEW_MS) {
        ESP_LOGW(TAG, "指令时间戳超前: skew=%lld ms (可能是重放)", (long long)skew);
        return false;
    }

    /* 2. 检查是否过期 */
    if (cmd->expire_time > 0 && current_time > cmd->expire_time) {
        ESP_LOGW(TAG, "指令已过期: expire=%llu, now=%llu",
                 (unsigned long long)cmd->expire_time,
                 (unsigned long long)current_time);
        return false;
    }

    return true;
}

/**
 * @brief 检查指令是否已处理过(防重放)
 */
bool is_command_duplicate(uint32_t cmd_id)
{
    taskENTER_CRITICAL(&g_cmd_cache_lock);
    for (int i = 0; i < PROCESSED_CMD_CACHE_SIZE; i++) {
        if (g_cmd_cache.cmd_ids[i] == cmd_id) {
            taskEXIT_CRITICAL(&g_cmd_cache_lock);
            ESP_LOGW(TAG, "检测到重复指令: cmd_id=%lu", (unsigned long)cmd_id);
            return true;
        }
    }
    /* 记录新指令ID */
    g_cmd_cache.cmd_ids[g_cmd_cache.write_idx] = cmd_id;
    g_cmd_cache.write_idx = (g_cmd_cache.write_idx + 1) % PROCESSED_CMD_CACHE_SIZE;
    taskEXIT_CRITICAL(&g_cmd_cache_lock);
    return false;
}

四、指令执行引擎

4.1 优先级队列与执行调度

#include "freertos/queue.h"

/* 指令执行上下文 */
typedef struct {
    remote_cmd_t *cmd;          /* 指令数据 */
    size_t         cmd_len;      /* 指令总长度 */
    TaskHandle_t   caller_task;  /* 调用方任务 */
} cmd_context_t;

/* 优先级队列 (4个优先级,每个一个队列) */
#define CMD_QUEUE_DEPTH  16
static QueueHandle_t g_cmd_queues[4] = {NULL};

/* 执行器注册表 */
typedef esp_err_t (*cmd_executor_fn)(const uint8_t *payload, uint16_t len,
                                       uint8_t *resp_data, uint16_t *resp_len);

typedef struct {
    uint8_t          cmd_type;
    cmd_executor_fn  executor;
    const char      *name;
    uint32_t         timeout_ms;   /* 执行超时 */
    uint8_t          max_concurrent; /* 最大并发数 */
    uint8_t          current_count;  /* 当前执行中数 */
} cmd_executor_entry_t;

/* 注册的执行器列表 */
#define MAX_EXECUTORS 16
static cmd_executor_entry_t g_executors[MAX_EXECUTORS];
static uint8_t g_executor_count = 0;
static SemaphoreHandle_t g_executor_mutex;

/**
 * @brief 注册指令执行器
 */
esp_err_t register_cmd_executor(uint8_t cmd_type, cmd_executor_fn fn,
                                 const char *name, uint32_t timeout_ms,
                                 uint8_t max_concurrent)
{
    xSemaphoreTake(g_executor_mutex, portMAX_DELAY);

    if (g_executor_count >= MAX_EXECUTORS) {
        xSemaphoreGive(g_executor_mutex);
        return ESP_ERR_NO_MEM;
    }

    g_executors[g_executor_count] = (cmd_executor_entry_t){
        .cmd_type = cmd_type,
        .executor = fn,
        .name = name,
        .timeout_ms = timeout_ms,
        .max_concurrent = max_concurrent,
        .current_count = 0,
    };
    g_executor_count++;

    xSemaphoreGive(g_executor_mutex);
    ESP_LOGI(TAG, "执行器注册: type=%d, name=%s, timeout=%lu ms",
             cmd_type, name, (unsigned long)timeout_ms);
    return ESP_OK;
}

/**
 * @brief 初始化指令执行引擎
 */
esp_err_t init_cmd_execution_engine(void)
{
    /* 创建4个优先级队列 */
    for (int i = 0; i < 4; i++) {
        g_cmd_queues[i] = xQueueCreate(CMD_QUEUE_DEPTH, sizeof(cmd_context_t));
        if (g_cmd_queues[i] == NULL) {
            ESP_LOGE(TAG, "指令队列 %d 创建失败", i);
            return ESP_ERR_NO_MEM;
        }
    }

    g_executor_mutex = xSemaphoreCreateMutex();

    /* 注册内置执行器 */
    register_cmd_executor(CMD_TYPE_ACTUATOR, execute_actuator_cmd,
                          "actuator", 5000, 1);
    register_cmd_executor(CMD_TYPE_CONFIG, execute_config_cmd,
                          "config", 3000, 1);
    register_cmd_executor(CMD_TYPE_QUERY, execute_query_cmd,
                          "query", 2000, 4);
    register_cmd_executor(CMD_TYPE_OTA, execute_ota_cmd,
                          "ota", 300000, 1);
    register_cmd_executor(CMD_TYPE_SYSTEM, execute_system_cmd,
                          "system", 10000, 1);

    ESP_LOGI(TAG, "指令执行引擎初始化完成");
    return ESP_OK;
}

/**
 * @brief 调度指令到对应优先级队列
 */
esp_err_t dispatch_command(remote_cmd_t *cmd, size_t cmd_len)
{
    cmd_context_t ctx = {
        .cmd = cmd,
        .cmd_len = cmd_len,
        .caller_task = xTaskGetCurrentTaskHandle(),
    };

    if (cmd->priority > 3) {
        cmd->priority = 3;
    }

    /* 紧急指令直接插队到队列头部 */
    if (cmd->priority == CMD_PRIORITY_URGENT) {
        BaseType_t ret = xQueueSendToFront(g_cmd_queues[cmd->priority],
                                            &ctx, pdMS_TO_TICKS(100));
        if (ret != pdPASS) {
            ESP_LOGE(TAG, "紧急指令入队失败");
            return ESP_FAIL;
        }
    } else {
        BaseType_t ret = xQueueSend(g_cmd_queues[cmd->priority],
                                     &ctx, pdMS_TO_TICKS(100));
        if (ret != pdPASS) {
            ESP_LOGW(TAG, "指令入队失败(队列满), priority=%d", cmd->priority);
            return ESP_FAIL;
        }
    }

    ESP_LOGD(TAG, "指令已调度: cmd_id=%lu, priority=%d, queue_depth=%d",
             (unsigned long)cmd->cmd_id, cmd->priority,
             (int)uxQueueMessagesWaiting(g_cmd_queues[cmd->priority]));
    return ESP_OK;
}

4.2 指令执行任务

/* 指令执行任务:从优先级队列取出并执行 */
void cmd_executor_task(void *arg)
{
    uint8_t priority = (uint8_t)(uintptr_t)arg;
    cmd_context_t ctx;
    uint8_t resp_buf[256];

    ESP_LOGI(TAG, "指令执行任务启动, priority=%d", priority);

    while (1) {
        /* 从对应优先级队列取出指令 */
        if (xQueueReceive(g_cmd_queues[priority], &ctx,
                           portMAX_DELAY) == pdPASS) {
            remote_cmd_t *cmd = ctx.cmd;

            ESP_LOGI(TAG, "开始执行指令: id=%lu, type=%d, priority=%d",
                     (unsigned long)cmd->cmd_id, cmd->cmd_type, priority);

            /* 查找执行器 */
            cmd_executor_entry_t *exec = find_executor(cmd->cmd_type);
            if (exec == NULL) {
                ESP_LOGE(TAG, "未找到执行器: type=%d", cmd->cmd_type);
                send_cmd_response(cmd->cmd_id, CMD_RESULT_NOT_SUPPORTED,
                                   NULL, 0);
                goto cleanup;
            }

            /* 并发数检查 */
            xSemaphoreTake(g_executor_mutex, portMAX_DELAY);
            if (exec->current_count >= exec->max_concurrent) {
                xSemaphoreGive(g_executor_mutex);
                ESP_LOGW(TAG, "执行器并发已满: %s (%d/%d)",
                         exec->name, exec->current_count, exec->max_concurrent);
                send_cmd_response(cmd->cmd_id, CMD_RESULT_BUSY, NULL, 0);
                goto cleanup;
            }
            exec->current_count++;
            xSemaphoreGive(g_executor_mutex);

            /* 执行指令 (带超时监控) */
            uint16_t resp_len = 0;
            uint64_t start_time = esp_timer_get_time();

            esp_err_t result = exec->executor(
                cmd->payload, cmd->payload_len,
                resp_buf, &resp_len);

            uint64_t elapsed_ms = (esp_timer_get_time() - start_time) / 1000;

            /* 更新并发计数 */
            xSemaphoreTake(g_executor_mutex, portMAX_DELAY);
            exec->current_count--;
            xSemaphoreGive(g_executor_mutex);

            /* 发送响应 */
            uint8_t result_code = (result == ESP_OK) ?
                CMD_RESULT_SUCCESS : CMD_RESULT_FAIL;
            send_cmd_response(cmd->cmd_id, result_code,
                               resp_buf, resp_len);

            ESP_LOGI(TAG, "指令执行完成: id=%lu, result=%d, 耗时=%llu ms",
                     (unsigned long)cmd->cmd_id, result_code,
                     (unsigned long long)elapsed_ms);

cleanup:
            /* 释放指令内存 */
            free(ctx.cmd);
        }
    }
}

/**
 * @brief 发送指令执行响应
 */
void send_cmd_response(uint32_t cmd_id, uint8_t result_code,
                        const uint8_t *data, uint16_t data_len)
{
    cmd_response_t resp = {
        .cmd_id = cmd_id,
        .result_code = result_code,
        .execute_time = get_synced_time_ms(),
        .data_len = data_len,
    };

    /* 构造JSON响应 */
    cJSON *json = cJSON_CreateObject();
    cJSON_AddNumberToObject(json, "cmd_id", cmd_id);
    cJSON_AddNumberToObject(json, "result", result_code);
    cJSON_AddNumberToObject(json, "timestamp", resp.execute_time);

    if (data && data_len > 0) {
        /* Base64编码返回数据 */
        size_t b64_len = 0;
        mbedtls_base64_encode(NULL, 0, &b64_len, data, data_len);
        char *b64_buf = malloc(b64_len + 1);
        if (b64_buf) {
            mbedtls_base64_encode((unsigned char *)b64_buf, b64_len + 1,
                                  &b64_len, data, data_len);
            b64_buf[b64_len] = '\0';
            cJSON_AddStringToObject(json, "data", b64_buf);
            free(b64_buf);
        }
    }

    char *json_str = cJSON_PrintUnformatted(json);
    if (json_str && g_mqtt_connected) {
        char topic[64];
        snprintf(topic, sizeof(topic), "device/resp/%s", g_device_id);
        esp_mqtt_client_publish(g_mqtt_client, topic, json_str, 0, 1, 0);
        ESP_LOGD(TAG, "响应已发送: %s", json_str);
    }

    cJSON_free(json_str);
    cJSON_Delete(json);
}

五、执行器实现示例

5.1 继电器控制执行器

/* 继电器GPIO映射表 */
typedef struct {
    uint8_t  device_id;
    uint8_t  gpio_pin;
    bool     active_low;    /* 低电平触发 */
    bool     current_state;
} relay_config_t;

static relay_config_t g_relays[] = {
    {1, GPIO_NUM_26, false, false},  /* 灌溉阀门1 */
    {2, GPIO_NUM_27, false, false},  /* 灌溉阀门2 */
    {3, GPIO_NUM_14, false, false},  /* 水泵 */
    {4, GPIO_NUM_12, true,  false},  /* 风机 */
};

#define RELAY_COUNT (sizeof(g_relays) / sizeof(g_relays[0]))

/**
 * @brief 执行器控制指令处理器
 * 负载格式(JSON): {"device":1,"action":"on","duration":0}
 */
esp_err_t execute_actuator_cmd(const uint8_t *payload, uint16_t len,
                                uint8_t *resp_data, uint16_t *resp_len)
{
    /* 解析JSON负载 */
    cJSON *json = cJSON_ParseWithLength((const char *)payload, len);
    if (json == NULL) {
        ESP_LOGE(TAG, "执行器指令JSON解析失败");
        return ESP_ERR_INVALID_ARG;
    }

    cJSON *device = cJSON_GetObjectItem(json, "device");
    cJSON *action = cJSON_GetObjectItem(json, "action");
    cJSON *duration = cJSON_GetObjectItem(json, "duration");

    if (!cJSON_IsNumber(device) || !cJSON_IsString(action)) {
        cJSON_Delete(json);
        return ESP_ERR_INVALID_ARG;
    }

    uint8_t dev_id = device->valueint;
    const char *action_str = action->valuestring;
    uint32_t dur_ms = cJSON_IsNumber(duration) ? duration->valuedouble * 1000 : 0;

    /* 查找继电器 */
    relay_config_t *relay = NULL;
    for (int i = 0; i < RELAY_COUNT; i++) {
        if (g_relays[i].device_id == dev_id) {
            relay = &g_relays[i];
            break;
        }
    }

    if (relay == NULL) {
        ESP_LOGE(TAG, "未找到设备: id=%d", dev_id);
        cJSON_Delete(json);
        return ESP_ERR_NOT_FOUND;
    }

    /* 执行动作 */
    bool target_state;
    if (strcmp(action_str, "on") == 0) {
        target_state = true;
    } else if (strcmp(action_str, "off") == 0) {
        target_state = false;
    } else if (strcmp(action_str, "toggle") == 0) {
        target_state = !relay->current_state;
    } else {
        ESP_LOGE(TAG, "未知动作: %s", action_str);
        cJSON_Delete(json);
        return ESP_ERR_INVALID_ARG;
    }

    /* 设置GPIO */
    int level = relay->active_low ? !target_state : target_state;
    gpio_set_level(relay->gpio_pin, level);
    relay->current_state = target_state;

    ESP_LOGI(TAG, "执行器控制: device=%d, action=%s, gpio=%d, level=%d",
             dev_id, action_str, relay->gpio_pin, level);

    /* 定时关闭 */
    if (dur_ms > 0 && target_state) {
        /* 创建定时关闭任务或使用软件定时器 */
        TimerHandle_t timer = xTimerCreate("relay_off",
                                            pdMS_TO_TICKS(dur_ms),
                                            pdFALSE,
                                            (void *)(uintptr_t)relay->gpio_pin,
                                            relay_off_timer_cb);
        if (timer) {
            xTimerStart(timer, 0);
            ESP_LOGI(TAG, "定时关闭: %lu ms后", (unsigned long)dur_ms);
        }
    }

    /* 构造响应 */
    cJSON *resp = cJSON_CreateObject();
    cJSON_AddNumberToObject(resp, "device", dev_id);
    cJSON_AddStringToObject(resp, "state", target_state ? "on" : "off");
    char *resp_str = cJSON_PrintUnformatted(resp);
    if (resp_str) {
        size_t slen = strlen(resp_str);
        if (slen < 256) {
            memcpy(resp_data, resp_str, slen);
            *resp_len = slen;
        }
        cJSON_free(resp_str);
    }
    cJSON_Delete(resp);
    cJSON_Delete(json);

    return ESP_OK;
}

/* 定时关闭回调 */
static void relay_off_timer_cb(TimerHandle_t timer)
{
    uint8_t pin = (uint8_t)(uintptr_t)pvTimerGetTimerID(timer);
    /* 查找对应的继电器并关闭 */
    for (int i = 0; i < RELAY_COUNT; i++) {
        if (g_relays[i].gpio_pin == pin) {
            int level = g_relays[i].active_low ? 1 : 0;
            gpio_set_level(pin, level);
            g_relays[i].current_state = false;
            ESP_LOGI(TAG, "定时关闭: device=%d, gpio=%d",
                     g_relays[i].device_id, pin);
            break;
        }
    }
    xTimerDelete(timer, 0);
}

六、离线指令缓存与重发

6.1 本地指令缓存

当设备网络断开时,收到的指令无法通过MQTT传输。此时需要在本地缓存指令,待网络恢复后执行。

/* 离线指令缓存 (存储到SPI Flash的NVS分区) */
#define OFFLINE_CMD_MAX  10

typedef struct {
    uint32_t cmd_id;
    uint64_t timestamp;
    uint16_t cmd_len;
    uint8_t  cmd_data[512];  /* 压缩后的指令数据 */
} offline_cmd_entry_t;

typedef struct {
    uint8_t              count;
    offline_cmd_entry_t  entries[OFFLINE_CMD_MAX];
} offline_cmd_cache_t;

static offline_cmd_cache_t g_offline_cache;
static bool g_offline_mode = false;

/**
 * @brief 缓存离线指令
 */
esp_err_t cache_offline_command(const remote_cmd_t *cmd, size_t cmd_len)
{
    if (g_offline_cache.count >= OFFLINE_CMD_MAX) {
        ESP_LOGW(TAG, "离线指令缓存已满,丢弃最旧指令");
        /* 移除最旧指令 */
        memmove(&g_offline_cache.entries[0],
                &g_offline_cache.entries[1],
                sizeof(offline_cmd_entry_t) * (OFFLINE_CMD_MAX - 1));
        g_offline_cache.count--;
    }

    /* 压缩并存储指令 */
    offline_cmd_entry_t *entry = &g_offline_cache.entries[g_offline_cache.count];
    entry->cmd_id = cmd->cmd_id;
    entry->timestamp = cmd->timestamp;

    size_t compressed_len = sizeof(entry->cmd_data);
    /* 这里可以添加压缩逻辑,如LZ4或Heatshrink */
    if (cmd_len > sizeof(entry->cmd_data)) {
        ESP_LOGE(TAG, "指令数据过大: %d > %d", (int)cmd_len,
                 (int)sizeof(entry->cmd_data));
        return ESP_ERR_NO_MEM;
    }
    memcpy(entry->cmd_data, cmd, cmd_len);
    entry->cmd_len = cmd_len;

    g_offline_cache.count++;
    ESP_LOGI(TAG, "离线指令已缓存: cmd_id=%lu, 缓存数=%d",
             (unsigned long)cmd->cmd_id, g_offline_cache.count);

    /* 持久化到NVS */
    save_offline_cache_to_nvs();

    return ESP_OK;
}

/**
 * @brief 网络恢复后执行缓存的指令
 */
void process_offline_cache(void)
{
    if (g_offline_cache.count == 0) {
        ESP_LOGI(TAG, "无离线指令需要处理");
        return;
    }

    ESP_LOGI(TAG, "开始处理 %d 条离线指令", g_offline_cache.count);

    for (int i = 0; i < g_offline_cache.count; i++) {
        offline_cmd_entry_t *entry = &g_offline_cache.entries[i];
        remote_cmd_t *cmd = (remote_cmd_t *)entry->cmd_data;

        /* 检查指令是否过期 */
        if (cmd->expire_time > 0 &&
            get_synced_time_ms() > cmd->expire_time) {
            ESP_LOGW(TAG, "离线指令已过期: cmd_id=%lu",
                     (unsigned long)cmd->cmd_id);
            send_cmd_response(cmd->cmd_id, CMD_RESULT_EXPIRED, NULL, 0);
            continue;
        }

        /* 重新验证签名 */
        if (!verify_command_signature(cmd, entry->cmd_len)) {
            ESP_LOGW(TAG, "离线指令签名验证失败: cmd_id=%lu",
                     (unsigned long)cmd->cmd_id);
            send_cmd_response(cmd->cmd_id, CMD_RESULT_UNAUTHORIZED, NULL, 0);
            continue;
        }

        /* 调度执行 */
        remote_cmd_t *cmd_copy = malloc(entry->cmd_len);
        if (cmd_copy) {
            memcpy(cmd_copy, cmd, entry->cmd_len);
            dispatch_command(cmd_copy, entry->cmd_len);
        }
    }

    /* 清空缓存 */
    g_offline_cache.count = 0;
    save_offline_cache_to_nvs();
    ESP_LOGI(TAG, "离线指令处理完成");
}

6.2 指令重发机制

/* 指令重发管理 */
typedef struct {
    uint32_t  cmd_id;
    uint8_t   retry_count;
    uint32_t  last_send_time;
    uint32_t  ack_timeout;
} pending_cmd_t;

#define MAX_PENDING_CMDS  8
#define MAX_RETRY_COUNT   3
#define RETRY_BASE_DELAY  2000  /* 2秒基础延迟 */

static pending_cmd_t g_pending_cmds[MAX_PENDING_CMDS];

/**
 * @brief 记录已发送指令(等待响应)
 */
void track_pending_command(uint32_t cmd_id)
{
    for (int i = 0; i < MAX_PENDING_CMDS; i++) {
        if (g_pending_cmds[i].cmd_id == 0) {
            g_pending_cmds[i].cmd_id = cmd_id;
            g_pending_cmds[i].retry_count = 0;
            g_pending_cmds[i].last_send_time = esp_timer_get_time() / 1000;
            g_pending_cmds[i].ack_timeout = RETRY_BASE_DELAY;
            return;
        }
    }
    ESP_LOGW(TAG, "待确认指令列表已满");
}

/**
 * @brief 收到响应后移除待确认指令
 */
void acknowledge_command(uint32_t cmd_id)
{
    for (int i = 0; i < MAX_PENDING_CMDS; i++) {
        if (g_pending_cmds[i].cmd_id == cmd_id) {
            memset(&g_pending_cmds[i], 0, sizeof(pending_cmd_t));
            ESP_LOGD(TAG, "指令确认: cmd_id=%lu", (unsigned long)cmd_id);
            return;
        }
    }
}

/**
 * @brief 检查超时未确认的指令并重发
 */
void check_pending_commands(void)
{
    uint32_t now = esp_timer_get_time() / 1000;

    for (int i = 0; i < MAX_PENDING_CMDS; i++) {
        pending_cmd_t *pending = &g_pending_cmds[i];
        if (pending->cmd_id == 0) continue;

        uint32_t elapsed = now - pending->last_send_time;
        if (elapsed > pending->ack_timeout) {
            if (pending->retry_count >= MAX_RETRY_COUNT) {
                ESP_LOGE(TAG, "指令重发次数耗尽: cmd_id=%lu",
                         (unsigned long)pending->cmd_id);
                /* 通知上层指令失败 */
                memset(pending, 0, sizeof(pending_cmd_t));
            } else {
                pending->retry_count++;
                /* 指数退避 */
                pending->ack_timeout = RETRY_BASE_DELAY *
                                       (1 << pending->retry_count);
                pending->last_send_time = now;
                ESP_LOGW(TAG, "指令重发: cmd_id=%lu, retry=%d, delay=%lu ms",
                         (unsigned long)pending->cmd_id,
                         pending->retry_count,
                         (unsigned long)pending->ack_timeout);
                /* 重新发送指令 */
                resend_command(pending->cmd_id);
            }
        }
    }
}

七、权限控制与审计

7.1 指令权限验证

/* 权限级别 */
typedef enum {
    AUTH_LEVEL_GUEST    = 0,  /* 访客: 只能查询 */
    AUTH_LEVEL_USER     = 1,  /* 用户: 查询+普通控制 */
    AUTH_LEVEL_OPERATOR = 2,  /* 操作员: 所有控制+配置 */
    AUTH_LEVEL_ADMIN    = 3,  /* 管理员: 所有操作包括系统命令 */
} auth_level_t;

/* 指令权限映射表 */
typedef struct {
    uint8_t  cmd_type;
    uint8_t  min_auth_level;
} cmd_permission_t;

static const cmd_permission_t g_permissions[] = {
    {CMD_TYPE_ACTUATOR, AUTH_LEVEL_USER},
    {CMD_TYPE_CONFIG,   AUTH_LEVEL_OPERATOR},
    {CMD_TYPE_QUERY,    AUTH_LEVEL_GUEST},
    {CMD_TYPE_OTA,      AUTH_LEVEL_ADMIN},
    {CMD_TYPE_DIAG,     AUTH_LEVEL_OPERATOR},
    {CMD_TYPE_SYSTEM,   AUTH_LEVEL_ADMIN},
};

/**
 * @brief 验证指令权限
 */
bool verify_command_permission(const remote_cmd_t *cmd, auth_level_t user_level)
{
    for (int i = 0; i < sizeof(g_permissions) / sizeof(g_permissions[0]); i++) {
        if (g_permissions[i].cmd_type == cmd->cmd_type) {
            if (user_level < g_permissions[i].min_auth_level) {
                ESP_LOGW(TAG, "权限不足: cmd_type=%d, user_level=%d, required=%d",
                         cmd->cmd_type, user_level, g_permissions[i].min_auth_level);
                return false;
            }
            return true;
        }
    }
    /* 未配置权限的指令默认拒绝 */
    ESP_LOGW(TAG, "未找到权限配置: cmd_type=%d", cmd->cmd_type);
    return false;
}

/* 指令审计日志 */
typedef struct {
    uint32_t cmd_id;
    uint8_t  cmd_type;
    uint8_t  result;
    uint64_t timestamp;
    uint32_t duration_ms;
} audit_entry_t;

#define AUDIT_LOG_SIZE 50
static audit_entry_t g_audit_log[AUDIT_LOG_SIZE];
static uint8_t g_audit_idx = 0;

void log_audit(uint32_t cmd_id, uint8_t cmd_type, uint8_t result,
               uint32_t duration_ms)
{
    audit_entry_t *entry = &g_audit_log[g_audit_idx];
    entry->cmd_id = cmd_id;
    entry->cmd_type = cmd_type;
    entry->result = result;
    entry->timestamp = get_synced_time_ms();
    entry->duration_ms = duration_ms;
    g_audit_idx = (g_audit_idx + 1) % AUDIT_LOG_SIZE;
}

八、完整系统启动流程

void app_main(void)
{
    ESP_LOGI(TAG, "沧州虎王科技 ESP32 远程控制系统启动");

    /* 1. 初始化NVS */
    esp_err_t ret = nvs_flash_init();
    if (ret != ESP_OK) {
        ESP_ERROR_CHECK(nvs_flash_erase());
        ESP_ERROR_CHECK(nvs_flash_init());
    }

    /* 2. 初始化网络 */
    init_wifi_or_4g();
    init_ntp_time_sync();

    /* 3. 加载安全配置 */
    load_device_secret();
    load_device_config();

    /* 4. 初始化指令执行引擎 */
    init_cmd_execution_engine();

    /* 5. 初始化硬件执行器 */
    init_actuator_gpios();

    /* 6. 初始化MQTT安全通道 */
    init_mqtt_secure_channel(g_device_id);
    esp_mqtt_client_start(g_mqtt_client);

    /* 7. 创建指令执行任务 */
    for (int i = 0; i < 4; i++) {
        char name[16];
        snprintf(name, sizeof(name), "cmd_exec_%d", i);
        xTaskCreatePinnedToCore(cmd_executor_task, name, 6144,
                                (void *)(uintptr_t)i, 5 - i, NULL, 0);
    }

    /* 8. 创建重发检查任务 */
    xTaskCreate(retry_check_task, "retry_chk", 2048, NULL, 3, NULL);

    ESP_LOGI(TAG, "远程控制系统启动完成");

    /* 主循环 */
    while (1) {
        vTaskDelay(pdMS_TO_TICKS(10000));
        /* 定期检查离线缓存 */
        if (g_mqtt_connected && g_offline_cache.count > 0) {
            process_offline_cache();
        }
    }
}

九、工程实践总结

远程控制通道的设计需要在安全性、可靠性和实时性之间取得平衡。我们团队的经验是:

安全优先:所有远程指令必须经过签名验证和时序检查,禁止明文传输控制指令。TLS双向认证是基础保障,HMAC签名防止指令篡改,时间戳防止重放攻击。

可靠性保障:MQTT QoS 1保证指令至少送达一次,指令ID去重防止重复执行,离线缓存保证断网期间不丢失指令,指数退避重发保证网络波动下的最终送达。

执行可控:优先级队列保证紧急指令优先执行,并发数限制防止系统过载,执行超时防止指令卡死,审计日志便于事后追溯。

在调试和运维方面,ESP32工具箱V2.0提供了指令模拟发送功能,可以在本地测试指令解析和执行逻辑。随身WiFi硬件调试工具(hardware.czkree.com)支持在现场通过串口直接发送测试指令,无需依赖云平台。所有指令执行记录会上报到物联网平台的审计中心,运维团队可以追踪每条指令的完整生命周期。


作者:沧州虎王科技技术团队
标签:物联网、嵌入式、ESP32
产品推荐:ESP32工具箱V2.0 | 随身WiFi硬件调试工具(hardware.czkree.com) | 物联网平台

Logo

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

更多推荐