1 /* $NetBSD: sys_pipe.c,v 1.185 2026/10/05 22:59:15 kre Exp $ */ 2 3 /*- 4 * Copyright (c) 2003, 2007, 2008, 2009, 2023 The NetBSD Foundation, Inc. 5 * All rights reserved. 6 * 7 * This code is derived from software contributed to The NetBSD Foundation 8 * by Paul Kranenburg, and by Andrew Doran. 9 * 10 * Redistribution and use in source and binary forms, with or without 11 * modification, are permitted provided that the following conditions 12 * are met: 13 * 1. Redistributions of source code must retain the above copyright 14 * notice, this list of conditions and the following disclaimer. 15 * 2. Redistributions in binary form must reproduce the above copyright 16 * notice, this list of conditions and the following disclaimer in the 17 * documentation and/or other materials provided with the distribution. 18 * 19 * THIS SOFTWARE IS PROVIDED BY THE NETBSD FOUNDATION, INC. AND CONTRIBUTORS 20 * ``AS IS'' AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED 21 * TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR 22 * PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE FOUNDATION OR CONTRIBUTORS 23 * BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR 24 * CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF 25 * SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS 26 * INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN 27 * CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) 28 * ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE 29 * POSSIBILITY OF SUCH DAMAGE. 30 */ 31 32 /* 33 * Copyright (c) 1996 John S. Dyson 34 * All rights reserved. 35 * 36 * Redistribution and use in source and binary forms, with or without 37 * modification, are permitted provided that the following conditions 38 * are met: 39 * 1. Redistributions of source code must retain the above copyright 40 * notice immediately at the beginning of the file, without modification, 41 * this list of conditions, and the following disclaimer. 42 * 2. Redistributions in binary form must reproduce the above copyright 43 * notice, this list of conditions and the following disclaimer in the 44 * documentation and/or other materials provided with the distribution. 45 * 3. Absolutely no warranty of function or purpose is made by the author 46 * John S. Dyson. 47 * 4. Modifications may be freely made to this file if the above conditions 48 * are met. 49 */ 50 51 /* 52 * This file contains a high-performance replacement for the socket-based 53 * pipes scheme originally used. It does not support all features of 54 * sockets, but does do everything that pipes normally do. 55 */ 56 57 #include <sys/cdefs.h> 58 __KERNEL_RCSID(0, "$NetBSD: sys_pipe.c,v 1.185 2026/10/05 22:59:15 kre Exp $"); 59 60 #include <sys/param.h> 61 #include <sys/types.h> 62 63 #include <sys/atomic.h> 64 #include <sys/fcntl.h> 65 #include <sys/file.h> 66 #include <sys/filedesc.h> 67 #include <sys/filio.h> 68 #include <sys/kauth.h> 69 #include <sys/kernel.h> 70 #include <sys/mount.h> 71 #include <sys/pipe.h> 72 #include <sys/poll.h> 73 #include <sys/proc.h> 74 #include <sys/sdt.h> 75 #include <sys/select.h> 76 #include <sys/signalvar.h> 77 #include <sys/stat.h> 78 #include <sys/syscallargs.h> 79 #include <sys/sysctl.h> 80 #include <sys/systm.h> 81 #include <sys/ttycom.h> 82 #include <sys/uio.h> 83 #include <sys/vnode.h> 84 85 static int pipe_read(file_t *, off_t *, struct uio *, kauth_cred_t, int); 86 static int pipe_write(file_t *, off_t *, struct uio *, kauth_cred_t, int); 87 static int pipe_close(file_t *); 88 static int pipe_poll(file_t *, int); 89 static int pipe_kqfilter(file_t *, struct knote *); 90 static int pipe_stat(file_t *, struct stat *); 91 static int pipe_ioctl(file_t *, u_long, void *); 92 static void pipe_restart(file_t *); 93 static int pipe_fpathconf(file_t *, int, register_t *); 94 static int pipe_posix_fadvise(file_t *, off_t, off_t, int); 95 96 static const struct fileops pipeops = { 97 .fo_name = "pipe", 98 .fo_read = pipe_read, 99 .fo_write = pipe_write, 100 .fo_ioctl = pipe_ioctl, 101 .fo_fcntl = fnullop_fcntl, 102 .fo_poll = pipe_poll, 103 .fo_stat = pipe_stat, 104 .fo_close = pipe_close, 105 .fo_kqfilter = pipe_kqfilter, 106 .fo_restart = pipe_restart, 107 .fo_fpathconf = pipe_fpathconf, 108 .fo_posix_fadvise = pipe_posix_fadvise, 109 }; 110 111 /* 112 * Default pipe buffer size(s), this can be kind-of large now because pipe 113 * space is pageable. The pipe code will try to maintain locality of 114 * reference for performance reasons, so small amounts of outstanding I/O 115 * will not wipe the cache. 116 */ 117 #define MINPIPESIZE (PIPE_SIZE / 3) 118 #define MAXPIPESIZE (2 * PIPE_SIZE / 3) 119 120 /* 121 * Limit the number of "big" pipes 122 */ 123 #define LIMITBIGPIPES 32 124 static u_int maxbigpipes __read_mostly = LIMITBIGPIPES; 125 static u_int nbigpipe = 0; 126 127 /* 128 * Amount of KVA consumed by pipe buffers. 129 */ 130 static u_int amountpipekva = 0; 131 132 static void pipeclose(struct file *, struct pipe *); 133 static void pipefree(struct pipe *); 134 static void pipe_free_kmem(struct pipe *); 135 static int pipe_create(struct pipe **, pool_cache_t, struct timespec *); 136 static int pipelock(struct pipe *, bool); 137 static inline void pipeunlock(struct pipe *); 138 static void pipeselwakeup(struct pipe *, int); 139 static int pipespace(struct pipe *, int); 140 static int pipe_ctor(void *, void *, int); 141 static void pipe_dtor(void *, void *); 142 143 static pool_cache_t pipe_wr_cache; 144 static pool_cache_t pipe_rd_cache; 145 146 void 147 pipe_init(void) 148 { 149 150 /* Writer side is not automatically allocated KVA. */ 151 pipe_wr_cache = pool_cache_init(sizeof(struct pipe), 0, 0, 0, "pipewr", 152 NULL, IPL_NONE, pipe_ctor, pipe_dtor, NULL); 153 KASSERT(pipe_wr_cache != NULL); 154 155 /* Reader side gets preallocated KVA. */ 156 pipe_rd_cache = pool_cache_init(sizeof(struct pipe), 0, 0, 0, "piperd", 157 NULL, IPL_NONE, pipe_ctor, pipe_dtor, (void *)1); 158 KASSERT(pipe_rd_cache != NULL); 159 } 160 161 static int 162 pipe_ctor(void *arg, void *obj, int flags) 163 { 164 struct pipe *pipe; 165 vaddr_t va; 166 167 pipe = obj; 168 169 memset(pipe, 0, sizeof(struct pipe)); 170 if (arg != NULL) { 171 /* Preallocate space. */ 172 va = uvm_km_alloc(kernel_map, PIPE_SIZE, 0, 173 UVM_KMF_PAGEABLE | UVM_KMF_WAITVA); 174 KASSERT(va != 0); 175 pipe->pipe_kmem = va; 176 atomic_add_int(&amountpipekva, PIPE_SIZE); 177 } 178 cv_init(&pipe->pipe_rcv, "pipe_rd"); 179 cv_init(&pipe->pipe_wcv, "pipe_wr"); 180 cv_init(&pipe->pipe_draincv, "pipe_drn"); 181 cv_init(&pipe->pipe_lkcv, "pipe_lk"); 182 selinit(&pipe->pipe_sel); 183 pipe->pipe_state = PIPE_SIGNALR; 184 185 return 0; 186 } 187 188 static void 189 pipe_dtor(void *arg, void *obj) 190 { 191 struct pipe *pipe; 192 193 pipe = obj; 194 195 cv_destroy(&pipe->pipe_rcv); 196 cv_destroy(&pipe->pipe_wcv); 197 cv_destroy(&pipe->pipe_draincv); 198 cv_destroy(&pipe->pipe_lkcv); 199 seldestroy(&pipe->pipe_sel); 200 if (pipe->pipe_kmem != 0) { 201 uvm_km_free(kernel_map, pipe->pipe_kmem, PIPE_SIZE, 202 UVM_KMF_PAGEABLE); 203 atomic_add_int(&amountpipekva, -PIPE_SIZE); 204 } 205 } 206 207 /* 208 * The pipe system call for the DTYPE_PIPE type of pipes 209 */ 210 int 211 pipe1(struct lwp *l, int *fildes, int flags) 212 { 213 struct pipe *rpipe, *wpipe; 214 struct timespec nt; 215 file_t *rf, *wf; 216 int fd, error; 217 proc_t *p; 218 219 if (flags & ~(O_CLOEXEC|O_CLOFORK|O_NONBLOCK|O_NOSIGPIPE)) 220 return SET_ERROR(EINVAL); 221 p = curproc; 222 rpipe = wpipe = NULL; 223 getnanotime(&nt); 224 if ((error = pipe_create(&rpipe, pipe_rd_cache, &nt)) || 225 (error = pipe_create(&wpipe, pipe_wr_cache, &nt))) { 226 goto free2; 227 } 228 rpipe->pipe_lock = mutex_obj_alloc(MUTEX_DEFAULT, IPL_NONE); 229 wpipe->pipe_lock = rpipe->pipe_lock; 230 mutex_obj_hold(wpipe->pipe_lock); 231 232 error = fd_allocfile(&rf, &fd); 233 if (error) 234 goto free2; 235 fildes[0] = fd; 236 237 error = fd_allocfile(&wf, &fd); 238 if (error) 239 goto free3; 240 fildes[1] = fd; 241 242 rf->f_flag = FREAD | flags; 243 rf->f_type = DTYPE_PIPE; 244 rf->f_pipe = rpipe; 245 rf->f_ops = &pipeops; 246 fd_set_exclose(l, fildes[0], (flags & O_CLOEXEC) != 0); 247 fd_set_foclose(l, fildes[0], (flags & O_CLOFORK) != 0); 248 249 wf->f_flag = FWRITE | flags; 250 wf->f_type = DTYPE_PIPE; 251 wf->f_pipe = wpipe; 252 wf->f_ops = &pipeops; 253 fd_set_exclose(l, fildes[1], (flags & O_CLOEXEC) != 0); 254 fd_set_foclose(l, fildes[1], (flags & O_CLOFORK) != 0); 255 256 rpipe->pipe_peer = wpipe; 257 wpipe->pipe_peer = rpipe; 258 259 fd_affix(p, rf, fildes[0]); 260 fd_affix(p, wf, fildes[1]); 261 return 0; 262 free3: 263 fd_abort(p, rf, fildes[0]); 264 free2: 265 if (wpipe) 266 pipefree(wpipe); 267 if (rpipe) 268 pipefree(rpipe); 269 270 return error; 271 } 272 273 /* 274 * Allocate kva for pipe circular buffer, the space is pageable 275 * This routine will 'realloc' the size of a pipe safely, if it fails 276 * it will retain the old buffer. 277 * If it fails it will return ENOMEM. 278 */ 279 static int 280 pipespace(struct pipe *pipe, int size) 281 { 282 void *buffer; 283 284 /* 285 * Allocate pageable virtual address space. Physical memory is 286 * allocated on demand. 287 */ 288 if (size == PIPE_SIZE && pipe->pipe_kmem != 0) { 289 buffer = (void *)pipe->pipe_kmem; 290 } else { 291 buffer = (void *)uvm_km_alloc(kernel_map, round_page(size), 292 0, UVM_KMF_PAGEABLE); 293 if (buffer == NULL) 294 return SET_ERROR(ENOMEM); 295 atomic_add_int(&amountpipekva, size); 296 } 297 298 /* free old resources if we're resizing */ 299 pipe_free_kmem(pipe); 300 pipe->pipe_buffer.buffer = buffer; 301 pipe->pipe_buffer.size = size; 302 pipe->pipe_buffer.in = 0; 303 pipe->pipe_buffer.out = 0; 304 pipe->pipe_buffer.cnt = 0; 305 return 0; 306 } 307 308 /* 309 * Initialize and allocate VM and memory for pipe. 310 */ 311 static int 312 pipe_create(struct pipe **pipep, pool_cache_t cache, struct timespec *nt) 313 { 314 struct pipe *pipe; 315 int error; 316 317 pipe = pool_cache_get(cache, PR_WAITOK); 318 KASSERT(pipe != NULL); 319 *pipep = pipe; 320 error = 0; 321 pipe->pipe_atime = pipe->pipe_mtime = pipe->pipe_btime = *nt; 322 pipe->pipe_lock = NULL; 323 if (cache == pipe_rd_cache) { 324 error = pipespace(pipe, PIPE_SIZE); 325 } else { 326 pipe->pipe_buffer.buffer = NULL; 327 pipe->pipe_buffer.size = 0; 328 pipe->pipe_buffer.in = 0; 329 pipe->pipe_buffer.out = 0; 330 pipe->pipe_buffer.cnt = 0; 331 } 332 return error; 333 } 334 335 /* 336 * Lock a pipe for I/O, blocking other access 337 * Called with pipe spin lock held. 338 */ 339 static int 340 pipelock(struct pipe *pipe, bool catch_p) 341 { 342 int error; 343 344 KASSERT(mutex_owned(pipe->pipe_lock)); 345 346 while (pipe->pipe_state & PIPE_LOCKFL) { 347 if (catch_p) { 348 error = cv_wait_sig(&pipe->pipe_lkcv, pipe->pipe_lock); 349 if (error != 0) { 350 return error; 351 } 352 } else 353 cv_wait(&pipe->pipe_lkcv, pipe->pipe_lock); 354 } 355 356 pipe->pipe_state |= PIPE_LOCKFL; 357 358 return 0; 359 } 360 361 /* 362 * unlock a pipe I/O lock 363 */ 364 static inline void 365 pipeunlock(struct pipe *pipe) 366 { 367 368 KASSERT(pipe->pipe_state & PIPE_LOCKFL); 369 370 pipe->pipe_state &= ~PIPE_LOCKFL; 371 cv_signal(&pipe->pipe_lkcv); 372 } 373 374 /* 375 * pipeselwakeup(pipe, code) 376 * 377 * Activity has happened on pipe's peer oncausing I/O to be 378 * available on pipe, so: 379 * 380 * 1. Wake any threads waiting in select/poll on pipe. 381 * 382 * 2. Deliver SIGIO to any process (group) configured to receive 383 * notifications about I/O on pipe. 384 * 385 * `code' is a siginfo_t si_code value in the POLL_* namespace for 386 * the type of notification the waiters will receive, and it 387 * should match the direction of the pipe -- POLL_OUT/POLL_ERR 388 * with the writer side, POLL_IN/POLL_HUP with the reader side. 389 */ 390 static void 391 pipeselwakeup(struct pipe *pipe, int code) 392 { 393 int band; 394 395 KASSERT(mutex_owned(pipe->pipe_lock)); 396 397 switch (code) { 398 case POLL_IN: 399 band = POLLIN|POLLRDNORM; 400 break; 401 case POLL_OUT: 402 band = POLLOUT|POLLWRNORM; 403 break; 404 case POLL_HUP: 405 band = POLLHUP; 406 break; 407 case POLL_ERR: 408 band = POLLERR; 409 break; 410 default: 411 band = 0; 412 #ifdef DIAGNOSTIC 413 printf("bad siginfo code %d in pipe notification.\n", code); 414 #endif 415 break; 416 } 417 418 selnotify(&pipe->pipe_sel, band, NOTE_SUBMIT); 419 420 if ((pipe->pipe_state & PIPE_ASYNC) == 0) 421 return; 422 423 fownsignal(pipe->pipe_pgid, SIGIO, code, band, pipe); 424 } 425 426 static int 427 pipe_read(file_t *fp, off_t *offset, struct uio *uio, kauth_cred_t cred, 428 int flags) 429 { 430 struct pipe *rpipe = fp->f_pipe; 431 struct pipe *wpipe; 432 struct pipebuf *bp = &rpipe->pipe_buffer; 433 kmutex_t *lock = rpipe->pipe_lock; 434 int error; 435 size_t nread = 0; 436 size_t size; 437 size_t ocnt; 438 unsigned int wakeup_state = 0; 439 440 /* 441 * Try to avoid locking the pipe if we have nothing to do. 442 * 443 * There are programs which share one pipe amongst multiple processes 444 * and perform non-blocking reads in parallel, even if the pipe is 445 * empty. This in particular is the case with BSD make, which when 446 * spawned with a high -j number can find itself with over half of the 447 * calls failing to find anything. 448 */ 449 if ((fp->f_flag & FNONBLOCK) != 0) { 450 if (__predict_false(uio->uio_resid == 0)) 451 return 0; 452 if (atomic_load_relaxed(&bp->cnt) == 0 && 453 (atomic_load_relaxed(&rpipe->pipe_state) & PIPE_EOF) == 0) 454 return SET_ERROR(EAGAIN); 455 } 456 457 mutex_enter(lock); 458 ++rpipe->pipe_busy; 459 ocnt = bp->cnt; 460 461 again: 462 error = pipelock(rpipe, true); 463 if (error) 464 goto unlocked_error; 465 466 while (uio->uio_resid) { 467 /* 468 * Normal pipe buffer receive. 469 */ 470 if (bp->cnt > 0) { 471 size = bp->size - bp->out; 472 if (size > bp->cnt) 473 size = bp->cnt; 474 if (size > uio->uio_resid) 475 size = uio->uio_resid; 476 477 mutex_exit(lock); 478 error = uiomove((char *)bp->buffer + bp->out, size, uio); 479 mutex_enter(lock); 480 if (error) 481 break; 482 483 bp->out += size; 484 if (bp->out >= bp->size) 485 bp->out = 0; 486 487 bp->cnt -= size; 488 489 /* 490 * If there is no more to read in the pipe, reset 491 * its pointers to the beginning. This improves 492 * cache hit stats. 493 */ 494 if (bp->cnt == 0) { 495 bp->in = 0; 496 bp->out = 0; 497 } 498 nread += size; 499 continue; 500 } 501 502 /* 503 * Break if some data was read. 504 */ 505 if (nread > 0) 506 break; 507 508 /* 509 * Detect EOF condition. 510 * Read returns 0 on EOF, no need to set error. 511 * 512 * XXX Why rpipe->pipe_state and not wpipe->pipe_state? 513 * XXX Distinguish reader-closed from writer-closed? 514 */ 515 if (rpipe->pipe_state & PIPE_EOF) 516 break; 517 518 /* 519 * Don't block on non-blocking I/O. 520 */ 521 if (fp->f_flag & FNONBLOCK) { 522 error = SET_ERROR(EAGAIN); 523 break; 524 } 525 526 /* 527 * Unlock the pipe buffer for our remaining processing. 528 * We will either break out with an error or we will 529 * sleep and relock to loop. 530 */ 531 pipeunlock(rpipe); 532 533 /* 534 * If the "write-side" is blocked, wake it up now. 535 */ 536 KASSERT((rpipe->pipe_state & PIPE_EOF) == 0); 537 KASSERT(rpipe->pipe_peer != NULL); 538 wpipe = rpipe->pipe_peer; 539 pipeselwakeup(wpipe, POLL_OUT); 540 cv_broadcast(&wpipe->pipe_wcv); 541 542 if (wakeup_state & PIPE_RESTART) { 543 error = SET_ERROR(ERESTART); 544 goto unlocked_error; 545 } 546 547 /* Now wait until the pipe is filled */ 548 error = cv_wait_sig(&rpipe->pipe_rcv, lock); 549 if (error != 0) 550 goto unlocked_error; 551 wakeup_state = rpipe->pipe_state; 552 goto again; 553 } 554 555 if (error == 0) 556 getnanotime(&rpipe->pipe_atime); 557 pipeunlock(rpipe); 558 559 unlocked_error: 560 --rpipe->pipe_busy; 561 if (rpipe->pipe_busy == 0) { 562 cv_broadcast(&rpipe->pipe_draincv); 563 } 564 if (bp->cnt < MINPIPESIZE) { 565 if ((wpipe = rpipe->pipe_peer) != NULL) 566 cv_broadcast(&wpipe->pipe_wcv); 567 } 568 569 /* 570 * If anything was read off the buffer, signal to the writer it's 571 * possible to write more data. Also send signal if we are here for the 572 * first time after last write. 573 */ 574 if ((bp->size - bp->cnt) >= PIPE_BUF 575 && (ocnt != bp->cnt || (rpipe->pipe_state & PIPE_SIGNALR))) { 576 if ((wpipe = rpipe->pipe_peer) != NULL) 577 pipeselwakeup(wpipe, POLL_OUT); 578 rpipe->pipe_state &= ~PIPE_SIGNALR; 579 } 580 581 mutex_exit(lock); 582 return error; 583 } 584 585 static int 586 pipe_write(file_t *fp, off_t *offset, struct uio *uio, kauth_cred_t cred, 587 int flags) 588 { 589 struct pipe *wpipe, *rpipe; 590 struct pipebuf *bp; 591 kmutex_t *lock; 592 int error; 593 unsigned int wakeup_state = 0; 594 595 /* We want to write to our peer */ 596 wpipe = fp->f_pipe; 597 lock = wpipe->pipe_lock; 598 error = 0; 599 600 mutex_enter(lock); 601 rpipe = wpipe->pipe_peer; 602 603 /* 604 * Detect loss of pipe read side, issue SIGPIPE if lost. 605 * 606 * After this, once we busy wpipe, the peer rpipe will remain 607 * stable (though may have PIPE_EOF set) until we unbusy it. 608 */ 609 if (rpipe == NULL || (rpipe->pipe_state & PIPE_EOF) != 0) { 610 mutex_exit(lock); 611 return SET_ERROR(EPIPE); 612 } 613 ++wpipe->pipe_busy; 614 615 /* Acquire the long-term pipe lock */ 616 if ((error = pipelock(rpipe, true)) != 0) { 617 --wpipe->pipe_busy; 618 if (wpipe->pipe_busy == 0) { 619 cv_broadcast(&wpipe->pipe_draincv); 620 } 621 mutex_exit(lock); 622 return error; 623 } 624 625 bp = &rpipe->pipe_buffer; 626 627 /* 628 * If it is advantageous to resize the pipe buffer, do so. 629 */ 630 if ((uio->uio_resid > PIPE_SIZE) && 631 (nbigpipe < maxbigpipes) && 632 (bp->size <= PIPE_SIZE) && (bp->cnt == 0)) { 633 634 if (pipespace(rpipe, BIG_PIPE_SIZE) == 0) 635 atomic_inc_uint(&nbigpipe); 636 } 637 638 while (uio->uio_resid) { 639 size_t space; 640 641 space = bp->size - bp->cnt; 642 643 /* Writes of size <= PIPE_BUF must be atomic. */ 644 if ((space < uio->uio_resid) && (uio->uio_resid <= PIPE_BUF)) 645 space = 0; 646 647 if (space > 0) { 648 int size; /* Transfer size */ 649 int segsize; /* first segment to transfer */ 650 651 /* 652 * Transfer size is minimum of uio transfer 653 * and free space in pipe buffer. 654 */ 655 if (space > uio->uio_resid) 656 size = uio->uio_resid; 657 else 658 size = space; 659 /* 660 * First segment to transfer is minimum of 661 * transfer size and contiguous space in 662 * pipe buffer. If first segment to transfer 663 * is less than the transfer size, we've got 664 * a wraparound in the buffer. 665 */ 666 segsize = bp->size - bp->in; 667 if (segsize > size) 668 segsize = size; 669 670 /* Transfer first segment */ 671 mutex_exit(lock); 672 error = uiomove((char *)bp->buffer + bp->in, segsize, 673 uio); 674 675 if (error == 0 && segsize < size) { 676 /* 677 * Transfer remaining part now, to 678 * support atomic writes. Wraparound 679 * happened. 680 */ 681 KASSERT(bp->in + segsize == bp->size); 682 error = uiomove(bp->buffer, 683 size - segsize, uio); 684 } 685 mutex_enter(lock); 686 if (error) 687 break; 688 689 bp->in += size; 690 if (bp->in >= bp->size) { 691 KASSERT(bp->in == size - segsize + bp->size); 692 bp->in = size - segsize; 693 } 694 695 bp->cnt += size; 696 KASSERT(bp->cnt <= bp->size); 697 wakeup_state = 0; 698 } else { 699 /* 700 * If the "read-side" has been blocked, wake it up now. 701 */ 702 cv_broadcast(&rpipe->pipe_rcv); 703 704 /* 705 * Don't block on non-blocking I/O. 706 */ 707 if (fp->f_flag & FNONBLOCK) { 708 error = SET_ERROR(EAGAIN); 709 break; 710 } 711 712 /* 713 * We have no more space and have something to offer, 714 * wake up select/poll. 715 */ 716 if (bp->cnt) 717 pipeselwakeup(rpipe, POLL_IN); 718 719 if (wakeup_state & PIPE_RESTART) { 720 error = SET_ERROR(ERESTART); 721 break; 722 } 723 724 /* 725 * If read side wants to go away, we just issue a signal 726 * to ourselves. 727 * 728 * XXX Shouldn't this happen before we uiomove anything? 729 * 730 * XXX Why rpipe->pipe_state and not wpipe->pipe_state? 731 * XXX Distinguish reader-closed from writer-closed? 732 */ 733 if (rpipe->pipe_state & PIPE_EOF) { 734 error = SET_ERROR(EPIPE); 735 break; 736 } 737 738 pipeunlock(rpipe); 739 error = cv_wait_sig(&wpipe->pipe_wcv, lock); 740 (void)pipelock(rpipe, false); 741 if (error != 0) 742 break; 743 wakeup_state = wpipe->pipe_state; 744 } 745 } 746 747 --wpipe->pipe_busy; 748 if (wpipe->pipe_busy == 0) { 749 cv_broadcast(&wpipe->pipe_draincv); 750 } 751 if (bp->cnt > 0) { 752 cv_broadcast(&rpipe->pipe_rcv); 753 } 754 755 /* 756 * Don't return EPIPE if I/O was successful 757 * 758 * XXX Shouldn't we avoid returning _any_ error if we 759 * transmitted _any_ positive number of bytes? Or does that 760 * happen downstream of here, and if so, why do we need to do 761 * that here? 762 */ 763 if (error == EPIPE && bp->cnt == 0 && uio->uio_resid == 0) 764 error = 0; 765 766 if (error == 0) 767 getnanotime(&rpipe->pipe_mtime); 768 769 /* 770 * We have something to offer, wake up select/poll. 771 */ 772 if (bp->cnt) 773 pipeselwakeup(rpipe, POLL_IN); 774 775 /* 776 * Arrange for next read(2) to do a signal. 777 */ 778 rpipe->pipe_state |= PIPE_SIGNALR; 779 780 pipeunlock(rpipe); 781 mutex_exit(lock); 782 return error; 783 } 784 785 /* 786 * We implement a very minimal set of ioctls for compatibility with sockets. 787 */ 788 int 789 pipe_ioctl(file_t *fp, u_long cmd, void *data) 790 { 791 struct pipe *pipe = fp->f_pipe; 792 kmutex_t *lock = pipe->pipe_lock; 793 794 switch (cmd) { 795 796 case FIONBIO: 797 return 0; 798 799 case FIOASYNC: 800 mutex_enter(lock); 801 if (*(int *)data) { 802 pipe->pipe_state |= PIPE_ASYNC; 803 } else { 804 pipe->pipe_state &= ~PIPE_ASYNC; 805 } 806 mutex_exit(lock); 807 return 0; 808 809 case FIONREAD: 810 mutex_enter(lock); 811 *(int *)data = pipe->pipe_buffer.cnt; 812 mutex_exit(lock); 813 return 0; 814 815 case FIONWRITE: 816 /* Look at other side */ 817 mutex_enter(lock); 818 pipe = pipe->pipe_peer; 819 if (pipe == NULL) 820 *(int *)data = 0; 821 else 822 *(int *)data = pipe->pipe_buffer.cnt; 823 mutex_exit(lock); 824 return 0; 825 826 case FIONSPACE: 827 /* Look at other side */ 828 mutex_enter(lock); 829 pipe = pipe->pipe_peer; 830 if (pipe == NULL) 831 *(int *)data = 0; 832 else 833 *(int *)data = pipe->pipe_buffer.size - 834 pipe->pipe_buffer.cnt; 835 mutex_exit(lock); 836 return 0; 837 838 case TIOCSPGRP: 839 case FIOSETOWN: 840 return fsetown(&pipe->pipe_pgid, cmd, data); 841 842 case TIOCGPGRP: 843 case FIOGETOWN: 844 return fgetown(pipe->pipe_pgid, cmd, data); 845 846 } 847 return EPASSTHROUGH; 848 } 849 850 int 851 pipe_poll(file_t *fp, int events) 852 { 853 struct pipe *pipe = fp->f_pipe; 854 struct pipe *ppipe; 855 int revents = 0; 856 857 mutex_enter(pipe->pipe_lock); 858 ppipe = pipe->pipe_peer; 859 860 if (fp->f_flag & FREAD) { 861 struct pipe *rpipe = pipe; 862 863 /* 864 * If the writer has been closed, then we can always 865 * read (possibly returning EOF) without blocking, so 866 * set POLLIN|POLLRDNORM if requested, and set POLLHUP 867 * unsolicited to notify reader of the fact. 868 * 869 * Otherwise, we can only read without blocking if 870 * there are bytes in the buffer. 871 */ 872 if (rpipe->pipe_state & PIPE_EOF) { 873 revents |= events & (POLLIN | POLLRDNORM); 874 revents |= POLLHUP; 875 } else if (rpipe->pipe_buffer.cnt > 0) { 876 revents |= events & (POLLIN | POLLRDNORM); 877 } 878 } else if (fp->f_flag & FWRITE) { 879 struct pipe *wpipe = pipe; 880 struct pipe *rpipe = ppipe; 881 882 /* 883 * If the reader has been closed, then any writes will 884 * immediately fail with EPIPE, so report 885 * POLLOUT|POLLWRNORM if requested and POLLERR 886 * unsolicited. 887 * 888 * Otherwise, we can only write without blocking if 889 * there are at least PIPE_BUF bytes free in the 890 * buffer. 891 */ 892 if (rpipe == NULL || (wpipe->pipe_state & PIPE_EOF) != 0) { 893 revents |= events & (POLLOUT | POLLWRNORM); 894 revents |= POLLERR; 895 } else if (rpipe->pipe_buffer.size - rpipe->pipe_buffer.cnt >= 896 PIPE_BUF) { 897 revents |= events & (POLLOUT | POLLWRNORM); 898 } 899 } else { 900 panic("file %p pipe %p invalid direction flag 0x%x", 901 fp, pipe, fp->f_flag); 902 } 903 904 if (revents == 0) 905 selrecord(curlwp, &pipe->pipe_sel); 906 mutex_exit(pipe->pipe_lock); 907 908 return revents; 909 } 910 911 static int 912 pipe_stat(file_t *fp, struct stat *ub) 913 { 914 struct pipe *pipe = fp->f_pipe; 915 916 mutex_enter(pipe->pipe_lock); 917 memset(ub, 0, sizeof(*ub)); 918 ub->st_mode = S_IFIFO | S_IRUSR | S_IWUSR; 919 ub->st_blksize = pipe->pipe_buffer.size; 920 if (ub->st_blksize == 0 && pipe->pipe_peer) 921 ub->st_blksize = pipe->pipe_peer->pipe_buffer.size; 922 ub->st_size = pipe->pipe_buffer.cnt; 923 ub->st_blocks = (ub->st_size) ? 1 : 0; 924 ub->st_atimespec = pipe->pipe_atime; 925 ub->st_mtimespec = pipe->pipe_mtime; 926 ub->st_ctimespec = ub->st_birthtimespec = pipe->pipe_btime; 927 ub->st_uid = kauth_cred_geteuid(fp->f_cred); 928 ub->st_gid = kauth_cred_getegid(fp->f_cred); 929 930 /* 931 * Left as 0: st_dev, st_ino, st_nlink, st_rdev, st_flags, st_gen. 932 * XXX (st_dev, st_ino) should be unique. 933 */ 934 mutex_exit(pipe->pipe_lock); 935 return 0; 936 } 937 938 static int 939 pipe_close(file_t *fp) 940 { 941 struct pipe *pipe = fp->f_pipe; 942 943 fp->f_pipe = NULL; 944 pipeclose(fp, pipe); 945 return 0; 946 } 947 948 static void 949 pipe_restart(file_t *fp) 950 { 951 struct pipe *pipe = fp->f_pipe; 952 953 /* 954 * Unblock blocked reads/writes in order to allow close() to complete. 955 * System calls return ERESTART so that the fd is revalidated. 956 * (Partial writes return the transfer length.) 957 */ 958 mutex_enter(pipe->pipe_lock); 959 pipe->pipe_state |= PIPE_RESTART; 960 /* 961 * At most one of these is in use at any time, depending on 962 * whether fp->f_flag has FREAD or FWRITE set, but there's no 963 * harm in waking both here. 964 */ 965 cv_broadcast(&pipe->pipe_rcv); 966 cv_broadcast(&pipe->pipe_wcv); 967 mutex_exit(pipe->pipe_lock); 968 } 969 970 static int 971 pipe_fpathconf(struct file *fp, int name, register_t *retval) 972 { 973 974 switch (name) { 975 case _PC_PIPE_BUF: 976 *retval = PIPE_BUF; 977 return 0; 978 default: 979 return SET_ERROR(EINVAL); 980 } 981 } 982 983 static int 984 pipe_posix_fadvise(struct file *fp, off_t offset, off_t len, int advice) 985 { 986 987 return SET_ERROR(ESPIPE); 988 } 989 990 static void 991 pipe_free_kmem(struct pipe *pipe) 992 { 993 994 if (pipe->pipe_buffer.buffer != NULL) { 995 if (pipe->pipe_buffer.size > PIPE_SIZE) { 996 atomic_dec_uint(&nbigpipe); 997 } 998 if (pipe->pipe_buffer.buffer != (void *)pipe->pipe_kmem) { 999 uvm_km_free(kernel_map, 1000 (vaddr_t)pipe->pipe_buffer.buffer, 1001 pipe->pipe_buffer.size, UVM_KMF_PAGEABLE); 1002 atomic_add_int(&amountpipekva, 1003 -pipe->pipe_buffer.size); 1004 } 1005 pipe->pipe_buffer.buffer = NULL; 1006 } 1007 } 1008 1009 /* 1010 * Shutdown the pipe. 1011 */ 1012 static void 1013 pipeclose(struct file *fp, struct pipe *pipe) 1014 { 1015 kmutex_t *lock; 1016 struct pipe *ppipe; 1017 1018 KASSERT(cv_is_valid(&pipe->pipe_rcv)); 1019 KASSERT(cv_is_valid(&pipe->pipe_wcv)); 1020 KASSERT(cv_is_valid(&pipe->pipe_draincv)); 1021 KASSERT(cv_is_valid(&pipe->pipe_lkcv)); 1022 1023 lock = pipe->pipe_lock; 1024 KASSERT(lock != NULL); 1025 1026 mutex_enter(lock); 1027 1028 /* 1029 * fd_close has issued .fo_restart to wake all waiters on this 1030 * side of the pipe, blocked new references, and waited for all 1031 * references to drain, so it should not be possible for there 1032 * to be any waiters remaining. (Only one of the condvars was 1033 * ever in use anyway depending on whether this is the reader 1034 * side or the writer side of the pipe.) 1035 */ 1036 KASSERT(!cv_has_waiters(&pipe->pipe_rcv)); 1037 KASSERT(!cv_has_waiters(&pipe->pipe_wcv)); 1038 1039 /* 1040 * There may, however, be threads waiting in select/poll for 1041 * I/O to be ready on this side of the pipe. Wake them (but 1042 * don't send SIGIO as pipeselwakeup does) so they can fail 1043 * with EBADF/POLLNVAL. 1044 */ 1045 selnotify(&pipe->pipe_sel, 0, NOTE_SUBMIT); 1046 1047 /* 1048 * If the other side is busy, wake it up saying that 1049 * we want to close it down, which will prevent peers 1050 * from starting new I/O. Once it is no longer busy, 1051 * disconnect it. 1052 */ 1053 KASSERT(pipe->pipe_peer != NULL || (pipe->pipe_state & PIPE_EOF) != 0); 1054 pipe->pipe_state |= PIPE_EOF; 1055 if ((ppipe = pipe->pipe_peer) != NULL) { 1056 if (fp->f_flag & FREAD) { 1057 struct pipe *wpipe = ppipe; 1058 1059 pipeselwakeup(wpipe, POLL_ERR); 1060 } else if (fp->f_flag & FWRITE) { 1061 struct pipe *rpipe = ppipe; 1062 1063 pipeselwakeup(rpipe, POLL_HUP); 1064 } else { 1065 panic("file %p pipe %p invalid direction flag 0x%x", 1066 fp, pipe, fp->f_flag); 1067 } 1068 1069 ppipe->pipe_state |= PIPE_EOF; 1070 if (ppipe->pipe_busy) { 1071 cv_broadcast(&ppipe->pipe_rcv); 1072 cv_broadcast(&ppipe->pipe_wcv); 1073 while (ppipe->pipe_busy) { 1074 cv_wait(&ppipe->pipe_draincv, lock); 1075 1076 /* 1077 * After the cv_wait, another thread 1078 * may have concurrently closed ppipe, 1079 * with two effects: 1080 * 1081 * 1. ppipe may now be invalid, so we 1082 * MUST NOT touch it. 1083 * 1084 * 2. pipe->pipe_peer may have been set 1085 * to null (under the common mutex), 1086 * so we can detect this case. 1087 */ 1088 KASSERT(pipe->pipe_peer == NULL || 1089 pipe->pipe_peer == ppipe); 1090 if ((ppipe = pipe->pipe_peer) == NULL) 1091 break; 1092 } 1093 } 1094 if (ppipe) 1095 ppipe->pipe_peer = NULL; 1096 } 1097 1098 /* 1099 * Any knote objects still left in the list are 1100 * the one attached by peer. Since no one will 1101 * traverse this list, we just clear it. 1102 * 1103 * XXX Exposes select/kqueue internals. 1104 */ 1105 SLIST_INIT(&pipe->pipe_sel.sel_klist); 1106 1107 KASSERT((pipe->pipe_state & PIPE_LOCKFL) == 0); 1108 mutex_exit(lock); 1109 1110 /* 1111 * Free resources. 1112 */ 1113 pipefree(pipe); 1114 } 1115 1116 static void 1117 pipefree(struct pipe *pipe) 1118 { 1119 1120 pipe->pipe_pgid = 0; 1121 pipe->pipe_state = PIPE_SIGNALR; 1122 pipe->pipe_peer = NULL; 1123 mutex_obj_free(pipe->pipe_lock); 1124 pipe->pipe_lock = NULL; 1125 pipe_free_kmem(pipe); 1126 if (pipe->pipe_kmem != 0) { 1127 pool_cache_put(pipe_rd_cache, pipe); 1128 } else { 1129 pool_cache_put(pipe_wr_cache, pipe); 1130 } 1131 } 1132 1133 static void 1134 filt_pipedetach(struct knote *kn) 1135 { 1136 struct pipe *pipe; 1137 kmutex_t *lock; 1138 1139 pipe = ((file_t *)kn->kn_obj)->f_pipe; 1140 lock = pipe->pipe_lock; 1141 1142 mutex_enter(lock); 1143 KASSERT(kn->kn_hook == pipe); 1144 selremove_knote(&pipe->pipe_sel, kn); 1145 mutex_exit(lock); 1146 } 1147 1148 static void 1149 filt_pipenodetach(struct knote *kn) 1150 { 1151 /* not attached, nothing to do */ 1152 } 1153 1154 static int 1155 filt_pipewrongend(struct knote *kn, long hint) 1156 { 1157 1158 /* Never ready! */ 1159 return 0; 1160 } 1161 1162 static int 1163 filt_piperead(struct knote *kn, long hint) 1164 { 1165 struct pipe *rpipe = ((file_t *)kn->kn_obj)->f_pipe; 1166 struct pipe *wpipe; 1167 int rv; 1168 1169 if ((hint & NOTE_SUBMIT) == 0) { 1170 mutex_enter(rpipe->pipe_lock); 1171 } else { 1172 KASSERT(mutex_owned(rpipe->pipe_lock)); 1173 } 1174 wpipe = rpipe->pipe_peer; 1175 kn->kn_data = rpipe->pipe_buffer.cnt; 1176 1177 if ((rpipe->pipe_state & PIPE_EOF) || 1178 (wpipe == NULL) || (wpipe->pipe_state & PIPE_EOF)) { 1179 knote_set_eof(kn, 0); 1180 rv = 1; 1181 } else { 1182 rv = kn->kn_data > 0; 1183 } 1184 1185 if ((hint & NOTE_SUBMIT) == 0) { 1186 mutex_exit(rpipe->pipe_lock); 1187 } else { 1188 KASSERT(mutex_owned(rpipe->pipe_lock)); 1189 } 1190 return rv; 1191 } 1192 1193 static int 1194 filt_pipewrite(struct knote *kn, long hint) 1195 { 1196 struct pipe *wpipe = ((file_t *)kn->kn_obj)->f_pipe; 1197 struct pipe *rpipe; 1198 int rv; 1199 1200 if ((hint & NOTE_SUBMIT) == 0) { 1201 mutex_enter(wpipe->pipe_lock); 1202 } else { 1203 KASSERT(mutex_owned(wpipe->pipe_lock)); 1204 } 1205 rpipe = wpipe->pipe_peer; 1206 1207 if ((rpipe == NULL) || (rpipe->pipe_state & PIPE_EOF)) { 1208 kn->kn_data = 0; 1209 knote_set_eof(kn, 0); 1210 rv = 1; 1211 } else { 1212 kn->kn_data = rpipe->pipe_buffer.size - rpipe->pipe_buffer.cnt; 1213 rv = kn->kn_data >= PIPE_BUF; 1214 } 1215 1216 if ((hint & NOTE_SUBMIT) == 0) { 1217 mutex_exit(wpipe->pipe_lock); 1218 } else { 1219 KASSERT(mutex_owned(wpipe->pipe_lock)); 1220 } 1221 return rv; 1222 } 1223 1224 static const struct filterops pipe_wrongendfiltops = { 1225 .f_flags = FILTEROP_ISFD | FILTEROP_MPSAFE, 1226 .f_attach = NULL, 1227 .f_detach = filt_pipenodetach, 1228 .f_event = filt_pipewrongend, 1229 }; 1230 1231 static const struct filterops pipe_rfiltops = { 1232 .f_flags = FILTEROP_ISFD | FILTEROP_MPSAFE, 1233 .f_attach = NULL, 1234 .f_detach = filt_pipedetach, 1235 .f_event = filt_piperead, 1236 }; 1237 1238 static const struct filterops pipe_wfiltops = { 1239 .f_flags = FILTEROP_ISFD | FILTEROP_MPSAFE, 1240 .f_attach = NULL, 1241 .f_detach = filt_pipedetach, 1242 .f_event = filt_pipewrite, 1243 }; 1244 1245 static int 1246 pipe_kqfilter(file_t *fp, struct knote *kn) 1247 { 1248 struct pipe *pipe; 1249 kmutex_t *lock; 1250 1251 pipe = ((file_t *)kn->kn_obj)->f_pipe; 1252 lock = pipe->pipe_lock; 1253 1254 mutex_enter(lock); 1255 1256 switch (kn->kn_filter) { 1257 case EVFILT_READ: 1258 if ((fp->f_flag & FREAD) == 0) { 1259 kn->kn_fop = &pipe_wrongendfiltops; 1260 mutex_exit(lock); 1261 return 0; 1262 } 1263 kn->kn_fop = &pipe_rfiltops; 1264 break; 1265 case EVFILT_WRITE: 1266 if ((fp->f_flag & FWRITE) == 0) { 1267 kn->kn_fop = &pipe_wrongendfiltops; 1268 mutex_exit(lock); 1269 return 0; 1270 } 1271 kn->kn_fop = &pipe_wfiltops; 1272 break; 1273 default: 1274 mutex_exit(lock); 1275 return SET_ERROR(EINVAL); 1276 } 1277 1278 kn->kn_hook = pipe; 1279 selrecord_knote(&pipe->pipe_sel, kn); 1280 mutex_exit(lock); 1281 1282 return 0; 1283 } 1284 1285 /* 1286 * Handle pipe sysctls. 1287 */ 1288 SYSCTL_SETUP(sysctl_kern_pipe_setup, "sysctl kern.pipe subtree setup") 1289 { 1290 1291 sysctl_createv(clog, 0, NULL, NULL, 1292 CTLFLAG_PERMANENT, 1293 CTLTYPE_NODE, "pipe", 1294 SYSCTL_DESCR("Pipe settings"), 1295 NULL, 0, NULL, 0, 1296 CTL_KERN, KERN_PIPE, CTL_EOL); 1297 1298 sysctl_createv(clog, 0, NULL, NULL, 1299 CTLFLAG_PERMANENT|CTLFLAG_READWRITE, 1300 CTLTYPE_INT, "maxbigpipes", 1301 SYSCTL_DESCR("Maximum number of \"big\" pipes"), 1302 NULL, 0, &maxbigpipes, 0, 1303 CTL_KERN, KERN_PIPE, KERN_PIPE_MAXBIGPIPES, CTL_EOL); 1304 sysctl_createv(clog, 0, NULL, NULL, 1305 CTLFLAG_PERMANENT, 1306 CTLTYPE_INT, "nbigpipes", 1307 SYSCTL_DESCR("Number of \"big\" pipes"), 1308 NULL, 0, &nbigpipe, 0, 1309 CTL_KERN, KERN_PIPE, KERN_PIPE_NBIGPIPES, CTL_EOL); 1310 sysctl_createv(clog, 0, NULL, NULL, 1311 CTLFLAG_PERMANENT, 1312 CTLTYPE_INT, "kvasize", 1313 SYSCTL_DESCR("Amount of kernel memory consumed by pipe " 1314 "buffers"), 1315 NULL, 0, &amountpipekva, 0, 1316 CTL_KERN, KERN_PIPE, KERN_PIPE_KVASIZE, CTL_EOL); 1317 } 1318