#include "websocket_handle.h" #include "mongoose.h" #include "pthread.h" #include "elog.h" #include "common.h" #include "sys.h" #include "lock.h" #include "file.h" #include "thread.h" #include "sqlite_handle.h" #include "json_handle.h" #include "cascade.h" #include "cfg.h" #define USE_WS2 #define WS_SEND_TLEN_MAX (1024*1024*10) #define WS_MAX 10 #define USE_ONE_CONN 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_tls_opts mg_tls_t; typedef struct { mg_mgr_t mgr; char wpath[200]; char root[200]; mg_tls_t tls; mg_conn_t *c; mg_conn_t *upg; int inited; int ws_cnt; pwrall_info_t pwrall; }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) || (!data) || (len<=0)) { return -1; } #ifdef USE_WS2 mg_ws_send2(wh->c, data, len, isBinary?WEBSOCKET_OP_BINARY:WEBSOCKET_OP_TEXT, WS_SEND_TLEN_MAX); #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, WS_SEND_TLEN_MAX); #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 int get_board_info(board_info_t *info) { return 0; } static int ws_upg_update(ws_handle_t *wh) { if((!wh->upg)) { return -1; } extern board_list_t *get_board_list(void); board_list_t *l=get_board_list(); char* json = board_to_json(l); mg_ws_send2(wh->upg, json, strlen(json), WEBSOCKET_OP_TEXT, WS_SEND_TLEN_MAX); return 0; } 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 ; _OverChnPwrAckInfo chInfo; GlobalDeviceManager* gdm= &__globalDeviceManage; GlobalDeviceManager* gdm2=&__globalDeviceManage2; GlobalDeviceManager* pdm=NULL; int r,r1,r2,r3; if(ws_is_online(wh)) { //overall info pwrall_info_t *pall=&wh->pwrall; 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) { pall->ph3 = 1; if (cur_dev_addr == 0) { lock_s_hold(LOCK_ID_POWER_UPDATE); r1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &pall->ch[0], gdm); r2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &pall->ch[1], gdm); r3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &pall->ch[2], gdm); lock_s_release(LOCK_ID_POWER_UPDATE); } else { cascade_lock(); r1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &pall->ch[0], gdm2); r2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &pall->ch[1], gdm2); r3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &pall->ch[2], gdm2); cascade_unlock(); } if (0 == r1 && 0 == r2 && 0 == r3) { json_str = over_all_pwr_monitor_Tree_AC_to_json(&pall->ch[0], &pall->ch[1], &pall->ch[2]); ws_send_data(wh, json_str, strlen(json_str)); cJSON_free((void *)json_str); } } else { pall->ph3 = 0; if(cur_dev_addr==0) { lock_s_hold(LOCK_ID_POWER_UPDATE); r = dev_search_latest_power_statistic_info(gdm, &pall->ch[0]); lock_s_release(LOCK_ID_POWER_UPDATE); } else { cascade_lock(); r = dev_search_latest_power_statistic_info(gdm2, &pall->ch[0]); cascade_unlock(); } if (0==r) { json_str = ws_over_status_ack_to_json(0, &pall->ch[0]); ws_send_data(wh, json_str, strlen(json_str)); cJSON_free(json_str); } } //channel info { INIT_LIST_HEAD(&chInfo.list); if(cur_dev_addr==0) { lock_s_hold(LOCK_ID_POWER_UPDATE); r = dev_search_latest_power_All_info(gdm,&chInfo, cur_dev_addr); lock_s_release(LOCK_ID_POWER_UPDATE); pdm = gdm; } else { cascade_lock(); r = dev_search_latest_power_All_info(gdm2,&chInfo, cur_dev_addr); cascade_unlock(); pdm = gdm2; } if (0 == r) { 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 { //lock_s_hold(LOCK_ID_SENSOR); json_str = ws_sensor_status_ack_to_json(&gdm->_globalSensorManger); //lock_s_release(LOCK_ID_SENSOR); if(json_str) { ws_send_data(wh, json_str, strlen(json_str)); cJSON_free(json_str); } } // 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); } //upgrade refresh } } 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; } if(c==wh->upg) { wh->upg = 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_ACCEPT: { if (c->fn_data) { mg_tls_init(c, &wh->tls); } } 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("/websocket/controlBoardUpgrade"), NULL)) { mg_ws_upgrade(c, hm, NULL); wh->upg = c; } 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; } } static void* ws_thread(void* arg) { int r; thread_handle_t *h=(thread_handle_t*)arg; ws_handle_t *wh=(ws_handle_t*)h->arg; const char *ws_addr="ws://[::]:6785"; const char *wss_addr="wss://[::]:6786"; //mg_log_set(MG_LL_DEBUG); mg_mgr_init(&wh->mgr); // Initialise event manager mg_http_listen(&wh->mgr, ws_addr, fn, NULL); // Create HTTP listener mg_http_listen(&wh->mgr, wss_addr, fn, (void*)1); wh->tls.cert = mg_str(file_load2("/root/run/app/cert/server.crt",0)); wh->tls.key = mg_str(file_load2("/root/run/app/cert/server.key",0)); 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, ws_thread, wh); return 0; } pwrall_info_t *websocket_get_pwrall(void) { return &wsHandle.pwrall; }