瀏覽代碼

QT与pdu消息队列通信功能添加

chenyuanjie 1 年之前
父節點
當前提交
69bfde75ba
共有 3 個文件被更改,包括 300 次插入 和 0 次删除
  1. 277 0
      pro/src/app.c
  2. 22 0
      pro/src/common.h
  3. 1 0
      pro/src/thread.h

+ 277 - 0
pro/src/app.c

@@ -31,6 +31,7 @@
 #include "pduMIB_trap.h"
 #include "Cellular.h"
 #include "wifi.h"
+#include <sys/msg.h>
 
 
 
@@ -2221,6 +2222,278 @@ static int shm_New_CtrlBoard(void *arg)
 }
 
 
+/**
+ * @brief 处理控制消息
+ *
+ * 根据传入的消息类型和控制指令,执行相应的操作
+ *
+ * @param arg 指向消息数据的指针
+ *
+ * @return 返回值 0 表示成功,-1 表示失败
+ */
+static int msg_read_Ctrl(void* arg)
+{
+    if (arg==NULL)
+    {
+        return -1;
+        /* code */
+    }
+    char buff[256]= {0};
+    MsgData* msg = (MsgData*)arg;
+    int ret=-1;
+    switch (msg->CtrlType)
+    {
+    case Msg_Ch_Ctrl:
+    {
+        int nChID = msg->nValue[0];
+        int nCtrl = msg->nValue[1];
+        GlobalPowerManger* _globalPowerMangerTemp = NULL ;
+        list_for_each_entry(_globalPowerMangerTemp, &__globalDeviceManage._globalPowerManger.list, list)
+        {
+            if (_globalPowerMangerTemp == NULL)
+                break;
+            if (nChID == _globalPowerMangerTemp->product_ch_id)
+            {
+                log_d("switch ctrl:%d ctrlstatus:%d ", _globalPowerMangerTemp->product_ch_id,nCtrl);
+                log_d("switch product_saddr:%d product_ch_addr:%d ", _globalPowerMangerTemp->product_saddr,_globalPowerMangerTemp->product_ch_addr);
+                ret = g_switch_set_all_chn_ctrl(&__globalDeviceManage._globalRelaySampManger, _globalPowerMangerTemp,
+                                                _globalPowerMangerTemp->product_saddr,
+                                                _globalPowerMangerTemp->product_ch_addr,
+                                                nCtrl, false);
+                if (ret < 0)
+                {
+                    log_d("switch product_saddr:%d product_ch_addr:%d ", _globalPowerMangerTemp->product_saddr,_globalPowerMangerTemp->product_ch_addr);
+                    ret = g_switch_set_all_chn_ctrl(&__globalDeviceManage._globalRelaySampManger, _globalPowerMangerTemp,
+                                                    _globalPowerMangerTemp->product_saddr,
+                                                    _globalPowerMangerTemp->product_ch_addr,
+                                                    nCtrl, false);
+                    if (ret < 0)
+                    {
+                        log_e("switch product_saddr:%d product_ch_addr:%d Error",_globalPowerMangerTemp->product_saddr,_globalPowerMangerTemp->product_ch_addr);
+                    }
+
+                }
+                if (nCtrl == 1)
+                {
+                    sprintf(buff, "$开启$|$通道$|%d", _globalPowerMangerTemp->product_ch_id);
+                    dev_insert_alarm_ctrl(__globalDeviceManage.db, _globalPowerMangerTemp->product_id, ALARM_TYPE_OPRATION, buff);
+                }
+                else
+                {
+                    sprintf(buff, "$开启$|$通道$|%d", _globalPowerMangerTemp->product_ch_id);
+                    dev_insert_alarm_ctrl(__globalDeviceManage.db, _globalPowerMangerTemp->product_id, ALARM_TYPE_OPRATION, buff);
+                }
+            }
+        }
+    }
+    break;
+    case Msg_All_Ctrl:
+    {   
+        //控制所有通道开关
+        int nCtrl = msg->nValue[0];
+        int nSendAddr = 0;
+        char buff[256];
+
+        GlobalPowerManger* _globalPowerMangerTemp = NULL ;
+        list_for_each_entry(_globalPowerMangerTemp, &__globalDeviceManage._globalPowerManger.list, list)
+        {
+            if(_globalPowerMangerTemp==NULL)
+                break;
+            if(_globalPowerMangerTemp->product_saddr==nSendAddr)continue;
+            nSendAddr=_globalPowerMangerTemp->product_saddr;
+            log_i("Open All address Start");
+            int rec=g_switch_set_all_ctrl(&__globalDeviceManage._globalRelaySampManger,_globalPowerMangerTemp->product_ch_type,_globalPowerMangerTemp->product_saddr,nCtrl);
+            if(rec<0)log_i("Open All address:%d rec:%d:%s",nSendAddr,rec,modbus_strerror(errno));
+            
+        }
+        if (nCtrl == 1)
+        {
+            sprintf(buff, "|$开启$|$所有$|$通道$");
+            dev_insert_alarm_ctrl(__globalDeviceManage.db, _globalPowerMangerTemp->product_id, ALARM_TYPE_OPRATION, buff);
+        }
+        else
+        {
+            sprintf(buff, "|$关闭$|$所有$|$通道$");
+            dev_insert_alarm_ctrl(__globalDeviceManage.db, _globalPowerMangerTemp->product_id, ALARM_TYPE_OPRATION, buff);
+        }
+    }
+    break;
+    default:
+        return -1; // 未识别类消息类型,返回-1表错误
+        break;
+    }
+    return 0;
+}
+
+/**
+ * @brief 初始化消息队列
+ *
+ * 根据给定的键值key尝试创建消息队列。如果队列不存在,则创建一个新的消息队列;如果队列已存在,则检查队列中是否有残留消息,
+ * 如果有残留消息或者没有等待该消息队列的进程,则清理该队列并重新创建。
+ *
+ * @param key 消息队列的键值
+ *
+ * @return 成功时返回消息队列的标识符,失败时返回-1
+ */
+int safe_msg_init(key_t key) 
+{
+    int msgid;
+    struct msqid_ds queue_info;
+    
+    // 尝试创建全新队列
+    msgid = msgget(key, 0666 | IPC_CREAT | IPC_EXCL);
+    if (msgid != -1) {
+        return msgid; // 全新创建成功
+    }
+    
+    // 若队列已存在
+    if (errno == EEXIST) {
+        // 获取现有队列
+        if ((msgid = msgget(key, 0666)) == -1) {
+            perror("msgget existing");
+            return -1;
+        }
+        
+        // 检查队列是否残留消息
+        if (msgctl(msgid, IPC_STAT, &queue_info) == -1) {
+            perror("msgctl stat");
+            return -1;
+        }
+        
+        // 判断是否需要清理
+        if (queue_info.msg_qnum > 0 || queue_info.msg_lspid == 0) {
+            log_w("Cleaning stale queue (msgs:%lu)\n", queue_info.msg_qnum);
+            
+            // 删除旧队列
+            if (msgctl(msgid, IPC_RMID, NULL) == -1) {
+                log_e("msgctl rmid");
+                return -1;
+            }
+            
+            // 重新创建
+            msgid = msgget(key, 0666 | IPC_CREAT | IPC_EXCL);
+            if (msgid == -1) {
+                log_e("msgget recreate");
+                return -1;
+            }
+            return msgid;
+        }
+        return msgid;
+    }
+    
+    log_e("msgget init");
+    return -1;
+}
+
+/**
+ * @brief 接收消息处理线程
+ *
+ * 该函数用于接收并处理消息。它会不断循环接收类型为2的消息,并在接收到消息后打印其内容。
+ * 如果接收过程中出现错误(除了ENOMSG错误),则会打印错误信息。
+ *
+ * @param arg 指向线程参数的指针
+ *
+ * @return 返回NULL
+ */
+static void *recv_handler(void *arg)
+{
+    GlobalPowerManger *_globalPowerMangerTemp = NULL;
+    GlobalDeviceManager *_globalDeviceManager = (GlobalDeviceManager *)&__globalDeviceManage;
+    Msg_thread_args *args = &_globalDeviceManager->_MSG_args;
+    // 初始化消息队列
+    key_t key;
+    int msgid;
+    // 生成键值前确保文件存在
+    const char *path = "/root/run/app/smartPDU_MSGQ";
+
+    while (args->needInit)
+    {
+        if (access(path, F_OK) == -1)
+        {
+            int fd = open(path, O_CREAT, 0666);
+            if (fd == -1)
+            {
+                log_e("smartPDU_MSGQ_create key file");
+                sleep(500); // 延时5秒后再次创建
+                continue;
+            }
+            close(fd);
+        }
+
+        if ((key = ftok(path, 'Q')) == -1)
+        {
+            log_e("smartPDU_MSGQ ftok");
+            sleep(500); // 延时5秒后再次创建
+            continue;
+        }
+        
+        // 安全初始化队列
+        if ((msgid = safe_msg_init(key)) == -1)
+        {
+            log_e("Message queue init failed\n");
+            sleep(500); // 延时5秒后再次创建
+            continue;
+        }
+        
+        log_i("Message queue init success ID=%d\n",msgid);
+        log_i("Message queue init success args-ID=%d\n",args->msgid);
+        args->msgid=msgid;
+        args->running=1;
+        MsgData msg;
+        while (args->running)
+        {
+            memset(&msg, 0, sizeof(MsgData));
+            ssize_t ret = msgrcv(args->msgid, &msg, sizeof(msg)-sizeof(long),
+                                 1,           // 接收类型为2的消息
+                                 IPC_NOWAIT); // 非阻塞模式
+
+            if (ret > 0)
+            {
+                log_i("Received from Qt Type= %d, val1= %d val2= %d\n", msg.MsgType,msg.nValue[0],msg.nValue[1]);
+                msg_read_Ctrl(&msg);
+            }
+            else if (errno != ENOMSG)
+            {
+                log_e("msgrcv error");
+                usleep(100000); // 无数据则延时100ms轮询间隔
+            }else{
+                usleep(100000); // 无数据则延时100ms轮询间隔
+            }
+        }
+    }
+    return NULL;
+}
+/**
+ * @brief 消息发送处理函数
+ *
+ * 该函数用于处理消息发送的逻辑。
+ *
+ * @param arg 函数参数,通常为空指针
+ * @param _msg 指向待发送的消息数据的指针
+ * 
+*/
+static void send_handler(void *arg, MsgData* _msg)
+{
+    if (_msg == NULL) {
+        return;
+    }
+    GlobalDeviceManager *_globalDeviceManager = (GlobalDeviceManager *)arg;
+    Msg_thread_args args = _globalDeviceManager->_MSG_args;    
+    if (args.running == 0)
+    {
+        return;
+    }
+    
+    MsgData msg = *_msg;
+    msg.MsgType = 1; // 设置消息类型为1
+
+    ssize_t ret = msgsnd(args.msgid, &msg, sizeof(msg)-sizeof(long), 0);
+    if (ret < 0) {
+        log_e("msgsnd failed, errno: %d", errno);
+        // 可根据需要添加重试逻辑
+    }
+}
+
 
 //#define USE_NTPD
 static void* ntp_thread(void *arg)
@@ -2705,6 +2978,10 @@ int app_init(void)
 
 #ifndef USE_UI
     thread_start(THREAD_ID_SHMW,   shm_thread,    _globalDeviceManager);
+
+    _globalDeviceManager->_MSG_args.needInit = true;
+
+    thread_start(THREAD_ID_RECV_MESSAGE_QUEUE, recv_handler, _globalDeviceManager);
 #endif
 
     thread_start(THREAD_ID_BREAKER_SCANNER, breaker_scanner_thread, _globalDeviceManager);

+ 22 - 0
pro/src/common.h

@@ -478,7 +478,28 @@ typedef struct
     int nSelectPh;//当前通道的三相ph参数:0:A相 1:B相 2:C相 -1通道总计
 } smartPDU_shm_data;
 
+//消息队列头结构体定义
+enum MsgCtrlType
+{
+    Msg_Ch_Ctrl=1,          //通道控制
+    Msg_All_Ctrl=2,         //整体控制
+    Msg_Threshold_Set=3,    //阈值设置
+    Msg_Threshold_Update=4, //阈值更新
+};
+typedef struct {
+    long MsgType;       //消息方向1:Qt->PDU 2:pdu->Qt
+    int  CtrlType;      //控制类型
+    int  nValue[2];
+//    double fValue[10];
+}MsgData;
 
+// 消息队列线程参数结构
+typedef struct {
+    int msgid;         // 消息队列ID
+    int running;       // 线程运行标志
+    bool needInit;
+    pthread_mutex_t lock;
+} Msg_thread_args;
 
 ///////////////////////////////////////////////////////////////////////////////////
 ///////////////////////////////////////////////////////////////////////////////////
@@ -1287,6 +1308,7 @@ typedef struct
     mail_info_t          mailInfo;
     mqtt_info_t          mqttInfo;
     Modbus_Manger_Tcp *md_tcp;
+    Msg_thread_args  _MSG_args;
 #ifdef NORTH_USER_SPECILS
     cascade_north_user_t north_user;
 #endif

+ 1 - 0
pro/src/thread.h

@@ -33,6 +33,7 @@ enum {
     THREAD_ID_BOARD,        //board upgrade
     THREAD_ID_UI,
     THREAD_ID_CELLULAR,
+    THREAD_ID_RECV_MESSAGE_QUEUE,//进程间通信队列处理线程
 
     THREAD_ID_MAX
 };