xref: /freebsd/contrib/unbound/util/tube.c (revision 7a789145f88a6aceacc59029a0cafe7de7aeefea)
1 /*
2  * util/tube.c - pipe service
3  *
4  * Copyright (c) 2008, NLnet Labs. All rights reserved.
5  *
6  * This software is open source.
7  *
8  * Redistribution and use in source and binary forms, with or without
9  * modification, are permitted provided that the following conditions
10  * are met:
11  *
12  * Redistributions of source code must retain the above copyright notice,
13  * this list of conditions and the following disclaimer.
14  *
15  * Redistributions in binary form must reproduce the above copyright notice,
16  * this list of conditions and the following disclaimer in the documentation
17  * and/or other materials provided with the distribution.
18  *
19  * Neither the name of the NLNET LABS nor the names of its contributors may
20  * be used to endorse or promote products derived from this software without
21  * specific prior written permission.
22  *
23  * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
24  * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
25  * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
26  * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
27  * HOLDER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
28  * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED
29  * TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
30  * PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
31  * LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
32  * NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
33  * SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
34  */
35 
36 /**
37  * \file
38  *
39  * This file contains pipe service functions.
40  */
41 #include "config.h"
42 #include "util/tube.h"
43 #include "util/log.h"
44 #include "util/net_help.h"
45 #include "util/netevent.h"
46 #include "util/fptr_wlist.h"
47 #include "util/ub_event.h"
48 #ifdef HAVE_POLL_H
49 #include <poll.h>
50 #endif
51 
52 #ifndef USE_WINSOCK
53 /* on unix */
54 
55 #ifndef HAVE_SOCKETPAIR
56 /** no socketpair() available, like on Minix 3.1.7, use pipe */
57 #define socketpair(f, t, p, sv) pipe(sv)
58 #endif /* HAVE_SOCKETPAIR */
59 
tube_create(void)60 struct tube* tube_create(void)
61 {
62 	struct tube* tube = (struct tube*)calloc(1, sizeof(*tube));
63 	int sv[2];
64 	if(!tube) {
65 		int err = errno;
66 		log_err("tube_create: out of memory");
67 		errno = err;
68 		return NULL;
69 	}
70 	tube->sr = -1;
71 	tube->sw = -1;
72 	if(socketpair(AF_UNIX, SOCK_STREAM, 0, sv) == -1) {
73 		int err = errno;
74 		log_err("socketpair: %s", strerror(errno));
75 		free(tube);
76 		errno = err;
77 		return NULL;
78 	}
79 	tube->sr = sv[0];
80 	tube->sw = sv[1];
81 	if(!fd_set_nonblock(tube->sr) || !fd_set_nonblock(tube->sw)) {
82 		int err = errno;
83 		log_err("tube: cannot set nonblocking");
84 		tube_delete(tube);
85 		errno = err;
86 		return NULL;
87 	}
88 	return tube;
89 }
90 
tube_delete(struct tube * tube)91 void tube_delete(struct tube* tube)
92 {
93 	if(!tube) return;
94 	tube_remove_bg_listen(tube);
95 	tube_remove_bg_write(tube);
96 	/* close fds after deleting commpoints, to be sure.
97 	 *            Also epoll does not like closing fd before event_del */
98 	tube_close_read(tube);
99 	tube_close_write(tube);
100 	free(tube);
101 }
102 
tube_close_read(struct tube * tube)103 void tube_close_read(struct tube* tube)
104 {
105 	if(tube->sr != -1) {
106 		close(tube->sr);
107 		tube->sr = -1;
108 	}
109 }
110 
tube_close_write(struct tube * tube)111 void tube_close_write(struct tube* tube)
112 {
113 	if(tube->sw != -1) {
114 		close(tube->sw);
115 		tube->sw = -1;
116 	}
117 }
118 
tube_remove_bg_listen(struct tube * tube)119 void tube_remove_bg_listen(struct tube* tube)
120 {
121 	if(tube->listen_com) {
122 		comm_point_delete(tube->listen_com);
123 		tube->listen_com = NULL;
124 	}
125 	free(tube->cmd_msg);
126 	tube->cmd_msg = NULL;
127 }
128 
tube_remove_bg_write(struct tube * tube)129 void tube_remove_bg_write(struct tube* tube)
130 {
131 	if(tube->res_com) {
132 		comm_point_delete(tube->res_com);
133 		tube->res_com = NULL;
134 	}
135 	if(tube->res_list) {
136 		struct tube_res_list* np, *p = tube->res_list;
137 		tube->res_list = NULL;
138 		tube->res_last = NULL;
139 		while(p) {
140 			np = p->next;
141 			free(p->buf);
142 			free(p);
143 			p = np;
144 		}
145 	}
146 }
147 
148 /** Drain the pipe of bytes. */
149 static void
fd_drain(int fd,uint32_t len)150 fd_drain(int fd, uint32_t len)
151 {
152 	uint8_t discard[256];
153 	uint32_t remaining = len;
154 	while(remaining > 0) {
155 		ssize_t n = read(fd, discard,
156 			remaining < sizeof(discard) ? remaining : sizeof(discard));
157 		if(n <= 0) break;
158 		remaining -= (uint32_t)n;
159 	}
160 }
161 
162 int
tube_handle_listen(struct comm_point * c,void * arg,int error,struct comm_reply * ATTR_UNUSED (reply_info))163 tube_handle_listen(struct comm_point* c, void* arg, int error,
164         struct comm_reply* ATTR_UNUSED(reply_info))
165 {
166 	struct tube* tube = (struct tube*)arg;
167 	ssize_t r;
168 	if(error != NETEVENT_NOERROR) {
169 		fptr_ok(fptr_whitelist_tube_listen(tube->listen_cb));
170 		(*tube->listen_cb)(tube, NULL, 0, error, tube->listen_arg);
171 		return 0;
172 	}
173 
174 	if(tube->cmd_read < sizeof(tube->cmd_len)) {
175 		/* complete reading the length of control msg */
176 		r = read(c->fd, ((uint8_t*)&tube->cmd_len) + tube->cmd_read,
177 			sizeof(tube->cmd_len) - tube->cmd_read);
178 		if(r==0) {
179 			/* error has happened or */
180 			/* parent closed pipe, must have exited somehow */
181 			fptr_ok(fptr_whitelist_tube_listen(tube->listen_cb));
182 			(*tube->listen_cb)(tube, NULL, 0, NETEVENT_CLOSED,
183 				tube->listen_arg);
184 			return 0;
185 		}
186 		if(r==-1) {
187 			if(errno != EAGAIN && errno != EINTR) {
188 				log_err("rpipe error: %s", strerror(errno));
189 			}
190 			/* nothing to read now, try later */
191 			return 0;
192 		}
193 		tube->cmd_read += r;
194 		if(tube->cmd_read < sizeof(tube->cmd_len)) {
195 			/* not complete, try later */
196 			return 0;
197 		}
198 		tube->cmd_msg = (uint8_t*)calloc(1, tube->cmd_len);
199 		if(!tube->cmd_msg) {
200 			log_err("malloc failure");
201 			/* Drain the remaining bytes, since they belong to this
202 			 * message. The next message starts after it. */
203 			fd_drain(c->fd, tube->cmd_len);
204 			tube->cmd_read = 0;
205 			return 0;
206 		}
207 	}
208 	/* cmd_len has been read, read remainder */
209 	r = read(c->fd, tube->cmd_msg+tube->cmd_read-sizeof(tube->cmd_len),
210 		tube->cmd_len - (tube->cmd_read - sizeof(tube->cmd_len)));
211 	if(r==0) {
212 		/* error has happened or */
213 		/* parent closed pipe, must have exited somehow */
214 		fptr_ok(fptr_whitelist_tube_listen(tube->listen_cb));
215 		(*tube->listen_cb)(tube, NULL, 0, NETEVENT_CLOSED,
216 			tube->listen_arg);
217 		return 0;
218 	}
219 	if(r==-1) {
220 		/* nothing to read now, try later */
221 		if(errno != EAGAIN && errno != EINTR) {
222 			log_err("rpipe error: %s", strerror(errno));
223 		}
224 		return 0;
225 	}
226 	tube->cmd_read += r;
227 	if(tube->cmd_read < sizeof(tube->cmd_len) + tube->cmd_len) {
228 		/* not complete, try later */
229 		return 0;
230 	}
231 	tube->cmd_read = 0;
232 
233 	fptr_ok(fptr_whitelist_tube_listen(tube->listen_cb));
234 	(*tube->listen_cb)(tube, tube->cmd_msg, tube->cmd_len,
235 		NETEVENT_NOERROR, tube->listen_arg);
236 		/* also frees the buf */
237 	tube->cmd_msg = NULL;
238 	return 0;
239 }
240 
241 int
tube_handle_write(struct comm_point * c,void * arg,int error,struct comm_reply * ATTR_UNUSED (reply_info))242 tube_handle_write(struct comm_point* c, void* arg, int error,
243         struct comm_reply* ATTR_UNUSED(reply_info))
244 {
245 	struct tube* tube = (struct tube*)arg;
246 	struct tube_res_list* item = tube->res_list;
247 	ssize_t r;
248 	if(error != NETEVENT_NOERROR) {
249 		log_err("tube_handle_write net error %d", error);
250 		return 0;
251 	}
252 
253 	if(!item) {
254 		comm_point_stop_listening(c);
255 		return 0;
256 	}
257 
258 	if(tube->res_write < sizeof(item->len)) {
259 		r = write(c->fd, ((uint8_t*)&item->len) + tube->res_write,
260 			sizeof(item->len) - tube->res_write);
261 		if(r == -1) {
262 			if(errno != EAGAIN && errno != EINTR) {
263 				log_err("wpipe error: %s", strerror(errno));
264 			}
265 			return 0; /* try again later */
266 		}
267 		if(r == 0) {
268 			/* error on pipe, must have exited somehow */
269 			/* cannot signal this to pipe user */
270 			return 0;
271 		}
272 		tube->res_write += r;
273 		if(tube->res_write < sizeof(item->len))
274 			return 0;
275 	}
276 	r = write(c->fd, item->buf + tube->res_write - sizeof(item->len),
277 		item->len - (tube->res_write - sizeof(item->len)));
278 	if(r == -1) {
279 		if(errno != EAGAIN && errno != EINTR) {
280 			log_err("wpipe error: %s", strerror(errno));
281 		}
282 		return 0; /* try again later */
283 	}
284 	if(r == 0) {
285 		/* error on pipe, must have exited somehow */
286 		/* cannot signal this to pipe user */
287 		return 0;
288 	}
289 	tube->res_write += r;
290 	if(tube->res_write < sizeof(item->len) + item->len)
291 		return 0;
292 	/* done this result, remove it */
293 	free(item->buf);
294 	item->buf = NULL;
295 	tube->res_list = tube->res_list->next;
296 	free(item);
297 	if(!tube->res_list) {
298 		tube->res_last = NULL;
299 		comm_point_stop_listening(c);
300 	}
301 	tube->res_write = 0;
302 	return 0;
303 }
304 
tube_write_msg(struct tube * tube,uint8_t * buf,uint32_t len,int nonblock)305 int tube_write_msg(struct tube* tube, uint8_t* buf, uint32_t len,
306         int nonblock)
307 {
308 	ssize_t r, d;
309 	int fd = tube->sw;
310 
311 	/* test */
312 	if(nonblock) {
313 		r = write(fd, &len, sizeof(len));
314 		if(r == -1) {
315 			if(errno==EINTR || errno==EAGAIN)
316 				return -1;
317 			log_err("tube msg write failed: %s", strerror(errno));
318 			return -1; /* can still continue, perhaps */
319 		}
320 	} else r = 0;
321 	if(!fd_set_block(fd))
322 		return 0;
323 	/* write remainder */
324 	d = r;
325 	while(d != (ssize_t)sizeof(len)) {
326 		if((r=write(fd, ((char*)&len)+d, sizeof(len)-d)) == -1) {
327 			if(errno == EAGAIN)
328 				continue; /* temporarily unavail: try again*/
329 			log_err("tube msg write failed: %s", strerror(errno));
330 			(void)fd_set_nonblock(fd);
331 			return 0;
332 		}
333 		d += r;
334 	}
335 	d = 0;
336 	while(d != (ssize_t)len) {
337 		if((r=write(fd, buf+d, len-d)) == -1) {
338 			if(errno == EAGAIN)
339 				continue; /* temporarily unavail: try again*/
340 			log_err("tube msg write failed: %s", strerror(errno));
341 			(void)fd_set_nonblock(fd);
342 			return 0;
343 		}
344 		d += r;
345 	}
346 	if(!fd_set_nonblock(fd))
347 		return 0;
348 	return 1;
349 }
350 
tube_read_msg(struct tube * tube,uint8_t ** buf,uint32_t * len,int nonblock)351 int tube_read_msg(struct tube* tube, uint8_t** buf, uint32_t* len,
352         int nonblock)
353 {
354 	ssize_t r, d;
355 	int fd = tube->sr;
356 
357 	/* test */
358 	*len = 0;
359 	if(nonblock) {
360 		r = read(fd, len, sizeof(*len));
361 		if(r == -1) {
362 			if(errno==EINTR || errno==EAGAIN)
363 				return -1;
364 			log_err("tube msg read failed: %s", strerror(errno));
365 			return -1; /* we can still continue, perhaps */
366 		}
367 		if(r == 0) /* EOF */
368 			return 0;
369 	} else r = 0;
370 	if(!fd_set_block(fd))
371 		return 0;
372 	/* read remainder */
373 	d = r;
374 	while(d != (ssize_t)sizeof(*len)) {
375 		if((r=read(fd, ((char*)len)+d, sizeof(*len)-d)) == -1) {
376 			log_err("tube msg read failed: %s", strerror(errno));
377 			(void)fd_set_nonblock(fd);
378 			return 0;
379 		}
380 		if(r == 0) /* EOF */ {
381 			(void)fd_set_nonblock(fd);
382 			return 0;
383 		}
384 		d += r;
385 	}
386 	if (*len >= 65536*2) {
387 		log_err("tube msg length %u is too big", (unsigned)*len);
388 		(void)fd_set_nonblock(fd);
389 		return 0;
390 	}
391 	*buf = (uint8_t*)malloc(*len);
392 	if(!*buf) {
393 		log_err("tube read out of memory");
394 		/* Drain the remaining bytes, since they belong to this
395 		 * message. The next message starts after it. */
396 		fd_drain(fd, *len);
397 		(void)fd_set_nonblock(fd);
398 		return 0;
399 	}
400 	d = 0;
401 	while(d < (ssize_t)*len) {
402 		if((r=read(fd, (*buf)+d, (size_t)((ssize_t)*len)-d)) == -1) {
403 			log_err("tube msg read failed: %s", strerror(errno));
404 			(void)fd_set_nonblock(fd);
405 			free(*buf);
406 			return 0;
407 		}
408 		if(r == 0) { /* EOF */
409 			(void)fd_set_nonblock(fd);
410 			free(*buf);
411 			return 0;
412 		}
413 		d += r;
414 	}
415 	if(!fd_set_nonblock(fd)) {
416 		free(*buf);
417 		return 0;
418 	}
419 	return 1;
420 }
421 
422 /** perform poll() on the fd */
423 static int
pollit(int fd,struct timeval * t)424 pollit(int fd, struct timeval* t)
425 {
426 	struct pollfd fds;
427 	int pret;
428 	int msec = -1;
429 	memset(&fds, 0, sizeof(fds));
430 	fds.fd = fd;
431 	fds.events = POLLIN | POLLERR | POLLHUP;
432 #ifndef S_SPLINT_S
433 	if(t)
434 		msec = t->tv_sec*1000 + t->tv_usec/1000;
435 #endif
436 
437 	pret = poll(&fds, 1, msec);
438 
439 	if(pret == -1)
440 		return 0;
441 	if(pret != 0)
442 		return 1;
443 	return 0;
444 }
445 
tube_poll(struct tube * tube)446 int tube_poll(struct tube* tube)
447 {
448 	struct timeval t;
449 	memset(&t, 0, sizeof(t));
450 	return pollit(tube->sr, &t);
451 }
452 
tube_wait(struct tube * tube)453 int tube_wait(struct tube* tube)
454 {
455 	return pollit(tube->sr, NULL);
456 }
457 
tube_wait_timeout(struct tube * tube,int msec)458 int tube_wait_timeout(struct tube* tube, int msec)
459 {
460 	int ret = 0;
461 
462 	while(1) {
463 		struct pollfd fds;
464 		memset(&fds, 0, sizeof(fds));
465 
466 		fds.fd = tube->sr;
467 		fds.events = POLLIN | POLLERR | POLLHUP;
468 		ret = poll(&fds, 1, msec);
469 
470 		if(ret == -1) {
471 			if(errno == EAGAIN || errno == EINTR)
472 				continue;
473 			return -1;
474 		}
475 		break;
476 	}
477 
478 	if(ret != 0)
479 		return 1;
480 	return 0;
481 }
482 
tube_read_fd(struct tube * tube)483 int tube_read_fd(struct tube* tube)
484 {
485 	return tube->sr;
486 }
487 
tube_setup_bg_listen(struct tube * tube,struct comm_base * base,tube_callback_type * cb,void * arg)488 int tube_setup_bg_listen(struct tube* tube, struct comm_base* base,
489         tube_callback_type* cb, void* arg)
490 {
491 	tube->listen_cb = cb;
492 	tube->listen_arg = arg;
493 	if(!(tube->listen_com = comm_point_create_raw(base, tube->sr,
494 		0, tube_handle_listen, tube))) {
495 		int err = errno;
496 		log_err("tube_setup_bg_l: commpoint creation failed");
497 		errno = err;
498 		return 0;
499 	}
500 	return 1;
501 }
502 
tube_setup_bg_write(struct tube * tube,struct comm_base * base)503 int tube_setup_bg_write(struct tube* tube, struct comm_base* base)
504 {
505 	if(!(tube->res_com = comm_point_create_raw(base, tube->sw,
506 		1, tube_handle_write, tube))) {
507 		int err = errno;
508 		log_err("tube_setup_bg_w: commpoint creation failed");
509 		errno = err;
510 		return 0;
511 	}
512 	return 1;
513 }
514 
tube_queue_item(struct tube * tube,uint8_t * msg,size_t len)515 int tube_queue_item(struct tube* tube, uint8_t* msg, size_t len)
516 {
517 	struct tube_res_list* item;
518 	if(!tube || !tube->res_com) return 0;
519 	item = (struct tube_res_list*)malloc(sizeof(*item));
520 	if(!item) {
521 		free(msg);
522 		log_err("out of memory for async answer");
523 		return 0;
524 	}
525 	item->buf = msg;
526 	item->len = len;
527 	item->next = NULL;
528 	/* add at back of list, since the first one may be partially written */
529 	if(tube->res_last)
530 		tube->res_last->next = item;
531 	else    tube->res_list = item;
532 	tube->res_last = item;
533 	if(tube->res_list == tube->res_last) {
534 		/* first added item, start the write process */
535 		comm_point_start_listening(tube->res_com, -1, -1);
536 	}
537 	return 1;
538 }
539 
tube_handle_signal(int ATTR_UNUSED (fd),short ATTR_UNUSED (events),void * ATTR_UNUSED (arg))540 void tube_handle_signal(int ATTR_UNUSED(fd), short ATTR_UNUSED(events),
541 	void* ATTR_UNUSED(arg))
542 {
543 	log_assert(0);
544 }
545 
546 #else /* USE_WINSOCK */
547 /* on windows */
548 
549 
tube_create(void)550 struct tube* tube_create(void)
551 {
552 	/* windows does not have forks like unix, so we only support
553 	 * threads on windows. And thus the pipe need only connect
554 	 * threads. We use a mutex and a list of datagrams. */
555 	struct tube* tube = (struct tube*)calloc(1, sizeof(*tube));
556 	if(!tube) {
557 		int err = errno;
558 		log_err("tube_create: out of memory");
559 		errno = err;
560 		return NULL;
561 	}
562 	tube->event = WSACreateEvent();
563 	if(tube->event == WSA_INVALID_EVENT) {
564 		free(tube);
565 		log_err("WSACreateEvent: %s", wsa_strerror(WSAGetLastError()));
566 		return NULL;
567 	}
568 	if(!WSAResetEvent(tube->event)) {
569 		log_err("WSAResetEvent: %s", wsa_strerror(WSAGetLastError()));
570 	}
571 	lock_basic_init(&tube->res_lock);
572 	verbose(VERB_ALGO, "tube created");
573 	return tube;
574 }
575 
tube_delete(struct tube * tube)576 void tube_delete(struct tube* tube)
577 {
578 	if(!tube) return;
579 	tube_remove_bg_listen(tube);
580 	tube_remove_bg_write(tube);
581 	tube_close_read(tube);
582 	tube_close_write(tube);
583 	if(!WSACloseEvent(tube->event))
584 		log_err("WSACloseEvent: %s", wsa_strerror(WSAGetLastError()));
585 	lock_basic_destroy(&tube->res_lock);
586 	verbose(VERB_ALGO, "tube deleted");
587 	free(tube);
588 }
589 
tube_close_read(struct tube * ATTR_UNUSED (tube))590 void tube_close_read(struct tube* ATTR_UNUSED(tube))
591 {
592 	verbose(VERB_ALGO, "tube close_read");
593 }
594 
tube_close_write(struct tube * ATTR_UNUSED (tube))595 void tube_close_write(struct tube* ATTR_UNUSED(tube))
596 {
597 	verbose(VERB_ALGO, "tube close_write");
598 	/* wake up waiting reader with an empty queue */
599 	if(!WSASetEvent(tube->event)) {
600 		log_err("WSASetEvent: %s", wsa_strerror(WSAGetLastError()));
601 	}
602 }
603 
tube_remove_bg_listen(struct tube * tube)604 void tube_remove_bg_listen(struct tube* tube)
605 {
606 	verbose(VERB_ALGO, "tube remove_bg_listen");
607 	if (tube->ev_listen != NULL) {
608 		ub_winsock_unregister_wsaevent(tube->ev_listen);
609 		tube->ev_listen = NULL;
610 	}
611 }
612 
tube_remove_bg_write(struct tube * tube)613 void tube_remove_bg_write(struct tube* tube)
614 {
615 	verbose(VERB_ALGO, "tube remove_bg_write");
616 	if(tube->res_list) {
617 		struct tube_res_list* np, *p = tube->res_list;
618 		tube->res_list = NULL;
619 		tube->res_last = NULL;
620 		while(p) {
621 			np = p->next;
622 			free(p->buf);
623 			free(p);
624 			p = np;
625 		}
626 	}
627 }
628 
tube_write_msg(struct tube * tube,uint8_t * buf,uint32_t len,int ATTR_UNUSED (nonblock))629 int tube_write_msg(struct tube* tube, uint8_t* buf, uint32_t len,
630         int ATTR_UNUSED(nonblock))
631 {
632 	uint8_t* a;
633 	verbose(VERB_ALGO, "tube write_msg len %d", (int)len);
634 	a = (uint8_t*)memdup(buf, len);
635 	if(!a) {
636 		log_err("out of memory in tube_write_msg");
637 		return 0;
638 	}
639 	/* always nonblocking, this pipe cannot get full */
640 	return tube_queue_item(tube, a, len);
641 }
642 
tube_read_msg(struct tube * tube,uint8_t ** buf,uint32_t * len,int nonblock)643 int tube_read_msg(struct tube* tube, uint8_t** buf, uint32_t* len,
644         int nonblock)
645 {
646 	struct tube_res_list* item = NULL;
647 	verbose(VERB_ALGO, "tube read_msg %s", nonblock?"nonblock":"blocking");
648 	*buf = NULL;
649 	if(!tube_poll(tube)) {
650 		verbose(VERB_ALGO, "tube read_msg nodata");
651 		/* nothing ready right now, wait if we want to */
652 		if(nonblock)
653 			return -1; /* would block waiting for items */
654 		if(!tube_wait(tube))
655 			return 0;
656 	}
657 	lock_basic_lock(&tube->res_lock);
658 	if(tube->res_list) {
659 		item = tube->res_list;
660 		tube->res_list = item->next;
661 		if(tube->res_last == item) {
662 			/* the list is now empty */
663 			tube->res_last = NULL;
664 			verbose(VERB_ALGO, "tube read_msg lastdata");
665 			if(!WSAResetEvent(tube->event)) {
666 				log_err("WSAResetEvent: %s",
667 					wsa_strerror(WSAGetLastError()));
668 			}
669 		}
670 	}
671 	lock_basic_unlock(&tube->res_lock);
672 	if(!item)
673 		return 0; /* would block waiting for items */
674 	*buf = item->buf;
675 	*len = item->len;
676 	free(item);
677 	verbose(VERB_ALGO, "tube read_msg len %d", (int)*len);
678 	return 1;
679 }
680 
tube_poll(struct tube * tube)681 int tube_poll(struct tube* tube)
682 {
683 	struct tube_res_list* item = NULL;
684 	lock_basic_lock(&tube->res_lock);
685 	item = tube->res_list;
686 	lock_basic_unlock(&tube->res_lock);
687 	if(item)
688 		return 1;
689 	return 0;
690 }
691 
tube_wait(struct tube * tube)692 int tube_wait(struct tube* tube)
693 {
694 	/* block on eventhandle */
695 	DWORD res = WSAWaitForMultipleEvents(
696 		1 /* one event in array */,
697 		&tube->event /* the event to wait for, our pipe signal */,
698 		0 /* wait for all events is false */,
699 		WSA_INFINITE /* wait, no timeout */,
700 		0 /* we are not alertable for IO completion routines */
701 		);
702 	if(res == WSA_WAIT_TIMEOUT) {
703 		return 0;
704 	}
705 	if(res == WSA_WAIT_IO_COMPLETION) {
706 		/* a bit unexpected, since we were not alertable */
707 		return 0;
708 	}
709 	return 1;
710 }
711 
tube_wait_timeout(struct tube * tube,int msec)712 int tube_wait_timeout(struct tube* tube, int msec)
713 {
714 	/* block on eventhandle */
715 	DWORD res = WSAWaitForMultipleEvents(
716 		1 /* one event in array */,
717 		&tube->event /* the event to wait for, our pipe signal */,
718 		0 /* wait for all events is false */,
719 		msec /* wait for timeout */,
720 		0 /* we are not alertable for IO completion routines */
721 		);
722 	if(res == WSA_WAIT_TIMEOUT) {
723 		return 0;
724 	}
725 	if(res == WSA_WAIT_IO_COMPLETION) {
726 		/* a bit unexpected, since we were not alertable */
727 		return -1;
728 	}
729 	return 1;
730 }
731 
tube_read_fd(struct tube * ATTR_UNUSED (tube))732 int tube_read_fd(struct tube* ATTR_UNUSED(tube))
733 {
734 	/* nothing sensible on Windows */
735 	return -1;
736 }
737 
738 int
tube_handle_listen(struct comm_point * ATTR_UNUSED (c),void * ATTR_UNUSED (arg),int ATTR_UNUSED (error),struct comm_reply * ATTR_UNUSED (reply_info))739 tube_handle_listen(struct comm_point* ATTR_UNUSED(c), void* ATTR_UNUSED(arg),
740 	int ATTR_UNUSED(error), struct comm_reply* ATTR_UNUSED(reply_info))
741 {
742 	log_assert(0);
743 	return 0;
744 }
745 
746 int
tube_handle_write(struct comm_point * ATTR_UNUSED (c),void * ATTR_UNUSED (arg),int ATTR_UNUSED (error),struct comm_reply * ATTR_UNUSED (reply_info))747 tube_handle_write(struct comm_point* ATTR_UNUSED(c), void* ATTR_UNUSED(arg),
748 	int ATTR_UNUSED(error), struct comm_reply* ATTR_UNUSED(reply_info))
749 {
750 	log_assert(0);
751 	return 0;
752 }
753 
tube_setup_bg_listen(struct tube * tube,struct comm_base * base,tube_callback_type * cb,void * arg)754 int tube_setup_bg_listen(struct tube* tube, struct comm_base* base,
755         tube_callback_type* cb, void* arg)
756 {
757 	tube->listen_cb = cb;
758 	tube->listen_arg = arg;
759 	if(!comm_base_internal(base))
760 		return 1; /* ignore when no comm base - testing */
761 	tube->ev_listen = ub_winsock_register_wsaevent(
762 	    comm_base_internal(base), tube->event, &tube_handle_signal, tube);
763 	return tube->ev_listen ? 1 : 0;
764 }
765 
tube_setup_bg_write(struct tube * ATTR_UNUSED (tube),struct comm_base * ATTR_UNUSED (base))766 int tube_setup_bg_write(struct tube* ATTR_UNUSED(tube),
767 	struct comm_base* ATTR_UNUSED(base))
768 {
769 	/* the queue item routine performs the signaling */
770 	return 1;
771 }
772 
tube_queue_item(struct tube * tube,uint8_t * msg,size_t len)773 int tube_queue_item(struct tube* tube, uint8_t* msg, size_t len)
774 {
775 	struct tube_res_list* item;
776 	if(!tube) return 0;
777 	item = (struct tube_res_list*)malloc(sizeof(*item));
778 	verbose(VERB_ALGO, "tube queue_item len %d", (int)len);
779 	if(!item) {
780 		free(msg);
781 		log_err("out of memory for async answer");
782 		return 0;
783 	}
784 	item->buf = msg;
785 	item->len = len;
786 	item->next = NULL;
787 	lock_basic_lock(&tube->res_lock);
788 	/* add at back of list, since the first one may be partially written */
789 	if(tube->res_last)
790 		tube->res_last->next = item;
791 	else    tube->res_list = item;
792 	tube->res_last = item;
793 	/* signal the eventhandle */
794 	if(!WSASetEvent(tube->event)) {
795 		log_err("WSASetEvent: %s", wsa_strerror(WSAGetLastError()));
796 	}
797 	lock_basic_unlock(&tube->res_lock);
798 	return 1;
799 }
800 
tube_handle_signal(int ATTR_UNUSED (fd),short ATTR_UNUSED (events),void * arg)801 void tube_handle_signal(int ATTR_UNUSED(fd), short ATTR_UNUSED(events),
802 	void* arg)
803 {
804 	struct tube* tube = (struct tube*)arg;
805 	uint8_t* buf;
806 	uint32_t len = 0;
807 	verbose(VERB_ALGO, "tube handle_signal");
808 	while(tube_poll(tube)) {
809 		if(tube_read_msg(tube, &buf, &len, 1)) {
810 			fptr_ok(fptr_whitelist_tube_listen(tube->listen_cb));
811 			(*tube->listen_cb)(tube, buf, len, NETEVENT_NOERROR,
812 				tube->listen_arg);
813 		}
814 	}
815 }
816 
817 #endif /* USE_WINSOCK */
818