Просмотр исходного кода

1、级联主机与从设备代码解耦
2、websocket发送数据方式修改

guohui лет назад: 2
Родитель
Сommit
1d718fad0f
7 измененных файлов с 207 добавлено и 223 удалено
  1. 15 17
      pro/src/appweb_handle.c
  2. 26 5
      pro/src/cascade.c
  3. 3 29
      pro/src/cascade.h
  4. 98 98
      pro/src/cascade_slave.c
  5. 0 1
      pro/src/common.h
  6. 62 70
      pro/src/websocket_handle.c
  7. 3 3
      pro/src/websocket_handle.h

+ 15 - 17
pro/src/appweb_handle.c

@@ -223,7 +223,7 @@ void* update_thread(void* arg)
     //采集线程
     while (h->quit==0)
     {
-        if(websocket_isok())
+        if(websocket_is_online())
         {
            // pthread_mutex_lock(&Websocketlock); 
             //推送总体信息
@@ -252,7 +252,7 @@ void* update_thread(void* arg)
                     }
                     json_str = over_all_pwr_monitor_Tree_AC_to_json(&l1, &l2, &l3);
 
-                    websocket_send(json_str, strlen(json_str));
+                    websocket_broadcast(json_str, strlen(json_str));
                     cJSON_free((void*)json_str);
                 }
                 else
@@ -269,7 +269,7 @@ void* update_thread(void* arg)
                     }
                     json_str = ws_over_status_ack_to_json(0, &allInfo);
 
-                    websocket_send(json_str, strlen(json_str));
+                    websocket_broadcast(json_str, strlen(json_str));
                     cJSON_free(json_str);
                 }                
             }
@@ -290,7 +290,7 @@ void* update_thread(void* arg)
                 }
 
                 json_str = ws_chn_status_ack_to_json(1,&chInfo);
-                websocket_send(json_str, strlen(json_str));
+                websocket_broadcast(json_str, strlen(json_str));
                  
                 //释放资源,一定要记得释放json字符串!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!!
                 cJSON_free((void*)json_str);
@@ -304,8 +304,8 @@ void* update_thread(void* arg)
             }
             //else if (gdm->view_page==3)
             {
-                json_str =ws_sensor_status_ack_to_json(&__globalDeviceManage._globalSensorManger);
-                websocket_send(json_str, strlen(json_str));
+                json_str = ws_sensor_status_ack_to_json(&__globalDeviceManage._globalSensorManger);
+                websocket_broadcast(json_str, strlen(json_str));
                 cJSON_free((void*)json_str);
             }            
             //usleep(400*1000);
@@ -313,9 +313,6 @@ void* update_thread(void* arg)
             //log_i("ws upload.");
         }
 
-        int npoll=websocket_poll();
-        // log_i("ws poll once!  status=%d\n", npoll);
-
         usleep(500*1000);
     }
 
@@ -1644,7 +1641,7 @@ static void sensorStatusMonitoring(void* conn)
                     return ;
                 }   
             }            
-            if(_over_all_sensor_request.cgqtd)
+              if(_over_all_sensor_request.cgqtd)
             {
                 if (strcmp(_over_all_sensor_request.cgqtd, "RS485") == 0)
                 {
@@ -1654,25 +1651,25 @@ static void sensorStatusMonitoring(void* conn)
                         strcpy(_globalSensorMangerTemp->sensor_port_name, _over_all_sensor_request.comSlogan);
                     }
                 }
-                else if (strcmp(_over_all_sensor_request.cgqtd, "IN1") == 0)
+                else if (strcmp(_over_all_sensor_request.cgqtd, "IN3") == 0)
                 {
                     _globalSensorMangerTemp->sensor_interface_type = 2;
                     _globalSensorMangerTemp->sensor_addr = 0;
                     strcpy(_globalSensorMangerTemp->sensor_port_name, _over_all_sensor_request.cgqtd);
                 }
-                else if (strcmp(_over_all_sensor_request.cgqtd, "IN2") == 0)
+                else if (strcmp(_over_all_sensor_request.cgqtd, "IN4") == 0)
                 {
                     _globalSensorMangerTemp->sensor_interface_type = 2;
                     _globalSensorMangerTemp->sensor_addr = 1;
                     strcpy(_globalSensorMangerTemp->sensor_port_name, _over_all_sensor_request.cgqtd);
                 }
-                else if (strcmp(_over_all_sensor_request.cgqtd, "IN3") == 0)
+                else if (strcmp(_over_all_sensor_request.cgqtd, "IN2") == 0)
                 {
                     _globalSensorMangerTemp->sensor_interface_type = 2;
                     _globalSensorMangerTemp->sensor_addr = 2;
                     strcpy(_globalSensorMangerTemp->sensor_port_name, _over_all_sensor_request.cgqtd);
                 }
-                else if (strcmp(_over_all_sensor_request.cgqtd, "IN4") == 0)
+                else if (strcmp(_over_all_sensor_request.cgqtd, "IN1") == 0)
                 {
                     _globalSensorMangerTemp->sensor_interface_type = 2;
                     _globalSensorMangerTemp->sensor_addr = 3;
@@ -2545,6 +2542,7 @@ static void serviceManagement(void *conn)
 #endif
 
     //转换json参数
+    memset(&_serviceManageRequestInfo, 0, sizeof(_serviceManageRequestInfo));
     cjson = json_to_service_manage(rx_body_buf,&_serviceManageRequestInfo);
     if(cjson==NULL)
     {
@@ -2874,8 +2872,8 @@ static void serviceManagement(void *conn)
         }
         if (bflag == false) return;
         
-        ack_str = service_sx_ack_to_json(_globalPowerMangerTemp->product_id, &__globalDeviceManage._globalPowerManger);
-        if (_serviceManageRequestInfo._posix_info.productId) free(_serviceManageRequestInfo._posix_info.productId);
+        //ack_str = service_sx_ack_to_json(_globalPowerMangerTemp->product_id, &__globalDeviceManage._globalPowerManger);
+        //if (_serviceManageRequestInfo._posix_info.productId) free(_serviceManageRequestInfo._posix_info.productId);
     }   
     else if(strcmp(_serviceManageRequestInfo.type,"sXTimeSave")==0)//保存时间间隔
     {
@@ -3246,7 +3244,7 @@ static void serviceManagement(void *conn)
         ack_str = ctrl_sts_ack_to_json();
     }
     
-  //  cJSON_Delete(cjson);
+    //cJSON_Delete(cjson);
     http_reply(conn, 200, ack_str);
     cJSON_free((void*)ack_str);
 }

+ 26 - 5
pro/src/cascade.c

@@ -29,7 +29,27 @@
     #define MB_MASTER       0
 #endif
 
+typedef struct {
+    
+    int                 inited;
+    pthread_mutex_t     mutex;          //used for list lock  
+    pthread_mutex_t     lock;
+    
+    Modbus_Manger       m;
+    int                 addr;
+    slave_t             slaves[CASCADE_MAX+1];
+
+    ModbusInfo_t        mInfo;
+    modbus_mapping_t    *map;
+    modbus_mapping_t    *map2;
+
+    cmd_data_t          cmd;
+    int                 scanAddr;
 
+    slave_info_t        sInfo;      //current slave infomation
+
+    handle_t            list;
+}cascade_handle_t;
 
 int cur_dev_addr=0;
 static cascade_handle_t casHandle={.inited=0};
@@ -217,7 +237,7 @@ static int mb_init(cascade_handle_t *cas, char *path, int type, int addr, uint32
 
             cas->addr = addr;
             g_modbus_set_slave(&cas->m, addr);
-            cascade_slave_map_init(cas);
+            cascade_slave_init();
         }
         else {
             cas->addr = 0;
@@ -563,10 +583,10 @@ static int mb_receive(cascade_handle_t *cas)
             LOGD("___ slave XXXXXXXXXXX\n");
 
             if(h.func==MODBUS_FC_READ_HOLDING_REGISTERS) {
-                cascade_slave_read(cas, h.reg, h.regcnt);
+                cascade_slave_read(h.reg, h.regcnt);
             }
             else if(h.func==MODBUS_FC_WRITE_SINGLE_REGISTER) {
-                cascade_slave_write(cas, h.reg, h.regcnt);
+                cascade_slave_write(h.reg, h.regcnt);
             }
 
             r = _mb_reply(cas, buff, rc, cas->map);
@@ -722,6 +742,7 @@ static int power_update(cascade_handle_t *cas)
     GlobalDeviceManager *dm2=get_dm2();
     slave_info_t *info=&cas->sInfo;
 
+    LOGD("______ power_update, %d\n", info->cnt);
     if(info->cnt==0 || list_empty(&dm2->_globalPowerManger.list)) {
         LOGE("___ sInfo.cnt is 0\n");
         return -1;
@@ -729,12 +750,12 @@ static int power_update(cascade_handle_t *cas)
 
     list_for_each_entry(tmp, &dm2->_globalPowerManger.list, list)
     {
-        int cnt2=0;
         if(tmp->product_ch_type==TREE_AC_TYPE) {
             if(list_empty(&tmp->list_Tree_AC)) {
                 continue;
             }
 
+            int cnt2=0;
             list_for_each_entry(tmp3,&tmp->list_Tree_AC,list_Tree_AC)
             {
                 tmp3->_PowerInfo = info->ch[cnt].pinfo[cnt2++].power;
@@ -742,8 +763,8 @@ static int power_update(cascade_handle_t *cas)
         }
         else {
             tmp->_PowerInfo = info->ch[cnt].pinfo[0].power;
-            cnt++;
         }
+        cnt++;
     }
     
     return 0;

+ 3 - 29
pro/src/cascade.h

@@ -189,29 +189,6 @@ typedef struct mmap_modbus{
 }mmap_modbus_t;
 //////////////////////////////////////
 
-typedef struct {
-    
-    int                 inited;
-    pthread_mutex_t     mutex;          //used for list lock  
-    pthread_mutex_t     lock;
-    
-    Modbus_Manger       m;
-    int                 addr;
-    slave_t             slaves[CASCADE_MAX+1];
-
-    ModbusInfo_t        mInfo;
-    modbus_mapping_t    *map;
-    modbus_mapping_t    *map2;
-
-    cmd_data_t          cmd;
-    int                 scanAddr;
-
-    slave_info_t        sInfo;      //current slave infomation
-
-    handle_t            list;
-}cascade_handle_t;
-
-
 extern int cur_dev_addr;
 
 int cascade_init(void);
@@ -220,7 +197,6 @@ int cascade_request(cmd_data_t *cmd);
 int cascade_get_all(_OverAllPwrAckInfo *all);
 int cascade_get_ch(_OverChnPwrAckInfo *ch);
 
-
 int cascade_get_dlist(dev_list_t *dl);
 int cascade_free_dlist(dev_list_t *dl);
 int cascade_set_modbus(ModbusInfo_t *info);
@@ -228,11 +204,9 @@ int cascade_set_modbus(ModbusInfo_t *info);
 int cascade_lock(void);
 int cascade_unlock(void);
 
-int cascade_slave_map_init(cascade_handle_t *cas);
-void cascade_slave_read(cascade_handle_t *cas, uint32_t addr,uint32_t lenth);
-void cascade_slave_write(cascade_handle_t *cas, uint32_t addr, uint16_t val);
-
-
+int cascade_slave_init(void);
+void cascade_slave_read(uint32_t addr,uint32_t lenth);
+void cascade_slave_write(uint32_t addr, uint16_t val);
 
 #endif
 

+ 98 - 98
pro/src/cascade_slave.c

@@ -117,14 +117,21 @@ static mmap_modbus_t modbus_three[] = {
     [TREE_W_ALL_CHN_STA     ]{.offset = TREE_ALL_SET_SWITCH_REG }
 };
 
-int cascade_slave_map_init(cascade_handle_t *cas)
+typedef struct { 
+    modbus_mapping_t    *map;
+}slave_handle_t;
+
+static slave_handle_t slHandle;
+
+int cascade_slave_init(void)
 {
-    GlobalDeviceManager* mgr = get_dm();
+    slave_handle_t *sh=&slHandle;
+    GlobalDeviceManager* mgr=get_dm();
 
     if(mgr->_global_device_info->product_pwr_type == SmartPDU_AC ||
         mgr->_global_device_info->product_pwr_type == SmartPDU_DC)       // single AC
     {
-        cas->map = modbus_mapping_new_start_address(0,0,0,0,DEV_CHN_REG,5000,0,0);     //3000 ge
+        sh->map = modbus_mapping_new_start_address(0,0,0,0,DEV_CHN_REG,5000,0,0);     //3000 ge
         modbus_single[DS_R_DEV_CHN          ].addr = &mgr->single_mmap->dev_chn;
         modbus_single[DS_R_OUT_INFO         ].addr = &mgr->single_mmap->all_ch;
         modbus_single[DS_R_IN_TOTAL         ].addr = &mgr->single_mmap->total_valtage;
@@ -139,8 +146,7 @@ int cascade_slave_map_init(cascade_handle_t *cas)
         printf("single dc address init 4000\n");
     }else
     {
-
-        cas->map =  modbus_mapping_new_start_address(0,0,0,0,TREE_R_DEV_CHN,7000,0,0);
+        sh->map =  modbus_mapping_new_start_address(0,0,0,0,TREE_R_DEV_CHN,7000,0,0);
         modbus_three[TREE_R_DEV_CHN         ].addr = &mgr->triphasic_mmap->dev_chn;
         modbus_three[TREE_R_INPUT_ALL_A     ].addr = &mgr->triphasic_mmap->i_total.pw_i_total[0];
         modbus_three[TREE_R_INPUT_ALL_B     ].addr = &mgr->triphasic_mmap->i_total.pw_i_total[1];
@@ -162,12 +168,10 @@ int cascade_slave_map_init(cascade_handle_t *cas)
 }
 
 
+static void slave_read(slave_handle_t *sh, uint32_t addr,int num,uint32_t reg_type,uint32_t type)
+{
+    GlobalDeviceManager* mgr=get_dm();
 
-
-
-static void slave_read(cascade_handle_t *cas,uint32_t addr,int num,uint32_t reg_type,uint32_t type)
-{   
-    GlobalDeviceManager* mgr = get_dm();
     if(mgr->_global_device_info->product_pwr_type == SmartPDU_AC || 
             mgr->_global_device_info->product_pwr_type == SmartPDU_DC){    
         if(addr >= DEV_CHN_REG && addr <= ALL_SET_SWITCH_REG)
@@ -175,26 +179,22 @@ static void slave_read(cascade_handle_t *cas,uint32_t addr,int num,uint32_t reg_
             if(type == TYPE_U16)
             {
                 uint16_t *base_addr = (uint16_t *)modbus_single[reg_type].addr;
-                uint32_t step = addr - cas->map->start_registers;
+                uint32_t step = addr - sh->map->start_registers;
                 base_addr = base_addr + (addr - modbus_single[reg_type].offset);
-                uint16_t *map_addr = (uint16_t *)&(cas->map->tab_registers[step]);
+                uint16_t *map_addr = (uint16_t *)&(sh->map->tab_registers[step]);
                 if(num < 126)
                 {
-                    pthread_mutex_lock(&cas->mutex);
                     memcpy(map_addr, base_addr,2*num);
-                    pthread_mutex_unlock(&cas->mutex);
                 }
             }else
             {
                 uint32_t *base_addr = (uint32_t *)modbus_single[reg_type].addr;
-                uint32_t step = addr - cas->map->start_registers;
+                uint32_t step = addr - sh->map->start_registers;
                 base_addr = base_addr + (addr - modbus_single[reg_type].offset);
-                uint32_t *map_addr = (uint32_t *)&(cas->map->tab_registers[step]);
+                uint32_t *map_addr = (uint32_t *)&(sh->map->tab_registers[step]);
                 if(num < 126)
                 {
-                    pthread_mutex_lock(&cas->mutex);
                     memcpy(map_addr, base_addr,2*num);
-                    pthread_mutex_unlock(&cas->mutex);
                 }
             }
         }
@@ -204,112 +204,112 @@ static void slave_read(cascade_handle_t *cas,uint32_t addr,int num,uint32_t reg_
             if(type == TYPE_U16)
             {
                 uint16_t *base_addr = (uint16_t *)modbus_three[reg_type].addr;
-                uint32_t step = addr - cas->map->start_registers;
+                uint32_t step = addr - sh->map->start_registers;
                 base_addr = base_addr + (addr - modbus_three[reg_type].offset);
-                uint16_t *map_addr = (uint16_t *)&(cas->map->tab_registers[step]);
+                uint16_t *map_addr = (uint16_t *)&(sh->map->tab_registers[step]);
                 if(num < 126)
                 {
-                    pthread_mutex_lock(&cas->mutex);
                     memcpy(map_addr, base_addr,2*num);
-                    pthread_mutex_unlock(&cas->mutex);
                 }
             }else
             {
                 uint32_t *base_addr = (uint32_t *)modbus_three[reg_type].addr;
-                uint32_t step = addr - cas->map->start_registers;
+                uint32_t step = addr - sh->map->start_registers;
                 base_addr = base_addr + (addr - modbus_three[reg_type].offset);
-                uint32_t *map_addr = (uint32_t *)&(cas->map->tab_registers[step]);
+                uint32_t *map_addr = (uint32_t *)&(sh->map->tab_registers[step]);
                 if(num < 126)
                 {
-                    pthread_mutex_lock(&cas->mutex);
                     memcpy(map_addr, base_addr,2*num);
-                    pthread_mutex_unlock(&cas->mutex);
                 }
             }
         }
     }
 }
 
-void cascade_slave_read(cascade_handle_t *cas,uint32_t addr,uint32_t lenth)
+
+void cascade_slave_read(uint32_t addr,uint32_t lenth)
 {
+    slave_handle_t *sh=&slHandle;
 
-        switch (addr)
-        {
-            case 4000:
-                slave_read(cas,addr,lenth,DS_R_DEV_CHN,TYPE_U16);
-            break;
-            case 4001 ... 4768:
-                slave_read(cas,addr,lenth,DS_R_OUT_INFO,TYPE_U32);
-            break;
-            case 4769 ... 4776:
-                slave_read(cas,addr,lenth,DS_R_IN_TOTAL,TYPE_U32);
-            break;
-            case 4777 ... 4778:
-                slave_read(cas,addr,lenth,DS_R_TEMP_HUM,TYPE_U32);
-            case 4779 ... 4842:
-                slave_read(cas,addr,lenth,DS_R_SWITCH_STA,TYPE_U16);
-            break;
-            case 4843 ... 4906:
-                slave_read(cas,addr,lenth,DS_R_OPEN_DELAY,TYPE_U16);
-            break;
-            case 4907 ... 4970:
-                slave_read(cas,addr,lenth,DS_R_CLOSE_DELAY,TYPE_U16);
-            break;
-            case 5000 ... 5063:
-                slave_read(cas,addr,lenth,DS_W_OPEN_DELAY,TYPE_U16);
-            break;
-            case 5064 ... 5127:
-                slave_read(cas,addr,lenth,DS_W_OPEN_DELAY,TYPE_U16);
-            break;
-            case 5128 ... 5191:
-                slave_read(cas,addr,lenth,DS_W_CHN_SWITCH_STA,TYPE_U16);
-            break;
-            case 5192 :
-                slave_read(cas,addr,lenth,DS_W_ALL_CHN_STA,TYPE_U16);
-            break;
+    switch (addr)
+    {
+        case 4000:
+            slave_read(sh,addr,lenth,DS_R_DEV_CHN,TYPE_U16);
+        break;
+        case 4001 ... 4768:
+            slave_read(sh,addr,lenth,DS_R_OUT_INFO,TYPE_U32);
+        break;
+        case 4769 ... 4776:
+            slave_read(sh,addr,lenth,DS_R_IN_TOTAL,TYPE_U32);
+        break;
+        case 4777 ... 4778:
+            slave_read(sh,addr,lenth,DS_R_TEMP_HUM,TYPE_U32);
+        case 4779 ... 4842:
+            slave_read(sh,addr,lenth,DS_R_SWITCH_STA,TYPE_U16);
+        break;
+        case 4843 ... 4906:
+            slave_read(sh,addr,lenth,DS_R_OPEN_DELAY,TYPE_U16);
+        break;
+        case 4907 ... 4970:
+            slave_read(sh,addr,lenth,DS_R_CLOSE_DELAY,TYPE_U16);
+        break;
+        case 5000 ... 5063:
+            slave_read(sh,addr,lenth,DS_W_OPEN_DELAY,TYPE_U16);
+        break;
+        case 5064 ... 5127:
+            slave_read(sh,addr,lenth,DS_W_OPEN_DELAY,TYPE_U16);
+        break;
+        case 5128 ... 5191:
+            slave_read(sh,addr,lenth,DS_W_CHN_SWITCH_STA,TYPE_U16);
+        break;
+        case 5192 :
+            slave_read(sh,addr,lenth,DS_W_ALL_CHN_STA,TYPE_U16);
+        break;
+
+        // triple
+        case 6000:
+            slave_read(sh,addr,lenth,TREE_R_DEV_CHN,TYPE_U16);
+        break;
+        case 6001 ... 6008:
+            slave_read(sh,addr,lenth,TREE_R_INPUT_ALL_A,TYPE_U32);
+        break;
+        case 6009 ... 6016:
+            slave_read(sh,addr,lenth,TREE_R_INPUT_ALL_B,TYPE_U32);
+        break;
+        case 6017 ... 6024:
+            slave_read(sh,addr,lenth,TREE_R_INPUT_ALL_C,TYPE_U32);
+        break;
+        case 6025 ... 6664:
+            slave_read(sh,addr,lenth,TREE_R_CHN_TOTAL,TYPE_U32);
+        break;
+        case 6665 ... 8584:
+            slave_read(sh,addr,lenth,TREE_R_CHN_PH_INFO,TYPE_U32);
+        break;
+        case 8585 ... 8586:
+            slave_read(sh,addr,lenth,TREE_R_TEMP_HUM,TYPE_U16);
+        break;
+        case 8587 ... 8650:
+            slave_read(sh,addr,lenth,TREE_R_SWITCH_STA,TYPE_U16);
+        break;
+        case 8651 ... 8714:
+            slave_read(sh,addr,lenth,TREE_R_OPEN_DELAY,TYPE_U16);
+        break;
+        case 8715 ... 8778:
+            slave_read(sh,addr,lenth,TREE_R_CLOSE_DELAY,TYPE_U16);
+        break;
 
-            // triple
-            case 6000:
-                slave_read(cas,addr,lenth,TREE_R_DEV_CHN,TYPE_U16);
-            break;
-            case 6001 ... 6008:
-                slave_read(cas,addr,lenth,TREE_R_INPUT_ALL_A,TYPE_U32);
-            break;
-            case 6009 ... 6016:
-                slave_read(cas,addr,lenth,TREE_R_INPUT_ALL_B,TYPE_U32);
-            break;
-            case 6017 ... 6024:
-                slave_read(cas,addr,lenth,TREE_R_INPUT_ALL_C,TYPE_U32);
-            break;
-            case 6025 ... 6664:
-                slave_read(cas,addr,lenth,TREE_R_CHN_TOTAL,TYPE_U32);
-            break;
-            case 6665 ... 8584:
-                slave_read(cas,addr,lenth,TREE_R_CHN_PH_INFO,TYPE_U32);
-            break;
-            case 8585 ... 8586:
-                slave_read(cas,addr,lenth,TREE_R_TEMP_HUM,TYPE_U16);
-            break;
-            case 8587 ... 8650:
-                slave_read(cas,addr,lenth,TREE_R_SWITCH_STA,TYPE_U16);
-            break;
-            case 8651 ... 8714:
-                slave_read(cas,addr,lenth,TREE_R_OPEN_DELAY,TYPE_U16);
-            break;
-            case 8715 ... 8778:
-                slave_read(cas,addr,lenth,TREE_R_CLOSE_DELAY,TYPE_U16);
-            break;
         default:
-            break;
-        }
+        break;
+    }
 }
 
 
-
-void cascade_slave_write(cascade_handle_t *cas, uint32_t addr, uint16_t val)
-{   
-    GlobalDeviceManager* mgr = get_dm();
+void cascade_slave_write(uint32_t addr, uint16_t val)
+{
+    slave_handle_t *sh=&slHandle;
+    GlobalDeviceManager* mgr=get_dm();
     uint16_t buff[128] = {0};
+
     if(mgr->_global_device_info->product_pwr_type == SmartPDU_AC 
             || mgr->_global_device_info->product_pwr_type == SmartPDU_DC)   
     {

+ 0 - 1
pro/src/common.h

@@ -826,7 +826,6 @@ int days_in_month(int year, int month);
 // 函数:增加一天到日期  
 void add_day_to_date(char *date_str, int *year, int *month, int *day);
 
-pthread_mutex_t WSlock; 
 pthread_mutex_t Httplock; 
 
 #endif

+ 62 - 70
pro/src/websocket_handle.c

@@ -15,16 +15,13 @@ typedef struct mg_http_message  mg_http_msg_t;
 
 
 typedef struct {
+    mg_mgr_t  mgr;
+    char      wpath[200];
+    char      root[200];
 
-  mg_mgr_t  mgr;
-  char      wpath[200];
-  char      root[200];
-
-  mg_conn_t *c;
-
-  int       inited;
-  pthread_mutex_t mutex;
-
+    int       inited;
+    int       ws_cnt;
+    pthread_mutex_t mutex;
 }ws_handle_t;
 static ws_handle_t wsHandle;
 
@@ -49,44 +46,38 @@ static void fn(mg_conn_t *c, int ev, void *ev_data)
 
     case MG_EV_CLOSE:
     {
-      if (wh->c)
+      if (c->is_websocket)
       {
-        if (c->id == wh->c->id)
-        {
-          wh->c = NULL;
+        if(wh->ws_cnt>0) {
+            wh->ws_cnt--;
         }
       }
-      log_i("ws upgrade del connection ID:%d\n", c->id);
+      //log_i("ws upgrade del connection ID:%d\n", c->id);
+    }
+    break;
+
+    case MG_EV_WS_OPEN:
+    {
+        wh->ws_cnt++;
     }
     break;
 
     case MG_EV_HTTP_MSG:
     {
-       static  int ws_cnt=0;
-      mg_http_msg_t *hm = (mg_http_msg_t *) ev_data;
-      if (mg_http_match_uri(hm, "/websocket/PDU-WebSocket")) 
-       {
-          mg_ws_upgrade(c, hm, NULL);
-          wh->c = c;
-           log_i("ws upgrade new connection ID:%d\n",c->id);
-      }
-      else if (mg_http_match_uri(hm, "/rest")) 
-      {
-          // Serve REST response
-          mg_http_reply(c, 200, "", "{\"result\": %d}\n", 123);
-      } 
-      else 
-      {
-           wh->c = NULL;
-           log_i("ws upgrade del connection!");
-      }
+        mg_http_msg_t *hm = (mg_http_msg_t *) ev_data;
+        if (mg_http_match_uri(hm, "/websocket/PDU-WebSocket")) {
+            mg_ws_upgrade(c, hm, NULL);
+        }
+        else if (mg_http_match_uri(hm, "/rest")) {
+            mg_http_reply(c, 200, "", "{\"result\": %d}\n", 123);
+        }
     }
     break;
 
     case MG_EV_WS_MSG:
     {
-     // struct mg_ws_message *wm = (struct mg_ws_message *) ev_data;
-     // mg_ws_send(c, wm->data.ptr, wm->data.len, WEBSOCKET_OP_TEXT);
+        // struct mg_ws_message *wm = (struct mg_ws_message *) ev_data;
+        // mg_ws_send(c, wm->data.ptr, wm->data.len, WEBSOCKET_OP_TEXT);
     }
     break;
   }
@@ -102,26 +93,26 @@ void* websocket_thread(void* arg)
     mg_http_listen(&wh->mgr, wh->wpath, fn, NULL);  // Create HTTP listener
     wh->inited = 1;
 
-    /*while(h->quit==0) {
-          pthread_mutex_lock(&WSlock);
-          mg_mgr_poll(&wh->mgr, 1000);                    // Infinite event loop
-          pthread_mutex_unlock(&WSlock);
-          usleep(100000);
+    while(h->quit==0) {
+        pthread_mutex_lock(&wh->mutex);
+        mg_mgr_poll(&wh->mgr, 200);                    // Infinite event loop
+        pthread_mutex_unlock(&wh->mutex);
     }
     wh->inited = 0;
     mg_mgr_free(&wh->mgr);
-*/
+
     pthread_exit(NULL);
 }
 
 
 int websocket_init(void) 
 {
-        // ³õʼ»¯»¥³âËø  
-  if (pthread_mutex_init(&WSlock, NULL) != 0) {  
-        return -1;  
-  }  
     ws_handle_t *wh=&wsHandle;
+
+    if (pthread_mutex_init(&wh->mutex, NULL) != 0) {  
+        log_e("websocket mutex init failed!\n");
+        return -1;  
+    }  
     
     memset(wh, 0, sizeof(ws_handle_t));
     read_ipaddr(wh);
@@ -131,57 +122,58 @@ int websocket_init(void)
 }
 
 
-int websocket_send(void *data, int len) 
+int websocket_send(void *conn, void *data, int len) 
 {
     int i;
     ws_handle_t *wh=&wsHandle;
 
-    if(!wh->inited) {
+    if(!wh->inited || !conn) {
         return -1;
     }
-    pthread_mutex_lock(&WSlock);
-    mg_ws_send(wh->c, data, len, WEBSOCKET_OP_TEXT);
-    pthread_mutex_unlock(&WSlock);
+    pthread_mutex_lock(&wh->mutex);
+    mg_ws_send(conn, data, len, WEBSOCKET_OP_TEXT);
+    pthread_mutex_unlock(&wh->mutex);
     return 0;
 }
 
-
-int websocket_isok(void) 
+int websocket_broadcast(void *data, int len) 
 {
+    int i;
+    mg_conn_t *c;
     ws_handle_t *wh=&wsHandle;
-    return (wh->c)?1:0;
-}
 
-int websocket_poll(void) 
-{
-    int r;
-
-    ws_handle_t *wh=&wsHandle;
-    if (!wh->inited)
-    {
-      return -1;
+    if(!wh->inited) {
+        return -1;
     }
 
-    pthread_mutex_lock(&WSlock); 
-	  mg_mgr_poll(&wh->mgr, 1000);                    // Infinite event loop
-    pthread_mutex_unlock(&WSlock); 
-
+    pthread_mutex_lock(&wh->mutex);
+    for (c=wh->mgr.conns; c!=NULL; c=c->next) {
+        mg_ws_send(c, data, len, WEBSOCKET_OP_TEXT);
+    }
+    pthread_mutex_unlock(&wh->mutex);
 
     return 0;
 }
 
+int websocket_is_online(void) 
+{
+    ws_handle_t *wh=&wsHandle;
+    return (wh->ws_cnt>0)?1:0;
+}
+
+
 int websocket_free(void) 
 {
     int r;
-
     ws_handle_t *wh=&wsHandle;
-    if (!wh->inited)
-    {
+
+    if (!wh->inited) {
       return -1;
     }
 
-    wh->inited = 0;
-    mg_mgr_free(&wh->mgr);
+    thread_stop(THREAD_ID_WS);
+    pthread_mutex_destroy(&wh->mutex);
+    
     return 0;
 }
 

+ 3 - 3
pro/src/websocket_handle.h

@@ -11,9 +11,9 @@ typedef struct {
 
 
 int websocket_init(void);
-int websocket_send(void *data, int len);
-int websocket_isok(void);
-int websocket_poll(void) ;
+int websocket_send(void *conn, void *data, int len);
+int websocket_broadcast(void *data, int len);
+int websocket_is_online(void);
 int websocket_free(void) ;
 #endif