#include <pthread.h>
#include <stdbool.h>
#include <stdio.h>
#include <stdlib.h>
#include <unistd.h>
#include <sys/socket.h>
#include <sys/epoll.h>
#include <arpa/inet.h>
#include <string.h>
#include <curl/curl.h>
#include <signal.h>
#include <openssl/md5.h>
#include "lz4.h"

#include "dy_utils/uthash/uthash.h"
#include "dy_utils/dy_common.h"
#include "dy_utils/dy_db.h"
#include "dy_utils/dy_ipc.h"
#include "dy_utils/cJSON.h"

#include "../proto_forward.h"
#include "../mqtt_emms2/mqtt_emms2.h"
#include "../cc_lc_com/define.h"
#include "../cc_lc_com/collector-api.h"
#include "../frpc_proxy/frpc_proxy.h"
#include "ota_upgrade.h"

static ota_param_t* g_ota_param = NULL;

ota_param_t *get_ota_param(void){
    return g_ota_param;
}

#if 0
static int check_file_exist(char *file_name) {
    
    char file_path[128] = "";	
	snprintf(file_path, sizeof(file_path), "%s/%s", OTA_UPGRADE_DIR, file_name);

    if (access(file_path, F_OK) != -1) {  // 检查文件是否存在
        // 文件存在，删除文件
        if (remove(file_path) == 0) {
            ems_syslog(LOG_NOTICE, "File deleted successfully");
        } else {
            ems_syslog(LOG_ERR, "Error deleting file");
            return -1;
        }
    } else {
        ems_syslog(LOG_NOTICE, "File does not exist");
    }

    return 0;
}
#endif

static void refresh_ota_info_file(cJSON* root){
    if (root == NULL) return;
   
    char filepath[256] = {0};
    snprintf(filepath, sizeof(filepath), "%s/%s", OTA_UPGRADE_DIR, OTA_INFO_FILE);

	char *string = cJSON_Print(root);
	if (string)
	{
        write_file_data_safe(filepath, string, strlen(string));
	}
	else{
		ems_syslog(LOG_ERR, "refresh ota_info.json error!!!!");
	}

    if (string != NULL) free(string);	
}

static int ota_mqtt_publish_message(mqtt_emms2_var_t* var, char* pub_topic, char* data, int data_len)
{
    if (var->base_param.disable_lz4)
    {
        return mqtt_session_publish(var->session, pub_topic, data, data_len);
    }
    else
    {
        int   ret             = 0;
        int   src_size        = data_len + 1;
        int   max_dst_size    = LZ4_compressBound(src_size);
        char* compressed_data = calloc(1, max_dst_size);

        const int compressed_data_size = LZ4_compress_default(data, compressed_data, src_size, max_dst_size);
        if (compressed_data_size <= 0)
        {
            ems_syslog(LOG_ERR, "LZ4_compress err");
            free(compressed_data);
            return -1;
        }
        // printf("We successfully compressed some data! Ratio: %.2f\n",
        //     (float) compressed_data_size/src_size);
        char topic_lz4[256] = {0};
        sprintf(topic_lz4, "%s/lz4/%d", pub_topic, data_len + 1);
        ems_syslog(LOG_INFO, "lz4 change topic [%s]", topic_lz4);
        ret = mqtt_session_publish(var->session, topic_lz4, compressed_data, compressed_data_size);
        free(compressed_data);
        return ret;
    }
}

static int mqtt_publish_simple_result(mqtt_emms2_var_t* var, char* type, char* func_id, int seq, int result, char* info_key, cJSON* info)
{
    int   ret     = -1;
    char* topic   = malloc(256);
    char* payload = malloc(2048);
    if (topic && payload)
    {
        snprintf(topic, 256, "emms2/%s/%s/%s", type, var->topic_sn, func_id);
        int len = snprintf(payload,
                           1024,
                           "{\"funcId\":\"%s\",\"lcSN\": \"%s\",\"seq\": %d,\"time\":%lu,\"result\": %d",
                           func_id,
                           var->topic_sn,
                           seq,
                           time(NULL),
                           result);
        if (info_key && info)
        {
            char* i = cJSON_Print(info);
            if (i)
            {
                len += snprintf(payload + len, 1024 - len, ",\"%s\":%s", info_key, i);
                free(i);
            }
        }
        len += snprintf(payload + len, 1024 - len, "}");

        ret = ota_mqtt_publish_message(var, topic, payload, len);
    }
    if (topic) free(topic);
    if (payload) free(payload);
    return ret;
}

int mqtt_otalog_publish(uint8_t levl, char *content)
{
    ota_param_t *ota_param = get_ota_param();
    mqtt_emms2_var_t *var  = (mqtt_emms2_var_t*)ota_param->var;
    {
        cJSON* _info = cJSON_CreateObject();
        cJSON_AddItemToObject(_info, "type", cJSON_CreateString("otalog"));
        cJSON_AddItemToObject(_info, "dev_no", cJSON_CreateString(ota_param->dev_no));
        cJSON_AddNumberToObject(_info, "code", levl);
        cJSON_AddItemToObject(_info, "message", cJSON_CreateString(content));
        if (mqtt_publish_simple_result(var, FUNCTION_TYPE_POST, FUNCTION_ID_OTACOURSE, ota_param->seq, 0, "info", _info) < 0)
        {
            ems_syslog(LOG_ERR, "mqtt publish message failed");
        }
        cJSON_Delete(_info);
    }
    return 0;
}

int mqtt_otaprogress_publish(float progress, uint8_t indx)
{
    ota_param_t *ota_param = get_ota_param();
    mqtt_emms2_var_t  *var = (mqtt_emms2_var_t*)ota_param->var;
    ota_info_t       *info = (ota_info_t *)ota_param->info;

    indx = (indx >= (info->file_num - 2)) ? (info->file_num - 2) : indx;
    {
        cJSON* _info = cJSON_CreateObject();
        cJSON_AddItemToObject(_info, "type", cJSON_CreateString("progress"));
        cJSON_AddItemToObject(_info, "dev_no", cJSON_CreateString(ota_param->dev_no));
        cJSON_AddNumberToObject(_info, "code", (int)progress);
        cJSON_AddItemToObject(_info, "message", cJSON_CreateString(info->tmp[indx]));
        if (mqtt_publish_simple_result(var, FUNCTION_TYPE_POST, FUNCTION_ID_OTACOURSE, ota_param->seq, 0, "info", _info) < 0)
        {
            ems_syslog(LOG_ERR, "mqtt publish message failed");
        }
        cJSON_Delete(_info);
    }
    return 0;
}

int mqtt_otaend_publish(mqtt_emms2_var_t *var, char *dev_no, int code, char *ver, int seq)
{
    if (var == NULL) return -1;

    {
        cJSON* _info = cJSON_CreateObject();
        cJSON_AddItemToObject(_info, "type", cJSON_CreateString("end"));
        cJSON_AddItemToObject(_info, "dev_no", cJSON_CreateString(dev_no));
        cJSON_AddNumberToObject(_info, "code", code);
        if (code == 0){  // 0 成功，尝试附加版本信息
            if (ver != NULL && strlen(ver) > 0)
                cJSON_AddItemToObject(_info, "version", cJSON_CreateString(ver));
            else
                cJSON_AddItemToObject(_info, "version", cJSON_CreateString(""));
        }
        if (mqtt_publish_simple_result(var, FUNCTION_TYPE_POST, FUNCTION_ID_OTACOURSE, seq, 0, "info", _info) < 0)
        {
            ems_syslog(LOG_ERR, "mqtt publish message failed");
        }
        cJSON_Delete(_info);
    }
    return 0;
}

static int mqtt_download_ota_callback(void* clientp, double dltotal, double dlnow, double ultotal, double ulnow)
{
    download_progress_t *progress = (download_progress_t*)clientp;
    if (!progress) return 0;

    time_t now = time(NULL);

    ota_param_t *ota_param = get_ota_param();
    mqtt_emms2_var_t *var  = (mqtt_emms2_var_t*)ota_param->var;

    //ems_syslog(LOG_ERR, "%s dltotal：%f, dlnow %f last_dlnow %f, now %ld, last_report_time %ld", progress->current_file, dltotal, dlnow, progress->last_dlnow, now, progress->last_report_time);
    // 触发条件：1.完成下载或2.进度变化+时间间隔
    if ((dltotal > 0 && dlnow == dltotal && progress->flag == 0) ||
        (now - progress->last_report_time >= 5 && dlnow > progress->last_dlnow))
    {
        if (dlnow == dltotal) progress->flag = 1; //

        cJSON* _info = cJSON_CreateObject();
        cJSON_AddItemToObject(_info, "type", cJSON_CreateString("download"));
        cJSON_AddItemToObject(_info, "dev_no", cJSON_CreateString(ota_param->dev_no));
        // 计算进度
        int percent = (dltotal > 0) ? (int)((dlnow / dltotal) * 100) :
                    (dlnow > 0) ? 99 : 0;
        cJSON_AddNumberToObject(_info, "code", percent);
        cJSON_AddItemToObject(_info, "message", cJSON_CreateString(progress->current_file));
        if (mqtt_publish_simple_result(var, FUNCTION_TYPE_POST, FUNCTION_ID_OTACOURSE, ota_param->seq, 0, "info", _info) < 0)
        {
            ems_syslog(LOG_ERR, "mqtt publish message failed");
        }
        progress->last_report_time = now;
        progress->last_dlnow = dlnow;
        cJSON_Delete(_info);
    }
    return 0;
}

static int ota_http_post_download_file(const char* url, const char* save_path, const char* file_name, ota_download_callback_t callback)
{
    download_progress_t *progress = calloc(1, sizeof(download_progress_t));
    if (!progress) return -1;

    strncpy(progress->current_file, file_name, sizeof(progress->current_file) - 1);
    progress->last_report_time = time(NULL);

    CURL*    curl = NULL;
    FILE*    fp   = NULL;
    CURLcode res;
    int      ret = -1;
    char     filepath[256];
    snprintf(filepath, sizeof(filepath), "%s/%s", save_path, file_name);

    curl_global_init(CURL_GLOBAL_DEFAULT);
    curl = curl_easy_init();
    if (!curl)
    {
        ems_syslog(LOG_ERR, "curl_easy_init failed");
        goto END;
    }

    fp = fopen(filepath, "wb");
    if (!fp)
    {
        ems_syslog(LOG_ERR, "fopen file failed:%s\n", filepath);
        goto END;
    }
    curl_easy_setopt(curl, CURLOPT_URL, url);
    curl_easy_setopt(curl, CURLOPT_WRITEDATA, fp);
    curl_easy_setopt(curl, CURLOPT_NOPROGRESS, 0L);
    curl_easy_setopt(curl, CURLOPT_FOLLOWLOCATION, 1L);
    curl_easy_setopt(curl, CURLOPT_SSL_VERIFYPEER, 0L);      // 禁用 SSL 验证
    curl_easy_setopt(curl, CURLOPT_SSL_VERIFYHOST, 0L);    
    curl_easy_setopt(curl, CURLOPT_VERBOSE, 1L);
    curl_easy_setopt(curl, CURLOPT_XFERINFODATA, progress);  // 传递状态对象
    if (callback)
        curl_easy_setopt(curl, CURLOPT_PROGRESSFUNCTION, callback);

    res = curl_easy_perform(curl);
    if (res == CURLE_OK)
    {
        long http_code = 0;
        curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &http_code);
        ret = (http_code == 200) ? 0 : -1;
    }
    else
    {
        ems_syslog(LOG_ERR, "download file failed:%s\n", curl_easy_strerror(res));
    }

END:
    if (fp) fclose(fp);
    if (curl) curl_easy_cleanup(curl);
    curl_global_cleanup();

    if (ret != 0 && access(filepath, F_OK) == 0) {
        remove(filepath);
    }

    free(progress);
    return ret;
}

static int download_and_checkmd5(ota_param_t *ota_param)
{ 
    int err_code = 0;
    ota_info_t *info = (ota_info_t *)ota_param->info;

    for (int i = 0; i < info->file_num; i++)
    {
        int ret = ota_http_post_download_file(info->file_urls[i], OTA_UPGRADE_DIR, info->file[i], mqtt_download_ota_callback); // libcurl
        if (ret == 0){
            // 下载成功
            ems_syslog(LOG_INFO, "%s download success...", info->file_urls[i]);
        }
        else{  
            ems_syslog(LOG_ERR, "%s download failed!!!", info->file_urls[i]);
            err_code = OTA_TERM_DOWNLOAD_FAIL;
            return err_code;
        }
/*
        int ret = download_file(OTA_UPGRADE_DIR, info->file_urls[i], 10); // wget命令行 下载远程http文件，不提供下载进度上报功能
        if (ret != 0)
        {
            ems_syslog(LOG_ERR, "%s download failed!!!", info->file_urls[i]);
            err_code = OTA_TERM_DOWNLOAD_FAIL;
            return err_code;
        }
*/
    }

    char file_path[256] = {0};
    for (int i = 0; i < info->file_num; i++)
    {
        snprintf(file_path, sizeof(file_path), "%s/%s", OTA_UPGRADE_DIR, info->file[i]);
        if (access(file_path, F_OK) != 0)
        {
            ems_syslog(LOG_ERR, "file %s not exist", file_path);
            err_code = OTA_TERM_INTERNAL_FAULT;
            return err_code;
        }

        int ret = check_md5sum_file(file_path, info->file_md5s[i]);
        if (ret != 0)
        {
            ems_syslog(LOG_ERR, "upgrade nodes check md5 failed!!!");
            err_code = OTA_TERM_CHECK_FAIL;
            return err_code;
        }
    }
    return err_code;
}

keyvalue_t *check_port_inpair(ota_info_t *info){
    argv_t *argv = info->argv;

    for (int i = 0; i < argv->kv_num; i++) {
        if (strcmp(argv->kv_pairs[i].key, "port") == 0) {
            return &argv->kv_pairs[i];
        }
    }
    return NULL;
}

static int parse_ota_user_def(char *input, argv_t **argv) {
    if (input == NULL) return -1;

    const char *delimiter_kv = ";"; // 键值对分隔符""
    const char delimiter_eq = '=';  // 键值分隔符

    char *str = strdup(input);
    if (!str) {
        ems_syslog(LOG_ERR, "strdup failed");
        goto end;
    }

    int pair_cnt = 0;
    char *token = strtok(str, delimiter_kv);
    while (token != NULL) {
        if (strchr(token, delimiter_eq) != NULL) {
            pair_cnt++;
        }
        token = strtok(NULL, delimiter_kv);
    }

    // 分配内存
    *argv = calloc(1, sizeof(argv_t) + sizeof(keyvalue_t) * pair_cnt);
    if (!(*argv)) {
        ems_syslog(LOG_ERR, "argv malloc failed");
        goto end;
    }
    (*argv)->kv_num = pair_cnt;
    (*argv)->user_def = strdup(input);
    ems_syslog(LOG_INFO, "user_def: %s kv_num: %d", (*argv)->user_def, (*argv)->kv_num);

    int cnt = 0;
    strcpy(str, input); // 重置字符串
    token = strtok(str, delimiter_kv);
    while (token != NULL && cnt < pair_cnt) {
        char *eq_pos = strchr(token, delimiter_eq);
        if (eq_pos) {
            // 分割键和值
            *eq_pos = '\0'; // 在 '=' 处截断
            strncpy((*argv)->kv_pairs[cnt].key, token, sizeof((*argv)->kv_pairs[cnt].key));
            strncpy((*argv)->kv_pairs[cnt].value, eq_pos + 1, sizeof((*argv)->kv_pairs[cnt].value));
            cnt++;
        }
        token = strtok(NULL, delimiter_kv);
    }
    
    return 0;
end:
    if (str != NULL) free(str);
    return -1;
}

static void ota_info_free(ota_info_t* info);
static int ota_info_init(ota_info_t* info, cJSON* ota)
{
    if (info == NULL || ota == NULL)
    {
        return -1;
    }
    memset(info, 0, sizeof(ota_info_t));

    int    timeout = 30; // 默认超时30min
    int    num     = 0;
    cJSON* obj     = cJSON_GetObjectItem(ota, "timeout");

    if (obj != NULL)
    {
        timeout = obj->valueint;
    }
    info->start = time(NULL);
    info->end   = info->start + timeout * 60;

    obj = cJSON_GetObjectItem(ota, "files");
    if (obj == NULL || cJSON_IsArray(obj) == false) goto end;

    num = cJSON_GetArraySize(obj);
    if (num == 0) 
        goto end;
    info->file      = (char **)calloc(num, sizeof(char *)); // calloc内存
    info->file_urls = (char **)calloc(num, sizeof(char *));
    info->file_md5s = (char **)calloc(num, sizeof(char *));
    for (int i = 0; i < num; i++) //三要素： file  url  md5 （缺一不可，否则会被抛弃）默认第一个是升级脚本，后面是升级固件（支持单次多个固件升级）
    {
        cJSON* item = cJSON_GetArrayItem(obj, i);
        cJSON* file = cJSON_GetObjectItem(item, "file");
        if (file == NULL || cJSON_IsString(file) == false) continue;

        cJSON* url = cJSON_GetObjectItem(item, "url");
        if (url == NULL || cJSON_IsString(url) == false) continue;

        cJSON* md5 = cJSON_GetObjectItem(item, "md5");
        if (md5 == NULL && cJSON_IsString(md5) == false) continue;
        
        info->file[info->file_num]      = strdup(file->valuestring);
        info->file_urls[info->file_num] = strdup(url->valuestring);
        info->file_md5s[info->file_num] = strdup(md5->valuestring);
        info->file_num++;//
    }

    obj = cJSON_GetObjectItem(ota, "script");
    if (obj == NULL || cJSON_IsString(obj) == false) goto end;
    strncpy(info->script, obj->valuestring, sizeof(info->script) - 1);
    info->script[sizeof(info->script) - 1] = '\0';

/*
    Example: "argv":["RS485_2"]
    obj = cJSON_GetObjectItem(ota, "argv");
    if (obj == NULL || cJSON_IsArray(obj) == false) goto end;

    int size = cJSON_GetArraySize(obj);
    info->argc = 0; 
    info->argv = (char **)calloc(size, sizeof(char *));
    for (int i = 0; i < size; i++)
    {
        cJSON* child = cJSON_GetArrayItem(obj, i);
        if (child != NULL && cJSON_IsString(child)) // 脚本启动扩展参数，必须是字符串
        {
            info->argv[info->argc++] = strdup(child->valuestring);
        }
    }
*/
    // Example: "argv":"port=RS485_2;"
    cJSON* argv = cJSON_GetObjectItem(ota, "argv");
    if (argv == NULL || cJSON_IsString(argv) == false) goto end;

    if (parse_ota_user_def(argv->valuestring, &info->argv) == -1)
       goto end;

    return 0;

end:
    ota_info_free(info);
    return -1;
}

static void ota_info_free(ota_info_t* info)
{
    if (info == NULL) return;

    if (info->file != NULL){
        for (int i = 0; i < info->file_num; i++)
        {
            if (info->file[i] != NULL) free(info->file[i]);
        }
        free(info->file);
    }
    if (info->file_urls != NULL){
        for (int i = 0; i < info->file_num; i++)
        {
            if (info->file_urls[i] != NULL) free(info->file_urls[i]);
        }
        free(info->file_urls);
    }
    if (info->file_md5s != NULL){
        for (int i = 0; i < info->file_num; i++)
        {
            if (info->file_md5s[i] != NULL) free(info->file_md5s[i]);
        }
        free(info->file_md5s);
    }
    /*
    if (info->argv != NULL){
        for (int i = 0; i < info->argc; i++)
        {
            if (info->argv[i] != NULL) free(info->argv[i]);
        }
        free(info->argv);
    }*/
    if (info->argv != NULL) free(info->argv);

    for (int i = 0; i < OTA_FILE_MAX; i++)
    {
        if (info->tmp[i] != NULL) free(info->tmp[i]);
    }

    free(info); // must free
}

static int register_ota_deal(ota_param_t *ota_param){

    int err_code = 0;

    ota_param->interact = script_interact_init();// 脚本交互结构体初(管道)始化
    if (ota_param->interact == NULL){
        ems_syslog(LOG_ERR, "script interact Init faild!!!");
        err_code = OTA_TERM_INTERNAL_FAULT;
        goto out;
    }

    ota_param->hd = register_ota_interface(ota_param->dev_no, ota_param->chan_pr, ota_param->dev_pr); // 注册ota_hd接口，屏蔽硬件接口差异
    if (ota_param->hd == NULL){
        ems_syslog(LOG_ERR, "ota interface Register faild!!!");
        err_code = OTA_TERM_INTERNAL_FAULT;
        goto out;
    }

    return 0;

out:
    return err_code;
}

static int ota_param_init(const cJSON* info, ota_param_t *ota_param)
{
    if (info == NULL || ota_param == NULL) return -1;

    int err_code = 0;
    proto_forward_t* pro_pr  = get_proto_forward_var();
    channel_t*       chan_pr = NULL;
    device_t*        dev_pr  = NULL;

    ota_param->info = (ota_info_t *)malloc(sizeof(ota_info_t)); // OTA参数
    if (ota_param->info == NULL){
        err_code = OTA_TERM_INTERNAL_FAULT;
        goto out;  
    }  
    memset(ota_param->info, 0, sizeof(ota_info_t));

    cJSON* ota = cJSON_GetObjectItem(info, "ota");
    if (ota == NULL) {err_code = OTA_TERM_PARSE_FAILED; goto out;}
    if (ota_info_init((ota_info_t *)ota_param->info, ota) < 0){
        ems_syslog(LOG_ERR, "ota info Analysis faild!!!");
        err_code = OTA_TERM_PARSE_FAILED;
        goto out;  
    }
    if (ota_param->same){ // 与通讯口不一致，需检查port
        if (check_port_inpair((ota_info_t *)ota_param->info) == NULL){
            ems_syslog(LOG_ERR, "is not communication port, and argv is missing the port!!!");
            err_code = OTA_TERM_PARSE_FAILED;
            goto out; 
        }
    }

    cJSON* obj = cJSON_GetObjectItem(info, "dev_no");
    if (obj == NULL || obj->valuestring == NULL || obj->valuestring[0] == '\0'){
        err_code = OTA_TERM_PARSE_FAILED;
        goto out; 
    }
    strncpy(ota_param->dev_no, obj->valuestring, sizeof(ota_param->dev_no) - 1);
    ota_param->dev_no[sizeof(ota_param->dev_no) - 1] = '\0';

    if (find_dev_channel_pr_by_dev_no(pro_pr, ota_param->dev_no, &chan_pr, &dev_pr) < 0){
        ems_syslog(LOG_ERR, "Find channel and dev faild!!!");
        err_code = OTA_TERM_INTERNAL_FAULT;
        goto out;
    }
    if (ota_param->same == 0 && chan_pr->proxy_state != PROXY_MODE_NONE){
        ems_syslog(LOG_WARNING, "channel is busy: %d", chan_pr->proxy_state);
        err_code = OTA_TERM_CHANNEL_BUSY;
        goto out;
    }

    ota_param->chan_pr = chan_pr;
    ota_param->dev_pr  = dev_pr;

    ems_syslog(LOG_INFO, "The ready condition is ok...");
    return 0;

out:
    return err_code;
}

void free_ota_param(ota_param_t *ota_param)
{
    if (ota_param == NULL) return;

    if (ota_param->interact != NULL){
        interact_t* interact = ota_param->interact;
        script_interact_free(interact);
    }

    if (ota_param->hd != NULL){
        ota_hd_t *hd = ota_param->hd;
        ota_deinit(hd);
    }

    if (ota_param->info != NULL){
        ota_info_t *info = (ota_info_t *)ota_param->info;
        ota_info_free(info);
    }
    free(ota_param); 

    g_ota_param = NULL; // NULL
}

extern int start_upgrade_process(ota_param_t *ota_param);
static void* ota_upgrade_start(void* arg)
{
    pthread_detach(pthread_self());

    int err_code = 0;
    ota_param_t *ota_param = (ota_param_t *)arg;
    mqtt_emms2_var_t* var = (mqtt_emms2_var_t*)ota_param->var;

    while (1) {
        err_code = register_ota_deal(ota_param);
        if (err_code != 0)
        {
            ems_syslog(LOG_ERR, "register ota other failed!!!");
            goto out;
        }

        // 下载并验证 && MD5信息
        err_code = download_and_checkmd5(ota_param);
        if (err_code != 0)
        {
            ems_syslog(LOG_ERR, "download or check md5 error, err_code: %d", err_code);
            goto out;
        }

        // MD5验证通过，执行升级动作 
        err_code = start_upgrade_process(ota_param);
        if (err_code != 0){
            ems_syslog(LOG_ERR, "start upgrade process failed!!!");
            goto out;
        }
        return NULL;
    }
out:
    ems_syslog(LOG_WARNING, "upgrade fail err_code: %d", err_code);
    mqtt_otaend_publish(var, ota_param->dev_no, err_code, NULL, ota_param->seq);
    free_ota_param(ota_param);
    return NULL;
}

static int is_process_running(pid_t pid) {
    if (kill(pid, 0) == 0){  // 发送信号0（不实际发送，仅检查权限）
        return 1;
    } 
    else{
        if (errno == EPERM){
            ems_syslog(LOG_ERR, "Process %d exists, but no permission to send signal.", pid);
            return 1;
        } 
        else if (errno == ESRCH){
            return 0; // 进程不存在
        } 
        else{
            ems_syslog(LOG_ERR, "kill");
            return -1;
        }
    }
}

// 杀死进程（返回0：成功，-1：失败）
static int kill_process(pid_t pid) {
    if (kill(pid, SIGTERM) == 0) {
        ems_syslog(LOG_INFO, "Sent SIGTERM to process %d", pid);
        return 0;
    } 
    else {
        ems_syslog(LOG_ERR, "kill SIGTERM failed");
        
        // 如果SIGTERM失败，尝试强制杀死
        if (kill(pid, SIGKILL) == 0){
            ems_syslog(LOG_INFO, "Sent SIGKILL to process %d", pid);
            return 0;
        } 
        else{
            ems_syslog(LOG_ERR, "kill SIGKILL failed");
            return -1;
        }
    }
}

static int ota_upgrade_stop(ota_param_t *ota_param)
{
    if (ota_param == NULL) return -1;

    int err_code = 0;
    interact_t *interact = ota_param->interact; 
    mqtt_emms2_var_t* var = (mqtt_emms2_var_t*)ota_param->var;
    
    if (ota_param->state == OTA_UPGRADING_NODE){
        // kill lua process
        if (is_process_running(interact->child_pid) == 1) {
            ems_syslog(LOG_INFO, "Process %d is running. Attempting to kill...", interact->child_pid);
            kill_process(interact->child_pid);
        }
    }
    else if (ota_param->state == OTA_UPGRADE_IDLE){
        err_code = OTA_TERM_READY_RUN;
        goto out;
    }
    return 0;

out:
    ems_syslog(LOG_WARNING, "upgrade stop err_code: %d", err_code);
    mqtt_otaend_publish(var, ota_param->dev_no, err_code, NULL, ota_param->seq);
    return -1;
}

static int handle_start_command(mqtt_emms2_var_t* var, cJSON* root) {
    int err_code = 0;
    cJSON* obj = NULL;

    cJSON* seq_json = cJSON_GetObjectItem(root, "seq");
    cJSON* same     = cJSON_GetObjectItem(root, "same");
    cJSON* cmd      = cJSON_GetObjectItem(root, "cmd");
    cJSON* info     = cJSON_GetObjectItem(root, "info");

    if (g_ota_param != NULL && g_ota_param->state == OTA_UPGRADING_NODE) {
        ems_syslog(LOG_WARNING, "upgrade already in progress...");
        err_code = OTA_TERM_ALREADY_RUN;
        goto out;
    }

    int control_mode = dev_get_dev_tag_int(DEV_NO_EMS, CONTROL_MODE);
    double act_power = dev_get_dev_tag_float(DEV_NO_EMS, TOTAL_ACTIVE_POWER);
    if (control_mode != 0 || act_power > 2 || act_power < -2){ // Manual || Active power
        err_code = OTA_TERM_NOMANUAL_OR_POWER;
        goto out;
    }

// @attention: 此部分必须在回调函数中解析，函数结束后，root 会被释放
#if 1 //-----------------------------------------------------------ota_param start
{
    g_ota_param = (ota_param_t*)malloc(sizeof(ota_param_t));
    if (g_ota_param == NULL) {
        ems_syslog(LOG_ERR, "memory malloc failed!!!");
        err_code = OTA_TERM_INTERNAL_FAULT;
        goto out;
    }
    memset(g_ota_param, 0, sizeof(ota_param_t));

    g_ota_param->state = OTA_UPGRADE_IDLE;
    g_ota_param->seq   = seq_json->valueint;
    g_ota_param->mode  = 0;  // 固定为串行升级
    g_ota_param->same  = same->valueint;
    strncpy(g_ota_param->cmd, cmd->valuestring, sizeof(g_ota_param->cmd) - 1);
    g_ota_param->cmd[sizeof(g_ota_param->cmd) - 1] = '\0';
    g_ota_param->var   = var;
}
    // info:{"dev_no":"PCS", "ota":{}}
    err_code = ota_param_init(info, g_ota_param);
    if (err_code != 0){
        ems_syslog(LOG_ERR, "Failed to initialize OTA parameters!!!");
        free_ota_param(g_ota_param);
        goto out;
    }
#endif //---------------------------------------------------------ota_param end

    pthread_t tid = {0};
    int ret = pthread_create(&tid, NULL, ota_upgrade_start, g_ota_param);
    if (ret != 0) {
        ems_syslog(LOG_ERR, "Failed to create ota_upgrade_start thread: %s", strerror(ret));
        err_code = OTA_TERM_INTERNAL_FAULT;
        goto out;
    }
    else{
        pthread_setname_np(tid, "ota_upgrade_start");
    }
    return 0;
    
out:
    obj = cJSON_GetObjectItem(info, "dev_no");
    mqtt_otaend_publish(var, obj->valuestring, err_code, NULL, seq_json->valueint);
    return -1;
}

static int handle_stop_command(mqtt_emms2_var_t* var, cJSON* root) {
    int err_code = 0;
    cJSON* obj = NULL;

    cJSON* seq_json = cJSON_GetObjectItem(root, "seq");
    cJSON* same     = cJSON_GetObjectItem(root, "same");
    cJSON* cmd      = cJSON_GetObjectItem(root, "cmd");
    cJSON* info     = cJSON_GetObjectItem(root, "info");

    if (g_ota_param == NULL) {
        ems_syslog(LOG_WARNING, "upgrade already stopped!!!");
        err_code = OTA_TERM_ALREADY_STOP;
        goto out;
    }

    {
        g_ota_param->seq   = seq_json->valueint;
        g_ota_param->mode  = 0;  // 固定为串行升级
        g_ota_param->same  = same->valueint;
        strncpy(g_ota_param->cmd, cmd->valuestring, sizeof(g_ota_param->cmd) - 1);
        g_ota_param->cmd[sizeof(g_ota_param->cmd) - 1] = '\0';
    }

    ota_upgrade_stop(g_ota_param);
    return 0;

out:
    obj = cJSON_GetObjectItem(info, "dev_no");
    mqtt_otaend_publish(var, obj->valuestring, err_code, NULL, seq_json->valueint);
    return -1;
}

// 注意：：：升级前需切手动模式，PCS关机功率写0，方可进行升级操作，升级成功后，再切自动运行策略
int ota_upgrade_cmd(mqtt_emms2_var_t* var, cJSON* root, int seq) {

    mqtt_publish_simple_result(var, FUNCTION_TYPE_SETRESP, FUNCTION_ID_OTAUPDATE, seq, 0, NULL, NULL); // ACK

    int err_code = 0; cJSON* obj = NULL;
    if (var == NULL || root == NULL) {
        ems_syslog(LOG_ERR, "Invalid arguments!!!");
        return -1;
    }

    signal(SIGPIPE, SIG_IGN);

    cJSON* seq_json = cJSON_GetObjectItem(root, "seq");
    cJSON* same     = cJSON_GetObjectItem(root, "same");
    cJSON* cmd      = cJSON_GetObjectItem(root, "cmd");
    cJSON* info     = cJSON_GetObjectItem(root, "info");
    if (seq_json == NULL || same == NULL || cmd == NULL || info == NULL) {
        ems_syslog(LOG_ERR, "missing required fields in JSON!!!");
        err_code = OTA_TERM_PARSE_FAILED;
        mqtt_otaend_publish(var, "", err_code, NULL, 0); // No method to sure dev_no and seq!!!
        return -1;
    }

    refresh_ota_info_file(root); //

    if (strcmp(cmd->valuestring, "start") == 0) {
        handle_start_command(var, root);
    } 
    else if (strcmp(cmd->valuestring, "stop") == 0) {
        handle_stop_command(var, root);
    } 
    else {
        ems_syslog(LOG_ERR, "unsupport command %s !!!", cmd->valuestring);
        err_code = OTA_TERM_UNSUPPORT_CMD;
        goto out;
    }
    return 0;

out:
    obj = cJSON_GetObjectItem(info, "dev_no");
    mqtt_otaend_publish(var, obj->valuestring, err_code, NULL, seq);
    return -1;
}
