mqtt.c 32 KB

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