mqtt.c 32 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246
  1. #include <regex.h>
  2. #include <time.h>
  3. #include "common.h"
  4. #include "cJSON.h"
  5. #include "mqtt.h"
  6. #include "thread.h"
  7. #include "mongoose.h"
  8. #include "elog.h"
  9. #include "lock.h"
  10. #include "sys.h"
  11. #include "Cellular.h"
  12. #include "xlist.h"
  13. #include "websocket_handle.h"
  14. #include "switch_ctrl.h"
  15. #include "sqlite_handle.h"
  16. #if(CHIP_TYPE == CHIP_T113s)
  17. #include <openssl/evp.h>
  18. #include <openssl/hmac.h>
  19. #endif
  20. #if 1
  21. #define LOGD log_d
  22. #define LOGE log_e
  23. #define LOGW log_w
  24. #else
  25. #define LOGD printf
  26. #define LOGE printf
  27. #define LOGW printf
  28. #endif
  29. #define CONN_PERIOD 5000 //ms
  30. #define SEND_PERIOD 5000 //ms
  31. #define BGET(flag,mask) ((flag)&(1<<(mask)))
  32. #define BSET(flag,mask) ((flag)|=(1<<(mask)))
  33. typedef struct mg_mgr mgr_t;
  34. typedef struct mg_timer mg_timer_t;
  35. typedef struct mg_mqtt_opts mg_opts_t;
  36. typedef struct mg_connection mg_conn_t;
  37. typedef struct {
  38. int id;
  39. char topic[256];
  40. char content[2048];
  41. }mqtt_pkt_t;
  42. typedef struct {
  43. mgr_t mgr;
  44. mg_conn_t *c;
  45. mqtt_server_t ser;
  46. pthread_t tid;
  47. int quit;
  48. void *h;
  49. }mqtt_conn_t;
  50. typedef struct {
  51. int inited;
  52. uint32_t prod_id;
  53. mqtt_conn_t conn[MQTT_SER_MAX];
  54. handle_t list;
  55. mg_timer_t *timer;
  56. GlobalDeviceManager *dm;
  57. }mqtt_handle_t;
  58. #ifdef USE_MQTT
  59. char *topic_sub[MQTT_SUB_MAX]={
  60. "/pdu/%d/control/power/all_channel",
  61. "/pdu/%d/control/power/+",
  62. "/pdu/%d/control/sensor/+",
  63. "/pdu/%d/control/device/restart",
  64. "/pdu/%d/control/device/reset",
  65. "/pdu/%d/control/service",
  66. };
  67. char *topic_pub[MQTT_PUB_MAX]={
  68. "/pdu/%d/info/device",
  69. "/pdu/%d/status/network",
  70. "/pdu/%d/status/power/all_channel",
  71. "/pdu/%d/status/power/%d",
  72. "/pdu/%d/status/sensor/%d",
  73. "/pdu/%d/alarm/network",
  74. "/pdu/%d/alarm/power",
  75. "/pdu/%d/alarm/sensor",
  76. "/pdu/%d/status/service",
  77. };
  78. const char *serv_str[SERV_MAX]={
  79. "ntp",
  80. "smtp",
  81. "message",
  82. "mqtt",
  83. "cloud",
  84. "telnet",
  85. "snmp_v1",
  86. "snmp_v2c",
  87. "snmp_v3",
  88. "snmp_trap",
  89. };
  90. static mqtt_handle_t mqHandle={0};
  91. static int my_recv(char *topic, char *data);
  92. static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data);
  93. static int send_stat(mqtt_conn_t *conn, int type);
  94. static int get_serv(char *json);
  95. static void* client_thread(void *arg);
  96. /////////////////////////////////////////////////////////
  97. typedef struct {
  98. char user[64];
  99. char passwd[64];
  100. }user_account_t;
  101. typedef struct {
  102. char *instanceId;
  103. char *host;
  104. char *topic;
  105. char *groupId;
  106. char *clientId;
  107. char *accessKey;
  108. char *secretKey;
  109. uint16_t port;
  110. }mqtt_account_t;
  111. static mqtt_account_t mqtt_aliyun={
  112. "post-cn-0w73xnchz01",
  113. "post-cn-0w73xnchz01.mqtt.aliyuncs.com",
  114. "SmartPDU",
  115. "GID_PDU",
  116. "smartPDU_23542352",
  117. "LTAI5tAfuRj1JkB8ZjDS2p3i",
  118. "p4FrWlrFim9zM9A5mYxjK99EqKrSHh",
  119. 8883,
  120. };
  121. static int get_account(mqtt_account_t *host, user_account_t *user)
  122. {
  123. unsigned int len=0;
  124. char tempData[100];
  125. char clientIdUrl[64];
  126. //username和 Password 签名模式下的设置方法,参考文档 https://help.aliyun.com/document_detail/48271.html?spm=a2c4g.11186623.6.553.217831c3BSFry7
  127. sprintf(clientIdUrl, "%s@@@%s", host->groupId, host->clientId);
  128. HMAC(EVP_sha1(), host->secretKey, strlen(host->secretKey), (const unsigned char*)clientIdUrl, strlen(clientIdUrl), (unsigned char*)tempData, &len);
  129. int passwdLen = EVP_EncodeBlock((unsigned char *) user->passwd, (const unsigned char*)tempData, len);
  130. user->passwd[passwdLen] = '\0';
  131. sprintf(user->user,"Signature|%s|%s", host->accessKey, host->instanceId);
  132. return 0;
  133. }
  134. static void set_status(mqtt_conn_t *conn, int flag)
  135. {
  136. mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
  137. int id=conn->ser.id;
  138. h->dm->mqttInfo.ser[id].status = flag;
  139. }
  140. static char *get_int_str(int n)
  141. {
  142. static char tmp[32];
  143. snprintf(tmp, sizeof(tmp), "%d", n);
  144. return tmp;
  145. }
  146. static char *get_float_str(float n)
  147. {
  148. static char tmp[32];
  149. snprintf(tmp, sizeof(tmp), "%0.3f", n);
  150. return tmp;
  151. }
  152. static int get_date_time(char *d, char *t)
  153. {
  154. time_t now;
  155. struct tm* tm;
  156. time(&now);
  157. tm = localtime(&now);
  158. if(!tm) {
  159. return -1;
  160. }
  161. sprintf(d, "%04d-%02d-%02d", tm->tm_year+1900, tm->tm_mon+1, tm->tm_mday);
  162. sprintf(t, "%02d:%02d:%02d", tm->tm_hour, tm->tm_min, tm->tm_sec);
  163. return 0;
  164. }
  165. static int sub_one(mg_conn_t *c, char *topic)
  166. {
  167. int r=-1;
  168. mg_opts_t opts={0};
  169. opts.topic = mg_str(topic);
  170. opts.qos = 1;
  171. mg_mqtt_sub(c, &opts);
  172. return 0;
  173. }
  174. static int my_sub(mg_conn_t *c, int prod_id)
  175. {
  176. int i;
  177. char topic[1024];
  178. for(i=0; i<MQTT_SUB_MAX; i++) {
  179. snprintf(topic, sizeof(topic), topic_sub[i], prod_id);
  180. sub_one(c, topic);
  181. }
  182. return 0;
  183. }
  184. static int pub_one(mg_conn_t *c, char *topic, char *data)
  185. {
  186. int r=-1;
  187. mg_opts_t opts={0};
  188. opts.topic = mg_str(topic);
  189. opts.message = mg_str(data);
  190. opts.qos = 1;
  191. opts.retain = false;
  192. mg_mqtt_pub(c, &opts);
  193. return 0;
  194. }
  195. static int my_pub(mqtt_conn_t *conn, char *topic, char *data)
  196. {
  197. if(conn->c) {
  198. pub_one(conn->c, topic, data);
  199. }
  200. return 0;
  201. }
  202. static GlobalPowerManger* get_power(mqtt_handle_t *h, int ch)
  203. {
  204. GlobalPowerManger *tmp=NULL;
  205. list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
  206. {
  207. if(tmp->product_ch_id==ch) {
  208. return tmp;
  209. }
  210. }
  211. return NULL;
  212. }
  213. static int set_ch(mqtt_handle_t *h, int ch, int flag)
  214. {
  215. int r;
  216. GlobalPowerManger *tmp=NULL;
  217. list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
  218. {
  219. if(tmp==NULL) {
  220. break;
  221. }
  222. if(tmp->product_ch_id==ch) {
  223. r = g_switch_set_all_chn_ctrl(&h->dm->_globalRelaySampManger, tmp, tmp->product_saddr, tmp->product_ch_addr, flag, false);
  224. if(r<0) {
  225. LOGE("___ g_switch_set_all_chn_ctrl failed, saddr: %d, ch_addr: %d\n", tmp->product_saddr, tmp->product_ch_addr);
  226. }
  227. break;
  228. }
  229. if(ch==0xff) {
  230. g_switch_set_all_ctrl(&h->dm->_globalRelaySampManger, tmp->product_ch_type, tmp->product_saddr, flag);
  231. }
  232. }
  233. return 0;
  234. }
  235. static GlobalSensorManger* get_sensor(mqtt_handle_t *h, int id)
  236. {
  237. GlobalSensorManger *tmp=NULL;
  238. lock_s_hold(LOCK_ID_SENSOR);
  239. list_for_each_entry(tmp, &h->dm->_globalSensorManger.list, list)
  240. {
  241. if(tmp->sensor_id==id) {
  242. return tmp;
  243. }
  244. }
  245. lock_s_release(LOCK_ID_SENSOR);
  246. return NULL;
  247. }
  248. //////////////////////////////////////////////////////////////
  249. static int get_url(mqtt_server_t *ser, char *url)
  250. {
  251. int port=1883;
  252. char *head="mqtts";
  253. if(!ser->server[0] || !ser->port[0] || !ser->mode) {
  254. return -1;
  255. }
  256. port = atoi(ser->port);
  257. if(port==1883) {
  258. head = "mqtt";
  259. }
  260. else if(port==8883) {
  261. head = "mqtts";
  262. }
  263. else if(port==8083) {
  264. head = "ws";
  265. }
  266. else if(port==8884) {
  267. head = "wss";
  268. }
  269. sprintf(url, "%s://%s:%d", head, ser->server, port);
  270. return 0;
  271. }
  272. static int start_one(mqtt_conn_t *conn)
  273. {
  274. int r;
  275. if(conn->tid==0) {
  276. conn->quit = 0;
  277. r = pthread_create(&conn->tid, NULL, client_thread, conn);
  278. }
  279. return r;
  280. }
  281. static int stop_one(mqtt_conn_t *conn)
  282. {
  283. if(conn->tid) {
  284. conn->quit = 1;
  285. pthread_join(conn->tid, NULL);
  286. conn->tid = 0;
  287. }
  288. return 0;
  289. }
  290. static int stop_all(mqtt_handle_t *h)
  291. {
  292. int i;
  293. for(i=0; i<MQTT_SER_MAX; i++) {
  294. stop_one(&h->conn[i]);
  295. }
  296. return 0;
  297. }
  298. static int mg_mqtt_bind(mg_conn_t *c)
  299. {
  300. int r=0;
  301. if(cellular_is_connected()) {
  302. r = sys_bind_nic(NIC_USB, (int)c->fd);
  303. }
  304. return r;
  305. }
  306. static int conn_one(mqtt_conn_t *conn)
  307. {
  308. int r=-1;
  309. mqtt_handle_t *h=&mqHandle;
  310. if(!conn->c) {
  311. char url[1024];
  312. if(conn->ser.plat==MQTT_PLAT_ALIYUN) {
  313. user_account_t user;
  314. mqtt_aliyun.clientId = conn->ser.cid;
  315. get_account(&mqtt_aliyun, &user);
  316. strcpy(conn->ser.user, user.user);
  317. strcpy(conn->ser.password, user.passwd);
  318. }
  319. mg_opts_t opts={
  320. .clean = true,
  321. .qos = 1,
  322. .version = 4,
  323. .keepalive = 60,
  324. .topic = mg_str("hello"),
  325. .message = mg_str("bye"),
  326. .client_id = mg_str(conn->ser.cid),
  327. .user = mg_str(conn->ser.user),
  328. .pass = mg_str(conn->ser.password),
  329. };
  330. get_url(&conn->ser, url);
  331. conn->c = mg_mqtt_connect(&conn->mgr, url, &opts, mqtt_fn, conn);
  332. if(conn->c) {
  333. //mg_mqtt_bind(conn->c);
  334. sprintf(conn->ser.cid, "smartPDU_%d\n", h->prod_id);
  335. my_sub(conn->c, h->prod_id);
  336. r = 0;
  337. }
  338. }
  339. return r;
  340. }
  341. static int my_check(mqtt_handle_t *h)
  342. {
  343. int i,j,r=-1;
  344. char url[1024];
  345. mqtt_server_t *ser=h->dm->mqttInfo.ser;
  346. for(i=0; i<MQTT_SER_MAX; i++) {
  347. r = get_url(&ser[i], url);
  348. if((r==0)) {
  349. if(h->conn[i].tid==0) {
  350. h->conn[i].ser = ser[i];
  351. h->conn[i].h = h;
  352. start_one(&h->conn[i]);
  353. }
  354. }
  355. else {
  356. if(h->conn[i].tid) {
  357. stop_one(&h->conn[i]);
  358. }
  359. }
  360. }
  361. return 0;
  362. }
  363. static int my_send(mqtt_conn_t *conn)
  364. {
  365. int r=-1;
  366. mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
  367. list_node_t *ln=NULL;
  368. r = xlist_take_node(h->list, &ln, 0);
  369. if(r==0) {
  370. mqtt_pkt_t *pkt=(mqtt_pkt_t*)ln->data.buf;
  371. my_pub(conn, pkt->topic, pkt->content);
  372. xlist_back_node(h->list, ln);
  373. }
  374. return r;
  375. }
  376. static int send_once(mqtt_conn_t *conn)
  377. {
  378. send_stat(conn, MQTT_PUB_INFO_DEVICE);
  379. send_stat(conn, MQTT_PUB_STAT_NETWORK);
  380. send_stat(conn, MQTT_PUB_STAT_SENSOR);
  381. send_stat(conn, MQTT_PUB_STAT_SERVICE);
  382. return 0;
  383. }
  384. static void send_period(mqtt_conn_t *conn)
  385. {
  386. send_stat(conn, MQTT_PUB_STAT_POWER_ALL);
  387. send_stat(conn, MQTT_PUB_STAT_POWER_CHN);
  388. }
  389. static void timer_conn_fn(void *arg)
  390. {
  391. mqtt_conn_t *conn=(mqtt_conn_t*)arg;
  392. conn_one(conn);
  393. }
  394. static void timer_period_fn(void *arg)
  395. {
  396. mqtt_conn_t *conn=(mqtt_conn_t*)arg;
  397. send_period(conn);
  398. }
  399. static void timer_add(mqtt_conn_t *conn)
  400. {
  401. mg_timer_add(&conn->mgr, CONN_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_conn_fn, conn);
  402. mg_timer_add(&conn->mgr, SEND_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period_fn, conn);
  403. }
  404. static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
  405. {
  406. mqtt_conn_t *conn=((mqtt_conn_t*)(c->fn_data));
  407. switch(ev) {
  408. case MG_EV_OPEN:
  409. {
  410. // c->is_hexdumping = 1;
  411. }
  412. break;
  413. case MG_EV_CONNECT:
  414. {
  415. char url[1024];
  416. get_url(&conn->ser, url);
  417. if (mg_url_is_ssl(url)) {
  418. struct mg_tls_opts opts = {.ca = mg_str(conn->ser.cert),
  419. .name = mg_url_host(conn->ser.server)};
  420. mg_tls_init(c, &opts);
  421. }
  422. }
  423. break;
  424. case MG_EV_MQTT_OPEN:
  425. {
  426. send_once(conn);
  427. set_status(conn, 1);
  428. }
  429. break;
  430. case MG_EV_POLL:
  431. {
  432. my_send(conn);
  433. }
  434. break;
  435. case MG_EV_MQTT_MSG:
  436. {
  437. struct mg_mqtt_message *mm=(struct mg_mqtt_message*)ev_data;
  438. if(mm && mm->topic.buf && mm->data.buf) {
  439. my_recv(mm->topic.buf, mm->data.buf);
  440. }
  441. }
  442. break;
  443. case MG_EV_ERROR:
  444. {
  445. //MG_ERROR(("___MG_EV_ERROR, %p %s", c->fd, (char *) ev_data));
  446. //memset(mc->flag, 0, sizeof(mc->flag));
  447. }
  448. break;
  449. case MG_EV_CLOSE:
  450. {
  451. conn->c = NULL;
  452. set_status(conn, 0);
  453. }
  454. break;
  455. }
  456. }
  457. static void* client_thread(void *arg)
  458. {
  459. int r;
  460. mqtt_conn_t *conn=(mqtt_conn_t*)arg;
  461. mg_mgr_init(&conn->mgr);
  462. conn->mgr.dns4.url = "udp://114.114.114.114:53";
  463. //conn->mgr.dns6.url = "udp://114.114.114.114:53";
  464. timer_add(conn);
  465. while(conn->quit==0) {
  466. mg_mgr_poll(&conn->mgr, 1000);
  467. }
  468. if(conn->c)mg_mqtt_disconnect(conn->c, NULL);
  469. mg_mgr_free(&conn->mgr);
  470. conn->c = NULL;
  471. pthread_exit(NULL);
  472. }
  473. #endif
  474. static void* mqtt_thread(void *arg)
  475. {
  476. #ifdef USE_MQTT
  477. int r;
  478. thread_handle_t *th=(thread_handle_t*)arg;
  479. mqtt_handle_t *h=(mqtt_handle_t*)th->arg;
  480. while(th->quit==0) {
  481. my_check(h);
  482. sleep(1);
  483. }
  484. stop_all(h);
  485. #endif
  486. pthread_exit(NULL);
  487. }
  488. //////////////////////////////////////////////////////////////////////////
  489. int mqtt_init(void)
  490. {
  491. int r=-1;
  492. #ifdef USE_MQTT
  493. mqtt_handle_t *h=&mqHandle;
  494. memset(h, 0, sizeof(mqtt_handle_t));
  495. //mg_log_set(MG_LL_DEBUG);
  496. h->dm = &__globalDeviceManage;
  497. h->prod_id = h->dm->_globalDevInfo.product.id;
  498. //sys_get_chip_id(&h->prod_id);
  499. list_cfg_t lc;
  500. lc.mode = LIST_FULL_FIFO;
  501. lc.max = 100;
  502. lc.log = 0;
  503. h->list = xlist_init(&lc);
  504. dev_mqtt_init(h->dm->db, &h->dm->mqttInfo);
  505. thread_start(THREAD_ID_MQTT, mqtt_thread, h);
  506. h->inited = 1; r = 0;
  507. #endif
  508. return r;
  509. }
  510. int mqtt_deinit(void)
  511. {
  512. int r=-1;
  513. #ifdef USE_MQTT
  514. mqtt_handle_t *h=&mqHandle;
  515. thread_stop(THREAD_ID_MQTT);
  516. xlist_free(h->list);
  517. h->inited = 0; r = 0;
  518. #endif
  519. return r;
  520. }
  521. static int send_stat(mqtt_conn_t *conn, int type)
  522. {
  523. int i,r=-1;
  524. #ifdef USE_MQTT
  525. float p1,p2,p3;
  526. char topic[256];
  527. char* content=NULL;
  528. char buf[20],date[40],time[40];
  529. mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
  530. if((!conn->c) || (type<0) || (type>=MQTT_PUB_MAX)) {
  531. return -1;
  532. }
  533. if((type!=MQTT_PUB_STAT_POWER_CHN) && (type!=MQTT_PUB_STAT_SENSOR)) {
  534. snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
  535. }
  536. switch(type) {
  537. case MQTT_PUB_INFO_DEVICE:
  538. {
  539. cJSON* root=cJSON_CreateObject();
  540. if(root) {
  541. cJSON_AddStringToObject(root,"id", get_int_str(h->prod_id));
  542. cJSON_AddStringToObject(root,"name", "Gowone smartPDU");
  543. cJSON_AddStringToObject(root,"type", "smartPDU AC");
  544. cJSON_AddStringToObject(root,"firmware", VERSION);
  545. cJSON_AddStringToObject(root,"channel_num", "8");
  546. cJSON_AddStringToObject(root,"sensor_num", "10");
  547. get_date_time(date, time);
  548. cJSON_AddStringToObject(root,"date", date);
  549. cJSON_AddStringToObject(root,"time", time);
  550. content = cJSON_Print(root);
  551. my_pub(conn, topic, content);
  552. cJSON_free(content);
  553. cJSON_Delete(root);
  554. }
  555. }
  556. break;
  557. case MQTT_PUB_STAT_NETWORK:
  558. {
  559. NetworkInfo_t nw4,nw6;
  560. int r4 = sys_get_net(&nw4, IP_V4);
  561. int r6 = sys_get_net(&nw6, IP_V6);
  562. cJSON* root=cJSON_CreateObject();
  563. if(root) {
  564. if(r4==0) {
  565. cJSON_AddStringToObject(root,"lan_ipv4_link", nw4.mode?"1":"0");
  566. cJSON_AddStringToObject(root,"lan_ipv4_address", nw4.ip_address);
  567. cJSON_AddStringToObject(root,"lan_ipv4_mask", nw4.mask);
  568. }
  569. else {
  570. cJSON_AddStringToObject(root,"lan_ipv4_link", "");
  571. cJSON_AddStringToObject(root,"lan_ipv4_address", "");
  572. cJSON_AddStringToObject(root,"lan_ipv4_mask", "");
  573. }
  574. if(r4==0) {
  575. cJSON_AddStringToObject(root,"lan_ipv6_link", nw6.mode?"1":"0");
  576. cJSON_AddStringToObject(root,"lan_ipv6_address", nw6.ip_address);
  577. cJSON_AddStringToObject(root,"lan_ipv6_subnet_length", nw6.mask);
  578. }
  579. else {
  580. cJSON_AddStringToObject(root,"lan_ipv6_link", "");
  581. cJSON_AddStringToObject(root,"lan_ipv6_address", "");
  582. cJSON_AddStringToObject(root,"lan_ipv6_subnet_length", "");
  583. }
  584. if(0) {
  585. cJSON_AddStringToObject(root,"wifi_ipv4_link", "");
  586. cJSON_AddStringToObject(root,"wifi_ipv4_address", "");
  587. cJSON_AddStringToObject(root,"wifi_ipv4_mask", "");
  588. cJSON_AddStringToObject(root,"wifi_ipv6_link", "");
  589. cJSON_AddStringToObject(root,"wifi_ipv6_address", "");
  590. cJSON_AddStringToObject(root,"wifi_ipv6_mask", "");
  591. }
  592. cJSON_AddStringToObject(root,"modbus_address", get_int_str(h->dm->_globalDevInfo.cascade.addr));
  593. cJSON_AddStringToObject(root,"modbus_baud", get_int_str(h->dm->_globalDevInfo.cascade.baudrate));
  594. cJSON_AddStringToObject(root,"modbus_mode", get_int_str(h->dm->_globalDevInfo.cascade.mode));
  595. cJSON_AddStringToObject(root,"vpn_enable", "0");
  596. get_date_time(date, time);
  597. cJSON_AddStringToObject(root,"date", date);
  598. cJSON_AddStringToObject(root,"time", time);
  599. content = cJSON_Print(root);
  600. my_pub(conn, topic, content);
  601. cJSON_free(content);
  602. cJSON_Delete(root);
  603. }
  604. }
  605. break;
  606. case MQTT_PUB_STAT_POWER_ALL:
  607. {
  608. int cnt=1;
  609. pwrall_info_t *pall=websocket_get_pwrall();
  610. if(pall->ph3) {
  611. cnt = 3;
  612. }
  613. for(i=0; i<cnt; i++) {
  614. _OverAllPwrAckInfo *info=&pall->ch[i];
  615. cJSON* root=cJSON_CreateObject();
  616. if(root) {
  617. sprintf(buf, "L%d", i+1);
  618. cJSON_AddStringToObject(root,"phase", buf);
  619. cJSON_AddStringToObject(root,"voltage", get_float_str(info->voltage));
  620. cJSON_AddStringToObject(root,"current", get_float_str(info->current));
  621. cJSON_AddStringToObject(root,"power", get_float_str(info->power/1000));
  622. cJSON_AddStringToObject(root,"consumption", get_float_str(info->consumption));
  623. p1 = info->power;
  624. p3 = (info->factor?(p1/info->factor):0.0f);
  625. p2 = p3 - p1;
  626. cJSON_AddStringToObject(root,"pactive_power", get_float_str(p1));
  627. cJSON_AddStringToObject(root,"reactive_power", get_float_str(p2));
  628. cJSON_AddStringToObject(root,"apparent_power", get_float_str(p3));
  629. get_date_time(date, time);
  630. cJSON_AddStringToObject(root,"date", date);
  631. cJSON_AddStringToObject(root,"time", time);
  632. content = cJSON_Print(root);
  633. my_pub(conn, topic, content);
  634. cJSON_free(content);
  635. cJSON_Delete(root);
  636. }
  637. }
  638. }
  639. break;
  640. case MQTT_PUB_STAT_POWER_CHN:
  641. {
  642. int quit=0;
  643. GlobalPowerManger *tmp=NULL;
  644. lock_s_hold(LOCK_ID_POWER_UPDATE);
  645. list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
  646. {
  647. if(tmp==NULL) {
  648. break;
  649. }
  650. cJSON* root=cJSON_CreateObject();
  651. if(root) {
  652. sprintf(buf, "CH%d", tmp->product_ch_id);
  653. cJSON_AddStringToObject(root,"name", buf);
  654. cJSON_AddStringToObject(root,"status", get_int_str(tmp->_PowerInfo.status));
  655. cJSON_AddStringToObject(root,"voltage", get_float_str(tmp->_PowerInfo.voltage));
  656. cJSON_AddStringToObject(root,"current", get_float_str(tmp->_PowerInfo.current));
  657. cJSON_AddStringToObject(root,"power", get_float_str(tmp->_PowerInfo.power));
  658. cJSON_AddStringToObject(root,"power_freq", get_float_str(tmp->_PowerInfo.freq));
  659. cJSON_AddStringToObject(root,"consumption", get_float_str(tmp->_PowerInfo.consumption));
  660. cJSON_AddStringToObject(root,"power_factor", get_float_str(tmp->_PowerInfo.factor));
  661. p1 = tmp->_PowerInfo.power;
  662. p3 = p1 / tmp->_PowerInfo.factor;
  663. p2 = p3 - p1;
  664. cJSON_AddStringToObject(root,"pactive_power", get_float_str(p1));
  665. cJSON_AddStringToObject(root,"reactive_power", get_float_str(p2));
  666. cJSON_AddStringToObject(root,"apparent_power", get_float_str(p3));
  667. get_date_time(date, time);
  668. cJSON_AddStringToObject(root,"date", date);
  669. cJSON_AddStringToObject(root,"time", time);
  670. snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->product_ch_id);
  671. content = cJSON_Print(root);
  672. my_pub(conn, topic, content);
  673. cJSON_free(content);
  674. cJSON_Delete(root);
  675. }
  676. }
  677. lock_s_release(LOCK_ID_POWER_UPDATE);
  678. }
  679. break;
  680. case MQTT_PUB_STAT_SENSOR:
  681. {
  682. GlobalSensorManger *tmp=NULL;
  683. lock_s_hold(LOCK_ID_SENSOR);
  684. list_for_each_entry(tmp, &h->dm->_globalSensorManger.list, list)
  685. {
  686. if(tmp==NULL) {
  687. break;
  688. }
  689. cJSON* root=cJSON_CreateObject();
  690. if(root) {
  691. cJSON_AddStringToObject(root,"name", "temp sensor");
  692. cJSON_AddStringToObject(root,"type", get_int_str(tmp->sensor_type));
  693. cJSON_AddStringToObject(root,"modbus_address", get_int_str(tmp->sensor_addr));
  694. cJSON_AddStringToObject(root,"node", get_int_str(tmp->sensor_node_number));
  695. cJSON_AddStringToObject(root,"status", get_int_str(tmp->sensor_status));
  696. cJSON_AddStringToObject(root,"value_num", get_int_str(tmp->sensor_val_count));
  697. cJSON_AddStringToObject(root,"value1", get_float_str(tmp->Cur_sensor_info.val1));
  698. cJSON_AddStringToObject(root,"value2", get_float_str(tmp->Cur_sensor_info.val2));
  699. get_date_time(date, time);
  700. cJSON_AddStringToObject(root,"date", date);
  701. cJSON_AddStringToObject(root,"time", time);
  702. snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->sensor_id);
  703. content = cJSON_Print(root);
  704. my_pub(conn, topic, content);
  705. cJSON_free(content);
  706. cJSON_Delete(root);
  707. }
  708. }
  709. lock_s_release(LOCK_ID_SENSOR);
  710. }
  711. break;
  712. case MQTT_PUB_STAT_SERVICE:
  713. {
  714. cJSON* root=cJSON_CreateObject();
  715. if(root) {
  716. int flag = get_serv(NULL);
  717. cJSON_AddStringToObject(root,"telnet", BGET(flag, SERV_TELNET)?"1":"0");
  718. cJSON_AddStringToObject(root,"smtp", BGET(flag, SERV_SMTP)?"1":"0");
  719. cJSON_AddStringToObject(root,"snmp_v1", BGET(flag, SERV_SNMP_V1)?"1":"0");
  720. cJSON_AddStringToObject(root,"snmp_v2c", BGET(flag, SERV_SNMP_V2C)?"1":"0");
  721. cJSON_AddStringToObject(root,"snmp_v3", BGET(flag, SERV_SNMP_V3)?"1":"0");
  722. cJSON_AddStringToObject(root,"snmp_trap", BGET(flag, SERV_SNMP_TRAP)?"1":"0");
  723. cJSON_AddStringToObject(root,"message", BGET(flag, SERV_MESG)?"1":"0");
  724. cJSON_AddStringToObject(root,"mqtt", BGET(flag, SERV_MQTT)?"1":"0");
  725. cJSON_AddStringToObject(root,"cloud", BGET(flag, SERV_CLOUD)?"1":"0");
  726. cJSON_AddStringToObject(root,"ntp", BGET(flag, SERV_NTP)?"1":"0");
  727. get_date_time(date, time);
  728. cJSON_AddStringToObject(root,"date", date);
  729. cJSON_AddStringToObject(root,"time", time);
  730. content = cJSON_Print(root);
  731. my_pub(conn, topic, content);
  732. cJSON_free(content);
  733. cJSON_Delete(root);
  734. }
  735. }
  736. break;
  737. default:
  738. return -1;
  739. }
  740. #endif
  741. return r;
  742. }
  743. static int post_alarm(int type, int subtype, char *content, alarm_para_t *para)
  744. {
  745. int r=-1;
  746. #ifdef USE_MQTT
  747. char topic[256];
  748. mqtt_handle_t *h=&mqHandle;
  749. char date[40],time[40];
  750. int pub_type;
  751. char *json_str=NULL;
  752. if(!h->inited) {
  753. return -1;
  754. }
  755. if(type==ALARM_TYPE_POWER) {
  756. pub_type = MQTT_PUB_ALARM_POWER;
  757. }
  758. else if(type==ALARM_TYPE_SENSOR) {
  759. pub_type = MQTT_PUB_ALARM_SENSOR;
  760. }
  761. else if(type==ALARM_TYPE_NETWORK) {
  762. pub_type = MQTT_PUB_ALARM_NETWORK;
  763. }
  764. else {
  765. return -1;
  766. }
  767. snprintf(topic, sizeof(topic), topic_pub[pub_type], h->prod_id);
  768. get_date_time(date, time);
  769. switch(pub_type) {
  770. case MQTT_PUB_ALARM_NETWORK:
  771. case MQTT_PUB_ALARM_SENSOR:
  772. case MQTT_PUB_ALARM_POWER:
  773. {
  774. cJSON* root=cJSON_CreateObject();
  775. if(root) {
  776. cJSON_AddStringToObject(root,"num", get_int_str(para->nalarm)); //alarm number, how to get?
  777. cJSON_AddStringToObject(root,"context", content);
  778. cJSON_AddStringToObject(root,"action", get_int_str(para->action));
  779. cJSON_AddStringToObject(root,"action_para", get_int_str(para->actionId));
  780. cJSON_AddStringToObject(root,"date", date);
  781. cJSON_AddStringToObject(root,"time", time);
  782. json_str = cJSON_Print(root);
  783. cJSON_Delete(root);
  784. }
  785. }
  786. break;
  787. }
  788. mqtt_pkt_t pkt;
  789. pkt.id = pub_type;
  790. snprintf(pkt.topic, sizeof(pkt.topic), "%s", topic);
  791. snprintf(pkt.content, sizeof(pkt.content), "%s", json_str);
  792. r = xlist_append(h->list, 0, &pkt, sizeof(pkt));
  793. cJSON_free(json_str);
  794. #endif
  795. return r;
  796. }
  797. int mqtt_post_alarm(int type, int subtype, char *content, alarm_para_t *para)
  798. {
  799. return post_alarm(type, subtype, content, para);
  800. }
  801. #ifdef USE_MQTT
  802. //////////////////////////////////////////////////////////
  803. static int get_cmd(char *json)
  804. {
  805. int cmd=-1;
  806. cJSON* cjson=cJSON_Parse(json);
  807. if(cjson) {
  808. cJSON* order=cJSON_GetObjectItem(cjson,"order");
  809. if(order && order->valuestring) {
  810. cmd = atoi(order->valuestring);
  811. }
  812. cJSON_Delete(cjson);
  813. }
  814. return cmd;
  815. }
  816. static int get_serv(char *json)
  817. {
  818. int i,flag=0;
  819. cJSON* tmp=NULL;
  820. if(json) {
  821. cJSON* cjson=cJSON_Parse(json);
  822. if(cjson) {
  823. for(i=0; i<SERV_MAX; i++) {
  824. tmp = cJSON_GetObjectItem(cjson, serv_str[i]);
  825. if(tmp && tmp->valuestring) {
  826. flag |= (atoi(tmp->valuestring)<<i);
  827. }
  828. }
  829. cJSON_Delete(cjson);
  830. }
  831. }
  832. else {
  833. if(thread_is_running(THREAD_ID_NTP)) {
  834. BSET(flag, SERV_NTP);
  835. }
  836. if(thread_is_running(THREAD_ID_MAIL)) {
  837. BSET(flag, SERV_SMTP);
  838. }
  839. if(thread_is_running(THREAD_ID_MQTT)) {
  840. BSET(flag, SERV_MQTT);
  841. }
  842. #if 0
  843. if(thread_is_running(SERV_MESG)) {
  844. BSET(flag, SERV_MESG);
  845. }
  846. if(thread_is_running(THREAD_ID_CLOUD)) {
  847. BSET(flag, SERV_CLOUD);
  848. }
  849. if(thread_is_running(THREAD_ID_TELNET)) {
  850. BSET(flag, SERV_TELNET);
  851. }
  852. #endif
  853. if(sys_is_running("snmpd")) {
  854. BSET(flag, SERV_NTP);
  855. }
  856. if(thread_is_running(THREAD_ID_NTP)) {
  857. BSET(flag, SERV_NTP);
  858. }
  859. }
  860. return flag;
  861. }
  862. static int set_serv(int flag)
  863. {
  864. if(BGET(flag,SERV_NTP)) {
  865. if(!thread_is_running(THREAD_ID_NTP)) {
  866. //thread_stop(THREAD_ID_NTP);
  867. }
  868. }
  869. else {
  870. if(thread_is_running(THREAD_ID_NTP)) {
  871. thread_stop(THREAD_ID_NTP);
  872. }
  873. }
  874. if(BGET(flag,SERV_SMTP)) {
  875. if(!thread_is_running(THREAD_ID_MAIL)) {
  876. //thread_stop(THREAD_ID_MAIL);
  877. }
  878. }
  879. else {
  880. if(thread_is_running(THREAD_ID_MAIL)) {
  881. thread_stop(THREAD_ID_MAIL);
  882. }
  883. }
  884. if(BGET(flag,SERV_MQTT)) {
  885. }
  886. else {
  887. }
  888. if(BGET(flag,SERV_MESG)) {
  889. }
  890. else {
  891. }
  892. if(BGET(flag,SERV_CLOUD)) {
  893. }
  894. else {
  895. }
  896. if(BGET(flag,SERV_TELNET)) {
  897. }
  898. else {
  899. }
  900. if(BGET(flag,SERV_SNMP_V1) || BGET(flag,SERV_SNMP_V2C) || BGET(flag,SERV_SNMP_V3) || BGET(flag,SERV_SNMP_TRAP)) {
  901. }
  902. else {
  903. }
  904. return 0;
  905. }
  906. static int my_recv(char *topic, char *data)
  907. {
  908. int i,r,cmd;
  909. #ifdef USE_MQTT
  910. char *p,temp[512];
  911. int flag=0;
  912. mqtt_handle_t *h=&mqHandle;
  913. for(i=0; i<MQTT_SUB_MAX; i++) {
  914. snprintf(temp, sizeof(temp), topic_sub[i], h->prod_id);
  915. p = strrchr(temp, '+');
  916. if(p) p[0] = 0;
  917. if(strstr(topic, temp)) {
  918. switch(i) {
  919. case MQTT_SUB_CMD_POWER_ALL:
  920. {
  921. cmd = get_cmd(data);
  922. LOGD("___ MQTT_SUB_CMD_POWER_ALL, %d\n", cmd);
  923. if(cmd==0 || cmd==1) {
  924. set_ch(h, 0xff, cmd);
  925. }
  926. }
  927. break;
  928. case MQTT_SUB_CMD_POWER_CHN:
  929. {
  930. int ch=atoi(topic+strlen(temp));
  931. cmd = get_cmd(data);
  932. LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d, %d\n", ch, cmd);
  933. if(cmd==0 || cmd==1) {
  934. set_ch(h, ch, cmd);
  935. }
  936. }
  937. break;
  938. case MQTT_SUB_CMD_SENSOR:
  939. {
  940. int id=atoi(topic+strlen(temp));
  941. cmd = get_cmd(data);
  942. switch(id) {
  943. case SENSOR_TYPE_TEMP_HUMI:
  944. break;
  945. case SENSOR_TYPE_SMOKE:
  946. break;
  947. case SENSOR_TYPE_WATER:
  948. break;
  949. case SENSOR_TYPE_ACCESS:
  950. break;
  951. case SENSOR_TYPE_GAS:
  952. break;
  953. case SENSOR_TYPE_AIR_PRESS:
  954. break;
  955. case SENSOR_TYPE_TEMP_DOUBLE:
  956. break;
  957. case SENSOR_TYPE_TEMP:
  958. break;
  959. case SENSOR_TYPE_LEAK:
  960. break;
  961. }
  962. LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d\n", id);
  963. }
  964. break;
  965. case MQTT_SUB_CMD_RESTART:
  966. {
  967. cmd = get_cmd(data);
  968. LOGD("___ MQTT_SUB_CMD_RESTART, %d\n", cmd);
  969. if(cmd==0) {
  970. system("shutdown");
  971. }
  972. else if(cmd==1) {
  973. system("reboot");
  974. }
  975. }
  976. break;
  977. case MQTT_SUB_CMD_RESET:
  978. {
  979. cmd = get_cmd(data);
  980. LOGD("___ MQTT_SUB_CMD_RESET, %d\n", cmd);
  981. if(cmd==1) {
  982. sys_set_factory();
  983. }
  984. }
  985. break;
  986. case MQTT_SUB_CMD_SERVICE:
  987. {
  988. flag = get_serv(data);
  989. LOGD("___ MQTT_SUB_CMD_SERVICE, 0x%08x\n", flag);
  990. set_serv(flag);
  991. }
  992. break;
  993. }
  994. }
  995. }
  996. #endif
  997. return 0;
  998. }
  999. #endif
  1000. int mqtt_test(void)
  1001. {
  1002. #ifdef USE_MQTT
  1003. int cnt=0;
  1004. mqtt_handle_t *h=&mqHandle;
  1005. mqtt_info_t mInfo={0};
  1006. mqtt_server_t ser={
  1007. .id = 0,
  1008. .mode = 1,
  1009. .cid = "mqtlx_e92334545tet3465",
  1010. .name = "serverX",
  1011. .server = "192.168.1.12",
  1012. .port = "1883",
  1013. //root
  1014. //root0219107X
  1015. .user = "gowone100",
  1016. .password = "gowone100",
  1017. //.user = "gowone101",
  1018. //.password = "gowone101",
  1019. };
  1020. mInfo.ser[0] = ser;
  1021. h->dm->mqttInfo = mInfo;
  1022. #endif
  1023. return 0;
  1024. }