mqtt.c 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604
  1. #include <regex.h>
  2. #include "common.h"
  3. #include "cJSON.h"
  4. #include "mqtt.h"
  5. #include "thread.h"
  6. #include "mongoose.h"
  7. #include "elog.h"
  8. #include "lock.h"
  9. #include "cfg.h"
  10. #if 0
  11. #define LOGD log_d
  12. #define LOGE log_e
  13. #define LOGW log_w
  14. #else
  15. #define LOGD printf
  16. #define LOGE printf
  17. #define LOGW printf
  18. #endif
  19. #ifdef USE_MQTT
  20. char *topic_sub[MQTT_SUB_MAX]={
  21. "/pdu/%d/control/power/all_channel",
  22. "/pdu/%d/control/power/%d",
  23. "/pdu/%d/control/sensor/%d",
  24. "/pdu/%d/control/device/restart",
  25. "/pdu/%d/control/device/reset",
  26. "/PDU/%d/control/service",
  27. };
  28. char *topic_pub[MQTT_PUB_MAX]={
  29. "/pdu/%d/info/device",
  30. "/PDU/%d/status/network",
  31. "/pdu/%d/status/power/all_channel",
  32. "/pdu/%d/status/power/%d",
  33. "/pdu/%d/status/sensor/%d",
  34. "/pdu/%d/alarm/network",
  35. "/pdu/%d/alarm/power",
  36. "/pdu/%d/alarm/sensor",
  37. "/PDU/%d/status/service",
  38. };
  39. typedef struct mg_mgr mgr_t;
  40. typedef struct mg_mqtt_opts mg_opts_t;
  41. typedef struct mg_connection mg_conn_t;
  42. typedef struct {
  43. uint8_t qos;
  44. uint8_t ver;
  45. uint8_t clean;
  46. uint8_t retain;
  47. }mqtt_para_t;
  48. typedef struct _mqtt_conn_t{
  49. mg_conn_t *c;
  50. mg_opts_t opts;
  51. mqtt_para_t para;
  52. mqtt_info_t *info;
  53. bool isover;
  54. struct _mqtt_conn_t *next;
  55. }mqtt_conn_t;
  56. typedef struct {
  57. mgr_t mgr;
  58. mqtt_conn_t *conn;
  59. int inited;
  60. int prod_id;
  61. GlobalDeviceManager *dm;
  62. }mqtt_handle_t;
  63. static mqtt_handle_t mqHandle={0};
  64. static int my_conn(mqtt_handle_t *h, mqtt_info_t *info);
  65. static int my_recv(char *topic, char *data);
  66. static int info_cmp(mqtt_info_t *a, mqtt_info_t *b)
  67. {
  68. return memcmp(a, b, sizeof(mqtt_info_t)-sizeof(mqtt_info_t*));
  69. }
  70. static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
  71. {
  72. mqtt_handle_t *h=&mqHandle;
  73. if (ev == MG_EV_OPEN) {
  74. // c->is_hexdumping = 1;
  75. } else if (ev == MG_EV_CONNECT) {
  76. if (mg_url_is_ssl(h->para.user.url)) {
  77. struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
  78. .name = mg_url_host(h->para.user.url)};
  79. mg_tls_init(c, &opts);
  80. }
  81. } else if (ev == MG_EV_ERROR) {
  82. // On error, log error message
  83. MG_ERROR(("%p %s", c->fd, (char *) ev_data));
  84. } else if (ev == MG_EV_MQTT_OPEN) {
  85. mg_opts_t opts={
  86. .client_id = mg_str(h->para.user.url),
  87. .user = mg_str(h->para.user.name),
  88. .pass = mg_str(h->para.user.pass),
  89. };
  90. size_t len=c->send.len;
  91. mg_mqtt_login(c, &opts);
  92. mg_ws_wrap(c, c->send.len - len, WEBSOCKET_OP_BINARY);
  93. } else if (ev == MG_EV_MQTT_MSG) {
  94. // When we receive MQTT message, print it
  95. struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
  96. //MG_INFO(("Received on %.*s : %.*s", (int) mm->topic.len, mm->topic.buf, (int) mm->data.len, mm->data.buf));
  97. my_recv(mm->topic.buf, mm->data.buf);
  98. }
  99. if (ev == MG_EV_ERROR || ev == MG_EV_CLOSE) {
  100. MG_INFO(("got event %d, stopping...", ev));
  101. *(bool *) c->fn_data = true; // Signal that we're done
  102. }
  103. }
  104. static void* mqtt_thread(void *arg)
  105. {
  106. int r;
  107. thread_handle_t *th=(thread_handle_t*)arg;
  108. mqtt_handle_t *h=(mqtt_handle_t*)th->arg;
  109. while(th->quit==0) {
  110. if(h->inited) {
  111. my_conn_check(h);
  112. lock_s_hold(LOCK_ID_MQTT);
  113. mg_mgr_poll(&h->mgr, 300);
  114. lock_s_release(LOCK_ID_MQTT);
  115. }
  116. }
  117. pthread_exit(NULL);
  118. }
  119. static int my_sub(mg_conn_t *c, int prod_id, int qos)
  120. {
  121. int i;
  122. mg_opts_t opts;
  123. char temp[512];
  124. memset(&opts, 0, sizeof(opts));
  125. opts.qos = qos;
  126. for(i=0; topic_sub[i]; i++) {
  127. sprintf(temp, topic_sub[i], prod_id);
  128. opts.topic = mg_str(temp);
  129. mg_mqtt_sub(c, &opts);
  130. }
  131. return 0;
  132. }
  133. static int find_conn(mqtt_handle_t *h, mqtt_info_t *info)
  134. {
  135. //
  136. return 0;
  137. }
  138. static int my_conn(mqtt_handle_t *h, mqtt_info_t *info)
  139. {
  140. int r=-1;
  141. mqtt_conn_t *conn,*c;
  142. h->opts.clean = true,
  143. h->opts.qos = h->para.conn.qos,
  144. h->opts.topic = mg_str(""),
  145. h->opts.version = h->para.conn.ver,
  146. h->opts.message = mg_str("bye");
  147. conn = h->conn;
  148. while(info) {
  149. if(!conn==NULL) {
  150. conn = calloc(1, sizeof(mqtt_conn_t));
  151. if(!conn) {
  152. LOGE("___my_conn calloc failed\n");
  153. break;
  154. }
  155. }
  156. conn->c = mg_mqtt_connect(&h->mgr, conn->para.user.url, &conn->opts, mqtt_fn, &conn->isover);
  157. if(conn>c) {
  158. my_sub(conn>c, h->prod_id, 1);
  159. }
  160. info = info->next;
  161. if(!info) {
  162. r = 0; break;
  163. }
  164. }
  165. return r;
  166. }
  167. static int my_disconn(mqtt_handle_t *h, mqtt_conn_t *c)
  168. {
  169. mqtt_conn_t *c1,*c2;
  170. c1 = c2 = h->conn;
  171. while(c1) {
  172. if(c1==c) {
  173. mg_mqtt_disconnect(c1, &c1->opts);
  174. break;
  175. }
  176. c1 = c1->next;
  177. }
  178. h->conn = NULL;
  179. }
  180. static int my_disconn_all(mqtt_handle_t *h)
  181. {
  182. mqtt_conn_t *c,*conn=h->conn;
  183. while(conn) {
  184. c = conn;
  185. mg_mqtt_disconnect(c, &c->opts);
  186. conn = conn->next;
  187. free(c);
  188. }
  189. h->conn = NULL;
  190. }
  191. static int my_conn_check(mqtt_handle_t *h)
  192. {
  193. int r=-1;
  194. mqtt_conn_t *conn=h->conn;
  195. while(conn) {
  196. if(conn->c && conn->isover) {
  197. conn->c = mg_mqtt_connect(&h->mgr, conn->para.user.url, &conn->opts, mqtt_fn, &conn->isover);
  198. if(conn>c) {
  199. my_sub(conn>c, h->prod_id, 1); r = 0;
  200. }
  201. }
  202. conn = conn->next;
  203. }
  204. return r;
  205. }
  206. #endif
  207. //////////////////////////////////////////////////////////////////////////
  208. int mqtt_init(void)
  209. {
  210. int r=-1;
  211. #ifdef USE_MQTT
  212. mqtt_handle_t *h=&mqHandle;
  213. memset(h, 0, sizeof(mqtt_handle_t));
  214. mg_mgr_init(&h->mgr);
  215. h->dm = &__globalDeviceManage;
  216. h->prod_id = h->dm->_globalDevInfo.product_id;
  217. //h->para.user = ;
  218. h->para.conn.ver = 4;
  219. h->para.conn.qos = 1;
  220. thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
  221. h->inited = 1;
  222. r = 0;
  223. #endif
  224. return r;
  225. }
  226. int mqtt_deinit(void)
  227. {
  228. int r=-1;
  229. #ifdef USE_MQTT
  230. mqtt_handle_t *h=&mqHandle;
  231. thread_stop(THREAD_ID_MQTT);
  232. mg_mgr_free(&h->mgr);
  233. h->inited = 0;
  234. r = 0;
  235. #endif
  236. return r;
  237. }
  238. int mqtt_conn(void)
  239. {
  240. int r=-1;
  241. #ifdef USE_MQTT
  242. mqtt_handle_t *h=&mqHandle;
  243. lock_s_hold(LOCK_ID_MQTT);
  244. if(!h->inited) {
  245. goto quit;
  246. }
  247. if(h->conn) {
  248. mg_mqtt_disconnect(h->conn, NULL);
  249. }
  250. my_conn(h);
  251. quit:
  252. lock_s_release(LOCK_ID_MQTT);
  253. r = h->conn?0:-1;
  254. #endif
  255. return r;
  256. }
  257. int mqtt_disconn(void)
  258. {
  259. int r=-1;
  260. #ifdef USE_MQTT
  261. mqtt_handle_t *h=&mqHandle;
  262. lock_s_hold(LOCK_ID_MQTT);
  263. if(!h->inited) {
  264. r = -1;
  265. goto quit;
  266. }
  267. mg_mqtt_disconnect(h->conn, NULL);
  268. r = 0;
  269. quit:
  270. lock_s_release(LOCK_ID_MQTT);
  271. #endif
  272. return r;
  273. }
  274. int mqtt_sub(char *topic, int qos)
  275. {
  276. int r=-1;
  277. #ifdef USE_MQTT
  278. mg_opts_t opts;
  279. mqtt_handle_t *h=&mqHandle;
  280. lock_s_hold(LOCK_ID_MQTT);
  281. if(!h->inited || !h->conn) {
  282. r = -1;
  283. goto quit;
  284. }
  285. memset(&opts, 0, sizeof(opts));
  286. opts.topic = mg_str(topic);
  287. opts.qos = qos;
  288. mg_mqtt_sub(h->conn, &opts);
  289. r = 0;
  290. quit:
  291. lock_s_release(LOCK_ID_MQTT);
  292. #endif
  293. return r;
  294. }
  295. int mqtt_pub(char *topic, int qos, char *data, int dlen)
  296. {
  297. int r=-1;
  298. #ifdef USE_MQTT
  299. mg_opts_t opts;
  300. mqtt_handle_t *h=&mqHandle;
  301. if(!h->inited) {
  302. return -1;
  303. }
  304. memset(&opts, 0, sizeof(opts));
  305. opts.topic = mg_str(topic);
  306. opts.qos = qos;
  307. opts.message.buf = data;
  308. opts.message.len = dlen;
  309. opts.retain = false;
  310. while(1) {
  311. lock_s_hold(LOCK_ID_MQTT);
  312. mg_mqtt_pub(h->conn, &opts);
  313. lock_s_release(LOCK_ID_MQTT);
  314. }
  315. #endif
  316. return r;
  317. }
  318. int mqtt_send(int type, int id, void *data)
  319. {
  320. int r=-1;
  321. #ifdef USE_MQTT
  322. char topic[512];
  323. char content[4096];
  324. mqtt_handle_t *h=&mqHandle;
  325. if(type<0 || type>=MQTT_PUB_MAX || !data) {
  326. return -1;
  327. }
  328. if(type==MQTT_PUB_STAT_POWER_CHN || type==MQTT_PUB_STAT_SENSOR) {
  329. snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, id);
  330. }
  331. else {
  332. snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
  333. }
  334. switch(type) {
  335. case MQTT_PUB_INFO_DEVICE:
  336. {
  337. //
  338. }
  339. break;
  340. case MQTT_PUB_STAT_NETWORK:
  341. {
  342. }
  343. break;
  344. case MQTT_PUB_STAT_POWER_ALL:
  345. {
  346. }
  347. break;
  348. case MQTT_PUB_STAT_POWER_CHN:
  349. {
  350. }
  351. break;
  352. case MQTT_PUB_STAT_SENSOR:
  353. {
  354. }
  355. break;
  356. case MQTT_PUB_ALARM_NETWORK:
  357. {
  358. }
  359. break;
  360. case MQTT_PUB_ALARM_POWER:
  361. {
  362. }
  363. break;
  364. case MQTT_PUB_ALARM_SENSOR:
  365. {
  366. }
  367. break;
  368. }
  369. //r = mqtt_pub(topic, );
  370. #endif
  371. return r;
  372. }
  373. //////////////////////////////////////////////////////////
  374. static int get_cmd(char *json)
  375. {
  376. int cmd=-1;
  377. cJSON* cjson=cJSON_Parse(json);
  378. if(cjson) {
  379. cJSON* order=cJSON_GetObjectItem(cjson,"order");
  380. if(order && order->valuestring) {
  381. cmd = atoi(order->valuestring);
  382. }
  383. cJSON_Delete(cjson);
  384. }
  385. return cmd;
  386. }
  387. enum {
  388. SERV_NTP=0,
  389. SERV_SMTP,
  390. SERV_MESG,
  391. SERV_MQTT,
  392. SERV_CLOUD,
  393. SERV_TELNET,
  394. SERV_SNMP_V1,
  395. SERV_SNMP_V2C,
  396. SERV_SNMP_V3,
  397. SERV_SNMP_TRAP,
  398. SERV_MAX
  399. };
  400. const char *serv_str[SERV_MAX]={
  401. "ntp",
  402. "smtp",
  403. "message",
  404. "mqtt",
  405. "cloud",
  406. "telnet",
  407. "snmp_v1",
  408. "snmp_v2c",
  409. "snmp_v3",
  410. "snmp_trap",
  411. };
  412. static uint32_t get_serv(char *json)
  413. {
  414. int i;
  415. cJSON* tmp=NULL;
  416. uint32_t flag=0;
  417. cJSON* cjson=cJSON_Parse(json);
  418. if(cjson) {
  419. for(i=0; i<SERV_MAX; i++) {
  420. tmp = cJSON_GetObjectItem(cjson, serv_str[i]);
  421. if(tmp && tmp->valuestring) {
  422. flag |= (atoi(tmp->valuestring)<<i);
  423. }
  424. }
  425. cJSON_Delete(cjson);
  426. }
  427. return flag;
  428. }
  429. static int my_recv(char *topic, char *data)
  430. {
  431. int i,r,cmd;
  432. #ifdef USE_MQTT
  433. char temp[2000];
  434. mqtt_handle_t *h=&mqHandle;
  435. for(i=0; i<MQTT_SUB_MAX; i++) {
  436. snprintf(temp, sizeof(temp), topic_sub[i], h->prod_id);
  437. if(strstr(topic, temp)) {
  438. switch(i) {
  439. case MQTT_SUB_CMD_POWER_ALL:
  440. {
  441. LOGD("___ MQTT_SUB_CMD_POWER_ALL\n");
  442. cmd = get_cmd(data);
  443. }
  444. break;
  445. case MQTT_SUB_CMD_POWER_CHN:
  446. {
  447. int ch=atoi(topic+strlen(temp));
  448. cmd = get_cmd(data);
  449. LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d\n", ch);
  450. }
  451. break;
  452. case MQTT_SUB_CMD_SENSOR:
  453. {
  454. int id=atoi(topic+strlen(temp));
  455. cmd = get_cmd(data);
  456. LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d\n", id);
  457. }
  458. break;
  459. case MQTT_SUB_CMD_RESTART:
  460. {
  461. LOGD("___ MQTT_SUB_CMD_RESTART\n");
  462. cmd = get_cmd(data);
  463. }
  464. break;
  465. case MQTT_SUB_CMD_RESET:
  466. {
  467. LOGD("___ MQTT_SUB_CMD_RESET\n");
  468. cmd = get_cmd(data);
  469. }
  470. break;
  471. case MQTT_SUB_CMD_SERVICE:
  472. {
  473. LOGD("___ MQTT_SUB_CMD_SERVICE\n");
  474. uint32_t flag=get_serv(data);
  475. }
  476. break;
  477. }
  478. }
  479. }
  480. #endif
  481. return 0;
  482. }