#include "websocket_handle.h" #include "mongoose.h" #include "pthread.h" #include "elog.h" #include "common.h" #include "sys.h" #include "thread.h" #include "sqlite_handle.h" #include "json_handle.h" #include "cascade.h" //#define USE_WS2 #define USE_ONE_CONN #define WS_MAX 10 typedef struct mg_mgr mg_mgr_t; typedef struct mg_connection mg_conn_t; typedef struct mg_ws_message mg_ws_msg_t; typedef struct mg_http_message mg_http_msg_t; typedef struct { mg_mgr_t mgr; char wpath[200]; char root[200]; mg_conn_t *c; int inited; int ws_cnt; }ws_handle_t; static ws_handle_t wsHandle; static int ws_is_online(ws_handle_t *wh) { #ifdef USE_ONE_CONN return (wh->c)?1:0; #else return (wh->ws_cnt>0)?1:0; #endif } static int ws_send(ws_handle_t *wh, void *data, int len, int isBinary) { if(!wh->c) { return -1; } #ifdef USE_WS2 mg_ws_send2(wh->c, data, len, isBinary?WEBSOCKET_OP_BINARY:WEBSOCKET_OP_TEXT); #else mg_ws_send(wh->c, data, len, isBinary?WEBSOCKET_OP_BINARY:WEBSOCKET_OP_TEXT); #endif return 0; } static int ws_broadcast(ws_handle_t *wh, void *data, int len, int isBinary) { int r=-1; mg_conn_t *c; for (c=wh->mgr.conns; c!=NULL; c=c->next) { #ifdef USE_WS2 mg_ws_send2(c, data, len, isBinary?WEBSOCKET_OP_BINARY:WEBSOCKET_OP_TEXT); #else mg_ws_send(c, data, len, isBinary?WEBSOCKET_OP_BINARY:WEBSOCKET_OP_TEXT); #endif r = 0; } return r; } static int ws_send_data(ws_handle_t *wh, void *data, int len) { #ifdef USE_ONE_CONN return ws_send(wh, data, len, 0); #else return ws_broadcast(wh, data, len, 0); #endif } static void read_ipaddr(ws_handle_t *h) { sprintf(h->wpath,"ws://%s:6785", getLocalIpAddress("eth0")); printf("____ ws path: %s\n", h->wpath); } static void web_update(ws_handle_t *wh) { GlobalSensorManger* _globalSensorMangerTemp; int ret = 0 ; GlobalSensorInfo _globalSensorInfo; struct tm* t ; struct timeval tv; struct timezone tz ; char* json_str = NULL ; GlobalPowerInfo _powerInfo; GlobalPowerManger* _globalPowerMangerTemp = NULL ; _OverAllPwrAckInfo allInfo; //设备电源信息 _OverChnPwrAckInfo chInfo; GlobalDeviceManager* gdm= &__globalDeviceManage; GlobalDeviceManager* gdm2=&__globalDeviceManage2; GlobalDeviceManager* pdm=NULL; if(ws_is_online(wh)) { //overall info if(gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_Tree || gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_One || gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_One_B) { int nRet_L1 = -1; int nRet_L2 = -1; int nRet_L3 = -1; _OverAllPwrAckInfo l1, l2, l3; if (cur_dev_addr == 0) { pthread_mutex_lock(&gdm->_power_update); nRet_L1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &l1, gdm); nRet_L2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &l2, gdm); nRet_L3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &l3, gdm); pthread_mutex_unlock(&gdm->_power_update); } else { cascade_lock(); nRet_L1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &l1, gdm2); nRet_L2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &l2, gdm2); nRet_L3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &l3, gdm2); cascade_unlock(); } if (0 == nRet_L1 && 0 == nRet_L2 && 0 == nRet_L3) { json_str = over_all_pwr_monitor_Tree_AC_to_json(&l1, &l2, &l3); ws_send_data(wh, json_str, strlen(json_str)); cJSON_free((void *)json_str); } } else { int nRetTotal = -1; if(cur_dev_addr==0) { pthread_mutex_lock(&gdm->_power_update); nRetTotal=dev_search_latest_power_statistic_info(gdm, &allInfo,true); pthread_mutex_unlock(&gdm->_power_update); } else { cascade_lock(); nRetTotal=dev_search_latest_power_statistic_info(gdm2, &allInfo,false); cascade_unlock(); } if (0==nRetTotal) { json_str = ws_over_status_ack_to_json(0, &allInfo); ws_send_data(wh, json_str, strlen(json_str)); cJSON_free(json_str); } } //channel info { INIT_LIST_HEAD(&chInfo.list); int nRetCh = -1; if(cur_dev_addr==0) { pthread_mutex_lock(&gdm->_power_update); //nRetCh=dev_search_latest_power_all_info_ip(gdm,&chInfo, cur_dev_addr); nRetCh=dev_search_latest_power_All_info(gdm,&chInfo, cur_dev_addr); pthread_mutex_unlock(&gdm->_power_update); pdm = gdm; } else { cascade_lock(); nRetCh=dev_search_latest_power_All_info(gdm2,&chInfo, cur_dev_addr); cascade_unlock(); pdm = gdm2; } if (0 == nRetCh) { json_str = ws_chn_status_ack_to_json(1, &chInfo, pdm); ws_send_data(wh, json_str, strlen(json_str)); cJSON_free((void *)json_str); } _OverChnPwrAckInfo *node, *next; list_for_each_entry_safe(node, next, &chInfo.list, list) { list_del(&node->list); free(node); } } //sensor info { pthread_mutex_lock(&__globalDeviceManage._sensor_update); json_str = ws_sensor_status_ack_to_json(&gdm->_globalSensorManger); pthread_mutex_unlock(&__globalDeviceManage._sensor_update); ws_send_data(wh, json_str, strlen(json_str)); cJSON_free(json_str); } #if 0 // // breaker info { if(cur_dev_addr==0) { json_str = breaker_status_ack_to_json(gdm,0); }else { json_str = breaker_status_ack_to_json(gdm2,cur_dev_addr); } ws_send_data(wh, json_str, strlen(json_str)); cJSON_free(json_str); } #endif } } static void fn(mg_conn_t *c, int ev, void *ev_data) { ws_handle_t *wh=&wsHandle; switch(ev) { case MG_EV_OPEN: { } break; case MG_EV_CLOSE: { if (c->is_websocket ) { if(wh->ws_cnt>0) { wh->ws_cnt--; } if(c==wh->c) { wh->c = NULL; } } 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: { mg_http_msg_t *hm = (mg_http_msg_t *) ev_data; if (mg_match(hm->uri, mg_str("/websocket/PDU-WebSocket"), NULL)) { mg_ws_upgrade(c, hm, NULL); wh->c = c; log_i("ws upgrade new connection ID:%d\n", c->id); } else if (mg_match(hm->uri, mg_str("/rest"), NULL)) { 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); } break; } } void* websocket_thread(void* arg) { int r; thread_handle_t *h=(thread_handle_t*)arg; ws_handle_t *wh=(ws_handle_t*)h->arg; const char *listen_addr="ws://[::]:6785"; mg_mgr_init(&wh->mgr); // Initialise event manager mg_http_listen(&wh->mgr, listen_addr, fn, NULL); // Create HTTP listener wh->inited = 1; while(h->quit==0) { if (wh->c) { web_update(wh); } mg_mgr_poll(&wh->mgr, 1000); // Infinite event loop usleep(500000); } wh->inited = 0; mg_mgr_free(&wh->mgr); pthread_exit(NULL); } int websocket_init(void) { ws_handle_t *wh=&wsHandle; memset(wh, 0, sizeof(ws_handle_t)); //read_ipaddr(wh); thread_start(THREAD_ID_WS, websocket_thread, wh, 30*MB, 0); return 0; }