#include "websocket_handle.h" #include "mongoose.h" #include "pthread.h" #include "elog.h" #include "common.h" #include "sys.h" #include "upg.h" #include "lock.h" #include "file.h" #include "thread.h" #include "sqlite_handle.h" #include "json_handle.h" #include "cascade.h" #include "cfg.h" #include "wifi.h" #include "Cellular.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; 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 ws_all_update(ws_handle_t *h) { char *json; int r,r1,r2,r3; GlobalDeviceManager* gdm= &__globalDeviceManage; GlobalDeviceManager* gdm2=&__globalDeviceManage2; GlobalDeviceManager* pdm=NULL; pwrall_info_t *pall=&h->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 = over_all_pwr_monitor_Tree_AC_to_json(&pall->ch[0], &pall->ch[1], &pall->ch[2]); ws_send_data(h, json, strlen(json)); cJSON_free(json); } } 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 = ws_over_status_ack_to_json(0, &pall->ch[0]); ws_send_data(h, json, strlen(json)); cJSON_free(json); } } return 0; } static int ws_chn_update(ws_handle_t *h) { int r; char *json; _OverChnPwrAckInfo chInfo; GlobalDeviceManager* gdm= &__globalDeviceManage; GlobalDeviceManager* gdm2=&__globalDeviceManage2; GlobalDeviceManager* pdm=NULL; 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 = ws_chn_status_ack_to_json(1, &chInfo, pdm); ws_send_data(h, json, strlen(json)); cJSON_free(json); } _OverChnPwrAckInfo *node, *next; list_for_each_entry_safe(node, next, &chInfo.list, list) { list_del(&node->list); free(node); } return 0; } static int ws_sensor_update(ws_handle_t *h) { char *json; GlobalDeviceManager* gdm= &__globalDeviceManage; //lock_s_hold(LOCK_ID_SENSOR); json = ws_sensor_status_ack_to_json(&gdm->_globalSensorManger); //lock_s_release(LOCK_ID_SENSOR); if(json) { ws_send_data(h, json, strlen(json)); cJSON_free(json); } return 0; } static int ws_breaker_update(ws_handle_t *h) { char *json; GlobalDeviceManager* gdm= &__globalDeviceManage; GlobalDeviceManager* gdm2=&__globalDeviceManage2; if(cur_dev_addr==0) { json = breaker_status_ack_to_json(gdm,0); } else { json = breaker_status_ack_to_json(gdm2,cur_dev_addr); } ws_send_data(h, json, strlen(json)); cJSON_free(json); return 0; } static int ws_threshold_update(ws_handle_t *h) { char *json; GlobalDeviceManager* dm= &__globalDeviceManage; if(cur_dev_addr==0) { json = threshold_status_ack_to_json(dm); ws_send_data(h, json, strlen(json)); cJSON_free(json); } return 0; } static int ws_upg_update(ws_handle_t *h) { handle_t l=upg_board_all_get(); char* json = board_to_json(l, 1); ws_send_data(h, json, strlen(json)); cJSON_free(json); return 0; } static int ws_Cellular_update(ws_handle_t *h) { char *json; if (__globalDeviceManage._globalDevInfo.cell.Cellular_enable) { json = communication_4G_ack_to_json(&__globalDeviceManage._globalDevInfo.cell, 1); if (json) { ws_send_data(h, json, strlen(json)); cJSON_free(json); } } return 0; } #if (CHIP_TYPE != CHIP_T113s) static int ws_WIFI_update(ws_handle_t *h) { char *json; if (__globalDeviceManage._globalDevInfo.wifi.WIFI_enable) { ap_stat_t stat; WIFIInfo_t *wm = &__globalDeviceManage._globalDevInfo.wifi; wifi_sta_stat(&stat); wm->Link_status = stat.link; strcpy(wm->WIFI_ip, stat.ipaddr); json = communication_WIFI_ack_to_json(wm, 1); if (json) { ws_send_data(h, json, strlen(json)); cJSON_free(json); } } return 0; } #endif static int ws_DS_update(ws_handle_t *h) { char *json; json = service_ds_all_ack_to_json("","","0",&__globalDeviceManage._globalPowerManger, 1); if (json) { ws_send_data(h, json, strlen(json)); cJSON_free(json); } return 0; } static int ws_MQTT_update(ws_handle_t *h) { char *json; json = service_mqtt_ack_to_json(&__globalDeviceManage.mqttInfo, 1); if (json) { ws_send_data(h, json, strlen(json)); cJSON_free(json); } return 0; } ///////////////////////////////////////////////////////////// static int ws_update(ws_handle_t *h) { int r; if(!ws_is_online(h)) { return -1; } //overall refresh r = ws_all_update(h); //channel refresh r = ws_chn_update(h); //sensor refresh r = ws_sensor_update(h); //breaker refresh r = ws_breaker_update(h); //threshold update //r = ws_threshold_update(); //cellular refresh r = ws_Cellular_update(h); //wifi refresh #if (CHIP_TYPE != CHIP_T113s) r = ws_WIFI_update(h); #endif //DS refresh r = ws_DS_update(h); //MQTT refresh r = ws_MQTT_update(h); //upgrade refresh r = ws_upg_update(h); return r; } 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_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("/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 timer_fn(void *arg) { ws_handle_t *h=(ws_handle_t*)arg; ws_update(h); } static void* ws_thread(void* arg) { int r=0; thread_handle_t *th=(thread_handle_t*)arg; ws_handle_t *h=(ws_handle_t*)th->arg; const char *ws_addr="ws://[::]:6785"; const char *wss_addr="wss://[::]:6786"; //mg_log_set(MG_LL_DEBUG); mg_mgr_init(&h->mgr); // Initialise event manager mg_http_listen(&h->mgr, ws_addr, fn, NULL); // Create HTTP listener mg_http_listen(&h->mgr, wss_addr, fn, (void*)1); #if (CHIP_TYPE == CHIP_T113s) h->tls.cert = mg_str(file_load2("/mnt/UDISK/app/cert/server.crt",0)); h->tls.key = mg_str(file_load2("/mnt/UDISK/app/cert/server.key",0)); #else h->tls.cert = mg_str(file_load2("/root/run/app/cert/server.crt",0)); h->tls.key = mg_str(file_load2("/root/run/app/cert/server.key",0)); #endif h->inited = 1; mg_timer_add(&h->mgr, 500, MG_TIMER_REPEAT, timer_fn, h); while(th->quit==0) { mg_mgr_poll(&h->mgr, 100); // Infinite event loop } h->inited = 0; mg_mgr_free(&h->mgr); pthread_exit(NULL); } int websocket_init(void) { ws_handle_t *wh=&wsHandle; memset(wh, 0, sizeof(ws_handle_t)); thread_start(THREAD_ID_WS, ws_thread, wh); return 0; } pwrall_info_t *websocket_get_pwrall(void) { return &wsHandle.pwrall; }