bus-unix.c 33 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361
  1. #include "precompile.h"
  2. #include "bus.h"
  3. #include "sockutil.h"
  4. #include "url.h"
  5. #include "spinlock.h"
  6. #include "list.h"
  7. #include "bus_internal.h"
  8. #include "evtpoll.h"
  9. #include <winpr/file.h>
  10. #include <winpr/pipe.h>
  11. #include <winpr/synch.h>
  12. #include <winpr/string.h>
  13. #include <sys/eventfd.h>
  14. #define TAG TOOLKIT_TAG("bus")
  15. #define BUS_RESULT_DATA 1 // ==BUS_TYPE_PACKET, callback: callback.on_pkt, no use
  16. #define BUS_RESULT_INFO 2 // ==BUS_TYPE_INFO, callback: callback.on_inf
  17. #define BUS_RESULT_EVENT 3 // ==BUS_TYPE_EVENT, callback: callback.on_evt , no use
  18. #define BUS_RESULT_SYSTEM 4 // ==BUS_TYPE_SYSTEM, callback: callback.on_sys
  19. #define BUS_RESULT_MSG 5 // send package msg, callback: callback.on_msg
  20. #define BUS_RESULT_UNKNOWN 6
  21. typedef struct msg_t {
  22. struct list_head entry;
  23. int type;
  24. int nparam;
  25. param_size_t* params;
  26. HANDLE evt;
  27. int evt_result;
  28. }msg_t;
  29. struct bus_endpt_t {
  30. int type;
  31. int epid;
  32. union {
  33. HANDLE pipe_handle;
  34. SOCKET sock_handle;
  35. };
  36. char* url;
  37. /*define here or iom_t area, this is definitely a problem for now.*/
  38. evtpoll_t* ep;
  39. int msg_fd;
  40. event_epoll_data_t* msg_sem;
  41. bus_endpt_callback callback;
  42. struct list_head msg_list;
  43. spinlock_t msg_lock;
  44. HANDLE tx_evt; //manually
  45. HANDLE rx_evt; //manually
  46. OVERLAPPED rx_overlapped;
  47. int rx_pending;
  48. int rx_pending_pkt_len;
  49. int rx_pending_pkt_uc_len;
  50. iobuffer_queue_t* rx_buf_queue;
  51. volatile int quit_flag;
  52. };
  53. static void free_msg(msg_t* msg)
  54. {
  55. free(msg->params);
  56. free(msg);
  57. }
  58. static __inline int bus_endpoint__data_is_handle(const bus_endpt_t* endpt, void* data)
  59. {
  60. if ((endpt->type == TYPE_TCP && data == &endpt->sock_handle) ||
  61. (endpt->type == TYPE_PIPE && data == &endpt->pipe_handle)) {
  62. return 1;
  63. }
  64. return 0;
  65. }
  66. static int to_result(int pkt_type)
  67. {
  68. switch (pkt_type) {
  69. case BUS_TYPE_EVENT:
  70. return BUS_RESULT_EVENT;
  71. case BUS_TYPE_SYSTEM:
  72. return BUS_RESULT_SYSTEM;
  73. case BUS_TYPE_PACKET:
  74. return BUS_RESULT_DATA;
  75. case BUS_TYPE_INFO:
  76. return BUS_RESULT_INFO;
  77. default:
  78. break;
  79. }
  80. return BUS_RESULT_UNKNOWN;
  81. }
  82. static HANDLE create_pipe_handle(const char* name)
  83. {
  84. char tmp[MAX_PATH];
  85. HANDLE pipe = INVALID_HANDLE_VALUE;
  86. sprintf(tmp, "\\\\.\\pipe\\%s", name);
  87. for (;;) {
  88. pipe = CreateFileA(tmp,
  89. GENERIC_READ | GENERIC_WRITE,
  90. 0,
  91. NULL,
  92. OPEN_EXISTING,
  93. FILE_FLAG_OVERLAPPED,
  94. NULL);
  95. if (pipe == INVALID_HANDLE_VALUE) {
  96. if (GetLastError() == ERROR_PIPE_BUSY) {
  97. if (WaitNamedPipeA(name, 20000))
  98. continue;
  99. }
  100. }
  101. break;
  102. }
  103. return pipe;
  104. }
  105. static SOCKET create_socket_handle(const char* ip, int port)
  106. {
  107. SOCKET fd;
  108. struct sockaddr_in addr;
  109. /*Warning: only front three param is effective*/
  110. fd = WSASocketA(AF_INET, SOCK_STREAM, IPPROTO_TCP, NULL, 0, WSA_FLAG_OVERLAPPED);
  111. if (fd == INVALID_SOCKET) {
  112. return fd;
  113. }
  114. {
  115. BOOL f = TRUE;
  116. u_long l = TRUE;
  117. setsockopt(fd, SOL_SOCKET, SO_DONTLINGER, (char*)&f, sizeof(f));
  118. /*non block sock*/
  119. _ioctlsocket(fd, FIONBIO, &l);
  120. }
  121. memset(&addr, 0, sizeof(addr));
  122. addr.sin_family = AF_INET;
  123. addr.sin_port = htons(port);
  124. addr.sin_addr.s_addr = inet_addr(ip);
  125. WLog_INFO(TAG, "fd(%d) start to connect...", fd);
  126. if (_connect(fd, (struct sockaddr*) & addr, sizeof(addr)) != 0) {
  127. if (errno == EINPROGRESS) {
  128. WLog_WARN(TAG, "in connect progress...");
  129. fd_set wr_set, ex_set;
  130. FD_ZERO(&wr_set);
  131. FD_ZERO(&ex_set);
  132. FD_SET(fd, &wr_set);
  133. FD_SET(fd, &ex_set);
  134. if (_select(fd + 1, NULL, &wr_set, &ex_set, NULL) > 0 && FD_ISSET(fd, &wr_set)) {
  135. return fd;
  136. }
  137. } else {
  138. WLog_ERR(TAG, "_connect failed : %d", errno);
  139. }
  140. closesocket(fd);
  141. return INVALID_SOCKET;
  142. }
  143. WLog_WARN(TAG, "this should be appear hardly!!!");
  144. return fd;
  145. }
  146. static int tcp_send_buf(bus_endpt_t* endpt, const char* buf, int n)
  147. {
  148. DWORD left = n;
  149. DWORD offset = 0;
  150. WLog_DBG(TAG, "==> fd(%d): tcp send buf len: %d", endpt->sock_handle, n);
  151. while (left > 0) {
  152. BOOL ret;
  153. WSABUF wsabuf;
  154. DWORD dwBytesTransfer;
  155. OVERLAPPED overlapped;
  156. memset(&overlapped, 0, sizeof(overlapped));
  157. overlapped.hEvent = endpt->tx_evt;
  158. ResetEvent(endpt->tx_evt);
  159. wsabuf.buf = (char*)buf + offset;
  160. wsabuf.len = left;
  161. ret = FALSE;
  162. n = _send(endpt->sock_handle, wsabuf.buf, wsabuf.len, 0);
  163. if (n == -1) {
  164. if (errno == EAGAIN || errno == EWOULDBLOCK) {
  165. fd_set wfds;
  166. struct timeval tv;
  167. int retval;
  168. FD_ZERO(&wfds);
  169. FD_SET(endpt->sock_handle, &wfds);
  170. tv.tv_sec = 5;
  171. tv.tv_usec = 0;
  172. retval = _select(endpt->sock_handle + 1, NULL, &wfds, NULL, &tv);
  173. if (retval == -1) {
  174. WLog_ERR(TAG, "select error, errno: %s", strerror(errno));
  175. }
  176. else if (retval && FD_ISSET(endpt->sock_handle, &wfds) > 0) {
  177. WLog_INFO(TAG, "can write");
  178. n = _send(endpt->sock_handle, wsabuf.buf, wsabuf.len, 0);
  179. if (n >= 0) {
  180. ret = TRUE;
  181. dwBytesTransfer = n;
  182. }
  183. } else {
  184. WLog_WARN(TAG, "write timeout");
  185. ret = TRUE;
  186. dwBytesTransfer = 0;
  187. }
  188. }
  189. else {
  190. WLog_ERR(TAG, "_send failed: %s", strerror(errno));
  191. }
  192. }
  193. else {
  194. ret = TRUE;
  195. dwBytesTransfer = n;
  196. }
  197. if (ret && dwBytesTransfer) {
  198. offset += dwBytesTransfer;
  199. left -= dwBytesTransfer;
  200. }
  201. else {
  202. WLog_DBG(TAG, "<== error fd(%d): tcp send buf left: %d", endpt->sock_handle, left);
  203. return -1;
  204. }
  205. }
  206. WLog_DBG(TAG, "<== fd(%d): tcp send buf len: %d", endpt->sock_handle, n);
  207. return 0;
  208. }
  209. static int pipe_send_buf(bus_endpt_t* endpt, const char* buf, int n)
  210. {
  211. DWORD left = n;
  212. DWORD offset = 0;
  213. while (left > 0) {
  214. BOOL ret;
  215. DWORD dwBytesTransfer;
  216. OVERLAPPED overlapped;
  217. memset(&overlapped, 0, sizeof(overlapped));
  218. overlapped.hEvent = endpt->tx_evt;
  219. ResetEvent(endpt->tx_evt);
  220. ret = WriteFile(endpt->pipe_handle, buf + offset, left, &dwBytesTransfer, &overlapped);
  221. if (!ret && GetLastError() == ERROR_IO_PENDING) {
  222. ret = GetOverlappedResult(endpt->pipe_handle, &overlapped, &dwBytesTransfer, TRUE);
  223. }
  224. if (ret && dwBytesTransfer) {
  225. offset += dwBytesTransfer;
  226. left -= dwBytesTransfer;
  227. }
  228. else {
  229. return -1;
  230. }
  231. }
  232. return 0;
  233. }
  234. static int send_buf(bus_endpt_t* endpt, const char* buf, int n)
  235. {
  236. if (endpt->type == TYPE_PIPE) {
  237. return pipe_send_buf(endpt, buf, n);
  238. }
  239. else if (endpt->type == TYPE_TCP) {
  240. return tcp_send_buf(endpt, buf, n);
  241. }
  242. else {
  243. return -1;
  244. }
  245. }
  246. static int send_pkt_raw(bus_endpt_t* endpt, iobuffer_t* pkt)
  247. {
  248. int pkt_len = iobuffer_get_length(pkt);
  249. int rc;
  250. iobuffer_write_head(pkt, IOBUF_T_I4, &pkt_len, 0);
  251. rc = send_buf(endpt, iobuffer_data(pkt, 0), iobuffer_get_length(pkt));
  252. return rc;
  253. }
  254. static int pipe_recv_buf(bus_endpt_t* endpt, char* buf, DWORD n)
  255. {
  256. DWORD left = n;
  257. DWORD offset = 0;
  258. while (left > 0) {
  259. BOOL ret;
  260. DWORD dwBytesTransfer;
  261. OVERLAPPED overlapped;
  262. memset(&overlapped, 0, sizeof(overlapped));
  263. overlapped.hEvent = endpt->rx_evt;
  264. ResetEvent(overlapped.hEvent);
  265. ret = ReadFile(endpt->pipe_handle, buf + offset, left, &dwBytesTransfer, &overlapped);
  266. if (!ret && GetLastError() == ERROR_IO_PENDING) {
  267. ret = GetOverlappedResult(endpt->pipe_handle, &overlapped, &dwBytesTransfer, TRUE);
  268. }
  269. if (ret && dwBytesTransfer) {
  270. offset += dwBytesTransfer;
  271. left -= dwBytesTransfer;
  272. }
  273. else {
  274. return -1;
  275. }
  276. }
  277. return 0;
  278. }
  279. static int tcp_recv_buf(bus_endpt_t* endpt, char* buf, DWORD n)
  280. {
  281. DWORD left = n;
  282. DWORD offset = 0;
  283. WLog_DBG(TAG, "==> fd(%d): tcp recv buf len: %d", endpt->sock_handle, n);
  284. while (left > 0) {
  285. BOOL ret;
  286. DWORD dwBytesTransfer;
  287. ret = FALSE;
  288. n = _recv(endpt->sock_handle, buf + offset, left, 0);
  289. if (n == -1) {
  290. if (errno == EAGAIN || errno == EWOULDBLOCK) {
  291. fd_set rfds;
  292. struct timeval tv;
  293. int retval;
  294. FD_ZERO(&rfds);
  295. FD_SET(endpt->sock_handle, &rfds);
  296. tv.tv_sec = 5;
  297. tv.tv_usec = 0;
  298. retval = _select(endpt->sock_handle + 1, &rfds, NULL, NULL, &tv);
  299. if (retval == -1) {
  300. WLog_ERR(TAG, "select failed: %s", strerror(errno));
  301. }
  302. else if (retval && FD_ISSET(endpt->sock_handle, &rfds) > 0) {
  303. WLog_INFO(TAG, "can read");
  304. n = _recv(endpt->sock_handle, buf + offset, left, 0);
  305. if (n >= 0) {
  306. ret = TRUE;
  307. dwBytesTransfer = n;
  308. }
  309. }
  310. else {
  311. WLog_WARN(TAG, "read timeout: %d, left: %d", retval, left);
  312. ret = TRUE;
  313. dwBytesTransfer = 0;
  314. }
  315. }
  316. }
  317. else {
  318. ret = TRUE;
  319. dwBytesTransfer = n;
  320. }
  321. if (ret && dwBytesTransfer) {
  322. offset += dwBytesTransfer;
  323. left -= dwBytesTransfer;
  324. }
  325. else if (!ret) {
  326. return -1;
  327. }
  328. }
  329. WLog_DBG(TAG, "<== fd(%d): tcp recv buf len: %d", endpt->sock_handle, n);
  330. return 0;
  331. }
  332. /*read data block, n is the dream length of data to read which will store in buf*/
  333. static int recv_buf(bus_endpt_t* endpt, char* buf, DWORD n)
  334. {
  335. if (endpt->type == TYPE_PIPE) {
  336. return pipe_recv_buf(endpt, buf, n);
  337. }
  338. else if (endpt->type == TYPE_TCP) {
  339. return tcp_recv_buf(endpt, buf, n);
  340. }
  341. else {
  342. return -1;
  343. }
  344. }
  345. static int recv_pkt_raw(bus_endpt_t* endpt, iobuffer_t** pkt)
  346. {
  347. int pkt_len;
  348. int rc = -1;
  349. rc = recv_buf(endpt, (char*)&pkt_len, 4);
  350. if (rc != 0)
  351. return rc;
  352. *pkt = iobuffer_create(-1, pkt_len);
  353. iobuffer_push_count(*pkt, pkt_len);
  354. if (pkt_len > 0) {
  355. rc = recv_buf(endpt, iobuffer_data(*pkt, 0), pkt_len);
  356. }
  357. if (rc < 0) {
  358. iobuffer_destroy(*pkt);
  359. *pkt = NULL;
  360. }
  361. return rc;
  362. }
  363. static int start_read_pkt(bus_endpt_t* endpt, iobuffer_t** p_pkt)
  364. {
  365. DWORD dwBytesTransferred;
  366. BOOL ret;
  367. int rc = 0;
  368. iobuffer_t* pkt = NULL;
  369. *p_pkt = NULL;
  370. WLog_DBG(TAG, "==>endpt(%d): start_read_pkt", endpt->epid);
  371. ResetEvent(endpt->rx_evt);
  372. WLog_DBG(TAG, "ResetEvent(endpt->rx_evt) ???");
  373. endpt->rx_pending_pkt_uc_len = 0;
  374. if (endpt->type == TYPE_PIPE) {
  375. ret = ReadFile(endpt->pipe_handle,
  376. &endpt->rx_pending_pkt_len,
  377. 4,
  378. &dwBytesTransferred,
  379. &endpt->rx_overlapped);
  380. }
  381. else if (endpt->type == TYPE_TCP) {
  382. do {
  383. WLog_DBG(TAG, "_recv ???");
  384. ret = _recv(endpt->sock_handle, (char*)&endpt->rx_pending_pkt_len, 4, 0);
  385. } while (ret < 0 && errno == EINTR);
  386. WLog_DBG(TAG, "_recv return: %d, %d", ret, errno);
  387. }
  388. else {
  389. WLog_ERR(TAG, "<== endpt(%d): start_read_pkt unkonwn type", endpt->epid);
  390. return -1;
  391. }
  392. if (ret >= 0) {
  393. dwBytesTransferred = ret;
  394. endpt->rx_pending_pkt_uc_len = ret;
  395. if (dwBytesTransferred == 0) {
  396. WLog_ERR(TAG, "<== endpt(%d): start_read_pkt dwBytesTransferred is error.", endpt->epid);
  397. return -1;
  398. }
  399. if (dwBytesTransferred < 4) {
  400. WLog_DBG(TAG, "endpt(%d): receive buffer less than dream len", endpt->epid);
  401. rc = recv_buf(endpt, (char*)&endpt->rx_pending_pkt_len + dwBytesTransferred, 4 - dwBytesTransferred);
  402. if (rc < 0) {
  403. WLog_ERR(TAG, "<== endpt(%d): start_read_pkt recv_buf error.", endpt->epid);
  404. return rc;
  405. }
  406. }
  407. WLog_DBG(TAG, "iobuffer_create ???");
  408. pkt = iobuffer_create(0, endpt->rx_pending_pkt_len);
  409. endpt->rx_pending_pkt_uc_len = 0;
  410. if (endpt->rx_pending_pkt_len > 0) {
  411. WLog_DBG(TAG, "endpt->rx_pending_pkt_len > 0 ???");
  412. rc = recv_buf(endpt, iobuffer_data(pkt, 0), endpt->rx_pending_pkt_len);
  413. if (rc < 0) {
  414. WLog_DBG(TAG, "iobuffer_destroy ???");
  415. iobuffer_destroy(pkt);
  416. WLog_ERR(TAG, "<== endpt(%d): start_read_pkt recv_buf error.", endpt->epid);
  417. return rc;
  418. }
  419. WLog_DBG(TAG, "iobuffer_push_count ???");
  420. iobuffer_push_count(pkt, endpt->rx_pending_pkt_len);
  421. }
  422. WLog_DBG(TAG, "*p_pkt = pkt ???");
  423. *p_pkt = pkt;
  424. WLog_ERR(TAG, "<== endpt(%d): start_read_pkt recv_buf succ.", endpt->epid);
  425. return 0;
  426. }
  427. else if(errno == EAGAIN || errno == EWOULDBLOCK) {
  428. WLog_DBG(TAG, "endpt(%d): set rx pending flag.", endpt->epid);
  429. endpt->rx_pending = 1;
  430. return 0;
  431. }
  432. WLog_ERR(TAG, "<== endpt(%d): start_read_pkt recv_buf failed.", endpt->epid);
  433. return -1;
  434. }
  435. static int read_left_pkt(bus_endpt_t* endpt, iobuffer_t** p_pkt)
  436. {
  437. BOOL ret;
  438. int rc;
  439. DWORD dwBytesTransferred;
  440. iobuffer_t* pkt = NULL;
  441. WLog_DBG(TAG, "==> fd(%d): read left pkt", endpt->sock_handle);
  442. if (endpt->type == TYPE_PIPE) {
  443. ret = GetOverlappedResult(endpt->pipe_handle, &endpt->rx_overlapped, &dwBytesTransferred, TRUE);
  444. }
  445. else if (endpt->type == TYPE_TCP) {
  446. do
  447. {
  448. ret = _recv(endpt->sock_handle,
  449. (char*)&endpt->rx_pending_pkt_len + endpt->rx_pending_pkt_uc_len,
  450. 4 - endpt->rx_pending_pkt_uc_len, 0);
  451. } while (ret == -1 && errno == EINTR);
  452. }
  453. if (ret < 0) {
  454. WLog_ERR(TAG, "<== fd(%d): read left pkt failed: ret %d, err: %s", endpt->sock_handle, ret, strerror(errno));
  455. return -1;
  456. }
  457. else if(ret == 0) {
  458. WLog_WARN(TAG, "<== fd(%d): peer socket close.", endpt->sock_handle);
  459. }
  460. if (dwBytesTransferred < 4) {
  461. rc = recv_buf(endpt, (char*)&endpt->rx_pending_pkt_len + dwBytesTransferred, 4 - dwBytesTransferred);
  462. if (rc < 0)
  463. return rc;
  464. }
  465. /*the first 4 bytes indicates the length of content and then read the buffer content.*/
  466. pkt = iobuffer_create(-1, endpt->rx_pending_pkt_len);
  467. WLog_DBG(TAG, "after read pkt len, start to read pkt's content: %d", endpt->rx_pending_pkt_len);
  468. rc = recv_buf(endpt, iobuffer_data(pkt, 0), endpt->rx_pending_pkt_len);
  469. if (rc < 0) {
  470. iobuffer_destroy(pkt);
  471. return rc;
  472. }
  473. iobuffer_push_count(pkt, endpt->rx_pending_pkt_len);
  474. *p_pkt = pkt;
  475. WLog_DBG(TAG, "endpt(%d): reset rx_pending and return!", endpt->epid);
  476. endpt->rx_pending = 0;
  477. return 0;
  478. }
  479. static int append_rx_pkt(bus_endpt_t* endpt, iobuffer_t* pkt)
  480. {
  481. int type;
  482. int read_state;
  483. read_state = iobuffer_get_read_state(pkt);
  484. iobuffer_read(pkt, IOBUF_T_I4, &type, 0);
  485. iobuffer_restore_read_state(pkt, read_state);
  486. if (type == BUS_TYPE_PACKET || type == BUS_TYPE_INFO || type == BUS_TYPE_EVENT || type == BUS_TYPE_SYSTEM) {
  487. iobuffer_queue_enqueue(endpt->rx_buf_queue, pkt);
  488. WLog_DBG(TAG, "<== append_rx_pkt finished");
  489. return 1;
  490. }
  491. else {
  492. WLog_DBG(TAG, "<== append_rx_pkt failed!");
  493. return -1;
  494. }
  495. }
  496. TOOLKIT_API int bus_endpt_create(const char* url, int epid, const bus_endpt_callback* callback, bus_endpt_t** p_endpt)
  497. {
  498. bus_endpt_t* endpt = NULL;
  499. char* tmp_url;
  500. url_fields uf;
  501. int rc;
  502. int v;
  503. iobuffer_t* buf = NULL;
  504. iobuffer_t* ans_buf = NULL;
  505. if (!url)
  506. return -1;
  507. tmp_url = _strdup(url);
  508. if (url_parse(tmp_url, &uf) < 0) {
  509. free(tmp_url);
  510. return -1;
  511. }
  512. endpt = ZALLOC_T(bus_endpt_t);
  513. endpt->sock_handle = -1;
  514. endpt->msg_fd = -1;
  515. endpt->ep = NULL;
  516. endpt->url = tmp_url;
  517. if (_stricmp(uf.scheme, "tcp") == 0) {
  518. endpt->type = TYPE_TCP;
  519. endpt->sock_handle = create_socket_handle(uf.host, uf.port);
  520. if (endpt->sock_handle == INVALID_SOCKET)
  521. goto on_error;
  522. WLog_INFO(TAG, "bus endpt socket fd: %d", endpt->sock_handle);
  523. }
  524. else if (_stricmp(uf.scheme, "pipe") == 0) {
  525. endpt->type = TYPE_PIPE;
  526. endpt->pipe_handle = create_pipe_handle(uf.host);
  527. if (endpt->pipe_handle == INVALID_HANDLE_VALUE)
  528. goto on_error;
  529. WLog_INFO(TAG, "bus endpt pipe fd: %d", endpt->sock_handle);
  530. }
  531. else {
  532. goto on_error;
  533. }
  534. endpt->ep = evtpoll_create();
  535. if (!endpt->ep) {
  536. WLog_ERR(TAG, "evtpoll create failed!");
  537. goto on_error;
  538. }
  539. endpt->epid = epid;
  540. endpt->tx_evt = CreateEventA(NULL, TRUE, FALSE, NULL);
  541. endpt->rx_evt = CreateEventA(NULL, TRUE, FALSE, NULL);
  542. endpt->rx_buf_queue = iobuffer_queue_create();
  543. endpt->msg_fd = eventfd(0, EFD_CLOEXEC | EFD_SEMAPHORE);
  544. if (endpt->msg_fd == -1) {
  545. WLog_ERR(TAG, "create event fd failed: %s(%d)", strerror(errno), errno);
  546. goto on_error;
  547. }
  548. INIT_LIST_HEAD(&endpt->msg_list);
  549. spinlock_init(&endpt->msg_lock);
  550. memcpy(&endpt->callback, callback, sizeof(bus_endpt_callback));
  551. buf = iobuffer_create(-1, -1);
  552. v = BUS_TYPE_ENDPT_REGISTER;
  553. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  554. v = endpt->epid;
  555. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  556. WLog_DBG(TAG, "start to send_pkt_raw");
  557. rc = send_pkt_raw(endpt, buf);
  558. WLog_INFO(TAG, "send_pkt_raw return %d", rc);
  559. if (rc != 0)
  560. goto on_error;
  561. rc = recv_pkt_raw(endpt, &ans_buf);
  562. WLog_INFO(TAG, "recv_pkt_raw return %d", rc);
  563. if (rc != 0) {
  564. goto on_error;
  565. }
  566. iobuffer_read(ans_buf, IOBUF_T_I4, &v, 0);
  567. iobuffer_read(ans_buf, IOBUF_T_I4, &rc, 0);
  568. if (rc != 0)
  569. goto on_error;
  570. url_free_fields(&uf);
  571. if (buf)
  572. iobuffer_destroy(buf);
  573. if (ans_buf)
  574. iobuffer_destroy(ans_buf);
  575. if (-1 == evtpoll_attach(endpt->ep, endpt->msg_fd)) {
  576. goto on_error;
  577. }
  578. if (evtpoll_subscribe(endpt->ep, EV_READ, endpt->msg_fd, &endpt->msg_fd, NULL)) {
  579. WLog_ERR(TAG, "epoll subscribe bus endpt eventfd failed.");
  580. goto on_error_msg;
  581. }
  582. /* subscribe read event [3/27/2020 Gifur] */
  583. if (evtpoll_attach(endpt->ep, endpt->sock_handle)) {
  584. WLog_ERR(TAG, "epoll attch bus endpt failed.");
  585. goto on_error_msg;
  586. }
  587. if (evtpoll_subscribe(endpt->ep, EV_READ, endpt->sock_handle, &endpt->sock_handle, NULL)) {
  588. WLog_ERR(TAG, "epoll subscribe bus endpt failed.");
  589. goto on_error_handle;
  590. }
  591. *p_endpt = endpt;
  592. return 0;
  593. on_error_handle:
  594. evtpoll_detach(endpt->ep, endpt->sock_handle);
  595. on_error_msg:
  596. evtpoll_detach(endpt->ep, endpt->msg_fd);
  597. on_error:
  598. if (endpt->type == TYPE_TCP) {
  599. closesocket(endpt->sock_handle);
  600. }
  601. else if (endpt->type == TYPE_PIPE) {
  602. CloseHandle(endpt->pipe_handle);
  603. }
  604. if (endpt->msg_fd > 0)
  605. close(endpt->msg_fd);
  606. if (endpt->tx_evt)
  607. CloseHandle(endpt->tx_evt);
  608. if (endpt->rx_evt)
  609. CloseHandle(endpt->rx_evt);
  610. if (endpt->rx_buf_queue)
  611. iobuffer_queue_destroy(endpt->rx_buf_queue);
  612. if (endpt->url)
  613. free(endpt->url);
  614. free(endpt);
  615. url_free_fields(&uf);
  616. if (buf)
  617. iobuffer_destroy(buf);
  618. if (ans_buf)
  619. iobuffer_destroy(ans_buf);
  620. if (endpt->ep != NULL) {
  621. evtpoll_destroy(endpt->ep);
  622. }
  623. return -1;
  624. }
  625. TOOLKIT_API void bus_endpt_destroy(bus_endpt_t* endpt)
  626. {
  627. int rc = -1;
  628. iobuffer_t* buf = NULL;
  629. iobuffer_t* ans_buf = NULL;
  630. int v;
  631. assert(endpt);
  632. buf = iobuffer_create(-1, -1);
  633. v = BUS_TYPE_ENDPT_UNREGISTER;
  634. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  635. v = endpt->epid;
  636. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  637. rc = send_pkt_raw(endpt, buf);
  638. if (rc != 0)
  639. goto on_error;
  640. rc = recv_pkt_raw(endpt, &ans_buf);
  641. if (rc != 0)
  642. goto on_error;
  643. iobuffer_read(ans_buf, IOBUF_T_I4, &v, 0);
  644. iobuffer_read(ans_buf, IOBUF_T_I4, &rc, 0);
  645. {
  646. msg_t* msg, * t;
  647. list_for_each_entry_safe(msg, t, &endpt->msg_list, msg_t, entry) {
  648. list_del(&msg->entry);
  649. if (msg->evt)
  650. CloseHandle(msg->evt);
  651. free(msg);
  652. }
  653. }
  654. on_error:
  655. if (buf)
  656. iobuffer_destroy(buf);
  657. if (ans_buf)
  658. iobuffer_destroy(ans_buf);
  659. if (endpt->type == TYPE_TCP) {
  660. closesocket(endpt->sock_handle);
  661. }
  662. else if (endpt->type == TYPE_PIPE) {
  663. CloseHandle(endpt->pipe_handle);
  664. }
  665. if (endpt->msg_fd)
  666. close(endpt->msg_fd);
  667. if (endpt->tx_evt)
  668. CloseHandle(endpt->tx_evt);
  669. if (endpt->rx_evt)
  670. CloseHandle(endpt->rx_evt);
  671. if (endpt->rx_buf_queue)
  672. iobuffer_queue_destroy(endpt->rx_buf_queue);
  673. assert(endpt->ep);
  674. evtpoll_destroy(endpt->ep);
  675. if (endpt->url)
  676. free(endpt->url);
  677. free(endpt);
  678. }
  679. // 1 : recv ok
  680. // 0 : time out
  681. // <0 : error
  682. static int bus_endpt_poll_internal(bus_endpt_t* endpt, int* result, int timeout)
  683. {
  684. iobuffer_t* pkt = NULL;
  685. int rc;
  686. BOOL ret;
  687. assert(endpt);
  688. // peek first packge type
  689. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  690. pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  691. }
  692. else {
  693. // no received package, try to receive one
  694. if (!endpt->rx_pending) {
  695. rc = start_read_pkt(endpt, &pkt);
  696. if (rc < 0) {
  697. WLog_ERR(TAG, "start read pkt failed.");
  698. return rc;
  699. }
  700. if (pkt) {
  701. WLog_INFO(TAG, "pkt has read");
  702. rc = append_rx_pkt(endpt, pkt); // append pkt to rx_buf_queue
  703. if (rc < 0) {
  704. iobuffer_destroy(pkt);
  705. return -1;
  706. }
  707. }
  708. }
  709. else {
  710. //WLog_ERR(TAG, "is pending now.");
  711. }
  712. // if receive is pending, wait for send or receive complete event
  713. if (!pkt) {
  714. //WLog_DBG(TAG, "wait msg sem or received event. tiemout: %d", timeout);
  715. {
  716. int nfds;
  717. int ret;
  718. int i;
  719. struct epoll_event events[MAX_EPOLL_EVENT];
  720. struct epoll_event* pe;
  721. nfds = evtpoll_wait(endpt->ep, events, MAX_EPOLL_EVENT, timeout);
  722. if (nfds == 0) {
  723. //WLog_DBG(TAG, "epoll wait timeout.");
  724. return 0; //timeout
  725. }
  726. if (nfds == -1) {
  727. return -1;
  728. }
  729. WLog_DBG(TAG, "epoll wait return nfd: %d", nfds);
  730. for (i = 0; i < nfds; ++i) {
  731. void* pdata = NULL;
  732. pe = events + i;
  733. WLog_INFO(TAG, "loop events[%d]::fd(0x%08X) OUT:%d, IN:%d", i, pe->data.fd,
  734. pe->events & EPOLLOUT ? 1 : 0, pe->events & EPOLLIN ? 1 : 0);
  735. assert(pe->events & EPOLLIN);
  736. ret = evtpoll_deal(endpt->ep, pe, &pdata, 0);
  737. if (!ret) {
  738. assert(pdata);
  739. if(bus_endpoint__data_is_handle(endpt, pdata))
  740. {
  741. rc = read_left_pkt(endpt, &pkt);
  742. if (rc < 0)
  743. return rc;
  744. if (pkt) {
  745. rc = append_rx_pkt(endpt, pkt);
  746. if (rc < 0) {
  747. iobuffer_destroy(pkt);
  748. return -1;
  749. }
  750. }
  751. }
  752. else if (pdata == &endpt->msg_fd) {
  753. uint64_t rdata;
  754. WLog_DBG(TAG, "message arrive.");
  755. do
  756. {
  757. ret = read(endpt->msg_fd, &rdata, sizeof rdata);
  758. } while (ret < 0 && errno == EINTR);
  759. if (ret < 0) {
  760. WLog_ERR(TAG, "read msg fd failed: %s", strerror(errno));
  761. abort();
  762. }
  763. *result = BUS_RESULT_MSG;
  764. return 1;
  765. }
  766. }
  767. }
  768. }
  769. }
  770. else {
  771. WLog_ERR(TAG, "pkt has readed");
  772. }
  773. }
  774. if (pkt) {
  775. int type;
  776. int read_state = iobuffer_get_read_state(pkt);
  777. iobuffer_read(pkt, IOBUF_T_I4, &type, 0);
  778. iobuffer_restore_read_state(pkt, read_state);
  779. *result = to_result(type);
  780. if (*result == BUS_RESULT_UNKNOWN) {
  781. WLog_ERR(TAG, "bug: unknown pkt type!");
  782. return -1;
  783. }
  784. return 1;
  785. }
  786. return -1;
  787. }
  788. static int recv_until(bus_endpt_t* endpt, int type, iobuffer_t** p_ansbuf)
  789. {
  790. int rc;
  791. iobuffer_t* ans_pkt = NULL;
  792. int ans_type;
  793. WLog_DBG(TAG, "==>endpt(%d): recv until type: 0x%08X", endpt->epid, type);
  794. for (;;) {
  795. if (!endpt->rx_pending) {
  796. rc = start_read_pkt(endpt, &ans_pkt);
  797. if (rc < 0) {
  798. break;
  799. }
  800. } else {
  801. WLog_DBG(TAG, "endpt(%d) is pending", endpt->epid);
  802. }
  803. if (!ans_pkt) {
  804. int nfds;
  805. int ret;
  806. int i, flag = 0;
  807. struct epoll_event events[MAX_EPOLL_EVENT];
  808. struct epoll_event* pe;
  809. nfds = evtpoll_wait(endpt->ep, events, MAX_EPOLL_EVENT, -1);
  810. if (nfds == 0) {
  811. continue;
  812. }
  813. if (nfds == -1) {
  814. return -1;
  815. }
  816. WLog_DBG(TAG, "epoll wait return nfd: %d", nfds);
  817. for (i = 0; i < nfds; ++i) {
  818. void* pdata = NULL;
  819. pe = events + i;
  820. ret = evtpoll_deal(endpt->ep, pe, &pdata, 0);
  821. if (!ret && bus_endpoint__data_is_handle(endpt, pdata)) {
  822. flag = 1;
  823. break;
  824. }
  825. }
  826. if (flag) {
  827. rc = read_left_pkt(endpt, &ans_pkt);
  828. if (rc < 0) {
  829. break;
  830. }
  831. }
  832. }
  833. if (ans_pkt) {
  834. int read_state = iobuffer_get_read_state(ans_pkt);
  835. iobuffer_read(ans_pkt, IOBUF_T_I4, &ans_type, 0);
  836. iobuffer_restore_read_state(ans_pkt, read_state);
  837. if (ans_type == type) {
  838. *p_ansbuf = ans_pkt;
  839. break;
  840. }
  841. else {
  842. rc = append_rx_pkt(endpt, ans_pkt);
  843. if (rc < 0) {
  844. iobuffer_destroy(ans_pkt);
  845. break;
  846. }
  847. else {
  848. ans_pkt = NULL;
  849. }
  850. }
  851. }
  852. }
  853. return rc;
  854. }
  855. static int recv_until_result(bus_endpt_t* endpt, int* p_result)
  856. {
  857. int rc;
  858. iobuffer_t* ans_pkt = NULL;
  859. int type, error;
  860. rc = recv_until(endpt, BUS_TYPE_ERROR, &ans_pkt);
  861. if (rc < 0)
  862. return rc;
  863. iobuffer_read(ans_pkt, IOBUF_T_I4, &type, 0);
  864. iobuffer_read(ans_pkt, IOBUF_T_I4, &error, 0);
  865. iobuffer_destroy(ans_pkt);
  866. *p_result = error;
  867. return rc;
  868. }
  869. static int recv_until_state(bus_endpt_t* endpt, int* p_state)
  870. {
  871. int rc;
  872. iobuffer_t* ans_pkt = NULL;
  873. int type, epid, state;
  874. rc = recv_until(endpt, BUS_TYPE_ENDPT_GET_STATE, &ans_pkt);
  875. if (rc < 0)
  876. return rc;
  877. iobuffer_read(ans_pkt, IOBUF_T_I4, &type, 0);
  878. iobuffer_read(ans_pkt, IOBUF_T_I4, &epid, 0);
  879. iobuffer_read(ans_pkt, IOBUF_T_I4, &state, 0);
  880. iobuffer_destroy(ans_pkt);
  881. WLog_DBG(TAG, "state address: 0x%08X", p_state);
  882. *p_state = state;
  883. return rc;
  884. }
  885. TOOLKIT_API int bus_endpt_send_pkt(bus_endpt_t* endpt, int epid, int type, iobuffer_t* pkt)
  886. {
  887. int t;
  888. int rc;
  889. int read_state;
  890. int write_state;
  891. int error;
  892. assert(endpt);
  893. read_state = iobuffer_get_read_state(pkt);
  894. write_state = iobuffer_get_write_state(pkt);
  895. t = epid; // remote epid
  896. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  897. t = endpt->epid; // local epid
  898. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  899. t = type; // user type
  900. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  901. t = BUS_TYPE_PACKET; // type
  902. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  903. rc = send_pkt_raw(endpt, pkt);
  904. iobuffer_restore_read_state(pkt, read_state);
  905. iobuffer_restore_write_state(pkt, write_state);
  906. if (rc < 0)
  907. return rc;
  908. rc = recv_until_result(endpt, &error);
  909. if (rc == 0 && error != 0)
  910. rc = error;
  911. return rc;
  912. }
  913. TOOLKIT_API int bus_endpt_send_info(bus_endpt_t* endpt, int epid, int type, iobuffer_t* pkt)
  914. {
  915. int t;
  916. int rc;
  917. int read_state;
  918. int write_state;
  919. assert(endpt);
  920. WLog_DBG(TAG, "==> endpt(%d) send info: %d, 0x%08X.", endpt->epid, epid, type);
  921. read_state = iobuffer_get_read_state(pkt);
  922. write_state = iobuffer_get_write_state(pkt);
  923. t = epid; // remote epid
  924. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  925. t = endpt->epid; // local epid
  926. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  927. t = type; // user type
  928. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  929. t = BUS_TYPE_INFO; // type
  930. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  931. rc = send_pkt_raw(endpt, pkt);
  932. iobuffer_restore_read_state(pkt, read_state);
  933. iobuffer_restore_write_state(pkt, write_state);
  934. return rc;
  935. }
  936. TOOLKIT_API int bus_endpt_bcast_evt(bus_endpt_t* endpt, int type, iobuffer_t* pkt)
  937. {
  938. int t;
  939. int rc;
  940. int read_state;
  941. int write_state;
  942. assert(endpt);
  943. read_state = iobuffer_get_read_state(pkt);
  944. write_state = iobuffer_get_write_state(pkt);
  945. t = endpt->epid;
  946. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  947. t = type;
  948. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  949. t = BUS_TYPE_EVENT;
  950. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  951. rc = send_pkt_raw(endpt, pkt);
  952. iobuffer_restore_read_state(pkt, read_state);
  953. iobuffer_restore_write_state(pkt, write_state);
  954. return rc;
  955. }
  956. static int bus_endpt_recv_pkt(bus_endpt_t* endpt, int* p_epid, int* p_type, iobuffer_t** p_pkt)
  957. {
  958. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  959. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  960. int read_state = iobuffer_get_read_state(pkt);
  961. int pkt_type, usr_type, from_epid, to_epid;
  962. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, 0);
  963. if (pkt_type == BUS_TYPE_PACKET || pkt_type == BUS_TYPE_INFO) {
  964. iobuffer_read(pkt, IOBUF_T_I4, &usr_type, 0);
  965. iobuffer_read(pkt, IOBUF_T_I4, &from_epid, 0);
  966. iobuffer_read(pkt, IOBUF_T_I4, &to_epid, 0);
  967. if (p_epid)
  968. *p_epid = from_epid;
  969. if (p_type)
  970. *p_type = usr_type;
  971. iobuffer_queue_deque(endpt->rx_buf_queue);
  972. if (p_pkt) {
  973. *p_pkt = pkt;
  974. }
  975. else {
  976. iobuffer_destroy(pkt);
  977. }
  978. return 0;
  979. }
  980. else {
  981. iobuffer_restore_read_state(pkt, read_state);
  982. }
  983. }
  984. return -1;
  985. }
  986. static int bus_endpt_recv_evt(bus_endpt_t* endpt, int* p_epid, int* p_type, iobuffer_t** p_pkt)
  987. {
  988. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  989. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  990. int read_state = iobuffer_get_read_state(pkt);
  991. int pkt_type, usr_type, from_epid;
  992. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, 0);
  993. if (pkt_type == BUS_TYPE_EVENT) {
  994. iobuffer_read(pkt, IOBUF_T_I4, &usr_type, 0);
  995. iobuffer_read(pkt, IOBUF_T_I4, &from_epid, 0);
  996. if (p_epid)
  997. *p_epid = from_epid;
  998. if (p_type)
  999. *p_type = usr_type;
  1000. iobuffer_queue_deque(endpt->rx_buf_queue);
  1001. if (p_pkt) {
  1002. *p_pkt = pkt;
  1003. }
  1004. else {
  1005. iobuffer_destroy(pkt);
  1006. }
  1007. return 0;
  1008. }
  1009. else {
  1010. iobuffer_restore_read_state(pkt, read_state);
  1011. }
  1012. }
  1013. return -1;
  1014. }
  1015. static int bus_endpt_recv_sys(bus_endpt_t* endpt, int* p_epid, int* p_state)
  1016. {
  1017. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  1018. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  1019. int read_state = iobuffer_get_read_state(pkt);
  1020. int pkt_type, epid, state;
  1021. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, 0);
  1022. if (pkt_type == BUS_TYPE_SYSTEM) {
  1023. iobuffer_read(pkt, IOBUF_T_I4, &epid, 0);
  1024. iobuffer_read(pkt, IOBUF_T_I4, &state, 0);
  1025. if (p_epid)
  1026. *p_epid = epid;
  1027. if (p_state)
  1028. *p_state = state;
  1029. iobuffer_queue_deque(endpt->rx_buf_queue);
  1030. iobuffer_destroy(pkt);
  1031. return 0;
  1032. }
  1033. else {
  1034. iobuffer_restore_read_state(pkt, read_state);
  1035. }
  1036. }
  1037. return -1;
  1038. }
  1039. static int bus_endpt_recv_msg(bus_endpt_t* endpt, msg_t** p_msg)
  1040. {
  1041. int rc = -1;
  1042. assert(endpt);
  1043. assert(p_msg);
  1044. spinlock_enter(&endpt->msg_lock, -1);
  1045. if (!list_empty(&endpt->msg_list)) {
  1046. msg_t* e = list_first_entry(&endpt->msg_list, msg_t, entry);
  1047. list_del(&e->entry);
  1048. rc = 0;
  1049. *p_msg = e;
  1050. }
  1051. spinlock_leave(&endpt->msg_lock);
  1052. return rc;
  1053. }
  1054. TOOLKIT_API int bus_endpt_get_state(bus_endpt_t* endpt, int epid, int* p_state)
  1055. {
  1056. iobuffer_t* buf = NULL;
  1057. int v;
  1058. int rc = -1;
  1059. assert(endpt);
  1060. buf = iobuffer_create(-1, -1);
  1061. v = BUS_TYPE_ENDPT_GET_STATE;
  1062. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  1063. v = epid;
  1064. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  1065. rc = send_pkt_raw(endpt, buf);
  1066. if (rc < 0) {
  1067. WLog_ERR(TAG, "send pkt raw failed.");
  1068. goto on_error;
  1069. }
  1070. rc = recv_until_state(endpt, p_state);
  1071. on_error:
  1072. if (buf)
  1073. iobuffer_destroy(buf);
  1074. return rc;
  1075. }
  1076. TOOLKIT_API int bus_endpt_post_msg(bus_endpt_t* endpt, int msg, int nparam, param_size_t params[])
  1077. {
  1078. msg_t* e;
  1079. int rc;
  1080. uint64_t wdata = 0;
  1081. assert(endpt);
  1082. WLog_DBG(TAG, "==> endpt(%d) post msg: %d", endpt->epid, msg);
  1083. e = MALLOC_T(msg_t);
  1084. e->type = msg;
  1085. e->nparam = nparam;
  1086. if (nparam) {
  1087. e->params = (param_size_t*)malloc(sizeof(param_size_t) * nparam);
  1088. memcpy(e->params, params, sizeof(param_size_t) * nparam);
  1089. }
  1090. else {
  1091. e->params = NULL;
  1092. }
  1093. e->evt = NULL;
  1094. spinlock_enter(&endpt->msg_lock, -1);
  1095. list_add_tail(&e->entry, &endpt->msg_list);
  1096. spinlock_leave(&endpt->msg_lock);
  1097. wdata = 1;
  1098. rc = write(endpt->msg_fd, &wdata, sizeof wdata);
  1099. if (rc == -1) {
  1100. WLog_ERR(TAG, "write to eventfd failed: %s(%d)", strerror(errno), errno);
  1101. return -1;
  1102. }
  1103. return 0;
  1104. }
  1105. TOOLKIT_API int bus_endpt_send_msg(bus_endpt_t* endpt, int msg, int nparam, param_size_t params[])
  1106. {
  1107. msg_t e;
  1108. int rc;
  1109. uint64_t wdata = 0;
  1110. assert(endpt);
  1111. WLog_DBG(TAG, "==> endpt(%d) send msg: epid %d, param counts %d", endpt->epid, msg, nparam);
  1112. e.type = msg;
  1113. e.nparam = nparam;
  1114. if (nparam) {
  1115. e.params = (param_size_t*)malloc(sizeof(param_size_t) * nparam);
  1116. memcpy(e.params, params, sizeof(param_size_t) * nparam);
  1117. }
  1118. else {
  1119. e.params = NULL;
  1120. }
  1121. e.evt_result = 0;
  1122. e.evt = CreateEventA(NULL, TRUE, FALSE, NULL);
  1123. assert(e.evt != NULL);
  1124. spinlock_enter(&endpt->msg_lock, -1);
  1125. list_add_tail(&e.entry, &endpt->msg_list);
  1126. spinlock_leave(&endpt->msg_lock);
  1127. wdata = 1;
  1128. rc = write(endpt->msg_fd, &wdata, sizeof wdata);
  1129. if (rc == -1) {
  1130. WLog_ERR(TAG, "write to eventfd failed: %s(%d)", strerror(errno), errno);
  1131. CloseHandle(e.evt);
  1132. if (nparam) {
  1133. free(e.params);
  1134. }
  1135. WLog_DBG(TAG, "<== error endpt(%d) send msg: %d", endpt->epid, msg);
  1136. return -1;
  1137. }
  1138. WaitForSingleObject(e.evt, INFINITE);
  1139. CloseHandle(e.evt);
  1140. if (nparam) {
  1141. free(e.params);
  1142. }
  1143. WLog_DBG(TAG, "<== endpt(%d) send msg: epid %d, evt_res: %d", endpt->epid, msg, e.evt_result);
  1144. return e.evt_result;
  1145. }
  1146. TOOLKIT_API int bus_endpt_get_epid(bus_endpt_t* endpt)
  1147. {
  1148. return endpt->epid;
  1149. }
  1150. TOOLKIT_API const char* bus_endpt_get_url(bus_endpt_t* endpt)
  1151. {
  1152. return endpt->url;
  1153. }
  1154. TOOLKIT_API int bus_endpt_poll(bus_endpt_t* endpt, int timeout)
  1155. {
  1156. int result;
  1157. int rc;
  1158. int epid, type, state;
  1159. iobuffer_t* pkt = NULL;
  1160. rc = bus_endpt_poll_internal(endpt, &result, timeout);
  1161. //WLog_DBG(TAG, "bus endpt poll internal: %d, %d", rc, result);
  1162. if (rc > 0) {
  1163. if (result == BUS_RESULT_DATA) {
  1164. bus_endpt_recv_pkt(endpt, &epid, &type, &pkt);
  1165. if (endpt->callback.on_pkt)
  1166. endpt->callback.on_pkt(endpt, epid, type, &pkt, endpt->callback.user_data);
  1167. if (pkt)
  1168. iobuffer_dec_ref(pkt);
  1169. }
  1170. else if (result == BUS_RESULT_INFO) {
  1171. bus_endpt_recv_pkt(endpt, &epid, &type, &pkt);
  1172. if (endpt->callback.on_inf)
  1173. endpt->callback.on_inf(endpt, epid, type, &pkt, endpt->callback.user_data);
  1174. if (pkt)
  1175. iobuffer_dec_ref(pkt);
  1176. }
  1177. else if (result == BUS_RESULT_EVENT) {
  1178. bus_endpt_recv_evt(endpt, &epid, &type, &pkt);
  1179. if (endpt->callback.on_evt)
  1180. endpt->callback.on_evt(endpt, epid, type, &pkt, endpt->callback.user_data);
  1181. if (pkt)
  1182. iobuffer_dec_ref(pkt);
  1183. }
  1184. else if (result == BUS_RESULT_SYSTEM) {
  1185. bus_endpt_recv_sys(endpt, &epid, &state);
  1186. if (endpt->callback.on_sys)
  1187. endpt->callback.on_sys(endpt, epid, state, endpt->callback.user_data);
  1188. }
  1189. else if (result == BUS_RESULT_MSG) {
  1190. msg_t* msg = NULL;
  1191. bus_endpt_recv_msg(endpt, &msg);
  1192. if (endpt->callback.on_msg) {
  1193. endpt->callback.on_msg(endpt,
  1194. msg->type,
  1195. msg->nparam,
  1196. msg->params,
  1197. msg->evt ? &msg->evt_result : NULL,
  1198. endpt->callback.user_data);
  1199. if (msg->evt) {
  1200. WLog_DBG(TAG, "after recv msg, send finished evt.");
  1201. SetEvent(msg->evt);
  1202. }
  1203. else {
  1204. WLog_DBG(TAG, "free msg");
  1205. free_msg(msg);
  1206. }
  1207. }
  1208. else {
  1209. if (msg->evt) {
  1210. msg->evt_result = -1;
  1211. WLog_DBG(TAG, "after on msg failed, send finished evt.");
  1212. SetEvent(msg->evt);
  1213. }
  1214. else {
  1215. WLog_DBG(TAG, "free msg");
  1216. free_msg(msg);
  1217. }
  1218. }
  1219. }
  1220. else {
  1221. assert(0);
  1222. rc = -1;
  1223. }
  1224. }
  1225. return rc;
  1226. }
  1227. TOOLKIT_API int bus_endpt_set_quit_flag(bus_endpt_t* endpt)
  1228. {
  1229. endpt->quit_flag = 1;
  1230. return 0;
  1231. }
  1232. TOOLKIT_API int bus_endpt_get_quit_flag(bus_endpt_t* endpt)
  1233. {
  1234. return endpt->quit_flag;
  1235. }