Home | History | Annotate | Line # | Download | only in util
      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 
     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 
     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 
    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 
    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 
    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 
    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
    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
    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
    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 
    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 
    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
    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 
    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 
    453 int tube_wait(struct tube* tube)
    454 {
    455 	return pollit(tube->sr, NULL);
    456 }
    457 
    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 
    483 int tube_read_fd(struct tube* tube)
    484 {
    485 	return tube->sr;
    486 }
    487 
    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 
    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 
    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 
    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 
    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 
    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 
    590 void tube_close_read(struct tube* ATTR_UNUSED(tube))
    591 {
    592 	verbose(VERB_ALGO, "tube close_read");
    593 }
    594 
    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 
    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 
    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 
    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 
    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 
    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 
    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 
    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 
    732 int tube_read_fd(struct tube* ATTR_UNUSED(tube))
    733 {
    734 	/* nothing sensible on Windows */
    735 	return -1;
    736 }
    737 
    738 int
    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
    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 
    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 
    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 
    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 
    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