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