bus-unix.c 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158
  1. #include "precompile.h"
  2. #include "bus.h"
  3. #include "sockutil.h"
  4. #include "url.h"
  5. #include "memutil.h"
  6. #include "spinlock.h"
  7. #include "list.h"
  8. #include "bus_internal.h"
  9. #include <winpr/file.h>
  10. #include <winpr/pipe.h>
  11. #include <winpr/synch.h>
  12. #include <winpr/string.h>
  13. #define TAG TOOLKIT_TAG("bus_unix")
  14. #define BUS_RESULT_DATA 1 // ==BUS_TYPE_PACKET, callback: callback.on_pkt, no use
  15. #define BUS_RESULT_INFO 2 // ==BUS_TYPE_INFO, callback: callback.on_inf
  16. #define BUS_RESULT_EVENT 3 // ==BUS_TYPE_EVENT, callback: callback.on_evt , no use
  17. #define BUS_RESULT_SYSTEM 4 // ==BUS_TYPE_SYSTEM, callback: callback.on_sys
  18. #define BUS_RESULT_MSG 5 // send package msg, callback: callback.on_msg
  19. #define BUS_RESULT_UNKNOWN 6
  20. typedef struct msg_t {
  21. struct list_head entry;
  22. int type;
  23. int nparam;
  24. int* params;
  25. HANDLE evt;
  26. int evt_result;
  27. }msg_t;
  28. struct bus_endpt_t {
  29. int type;
  30. int epid;
  31. union {
  32. HANDLE pipe_handle;
  33. SOCKET sock_handle;
  34. };
  35. char* url;
  36. bus_endpt_callback callback;
  37. struct list_head msg_list;
  38. spinlock_t msg_lock;
  39. HANDLE msg_sem;
  40. HANDLE tx_evt;
  41. HANDLE rx_evt;
  42. OVERLAPPED rx_overlapped;
  43. int rx_pending;
  44. int rx_pending_pkt_len;
  45. iobuffer_queue_t* rx_buf_queue;
  46. volatile int quit_flag;
  47. };
  48. static void free_msg(msg_t* msg)
  49. {
  50. free(msg->params);
  51. free(msg);
  52. }
  53. static int to_result(int pkt_type)
  54. {
  55. switch (pkt_type) {
  56. case BUS_TYPE_EVENT:
  57. return BUS_RESULT_EVENT;
  58. case BUS_TYPE_SYSTEM:
  59. return BUS_RESULT_SYSTEM;
  60. case BUS_TYPE_PACKET:
  61. return BUS_RESULT_DATA;
  62. case BUS_TYPE_INFO:
  63. return BUS_RESULT_INFO;
  64. default:
  65. break;
  66. }
  67. return BUS_RESULT_UNKNOWN;
  68. }
  69. static HANDLE create_pipe_handle(const char* name)
  70. {
  71. char tmp[MAX_PATH];
  72. HANDLE pipe = INVALID_HANDLE_VALUE;
  73. sprintf(tmp, "\\\\.\\pipe\\%s", name);
  74. for (;;) {
  75. pipe = CreateFileA(tmp,
  76. GENERIC_READ | GENERIC_WRITE,
  77. 0,
  78. NULL,
  79. OPEN_EXISTING,
  80. FILE_FLAG_OVERLAPPED,
  81. NULL);
  82. if (pipe == INVALID_HANDLE_VALUE) {
  83. if (GetLastError() == ERROR_PIPE_BUSY) {
  84. if (WaitNamedPipeA(name, 20000))
  85. continue;
  86. }
  87. }
  88. break;
  89. }
  90. return pipe;
  91. }
  92. static SOCKET create_socket_handle(const char* ip, int port)
  93. {
  94. SOCKET fd;
  95. struct sockaddr_in addr;
  96. fd = WSASocketA(AF_INET, SOCK_STREAM, IPPROTO_TCP, NULL, 0, WSA_FLAG_OVERLAPPED);
  97. if (fd == INVALID_SOCKET)
  98. return fd;
  99. {
  100. BOOL f = TRUE;
  101. u_long l = TRUE;
  102. setsockopt(fd, SOL_SOCKET, SO_DONTLINGER, (char*)&f, sizeof(f));
  103. //setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, (char*)&f, sizeof(f));
  104. _ioctlsocket(fd, FIONBIO, &l);
  105. }
  106. memset(&addr, 0, sizeof(addr));
  107. addr.sin_family = AF_INET;
  108. addr.sin_port = htons(port);
  109. addr.sin_addr.s_addr = inet_addr(ip);
  110. if (connect(fd, (struct sockaddr*) & addr, sizeof(addr)) != 0) {
  111. if (WSAGetLastError() == WSAEWOULDBLOCK) {
  112. fd_set wr_set, ex_set;
  113. FD_ZERO(&wr_set);
  114. FD_ZERO(&ex_set);
  115. FD_SET(fd, &wr_set);
  116. FD_SET(fd, &ex_set);
  117. if (select(fd, NULL, &wr_set, &ex_set, NULL) > 0 && FD_ISSET(fd, &wr_set))
  118. return fd;
  119. }
  120. }
  121. if (fd != INVALID_SOCKET)
  122. closesocket(fd);
  123. //return INVALID_SOCKET;//{bug}
  124. return fd;
  125. }
  126. static int tcp_send_buf(bus_endpt_t* endpt, const char* buf, int n)
  127. {
  128. DWORD left = n;
  129. DWORD offset = 0;
  130. while (left > 0) {
  131. BOOL ret;
  132. WSABUF wsabuf;
  133. DWORD dwBytesTransfer;
  134. OVERLAPPED overlapped;
  135. memset(&overlapped, 0, sizeof(overlapped));
  136. overlapped.hEvent = endpt->tx_evt;
  137. ResetEvent(endpt->tx_evt);
  138. wsabuf.buf = (char*)buf + offset;
  139. wsabuf.len = left;
  140. ret = WSASend(endpt->sock_handle, &wsabuf, 1, &dwBytesTransfer, 0, &overlapped, NULL) == 0;
  141. if (!ret && WSAGetLastError() == WSA_IO_PENDING) {
  142. DWORD dwFlags = 0;
  143. ret = WSAGetOverlappedResult(endpt->sock_handle, &overlapped, &dwBytesTransfer, TRUE, &dwFlags);
  144. }
  145. if (ret && dwBytesTransfer) {
  146. offset += dwBytesTransfer;
  147. left -= dwBytesTransfer;
  148. }
  149. else {
  150. return -1;
  151. }
  152. }
  153. return 0;
  154. }
  155. static int pipe_send_buf(bus_endpt_t* endpt, const char* buf, int n)
  156. {
  157. DWORD left = n;
  158. DWORD offset = 0;
  159. while (left > 0) {
  160. BOOL ret;
  161. DWORD dwBytesTransfer;
  162. OVERLAPPED overlapped;
  163. memset(&overlapped, 0, sizeof(overlapped));
  164. overlapped.hEvent = endpt->tx_evt;
  165. ResetEvent(endpt->tx_evt);
  166. ret = WriteFile(endpt->pipe_handle, buf + offset, left, &dwBytesTransfer, &overlapped);
  167. if (!ret && GetLastError() == ERROR_IO_PENDING) {
  168. ret = GetOverlappedResult(endpt->pipe_handle, &overlapped, &dwBytesTransfer, TRUE);
  169. }
  170. if (ret && dwBytesTransfer) {
  171. offset += dwBytesTransfer;
  172. left -= dwBytesTransfer;
  173. }
  174. else {
  175. return -1;
  176. }
  177. }
  178. return 0;
  179. }
  180. static int send_buf(bus_endpt_t* endpt, const char* buf, int n)
  181. {
  182. if (endpt->type == TYPE_PIPE) {
  183. return pipe_send_buf(endpt, buf, n);
  184. }
  185. else if (endpt->type == TYPE_TCP) {
  186. return tcp_send_buf(endpt, buf, n);
  187. }
  188. else {
  189. return -1;
  190. }
  191. }
  192. static int send_pkt_raw(bus_endpt_t* endpt, iobuffer_t* pkt)
  193. {
  194. int pkt_len = iobuffer_get_length(pkt);
  195. int rc;
  196. iobuffer_write_head(pkt, IOBUF_T_I4, &pkt_len, 0);
  197. rc = send_buf(endpt, iobuffer_data(pkt, 0), iobuffer_get_length(pkt));
  198. return rc;
  199. }
  200. static int pipe_recv_buf(bus_endpt_t* endpt, char* buf, DWORD n)
  201. {
  202. DWORD left = n;
  203. DWORD offset = 0;
  204. while (left > 0) {
  205. BOOL ret;
  206. DWORD dwBytesTransfer;
  207. OVERLAPPED overlapped;
  208. memset(&overlapped, 0, sizeof(overlapped));
  209. overlapped.hEvent = endpt->rx_evt;
  210. ResetEvent(overlapped.hEvent);
  211. ret = ReadFile(endpt->pipe_handle, buf + offset, left, &dwBytesTransfer, &overlapped);
  212. if (!ret && GetLastError() == ERROR_IO_PENDING) {
  213. ret = GetOverlappedResult(endpt->pipe_handle, &overlapped, &dwBytesTransfer, TRUE);
  214. }
  215. if (ret && dwBytesTransfer) {
  216. offset += dwBytesTransfer;
  217. left -= dwBytesTransfer;
  218. }
  219. else {
  220. return -1;
  221. }
  222. }
  223. return 0;
  224. }
  225. static int tcp_recv_buf(bus_endpt_t* endpt, char* buf, DWORD n)
  226. {
  227. DWORD left = n;
  228. DWORD offset = 0;
  229. while (left > 0) {
  230. BOOL ret;
  231. DWORD dwFlags = 0;
  232. WSABUF wsabuf;
  233. DWORD dwBytesTransfer;
  234. OVERLAPPED overlapped;
  235. memset(&overlapped, 0, sizeof(overlapped));
  236. overlapped.hEvent = endpt->rx_evt;
  237. wsabuf.buf = buf + offset;
  238. wsabuf.len = left;
  239. ResetEvent(overlapped.hEvent);
  240. ret = WSARecv(endpt->sock_handle, &wsabuf, 1, &dwBytesTransfer, &dwFlags, &overlapped, NULL) == 0;
  241. if (!ret && WSAGetLastError() == WSA_IO_PENDING) {
  242. ret = WSAGetOverlappedResult(endpt->sock_handle, &overlapped, &dwBytesTransfer, TRUE, &dwFlags);
  243. }
  244. if (ret && dwBytesTransfer) {
  245. offset += dwBytesTransfer;
  246. left -= dwBytesTransfer;
  247. }
  248. else {
  249. return -1;
  250. }
  251. }
  252. return 0;
  253. }
  254. static int recv_buf(bus_endpt_t* endpt, char* buf, DWORD n)
  255. {
  256. if (endpt->type == TYPE_PIPE) {
  257. return pipe_recv_buf(endpt, buf, n);
  258. }
  259. else if (endpt->type == TYPE_TCP) {
  260. return tcp_recv_buf(endpt, buf, n);
  261. }
  262. else {
  263. return -1;
  264. }
  265. }
  266. static int recv_pkt_raw(bus_endpt_t* endpt, iobuffer_t** pkt)
  267. {
  268. int pkt_len;
  269. int rc = -1;
  270. rc = recv_buf(endpt, (char*)&pkt_len, 4);
  271. if (rc != 0)
  272. return rc;
  273. *pkt = iobuffer_create(-1, pkt_len);
  274. iobuffer_push_count(*pkt, pkt_len);
  275. if (pkt_len > 0) {
  276. rc = recv_buf(endpt, iobuffer_data(*pkt, 0), pkt_len);
  277. }
  278. if (rc < 0) {
  279. iobuffer_destroy(*pkt);
  280. *pkt = NULL;
  281. }
  282. return rc;
  283. }
  284. static int start_read_pkt(bus_endpt_t* endpt, iobuffer_t** p_pkt)
  285. {
  286. DWORD dwBytesTransferred;
  287. BOOL ret;
  288. int rc = 0;
  289. iobuffer_t* pkt = NULL;
  290. *p_pkt = NULL;
  291. ResetEvent(endpt->rx_evt);
  292. memset(&endpt->rx_overlapped, 0, sizeof(OVERLAPPED));
  293. endpt->rx_overlapped.hEvent = endpt->rx_evt;
  294. if (endpt->type == TYPE_PIPE) {
  295. ret = ReadFile(endpt->pipe_handle,
  296. &endpt->rx_pending_pkt_len,
  297. 4,
  298. &dwBytesTransferred,
  299. &endpt->rx_overlapped);
  300. }
  301. else if (endpt->type == TYPE_TCP) {
  302. DWORD dwFlags = 0;
  303. WSABUF wsabuf;
  304. wsabuf.buf = (char*)&endpt->rx_pending_pkt_len;
  305. wsabuf.len = 4;
  306. ret = WSARecv(endpt->sock_handle,
  307. &wsabuf,
  308. 1,
  309. &dwBytesTransferred,
  310. &dwFlags,
  311. &endpt->rx_overlapped,
  312. NULL) == 0;
  313. }
  314. else {
  315. return -1;
  316. }
  317. if (ret) {
  318. if (dwBytesTransferred == 0)
  319. return -1;
  320. if (dwBytesTransferred < 4) {
  321. rc = recv_buf(endpt, (char*)&endpt->rx_pending_pkt_len + dwBytesTransferred, 4 - dwBytesTransferred);
  322. if (rc < 0)
  323. return rc;
  324. }
  325. pkt = iobuffer_create(0, endpt->rx_pending_pkt_len);
  326. if (endpt->rx_pending_pkt_len > 0) {
  327. rc = recv_buf(endpt, iobuffer_data(pkt, 0), endpt->rx_pending_pkt_len);
  328. if (rc < 0) {
  329. iobuffer_destroy(pkt);
  330. return rc;
  331. }
  332. iobuffer_push_count(pkt, endpt->rx_pending_pkt_len);
  333. }
  334. *p_pkt = pkt;
  335. }
  336. else {
  337. if (WSAGetLastError() == WSA_IO_PENDING) {
  338. endpt->rx_pending = 1;
  339. }
  340. else {
  341. return -1;
  342. }
  343. }
  344. return 0;
  345. }
  346. static int read_left_pkt(bus_endpt_t* endpt, iobuffer_t** p_pkt)
  347. {
  348. BOOL ret;
  349. int rc;
  350. DWORD dwBytesTransferred;
  351. iobuffer_t* pkt = NULL;
  352. if (endpt->type == TYPE_PIPE) {
  353. ret = GetOverlappedResult(endpt->pipe_handle, &endpt->rx_overlapped, &dwBytesTransferred, TRUE);
  354. }
  355. else if (endpt->type == TYPE_TCP) {
  356. DWORD dwFlags = 0;
  357. ret = WSAGetOverlappedResult(endpt->sock_handle, &endpt->rx_overlapped, &dwBytesTransferred, TRUE, &dwFlags);
  358. }
  359. if (!ret || dwBytesTransferred == 0) {
  360. DWORD dwError = GetLastError();
  361. return -1;
  362. }
  363. if (dwBytesTransferred < 4) {
  364. rc = recv_buf(endpt, (char*)&endpt->rx_pending_pkt_len + dwBytesTransferred, 4 - dwBytesTransferred);
  365. if (rc < 0)
  366. return rc;
  367. }
  368. pkt = iobuffer_create(-1, endpt->rx_pending_pkt_len);
  369. rc = recv_buf(endpt, iobuffer_data(pkt, 0), endpt->rx_pending_pkt_len);
  370. if (rc < 0) {
  371. iobuffer_destroy(pkt);
  372. return rc;
  373. }
  374. iobuffer_push_count(pkt, endpt->rx_pending_pkt_len);
  375. *p_pkt = pkt;
  376. endpt->rx_pending = 0;
  377. return 0;
  378. }
  379. static int append_rx_pkt(bus_endpt_t* endpt, iobuffer_t* pkt)
  380. {
  381. int type;
  382. int read_state;
  383. read_state = iobuffer_get_read_state(pkt);
  384. iobuffer_read(pkt, IOBUF_T_I4, &type, 0);
  385. iobuffer_restore_read_state(pkt, read_state);
  386. if (type == BUS_TYPE_PACKET || type == BUS_TYPE_INFO || type == BUS_TYPE_EVENT || type == BUS_TYPE_SYSTEM) {
  387. iobuffer_queue_enqueue(endpt->rx_buf_queue, pkt);
  388. return 1;
  389. }
  390. else {
  391. return -1;
  392. }
  393. }
  394. TOOLKIT_API int bus_endpt_create(const char* url, int epid, const bus_endpt_callback* callback, bus_endpt_t** p_endpt)
  395. {
  396. bus_endpt_t* endpt = NULL;
  397. char* tmp_url;
  398. url_fields uf;
  399. int rc;
  400. int v;
  401. iobuffer_t* buf = NULL;
  402. iobuffer_t* ans_buf = NULL;
  403. if (!url)
  404. return -1;
  405. tmp_url = _strdup(url);
  406. if (url_parse(tmp_url, &uf) < 0) {
  407. free(tmp_url);
  408. return -1;
  409. }
  410. endpt = ZALLOC_T(bus_endpt_t);
  411. endpt->sock_handle = -1;
  412. endpt->url = tmp_url;
  413. if (_stricmp(uf.scheme, "tcp") == 0) {
  414. endpt->type = TYPE_TCP;
  415. endpt->sock_handle = create_socket_handle(uf.host, uf.port);
  416. if (endpt->sock_handle == INVALID_SOCKET)
  417. goto on_error;
  418. }
  419. else if (_stricmp(uf.scheme, "pipe") == 0) {
  420. endpt->type = TYPE_PIPE;
  421. endpt->pipe_handle = create_pipe_handle(uf.host);
  422. if (endpt->pipe_handle == INVALID_HANDLE_VALUE)
  423. goto on_error;
  424. }
  425. else {
  426. goto on_error;
  427. }
  428. endpt->epid = epid;
  429. endpt->tx_evt = CreateEventA(NULL, TRUE, FALSE, NULL);
  430. endpt->rx_evt = CreateEventA(NULL, TRUE, FALSE, NULL);
  431. endpt->rx_buf_queue = iobuffer_queue_create();
  432. endpt->msg_sem = CreateSemaphoreA(NULL, 0, 0x7fffffff, NULL);
  433. INIT_LIST_HEAD(&endpt->msg_list);
  434. spinlock_init(&endpt->msg_lock);
  435. memcpy(&endpt->callback, callback, sizeof(bus_endpt_callback));
  436. buf = iobuffer_create(-1, -1);
  437. v = BUS_TYPE_ENDPT_REGISTER;
  438. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  439. v = endpt->epid;
  440. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  441. rc = send_pkt_raw(endpt, buf);
  442. if (rc != 0)
  443. goto on_error;
  444. rc = recv_pkt_raw(endpt, &ans_buf);
  445. if (rc != 0) {
  446. DWORD dwError = GetLastError();
  447. goto on_error;
  448. }
  449. iobuffer_read(ans_buf, IOBUF_T_I4, &v, 0);
  450. iobuffer_read(ans_buf, IOBUF_T_I4, &rc, 0);
  451. if (rc != 0)
  452. goto on_error;
  453. url_free_fields(&uf);
  454. if (buf)
  455. iobuffer_destroy(buf);
  456. if (ans_buf)
  457. iobuffer_destroy(ans_buf);
  458. *p_endpt = endpt;
  459. return 0;
  460. on_error:
  461. if (endpt->type == TYPE_TCP) {
  462. closesocket(endpt->sock_handle);
  463. }
  464. else if (endpt->type == TYPE_PIPE) {
  465. CloseHandle(endpt->pipe_handle);
  466. }
  467. if (endpt->msg_sem)
  468. CloseHandle(endpt->msg_sem);
  469. if (endpt->tx_evt)
  470. CloseHandle(endpt->tx_evt);
  471. if (endpt->rx_evt)
  472. CloseHandle(endpt->rx_evt);
  473. if (endpt->rx_buf_queue)
  474. iobuffer_queue_destroy(endpt->rx_buf_queue);
  475. if (endpt->url)
  476. free(endpt->url);
  477. free(endpt);
  478. url_free_fields(&uf);
  479. if (buf)
  480. iobuffer_destroy(buf);
  481. if (ans_buf)
  482. iobuffer_destroy(ans_buf);
  483. return -1;
  484. }
  485. TOOLKIT_API void bus_endpt_destroy(bus_endpt_t* endpt)
  486. {
  487. int rc = -1;
  488. iobuffer_t* buf = NULL;
  489. iobuffer_t* ans_buf = NULL;
  490. int v;
  491. assert(endpt);
  492. buf = iobuffer_create(-1, -1);
  493. v = BUS_TYPE_ENDPT_UNREGISTER;
  494. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  495. v = endpt->epid;
  496. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  497. rc = send_pkt_raw(endpt, buf);
  498. if (rc != 0)
  499. goto on_error;
  500. rc = recv_pkt_raw(endpt, &ans_buf);
  501. if (rc != 0)
  502. goto on_error;
  503. iobuffer_read(ans_buf, IOBUF_T_I4, &v, 0);
  504. iobuffer_read(ans_buf, IOBUF_T_I4, &rc, 0);
  505. {
  506. msg_t* msg, * t;
  507. list_for_each_entry_safe(msg, t, &endpt->msg_list, msg_t, entry) {
  508. list_del(&msg->entry);
  509. if (msg->evt)
  510. CloseHandle(msg->evt);
  511. free(msg);
  512. }
  513. }
  514. on_error:
  515. if (buf)
  516. iobuffer_destroy(buf);
  517. if (ans_buf)
  518. iobuffer_destroy(ans_buf);
  519. if (endpt->type == TYPE_TCP) {
  520. closesocket(endpt->sock_handle);
  521. }
  522. else if (endpt->type == TYPE_PIPE) {
  523. CloseHandle(endpt->pipe_handle);
  524. }
  525. if (endpt->msg_sem)
  526. CloseHandle(endpt->msg_sem);
  527. if (endpt->tx_evt)
  528. CloseHandle(endpt->tx_evt);
  529. if (endpt->rx_evt)
  530. CloseHandle(endpt->rx_evt);
  531. if (endpt->rx_buf_queue)
  532. iobuffer_queue_destroy(endpt->rx_buf_queue);
  533. if (endpt->url)
  534. free(endpt->url);
  535. free(endpt);
  536. }
  537. // 1 : recv ok
  538. // 0 : time out
  539. // <0 : error
  540. static int bus_endpt_poll_internal(bus_endpt_t* endpt, int* result, int timeout)
  541. {
  542. iobuffer_t* pkt = NULL;
  543. int rc;
  544. BOOL ret;
  545. assert(endpt);
  546. // peek first packge type
  547. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  548. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  549. int pkt_type;
  550. int read_state = iobuffer_get_read_state(pkt);
  551. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, NULL);
  552. iobuffer_restore_read_state(pkt, read_state);
  553. *result = to_result(pkt_type);
  554. if (*result == BUS_RESULT_UNKNOWN) {
  555. OutputDebugStringA("bug: unknown pkt type!\n");
  556. return -1;
  557. }
  558. return 1;
  559. }
  560. // no received package, try to receive one
  561. if (!endpt->rx_pending) {
  562. rc = start_read_pkt(endpt, &pkt);
  563. if (rc < 0)
  564. return rc;
  565. if (pkt) {
  566. OutputDebugStringA("pkt has read\n");
  567. rc = append_rx_pkt(endpt, pkt);
  568. if (rc < 0) {
  569. iobuffer_destroy(pkt);
  570. return -1;
  571. }
  572. }
  573. else {
  574. OutputDebugStringA("pending\n");
  575. }
  576. }
  577. // if receive is pending, wait for send or receive complete event
  578. if (!pkt) {
  579. HANDLE hs[] = { endpt->msg_sem, endpt->rx_evt };
  580. ret = WaitForMultipleObjects(array_size(hs), &hs[0], FALSE, (DWORD)timeout);
  581. if (ret == WAIT_TIMEOUT) {
  582. return 0;
  583. }
  584. else if (ret == WAIT_OBJECT_0) {
  585. *result = BUS_RESULT_MSG; // indicate send package event
  586. return 1;
  587. }
  588. else if (ret == WAIT_OBJECT_0 + 1) {
  589. rc = read_left_pkt(endpt, &pkt);
  590. if (rc < 0)
  591. return rc;
  592. if (pkt) {
  593. rc = append_rx_pkt(endpt, pkt);
  594. if (rc < 0) {
  595. iobuffer_destroy(pkt);
  596. return -1;
  597. }
  598. }
  599. }
  600. }
  601. else {
  602. OutputDebugStringA("pkt has readed\n");
  603. }
  604. if (pkt) {
  605. int type;
  606. int read_state = iobuffer_get_read_state(pkt);
  607. iobuffer_read(pkt, IOBUF_T_I4, &type, 0);
  608. iobuffer_restore_read_state(pkt, read_state);
  609. *result = to_result(type);
  610. if (*result == BUS_RESULT_UNKNOWN) {
  611. OutputDebugStringA("bug: unknown pkt type!\n");
  612. return -1;
  613. }
  614. }
  615. else {
  616. return -1;
  617. }
  618. return 1;
  619. }
  620. static int recv_until(bus_endpt_t* endpt, int type, iobuffer_t** p_ansbuf)
  621. {
  622. int rc;
  623. iobuffer_t* ans_pkt = NULL;
  624. int ans_type;
  625. for (;;) {
  626. if (!endpt->rx_pending) {
  627. rc = start_read_pkt(endpt, &ans_pkt);
  628. if (rc < 0) {
  629. DWORD dwError = WSAGetLastError();
  630. break;
  631. }
  632. }
  633. if (!ans_pkt) {
  634. DWORD ret = WaitForSingleObject(endpt->rx_evt, INFINITE);
  635. if (ret != WAIT_OBJECT_0)
  636. return -1;
  637. rc = read_left_pkt(endpt, &ans_pkt);
  638. if (rc < 0)
  639. break;
  640. }
  641. if (ans_pkt) {
  642. int read_state = iobuffer_get_read_state(ans_pkt);
  643. iobuffer_read(ans_pkt, IOBUF_T_I4, &ans_type, 0);
  644. iobuffer_restore_read_state(ans_pkt, read_state);
  645. if (ans_type == type) {
  646. *p_ansbuf = ans_pkt;
  647. break;
  648. }
  649. else {
  650. rc = append_rx_pkt(endpt, ans_pkt);
  651. if (rc < 0) {
  652. iobuffer_destroy(ans_pkt);
  653. break;
  654. }
  655. else {
  656. ans_pkt = NULL;
  657. }
  658. }
  659. }
  660. }
  661. return rc;
  662. }
  663. static int recv_until_result(bus_endpt_t* endpt, int* p_result)
  664. {
  665. int rc;
  666. iobuffer_t* ans_pkt = NULL;
  667. int type, error;
  668. rc = recv_until(endpt, BUS_TYPE_ERROR, &ans_pkt);
  669. if (rc < 0)
  670. return rc;
  671. iobuffer_read(ans_pkt, IOBUF_T_I4, &type, 0);
  672. iobuffer_read(ans_pkt, IOBUF_T_I4, &error, 0);
  673. iobuffer_destroy(ans_pkt);
  674. *p_result = error;
  675. return rc;
  676. }
  677. static int recv_until_state(bus_endpt_t* endpt, int* p_state)
  678. {
  679. int rc;
  680. iobuffer_t* ans_pkt = NULL;
  681. int type, epid, state;
  682. rc = recv_until(endpt, BUS_TYPE_ENDPT_GET_STATE, &ans_pkt);
  683. if (rc < 0)
  684. return rc;
  685. iobuffer_read(ans_pkt, IOBUF_T_I4, &type, 0);
  686. iobuffer_read(ans_pkt, IOBUF_T_I4, &epid, 0);
  687. iobuffer_read(ans_pkt, IOBUF_T_I4, &state, 0);
  688. iobuffer_destroy(ans_pkt);
  689. *p_state = state;
  690. return rc;
  691. }
  692. TOOLKIT_API int bus_endpt_send_pkt(bus_endpt_t* endpt, int epid, int type, iobuffer_t* pkt)
  693. {
  694. int t;
  695. int rc;
  696. int read_state;
  697. int write_state;
  698. int error;
  699. assert(endpt);
  700. read_state = iobuffer_get_read_state(pkt);
  701. write_state = iobuffer_get_write_state(pkt);
  702. t = epid; // remote epid
  703. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  704. t = endpt->epid; // local epid
  705. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  706. t = type; // user type
  707. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  708. t = BUS_TYPE_PACKET; // type
  709. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  710. rc = send_pkt_raw(endpt, pkt);
  711. iobuffer_restore_read_state(pkt, read_state);
  712. iobuffer_restore_write_state(pkt, write_state);
  713. if (rc < 0)
  714. return rc;
  715. rc = recv_until_result(endpt, &error);
  716. if (rc == 0 && error != 0)
  717. rc = error;
  718. return rc;
  719. }
  720. TOOLKIT_API int bus_endpt_send_info(bus_endpt_t* endpt, int epid, int type, iobuffer_t* pkt)
  721. {
  722. int t;
  723. int rc;
  724. int read_state;
  725. int write_state;
  726. assert(endpt);
  727. read_state = iobuffer_get_read_state(pkt);
  728. write_state = iobuffer_get_write_state(pkt);
  729. t = epid; // remote epid
  730. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  731. t = endpt->epid; // local epid
  732. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  733. t = type; // user type
  734. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  735. t = BUS_TYPE_INFO; // type
  736. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  737. rc = send_pkt_raw(endpt, pkt);
  738. iobuffer_restore_read_state(pkt, read_state);
  739. iobuffer_restore_write_state(pkt, write_state);
  740. return rc;
  741. }
  742. TOOLKIT_API int bus_endpt_bcast_evt(bus_endpt_t* endpt, int type, iobuffer_t* pkt)
  743. {
  744. int t;
  745. int rc;
  746. int read_state;
  747. int write_state;
  748. assert(endpt);
  749. read_state = iobuffer_get_read_state(pkt);
  750. write_state = iobuffer_get_write_state(pkt);
  751. t = endpt->epid;
  752. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  753. t = type;
  754. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  755. t = BUS_TYPE_EVENT;
  756. iobuffer_write_head(pkt, IOBUF_T_I4, &t, 0);
  757. rc = send_pkt_raw(endpt, pkt);
  758. iobuffer_restore_read_state(pkt, read_state);
  759. iobuffer_restore_write_state(pkt, write_state);
  760. return rc;
  761. }
  762. static int bus_endpt_recv_pkt(bus_endpt_t* endpt, int* p_epid, int* p_type, iobuffer_t** p_pkt)
  763. {
  764. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  765. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  766. int read_state = iobuffer_get_read_state(pkt);
  767. int pkt_type, usr_type, from_epid, to_epid;
  768. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, 0);
  769. if (pkt_type == BUS_TYPE_PACKET || pkt_type == BUS_TYPE_INFO) {
  770. iobuffer_read(pkt, IOBUF_T_I4, &usr_type, 0);
  771. iobuffer_read(pkt, IOBUF_T_I4, &from_epid, 0);
  772. iobuffer_read(pkt, IOBUF_T_I4, &to_epid, 0);
  773. if (p_epid)
  774. *p_epid = from_epid;
  775. if (p_type)
  776. *p_type = usr_type;
  777. iobuffer_queue_deque(endpt->rx_buf_queue);
  778. if (p_pkt) {
  779. *p_pkt = pkt;
  780. }
  781. else {
  782. iobuffer_destroy(pkt);
  783. }
  784. return 0;
  785. }
  786. else {
  787. iobuffer_restore_read_state(pkt, read_state);
  788. }
  789. }
  790. return -1;
  791. }
  792. static int bus_endpt_recv_evt(bus_endpt_t* endpt, int* p_epid, int* p_type, iobuffer_t** p_pkt)
  793. {
  794. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  795. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  796. int read_state = iobuffer_get_read_state(pkt);
  797. int pkt_type, usr_type, from_epid;
  798. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, 0);
  799. if (pkt_type == BUS_TYPE_EVENT) {
  800. iobuffer_read(pkt, IOBUF_T_I4, &usr_type, 0);
  801. iobuffer_read(pkt, IOBUF_T_I4, &from_epid, 0);
  802. if (p_epid)
  803. *p_epid = from_epid;
  804. if (p_type)
  805. *p_type = usr_type;
  806. iobuffer_queue_deque(endpt->rx_buf_queue);
  807. if (p_pkt) {
  808. *p_pkt = pkt;
  809. }
  810. else {
  811. iobuffer_destroy(pkt);
  812. }
  813. return 0;
  814. }
  815. else {
  816. iobuffer_restore_read_state(pkt, read_state);
  817. }
  818. }
  819. return -1;
  820. }
  821. static int bus_endpt_recv_sys(bus_endpt_t* endpt, int* p_epid, int* p_state)
  822. {
  823. if (iobuffer_queue_count(endpt->rx_buf_queue) > 0) {
  824. iobuffer_t* pkt = iobuffer_queue_head(endpt->rx_buf_queue);
  825. int read_state = iobuffer_get_read_state(pkt);
  826. int pkt_type, epid, state;
  827. iobuffer_read(pkt, IOBUF_T_I4, &pkt_type, 0);
  828. if (pkt_type == BUS_TYPE_SYSTEM) {
  829. iobuffer_read(pkt, IOBUF_T_I4, &epid, 0);
  830. iobuffer_read(pkt, IOBUF_T_I4, &state, 0);
  831. if (p_epid)
  832. *p_epid = epid;
  833. if (p_state)
  834. *p_state = state;
  835. iobuffer_queue_deque(endpt->rx_buf_queue);
  836. iobuffer_destroy(pkt);
  837. return 0;
  838. }
  839. else {
  840. iobuffer_restore_read_state(pkt, read_state);
  841. }
  842. }
  843. return -1;
  844. }
  845. static int bus_endpt_recv_msg(bus_endpt_t* endpt, msg_t** p_msg)
  846. {
  847. int rc = -1;
  848. assert(endpt);
  849. assert(p_msg);
  850. spinlock_enter(&endpt->msg_lock, -1);
  851. if (!list_empty(&endpt->msg_list)) {
  852. msg_t* e = list_first_entry(&endpt->msg_list, msg_t, entry);
  853. list_del(&e->entry);
  854. rc = 0;
  855. *p_msg = e;
  856. }
  857. spinlock_leave(&endpt->msg_lock);
  858. return rc;
  859. }
  860. TOOLKIT_API int bus_endpt_get_state(bus_endpt_t* endpt, int epid, int* p_state)
  861. {
  862. iobuffer_t* buf = NULL;
  863. int v;
  864. int rc = -1;
  865. assert(endpt);
  866. buf = iobuffer_create(-1, -1);
  867. v = BUS_TYPE_ENDPT_GET_STATE;
  868. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  869. v = epid;
  870. iobuffer_write(buf, IOBUF_T_I4, &v, 0);
  871. rc = send_pkt_raw(endpt, buf);
  872. if (rc < 0)
  873. goto on_error;
  874. rc = recv_until_state(endpt, p_state);
  875. on_error:
  876. if (buf)
  877. iobuffer_destroy(buf);
  878. return rc;
  879. }
  880. TOOLKIT_API int bus_endpt_post_msg(bus_endpt_t* endpt, int msg, int nparam, int params[])
  881. {
  882. msg_t* e;
  883. assert(endpt);
  884. e = MALLOC_T(msg_t);
  885. e->type = msg;
  886. e->nparam = nparam;
  887. if (nparam) {
  888. e->params = (int*)malloc(sizeof(int) * nparam);
  889. memcpy(e->params, params, sizeof(int) * nparam);
  890. }
  891. else {
  892. e->params = NULL;
  893. }
  894. e->evt = NULL;
  895. spinlock_enter(&endpt->msg_lock, -1);
  896. list_add_tail(&e->entry, &endpt->msg_list);
  897. spinlock_leave(&endpt->msg_lock);
  898. ReleaseSemaphore(endpt->msg_sem, 1, NULL);
  899. return 0;
  900. }
  901. TOOLKIT_API int bus_endpt_send_msg(bus_endpt_t* endpt, int msg, int nparam, int params[])
  902. {
  903. msg_t e;
  904. assert(endpt);
  905. e.type = msg;
  906. e.nparam = nparam;
  907. if (nparam) {
  908. e.params = (int*)malloc(sizeof(int) * nparam);
  909. memcpy(e.params, params, sizeof(int) * nparam);
  910. }
  911. else {
  912. e.params = NULL;
  913. }
  914. e.evt_result = 0;
  915. e.evt = CreateEventA(NULL, TRUE, FALSE, NULL);
  916. spinlock_enter(&endpt->msg_lock, -1);
  917. list_add_tail(&e.entry, &endpt->msg_list);
  918. spinlock_leave(&endpt->msg_lock);
  919. ReleaseSemaphore(endpt->msg_sem, 1, NULL);
  920. WaitForSingleObject(e.evt, INFINITE);
  921. CloseHandle(e.evt);
  922. if (nparam) {
  923. free(e.params);
  924. }
  925. return e.evt_result;
  926. }
  927. TOOLKIT_API int bus_endpt_get_epid(bus_endpt_t* endpt)
  928. {
  929. return endpt->epid;
  930. }
  931. TOOLKIT_API const char* bus_endpt_get_url(bus_endpt_t* endpt)
  932. {
  933. return endpt->url;
  934. }
  935. TOOLKIT_API int bus_endpt_poll(bus_endpt_t* endpt, int timeout)
  936. {
  937. int result;
  938. int rc;
  939. int epid, type, state;
  940. iobuffer_t* pkt = NULL;
  941. rc = bus_endpt_poll_internal(endpt, &result, timeout);
  942. if (rc > 0) {
  943. if (result == BUS_RESULT_DATA) {
  944. bus_endpt_recv_pkt(endpt, &epid, &type, &pkt);
  945. if (endpt->callback.on_pkt)
  946. endpt->callback.on_pkt(endpt, epid, type, &pkt, endpt->callback.user_data);
  947. if (pkt)
  948. iobuffer_dec_ref(pkt);
  949. }
  950. else if (result == BUS_RESULT_INFO) {
  951. bus_endpt_recv_pkt(endpt, &epid, &type, &pkt);
  952. if (endpt->callback.on_inf)
  953. endpt->callback.on_inf(endpt, epid, type, &pkt, endpt->callback.user_data);
  954. if (pkt)
  955. iobuffer_dec_ref(pkt);
  956. }
  957. else if (result == BUS_RESULT_EVENT) {
  958. bus_endpt_recv_evt(endpt, &epid, &type, &pkt);
  959. if (endpt->callback.on_evt)
  960. endpt->callback.on_evt(endpt, epid, type, &pkt, endpt->callback.user_data);
  961. if (pkt)
  962. iobuffer_dec_ref(pkt);
  963. }
  964. else if (result == BUS_RESULT_SYSTEM) {
  965. bus_endpt_recv_sys(endpt, &epid, &state);
  966. if (endpt->callback.on_sys)
  967. endpt->callback.on_sys(endpt, epid, state, endpt->callback.user_data);
  968. }
  969. else if (result == BUS_RESULT_MSG) {
  970. msg_t* msg = NULL;
  971. bus_endpt_recv_msg(endpt, &msg);
  972. if (endpt->callback.on_msg) {
  973. endpt->callback.on_msg(endpt,
  974. msg->type,
  975. msg->nparam,
  976. msg->params,
  977. msg->evt ? &msg->evt_result : NULL,
  978. endpt->callback.user_data);
  979. if (msg->evt)
  980. SetEvent(msg->evt);
  981. else
  982. free_msg(msg);
  983. }
  984. else {
  985. if (msg->evt) {
  986. msg->evt_result = -1;
  987. SetEvent(msg->evt);
  988. }
  989. else {
  990. free_msg(msg);
  991. }
  992. }
  993. }
  994. else {
  995. assert(0);
  996. rc = -1;
  997. }
  998. }
  999. return rc;
  1000. }
  1001. TOOLKIT_API int bus_endpt_set_quit_flag(bus_endpt_t* endpt)
  1002. {
  1003. endpt->quit_flag = 1;
  1004. return 0;
  1005. }
  1006. TOOLKIT_API int bus_endpt_get_quit_flag(bus_endpt_t* endpt)
  1007. {
  1008. return endpt->quit_flag;
  1009. }