From mboxrd@z Thu Jan 1 00:00:00 1970 Received: from mail-wm2-f9.google.com (mail-wm2-f9.google.com [74.125.225.137]) (using TLSv1.2 with cipher ECDHE-RSA-AES128-GCM-SHA256 (128/128 bits)) (No client certificate requested) by smtp.subspace.kernel.org (Postfix) with ESMTPS id 16BB0395AF0 for ; Fri, 25 Sep 2026 04:55:46 +0000 (UTC) Authentication-Results: smtp.subspace.kernel.org; arc=none smtp.client-ip=74.125.225.137 ARC-Seal:i=1; a=rsa-sha256; d=subspace.kernel.org; s=arc-20240116; t=1790312149; cv=none; b=ZFNSLNhDCFiuaBq+TZFCHJI+g18Z/GeqVPgjNSKmeKmX0y94qltngPeQZO8rFPoi2zNgCy/MH0PfqC1swH4QKTd+qfucIzmkYhdYLM8Kp3y1IX8xQWbynwkSHAM2wlVKETiViyKlUXdvegLRVGIydVAv3uE7FDoLERyO6dbL3HY= ARC-Message-Signature:i=1; a=rsa-sha256; d=subspace.kernel.org; s=arc-20240116; t=1790312149; c=relaxed/simple; bh=RmtRCL9EMKDN+kHvQhS2FNW6gwOmjNuVjs4gbV7DAXc=; h=From:To:Cc:Subject:Date:Message-ID:In-Reply-To:References: MIME-Version; b=AWGY2kh/ecEWd4c2JH375p8sr0GbZRBEh1Gw0LTuYhvbCQYOQdN4Q714MMgdFXH0IuBomPfcZVaJOMVAbOJ8ZhgLbaZCoDw8GubhsDTIgWlFFzP6vCnaDmKWYBiXUHOtdIdZo/ZdsPO4vzYJap98c+3xGSlvPxHlg9vL0w4QXEM= ARC-Authentication-Results:i=1; smtp.subspace.kernel.org; dmarc=pass (p=none dis=none) header.from=gmail.com; spf=pass smtp.mailfrom=gmail.com; dkim=pass (2048-bit key) header.d=gmail.com header.i=@gmail.com header.b=NtnuY81X; arc=none smtp.client-ip=74.125.225.137 Authentication-Results: smtp.subspace.kernel.org; dmarc=pass (p=none dis=none) header.from=gmail.com Authentication-Results: smtp.subspace.kernel.org; spf=pass smtp.mailfrom=gmail.com Authentication-Results: smtp.subspace.kernel.org; dkim=pass (2048-bit key) header.d=gmail.com header.i=@gmail.com header.b="NtnuY81X" Received: by mail-wm2-f9.google.com with SMTP id 5b1f17b1804b1-49e6e43b9d8so1724785e9.1 for ; Thu, 24 Sep 2026 21:55:46 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1790312145; x=1790916945; darn=vger.kernel.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=xlBGJFUG1mLcTy/zKM82sg2p3GnxPZeP72XcfVV+1Rg=; b=NtnuY81Xsm9a/keyD/mva6h817Z2Fo58bCxK/pRmljsYDd2ZBPAaV4pgpCFU/eoGfD DG0Vzg+C3rNVtcZOvaZvNN9BRJftKvgAToO+4Mej1RfrfNDt5w3UNwW9giHZKzqYWDxU L8vGCZ9goa0bHtsJq8VG8YITgGIDeL9YvZPxKWaz05M0d2cOG2WNoNOqmpWwCoYh0sEK W951PABK8FuteqzgFwkWzkd5u//YtTFqy/8GLV/48xxuulInuXnU/1DlAWHS4hKWM/2p ODgXTUSVnIXMWtG99y1bNfR6kH59umdIbbCnY9CECcuiGqJks77zR/llRDYUBy31Jx5r Lckw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20260707; t=1790312145; x=1790916945; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=xlBGJFUG1mLcTy/zKM82sg2p3GnxPZeP72XcfVV+1Rg=; b=vf0gLsUP3G2FHl2OKpfygm3qo3KDZGaxtxGt8CI4uddlZ0ExXVLelTx20QschrH9TT NSqC31SGWQCtJoQCMP8UYqU50bEILkxLCQ1aqX9OT+tIM6mrIELOyKNeGXaiFpQeY3Uz kvKaB3uIs1yo48bOZLRaB4crQv/YMwZ5apl2ftDcgsdOrwW37C9bV5surO/H36o4dTNt NM6/OoxyObn72tolaUXI5JOuKm7h6oH+J08OJ/NyX1Bol0yeYqvHk8ZTFYRZoeIzRZp/ OAiQ1J9arz+tZTb/Wxjhk0TzZzWgUSi2172QuJ3KL9Ns0IWl46Xjf9GsNiDU0isqbEHw uMYQ== X-Gm-Message-State: AFuF++kQ91Oer4Gzw9U5hSRNhtUb7606H6SHPvUKYtUCofsgR4GfOk/S TeD7JoCsXLfy1G/HaT7KdNX4zFM29jc2pj1ybVxDdPVjJdVI4EA62gdn8hFG+X7ugF4= X-Gm-Gg: AYBFou2fyGPp11elUREoM21dqUAzTacBz2T8KP+szD3KP2k1mVeRlG1GbWvnhfJ5IUS s5FNIjbgdAOSaGGxtIJeFXmRShpXcKHgg6/yX19quKihJ2lL01Xa7ICu/n4Yn4PM1foi/UcbUjf vYPCCR9tTum8inOkVfgXT0q6DG79CSrUUCD1HyeRiFLp9tK96/K6Ho1RSDGiRnhUQyjBBhNple6 H4w1Y5Ig2GsXOeC6Wm9VlDCOzIkpFX5j/sUrE3PhgyH9K3ZWLPFI+bilg74BmeWbQCKLRjr0tVv ZApmo+gz63z32M/x0DE5n2NPjfeP5t5xia8eJDYk6FcRskqlEq8V4yVvRMZEtxmzJo0w/uxmAi9 wrxKCk/D0OKOgdAd/3aHplZC7g6dL/gNhEEwa+xHFb7Nl2mTp2gPE8RrtF97OyIaexCyghbvCXn atBXAuAf+5sgJnRunreKrLyDZm40FJphVDZ8Rf/SnPdtqBbDMgrId9Wq5PMJRFBkO7LKtpDk74X iWsy9EyEeIgin3tBs5xUurHn7uPwhuv1MrrqV7EnirbwJoaRDcw0q9YRwapHetVBH51uzstlPro BeLNRcMnSUrEG8oJurS77EviLhqaCLRzDvMVog== X-Received: by 2002:a05:600c:1389:b0:49f:dcc8:c735 with SMTP id 5b1f17b1804b1-49fe66f4c7dmr73162095e9.18.1790312144944; Thu, 24 Sep 2026 21:55:44 -0700 (PDT) Received: from localhost (nat-icclus-192-26-29-3.epfl.ch. [192.26.29.3]) by smtp.gmail.com with ESMTPSA id 5b1f17b1804b1-49ff06b6108sm33956605e9.8.2026.09.24.21.55.43 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 24 Sep 2026 21:55:43 -0700 (PDT) From: Kumar Kartikeya Dwivedi To: bpf@vger.kernel.org Cc: Alexei Starovoitov , Andrii Nakryiko , Daniel Borkmann , Eduard Zingerman , Emil Tsalapatis , Tejun Heo , kkd@meta.com, kernel-team@meta.com Subject: [PATCH bpf-next v3 2/5] bpf: Add file descriptor interface for program streams Date: Fri, 25 Sep 2026 06:55:29 +0200 Message-ID: <20260925045536.1480933-3-memxor@gmail.com> X-Mailer: git-send-email 2.53.0 In-Reply-To: <20260925045536.1480933-1-memxor@gmail.com> References: <20260925045536.1480933-1-memxor@gmail.com> Precedence: bulk X-Mailing-List: bpf@vger.kernel.org List-Id: List-Subscribe: List-Unsubscribe: MIME-Version: 1.0 X-Developer-Signature: v=1; a=openpgp-sha256; l=23815; i=memxor@gmail.com; h=from:subject; bh=RmtRCL9EMKDN+kHvQhS2FNW6gwOmjNuVjs4gbV7DAXc=; b=owGbwMvMwCXmrmtenRyi38x4Wi2JIWvrvym2jxrN2WcG/WdXWyqtPOFnxe3vEu5i/rmcRs8kX 81Zc46xo5SFQYyLQVZMkaXk/z4m4xOVvwNtl3HDzGFlAhnCwMUpABNxNmL4n8t/THHJjfldskG/ xdsVuKUEeT6/NVOZIXp2En/Omvjj9xn+h7VyMmoeYOFLurqYxVyxaYtst5Dvt10zJ/UG/JhxJvc qLwA= X-Developer-Key: i=memxor@gmail.com; a=openpgp; fpr=B34BD741DE8494B76E2F717880EF20021D46C59B Content-Transfer-Encoding: 8bit The existing BPF_PROG_STREAM_READ_BY_FD command only supports polling a program stream through repeated bpf() calls. It cannot block for new data or integrate with poll-based event loops. Add BPF_PROG_STREAM_OPEN to return a read-only, close-on-exec file descriptor for a selected program stream. Reads block by default and BPF_F_STREAM_NONBLOCK, the only accepted flag, provides non-blocking behavior. poll reports readable data and reports hangup once the program has been freed. Like pipes and sockets, the descriptor is not seekable and lseek fails with ESPIPE. A stream descriptor deliberately does not retain the program. Move each stream into a separately refcounted allocation so program teardown can mark it dead and wake descriptor users while outstanding descriptors drain buffered data safely. Readers sample the dead flag before looking for data, so EOF is reported only when the stream was already dead before it was found empty; data published right before teardown is never skipped. Only programs loaded through BPF_PROG_LOAD get streams. Classic BPF filters, JIT subprograms and shim programs never write to one, and kernel-side writers already resolve a subprogram to its main program, so those programs no longer carry stream state. Readiness needs its own counter. Stream capacity is charged before allocation and before an element is published to the stream log, so using that reservation as the read and poll condition can report readable data while no element exists: a blocking reader retries instead of sleeping and a lone non-blocking reader can see POLLIN followed by EAGAIN. Publish bytes with release ordering after adding elements to the lockless log, use acquire loads before consuming them or reporting readiness, limit each read to its readable snapshot and subtract only bytes actually copied. This keeps the aggregate count correct even when concurrent publishers update it out of publication order. With several readers on one stream, readiness remains advisory, as it is for pipes. The capacity counter is kept solely for enforcing the stream size limit. Wakeups are always deferred through irq_work. Stream writers run in whatever context the program runs in: NMI context for perf_event programs, sections with interrupts disabled inside bpf_spin_lock or rqspinlock critical sections since bpf_stream_vprintk() is KF_SPINLOCK_SAFE, and tracing programs attached anywhere in the kernel, including inside the wait queue and epoll code itself. Waking waiters directly from there can deadlock, and no cheap context check covers every case: on PREEMPT_RT, spinlock_t sections do not disable interrupts, so in_nmi() or irqs_disabled() cannot tell such a program apart from a benign one. Queue an irq_work item instead, as bpf_ringbuf does. Queue it only when a publication turns an empty stream readable. Readers block and pollers wait only after finding the stream empty, and the readable count never drops below zero because each read is bounded by its snapshot, so the first publication after such an observation is the one that makes the count positive, and it is the one that queues the wakeup. Publications into a stream that already holds data raise no interrupt, so a program that prints while nobody drains its stream pays for a single irq_work until the stream is emptied again. This matches bpf_ringbuf, which notifies only once the consumer has caught up. Blocking readers and level-triggered pollers re-check the readable count before waiting, so they cannot miss data, and edge-triggered epoll consumers drain until EAGAIN before waiting again, as epoll(7) requires. Synchronize pending work before releasing the final stream reference so the callback cannot outlive the stream, but only when the work was ever queued: irq_work_sync() waits for an RCU grace period on PREEMPT_RT and on architectures without an irq_work interrupt. Signed-off-by: Kumar Kartikeya Dwivedi --- include/linux/bpf.h | 14 ++- include/uapi/linux/bpf.h | 39 ++++++ kernel/bpf/core.c | 8 +- kernel/bpf/stream.c | 209 ++++++++++++++++++++++++++++++--- kernel/bpf/syscall.c | 29 +++++ tools/include/uapi/linux/bpf.h | 39 ++++++ 6 files changed, 315 insertions(+), 23 deletions(-) diff --git a/include/linux/bpf.h b/include/linux/bpf.h index 8594c8aff745..b2b4a769bcca 100644 --- a/include/linux/bpf.h +++ b/include/linux/bpf.h @@ -17,6 +17,7 @@ #include #include #include +#include #include #include #include @@ -1771,12 +1772,18 @@ enum { }; struct bpf_stream { - atomic_t capacity; + refcount_t refcnt; + atomic_t capacity; /* bytes reserved against the stream limit */ + atomic_t readable; /* published bytes available to readers */ struct llist_head log; /* list of in-flight stream elements in LIFO order */ struct mutex lock; /* lock protecting backlog_{head,tail} */ struct llist_node *backlog_head; /* list of in-flight stream elements in FIFO order */ struct llist_node *backlog_tail; /* tail of the list above */ + wait_queue_head_t waitq; + struct irq_work notify_work; + bool notify_used; /* notify_work was queued at least once */ + bool dead; }; struct bpf_stream_stage { @@ -1915,7 +1922,7 @@ struct bpf_prog_aux { struct work_struct work; struct rcu_head rcu; }; - struct bpf_stream stream[2]; + struct bpf_stream *stream[2]; struct mutex st_ops_assoc_mutex; struct bpf_map __rcu *st_ops_assoc; }; @@ -4224,9 +4231,10 @@ void bpf_bprintf_cleanup(struct bpf_bprintf_data *data); int bpf_try_get_buffers(struct bpf_bprintf_buffers **bufs); void bpf_put_buffers(void); -void bpf_prog_stream_init(struct bpf_prog *prog); +int bpf_prog_stream_init(struct bpf_prog *prog, gfp_t gfp_extra_flags); void bpf_prog_stream_free(struct bpf_prog *prog); int bpf_prog_stream_read(struct bpf_prog *prog, enum bpf_stream_id stream_id, void __user *buf, u32 len); +int bpf_prog_stream_new_fd(struct bpf_prog *prog, enum bpf_stream_id stream_id, u32 flags); void bpf_stream_stage_init(struct bpf_stream_stage *ss); void bpf_stream_stage_free(struct bpf_stream_stage *ss); __printf(2, 3) diff --git a/include/uapi/linux/bpf.h b/include/uapi/linux/bpf.h index 0aaa54359aeb..4687c3310996 100644 --- a/include/uapi/linux/bpf.h +++ b/include/uapi/linux/bpf.h @@ -936,6 +936,33 @@ union bpf_iter_link_info { * 0 on success or -1 if an error occurred (in which case, * *errno* is set appropriately). * + * BPF_PROG_STREAM_OPEN + * Description + * Open a file descriptor for one of the BPF streams associated + * with the program identified by *prog_fd*. The stream is selected + * by *stream_id*. + * + * The returned file descriptor supports **read**\ (2) and + * **poll**\ (2). Reads block while the stream is empty unless + * **BPF_F_STREAM_NONBLOCK** is specified in *flags*. A non-blocking + * read of an empty stream fails with **EAGAIN**. + * + * **poll**\ (2) reports **POLLIN** when data is available and + * **POLLHUP** once the program has been freed, that is, after every + * reference to it, including links and other file descriptors, has + * been dropped. Hangup may lag the final release because program + * teardown is deferred. Buffered data remains readable after + * **POLLHUP** and a read returns zero after all such data has been + * consumed. + * + * The file descriptor is read-only and has the close-on-exec flag + * set. It is not seekable and **lseek**\ (2) fails with **ESPIPE**. + * *flags* may only contain **BPF_F_STREAM_NONBLOCK**. + * + * Return + * A new file descriptor (a nonnegative integer), or -1 if an + * error occurred (in which case, *errno* is set appropriately). + * * NOTES * eBPF objects (maps and programs) can be shared between processes. * @@ -993,6 +1020,7 @@ enum bpf_cmd { BPF_TOKEN_CREATE, BPF_PROG_STREAM_READ_BY_FD, BPF_PROG_ASSOC_STRUCT_OPS, + BPF_PROG_STREAM_OPEN, __MAX_BPF_CMD, BPF_COMMON_ATTRS = 1 << 16, /* Indicate carrying syscall common attrs. */ }; @@ -1524,6 +1552,11 @@ enum { BPF_STREAM_STDERR = 2, }; +/* flags for BPF_PROG_STREAM_OPEN command */ +enum { + BPF_F_STREAM_NONBLOCK = (1U << 0), +}; + union bpf_attr { struct { /* anonymous struct used by BPF_MAP_CREATE command */ __u32 map_type; /* one of enum bpf_map_type */ @@ -1950,6 +1983,12 @@ union bpf_attr { __u32 flags; } prog_assoc_struct_ops; + struct { + __u32 prog_fd; + __u32 stream_id; + __u32 flags; + } prog_stream_open; + } __attribute__((aligned(8))); /* The description below is an attempt at providing documentation to eBPF diff --git a/kernel/bpf/core.c b/kernel/bpf/core.c index aa72ad3ec689..d3b8b626ec0f 100644 --- a/kernel/bpf/core.c +++ b/kernel/bpf/core.c @@ -142,10 +142,6 @@ struct bpf_prog *bpf_prog_alloc_no_stats(unsigned int size, gfp_t gfp_extra_flag mutex_init(&fp->aux->dst_mutex); mutex_init(&fp->aux->st_ops_assoc_mutex); -#ifdef CONFIG_BPF_SYSCALL - bpf_prog_stream_init(fp); -#endif - return fp; } @@ -289,6 +285,9 @@ struct bpf_prog *bpf_prog_realloc(struct bpf_prog *fp_old, unsigned int size, void __bpf_prog_free(struct bpf_prog *fp) { if (fp->aux) { +#ifdef CONFIG_BPF_SYSCALL + bpf_prog_stream_free(fp); +#endif mutex_destroy(&fp->aux->used_maps_mutex); mutex_destroy(&fp->aux->dst_mutex); mutex_destroy(&fp->aux->st_ops_assoc_mutex); @@ -3086,7 +3085,6 @@ static void bpf_prog_free_deferred(struct work_struct *work) aux = container_of(work, struct bpf_prog_aux, work); #ifdef CONFIG_BPF_SYSCALL bpf_free_kfunc_btf_tab(aux->kfunc_btf_tab); - bpf_prog_stream_free(aux->prog); #endif #ifdef CONFIG_CGROUP_BPF if (aux->cgroup_atype != CGROUP_BPF_ATTACH_TYPE_INVALID) diff --git a/kernel/bpf/stream.c b/kernel/bpf/stream.c index 9829d1ebce70..83e803ea54ef 100644 --- a/kernel/bpf/stream.c +++ b/kernel/bpf/stream.c @@ -2,11 +2,15 @@ /* Copyright (c) 2025 Meta Platforms, Inc. and affiliates. */ #include +#include #include #include #include +#include #include #include +#include +#include static void bpf_stream_elem_init(struct bpf_stream_elem *elem, int len) { @@ -73,6 +77,51 @@ static void bpf_stream_release_capacity(struct bpf_stream *stream, int len) atomic_sub(len, &stream->capacity); } +static void bpf_stream_notify(struct irq_work *work) +{ + struct bpf_stream *stream = container_of(work, struct bpf_stream, notify_work); + + /* + * Writers run in arbitrary program contexts, including NMI and regions + * that already hold wait queue or epoll locks. Wake waiters from + * irq_work instead, where taking those locks is safe. + */ + wake_up_interruptible_poll(&stream->waitq, EPOLLIN | EPOLLRDNORM); +} + +static void bpf_stream_queue_notify(struct bpf_stream *stream) +{ + /* + * Record that the work has been used so that teardown only pays for + * irq_work_sync(), which may wait for an RCU grace period, when a + * callback could actually be in flight. + */ + if (!READ_ONCE(stream->notify_used)) + WRITE_ONCE(stream->notify_used, true); + irq_work_queue(&stream->notify_work); +} + +static int bpf_stream_readable_bytes(struct bpf_stream *stream) +{ + return atomic_read_acquire(&stream->readable); +} + +static void bpf_stream_publish(struct bpf_stream *stream, int len) +{ + /* + * Pairs with atomic_read_acquire() in bpf_stream_readable_bytes(). + * + * Notify only when the stream turns from empty to readable. Readers + * block and pollers wait only after finding it empty, and a read never + * takes the count below zero, so the publication that makes it positive + * is the one they wait for. Publishing into a stream that already holds + * data would only raise an interrupt for waiters that are being woken + * already, or for nobody at all. + */ + if (atomic_add_return_release(len, &stream->readable) == len) + bpf_stream_queue_notify(stream); +} + static int bpf_stream_push_str(struct bpf_stream *stream, const char *str, int len) { int ret; @@ -88,6 +137,8 @@ static int bpf_stream_push_str(struct bpf_stream *stream, const char *str, int l ret = __bpf_stream_push_str(&stream->log, str, len); if (ret) bpf_stream_release_capacity(stream, len); + else + bpf_stream_publish(stream, len); return ret; } @@ -96,7 +147,7 @@ static struct bpf_stream *bpf_stream_get(enum bpf_stream_id stream_id, struct bp { if (stream_id != BPF_STDOUT && stream_id != BPF_STDERR) return NULL; - return &aux->stream[stream_id - 1]; + return aux->stream[stream_id - 1]; } static void bpf_stream_free_elem(struct bpf_stream_elem *elem) @@ -164,14 +215,16 @@ static bool bpf_stream_consume_elem(struct bpf_stream_elem *elem, int *len) static int bpf_stream_read(struct bpf_stream *stream, void __user *buf, int len) { - int rem_len = len, cons_len, ret = 0; + int read_len, rem_len, cons_len, ret = 0; struct bpf_stream_elem *elem = NULL; struct llist_node *node; mutex_lock(&stream->lock); + read_len = min(len, bpf_stream_readable_bytes(stream)); + rem_len = read_len; while (rem_len) { - int pos = len - rem_len; + int pos = read_len - rem_len; int chunk, n; bool cont; @@ -193,7 +246,7 @@ static int bpf_stream_read(struct bpf_stream *stream, void __user *buf, int len) /* Keep any successfully copied bytes; -EFAULT only if none. */ elem->consumed_len -= n; rem_len += n; - ret = (len == rem_len) ? -EFAULT : 0; + ret = (read_len == rem_len) ? -EFAULT : 0; break; } @@ -204,8 +257,9 @@ static int bpf_stream_read(struct bpf_stream *stream, void __user *buf, int len) bpf_stream_free_elem(elem); } + atomic_sub(read_len - rem_len, &stream->readable); mutex_unlock(&stream->lock); - return ret ? ret : len - rem_len; + return ret ? ret : read_len - rem_len; } int bpf_prog_stream_read(struct bpf_prog *prog, enum bpf_stream_id stream_id, void __user *buf, u32 len) @@ -220,6 +274,111 @@ int bpf_prog_stream_read(struct bpf_prog *prog, enum bpf_stream_id stream_id, vo return bpf_stream_read(stream, buf, len); } +static bool bpf_stream_has_data(struct bpf_stream *stream) +{ + return bpf_stream_readable_bytes(stream) > 0; +} + +static void bpf_stream_put(struct bpf_stream *stream) +{ + if (refcount_dec_and_test(&stream->refcnt)) { + struct llist_node *list; + + /* Only a stream that ever queued its work can have a callback in flight. */ + if (READ_ONCE(stream->notify_used)) + irq_work_sync(&stream->notify_work); + list = llist_del_all(&stream->log); + bpf_stream_free_list(list); + bpf_stream_free_list(stream->backlog_head); + mutex_destroy(&stream->lock); + kfree(stream); + } +} + +static int bpf_stream_release(struct inode *inode, struct file *file) +{ + bpf_stream_put(file->private_data); + return 0; +} + +static ssize_t bpf_stream_file_read(struct file *file, char __user *buf, size_t len, + loff_t *ppos) +{ + struct bpf_stream *stream = file->private_data; + bool dead; + int ret; + + if (!len) + return 0; + + for (;;) { + /* + * Sample teardown state before looking for data. Nothing is + * published once the program is gone, so finding the stream + * empty after observing dead means EOF. The opposite order could + * report EOF while data published just before teardown is still + * buffered. + */ + dead = smp_load_acquire(&stream->dead); + ret = bpf_stream_read(stream, buf, len); + if (ret) + return ret; + if (dead) + return 0; + if (file->f_flags & O_NONBLOCK) + return -EAGAIN; + + ret = wait_event_interruptible(stream->waitq, + bpf_stream_has_data(stream) || + READ_ONCE(stream->dead)); + if (ret) + return ret; + } +} + +static __poll_t bpf_stream_poll(struct file *file, struct poll_table_struct *pts) +{ + struct bpf_stream *stream = file->private_data; + __poll_t events = 0; + + /* + * poll_wait() only registers the wait queue callback. Register before + * checking persistent state so a concurrent publication or teardown is + * observed either by the callback or by the checks below. + */ + poll_wait(file, &stream->waitq, pts); + if (bpf_stream_has_data(stream)) + events |= EPOLLIN | EPOLLRDNORM; + if (READ_ONCE(stream->dead)) + events |= EPOLLHUP; + return events; +} + +static const struct file_operations bpf_stream_fops = { + .release = bpf_stream_release, + .read = bpf_stream_file_read, + .poll = bpf_stream_poll, +}; + +int bpf_prog_stream_new_fd(struct bpf_prog *prog, enum bpf_stream_id stream_id, u32 flags) +{ + struct bpf_stream *stream; + int fd_flags = O_RDONLY | O_CLOEXEC; + int fd; + + stream = bpf_stream_get(stream_id, prog->aux); + if (!stream) + return -ENOENT; + if (flags & BPF_F_STREAM_NONBLOCK) + fd_flags |= O_NONBLOCK; + + refcount_inc(&stream->refcnt); + fd = anon_inode_getfd("bpf-stream", &bpf_stream_fops, stream, fd_flags); + if (fd < 0) + bpf_stream_put(stream); + return fd; +} + __bpf_kfunc_start_defs(); /* @@ -287,28 +446,47 @@ __bpf_kfunc_end_defs(); /* Added kfunc to common_btf_ids */ -void bpf_prog_stream_init(struct bpf_prog *prog) +int bpf_prog_stream_init(struct bpf_prog *prog, gfp_t gfp_extra_flags) { int i; for (i = 0; i < ARRAY_SIZE(prog->aux->stream); i++) { - atomic_set(&prog->aux->stream[i].capacity, 0); - init_llist_head(&prog->aux->stream[i].log); - mutex_init(&prog->aux->stream[i].lock); - prog->aux->stream[i].backlog_head = NULL; - prog->aux->stream[i].backlog_tail = NULL; + struct bpf_stream *stream; + + /* On failure, bpf_prog_stream_free() releases the streams allocated so far. */ + stream = kzalloc_obj(*stream, + bpf_memcg_flags(GFP_KERNEL | gfp_extra_flags)); + if (!stream) + return -ENOMEM; + + refcount_set(&stream->refcnt, 1); + init_llist_head(&stream->log); + mutex_init(&stream->lock); + init_waitqueue_head(&stream->waitq); + init_irq_work(&stream->notify_work, bpf_stream_notify); + prog->aux->stream[i] = stream; } + return 0; } void bpf_prog_stream_free(struct bpf_prog *prog) { - struct llist_node *list; int i; for (i = 0; i < ARRAY_SIZE(prog->aux->stream); i++) { - list = llist_del_all(&prog->aux->stream[i].log); - bpf_stream_free_list(list); - bpf_stream_free_list(prog->aux->stream[i].backlog_head); + struct bpf_stream *stream = prog->aux->stream[i]; + + if (!stream) + continue; + /* + * Pairs with smp_load_acquire() in bpf_stream_file_read(): every + * publication precedes the dead flag, so a reader that observes + * it also observes all buffered data. + */ + smp_store_release(&stream->dead, true); + wake_up_interruptible_poll(&stream->waitq, EPOLLHUP); + bpf_stream_put(stream); + prog->aux->stream[i] = NULL; } } @@ -371,6 +549,7 @@ int bpf_stream_stage_commit(struct bpf_stream_stage *ss, struct bpf_prog *prog, list = tail; } llist_add_batch(head, tail, &stream->log); + bpf_stream_publish(stream, ss->len); return 0; } diff --git a/kernel/bpf/syscall.c b/kernel/bpf/syscall.c index e1bcfd43f54f..ac52f4ae414c 100644 --- a/kernel/bpf/syscall.c +++ b/kernel/bpf/syscall.c @@ -3080,6 +3080,10 @@ static int bpf_prog_load(union bpf_attr *attr, bpfptr_t uattr, struct bpf_log_at prog->aux->user = get_current_user(); prog->len = attr->insn_cnt; + err = bpf_prog_stream_init(prog, GFP_USER); + if (err) + goto free_prog; + err = -EFAULT; if (copy_from_bpfptr(prog->insns, make_bpfptr(attr->insns, uattr.is_kernel), @@ -6316,6 +6320,28 @@ static int prog_assoc_struct_ops(union bpf_attr *attr) return ret; } +#define BPF_PROG_STREAM_OPEN_LAST_FIELD prog_stream_open.flags + +static int prog_stream_open(union bpf_attr *attr) +{ + struct bpf_prog *prog; + u32 flags = attr->prog_stream_open.flags; + int ret; + + if (CHECK_ATTR(BPF_PROG_STREAM_OPEN)) + return -EINVAL; + if (flags & ~BPF_F_STREAM_NONBLOCK) + return -EINVAL; + + prog = bpf_prog_get(attr->prog_stream_open.prog_fd); + if (IS_ERR(prog)) + return PTR_ERR(prog); + + ret = bpf_prog_stream_new_fd(prog, attr->prog_stream_open.stream_id, flags); + bpf_prog_put(prog); + return ret; +} + static int __sys_bpf(enum bpf_cmd cmd, bpfptr_t uattr, unsigned int size, bpfptr_t uattr_common, unsigned int size_common) { @@ -6488,6 +6514,9 @@ static int __sys_bpf(enum bpf_cmd cmd, bpfptr_t uattr, unsigned int size, case BPF_PROG_ASSOC_STRUCT_OPS: err = prog_assoc_struct_ops(&attr); break; + case BPF_PROG_STREAM_OPEN: + err = prog_stream_open(&attr); + break; default: err = -EINVAL; break; diff --git a/tools/include/uapi/linux/bpf.h b/tools/include/uapi/linux/bpf.h index 0aaa54359aeb..4687c3310996 100644 --- a/tools/include/uapi/linux/bpf.h +++ b/tools/include/uapi/linux/bpf.h @@ -936,6 +936,33 @@ union bpf_iter_link_info { * 0 on success or -1 if an error occurred (in which case, * *errno* is set appropriately). * + * BPF_PROG_STREAM_OPEN + * Description + * Open a file descriptor for one of the BPF streams associated + * with the program identified by *prog_fd*. The stream is selected + * by *stream_id*. + * + * The returned file descriptor supports **read**\ (2) and + * **poll**\ (2). Reads block while the stream is empty unless + * **BPF_F_STREAM_NONBLOCK** is specified in *flags*. A non-blocking + * read of an empty stream fails with **EAGAIN**. + * + * **poll**\ (2) reports **POLLIN** when data is available and + * **POLLHUP** once the program has been freed, that is, after every + * reference to it, including links and other file descriptors, has + * been dropped. Hangup may lag the final release because program + * teardown is deferred. Buffered data remains readable after + * **POLLHUP** and a read returns zero after all such data has been + * consumed. + * + * The file descriptor is read-only and has the close-on-exec flag + * set. It is not seekable and **lseek**\ (2) fails with **ESPIPE**. + * *flags* may only contain **BPF_F_STREAM_NONBLOCK**. + * + * Return + * A new file descriptor (a nonnegative integer), or -1 if an + * error occurred (in which case, *errno* is set appropriately). + * * NOTES * eBPF objects (maps and programs) can be shared between processes. * @@ -993,6 +1020,7 @@ enum bpf_cmd { BPF_TOKEN_CREATE, BPF_PROG_STREAM_READ_BY_FD, BPF_PROG_ASSOC_STRUCT_OPS, + BPF_PROG_STREAM_OPEN, __MAX_BPF_CMD, BPF_COMMON_ATTRS = 1 << 16, /* Indicate carrying syscall common attrs. */ }; @@ -1524,6 +1552,11 @@ enum { BPF_STREAM_STDERR = 2, }; +/* flags for BPF_PROG_STREAM_OPEN command */ +enum { + BPF_F_STREAM_NONBLOCK = (1U << 0), +}; + union bpf_attr { struct { /* anonymous struct used by BPF_MAP_CREATE command */ __u32 map_type; /* one of enum bpf_map_type */ @@ -1950,6 +1983,12 @@ union bpf_attr { __u32 flags; } prog_assoc_struct_ops; + struct { + __u32 prog_fd; + __u32 stream_id; + __u32 flags; + } prog_stream_open; + } __attribute__((aligned(8))); /* The description below is an attempt at providing documentation to eBPF -- 2.53.0