* [RFC v3 2/3] tools/ufq_iosched: add BPF example scheduler and build scaffolding
2026-10-03 4:27 [RFC v3 0/3] block: Introduce a BPF-based I/O scheduler Kaitao Cheng
@ 2026-10-03 4:27 ` Kaitao Cheng
2026-10-03 4:27 ` [RFC v3 1/3] block: Introduce the UFQ I/O scheduler Kaitao Cheng
` (2 subsequent siblings)
3 siblings, 0 replies; 6+ messages in thread
From: Kaitao Cheng @ 2026-10-03 4:27 UTC (permalink / raw)
To: Jens Axboe; +Cc: linux-block, linux-kernel, bpf, Kaitao Cheng
From: Kaitao Cheng <chengkaitao@kylinos.cn>
Add ufq_iosched as a simple example for the UFQ block I/O scheduler,
In the ufq_simple example, we implement the eBPF struct_ops hooks the
kernel exposes so we can exercise and validate the behavior and stability
of the kernel UFQ scheduling framework. The Makefile and directory
layout are modeled after sched_ext.
This mirrors the sched_ext examples pattern so developers can experiment
with user-defined queueing policies on top of IOSCHED_UFQ.
Signed-off-by: Kaitao Cheng <chengkaitao@kylinos.cn>
---
tools/ufq_iosched/.gitignore | 2 +
tools/ufq_iosched/Makefile | 262 ++++++++
tools/ufq_iosched/README.md | 145 +++++
.../include/bpf-compat/gnu/stubs.h | 7 +
tools/ufq_iosched/include/ufq/common.bpf.h | 83 +++
tools/ufq_iosched/include/ufq/common.h | 87 +++
tools/ufq_iosched/include/ufq/simple_stat.h | 23 +
tools/ufq_iosched/ufq_simple.bpf.c | 604 ++++++++++++++++++
tools/ufq_iosched/ufq_simple.c | 120 ++++
9 files changed, 1333 insertions(+)
create mode 100644 tools/ufq_iosched/.gitignore
create mode 100644 tools/ufq_iosched/Makefile
create mode 100644 tools/ufq_iosched/README.md
create mode 100644 tools/ufq_iosched/include/bpf-compat/gnu/stubs.h
create mode 100644 tools/ufq_iosched/include/ufq/common.bpf.h
create mode 100644 tools/ufq_iosched/include/ufq/common.h
create mode 100644 tools/ufq_iosched/include/ufq/simple_stat.h
create mode 100644 tools/ufq_iosched/ufq_simple.bpf.c
create mode 100644 tools/ufq_iosched/ufq_simple.c
diff --git a/tools/ufq_iosched/.gitignore b/tools/ufq_iosched/.gitignore
new file mode 100644
index 000000000000..d6264fe1c8cd
--- /dev/null
+++ b/tools/ufq_iosched/.gitignore
@@ -0,0 +1,2 @@
+tools/
+build/
diff --git a/tools/ufq_iosched/Makefile b/tools/ufq_iosched/Makefile
new file mode 100644
index 000000000000..7dc37d9172aa
--- /dev/null
+++ b/tools/ufq_iosched/Makefile
@@ -0,0 +1,262 @@
+# SPDX-License-Identifier: GPL-2.0
+# Copyright (c) 2026 KylinSoft Corporation.
+# Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+include ../build/Build.include
+include ../scripts/Makefile.arch
+include ../scripts/Makefile.include
+
+all: all_targets
+
+ifneq ($(LLVM),)
+ifneq ($(filter %/,$(LLVM)),)
+LLVM_PREFIX := $(LLVM)
+else ifneq ($(filter -%,$(LLVM)),)
+LLVM_SUFFIX := $(LLVM)
+endif
+
+CLANG_TARGET_FLAGS_arm := arm-linux-gnueabi
+CLANG_TARGET_FLAGS_arm64 := aarch64-linux-gnu
+CLANG_TARGET_FLAGS_hexagon := hexagon-linux-musl
+CLANG_TARGET_FLAGS_m68k := m68k-linux-gnu
+CLANG_TARGET_FLAGS_mips := mipsel-linux-gnu
+CLANG_TARGET_FLAGS_powerpc := powerpc64le-linux-gnu
+CLANG_TARGET_FLAGS_riscv := riscv64-linux-gnu
+CLANG_TARGET_FLAGS_s390 := s390x-linux-gnu
+CLANG_TARGET_FLAGS_x86 := x86_64-linux-gnu
+CLANG_TARGET_FLAGS := $(CLANG_TARGET_FLAGS_$(ARCH))
+
+ifeq ($(CROSS_COMPILE),)
+ifeq ($(CLANG_TARGET_FLAGS),)
+$(error Specify CROSS_COMPILE or add '--target=' option to lib.mk)
+else
+CLANG_FLAGS += --target=$(CLANG_TARGET_FLAGS)
+endif # CLANG_TARGET_FLAGS
+else
+CLANG_FLAGS += --target=$(notdir $(CROSS_COMPILE:%-=%))
+endif # CROSS_COMPILE
+
+CC := $(LLVM_PREFIX)clang$(LLVM_SUFFIX) $(CLANG_FLAGS) -fintegrated-as
+else
+CC := $(CROSS_COMPILE)gcc
+endif # LLVM
+
+CURDIR := $(abspath .)
+TOOLSDIR := $(abspath ..)
+LIBDIR := $(TOOLSDIR)/lib
+BPFDIR := $(LIBDIR)/bpf
+TOOLSINCDIR := $(TOOLSDIR)/include
+BPFTOOLDIR := $(TOOLSDIR)/bpf/bpftool
+APIDIR := $(TOOLSINCDIR)/uapi
+GENDIR := $(abspath ../../include/generated)
+GENHDR := $(GENDIR)/autoconf.h
+
+ifeq ($(O),)
+OUTPUT_DIR := $(CURDIR)/build
+else
+OUTPUT_DIR := $(O)/build
+endif # O
+OBJ_DIR := $(OUTPUT_DIR)/obj
+INCLUDE_DIR := $(OUTPUT_DIR)/include
+BPFOBJ_DIR := $(OBJ_DIR)/libbpf
+UFQOBJ_DIR := $(OBJ_DIR)/ufq_iosched
+BINDIR := $(OUTPUT_DIR)/bin
+BPFOBJ := $(BPFOBJ_DIR)/libbpf.a
+ifneq ($(CROSS_COMPILE),)
+HOST_BUILD_DIR := $(OBJ_DIR)/host/obj
+HOST_OUTPUT_DIR := $(OBJ_DIR)/host
+HOST_INCLUDE_DIR := $(HOST_OUTPUT_DIR)/include
+else
+HOST_BUILD_DIR := $(OBJ_DIR)
+HOST_OUTPUT_DIR := $(OUTPUT_DIR)
+HOST_INCLUDE_DIR := $(INCLUDE_DIR)
+endif
+HOST_BPFOBJ := $(HOST_BUILD_DIR)/libbpf/libbpf.a
+RESOLVE_BTFIDS := $(HOST_BUILD_DIR)/resolve_btfids/resolve_btfids
+DEFAULT_BPFTOOL := $(HOST_OUTPUT_DIR)/sbin/bpftool
+
+VMLINUX_BTF_PATHS ?= $(if $(O),$(O)/vmlinux) \
+ $(if $(KBUILD_OUTPUT),$(KBUILD_OUTPUT)/vmlinux) \
+ ../../vmlinux \
+ /sys/kernel/btf/vmlinux \
+ /boot/vmlinux-$(shell uname -r)
+VMLINUX_BTF ?= $(abspath $(firstword $(wildcard $(VMLINUX_BTF_PATHS))))
+ifeq ($(VMLINUX_BTF),)
+$(error Cannot find a vmlinux for VMLINUX_BTF at any of "$(VMLINUX_BTF_PATHS)")
+endif
+
+BPFTOOL ?= $(DEFAULT_BPFTOOL)
+
+ifneq ($(wildcard $(GENHDR)),)
+ GENFLAGS := -DHAVE_GENHDR
+endif
+
+CFLAGS += -g -O2 -rdynamic -pthread -Wall -Werror $(GENFLAGS) \
+ -I$(INCLUDE_DIR) -I$(GENDIR) -I$(LIBDIR) \
+ -I$(TOOLSINCDIR) -I$(APIDIR) -I$(CURDIR)/include
+
+# Silence some warnings when compiled with clang
+ifneq ($(LLVM),)
+CFLAGS += -Wno-unused-command-line-argument
+endif
+
+LDFLAGS += -lelf -lz -lpthread
+
+IS_LITTLE_ENDIAN = $(shell $(CC) -dM -E - </dev/null | \
+ grep 'define __BYTE_ORDER__ __ORDER_LITTLE_ENDIAN__')
+
+# Get Clang's default includes on this system, as opposed to those seen by
+# '-target bpf'. This fixes "missing" files on some architectures/distros,
+# such as asm/byteorder.h, asm/socket.h, asm/sockios.h, sys/cdefs.h etc.
+#
+# Use '-idirafter': Don't interfere with include mechanics except where the
+# build would have failed anyways.
+define get_sys_includes
+$(shell $(1) -v -E - </dev/null 2>&1 \
+ | sed -n '/<...> search starts here:/,/End of search list./{ s| \(/.*\)|-idirafter \1|p }') \
+$(shell $(1) -dM -E - </dev/null | grep '__riscv_xlen ' | awk '{printf("-D__riscv_xlen=%d -D__BITS_PER_LONG=%d", $$3, $$3)}')
+endef
+
+BPF_CFLAGS = -g -D__TARGET_ARCH_$(SRCARCH) \
+ $(if $(IS_LITTLE_ENDIAN),-mlittle-endian,-mbig-endian) \
+ -I$(CURDIR)/include -I$(CURDIR)/include/bpf-compat \
+ -I$(INCLUDE_DIR) -I$(APIDIR) \
+ -I../../include \
+ $(call get_sys_includes,$(CLANG)) \
+ -Wall -Wno-compare-distinct-pointer-types \
+ -Wno-microsoft-anon-tag \
+ -fms-extensions \
+ -O2 -mcpu=v3
+
+# sort removes libbpf duplicates when not cross-building
+MAKE_DIRS := $(sort $(OBJ_DIR)/libbpf $(HOST_BUILD_DIR)/libbpf \
+ $(HOST_BUILD_DIR)/bpftool $(HOST_BUILD_DIR)/resolve_btfids \
+ $(INCLUDE_DIR) $(UFQOBJ_DIR) $(BINDIR))
+
+$(MAKE_DIRS):
+ $(call msg,MKDIR,,$@)
+ $(Q)mkdir -p $@
+
+ifneq ($(CROSS_COMPILE),)
+$(BPFOBJ): $(wildcard $(BPFDIR)/*.[ch] $(BPFDIR)/Makefile) \
+ $(APIDIR)/linux/bpf.h \
+ | $(OBJ_DIR)/libbpf
+ $(Q)$(MAKE) $(submake_extras) CROSS_COMPILE=$(CROSS_COMPILE) \
+ -C $(BPFDIR) OUTPUT=$(OBJ_DIR)/libbpf/ \
+ EXTRA_CFLAGS='-g -O0 -fPIC' \
+ LDFLAGS="$(LDFLAGS)" \
+ DESTDIR=$(OUTPUT_DIR) prefix= all install_headers
+endif
+
+$(HOST_BPFOBJ): $(wildcard $(BPFDIR)/*.[ch] $(BPFDIR)/Makefile) \
+ $(APIDIR)/linux/bpf.h \
+ | $(HOST_BUILD_DIR)/libbpf
+ $(Q)$(MAKE) $(submake_extras) -C $(BPFDIR) \
+ OUTPUT=$(HOST_BUILD_DIR)/libbpf/ \
+ ARCH= CROSS_COMPILE= CC="$(HOSTCC)" LD=$(HOSTLD) \
+ EXTRA_CFLAGS='-g -O0 -fPIC' \
+ DESTDIR=$(HOST_OUTPUT_DIR) prefix= all install_headers
+
+$(DEFAULT_BPFTOOL): $(wildcard $(BPFTOOLDIR)/*.[ch] $(BPFTOOLDIR)/Makefile) \
+ $(HOST_BPFOBJ) | $(HOST_BUILD_DIR)/bpftool
+ $(Q)$(MAKE) $(submake_extras) -C $(BPFTOOLDIR) \
+ ARCH= CROSS_COMPILE= CC="$(HOSTCC)" LD=$(HOSTLD) \
+ EXTRA_CFLAGS='-g -O0' \
+ OUTPUT=$(HOST_BUILD_DIR)/bpftool/ \
+ LIBBPF_OUTPUT=$(HOST_BUILD_DIR)/libbpf/ \
+ LIBBPF_DESTDIR=$(HOST_OUTPUT_DIR)/ \
+ prefix= DESTDIR=$(HOST_OUTPUT_DIR)/ install-bin
+
+$(INCLUDE_DIR)/vmlinux.h: $(VMLINUX_BTF) $(BPFTOOL) | $(INCLUDE_DIR)
+ifeq ($(VMLINUX_H),)
+ $(call msg,GEN,,$@)
+ $(Q)$(BPFTOOL) btf dump file $(VMLINUX_BTF) format c > $@
+else
+ $(call msg,CP,,$@)
+ $(Q)cp "$(VMLINUX_H)" $@
+endif
+
+$(UFQOBJ_DIR)/%.bpf.o: %.bpf.c $(INCLUDE_DIR)/vmlinux.h include/ufq/*.h \
+ | $(BPFOBJ) $(UFQOBJ_DIR)
+ $(call msg,CLNG-BPF,,$(notdir $@))
+ $(Q)$(CLANG) $(BPF_CFLAGS) -target bpf -c $< -o $@
+
+$(INCLUDE_DIR)/%.bpf.skel.h: $(UFQOBJ_DIR)/%.bpf.o $(INCLUDE_DIR)/vmlinux.h $(BPFTOOL)
+ $(eval sched=$(notdir $@))
+ $(call msg,GEN-SKEL,,$(sched))
+ $(Q)$(BPFTOOL) gen object $(<:.o=.linked1.o) $<
+ $(Q)$(BPFTOOL) gen object $(<:.o=.linked2.o) $(<:.o=.linked1.o)
+ $(Q)$(BPFTOOL) gen object $(<:.o=.linked3.o) $(<:.o=.linked2.o)
+ $(Q)diff $(<:.o=.linked2.o) $(<:.o=.linked3.o)
+ $(Q)$(BPFTOOL) gen skeleton $(<:.o=.linked3.o) name $(subst .bpf.skel.h,,$(sched)) > $@
+ $(Q)$(BPFTOOL) gen subskeleton $(<:.o=.linked3.o) name $(subst .bpf.skel.h,,$(sched)) > $(@:.skel.h=.subskel.h)
+
+UFQ_COMMON_DEPS := include/ufq/common.h include/ufq/simple_stat.h | $(BINDIR)
+
+c-sched-targets = ufq_simple
+
+$(addprefix $(BINDIR)/,$(c-sched-targets)): \
+ $(BINDIR)/%: \
+ $(filter-out %.bpf.c,%.c) \
+ $(INCLUDE_DIR)/%.bpf.skel.h \
+ $(UFQ_COMMON_DEPS)
+ $(eval sched=$(notdir $@))
+ $(CC) $(CFLAGS) -c $(sched).c -o $(UFQOBJ_DIR)/$(sched).o
+ $(CC) -o $@ $(UFQOBJ_DIR)/$(sched).o $(BPFOBJ) $(LDFLAGS)
+
+$(c-sched-targets): %: $(BINDIR)/%
+
+install: all
+ $(Q)mkdir -p $(DESTDIR)/usr/local/bin/
+ $(Q)cp $(BINDIR)/* $(DESTDIR)/usr/local/bin/
+
+clean:
+ rm -rf $(OUTPUT_DIR) $(HOST_OUTPUT_DIR)
+ rm -f *.o *.bpf.o *.bpf.skel.h *.bpf.subskel.h
+ rm -f $(c-sched-targets)
+
+help:
+ @echo 'Building targets'
+ @echo '================'
+ @echo ''
+ @echo ' all - Compile all schedulers'
+ @echo ''
+ @echo 'Alternatively, you may compile individual schedulers:'
+ @echo ''
+ @printf ' %s\n' $(c-sched-targets)
+ @echo ''
+ @echo 'For any scheduler build target, you may specify an alternative'
+ @echo 'build output path with the O= environment variable. For example:'
+ @echo ''
+ @echo ' O=/tmp/ufq_iosched make all'
+ @echo ''
+ @echo 'will compile all schedulers, and emit the build artifacts to'
+ @echo '/tmp/ufq_iosched/build.'
+ @echo ''
+ @echo ''
+ @echo 'Installing targets'
+ @echo '=================='
+ @echo ''
+ @echo ' install - Compile and install all schedulers to /usr/bin.'
+ @echo ' You may specify the DESTDIR= environment variable'
+ @echo ' to indicate a prefix for /usr/bin. For example:'
+ @echo ''
+ @echo ' DESTDIR=/tmp/ufq_iosched make install'
+ @echo ''
+ @echo ' will build the schedulers in CWD/build, and'
+ @echo ' install the schedulers to /tmp/ufq_iosched/usr/bin.'
+ @echo ''
+ @echo ''
+ @echo 'Cleaning targets'
+ @echo '================'
+ @echo ''
+ @echo ' clean - Remove all generated files'
+
+all_targets: $(c-sched-targets)
+
+.PHONY: all all_targets $(c-sched-targets) clean help
+
+# delete failed targets
+.DELETE_ON_ERROR:
+
+# keep intermediate (.bpf.skel.h, .bpf.o, etc) targets
+.SECONDARY:
diff --git a/tools/ufq_iosched/README.md b/tools/ufq_iosched/README.md
new file mode 100644
index 000000000000..1fa305fb576e
--- /dev/null
+++ b/tools/ufq_iosched/README.md
@@ -0,0 +1,145 @@
+UFQ IOSCHED EXAMPLE SCHEDULERS
+============================
+
+# Introduction
+
+This directory contains a simple example of the ufq IO scheduler. It is meant
+to illustrate the different kinds of IO schedulers you can build with ufq;
+new schedulers will be added as the project evolves across releases.
+
+# Compiling the examples
+
+There are a few toolchain dependencies for compiling the example schedulers.
+
+## Toolchain dependencies
+
+1. clang >= 16.0.0
+
+The schedulers are BPF programs, and therefore must be compiled with clang. gcc
+is actively working on adding a BPF backend compiler as well, but are still
+missing some features such as BTF type tags which are necessary for using
+kptrs.
+
+2. pahole >= 1.25
+
+You may need pahole in order to generate BTF from DWARF.
+
+3. rust >= 1.85.0
+
+Rust schedulers uses features present in the rust toolchain >= 1.85.0. You
+should be able to use the stable build from rustup.
+
+There are other requirements as well, such as make, but these are the main /
+non-trivial ones.
+
+## Compiling the kernel
+
+In order to run a ufq scheduler, you'll have to run a kernel compiled
+with the patches in this repository, and with a minimum set of necessary
+Kconfig options:
+
+```
+CONFIG_BPF=y
+CONFIG_IOSCHED_UFQ=y
+CONFIG_BPF_SYSCALL=y
+CONFIG_BPF_JIT=y
+CONFIG_DEBUG_INFO_BTF=y
+```
+
+It's also recommended that you also include the following Kconfig options:
+
+```
+CONFIG_BPF_JIT_ALWAYS_ON=y
+CONFIG_BPF_JIT_DEFAULT_ON=y
+CONFIG_PAHOLE_HAS_SPLIT_BTF=y
+CONFIG_PAHOLE_HAS_BTF_TAG=y
+```
+
+There is a `Kconfig` file in this directory whose contents you can append to
+your local `.config` file, as long as there are no conflicts with any existing
+options in the file.
+
+## Getting a vmlinux.h file
+
+You may notice that most of the example schedulers include a "vmlinux.h" file.
+This is a large, auto-generated header file that contains all of the types
+defined in some vmlinux binary that was compiled with
+[BTF](https://docs.kernel.org/bpf/btf.html) (i.e. with the BTF-related Kconfig
+options specified above).
+
+The header file is created using `bpftool`, by passing it a vmlinux binary
+compiled with BTF as follows:
+
+```bash
+$ bpftool btf dump file /path/to/vmlinux format c > vmlinux.h
+```
+
+`bpftool` analyzes all of the BTF encodings in the binary, and produces a
+header file that can be included by BPF programs to access those types. For
+example, using vmlinux.h allows a scheduler to access fields defined directly
+in vmlinux
+
+The scheduler build system will generate this vmlinux.h file as part of the
+scheduler build pipeline. It looks for a vmlinux file in the following
+dependency order:
+
+1. If the O= environment variable is defined, at `$O/vmlinux`
+2. If the KBUILD_OUTPUT= environment variable is defined, at
+ `$KBUILD_OUTPUT/vmlinux`
+3. At `../../vmlinux` (i.e. at the root of the kernel tree where you're
+ compiling the schedulers)
+4. `/sys/kernel/btf/vmlinux`
+5. `/boot/vmlinux-$(uname -r)`
+
+In other words, if you have compiled a kernel in your local repo, its vmlinux
+file will be used to generate vmlinux.h. Otherwise, it will be the vmlinux of
+the kernel you're currently running on. This means that if you're running on a
+kernel with ufq support, you may not need to compile a local kernel at
+all.
+
+### Aside on CO-RE
+
+One of the cooler features of BPF is that it supports
+[CO-RE](https://nakryiko.com/posts/bpf-core-reference-guide/) (Compile Once Run
+Everywhere). This feature allows you to reference fields inside of structs with
+types defined internal to the kernel, and not have to recompile if you load the
+BPF program on a different kernel with the field at a different offset.
+
+## Compiling the schedulers
+
+Once you have your toolchain setup, and a vmlinux that can be used to generate
+a full vmlinux.h file, you can compile the schedulers using `make`. The build
+vendors libbpf and bpftool under `build/`, then compiles BPF objects with
+`clang` and the userspace loader (for example `ufq_simple.c`) with `gcc`.
+
+```bash
+$ make -j($nproc)
+```
+
+## ufq_simple
+
+A simple IO scheduler that provides an example of a minimal ufq scheduler.
+Populates commonly used kernel-exposed BPF interfaces for testing the UFQ
+scheduler framework in the kernel.
+
+### bpftool feature detection (`llvm`, `libcap`, `libbfd`)
+
+While building bpftool, you may see lines similar to:
+
+```
+Auto-detecting system features:
+... clang-bpf-co-re: [ on ]
+... llvm: [ OFF ]
+... libcap: [ OFF ]
+... libbfd: [ OFF ]
+```
+
+This matches a typical minimal build environment: CO-RE support is on, while
+bpftool's optional links to LLVM (for some disassembly paths), libcap, and
+libbfd are off. **`llvm: [ OFF ]` is not a build failure** for these examples;
+the scheduler build can still complete (libbpf static library, bpftool,
+`vmlinux.h` generation, BPF `.o` files, and final `ufq_simple` link with `gcc`).
+
+If `libcap` or `libbfd` show `[ OFF ]`, you only miss optional bpftool features
+that depend on those libraries; you do not need them for compiling and loading
+the schedulers in this directory.
diff --git a/tools/ufq_iosched/include/bpf-compat/gnu/stubs.h b/tools/ufq_iosched/include/bpf-compat/gnu/stubs.h
new file mode 100644
index 000000000000..2a8c8aced06a
--- /dev/null
+++ b/tools/ufq_iosched/include/bpf-compat/gnu/stubs.h
@@ -0,0 +1,7 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * BPF builds can reach glibc headers through Clang's host include paths.
+ * On x86, the BPF target leaves __x86_64__ undefined, causing gnu/stubs.h
+ * to include stubs-32.h, which may be absent without 32-bit glibc headers.
+ * Put this empty substitute before the system headers in BPF include paths.
+ */
diff --git a/tools/ufq_iosched/include/ufq/common.bpf.h b/tools/ufq_iosched/include/ufq/common.bpf.h
new file mode 100644
index 000000000000..d3d0b4aec6a7
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/common.bpf.h
@@ -0,0 +1,83 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#ifndef __UFQ_COMMON_BPF_H
+#define __UFQ_COMMON_BPF_H
+
+#ifdef LSP
+#define __bpf__
+#include "../vmlinux/vmlinux.h"
+#else
+#include "vmlinux.h"
+#endif
+
+#include <bpf/bpf_helpers.h>
+#include <bpf/bpf_tracing.h>
+#include <bpf/bpf_core_read.h>
+#include <asm-generic/errno.h>
+#include "simple_stat.h"
+
+#define BPF_STRUCT_OPS(name, args...) \
+ SEC("struct_ops/" #name) BPF_PROG(name, ##args)
+
+/* Define struct ufq_iosched_ops for .struct_ops.link in the BPF object */
+#define UFQ_OPS_DEFINE(__name, ...) \
+ SEC(".struct_ops.link") \
+ struct ufq_iosched_ops __name = { \
+ __VA_ARGS__, \
+ }
+
+/* Associate each graph root with its node type and member for the verifier. */
+#define __contains(name, node) __attribute__((btf_decl_tag("contains:" #name ":" #node)))
+
+struct request *bpf_request_acquire(struct request *rq) __ksym;
+void bpf_request_release(struct request *rq) __ksym;
+bool bpf_request_bio_try_merge(struct request *rq, struct bio *bio,
+ unsigned int nr_segs) __ksym;
+struct request *bpf_request_try_merge(struct request *rq, struct request *next) __ksym;
+/* BPF cannot write elv.priv[1] directly; this kfunc is KF_SPINLOCK_SAFE. */
+void bpf_request_set_elv_priv1(struct request *rq, u64 val) __ksym;
+
+void *bpf_obj_new_impl(__u64 local_type_id, void *meta) __ksym;
+void bpf_obj_drop_impl(void *kptr, void *meta) __ksym;
+
+#define bpf_obj_new(type) ((type *)bpf_obj_new_impl(bpf_core_type_id_local(type), NULL))
+#define bpf_obj_drop(kptr) bpf_obj_drop_impl(kptr, NULL)
+
+int bpf_list_push_front_impl(struct bpf_list_head *head,
+ struct bpf_list_node *node,
+ void *meta, __u64 off) __ksym;
+#define bpf_list_push_front(head, node) bpf_list_push_front_impl(head, node, NULL, 0)
+
+int bpf_list_push_back_impl(struct bpf_list_head *head,
+ struct bpf_list_node *node,
+ void *meta, __u64 off) __ksym;
+#define bpf_list_push_back(head, node) bpf_list_push_back_impl(head, node, NULL, 0)
+
+struct bpf_list_node *bpf_list_pop_front(struct bpf_list_head *head) __ksym;
+struct bpf_list_node *bpf_list_pop_back(struct bpf_list_head *head) __ksym;
+struct bpf_list_node *bpf_list_front(struct bpf_list_head *head) __ksym;
+bool bpf_list_empty(struct bpf_list_head *head) __ksym;
+struct bpf_list_node *bpf_list_del(struct bpf_list_head *head,
+ struct bpf_list_node *node) __ksym;
+
+struct bpf_rb_node *bpf_rbtree_remove(struct bpf_rb_root *root,
+ struct bpf_rb_node *node) __ksym;
+int bpf_rbtree_add_impl(struct bpf_rb_root *root, struct bpf_rb_node *node,
+ bool (less)(struct bpf_rb_node *a, const struct bpf_rb_node *b),
+ void *meta, __u64 off) __ksym;
+#define bpf_rbtree_add(head, node, less) bpf_rbtree_add_impl(head, node, less, NULL, 0)
+
+struct bpf_rb_node *bpf_rbtree_first(struct bpf_rb_root *root) __ksym;
+struct bpf_rb_node *bpf_rbtree_root(struct bpf_rb_root *root) __ksym;
+struct bpf_rb_node *bpf_rbtree_left(struct bpf_rb_root *root,
+ struct bpf_rb_node *node) __ksym;
+struct bpf_rb_node *bpf_rbtree_right(struct bpf_rb_root *root,
+ struct bpf_rb_node *node) __ksym;
+
+void *bpf_refcount_acquire_impl(void *kptr, void *meta) __ksym;
+#define bpf_refcount_acquire(kptr) bpf_refcount_acquire_impl(kptr, NULL)
+
+#endif /* __UFQ_COMMON_BPF_H */
diff --git a/tools/ufq_iosched/include/ufq/common.h b/tools/ufq_iosched/include/ufq/common.h
new file mode 100644
index 000000000000..f6f43ad8efa0
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/common.h
@@ -0,0 +1,87 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#ifndef __UFQ_IOSCHED_COMMON_H
+#define __UFQ_IOSCHED_COMMON_H
+
+#ifdef __KERNEL__
+#error "Should not be included by BPF programs"
+#endif
+
+#include <stdarg.h>
+#include <stdio.h>
+#include <stdlib.h>
+#include <stdint.h>
+#include <errno.h>
+#include <bpf/bpf.h>
+#include "simple_stat.h"
+
+typedef uint8_t u8;
+typedef uint16_t u16;
+typedef uint32_t u32;
+typedef uint64_t u64;
+typedef int8_t s8;
+typedef int16_t s16;
+typedef int32_t s32;
+typedef int64_t s64;
+
+#define UFQ_ERR(__fmt, ...) \
+ do { \
+ fprintf(stderr, "[UFQ_ERR] %s:%d", __FILE__, __LINE__); \
+ if (errno) \
+ fprintf(stderr, " (%s)\n", strerror(errno)); \
+ else \
+ fprintf(stderr, "\n"); \
+ fprintf(stderr, __fmt __VA_OPT__(,) __VA_ARGS__); \
+ fprintf(stderr, "\n"); \
+ \
+ exit(EXIT_FAILURE); \
+ } while (0)
+
+#define UFQ_ERR_IF(__cond, __fmt, ...) \
+ do { \
+ if (__cond) \
+ UFQ_ERR((__fmt) __VA_OPT__(,) __VA_ARGS__); \
+ } while (0)
+
+/*
+ * Pair these skeleton helpers with UFQ_OPS_DEFINE() in common.bpf.h.
+ */
+#define UFQ_OPS_OPEN(__ops_name, __ufq_name) ({ \
+ struct __ufq_name *__skel; \
+ \
+ __skel = __ufq_name##__open(); \
+ UFQ_ERR_IF(!__skel, "Could not open " #__ufq_name); \
+ (void)__skel->maps.__ops_name; \
+ __skel; \
+})
+
+#define UFQ_OPS_LOAD(__skel, __ops_name, __ufq_name) ({ \
+ (void)(__skel)->maps.__ops_name; \
+ UFQ_ERR_IF(__ufq_name##__load((__skel)), "Failed to load skel"); \
+})
+
+/*
+ * libbpf 1.5 and later can autoattach struct_ops maps through the skeleton.
+ * Disable autoattach because UFQ_OPS_ATTACH() attaches the map explicitly.
+ */
+#if LIBBPF_MAJOR_VERSION > 1 || \
+ (LIBBPF_MAJOR_VERSION == 1 && LIBBPF_MINOR_VERSION >= 5)
+#define __UFQ_OPS_DISABLE_AUTOATTACH(__skel, __ops_name) \
+ bpf_map__set_autoattach((__skel)->maps.__ops_name, false)
+#else
+#define __UFQ_OPS_DISABLE_AUTOATTACH(__skel, __ops_name) do {} while (0)
+#endif
+
+#define UFQ_OPS_ATTACH(__skel, __ops_name, __ufq_name) ({ \
+ struct bpf_link *__link; \
+ __UFQ_OPS_DISABLE_AUTOATTACH(__skel, __ops_name); \
+ UFQ_ERR_IF(__ufq_name##__attach((__skel)), "Failed to attach skel"); \
+ __link = bpf_map__attach_struct_ops((__skel)->maps.__ops_name); \
+ UFQ_ERR_IF(!__link, "Failed to attach struct_ops"); \
+ __link; \
+})
+
+#endif /* __UFQ_IOSCHED_COMMON_H */
diff --git a/tools/ufq_iosched/include/ufq/simple_stat.h b/tools/ufq_iosched/include/ufq/simple_stat.h
new file mode 100644
index 000000000000..c26d0ea57a69
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/simple_stat.h
@@ -0,0 +1,23 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#ifndef __UFQ_SIMPLE_STAT_H
+#define __UFQ_SIMPLE_STAT_H
+
+enum ufq_simp_stat_index {
+ UFQ_SIMP_INSERT_CNT,
+ UFQ_SIMP_INSERT_SIZE,
+ UFQ_SIMP_DISPATCH_CNT,
+ UFQ_SIMP_DISPATCH_SIZE,
+ UFQ_SIMP_RQMERGE_CNT,
+ UFQ_SIMP_RQMERGE_SIZE,
+ UFQ_SIMP_BIOMERGE_CNT,
+ UFQ_SIMP_BIOMERGE_SIZE,
+ UFQ_SIMP_FINISH_CNT,
+ UFQ_SIMP_FINISH_SIZE,
+ UFQ_SIMP_STAT_MAX,
+};
+
+#endif /* __UFQ_SIMPLE_STAT_H */
diff --git a/tools/ufq_iosched/ufq_simple.bpf.c b/tools/ufq_iosched/ufq_simple.bpf.c
new file mode 100644
index 000000000000..f49ef0d19258
--- /dev/null
+++ b/tools/ufq_iosched/ufq_simple.bpf.c
@@ -0,0 +1,604 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#include <ufq/common.bpf.h>
+
+char _license[] SEC("license") = "GPL";
+
+#define UFQ_DISK_SUM 20
+#define BLK_MQ_INSERT_AT_HEAD 0x01
+#define REQ_OP_MASK ((1 << 8) - 1)
+#define SECTOR_SHIFT 9
+#define UFQ_LOOP_MAX 100
+
+struct {
+ __uint(type, BPF_MAP_TYPE_PERCPU_ARRAY);
+ __uint(max_entries, UFQ_SIMP_STAT_MAX);
+ __uint(key_size, sizeof(u32));
+ __uint(value_size, sizeof(u64));
+} stats SEC(".maps");
+
+enum ufq_simp_data_dir {
+ UFQ_SIMP_READ,
+ UFQ_SIMP_WRITE,
+ UFQ_SIMP_DIR_COUNT
+};
+
+struct queue_list_node {
+ struct bpf_list_node node;
+ struct request __kptr * req;
+};
+
+struct sort_tree_node {
+ struct bpf_refcount ref;
+ struct bpf_rb_node rb_node;
+ struct request __kptr * req;
+ u64 key;
+};
+
+struct ufq_simple_data {
+ struct bpf_rb_root sort_tree_read __contains(sort_tree_node, rb_node);
+ struct bpf_rb_root sort_tree_write __contains(sort_tree_node, rb_node);
+ struct bpf_list_head dispatch __contains(queue_list_node, node);
+ struct bpf_spin_lock lock;
+};
+
+struct {
+ __uint(type, BPF_MAP_TYPE_HASH);
+ __uint(max_entries, UFQ_DISK_SUM);
+ __type(key, s32);
+ __type(value, struct ufq_simple_data);
+} ufq_map SEC(".maps");
+
+static void stat_add(u32 idx, u32 val)
+{
+ u64 *cnt_p = bpf_map_lookup_elem(&stats, &idx);
+
+ if (cnt_p)
+ __sync_fetch_and_add(cnt_p, val);
+}
+
+static void stat_sub(u32 idx, u32 val)
+{
+ u64 *cnt_p = bpf_map_lookup_elem(&stats, &idx);
+
+ if (cnt_p)
+ __sync_fetch_and_sub(cnt_p, val);
+}
+
+static bool sort_tree_less(struct bpf_rb_node *a, const struct bpf_rb_node *b)
+{
+ struct sort_tree_node *node_a, *node_b;
+
+ node_a = container_of(a, struct sort_tree_node, rb_node);
+ node_b = container_of(b, struct sort_tree_node, rb_node);
+
+ return node_a->key < node_b->key;
+}
+
+static struct ufq_simple_data *dd_init_sched(struct request_queue *q)
+{
+ struct ufq_simple_data ufq_sd = {}, *ufq_sp;
+ int ret, id = q->id;
+
+ bpf_printk("ufq_simple init sched!");
+ ret = bpf_map_update_elem(&ufq_map, &id, &ufq_sd, BPF_NOEXIST);
+ if (ret && ret != -EEXIST) {
+ bpf_printk("ufq_simple/init_sched: update ufq_map err %d", ret);
+ return NULL;
+ }
+
+ ufq_sp = bpf_map_lookup_elem(&ufq_map, &id);
+ if (!ufq_sp) {
+ bpf_printk("ufq_simple/init_sched: lookup queue id %d in ufq_map failed", id);
+ return NULL;
+ }
+
+ return ufq_sp;
+}
+
+int BPF_STRUCT_OPS(ufq_simple_init_sched, struct request_queue *q)
+{
+ if (dd_init_sched(q))
+ return 0;
+ else
+ return -EPERM;
+}
+
+int BPF_STRUCT_OPS(ufq_simple_exit_sched, struct request_queue *q)
+{
+ int id = q->id;
+
+ bpf_printk("ufq_simple exit sched!");
+ bpf_map_delete_elem(&ufq_map, &id);
+ return 0;
+}
+
+int BPF_STRUCT_OPS(ufq_simple_insert_req, struct request_queue *q,
+ struct request *rq, blk_insert_t flags,
+ struct list_head *freeq)
+{
+ struct ufq_simple_data *ufq_sd;
+ struct request *acquired, *old;
+ struct queue_list_node *qnode;
+ struct sort_tree_node *snode;
+ int id = q->id, ret = 0;
+ enum ufq_simp_data_dir dir = ((rq->cmd_flags & REQ_OP_MASK) & 1) ?
+ UFQ_SIMP_WRITE : UFQ_SIMP_READ;
+
+ ufq_sd = bpf_map_lookup_elem(&ufq_map, &id);
+ if (!ufq_sd) {
+ ufq_sd = dd_init_sched(q);
+ if (!ufq_sd) {
+ bpf_printk("ufq_simple/insert_req: dd_init_sched failed");
+ return -EPERM;
+ }
+ }
+
+ if (flags & BLK_MQ_INSERT_AT_HEAD) {
+ qnode = bpf_obj_new(typeof(*qnode));
+ if (!qnode) {
+ bpf_printk("ufq_simple/insert_req: qnode alloc failed");
+ return -ENOMEM;
+ }
+
+ acquired = bpf_request_acquire(rq);
+ if (!acquired) {
+ bpf_obj_drop(qnode);
+ bpf_printk("ufq_simple/head-insert_req: request_acquire failed");
+ return -EPERM;
+ }
+
+ old = bpf_kptr_xchg(&qnode->req, acquired);
+ if (old)
+ bpf_request_release(old);
+
+ bpf_spin_lock(&ufq_sd->lock);
+ ret = bpf_list_push_back(&ufq_sd->dispatch, &qnode->node);
+ bpf_spin_unlock(&ufq_sd->lock);
+ } else {
+ snode = bpf_obj_new(typeof(*snode));
+ if (!snode) {
+ bpf_printk("ufq_simple/insert_req: sort_tree_node alloc failed");
+ return -ENOMEM;
+ }
+
+ snode->key = rq->__sector;
+
+ acquired = bpf_request_acquire(rq);
+ if (!acquired) {
+ bpf_obj_drop(snode);
+ bpf_printk("ufq_simple/insert_req: bpf_request_acquire failed");
+ return -EPERM;
+ }
+
+ old = bpf_kptr_xchg(&snode->req, acquired);
+ if (old)
+ bpf_request_release(old);
+
+ bpf_spin_lock(&ufq_sd->lock);
+ if (dir == UFQ_SIMP_READ)
+ bpf_rbtree_add(&ufq_sd->sort_tree_read, &snode->rb_node, sort_tree_less);
+ else
+ bpf_rbtree_add(&ufq_sd->sort_tree_write, &snode->rb_node, sort_tree_less);
+ bpf_spin_unlock(&ufq_sd->lock);
+ }
+
+ if (!ret) {
+ stat_add(UFQ_SIMP_INSERT_CNT, 1);
+ stat_add(UFQ_SIMP_INSERT_SIZE, rq->__data_len);
+ }
+ return ret;
+}
+
+struct request *BPF_STRUCT_OPS(ufq_simple_dispatch_req, struct request_queue *q)
+{
+ struct bpf_rb_node *rb_node = NULL;
+ struct bpf_list_node *list_node;
+ struct ufq_simple_data *ufq_sd;
+ struct queue_list_node *qnode;
+ struct sort_tree_node *snode;
+ struct request *rq = NULL;
+ int id = q->id;
+
+ ufq_sd = bpf_map_lookup_elem(&ufq_map, &id);
+ if (!ufq_sd) {
+ bpf_printk("ufq_simple/dispatch_req: ufq_map lookup %d failed", id);
+ return NULL;
+ }
+
+ bpf_spin_lock(&ufq_sd->lock);
+ list_node = bpf_list_pop_front(&ufq_sd->dispatch);
+
+ if (list_node) {
+ qnode = container_of(list_node, struct queue_list_node, node);
+ rq = bpf_kptr_xchg(&qnode->req, NULL);
+ bpf_spin_unlock(&ufq_sd->lock);
+ bpf_obj_drop(qnode);
+ } else {
+ rb_node = bpf_rbtree_first(&ufq_sd->sort_tree_read);
+ if (rb_node) {
+ rb_node = bpf_rbtree_remove(&ufq_sd->sort_tree_read, rb_node);
+ } else {
+ rb_node = bpf_rbtree_first(&ufq_sd->sort_tree_write);
+ if (rb_node)
+ rb_node = bpf_rbtree_remove(&ufq_sd->sort_tree_write, rb_node);
+ }
+
+ if (!rb_node) {
+ bpf_spin_unlock(&ufq_sd->lock);
+ goto out;
+ }
+
+ snode = container_of(rb_node, struct sort_tree_node, rb_node);
+ rq = bpf_kptr_xchg(&snode->req, NULL);
+ bpf_spin_unlock(&ufq_sd->lock);
+ bpf_obj_drop(snode);
+ }
+ if (!rq)
+ bpf_printk("ufq_simple/dispatch_req: no request to dispatch");
+
+out:
+ if (rq) {
+ stat_add(UFQ_SIMP_DISPATCH_CNT, 1);
+ stat_add(UFQ_SIMP_DISPATCH_SIZE, rq->__data_len);
+ }
+
+ return rq;
+}
+
+bool BPF_STRUCT_OPS(ufq_simple_has_req, struct request_queue *q, int rqs_count)
+{
+ struct ufq_simple_data *ufq_sd;
+ int id = q->id;
+ bool has;
+
+ ufq_sd = bpf_map_lookup_elem(&ufq_map, &id);
+ if (!ufq_sd) {
+ bpf_printk("ufq_simple/has_req: ufq_map lookup %d failed", id);
+ return false;
+ }
+
+ bpf_spin_lock(&ufq_sd->lock);
+ has = !bpf_list_empty(&ufq_sd->dispatch) ||
+ bpf_rbtree_root(&ufq_sd->sort_tree_read) ||
+ bpf_rbtree_root(&ufq_sd->sort_tree_write);
+ bpf_spin_unlock(&ufq_sd->lock);
+
+ return has;
+}
+
+void BPF_STRUCT_OPS(ufq_simple_finish_req, struct request *rq)
+{
+ if (rq) {
+ stat_add(UFQ_SIMP_FINISH_CNT, 1);
+ stat_add(UFQ_SIMP_FINISH_SIZE, (u64)rq->stats_sectors << SECTOR_SHIFT);
+ }
+}
+
+struct request *BPF_STRUCT_OPS(ufq_simple_merge_req, struct request_queue *q,
+ struct request *rq, int *type)
+{
+ sector_t rq_start, rq_end, other_start, other_end;
+ enum elv_merge mt = ELEVATOR_NO_MERGE;
+ struct sort_tree_node *snode = NULL;
+ struct bpf_rb_node *rb_node = NULL;
+ struct blk_mq_hw_ctx *targ_mq_hctx;
+ struct blk_mq_ctx *targ_mq_ctx;
+ struct ufq_simple_data *ufq_sd;
+ struct request *targ = NULL;
+ enum ufq_simp_data_dir dir;
+ struct bpf_rb_root *tree;
+ int id = q->id, count = 0;
+
+ *type = ELEVATOR_NO_MERGE;
+ dir = ((rq->cmd_flags & REQ_OP_MASK) & 1) ? UFQ_SIMP_WRITE : UFQ_SIMP_READ;
+ ufq_sd = bpf_map_lookup_elem(&ufq_map, &id);
+ if (!ufq_sd)
+ return NULL;
+
+ rq_start = rq->__sector;
+ rq_end = rq_start + (rq->__data_len >> SECTOR_SHIFT);
+
+ if (dir == UFQ_SIMP_READ)
+ tree = &ufq_sd->sort_tree_read;
+ else
+ tree = &ufq_sd->sort_tree_write;
+
+ bpf_spin_lock(&ufq_sd->lock);
+ rb_node = bpf_rbtree_root(tree);
+ if (!rb_node) {
+ bpf_spin_unlock(&ufq_sd->lock);
+ return NULL;
+ }
+
+ while (mt == ELEVATOR_NO_MERGE && rb_node && count < UFQ_LOOP_MAX) {
+ count++;
+ snode = container_of(rb_node, struct sort_tree_node, rb_node);
+ targ = bpf_kptr_xchg(&snode->req, NULL);
+ if (!targ)
+ break;
+
+ other_start = targ->__sector;
+ other_end = other_start + (targ->__data_len >> SECTOR_SHIFT);
+ targ_mq_ctx = targ->mq_ctx;
+ targ_mq_hctx = targ->mq_hctx;
+
+ targ = bpf_kptr_xchg(&snode->req, targ);
+ if (targ) {
+ bpf_spin_unlock(&ufq_sd->lock);
+ bpf_request_release(targ);
+ return NULL;
+ }
+
+ if (rq_start > other_end)
+ rb_node = bpf_rbtree_right(tree, rb_node);
+ else if (rq_end < other_start)
+ rb_node = bpf_rbtree_left(tree, rb_node);
+ else if (rq_end == other_start)
+ mt = ELEVATOR_FRONT_MERGE;
+ else if (other_end == rq_start)
+ mt = ELEVATOR_BACK_MERGE;
+ else
+ break;
+
+ if (mt) {
+ if (rq->mq_ctx != targ_mq_ctx || rq->mq_hctx != targ_mq_hctx) {
+ mt = ELEVATOR_NO_MERGE;
+ break;
+ }
+
+ rb_node = bpf_rbtree_remove(tree, rb_node);
+ if (rb_node) {
+ snode = container_of(rb_node,
+ struct sort_tree_node, rb_node);
+ targ = bpf_kptr_xchg(&snode->req, NULL);
+ bpf_spin_unlock(&ufq_sd->lock);
+ if (targ) {
+ *type = mt;
+ stat_add(UFQ_SIMP_RQMERGE_CNT, 1);
+ stat_add(UFQ_SIMP_RQMERGE_SIZE, targ->__data_len);
+ stat_sub(UFQ_SIMP_INSERT_CNT, 1);
+ stat_sub(UFQ_SIMP_INSERT_SIZE, targ->__data_len);
+ }
+
+ bpf_obj_drop(snode);
+ } else {
+ bpf_spin_unlock(&ufq_sd->lock);
+ *type = ELEVATOR_NO_MERGE;
+ }
+ return targ;
+ }
+ }
+ bpf_spin_unlock(&ufq_sd->lock);
+
+ return NULL;
+}
+
+static struct request *merge_bio_left_unlock(struct ufq_simple_data *ufq_sd,
+ struct bpf_rb_root *tree,
+ struct sort_tree_node *snode,
+ struct request *cand)
+{
+ struct sort_tree_node *left_node = NULL, *owned;
+ sector_t cand_start, left_start, left_end;
+ struct bpf_rb_node *tmp, *removed = NULL;
+ struct request *left_rq = NULL;
+ bool merged = false;
+
+ cand_start = cand->__sector;
+ tmp = bpf_rbtree_left(tree, &snode->rb_node);
+ if (!tmp)
+ goto end;
+
+ left_node = container_of(tmp, struct sort_tree_node, rb_node);
+ if (!left_node)
+ goto end;
+
+ left_rq = bpf_kptr_xchg(&left_node->req, NULL);
+ if (!left_rq)
+ goto end;
+
+ left_start = left_rq->__sector;
+ left_end = left_start + (left_rq->__data_len >> SECTOR_SHIFT);
+
+ if (left_end == cand_start)
+ merged = !!bpf_request_try_merge(left_rq, cand);
+
+ if (merged) {
+ removed = bpf_rbtree_remove(tree, &snode->rb_node);
+ left_rq = bpf_kptr_xchg(&left_node->req, left_rq);
+ bpf_spin_unlock(&ufq_sd->lock);
+ if (removed) {
+ struct sort_tree_node *drop = container_of(removed,
+ struct sort_tree_node, rb_node);
+ bpf_obj_drop(drop);
+ }
+ stat_add(UFQ_SIMP_RQMERGE_CNT, 1);
+ stat_add(UFQ_SIMP_RQMERGE_SIZE, cand->__data_len);
+
+ if (left_rq)
+ bpf_request_release(left_rq);
+
+ /* Transfer the owned kptr, not the untracked kfunc result. */
+ return cand;
+ }
+
+ left_rq = bpf_kptr_xchg(&left_node->req, left_rq);
+
+end:
+ cand = bpf_kptr_xchg(&snode->req, cand);
+ /*
+ * A front bio merge changes the request sector. Reposition its node
+ * before another tree walk compares that sector with the sort key.
+ */
+ if (snode->key != cand_start) {
+ removed = bpf_rbtree_remove(tree, &snode->rb_node);
+ if (removed) {
+ owned = container_of(removed, struct sort_tree_node,
+ rb_node);
+ owned->key = cand_start;
+ bpf_rbtree_add(tree, &owned->rb_node, sort_tree_less);
+ }
+ }
+ bpf_spin_unlock(&ufq_sd->lock);
+ if (left_rq)
+ bpf_request_release(left_rq);
+
+ if (cand)
+ bpf_request_release(cand);
+
+ return NULL;
+}
+
+static struct request *merge_bio_right_unlock(struct ufq_simple_data *ufq_sd,
+ struct bpf_rb_root *tree,
+ struct sort_tree_node *snode,
+ struct request *cand)
+{
+ struct bpf_rb_node *right_rb, *removed = NULL;
+ struct sort_tree_node *right_node = NULL;
+ struct request *right_rq = NULL;
+ sector_t cand_end, right_start;
+ bool merged = false;
+
+ cand_end = cand->__sector + (cand->__data_len >> SECTOR_SHIFT);
+
+ right_rb = bpf_rbtree_right(tree, &snode->rb_node);
+ if (!right_rb)
+ goto end;
+
+ right_node = container_of(right_rb, struct sort_tree_node, rb_node);
+ if (!right_node)
+ goto end;
+
+ right_rq = bpf_kptr_xchg(&right_node->req, NULL);
+ if (!right_rq)
+ goto end;
+
+ right_start = right_rq->__sector;
+ if (cand_end == right_start)
+ merged = !!bpf_request_try_merge(cand, right_rq);
+
+ if (merged) {
+ removed = bpf_rbtree_remove(tree, right_rb);
+ cand = bpf_kptr_xchg(&snode->req, cand);
+ bpf_spin_unlock(&ufq_sd->lock);
+ if (removed) {
+ struct sort_tree_node *drop = container_of(removed,
+ struct sort_tree_node, rb_node);
+ bpf_obj_drop(drop);
+ }
+ stat_add(UFQ_SIMP_RQMERGE_CNT, 1);
+ stat_add(UFQ_SIMP_RQMERGE_SIZE, right_rq->__data_len);
+
+ if (cand)
+ bpf_request_release(cand);
+
+ /* Transfer the owned kptr, not the untracked kfunc result. */
+ return right_rq;
+ }
+
+ right_rq = bpf_kptr_xchg(&right_node->req, right_rq);
+
+end:
+ cand = bpf_kptr_xchg(&snode->req, cand);
+ bpf_spin_unlock(&ufq_sd->lock);
+ if (right_rq)
+ bpf_request_release(right_rq);
+
+ if (cand)
+ bpf_request_release(cand);
+
+ return NULL;
+}
+
+struct request *BPF_STRUCT_OPS(ufq_simple_merge_bio,
+ struct request_queue *q, struct bio *bio,
+ unsigned int nr_segs, bool *merged)
+{
+ struct request *cand = NULL, *old, *free = NULL;
+ sector_t start, end, cand_start, cand_end;
+ struct sort_tree_node *snode = NULL;
+ struct bpf_rb_node *rb_node = NULL;
+ struct ufq_simple_data *ufq_sd;
+ enum ufq_simp_data_dir dir;
+ struct bpf_rb_root *tree;
+ int id = q->id, count = 0;
+
+ if (!merged)
+ return NULL;
+ start = bio->bi_iter.bi_sector;
+ end = start + (bio->bi_iter.bi_size >> SECTOR_SHIFT);
+ dir = ((bio->bi_opf & REQ_OP_MASK) & 1) ? UFQ_SIMP_WRITE : UFQ_SIMP_READ;
+ ufq_sd = bpf_map_lookup_elem(&ufq_map, &id);
+ if (!ufq_sd)
+ return NULL;
+
+ bpf_spin_lock(&ufq_sd->lock);
+ if (dir == UFQ_SIMP_READ)
+ tree = &ufq_sd->sort_tree_read;
+ else
+ tree = &ufq_sd->sort_tree_write;
+
+ rb_node = bpf_rbtree_root(tree);
+ while (rb_node && count < UFQ_LOOP_MAX) {
+ count++;
+ snode = container_of(rb_node, struct sort_tree_node, rb_node);
+ cand = bpf_kptr_xchg(&snode->req, NULL);
+ if (!cand)
+ break;
+
+ cand_start = cand->__sector;
+ cand_end = cand_start + (cand->__data_len >> SECTOR_SHIFT);
+
+ if (end < cand_start) {
+ rb_node = bpf_rbtree_left(tree, rb_node);
+ } else if (start > cand_end) {
+ rb_node = bpf_rbtree_right(tree, rb_node);
+ } else if (cand_start == end) {
+ if (bpf_request_bio_try_merge(cand, bio, nr_segs)) {
+ *merged = true;
+ free = merge_bio_left_unlock(ufq_sd, tree, snode, cand);
+ stat_add(UFQ_SIMP_BIOMERGE_CNT, 1);
+ stat_add(UFQ_SIMP_BIOMERGE_SIZE, bio->bi_iter.bi_size);
+ return free;
+ }
+ rb_node = NULL;
+ } else if (cand_end == start) {
+ if (bpf_request_bio_try_merge(cand, bio, nr_segs)) {
+ *merged = true;
+ free = merge_bio_right_unlock(ufq_sd, tree, snode, cand);
+ stat_add(UFQ_SIMP_BIOMERGE_CNT, 1);
+ stat_add(UFQ_SIMP_BIOMERGE_SIZE, bio->bi_iter.bi_size);
+ return free;
+ }
+ rb_node = NULL;
+ } else {
+ rb_node = NULL;
+ }
+
+ old = bpf_kptr_xchg(&snode->req, cand);
+ if (old) {
+ bpf_spin_unlock(&ufq_sd->lock);
+ bpf_request_release(old);
+ return NULL;
+ }
+ }
+
+ bpf_spin_unlock(&ufq_sd->lock);
+ return NULL;
+}
+
+UFQ_OPS_DEFINE(ufq_simple_ops,
+ .init_sched = (void *)ufq_simple_init_sched,
+ .exit_sched = (void *)ufq_simple_exit_sched,
+ .insert_req = (void *)ufq_simple_insert_req,
+ .dispatch_req = (void *)ufq_simple_dispatch_req,
+ .has_req = (void *)ufq_simple_has_req,
+ .finish_req = (void *)ufq_simple_finish_req,
+ .merge_req = (void *)ufq_simple_merge_req,
+ .merge_bio = (void *)ufq_simple_merge_bio,
+ .name = "ufq_simple");
diff --git a/tools/ufq_iosched/ufq_simple.c b/tools/ufq_iosched/ufq_simple.c
new file mode 100644
index 000000000000..0e7b018699a2
--- /dev/null
+++ b/tools/ufq_iosched/ufq_simple.c
@@ -0,0 +1,120 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#include <stdio.h>
+#include <stdlib.h>
+#include <string.h>
+#include <unistd.h>
+#include <signal.h>
+#include <libgen.h>
+#include <bpf/bpf.h>
+#include <ufq/common.h>
+#include "ufq_simple.bpf.skel.h"
+
+const char help_fmt[] =
+"A simple ufq scheduler.\n"
+"\n"
+"Usage: %s [-v] [-d] [-h]\n"
+"\n"
+" -v Print version\n"
+" -d Print libbpf debug messages\n"
+" -h Display this help and exit\n";
+
+#define UFQ_SIMPLE_VERSION "0.1.0"
+#define TIME_INTERVAL 3
+static bool verbose;
+static volatile int exit_req;
+__u64 old_stats[UFQ_SIMP_STAT_MAX];
+
+static int libbpf_print_fn(enum libbpf_print_level level, const char *format, va_list args)
+{
+ if (level == LIBBPF_DEBUG && !verbose)
+ return 0;
+ return vfprintf(stderr, format, args);
+}
+
+static void sigint_handler(int simple)
+{
+ exit_req = 1;
+}
+
+static void read_stats(struct ufq_simple *skel, __u64 *stats)
+{
+ int nr_cpus = libbpf_num_possible_cpus();
+ __u64 cnts[UFQ_SIMP_STAT_MAX][nr_cpus];
+ __u32 idx;
+
+ memset(stats, 0, sizeof(stats[0]) * UFQ_SIMP_STAT_MAX);
+
+ for (idx = 0; idx < UFQ_SIMP_STAT_MAX; idx++) {
+ int ret, cpu;
+
+ ret = bpf_map_lookup_elem(bpf_map__fd(skel->maps.stats),
+ &idx, cnts[idx]);
+ if (ret < 0)
+ continue;
+ for (cpu = 0; cpu < nr_cpus; cpu++)
+ stats[idx] += cnts[idx][cpu];
+ }
+}
+
+int main(int argc, char **argv)
+{
+ struct ufq_simple *skel;
+ struct bpf_link *link;
+ __u32 opt;
+
+ libbpf_set_print(libbpf_print_fn);
+ signal(SIGINT, sigint_handler);
+ signal(SIGTERM, sigint_handler);
+
+ skel = UFQ_OPS_OPEN(ufq_simple_ops, ufq_simple);
+
+ while ((opt = getopt(argc, argv, "vdh")) != -1) {
+ switch (opt) {
+ case 'v':
+ printf("ufq_simple version: %s\n", UFQ_SIMPLE_VERSION);
+ return 0;
+ case 'd':
+ verbose = true;
+ break;
+ default:
+ fprintf(stderr, help_fmt, basename(argv[0]));
+ return opt != 'h';
+ }
+ }
+
+ UFQ_OPS_LOAD(skel, ufq_simple_ops, ufq_simple);
+ link = UFQ_OPS_ATTACH(skel, ufq_simple_ops, ufq_simple);
+
+ printf("ufq_simple loop ...\n");
+ while (!exit_req) {
+ __u64 stats[UFQ_SIMP_STAT_MAX];
+
+ printf("--------------------------------\n");
+ read_stats(skel, stats);
+ printf("bps:%lluk iops:%llu\n",
+ (stats[UFQ_SIMP_FINISH_SIZE] -
+ old_stats[UFQ_SIMP_FINISH_SIZE]) / 1024 / TIME_INTERVAL,
+ (stats[UFQ_SIMP_FINISH_CNT] -
+ old_stats[UFQ_SIMP_FINISH_CNT]) / TIME_INTERVAL);
+ printf("(insert: cnt=%llu size=%llu)\n",
+ stats[UFQ_SIMP_INSERT_CNT], stats[UFQ_SIMP_INSERT_SIZE]);
+ printf("(rqmerge: cnt=%llu size=%llu) (biomerge: cnt=%llu size=%llu)\n",
+ stats[UFQ_SIMP_RQMERGE_CNT], stats[UFQ_SIMP_RQMERGE_SIZE],
+ stats[UFQ_SIMP_BIOMERGE_CNT], stats[UFQ_SIMP_BIOMERGE_SIZE]);
+ printf("(dispatch: cnt=%llu size=%llu) (finish: cnt=%llu size=%llu)\n",
+ stats[UFQ_SIMP_DISPATCH_CNT], stats[UFQ_SIMP_DISPATCH_SIZE],
+ stats[UFQ_SIMP_FINISH_CNT], stats[UFQ_SIMP_FINISH_SIZE]);
+ memcpy(old_stats, stats, sizeof(old_stats));
+ sleep(TIME_INTERVAL);
+ }
+
+ printf("ufq_simple loop exit ...\n");
+ bpf_link__destroy(link);
+ ufq_simple__destroy(skel);
+
+ return 0;
+}
--
2.53.0
^ permalink raw reply [flat|nested] 6+ messages in thread* [RFC v3 1/3] block: Introduce the UFQ I/O scheduler
2026-10-03 4:27 [RFC v3 0/3] block: Introduce a BPF-based I/O scheduler Kaitao Cheng
2026-10-03 4:27 ` [RFC v3 2/3] tools/ufq_iosched: add BPF example scheduler and build scaffolding Kaitao Cheng
@ 2026-10-03 4:27 ` Kaitao Cheng
2026-10-03 4:27 ` [RFC v3 3/3] tools/ufq_iosched: add PFQ eBPF " Kaitao Cheng
2026-10-03 9:36 ` [RFC v3 0/3] block: Introduce a BPF-based " Alexei Starovoitov
3 siblings, 0 replies; 6+ messages in thread
From: Kaitao Cheng @ 2026-10-03 4:27 UTC (permalink / raw)
To: Jens Axboe; +Cc: linux-block, linux-kernel, bpf, Kaitao Cheng
From: Kaitao Cheng <chengkaitao@kylinos.cn>
Introduce IOSCHED_UFQ, a blk-mq elevator ("ufq: User-programmable
Flexible Queueing") whose policy is supplied by an eBPF program via
struct_ops (insert, dispatch, merge, finish, etc.).
When no eBPF program is attached, the UFQ I/O scheduler uses a simple,
per-ctx queueing policy (similar to none). After an eBPF program is
attached, the user-defined scheduling policy replaces UFQ’s built-in
queueing policy, while per-ctx queues remain available as a fallback
mechanism.
Signed-off-by: Kaitao Cheng <chengkaitao@kylinos.cn>
---
block/Kconfig.iosched | 8 +
block/Makefile | 1 +
block/blk-merge.c | 28 +-
block/blk-mq.c | 8 +-
block/blk-mq.h | 2 +-
block/blk.h | 5 +
block/ufq-bpfops.c | 267 ++++++++++++++
block/ufq-iosched.c | 789 ++++++++++++++++++++++++++++++++++++++++++
block/ufq-iosched.h | 79 +++++
block/ufq-kfunc.c | 169 +++++++++
10 files changed, 1343 insertions(+), 13 deletions(-)
create mode 100644 block/ufq-bpfops.c
create mode 100644 block/ufq-iosched.c
create mode 100644 block/ufq-iosched.h
create mode 100644 block/ufq-kfunc.c
diff --git a/block/Kconfig.iosched b/block/Kconfig.iosched
index 27f11320b8d1..56afc425cc52 100644
--- a/block/Kconfig.iosched
+++ b/block/Kconfig.iosched
@@ -44,4 +44,12 @@ config BFQ_CGROUP_DEBUG
Enable some debugging help. Currently it exports additional stat
files in a cgroup which can be useful for debugging.
+config IOSCHED_UFQ
+ tristate "UFQ I/O scheduler"
+ default y
+ help
+ The UFQ I/O scheduler is a programmable I/O scheduler. When
+ enabled, an out-of-kernel I/O scheduler based on eBPF can be
+ designed to interact with it, leveraging its customizable
+ hooks to redefine I/O scheduling policies.
endmenu
diff --git a/block/Makefile b/block/Makefile
index e7bd320e3d69..3b627c63e749 100644
--- a/block/Makefile
+++ b/block/Makefile
@@ -27,6 +27,7 @@ obj-$(CONFIG_MQ_IOSCHED_DEADLINE) += mq-deadline.o
obj-$(CONFIG_MQ_IOSCHED_KYBER) += kyber-iosched.o
bfq-y := bfq-iosched.o bfq-wf2q.o bfq-cgroup.o
obj-$(CONFIG_IOSCHED_BFQ) += bfq.o
+obj-$(CONFIG_IOSCHED_UFQ) += ufq-iosched.o ufq-bpfops.o ufq-kfunc.o
obj-$(CONFIG_BLK_DEV_INTEGRITY) += bio-integrity.o blk-integrity.o t10-pi.o \
bio-integrity-auto.o bio-integrity-fs.o
diff --git a/block/blk-merge.c b/block/blk-merge.c
index 258a726071d1..0a89b476ce71 100644
--- a/block/blk-merge.c
+++ b/block/blk-merge.c
@@ -771,8 +771,8 @@ u8 bio_seg_gap(struct request_queue *q, struct bio *prev, struct bio *next,
* For non-mq, this has to be called with the request spinlock acquired.
* For mq with scheduling, the appropriate queue wide lock should be held.
*/
-static struct request *attempt_merge(struct request_queue *q,
- struct request *req, struct request *next)
+static struct request *attempt_merge(struct request_queue *q, struct request *req,
+ struct request *next, bool nohash)
{
if (!rq_mergeable(req) || !rq_mergeable(next))
return NULL;
@@ -839,7 +839,7 @@ static struct request *attempt_merge(struct request_queue *q,
req->__data_len += blk_rq_bytes(next);
- if (!blk_discard_mergable(req))
+ if (!nohash && !blk_discard_mergable(req))
elv_merge_requests(q, req, next);
blk_crypto_rq_put_keyslot(next);
@@ -865,7 +865,7 @@ static struct request *attempt_back_merge(struct request_queue *q,
struct request *next = elv_latter_request(q, rq);
if (next)
- return attempt_merge(q, rq, next);
+ return attempt_merge(q, rq, next, false);
return NULL;
}
@@ -876,11 +876,17 @@ static struct request *attempt_front_merge(struct request_queue *q,
struct request *prev = elv_former_request(q, rq);
if (prev)
- return attempt_merge(q, prev, rq);
+ return attempt_merge(q, prev, rq, false);
return NULL;
}
+struct request *bpf_attempt_merge(struct request_queue *q, struct request *rq,
+ struct request *next)
+{
+ return attempt_merge(q, rq, next, true);
+}
+
/*
* Try to merge 'next' into 'rq'. Return true if the merge happened, false
* otherwise. The caller is responsible for freeing 'next' if the merge
@@ -889,7 +895,7 @@ static struct request *attempt_front_merge(struct request_queue *q,
bool blk_attempt_req_merge(struct request_queue *q, struct request *rq,
struct request *next)
{
- return attempt_merge(q, rq, next);
+ return attempt_merge(q, rq, next, false);
}
bool blk_rq_merge_ok(struct request *rq, struct bio *bio)
@@ -1032,11 +1038,11 @@ static enum bio_merge_status bio_attempt_discard_merge(struct request_queue *q,
return BIO_MERGE_FAILED;
}
-static enum bio_merge_status blk_attempt_bio_merge(struct request_queue *q,
- struct request *rq,
- struct bio *bio,
- unsigned int nr_segs,
- bool sched_allow_merge)
+enum bio_merge_status blk_attempt_bio_merge(struct request_queue *q,
+ struct request *rq,
+ struct bio *bio,
+ unsigned int nr_segs,
+ bool sched_allow_merge)
{
if (!blk_rq_merge_ok(rq, bio))
return BIO_MERGE_NONE;
diff --git a/block/blk-mq.c b/block/blk-mq.c
index a26a11c73ee3..533936beaf38 100644
--- a/block/blk-mq.c
+++ b/block/blk-mq.c
@@ -796,7 +796,7 @@ static void blk_mq_finish_request(struct request *rq)
}
}
-static void __blk_mq_free_request(struct request *rq)
+void __blk_mq_free_request(struct request *rq)
{
struct request_queue *q = rq->q;
struct blk_mq_ctx *ctx = rq->mq_ctx;
@@ -1813,6 +1813,12 @@ static bool dispatch_rq_from_ctx(struct sbitmap *sb, unsigned int bitnr,
if (list_empty(&ctx->rq_lists[type]))
sbitmap_clear_bit(sb, bitnr);
}
+
+ if (dispatch_data->rq) {
+ dispatch_data->rq->rq_flags |= RQF_STARTED;
+ if (hctx->queue->last_merge == dispatch_data->rq)
+ hctx->queue->last_merge = NULL;
+ }
spin_unlock(&ctx->lock);
return !dispatch_data->rq;
diff --git a/block/blk-mq.h b/block/blk-mq.h
index aa15d31aaae9..3f85cae7bf57 100644
--- a/block/blk-mq.h
+++ b/block/blk-mq.h
@@ -56,7 +56,7 @@ void blk_mq_flush_busy_ctxs(struct blk_mq_hw_ctx *hctx, struct list_head *list);
struct request *blk_mq_dequeue_from_ctx(struct blk_mq_hw_ctx *hctx,
struct blk_mq_ctx *start);
void blk_mq_put_rq_ref(struct request *rq);
-
+void __blk_mq_free_request(struct request *rq);
/*
* Internal helpers for allocating/freeing the request map
*/
diff --git a/block/blk.h b/block/blk.h
index 50abfd932886..54d5a860d5e5 100644
--- a/block/blk.h
+++ b/block/blk.h
@@ -340,6 +340,9 @@ enum bio_merge_status {
enum bio_merge_status bio_attempt_back_merge(struct request *req,
struct bio *bio, unsigned int nr_segs);
+enum bio_merge_status blk_attempt_bio_merge(struct request_queue *q,
+ struct request *rq, struct bio *bio, unsigned int nr_segs,
+ bool sched_allow_merge);
bool blk_attempt_plug_merge(struct request_queue *q, struct bio *bio,
unsigned int nr_segs);
bool blk_bio_list_merge(struct request_queue *q, struct list_head *list,
@@ -473,6 +476,8 @@ static inline unsigned get_max_segment_size(const struct queue_limits *lim,
int ll_back_merge_fn(struct request *req, struct bio *bio,
unsigned int nr_segs);
+struct request *bpf_attempt_merge(struct request_queue *q, struct request *rq,
+ struct request *next);
bool blk_attempt_req_merge(struct request_queue *q, struct request *rq,
struct request *next);
unsigned int blk_recalc_rq_segments(struct request *rq);
diff --git a/block/ufq-bpfops.c b/block/ufq-bpfops.c
new file mode 100644
index 000000000000..5efcfcaa9cb8
--- /dev/null
+++ b/block/ufq-bpfops.c
@@ -0,0 +1,267 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#include <linux/init.h>
+#include <linux/types.h>
+#include <linux/bpf_verifier.h>
+#include <linux/bpf.h>
+#include <linux/btf.h>
+#include <linux/btf_ids.h>
+#include <linux/string.h>
+#include <linux/wait.h>
+#include <linux/rcupdate.h>
+#include <linux/mutex.h>
+#include "ufq-iosched.h"
+
+struct ufq_iosched_ops ufq_ops;
+#define UFQ_BPFOPS_ENABLED BIT(30)
+/* The enable bit and user count must change as one atomic state. */
+static atomic_t ufq_bpfops_state;
+static DECLARE_WAIT_QUEUE_HEAD(ufq_bpfops_wq);
+/* Serialize the whole attach/detach, including queue freeze and restart. */
+static DEFINE_MUTEX(ufq_bpfops_lock);
+static void *ufq_bpfops_owner;
+
+const struct ufq_iosched_ops *ufq_bpfops_tryget(void)
+{
+ int state = atomic_read(&ufq_bpfops_state);
+
+ do {
+ if (!(state & UFQ_BPFOPS_ENABLED))
+ return NULL;
+ if (WARN_ON_ONCE((state & ~UFQ_BPFOPS_ENABLED) ==
+ UFQ_BPFOPS_ENABLED - 1))
+ return NULL;
+ } while (!atomic_try_cmpxchg(&ufq_bpfops_state, &state, state + 1));
+
+ return &ufq_ops;
+}
+
+void ufq_bpfops_put(void)
+{
+ if (atomic_dec_and_test(&ufq_bpfops_state))
+ wake_up_all(&ufq_bpfops_wq);
+}
+
+bool ufq_bpfops_draining(void)
+{
+ int state = atomic_read(&ufq_bpfops_state);
+
+ return !(state & UFQ_BPFOPS_ENABLED) && state;
+}
+
+static const struct bpf_func_proto *
+bpf_ufq_get_func_proto(enum bpf_func_id func_id, const struct bpf_prog *prog)
+{
+ return bpf_base_func_proto(func_id, prog);
+}
+
+static bool bpf_ufq_is_valid_access(int off, int size,
+ enum bpf_access_type type,
+ const struct bpf_prog *prog,
+ struct bpf_insn_access_aux *info)
+{
+ if (type != BPF_READ)
+ return false;
+ if (off < 0 || off >= sizeof(__u64) * MAX_BPF_FUNC_ARGS)
+ return false;
+ if (off % size != 0)
+ return false;
+
+ /*
+ * btf_ctx_access() treats pointers that are not "pointer to struct"
+ * as scalars (no reg_type). These two output arguments point to a
+ * single int or bool on the kernel stack. PTR_TO_BUF would permit
+ * unbounded offsets for struct_ops programs, so give the verifier the
+ * actual pointee size instead.
+ */
+ if (size == sizeof(__u64) && prog->aux->attach_func_name &&
+ ((!strcmp(prog->aux->attach_func_name, "merge_req") && off == 16) ||
+ (!strcmp(prog->aux->attach_func_name, "merge_bio") && off == 24))) {
+ if (!btf_ctx_access(off, size, type, prog, info))
+ return false;
+ info->reg_type = PTR_TO_BUF;
+ return true;
+ }
+
+ return btf_ctx_access(off, size, type, prog, info);
+}
+
+static const struct bpf_verifier_ops bpf_ufq_verifier_ops = {
+ .get_func_proto = bpf_ufq_get_func_proto,
+ .is_valid_access = bpf_ufq_is_valid_access,
+};
+
+static int bpf_ufq_init_member(const struct btf_type *t,
+ const struct btf_member *member,
+ void *kdata, const void *udata)
+{
+ const struct ufq_iosched_ops *uops = udata;
+ struct ufq_iosched_ops *ops = kdata;
+ u32 moff = __btf_member_bit_offset(t, member) / 8;
+ int ret;
+
+ switch (moff) {
+ case offsetof(struct ufq_iosched_ops, name):
+ ret = bpf_obj_name_cpy(ops->name, uops->name,
+ sizeof(ops->name));
+ if (ret < 0)
+ return ret;
+ if (ret == 0)
+ return -EINVAL;
+ return 1;
+ /* other var adding .... */
+ }
+
+ return 0;
+}
+
+static int bpf_ufq_check_member(const struct btf_type *t,
+ const struct btf_member *member,
+ const struct bpf_prog *prog)
+{
+ return 0;
+}
+
+static int bpf_ufq_enable(void *ops)
+{
+ ufq_ops = *(struct ufq_iosched_ops *)ops;
+ atomic_set_release(&ufq_bpfops_state, UFQ_BPFOPS_ENABLED);
+ return 0;
+}
+
+static void bpf_ufq_disable(struct ufq_iosched_ops *ops)
+{
+ atomic_fetch_andnot(UFQ_BPFOPS_ENABLED, &ufq_bpfops_state);
+ wait_event(ufq_bpfops_wq, !atomic_read(&ufq_bpfops_state));
+ memset(&ufq_ops, 0, sizeof(ufq_ops));
+}
+
+static int bpf_ufq_reg(void *kdata, struct bpf_link *link)
+{
+ int ret;
+
+ mutex_lock(&ufq_bpfops_lock);
+ if (ufq_bpfops_owner) {
+ ret = -EBUSY;
+ goto out;
+ }
+
+ ret = ufq_prepare_bpf_attach(bpf_ufq_enable, kdata);
+ if (!ret)
+ ufq_bpfops_owner = kdata;
+out:
+ mutex_unlock(&ufq_bpfops_lock);
+ return ret;
+}
+
+static void bpf_ufq_unreg(void *kdata, struct bpf_link *link)
+{
+ mutex_lock(&ufq_bpfops_lock);
+ if (WARN_ON_ONCE(ufq_bpfops_owner != kdata))
+ goto out;
+
+ bpf_ufq_disable(kdata);
+ ufq_kick_all_hw_queues();
+ ufq_bpfops_owner = NULL;
+out:
+ mutex_unlock(&ufq_bpfops_lock);
+}
+
+static int bpf_ufq_init(struct btf *btf)
+{
+ return 0;
+}
+
+static int bpf_ufq_update(void *kdata, void *old_kdata, struct bpf_link *link)
+{
+ /*
+ * UFQ does not support live-updating an already-attached BPF scheduler:
+ * partial failure during callback setup (e.g. init_sched) would be hard
+ * to reason about, and update can race with unregister/teardown.
+ */
+ return -EOPNOTSUPP;
+}
+
+static int bpf_ufq_validate(void *kdata)
+{
+ return 0;
+}
+
+static int init_sched_stub(struct request_queue *q)
+{
+ return -EPERM;
+}
+
+static int exit_sched_stub(struct request_queue *q)
+{
+ return -EPERM;
+}
+
+static int insert_req_stub(struct request_queue *q, struct request *rq,
+ blk_insert_t flags)
+{
+ return 0;
+}
+
+static struct request *dispatch_req_stub(struct request_queue *q)
+{
+ return NULL;
+}
+
+static bool has_req_stub(struct request_queue *q, int rqs_count)
+{
+ return rqs_count > 0;
+}
+
+static void finish_req_stub(struct request *rq)
+{
+}
+
+static struct request *merge_req_stub(struct request_queue *q, struct request *rq,
+ int *type)
+{
+ *type = ELEVATOR_NO_MERGE;
+ return NULL;
+}
+
+static struct request *merge_bio_stub(struct request_queue *q, struct bio *bio,
+ unsigned int nr_segs, bool *merged)
+{
+ if (merged)
+ *merged = false;
+
+ return NULL;
+}
+
+static struct ufq_iosched_ops __bpf_ops_ufq_ops = {
+ .init_sched = init_sched_stub,
+ .exit_sched = exit_sched_stub,
+ .insert_req = insert_req_stub,
+ .dispatch_req = dispatch_req_stub,
+ .has_req = has_req_stub,
+ .merge_req = merge_req_stub,
+ .finish_req = finish_req_stub,
+ .merge_bio = merge_bio_stub,
+};
+
+static struct bpf_struct_ops bpf_iosched_ufq_ops = {
+ .verifier_ops = &bpf_ufq_verifier_ops,
+ .reg = bpf_ufq_reg,
+ .unreg = bpf_ufq_unreg,
+ .check_member = bpf_ufq_check_member,
+ .init_member = bpf_ufq_init_member,
+ .init = bpf_ufq_init,
+ .update = bpf_ufq_update,
+ .validate = bpf_ufq_validate,
+ .name = "ufq_iosched_ops",
+ .owner = THIS_MODULE,
+ .cfi_stubs = &__bpf_ops_ufq_ops
+};
+
+int bpf_ufq_ops_init(void)
+{
+ return register_bpf_struct_ops(&bpf_iosched_ufq_ops, ufq_iosched_ops);
+}
diff --git a/block/ufq-iosched.c b/block/ufq-iosched.c
new file mode 100644
index 000000000000..f7964ccfc1ee
--- /dev/null
+++ b/block/ufq-iosched.c
@@ -0,0 +1,789 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#include <linux/kernel.h>
+#include <linux/fs.h>
+#include <linux/blkdev.h>
+#include <linux/bio.h>
+#include <linux/module.h>
+#include <linux/slab.h>
+#include <linux/init.h>
+#include <linux/compiler.h>
+#include <linux/sbitmap.h>
+#include <linux/workqueue.h>
+
+#include <trace/events/block.h>
+
+#include "elevator.h"
+#include "blk.h"
+#include "blk-mq.h"
+#include "blk-mq-sched.h"
+#include "blk-mq-debugfs.h"
+#include "ufq-iosched.h"
+
+static DEFINE_MUTEX(ufq_active_queues_lock);
+static LIST_HEAD(ufq_active_queues);
+
+enum ufq_priv_state {
+ UFQ_PRIV_NOT_IN_SCHED = 0,
+ UFQ_PRIV_IN_BPF = 1,
+ UFQ_PRIV_IN_UFQ = 2,
+ UFQ_PRIV_IN_SCHED = 3,
+};
+
+/* has_work can run from completion softirq, so all users need irq-safe locking. */
+static inline bool ufq_insert_err_pending(struct ufq_data *ufq)
+{
+ unsigned long flags;
+ bool pending;
+
+ spin_lock_irqsave(&ufq->insert_err_lock, flags);
+ pending = !list_empty(&ufq->insert_err_list);
+ spin_unlock_irqrestore(&ufq->insert_err_lock, flags);
+ return pending;
+}
+
+/*
+ * Move @rq off ctx->rq_lists onto ufq->insert_err_list after BPF insert_req
+ * failed. Keeps BPF-owned (still on ctx) rqs unreachable to err drain.
+ */
+static bool ufq_move_to_insert_err(struct ufq_data *ufq, struct request *rq)
+{
+ struct blk_mq_ctx *ctx = rq->mq_ctx;
+ struct blk_mq_hw_ctx *hctx = rq->mq_hctx;
+ enum hctx_type type = hctx->type;
+ unsigned long flags;
+ int bit;
+
+ spin_lock(&ctx->lock);
+ if (!((uintptr_t)rq->elv.priv[0] & UFQ_PRIV_IN_BPF) ||
+ rq->rq_flags & RQF_STARTED || list_empty(&rq->queuelist)) {
+ spin_unlock(&ctx->lock);
+ return false;
+ }
+ list_del_init(&rq->queuelist);
+ bit = ctx->index_hw[type];
+ if (list_empty(&ctx->rq_lists[type]))
+ sbitmap_clear_bit(&hctx->ctx_map, bit);
+ if (hctx->queue->last_merge == rq)
+ hctx->queue->last_merge = NULL;
+ rq->elv.priv[0] = (void *)((uintptr_t)rq->elv.priv[0]
+ & ~UFQ_PRIV_IN_BPF);
+ spin_unlock(&ctx->lock);
+
+ spin_lock_irqsave(&ufq->insert_err_lock, flags);
+ list_add_tail(&rq->queuelist, &ufq->insert_err_list);
+ spin_unlock_irqrestore(&ufq->insert_err_lock, flags);
+ return true;
+}
+
+/* A BPF policy may accidentally pop a request belonging to another queue. */
+static void ufq_recover_bad_dispatch(struct request *rq)
+{
+ struct request_queue *q = rq->q;
+ struct ufq_data *owner;
+
+ if (!q || !rq->mq_ctx || !rq->mq_hctx ||
+ blk_queue_enter(q, BLK_MQ_REQ_NOWAIT))
+ return;
+ if (ufq_is_queue(q)) {
+ owner = q->elevator->elevator_data;
+ if (ufq_move_to_insert_err(owner, rq))
+ blk_mq_run_hw_queues(q, true);
+ }
+ blk_queue_exit(q);
+}
+
+/* Pop one insert-fail rq; sets RQF_STARTED like blk_mq_dequeue_from_ctx. */
+static struct request *ufq_dispatch_insert_err(struct ufq_data *ufq)
+{
+ struct request *rq;
+ unsigned long flags;
+
+ spin_lock_irqsave(&ufq->insert_err_lock, flags);
+ rq = list_first_entry_or_null(&ufq->insert_err_list, struct request,
+ queuelist);
+ if (rq) {
+ list_del_init(&rq->queuelist);
+ rq->rq_flags |= RQF_STARTED;
+ if (ufq->q->last_merge == rq)
+ ufq->q->last_merge = NULL;
+ }
+ spin_unlock_irqrestore(&ufq->insert_err_lock, flags);
+
+ if (rq)
+ atomic_inc(&ufq->insert_err_handled);
+ return rq;
+}
+
+static struct request *ufq_dispatch_request(struct blk_mq_hw_ctx *hctx)
+{
+ struct ufq_data *ufq = hctx->queue->elevator->elevator_data;
+ const struct ufq_iosched_ops *ops;
+ struct blk_mq_ctx *ctx;
+ struct request *rq = NULL;
+ unsigned short idx;
+
+ /* Failed inserts are no longer on ctx lists, even after BPF detach. */
+ rq = ufq_dispatch_insert_err(ufq);
+ if (rq) {
+ atomic_dec(&ufq->rqs_count);
+ return rq;
+ }
+
+ ops = ufq_bpfops_tryget();
+ if (ops && ops->dispatch_req) {
+ rq = ops->dispatch_req(hctx->queue);
+ if (!rq) {
+ ufq_bpfops_put();
+ atomic_inc(&ufq->ops_stats.dispatch_null_count);
+ return NULL;
+ }
+
+ if (unlikely(rq->q != hctx->queue || !rq->mq_ctx ||
+ !rq->mq_hctx)) {
+ atomic_inc(&ufq->ops_stats.dispatch_bad_state_count);
+ ufq_recover_bad_dispatch(rq);
+ if (WARN_ON_ONCE(req_ref_put_and_test(rq)))
+ __blk_mq_free_request(rq);
+ ufq_bpfops_put();
+ return NULL;
+ }
+
+ /*
+ * The BPF insert_req callback bumps the request's reference
+ * count; dispatch_req returns that same request with an extra
+ * reference held. The kernel must put that reference here,
+ * and the request's refcount is always greater than zero at
+ * this point.
+ */
+ if (WARN_ON_ONCE(req_ref_put_and_test(rq))) {
+ __blk_mq_free_request(rq);
+ ufq_bpfops_put();
+ return NULL;
+ }
+
+ ctx = rq->mq_ctx;
+ spin_lock(&ctx->lock);
+ if (unlikely(blk_mq_rq_state(rq) != MQ_RQ_IDLE ||
+ (rq->rq_flags & RQF_STARTED) ||
+ list_empty(&rq->queuelist))) {
+ atomic_inc(&ufq->ops_stats.dispatch_bad_state_count);
+ spin_unlock(&ctx->lock);
+ ufq_bpfops_put();
+ return NULL;
+ }
+ list_del_init(&rq->queuelist);
+ rq->rq_flags |= RQF_STARTED;
+ if (hctx->queue->last_merge == rq)
+ hctx->queue->last_merge = NULL;
+ if (list_empty(&ctx->rq_lists[rq->mq_hctx->type]))
+ sbitmap_clear_bit(&rq->mq_hctx->ctx_map,
+ ctx->index_hw[rq->mq_hctx->type]);
+ spin_unlock(&ctx->lock);
+ atomic_inc(&ufq->ops_stats.dispatch_ok_count);
+ atomic64_add(blk_rq_sectors(rq), &ufq->ops_stats.dispatch_ok_sectors);
+ rq->elv.priv[0] = (void *)((uintptr_t)rq->elv.priv[0]
+ & ~UFQ_PRIV_IN_UFQ);
+ ufq_bpfops_put();
+ } else {
+ if (ops)
+ ufq_bpfops_put();
+ /* Do not dispatch an in-progress BPF insertion via the fallback. */
+ if (!ops && ufq_bpfops_draining())
+ return NULL;
+ ctx = READ_ONCE(hctx->dispatch_from);
+ rq = blk_mq_dequeue_from_ctx(hctx, ctx);
+ if (rq) {
+ idx = rq->mq_ctx->index_hw[hctx->type];
+ if (++idx == hctx->nr_ctx)
+ idx = 0;
+ WRITE_ONCE(hctx->dispatch_from, hctx->ctxs[idx]);
+ }
+ }
+
+ if (rq)
+ atomic_dec(&ufq->rqs_count);
+ return rq;
+}
+
+/*
+ * Called by __blk_mq_alloc_request(). The shallow_depth value set by this
+ * function is used by __blk_mq_get_tag().
+ */
+static void ufq_limit_depth(blk_opf_t opf, struct blk_mq_alloc_data *data)
+{
+ struct ufq_data *ufq = data->q->elevator->elevator_data;
+
+ /* Do not throttle synchronous reads. */
+ if (op_is_sync(opf) && !op_is_write(opf))
+ return;
+
+ /*
+ * Throttle asynchronous requests and writes such that these requests
+ * do not block the allocation of synchronous requests.
+ */
+ data->shallow_depth = ufq->async_depth;
+}
+
+static void ufq_depth_updated(struct request_queue *q)
+{
+ struct ufq_data *ufq = q->elevator->elevator_data;
+
+ ufq->async_depth = q->nr_requests;
+ q->async_depth = q->nr_requests;
+ blk_mq_set_min_shallow_depth(q, 1);
+}
+
+static int ufq_init_sched(struct request_queue *q, struct elevator_queue *eq)
+{
+ const struct ufq_iosched_ops *ops;
+ struct ufq_data *ufq;
+
+ ufq = kzalloc_node(sizeof(*ufq), GFP_KERNEL, q->node);
+ if (!ufq)
+ return -ENOMEM;
+
+ eq->elevator_data = ufq;
+ ufq->q = q;
+ INIT_LIST_HEAD(&ufq->active_node);
+ INIT_LIST_HEAD(&ufq->insert_err_list);
+ spin_lock_init(&ufq->insert_err_lock);
+
+ blk_queue_flag_set(QUEUE_FLAG_SQ_SCHED, q);
+ q->elevator = eq;
+
+ q->async_depth = q->nr_requests;
+ ufq->async_depth = q->nr_requests;
+
+ ops = ufq_bpfops_tryget();
+ if (ops) {
+ if (ops->init_sched)
+ ops->init_sched(q);
+ ufq_bpfops_put();
+ }
+
+ mutex_lock(&ufq_active_queues_lock);
+ list_add_tail(&ufq->active_node, &ufq_active_queues);
+ mutex_unlock(&ufq_active_queues_lock);
+
+ ufq_depth_updated(q);
+ return 0;
+}
+
+static void ufq_exit_sched(struct elevator_queue *e)
+{
+ const struct ufq_iosched_ops *ops;
+ struct ufq_data *ufq = e->elevator_data;
+
+ ops = ufq_bpfops_tryget();
+ if (ops) {
+ if (ops->exit_sched)
+ ops->exit_sched(ufq->q);
+ ufq_bpfops_put();
+ }
+
+ mutex_lock(&ufq_active_queues_lock);
+ if (!list_empty(&ufq->active_node))
+ list_del_init(&ufq->active_node);
+ mutex_unlock(&ufq_active_queues_lock);
+
+ WARN_ON_ONCE(atomic_read(&ufq->rqs_count));
+ WARN_ON_ONCE(!list_empty(&ufq->insert_err_list));
+
+ kfree(ufq);
+ e->elevator_data = NULL;
+}
+
+void ufq_kick_all_hw_queues(void)
+{
+ struct ufq_data *ufq;
+
+ mutex_lock(&ufq_active_queues_lock);
+ list_for_each_entry(ufq, &ufq_active_queues, active_node)
+ blk_mq_run_hw_queues(ufq->q, true);
+ mutex_unlock(&ufq_active_queues_lock);
+}
+
+static int ufq_drain_ctx_rqs(struct ufq_data *ufq)
+{
+ struct request_queue *q = ufq->q;
+ unsigned long deadline = jiffies + 8 * HZ;
+
+ while (atomic_read(&ufq->rqs_count) > 0 && time_before(jiffies, deadline)) {
+ blk_mq_run_hw_queues(q, false);
+ if (atomic_read(&ufq->rqs_count) > 0) {
+ struct blk_mq_hw_ctx *hctx;
+ unsigned long i;
+
+ queue_for_each_hw_ctx(q, hctx, i)
+ flush_delayed_work(&hctx->run_work);
+ flush_delayed_work(&q->requeue_work);
+ }
+ cond_resched();
+ }
+
+ if (atomic_read(&ufq->rqs_count) > 0) {
+ pr_warn_ratelimited("ufq: drain timeout (%d rqs) before BPF attach\n",
+ atomic_read(&ufq->rqs_count));
+ return -EBUSY;
+ }
+ return 0;
+}
+
+/*
+ * Mirror elevator_change(): freeze each queue, cancel mq dispatch work,
+ * then drain software-ctx requests while BPF callbacks are still off.
+ * @enable runs with all those queues still frozen so new ctx backlog cannot
+ * race ahead of turning BPF dispatch on.
+ */
+int ufq_prepare_bpf_attach(int (*enable)(void *kdata), void *kdata)
+{
+ struct ufq_data *ufq;
+ unsigned int memflags;
+ int frozen = 0, ret = 0;
+
+ mutex_lock(&ufq_active_queues_lock);
+ if (list_empty(&ufq_active_queues)) {
+ mutex_unlock(&ufq_active_queues_lock);
+ return enable(kdata);
+ }
+
+ memflags = memalloc_noio_save();
+ list_for_each_entry(ufq, &ufq_active_queues, active_node) {
+ /* Count the freeze before waiting so a timeout unwinds it too. */
+ blk_freeze_queue_start(ufq->q);
+ frozen++;
+ if (!blk_mq_freeze_queue_wait_timeout(ufq->q, 8 * HZ)) {
+ pr_warn_ratelimited("ufq: freeze timeout (%d queued rqs) before BPF attach\n",
+ atomic_read(&ufq->rqs_count));
+ ret = -EBUSY;
+ goto unfreeze;
+ }
+ blk_mq_cancel_work_sync(ufq->q);
+ ret = ufq_drain_ctx_rqs(ufq);
+ if (ret)
+ goto unfreeze;
+ }
+
+ ret = enable(kdata);
+unfreeze:
+ list_for_each_entry(ufq, &ufq_active_queues, active_node) {
+ if (!frozen--)
+ break;
+ blk_mq_unfreeze_queue_nomemrestore(ufq->q);
+ }
+ memalloc_noio_restore(memflags);
+ mutex_unlock(&ufq_active_queues_lock);
+ return ret;
+}
+
+static bool ufq_bio_merge(struct request_queue *q, struct bio *bio,
+ unsigned int nr_segs)
+{
+ struct ufq_data *ufq = q->elevator->elevator_data;
+ const struct ufq_iosched_ops *ops;
+ struct request *rq = NULL, *last;
+ enum bio_merge_status mstat;
+ struct blk_mq_ctx *ctx;
+ bool ret = false;
+
+ /*
+ * Levels of merges:
+ * nomerges: No merges at all attempted
+ * noxmerges: Only simple one-hit cache try
+ * merges: All merge tries attempted
+ */
+ if (blk_queue_nomerges(q) || !bio_mergeable(bio))
+ return false;
+
+ last = q->last_merge;
+ if (last) {
+ ctx = last->mq_ctx;
+ spin_lock(&ctx->lock);
+ if (last == q->last_merge && !list_empty(&last->queuelist)
+ && elv_bio_merge_ok(last, bio)) {
+ mstat = blk_attempt_bio_merge(q, last, bio, nr_segs, true);
+ if (mstat == BIO_MERGE_OK) {
+ spin_unlock(&ctx->lock);
+ atomic_inc(&ufq->ops_stats.merge_bio_ok_count);
+ atomic64_add(bio->bi_iter.bi_size >> SECTOR_SHIFT,
+ &ufq->ops_stats.merge_bio_ok_sectors);
+ return true;
+ }
+ if (mstat == BIO_MERGE_FAILED) {
+ spin_unlock(&ctx->lock);
+ return false;
+ }
+ }
+ spin_unlock(&ctx->lock);
+ }
+
+ if (blk_queue_noxmerges(q))
+ return false;
+
+ ops = ufq_bpfops_tryget();
+ if (ops) {
+ if (ops->merge_bio) {
+ rq = ops->merge_bio(q, bio, nr_segs, &ret);
+ if (ret) {
+ atomic_inc(&ufq->ops_stats.merge_bio_ok_count);
+ atomic64_add(bio->bi_iter.bi_size >> SECTOR_SHIFT,
+ &ufq->ops_stats.merge_bio_ok_sectors);
+ } else {
+ ufq_bpfops_put();
+ return false;
+ }
+
+ if (rq) {
+ spin_lock(&rq->mq_ctx->lock);
+ if (!list_empty(&rq->queuelist)) {
+ list_del_init(&rq->queuelist);
+ atomic_dec(&ufq->rqs_count);
+ }
+ spin_unlock(&rq->mq_ctx->lock);
+ /* merge_bio transfers an extra BPF-owned reference. */
+ WARN_ON_ONCE(req_ref_put_and_test(rq));
+ blk_mq_free_request(rq);
+ ufq_bpfops_put();
+ atomic_inc(&ufq->ops_stats.merge_request_ok_count);
+ atomic64_add(bio->bi_iter.bi_size >> SECTOR_SHIFT,
+ &ufq->ops_stats.merge_request_ok_sectors);
+ return ret;
+ }
+ }
+ ufq_bpfops_put();
+ }
+
+ return ret;
+}
+
+static enum elv_merge ufq_try_insert_merge(struct request_queue *q,
+ struct request **new)
+{
+ const struct ufq_iosched_ops *ops;
+ struct request *target = NULL, *free = NULL, *last, *rq = *new;
+ struct ufq_data *ufq = q->elevator->elevator_data;
+ enum elv_merge type = ELEVATOR_NO_MERGE;
+ int merge_type = ELEVATOR_NO_MERGE;
+
+ if (!rq_mergeable(rq))
+ return ELEVATOR_NO_MERGE;
+
+ if (blk_queue_nomerges(q))
+ return ELEVATOR_NO_MERGE;
+
+ last = q->last_merge;
+ if (last) {
+ spin_lock(&last->mq_ctx->lock);
+ if (last == q->last_merge && !list_empty(&last->queuelist)
+ && bpf_attempt_merge(q, last, rq)) {
+ spin_unlock(&last->mq_ctx->lock);
+ type = ELEVATOR_BACK_MERGE;
+ free = rq;
+ *new = NULL;
+ goto end;
+ }
+ spin_unlock(&last->mq_ctx->lock);
+ }
+
+ if (blk_queue_noxmerges(q))
+ return ELEVATOR_NO_MERGE;
+
+ ops = ufq_bpfops_tryget();
+ if (ops && ops->merge_req) {
+ target = ops->merge_req(q, rq, &merge_type);
+ type = (enum elv_merge)merge_type;
+ }
+
+ if (target && WARN_ON_ONCE(req_ref_put_and_test(target))) {
+ __blk_mq_free_request(target);
+ ufq_bpfops_put();
+ return ELEVATOR_NO_MERGE;
+ }
+
+ if (type == ELEVATOR_NO_MERGE || !target) {
+ if (ops)
+ ufq_bpfops_put();
+ return ELEVATOR_NO_MERGE;
+ } else if (type == ELEVATOR_FRONT_MERGE) {
+ if (rq->mq_ctx != target->mq_ctx || rq->mq_hctx != target->mq_hctx)
+ goto rollback;
+ spin_lock(&target->mq_ctx->lock);
+ free = bpf_attempt_merge(q, rq, target);
+ if (!free) {
+ spin_unlock(&target->mq_ctx->lock);
+ goto rollback;
+ }
+ rq->elv.priv[0] = (void *)((uintptr_t)rq->elv.priv[0]
+ | UFQ_PRIV_IN_UFQ);
+ list_replace_init(&target->queuelist, &rq->queuelist);
+ rq->fifo_time = target->fifo_time;
+ q->last_merge = rq;
+ } else if (type == ELEVATOR_BACK_MERGE) {
+ spin_lock(&target->mq_ctx->lock);
+ free = bpf_attempt_merge(q, target, rq);
+ if (!free) {
+ spin_unlock(&target->mq_ctx->lock);
+ goto rollback;
+ }
+ *new = target;
+ q->last_merge = target;
+ }
+
+ spin_unlock(&target->mq_ctx->lock);
+ if (ops)
+ ufq_bpfops_put();
+end:
+ atomic_inc(&ufq->ops_stats.merge_request_ok_count);
+ atomic64_add(blk_rq_sectors(free), &ufq->ops_stats.merge_request_ok_sectors);
+ blk_mq_free_request(free);
+ return type;
+
+rollback:
+ if (ops) {
+ if (ops->insert_req && ops->insert_req(q, target, 0)) {
+ pr_err("ufq-iosched: rollback insert_req error\n");
+ if (ufq_move_to_insert_err(ufq, target))
+ atomic_inc(&ufq->ops_stats.insert_err_count);
+ }
+ ufq_bpfops_put();
+ }
+
+ return ELEVATOR_NO_MERGE;
+}
+
+static void ufq_insert_requests(struct blk_mq_hw_ctx *hctx,
+ struct list_head *list,
+ blk_insert_t flags)
+{
+ struct request_queue *q = hctx->queue;
+ struct ufq_data *ufq = q->elevator->elevator_data;
+ const struct ufq_iosched_ops *ops;
+ struct blk_mq_ctx *ctx;
+ enum elv_merge type;
+ int bit, ret = 0;
+
+ ops = ufq_bpfops_tryget();
+
+ while (!list_empty(list)) {
+ struct request *rq;
+
+ rq = list_first_entry(list, struct request, queuelist);
+ list_del_init(&rq->queuelist);
+
+ type = ufq_try_insert_merge(q, &rq);
+ /* Hold rq across publication and any concurrent dispatch. */
+ if (rq && WARN_ON_ONCE(!req_ref_inc_not_zero(rq)))
+ continue;
+ if (type == ELEVATOR_NO_MERGE) {
+ rq->fifo_time = jiffies;
+ ctx = rq->mq_ctx;
+ rq->elv.priv[0] = (void *)((uintptr_t)rq->elv.priv[0]
+ | UFQ_PRIV_IN_UFQ);
+ spin_lock(&ctx->lock);
+ if (flags & BLK_MQ_INSERT_AT_HEAD)
+ list_add(&rq->queuelist, &ctx->rq_lists[hctx->type]);
+ else
+ list_add_tail(&rq->queuelist,
+ &ctx->rq_lists[hctx->type]);
+
+ bit = ctx->index_hw[hctx->type];
+ if (!sbitmap_test_bit(&hctx->ctx_map, bit))
+ sbitmap_set_bit(&hctx->ctx_map, bit);
+ q->last_merge = rq;
+ spin_unlock(&ctx->lock);
+ atomic_inc(&ufq->rqs_count);
+ }
+
+ if (ops && rq && ops->insert_req) {
+ rq->elv.priv[0] = (void *)((uintptr_t)rq->elv.priv[0]
+ | UFQ_PRIV_IN_BPF);
+ ret = ops->insert_req(q, rq, flags);
+ if (ret) {
+ pr_err("ufq-iosched: bpf insert_req error (%d)\n", ret);
+ /*
+ * Leave BPF-owned peers on ctx; park this rq
+ * on insert_err_list for dedicated drain.
+ */
+ if (ufq_move_to_insert_err(ufq, rq))
+ atomic_inc(&ufq->ops_stats.insert_err_count);
+ } else {
+ atomic_inc(&ufq->ops_stats.insert_ok_count);
+ atomic64_add(blk_rq_sectors(rq), &ufq->ops_stats.insert_ok_sectors);
+ }
+ }
+ if (rq && req_ref_put_and_test(rq))
+ __blk_mq_free_request(rq);
+ }
+
+ if (ops)
+ ufq_bpfops_put();
+}
+
+static void ufq_prepare_request(struct request *rq)
+{
+ rq->elv.priv[0] = (void *)(uintptr_t)UFQ_PRIV_NOT_IN_SCHED;
+}
+
+static void ufq_finish_request(struct request *rq)
+{
+ const struct ufq_iosched_ops *ops;
+ struct ufq_data *ufq = rq->q->elevator->elevator_data;
+
+ /*
+ * The block layer core may call ufq_finish_request() without having
+ * called ufq_insert_requests(). Skip requests that bypassed I/O
+ * scheduling.
+ */
+ if (!((uintptr_t)rq->elv.priv[0] & UFQ_PRIV_IN_BPF))
+ return;
+
+ ops = ufq_bpfops_tryget();
+ if (ops) {
+ if (ops->finish_req)
+ ops->finish_req(rq);
+ ufq_bpfops_put();
+ }
+
+ atomic_inc(&ufq->ops_stats.finish_ok_count);
+ atomic64_add(blk_rq_stats_sectors(rq), &ufq->ops_stats.finish_ok_sectors);
+}
+
+static bool ufq_has_work(struct blk_mq_hw_ctx *hctx)
+{
+ const struct ufq_iosched_ops *ops;
+ struct ufq_data *ufq = hctx->queue->elevator->elevator_data;
+ int rqs_count = atomic_read(&ufq->rqs_count);
+ int policy_has_work;
+
+ ops = ufq_bpfops_tryget();
+ if (!ops)
+ return rqs_count > 0;
+
+ policy_has_work = ops->has_req ?
+ ops->has_req(hctx->queue, rqs_count) : rqs_count;
+ ufq_bpfops_put();
+ if (!policy_has_work && rqs_count > 0)
+ atomic_inc(&ufq->ops_stats.has_work_false_pending_count);
+
+ return policy_has_work || ufq_insert_err_pending(ufq);
+}
+
+#ifdef CONFIG_BLK_DEBUG_FS
+static int ufq_ops_stats_show(void *data, struct seq_file *m)
+{
+ struct request_queue *q = data;
+ struct ufq_data *ufq = q->elevator->elevator_data;
+ struct ufq_ops_stats *s = &ufq->ops_stats;
+
+ /* for debug */
+ seq_printf(m, "dispatch_ok_count %d\n",
+ atomic_read(&s->dispatch_ok_count));
+ seq_printf(m, "dispatch_ok_sectors %lld\n",
+ (long long)atomic64_read(&s->dispatch_ok_sectors));
+ seq_printf(m, "dispatch_null_count %d\n",
+ atomic_read(&s->dispatch_null_count));
+ seq_printf(m, "dispatch_bad_state_count %d\n",
+ atomic_read(&s->dispatch_bad_state_count));
+ seq_printf(m, "has_work_false_pending_count %d\n",
+ atomic_read(&s->has_work_false_pending_count));
+ seq_printf(m, "insert_ok_count %d\n",
+ atomic_read(&s->insert_ok_count));
+ seq_printf(m, "insert_ok_sectors %lld\n",
+ (long long)atomic64_read(&s->insert_ok_sectors));
+ seq_printf(m, "insert_err_count %d\n",
+ atomic_read(&s->insert_err_count));
+ seq_printf(m, "insert_err_handled %d\n",
+ atomic_read(&ufq->insert_err_handled));
+ seq_printf(m, "rqs_count %d\n",
+ atomic_read(&ufq->rqs_count));
+ seq_printf(m, "merge_req_ok_count %d\n",
+ atomic_read(&s->merge_request_ok_count));
+ seq_printf(m, "merge_req_ok_sectors %lld\n",
+ (long long)atomic64_read(&s->merge_request_ok_sectors));
+ seq_printf(m, "merge_bio_ok_count %d\n",
+ atomic_read(&s->merge_bio_ok_count));
+ seq_printf(m, "merge_bio_ok_sectors %lld\n",
+ (long long)atomic64_read(&s->merge_bio_ok_sectors));
+ seq_printf(m, "finish_ok_count %d\n",
+ atomic_read(&s->finish_ok_count));
+ seq_printf(m, "finish_ok_sectors %lld\n",
+ (long long)atomic64_read(&s->finish_ok_sectors));
+ return 0;
+}
+
+static const struct blk_mq_debugfs_attr ufq_iosched_debugfs_attrs[] = {
+ {"ops_stats", 0400, ufq_ops_stats_show},
+ {},
+};
+#endif
+
+static struct elevator_type ufq_iosched_mq = {
+ .ops = {
+ .depth_updated = ufq_depth_updated,
+ .limit_depth = ufq_limit_depth,
+ .insert_requests = ufq_insert_requests,
+ .dispatch_request = ufq_dispatch_request,
+ .prepare_request = ufq_prepare_request,
+ .finish_request = ufq_finish_request,
+ .bio_merge = ufq_bio_merge,
+ .has_work = ufq_has_work,
+ .init_sched = ufq_init_sched,
+ .exit_sched = ufq_exit_sched,
+ },
+
+#ifdef CONFIG_BLK_DEBUG_FS
+ .queue_debugfs_attrs = ufq_iosched_debugfs_attrs,
+#endif
+ .elevator_name = "ufq",
+ .elevator_alias = "ufq_iosched",
+ .elevator_owner = THIS_MODULE,
+};
+bool ufq_is_queue(struct request_queue *q)
+{
+ struct elevator_queue *e = READ_ONCE(q->elevator);
+
+ /* Caller must hold a queue usage reference. */
+ return e && e->type == &ufq_iosched_mq;
+}
+
+static int __init ufq_init(void)
+{
+ int ret;
+
+ ret = elv_register(&ufq_iosched_mq);
+ if (ret)
+ return ret;
+
+ ret = bpf_ufq_kfunc_init();
+ if (ret) {
+ pr_err("ufq-iosched: Failed to register kfunc sets (%d)\n", ret);
+ elv_unregister(&ufq_iosched_mq);
+ return ret;
+ }
+
+ ret = bpf_ufq_ops_init();
+ if (ret) {
+ pr_err("ufq-iosched: Failed to register struct_ops (%d)\n", ret);
+ elv_unregister(&ufq_iosched_mq);
+ return ret;
+ }
+
+ return 0;
+}
+
+static void __exit ufq_exit(void)
+{
+ elv_unregister(&ufq_iosched_mq);
+}
+
+module_init(ufq_init);
+module_exit(ufq_exit);
+
+MODULE_LICENSE("GPL");
+MODULE_ALIAS("ufq-iosched");
+MODULE_AUTHOR("Kaitao Cheng <chengkaitao@kylinos.cn>");
+MODULE_DESCRIPTION("User-programmable Flexible Queueing");
diff --git a/block/ufq-iosched.h b/block/ufq-iosched.h
new file mode 100644
index 000000000000..177c192dab86
--- /dev/null
+++ b/block/ufq-iosched.h
@@ -0,0 +1,79 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#ifndef _BLOCK_UFQ_IOSCHED_H
+#define _BLOCK_UFQ_IOSCHED_H
+
+#include <linux/types.h>
+#include "elevator.h"
+#include "blk-mq.h"
+
+#ifndef BPF_IOSCHED_NAME_MAX
+#define BPF_IOSCHED_NAME_MAX 16
+#endif
+
+/* For testing and debugging */
+struct ufq_ops_stats {
+ atomic_t dispatch_ok_count;
+ atomic64_t dispatch_ok_sectors;
+ atomic_t dispatch_null_count;
+ atomic_t dispatch_bad_state_count;
+ atomic_t has_work_false_pending_count;
+ atomic_t insert_ok_count;
+ atomic64_t insert_ok_sectors;
+ atomic_t insert_err_count;
+ atomic_t merge_request_ok_count;
+ atomic64_t merge_request_ok_sectors;
+ atomic_t merge_bio_ok_count;
+ atomic64_t merge_bio_ok_sectors;
+ atomic_t finish_ok_count;
+ atomic64_t finish_ok_sectors;
+};
+
+/*
+ * merge_req, merge_bio and dispatch_req return one BPF-owned request
+ * reference. UFQ consumes it separately from the block layer's reference.
+ */
+struct ufq_iosched_ops {
+ int (*init_sched)(struct request_queue *q);
+ int (*exit_sched)(struct request_queue *q);
+ bool (*has_req)(struct request_queue *q, int rqs_count);
+ int (*insert_req)(struct request_queue *q, struct request *rq,
+ blk_insert_t flags);
+ void (*finish_req)(struct request *rq);
+ struct request *(*merge_req)(struct request_queue *q, struct request *rq,
+ int *type);
+ struct request *(*merge_bio)(struct request_queue *q, struct bio *bio,
+ unsigned int nr_segs, bool *merged);
+ struct request *(*dispatch_req)(struct request_queue *q);
+ char name[BPF_IOSCHED_NAME_MAX];
+};
+
+struct ufq_data {
+ struct request_queue *q;
+ u32 async_depth;
+ atomic_t rqs_count;
+ /*
+ * BPF insert_req failures: rq is moved off ctx onto insert_err_list
+ * (insert_err_lock) so err drain never steals BPF-owned ctx rqs.
+ * insert_err_handled counts successful err-list dispatches (debugfs).
+ */
+ spinlock_t insert_err_lock;
+ struct list_head insert_err_list;
+ atomic_t insert_err_handled;
+ struct list_head active_node;
+ struct ufq_ops_stats ops_stats;
+};
+
+const struct ufq_iosched_ops *ufq_bpfops_tryget(void);
+void ufq_bpfops_put(void);
+void ufq_kick_all_hw_queues(void);
+bool ufq_bpfops_draining(void);
+bool ufq_is_queue(struct request_queue *q);
+int ufq_prepare_bpf_attach(int (*enable)(void *kdata), void *kdata);
+int bpf_ufq_ops_init(void);
+int bpf_ufq_kfunc_init(void);
+
+#endif /* _BLOCK_UFQ_IOSCHED_H */
diff --git a/block/ufq-kfunc.c b/block/ufq-kfunc.c
new file mode 100644
index 000000000000..843aeb0e67cf
--- /dev/null
+++ b/block/ufq-kfunc.c
@@ -0,0 +1,169 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ */
+#include <linux/init.h>
+#include <linux/types.h>
+#include <linux/bpf_verifier.h>
+#include <linux/bpf.h>
+#include <linux/btf.h>
+#include <linux/btf_ids.h>
+#include <trace/events/block.h>
+#include "blk.h"
+#include "ufq-iosched.h"
+
+__bpf_kfunc_start_defs();
+
+__bpf_kfunc struct request *bpf_request_acquire(struct request *rq)
+{
+ if (req_ref_inc_not_zero(rq))
+ return rq;
+ return NULL;
+}
+
+__bpf_kfunc void bpf_request_release(struct request *rq)
+{
+ if (req_ref_put_and_test(rq))
+ __blk_mq_free_request(rq);
+}
+
+__bpf_kfunc bool bpf_request_bio_try_merge(struct request *rq, struct bio *bio,
+ unsigned int nr_segs)
+{
+ struct request_queue *q;
+ struct blk_mq_ctx *ctx;
+ bool merged = false;
+
+ if (!rq || !bio || !rq->q)
+ return false;
+
+ q = rq->q;
+ if (blk_queue_enter(q, BLK_MQ_REQ_NOWAIT))
+ return false;
+ ctx = rq->mq_ctx;
+ if (!ufq_is_queue(q) || !ctx || !bio->bi_bdev ||
+ !bio->bi_bdev->bd_disk || bio->bi_bdev->bd_disk->queue != q)
+ goto out;
+
+ /*
+ * KF_SPINLOCK_SAFE callers (BPF merge_bio) often already hold a scheduler
+ * lock taken with irqsave (e.g. PFQ dd->lock). Blocking on ctx->lock
+ * nesting that order deadlocks / hard-locks under load. Trylock: skip
+ * the merge if ctx is busy; bio will be issued as a new request.
+ */
+ if (!spin_trylock(&ctx->lock))
+ goto out;
+
+ merged = blk_attempt_bio_merge(q, rq, bio, nr_segs, true) == BIO_MERGE_OK;
+ spin_unlock(&ctx->lock);
+out:
+ blk_queue_exit(q);
+ return merged;
+}
+
+__bpf_kfunc struct request *bpf_request_try_merge(struct request *rq, struct request *next)
+{
+ struct request_queue *q;
+ struct blk_mq_ctx *ctx;
+ struct ufq_data *ufq;
+ struct request *free = NULL;
+
+ if (!rq || !next || !rq->q || rq->q != next->q)
+ return NULL;
+
+ q = rq->q;
+ if (blk_queue_enter(q, BLK_MQ_REQ_NOWAIT))
+ return NULL;
+ if (!ufq_is_queue(q))
+ goto out;
+
+ ufq = q->elevator->elevator_data;
+ if (rq->mq_ctx != next->mq_ctx || rq->mq_hctx != next->mq_hctx)
+ goto out;
+
+ ctx = rq->mq_ctx;
+ if (!ctx)
+ goto out;
+
+ /* Same nesting rule as bpf_request_bio_try_merge (see comment there). */
+ if (!spin_trylock(&ctx->lock))
+ goto out;
+
+ free = bpf_attempt_merge(q, rq, next);
+ if (free) {
+ if (q->last_merge == free)
+ q->last_merge = NULL;
+ list_del_init(&free->queuelist);
+ atomic_dec(&ufq->rqs_count);
+ }
+ spin_unlock(&ctx->lock);
+out:
+ blk_queue_exit(q);
+ return free;
+}
+
+/*
+ * Elevator private slot write for UFQ BPF.
+ *
+ * BPF cannot STX into request->elv.priv[] (verifier: "only read is supported"
+ * without btf_struct_access). Reads are fine from BPF; only the store needs
+ * this KF_SPINLOCK_SAFE kfunc (same pattern as bpf_request_try_merge).
+ *
+ * @val is a raw address / tagged qid scalar from BPF (not an owning ref).
+ */
+__bpf_kfunc void bpf_request_set_elv_priv1(struct request *rq, u64 val)
+{
+ if (!rq)
+ return;
+ rq->elv.priv[1] = (void *)(uintptr_t)val;
+}
+
+__bpf_kfunc_end_defs();
+
+#if defined(CONFIG_X86_KERNEL_IBT)
+static const void * const __used __section(".discard.ibt_endbr_noseal")
+__ibt_noseal_bpf_request_release = (void *)bpf_request_release;
+#endif
+
+BTF_KFUNCS_START(ufq_kfunc_set_ops)
+BTF_ID_FLAGS(func, bpf_request_acquire, KF_ACQUIRE | KF_RET_NULL)
+BTF_ID_FLAGS(func, bpf_request_release, KF_RELEASE)
+BTF_ID_FLAGS(func, bpf_request_bio_try_merge, KF_SPINLOCK_SAFE)
+BTF_ID_FLAGS(func, bpf_request_try_merge, KF_SPINLOCK_SAFE)
+BTF_ID_FLAGS(func, bpf_request_set_elv_priv1, KF_SPINLOCK_SAFE)
+BTF_KFUNCS_END(ufq_kfunc_set_ops)
+
+static const struct btf_kfunc_id_set bpf_ufq_kfunc_set = {
+ .owner = THIS_MODULE,
+ .set = &ufq_kfunc_set_ops,
+};
+
+BTF_ID_LIST(bpf_ufq_dtor_kfunc_ids)
+BTF_ID(struct, request)
+BTF_ID(func, bpf_request_release)
+
+int bpf_ufq_kfunc_init(void)
+{
+ int ret;
+ const struct btf_id_dtor_kfunc bpf_ufq_dtor_kfunc[] = {
+ {
+ .btf_id = bpf_ufq_dtor_kfunc_ids[0],
+ .kfunc_btf_id = bpf_ufq_dtor_kfunc_ids[1]
+ },
+ };
+
+ ret = register_btf_kfunc_id_set(BPF_PROG_TYPE_STRUCT_OPS, &bpf_ufq_kfunc_set);
+ if (ret)
+ return ret;
+ ret = register_btf_kfunc_id_set(BPF_PROG_TYPE_SYSCALL, &bpf_ufq_kfunc_set);
+ if (ret)
+ return ret;
+ ret = register_btf_id_dtor_kfuncs(bpf_ufq_dtor_kfunc,
+ ARRAY_SIZE(bpf_ufq_dtor_kfunc),
+ THIS_MODULE);
+ if (ret)
+ return ret;
+
+ return 0;
+}
--
2.53.0
^ permalink raw reply [flat|nested] 6+ messages in thread* [RFC v3 3/3] tools/ufq_iosched: add PFQ eBPF I/O scheduler
2026-10-03 4:27 [RFC v3 0/3] block: Introduce a BPF-based I/O scheduler Kaitao Cheng
2026-10-03 4:27 ` [RFC v3 2/3] tools/ufq_iosched: add BPF example scheduler and build scaffolding Kaitao Cheng
2026-10-03 4:27 ` [RFC v3 1/3] block: Introduce the UFQ I/O scheduler Kaitao Cheng
@ 2026-10-03 4:27 ` Kaitao Cheng
2026-10-03 9:36 ` [RFC v3 0/3] block: Introduce a BPF-based " Alexei Starovoitov
3 siblings, 0 replies; 6+ messages in thread
From: Kaitao Cheng @ 2026-10-03 4:27 UTC (permalink / raw)
To: Jens Axboe; +Cc: linux-block, linux-kernel, bpf, Kaitao Cheng, Li Youhong
From: Kaitao Cheng <chengkaitao@kylinos.cn>
Add a Priority Fair Queue (PFQ) eBPF I/O scheduler to the UFQ framework,
providing weighted fair scheduling and preferential service for
configured interactive threads.
Organize requests into some logical queues per disk by priority class,
priority level and I/O type. Select queues by virtual time and charge
service according to request size and queue weight, with bounded dispatch
batches. Within each queue, serve expired FIFO requests first, otherwise
prefer metadata and shorter seek distances. Dispatch requests inserted
with BLK_MQ_INSERT_AT_HEAD before the fair queues.
Support adjacent request and bio merging through sector-ordered indexes.
Give interactive read queues additional weight and allow empty queues to
briefly retain their service slot, with idle expiry checked by subsequent
scheduler callbacks.
Add a userspace loader to initialize scheduling tunables, configure
interactive threads by exact task comm using repeated -i options, and
report throughput, IOPS, queue state and scheduler event counters.
This is an intermediate version under active development and continued
iteration. It has not yet been used in production.
Signed-off-by: Kaitao Cheng <chengkaitao@kylinos.cn>
Signed-off-by: Li Youhong <liyouhong@kylinos.cn>
---
tools/ufq_iosched/Makefile | 2 +-
tools/ufq_iosched/README.md | 55 +-
tools/ufq_iosched/include/ufq/pfq.bpf.h | 125 +
tools/ufq_iosched/include/ufq/pfq_disk.h | 38 +
tools/ufq_iosched/include/ufq/pfq_stat.h | 44 +
tools/ufq_iosched/include/ufq/pfq_tunable.h | 23 +
tools/ufq_iosched/pfq.bpf.c | 2736 +++++++++++++++++++
tools/ufq_iosched/pfq.c | 329 +++
8 files changed, 3350 insertions(+), 2 deletions(-)
create mode 100644 tools/ufq_iosched/include/ufq/pfq.bpf.h
create mode 100644 tools/ufq_iosched/include/ufq/pfq_disk.h
create mode 100644 tools/ufq_iosched/include/ufq/pfq_stat.h
create mode 100644 tools/ufq_iosched/include/ufq/pfq_tunable.h
create mode 100644 tools/ufq_iosched/pfq.bpf.c
create mode 100644 tools/ufq_iosched/pfq.c
diff --git a/tools/ufq_iosched/Makefile b/tools/ufq_iosched/Makefile
index 7dc37d9172aa..f8723c139b65 100644
--- a/tools/ufq_iosched/Makefile
+++ b/tools/ufq_iosched/Makefile
@@ -192,7 +192,7 @@ $(INCLUDE_DIR)/%.bpf.skel.h: $(UFQOBJ_DIR)/%.bpf.o $(INCLUDE_DIR)/vmlinux.h $(BP
UFQ_COMMON_DEPS := include/ufq/common.h include/ufq/simple_stat.h | $(BINDIR)
-c-sched-targets = ufq_simple
+c-sched-targets = ufq_simple pfq
$(addprefix $(BINDIR)/,$(c-sched-targets)): \
$(BINDIR)/%: \
diff --git a/tools/ufq_iosched/README.md b/tools/ufq_iosched/README.md
index 1fa305fb576e..58f7f24a7f88 100644
--- a/tools/ufq_iosched/README.md
+++ b/tools/ufq_iosched/README.md
@@ -122,7 +122,60 @@ A simple IO scheduler that provides an example of a minimal ufq scheduler.
Populates commonly used kernel-exposed BPF interfaces for testing the UFQ
scheduler framework in the kernel.
-### bpftool feature detection (`llvm`, `libcap`, `libbfd`)
+## pfq
+
+Priority Fair Queue (PFQ) is an eBPF I/O scheduler built on the UFQ framework.
+It provides weighted fair scheduling, request merging, and preferential
+service for configured interactive threads.
+
+UFQ handles block-layer integration and request lifecycle management, calling
+PFQ's initialization, insertion, dispatch, and completion callbacks through
+`struct_ops`. PFQ implements the scheduling policy; the userspace loader
+loads the BPF program, configures it, and attaches the callbacks.
+
+### Scheduling policy
+
+Each disk has 96 logical queues, grouped by priority class, priority level,
+and I/O type (read, write, or synchronous write). Requests with the same
+classification share a queue. Scheduling has two levels:
+
+- **Between queues:** a service tree orders queues by virtual time. When
+ selecting a new queue, PFQ chooses the one with the lowest virtual time.
+ Each dispatch increases that queue's virtual time in proportion to the
+ request's sector count divided by the queue's weight. A higher weight
+ therefore gives more service. A queue can dispatch a bounded batch of
+ requests before participating in selection again.
+- **Within a queue:** expired FIFO requests take precedence. Otherwise, PFQ
+ selects candidates by sector position, preferring metadata and shorter seek
+ distances. The sector index also supports adjacent-request merging.
+
+Requests marked `BLK_MQ_INSERT_AT_HEAD` enter a separate list and are
+dispatched before requests in the fair queues.
+
+Configured interactive threads use the SPECIAL class, whose read requests
+receive additional weight. An empty SPECIAL read queue may briefly retain
+its service slot to wait for more requests. Later scheduler callbacks check
+when this idle hold expires.
+
+### Building and running
+
+Use a kernel built with `CONFIG_IOSCHED_UFQ=y`. From this directory, run:
+
+```bash
+$ make pfq
+$ sudo ./build/bin/pfq
+```
+
+Repeat `-i NAME` to configure interactive threads:
+
+```bash
+$ sudo ./build/bin/pfq -i firefox -i code
+```
+
+Names are matched exactly and case-sensitively against the current thread's
+`comm`. Without `-i`, requests use normal I/O priority classification.
+
+## bpftool feature detection (`llvm`, `libcap`, `libbfd`)
While building bpftool, you may see lines similar to:
diff --git a/tools/ufq_iosched/include/ufq/pfq.bpf.h b/tools/ufq_iosched/include/ufq/pfq.bpf.h
new file mode 100644
index 000000000000..34084b160930
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/pfq.bpf.h
@@ -0,0 +1,125 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ * Copyright (c) 2026 Li Youhong <liyouhong@kylinos.cn>
+ *
+ * PFQ eBPF scheduler constants shared by pfq.bpf.c and pfq.c.
+ */
+#ifndef __PFQ_BPF_H
+#define __PFQ_BPF_H
+
+/* Task comm includes the terminating NUL, matching TASK_COMM_LEN. */
+#define PFQ_COMM_LEN 16
+#define PFQ_INTERACTIVE_MAX 64
+
+enum pfq_init_state {
+ PFQ_UNINITIALIZED,
+ PFQ_INITIALIZING,
+ PFQ_INITIALIZED,
+};
+
+struct pfq_comm_key {
+ char name[PFQ_COMM_LEN];
+};
+
+#define PFQ_PRIO_LEVELS 8
+#define PFQ_IOCLASS_NUM 4
+#define PFQ_QUEUE_TYPES 3
+#define PFQ_TOTAL_QUEUES \
+ (PFQ_IOCLASS_NUM * PFQ_PRIO_LEVELS * PFQ_QUEUE_TYPES)
+
+/*
+ * Share request and FIFO roots by I/O type to stay within BTF_FIELDS_MAX.
+ * Each tree groups logical queues by qid, then orders requests by sector
+ * or FIFO deadline.
+ */
+#define PFQ_RQ_TREES PFQ_QUEUE_TYPES
+
+#define PFQ_QUEUE_MAP_MAX \
+ (PFQ_DISK_MAP_MAX * PFQ_TOTAL_QUEUES)
+
+#define PFQ_INDEX(class, prio, type) \
+ (((class) * PFQ_PRIO_LEVELS * PFQ_QUEUE_TYPES) + \
+ ((prio) * PFQ_QUEUE_TYPES) + (type))
+
+#define PFQ_GET_CLASS(index) \
+ ((index) / (PFQ_PRIO_LEVELS * PFQ_QUEUE_TYPES))
+#define PFQ_GET_PRIO(index) \
+ (((index) % (PFQ_PRIO_LEVELS * PFQ_QUEUE_TYPES)) / PFQ_QUEUE_TYPES)
+#define PFQ_GET_TYPE(index) \
+ ((index) % PFQ_QUEUE_TYPES)
+
+#define PFQ_QID_TREE_INDEX(qid) PFQ_GET_TYPE(qid)
+
+#define PFQ_DEFAULT_QUEUE_INDEX(type) \
+ PFQ_INDEX(2 /* PFQ_BE_CLASS */, 4 /* IOPRIO_BE_NORM */, (type))
+
+#define PFQ_WEIGHT_MULT_SPECIAL 400
+#define PFQ_PRIO_LEVEL_COEFF_SP_READ 200
+#define PFQ_PRIO_LEVEL_COEFF_READ 100
+#define PFQ_PRIO_LEVEL_COEFF_WRITE 5
+
+#define PFQ_VTIME_SHIFT 10
+
+#define PFQ_WEIGHT_BASE_DEFAULT 100
+#define PFQ_BATCH_READ_DEFAULT 8
+#define PFQ_BATCH_SYNC_WRITE_DEFAULT 4
+#define PFQ_BATCH_WRITE_DEFAULT 2
+#define PFQ_IDLE_DELAY_MIN_MS_DEFAULT 8
+#define PFQ_IDLE_DELAY_MAX_MS_DEFAULT 20
+
+#define PFQ_READ_EXPIRE_NS (500000000ULL)
+#define PFQ_SYNC_WRITE_EXPIRE_NS (2000000000ULL)
+#define PFQ_WRITE_EXPIRE_NS (5000000000ULL)
+
+#define PFQ_LOOP_MAX 128
+#define PFQ_RB_HEIGHT_MAX 24 /* verifier search bound */
+#define PFQ_DISK_MAP_MAX 32
+#define PFQ_QID_NONE 0xffffffffU
+
+#define PFQ_IOPRIO_CLASS_SHIFT 13
+
+#define PFQ_IOPRIO_CLASS_NONE 0
+#define PFQ_IOPRIO_CLASS_RT 1
+#define PFQ_IOPRIO_CLASS_BE 2
+#define PFQ_IOPRIO_CLASS_IDLE 3
+#define PFQ_IOPRIO_BE_NORM 4
+
+#define PFQ_SPECIAL_CLASS 0
+#define PFQ_RT_CLASS 1
+#define PFQ_BE_CLASS 2
+#define PFQ_IDLE_CLASS 3
+
+#define PFQ_READ 0
+#define PFQ_WRITE 1
+#define PFQ_SYNC_WRITE 2
+
+#define BLK_MQ_INSERT_AT_HEAD 0x01
+#define REQ_OP_MASK ((1 << 8) - 1)
+#define SECTOR_SHIFT 9
+
+/* BPF-only masks: derive bit positions from the generated kernel enums. */
+#ifdef __VMLINUX_H__
+#define PFQ_REQ_SYNC (1ULL << __REQ_SYNC)
+#define PFQ_REQ_META (1ULL << __REQ_META)
+#define PFQ_REQ_OP_FLUSH REQ_OP_FLUSH
+#define PFQ_REQ_OP_ZONE_APPEND REQ_OP_ZONE_APPEND
+#define PFQ_REQ_OP_WRITE_ZEROES REQ_OP_WRITE_ZEROES
+#define PFQ_REQ_OP_DRV_IN REQ_OP_DRV_IN
+#define PFQ_REQ_OP_DRV_OUT REQ_OP_DRV_OUT
+
+#define PFQ_REQ_NOMERGE (1ULL << __REQ_NOMERGE)
+#define PFQ_REQ_FUA (1ULL << __REQ_FUA)
+#define PFQ_REQ_PREFLUSH (1ULL << __REQ_PREFLUSH)
+#define PFQ_REQ_NOMERGE_FLAGS \
+ (PFQ_REQ_NOMERGE | PFQ_REQ_PREFLUSH | PFQ_REQ_FUA)
+
+#define PFQ_RQF_STARTED (1U << __RQF_STARTED)
+#define PFQ_RQF_FLUSH_SEQ (1U << __RQF_FLUSH_SEQ)
+#define PFQ_RQF_SPECIAL_PAYLOAD (1U << __RQF_SPECIAL_PAYLOAD)
+#define PFQ_RQF_NOMERGE_FLAGS \
+ (PFQ_RQF_STARTED | PFQ_RQF_FLUSH_SEQ | PFQ_RQF_SPECIAL_PAYLOAD)
+#endif /* __VMLINUX_H__ */
+
+#endif /* __PFQ_BPF_H */
diff --git a/tools/ufq_iosched/include/ufq/pfq_disk.h b/tools/ufq_iosched/include/ufq/pfq_disk.h
new file mode 100644
index 000000000000..d74ead1d3e5e
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/pfq_disk.h
@@ -0,0 +1,38 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ * Copyright (c) 2026 Li Youhong <liyouhong@kylinos.cn>
+ *
+ * Userspace prefix of struct pfq_disk_data through nr_queued_total.
+ * Keep field types and order in sync with pfq.bpf.c.
+ *
+ * Lookup must use a buffer of bpf_map__value_size(); this prefix is safe
+ * to interpret when the map value is at least sizeof(*this).
+ */
+#ifndef __PFQ_DISK_H
+#define __PFQ_DISK_H
+
+#include <linux/bpf.h>
+#include <stdint.h>
+#include <ufq/pfq.bpf.h>
+
+typedef int32_t s32;
+typedef uint32_t u32;
+typedef uint64_t u64;
+
+struct pfq_disk_data_user {
+ struct bpf_spin_lock lock;
+ struct bpf_list_head dispatch;
+ struct bpf_rb_root service_tree;
+ struct bpf_rb_root rq_trees[PFQ_RQ_TREES];
+ struct bpf_rb_root fifo_trees[PFQ_RQ_TREES];
+ s32 disk_id;
+ u64 vtime;
+ u64 min_vtime;
+ u64 last_sector;
+ u32 nr_active;
+ u32 nr_queued_total;
+};
+
+#endif /* __PFQ_DISK_H */
diff --git a/tools/ufq_iosched/include/ufq/pfq_stat.h b/tools/ufq_iosched/include/ufq/pfq_stat.h
new file mode 100644
index 000000000000..3a920c659292
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/pfq_stat.h
@@ -0,0 +1,44 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ * Copyright (c) 2026 Li Youhong <liyouhong@kylinos.cn>
+ *
+ * Per-disk, per-CPU PFQ event counters. Disk exit removes the counters;
+ * initialization reuses an existing entry without resetting it.
+ */
+#ifndef __PFQ_STAT_H
+#define __PFQ_STAT_H
+
+enum pfq_stat_idx {
+ PFQ_STAT_INSERT_CNT = 0,
+ PFQ_STAT_INSERT_SIZE,
+ PFQ_STAT_INSERT_ERR,
+ PFQ_STAT_AT_HEAD_CNT,
+ PFQ_STAT_AT_HEAD_SIZE,
+ PFQ_STAT_INTERACTIVE,
+ PFQ_STAT_RQMERGE_CNT,
+ PFQ_STAT_RQMERGE_SIZE,
+ PFQ_STAT_BIOMERGE_CNT,
+ PFQ_STAT_BIOMERGE_SIZE,
+ PFQ_STAT_DISPATCH_CNT,
+ PFQ_STAT_DISPATCH_SIZE,
+ PFQ_STAT_DISPATCH_AT_HEAD_CNT,
+ PFQ_STAT_DISPATCH_AT_HEAD_SIZE,
+ PFQ_STAT_FINISH_CNT,
+ PFQ_STAT_FINISH_SIZE,
+ PFQ_STAT_IDLE_HOLD,
+ PFQ_STAT_MAX,
+};
+
+/*
+ * INSERT counters count successful normal inserts before merge deductions.
+ * RQMERGE counters count queued candidates removed by merge_req, before the
+ * kernel attempts the merge. A rejected candidate is counted again on insert.
+ * For each disk, sum all CPUs before subtracting RQMERGE counters.
+ */
+struct pfq_stats {
+ __u64 counters[PFQ_STAT_MAX];
+};
+
+#endif /* __PFQ_STAT_H */
diff --git a/tools/ufq_iosched/include/ufq/pfq_tunable.h b/tools/ufq_iosched/include/ufq/pfq_tunable.h
new file mode 100644
index 000000000000..edd9ea245282
--- /dev/null
+++ b/tools/ufq_iosched/include/ufq/pfq_tunable.h
@@ -0,0 +1,23 @@
+/* SPDX-License-Identifier: GPL-2.0 */
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ * Copyright (c) 2026 Li Youhong <liyouhong@kylinos.cn>
+ */
+#ifndef __PFQ_TUNABLE_H
+#define __PFQ_TUNABLE_H
+
+#include <ufq/pfq.bpf.h>
+
+#define PFQ_TUNABLE_KEY 0
+
+struct pfq_tunables {
+ __u32 weight_base;
+ __u32 idle_delay_min_ms;
+ __u32 idle_delay_max_ms;
+ __u32 batch_read;
+ __u32 batch_sync_write;
+ __u32 batch_write;
+};
+
+#endif /* __PFQ_TUNABLE_H */
diff --git a/tools/ufq_iosched/pfq.bpf.c b/tools/ufq_iosched/pfq.bpf.c
new file mode 100644
index 000000000000..c1aa8650138e
--- /dev/null
+++ b/tools/ufq_iosched/pfq.bpf.c
@@ -0,0 +1,2736 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ * Copyright (c) 2026 Li Youhong <liyouhong@kylinos.cn>
+ *
+ * PFQ (Priority Fair Queue) eBPF scheduler for the UFQ iosched framework.
+ *
+ * Each disk has 96 logical queues sharing request and FIFO indexes by I/O
+ * type. The service tree orders queues by virtual time; dispatch serves
+ * bounded batches with FIFO expiry taking precedence over seek distance.
+ * Head inserts bypass fair queueing.
+ *
+ * dd->lock protects graph roots and mutable scheduling state. Map lookups,
+ * allocation and request acquisition run outside the lock. Detached and
+ * temporary owning references are released after unlocking.
+ *
+ * Empty SPECIAL+READ queues can retain their service slot until a later
+ * callback checks the idle deadline. UFQ provides no timer for this hold.
+ */
+#include <ufq/common.bpf.h>
+#include <ufq/pfq.bpf.h>
+#include <ufq/pfq_stat.h>
+#include <ufq/pfq_tunable.h>
+
+char _license[] SEC("license") = "GPL";
+
+/* Unlocked qid reads are prefetch hints, revalidated under dd->lock. */
+#define READ_ONCE(x) (*(volatile typeof(x) *)&(x))
+#define WRITE_ONCE(x, value) (*(volatile typeof(x) *)&(x) = (value))
+
+/*
+ * The request and FIFO trees each own a reference to the same wrapper.
+ * The wrapper owns one request reference, independent of its tree references.
+ */
+struct pfq_rq_core {
+ struct bpf_refcount ref;
+ struct bpf_rb_node rb_node; /* request-tree link */
+ struct bpf_rb_node fifo_rb_node; /* FIFO-tree link */
+ struct request __kptr *req; /* owned request (bpf_request_acquire) */
+ u64 fifo_deadline_ns; /* boot-time FIFO deadline */
+ u64 sector; /* cached req->__sector for rq_less */
+ u32 qid; /* logical queue index */
+ bool is_meta; /* cached REQ_META for seek choose */
+};
+
+/*
+ * Urgent dispatch entry: BLK_MQ_INSERT_AT_HEAD bypasses fair queueing and is
+ * drained from dd->dispatch before the service tree.
+ */
+struct pfq_dispatch_node {
+ struct bpf_list_node node;
+ struct request __kptr *req;
+};
+
+/*
+ * The queue map owns each logical queue for the disk lifetime. The service
+ * tree owns another reference while the queue is active or held idle.
+ * Mutable scheduling fields are protected by dd->lock.
+ */
+struct pfq_queue_data {
+ struct bpf_refcount ref;
+ struct bpf_rb_node svc_rb_node; /* link in dd->service_tree */
+ u64 vtime; /* virtual time charged on dispatch */
+ u64 weight; /* virtual-time divisor */
+ u32 nr_queued; /* requests in this logical queue */
+ bool in_tree; /* in dd->service_tree; dd->lock only */
+ u32 qid; /* index 0..95; tree tie-break */
+ u8 pq_class; /* PFQ_*_CLASS */
+ u8 priority; /* 0..7 from IOPRIO_PRIO_LEVEL */
+ u8 io_type; /* PFQ_READ / WRITE / SYNC_WRITE */
+};
+
+/* Prefetch queue unlocked; detach at most one expired hold under dd->lock. */
+struct pfq_idle_cleanup {
+ struct pfq_queue_data *queue;
+ struct pfq_queue_data *service;
+};
+
+static void pfq_idle_cleanup(struct pfq_idle_cleanup *cleanup)
+{
+ if (cleanup->service)
+ bpf_obj_drop(cleanup->service);
+ if (cleanup->queue)
+ bpf_obj_drop(cleanup->queue);
+}
+
+/* Detached dispatch references for release after unlocking dd->lock. */
+struct pfq_dispatch_cleanup {
+ struct pfq_queue_data *empty_service;
+ struct pfq_queue_data *old_service;
+ struct pfq_rq_core *fifo_node;
+ struct pfq_rq_core *rq_node;
+ struct pfq_stats *stats;
+};
+
+/*
+ * dd->lock also protects mutable fields in the referenced queue objects.
+ */
+struct pfq_disk_data {
+ struct bpf_spin_lock lock; /* sole lock for all scheduler state */
+ struct bpf_list_head dispatch __contains(pfq_dispatch_node, node);
+ struct bpf_rb_root service_tree __contains(pfq_queue_data, svc_rb_node);
+ struct bpf_rb_root rq_trees[PFQ_RQ_TREES]
+ __contains(pfq_rq_core, rb_node);
+ struct bpf_rb_root fifo_trees[PFQ_RQ_TREES]
+ __contains(pfq_rq_core, fifo_rb_node);
+ s32 disk_id; /* request_queue::id */
+ u64 vtime; /* global virtual time floor */
+ u64 min_vtime; /* cached service-tree minimum */
+ u64 last_sector; /* end sector of last fair dispatch */
+ u32 nr_active; /* active or idle-held queues */
+ u32 nr_queued_total; /* pending rq in fair queues */
+ u32 in_service_dispatched; /* dispatches in current batch */
+ u32 in_service_qid; /* queue being batch-dispatched */
+ u64 idle_deadline_ns; /* checked by later callbacks */
+ u32 idle_queue_qid; /* reserved service slot */
+ bool waiting_idle; /* SPECIAL+READ idle hold active */
+ u64 fifo_expire_ns[3]; /* FIFO windows by I/O type */
+ u32 max_batch[3]; /* batch limit per io_type */
+ u32 weight_base; /* base weight from tunables */
+ u32 idle_delay_min_ms; /* SPECIAL+READ idle hold lower bound */
+ u32 idle_delay_max_ms; /* SPECIAL+READ idle hold upper bound */
+ int init_state; /* enum pfq_init_state; atomic access */
+};
+
+struct pfq_queue_handle {
+ struct pfq_queue_data __kptr *queue;
+};
+
+struct pfq_queue_key {
+ s32 disk_id;
+ u32 qid;
+};
+
+/* Exact, case-sensitive task names configured before struct_ops attach. */
+struct {
+ __uint(type, BPF_MAP_TYPE_HASH);
+ __uint(max_entries, PFQ_INTERACTIVE_MAX);
+ __type(key, struct pfq_comm_key);
+ __type(value, u8);
+} pfq_interactive_comms SEC(".maps");
+
+/*
+ * Per-disk state: key = request_queue::id; value = pfq_disk_data.
+ * Dynamic hash elements run their BTF field destructor after deletion,
+ * releasing graph roots and node-owned kptrs; reused IDs get fresh elements.
+ */
+struct {
+ __uint(type, BPF_MAP_TYPE_HASH);
+ __uint(map_flags, BPF_F_NO_PREALLOC);
+ __uint(max_entries, PFQ_DISK_MAP_MAX);
+ __type(key, s32);
+ __type(value, struct pfq_disk_data);
+} pfq_map SEC(".maps");
+
+/* Per-logical-queue kptr: key = (disk_id, qid). */
+struct {
+ __uint(type, BPF_MAP_TYPE_HASH);
+ __uint(max_entries, PFQ_QUEUE_MAP_MAX);
+ __type(key, struct pfq_queue_key);
+ __type(value, struct pfq_queue_handle);
+} pfq_queue_map SEC(".maps");
+
+/*
+ * Initialization tunables (single entry, key 0). Userspace writes them
+ * before attach; pfq_do_init_sched snapshots them for each disk.
+ */
+struct {
+ __uint(type, BPF_MAP_TYPE_ARRAY);
+ __uint(max_entries, 1);
+ __type(key, u32);
+ __type(value, struct pfq_tunables);
+} pfq_tunables SEC(".maps");
+
+/*
+ * Keep the disk template off the limited BPF stack. Per-CPU storage
+ * separates concurrent initializers running on different CPUs.
+ */
+struct {
+ __uint(type, BPF_MAP_TYPE_PERCPU_ARRAY);
+ __uint(max_entries, 1);
+ __type(key, u32);
+ __type(value, struct pfq_disk_data);
+} pfq_disk_scratch SEC(".maps");
+
+/* One independent set of CPU counters per pfq_map disk entry. */
+struct {
+ __uint(type, BPF_MAP_TYPE_PERCPU_HASH);
+ __uint(max_entries, PFQ_DISK_MAP_MAX);
+ __type(key, s32);
+ __type(value, struct pfq_stats);
+} pfq_stats_map SEC(".maps");
+
+/* Lookup before taking dd->lock; map helpers are forbidden while locked. */
+static __always_inline struct pfq_stats *pfq_stats_lookup(int disk_id)
+{
+ return bpf_map_lookup_elem(&pfq_stats_map, &disk_id);
+}
+
+/* A pre-fetched per-CPU pointer is usable with or without dd->lock. */
+static __always_inline void pfq_stat_add(struct pfq_stats *stats, u32 idx,
+ u64 val)
+{
+ if (stats && idx < PFQ_STAT_MAX)
+ /* Also protect against same-CPU interrupt/reentrant writers. */
+ __sync_fetch_and_add(&stats->counters[idx], val);
+}
+
+/* Boot time, including suspend, for FIFO and idle deadlines. */
+static __always_inline u64 pfq_now_ns(void)
+{
+ return bpf_ktime_get_boot_ns();
+}
+
+#ifndef NSEC_PER_SEC
+#define NSEC_PER_SEC 1000000000ULL
+#endif
+#ifndef CONFIG_HZ
+#define CONFIG_HZ 1000
+#endif
+#define PFQ_NS_PER_JIFFY (NSEC_PER_SEC / CONFIG_HZ)
+
+static __always_inline u64 pfq_jiffies_to_deadline_ns(unsigned long jdeadline,
+ u64 now_ns)
+{
+ u64 jnow = bpf_jiffies64();
+ u64 delta_j = (u64)jdeadline > jnow ? (u64)jdeadline - jnow : 0;
+
+ return now_ns + delta_j * PFQ_NS_PER_JIFFY;
+}
+
+static __always_inline bool pfq_qid_valid(u32 qid)
+{
+ return qid < PFQ_TOTAL_QUEUES;
+}
+
+/* Return a non-owning queue pointer; caller must not hold dd->lock. */
+static __always_inline struct pfq_queue_data *
+pfq_queue_lookup(s32 disk_id, u32 qid)
+{
+ struct pfq_queue_key key = {};
+ struct pfq_queue_handle *hp;
+
+ if (!pfq_qid_valid(qid))
+ return NULL;
+
+ key.disk_id = disk_id;
+ key.qid = qid;
+ hp = bpf_map_lookup_elem(&pfq_queue_map, &key);
+ if (!hp || !hp->queue)
+ return NULL;
+ return hp->queue;
+}
+
+/*
+ * Initialization-only lookup: drop dd->lock for the map helper and return
+ * with it held. The INITIALIZING owner keeps queue entries alive across
+ * this unlocked interval. The returned pointer is non-owning.
+ */
+static __always_inline struct pfq_queue_data *
+pfq_queue_scalar(struct pfq_disk_data *dd, u32 qid)
+{
+ struct pfq_queue_data *queue;
+ s32 disk_id;
+
+ if (!dd || !pfq_qid_valid(qid))
+ return NULL;
+
+ disk_id = dd->disk_id;
+ bpf_spin_unlock(&dd->lock);
+ queue = pfq_queue_lookup(disk_id, qid);
+ bpf_spin_lock(&dd->lock);
+ return queue;
+}
+
+/* Map lookup + owning ref. Caller must not hold dd->lock. */
+static __always_inline struct pfq_queue_data *
+pfq_queue_acquire(s32 disk_id, u32 qid)
+{
+ struct pfq_queue_data *queue = NULL;
+ struct pfq_queue_key key = {};
+ struct pfq_queue_handle *hp;
+
+ if (!pfq_qid_valid(qid))
+ return NULL;
+
+ key.disk_id = disk_id;
+ key.qid = qid;
+ hp = bpf_map_lookup_elem(&pfq_queue_map, &key);
+ if (hp && hp->queue)
+ queue = bpf_refcount_acquire(hp->queue);
+ return queue;
+}
+
+/*
+ * Select roots at constant offsets: graph kfuncs reject a root pointer at
+ * a variable offset into the map value. Caller holds dd->lock.
+ */
+static __always_inline struct bpf_rb_root *
+pfq_rq_tree(struct pfq_disk_data *dd, u32 qid)
+{
+ u8 ty;
+
+ if (!pfq_qid_valid(qid))
+ return NULL;
+ ty = PFQ_GET_TYPE(qid);
+ if (ty == PFQ_READ)
+ return &dd->rq_trees[PFQ_READ];
+ if (ty == PFQ_WRITE)
+ return &dd->rq_trees[PFQ_WRITE];
+ if (ty == PFQ_SYNC_WRITE)
+ return &dd->rq_trees[PFQ_SYNC_WRITE];
+ return NULL;
+}
+
+static __always_inline struct bpf_rb_root *
+pfq_fifo_tree(struct pfq_disk_data *dd, u32 qid)
+{
+ u8 ty;
+
+ if (!pfq_qid_valid(qid))
+ return NULL;
+ ty = PFQ_GET_TYPE(qid);
+ if (ty == PFQ_READ)
+ return &dd->fifo_trees[PFQ_READ];
+ if (ty == PFQ_WRITE)
+ return &dd->fifo_trees[PFQ_WRITE];
+ if (ty == PFQ_SYNC_WRITE)
+ return &dd->fifo_trees[PFQ_SYNC_WRITE];
+ return NULL;
+}
+
+/*
+ * Lock-held detach helpers: bpf_rbtree_remove returns owning refs that must
+ * be bpf_obj_drop'd only after releasing bpf_spin_lock.
+ */
+static __always_inline struct pfq_rq_core *
+pfq_rq_tree_remove_take(struct bpf_rb_root *root, struct bpf_rb_node *node)
+{
+ struct bpf_rb_node *rb;
+
+ if (!root || !node)
+ return NULL;
+ rb = bpf_rbtree_remove(root, node);
+ if (!rb)
+ return NULL;
+ return container_of(rb, struct pfq_rq_core, rb_node);
+}
+
+static __always_inline struct pfq_rq_core *
+pfq_fifo_tree_remove_take(struct bpf_rb_root *root, struct bpf_rb_node *node)
+{
+ struct bpf_rb_node *rb;
+
+ if (!root || !node)
+ return NULL;
+ rb = bpf_rbtree_remove(root, node);
+ if (!rb)
+ return NULL;
+ return container_of(rb, struct pfq_rq_core, fifo_rb_node);
+}
+
+/**
+ * pfq_dispatch_push_dn - Append an urgent dispatch node
+ * @dispatch: urgent dispatch list protected by the disk lock
+ * @dn: node whose owning reference is transferred to the list
+ *
+ * bpf_list_push_back consumes the node reference even on failure. The
+ * caller must not drop dn or restore dn->req after a failed push.
+ *
+ * Context: Caller holds the disk lock.
+ *
+ * Return: 0 on success, or -EINVAL for invalid arguments or a failed push.
+ */
+static __noinline int
+pfq_dispatch_push_dn(struct bpf_list_head *dispatch,
+ struct pfq_dispatch_node *dn)
+{
+ if (!dispatch || !dn)
+ return -EINVAL;
+
+ if (bpf_list_push_back(dispatch, &dn->node))
+ return -EINVAL;
+
+ return 0;
+}
+
+/*
+ * Return the displaced request reference for release after unlocking.
+ * If core is NULL, return req so the caller retains ownership.
+ */
+static __always_inline struct request *
+pfq_core_put_req_stash(struct pfq_rq_core *core, struct request *req)
+{
+ if (!req)
+ return NULL;
+ if (!core)
+ return req;
+ return bpf_kptr_xchg(&core->req, req);
+}
+
+/*
+ * Map updates skip kptr fields, so transfer the queue reference with
+ * bpf_kptr_xchg. On success the map owns it; on failure the caller does.
+ */
+static int pfq_queue_map_insert(s32 disk_id, u32 qid,
+ struct pfq_queue_data *queue)
+{
+ struct pfq_queue_handle handle = {}, *hp;
+ struct pfq_queue_key key = {};
+ struct pfq_queue_data *old;
+ int ret;
+
+ key.disk_id = disk_id;
+ key.qid = qid;
+ ret = bpf_map_update_elem(&pfq_queue_map, &key, &handle, BPF_NOEXIST);
+ if (ret)
+ return ret;
+
+ hp = bpf_map_lookup_elem(&pfq_queue_map, &key);
+ if (!hp)
+ return -ENOENT;
+
+ old = bpf_kptr_xchg(&hp->queue, queue);
+ if (old)
+ bpf_obj_drop(old);
+ return 0;
+}
+
+static void pfq_queues_destroy(s32 disk_id)
+{
+ struct pfq_queue_key key = {};
+ struct pfq_queue_handle *hp;
+ struct pfq_queue_data __kptr *queue;
+ u32 i;
+
+ key.disk_id = disk_id;
+ for (i = 0; i < PFQ_TOTAL_QUEUES; i++) {
+ key.qid = i;
+ hp = bpf_map_lookup_elem(&pfq_queue_map, &key);
+ if (hp) {
+ queue = bpf_kptr_xchg(&hp->queue, NULL);
+ bpf_map_delete_elem(&pfq_queue_map, &key);
+ if (queue)
+ bpf_obj_drop(queue);
+ } else {
+ bpf_map_delete_elem(&pfq_queue_map, &key);
+ }
+ }
+}
+
+/**
+ * pfq_ioprio_to_class - Map kernel I/O priority class to PFQ class
+ * @ioprio: encoded I/O priority, including class and level
+ *
+ * Unspecified or invalid I/O priority classes use best-effort service.
+ * Task-name overrides are applied separately by pfq_class_from_*().
+ *
+ * Return: PFQ class; unspecified or invalid classes use PFQ_BE_CLASS.
+ */
+static __always_inline u8 pfq_ioprio_to_class(u16 ioprio)
+{
+ u8 iclass = (ioprio >> PFQ_IOPRIO_CLASS_SHIFT) & 0x7;
+
+ if (iclass <= PFQ_IOPRIO_CLASS_IDLE) {
+ switch (iclass) {
+ case PFQ_IOPRIO_CLASS_RT:
+ return PFQ_RT_CLASS;
+ case PFQ_IOPRIO_CLASS_BE:
+ return PFQ_BE_CLASS;
+ case PFQ_IOPRIO_CLASS_IDLE:
+ return PFQ_IDLE_CLASS;
+ default:
+ return PFQ_BE_CLASS;
+ }
+ }
+ return PFQ_BE_CLASS;
+}
+
+/*
+ * Match the current thread's comm before taking dd->lock. Worker-submitted
+ * I/O is classified by the worker name, not the original userspace submitter.
+ */
+static __always_inline bool pfq_current_is_interactive(void)
+{
+ struct pfq_comm_key key = {};
+ u8 *enabled;
+
+ if (bpf_get_current_comm(key.name, sizeof(key.name)))
+ return false;
+ enabled = bpf_map_lookup_elem(&pfq_interactive_comms, &key);
+ return enabled && *enabled;
+}
+
+/**
+ * pfq_class_from_ioprio - Resolve the request priority class
+ * @ioprio: encoded request I/O priority
+ * @interactive: whether the current thread matches a configured comm
+ *
+ * Return: PFQ_SPECIAL_CLASS for an interactive thread, otherwise the I/O
+ * class.
+ */
+static __always_inline u8 pfq_class_from_ioprio(u16 ioprio, bool interactive)
+{
+ if (interactive)
+ return PFQ_SPECIAL_CLASS;
+ return pfq_ioprio_to_class(ioprio);
+}
+
+/**
+ * pfq_class_from_bio - Resolve the bio priority class
+ * @bio: bio supplying the I/O priority
+ * @interactive: whether the current thread matches a configured comm
+ *
+ * Return: PFQ_SPECIAL_CLASS for an interactive thread, otherwise the I/O
+ * class.
+ */
+static __always_inline u8 pfq_class_from_bio(struct bio *bio, bool interactive)
+{
+ if (interactive)
+ return PFQ_SPECIAL_CLASS;
+ return pfq_ioprio_to_class(bio->bi_ioprio);
+}
+
+static __always_inline u8 pfq_ioprio_level(u16 ioprio)
+{
+ return ioprio & 0x7;
+}
+
+/**
+ * pfq_rq_io_type - Classify request direction and synchrony
+ * @rq: request supplying the operation and REQ_SYNC flag
+ *
+ * Match rq_data_dir(): odd opcodes use write queues, with REQ_SYNC
+ * selecting the synchronous-write queue.
+ *
+ * Return: PFQ_READ, PFQ_WRITE or PFQ_SYNC_WRITE.
+ */
+static __always_inline u8 pfq_rq_io_type(struct request *rq)
+{
+ if ((rq->cmd_flags & REQ_OP_MASK) & 1)
+ return ((rq->cmd_flags & PFQ_REQ_SYNC) ? PFQ_SYNC_WRITE : PFQ_WRITE);
+ return PFQ_READ;
+}
+
+static __always_inline u8 pfq_bio_io_type(struct bio *bio)
+{
+ if ((bio->bi_opf & REQ_OP_MASK) & 1)
+ return ((bio->bi_opf & PFQ_REQ_SYNC) ? PFQ_SYNC_WRITE : PFQ_WRITE);
+ return PFQ_READ;
+}
+
+/*
+ * Index fifo_expire_ns / max_batch with constant offsets only.
+ * Even a clamped variable index is rejected on PTR_TO_BTF_ID.
+ */
+static __always_inline u64 pfq_fifo_expire_ns(struct pfq_disk_data *dd, u8 io_type)
+{
+ if (io_type == PFQ_READ)
+ return dd->fifo_expire_ns[PFQ_READ];
+ if (io_type == PFQ_SYNC_WRITE)
+ return dd->fifo_expire_ns[PFQ_SYNC_WRITE];
+ return dd->fifo_expire_ns[PFQ_WRITE];
+}
+
+static __always_inline u32 pfq_max_batch(struct pfq_disk_data *dd, u8 io_type)
+{
+ if (io_type == PFQ_READ)
+ return dd->max_batch[PFQ_READ];
+ if (io_type == PFQ_SYNC_WRITE)
+ return dd->max_batch[PFQ_SYNC_WRITE];
+ return dd->max_batch[PFQ_WRITE];
+}
+
+/**
+ * pfq_queue_index - Compute the logical queue index
+ * @pq_class: PFQ priority class
+ * @level: priority level; invalid levels use IOPRIO_BE_NORM
+ * @io_type: PFQ_READ, PFQ_WRITE or PFQ_SYNC_WRITE
+ *
+ * Return: Flat index for the class, priority level and I/O type.
+ */
+static __always_inline u32 pfq_queue_index(u8 pq_class, u8 level, u8 io_type)
+{
+ if (level >= PFQ_PRIO_LEVELS)
+ level = PFQ_IOPRIO_BE_NORM;
+ return PFQ_INDEX(pq_class, level, io_type);
+}
+
+/**
+ * pfq_rq_ioprio - Read request priority with a best-effort fallback
+ * @rq: request whose first bio supplies the I/O priority
+ *
+ * Internal requests without a bio use normal best-effort priority.
+ *
+ * Return: Encoded bio priority, or normal best-effort priority without a
+ * bio.
+ */
+static __always_inline u16 pfq_rq_ioprio(struct request *rq)
+{
+ if (rq->bio)
+ return rq->bio->bi_ioprio;
+ return (PFQ_IOPRIO_CLASS_BE << PFQ_IOPRIO_CLASS_SHIFT) | PFQ_IOPRIO_BE_NORM;
+}
+
+static __always_inline u32 pfq_rq_qid(struct request *rq, bool interactive)
+{
+ u16 ioprio = pfq_rq_ioprio(rq);
+
+ return pfq_queue_index(pfq_class_from_ioprio(ioprio, interactive),
+ pfq_ioprio_level(ioprio),
+ pfq_rq_io_type(rq));
+}
+
+static __always_inline u32 pfq_bio_qid(struct bio *bio, bool interactive)
+{
+ return pfq_queue_index(pfq_class_from_bio(bio, interactive),
+ pfq_ioprio_level(bio->bi_ioprio),
+ pfq_bio_io_type(bio));
+}
+
+/**
+ * pfq_prio_coeff - Choose the per-level weight increment
+ * @queue: queue supplying the priority class and I/O type
+ *
+ * Return: Read or write coefficient, with a larger value for SPECIAL
+ * reads.
+ */
+static __always_inline u32 pfq_prio_coeff(struct pfq_queue_data *queue)
+{
+ if (queue->pq_class == PFQ_SPECIAL_CLASS && queue->io_type == PFQ_READ)
+ return PFQ_PRIO_LEVEL_COEFF_SP_READ;
+ if (queue->io_type == PFQ_READ)
+ return PFQ_PRIO_LEVEL_COEFF_READ;
+ return PFQ_PRIO_LEVEL_COEFF_WRITE;
+}
+
+/**
+ * pfq_calc_weight - Calculate the virtual-time weight
+ * @dd: disk state supplying the configured base weight
+ * @queue: queue supplying class, priority level and I/O type
+ *
+ * Higher weights slow virtual-time growth and give more service.
+ * SPECIAL reads receive an additional percentage multiplier.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: Queue weight, with a minimum of one.
+ */
+static __always_inline u64 pfq_calc_weight(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ u64 base = dd->weight_base, final;
+
+ base += (PFQ_PRIO_LEVELS - 1 - queue->priority) * pfq_prio_coeff(queue);
+ final = base;
+ if (queue->pq_class == PFQ_SPECIAL_CLASS && queue->io_type == PFQ_READ)
+ final = base * PFQ_WEIGHT_MULT_SPECIAL / 100;
+ return final ? final : 1;
+}
+
+/**
+ * pfq_apply_tunables - Apply a prefetched disk configuration
+ * @dd: disk being configured
+ * @tun: snapshot read from pfq_tunables before taking the lock
+ *
+ * Zero tunables select defaults. Apply the prefetched snapshot under
+ * dd->lock before initializing queue weights.
+ *
+ * Context: Caller holds @dd->lock.
+ */
+static void pfq_apply_tunables(struct pfq_disk_data *dd,
+ const struct pfq_tunables *tun)
+{
+ if (tun->weight_base)
+ dd->weight_base = tun->weight_base;
+ else
+ dd->weight_base = PFQ_WEIGHT_BASE_DEFAULT;
+
+ if (tun->idle_delay_min_ms)
+ dd->idle_delay_min_ms = tun->idle_delay_min_ms;
+ else
+ dd->idle_delay_min_ms = PFQ_IDLE_DELAY_MIN_MS_DEFAULT;
+
+ if (tun->idle_delay_max_ms)
+ dd->idle_delay_max_ms = tun->idle_delay_max_ms;
+ else
+ dd->idle_delay_max_ms = PFQ_IDLE_DELAY_MAX_MS_DEFAULT;
+
+ if (tun->batch_read)
+ dd->max_batch[PFQ_READ] = tun->batch_read;
+ else
+ dd->max_batch[PFQ_READ] = PFQ_BATCH_READ_DEFAULT;
+
+ if (tun->batch_sync_write)
+ dd->max_batch[PFQ_SYNC_WRITE] = tun->batch_sync_write;
+ else
+ dd->max_batch[PFQ_SYNC_WRITE] = PFQ_BATCH_SYNC_WRITE_DEFAULT;
+
+ if (tun->batch_write)
+ dd->max_batch[PFQ_WRITE] = tun->batch_write;
+ else
+ dd->max_batch[PFQ_WRITE] = PFQ_BATCH_WRITE_DEFAULT;
+}
+
+static __always_inline u64 pfq_distance(sector_t pos, u64 last_sector)
+{
+ if (pos >= last_sector)
+ return pos - last_sector;
+ return last_sector - pos;
+}
+
+/* Service tree ordering: lowest vtime first, tie-break by qid. */
+static bool svc_less(struct bpf_rb_node *a, const struct bpf_rb_node *b)
+{
+ struct pfq_queue_data *qa, *qb;
+
+ qa = container_of(a, struct pfq_queue_data, svc_rb_node);
+ qb = container_of(b, struct pfq_queue_data, svc_rb_node);
+ if (qa->vtime != qb->vtime)
+ return qa->vtime < qb->vtime;
+ return qa->qid < qb->qid;
+}
+
+/*
+ * Order by (qid, sector) using cached fields. Avoid request kptr exchanges
+ * in the comparator; request release is not KF_SPINLOCK_SAFE.
+ */
+static bool rq_less(struct bpf_rb_node *a, const struct bpf_rb_node *b)
+{
+ struct pfq_rq_core *na, *nb;
+
+ na = container_of(a, struct pfq_rq_core, rb_node);
+ nb = container_of(b, struct pfq_rq_core, rb_node);
+ if (na->qid != nb->qid)
+ return na->qid < nb->qid;
+ return na->sector < nb->sector;
+}
+
+/* io_type FIFO tree: (qid, deadline, sector). */
+static bool fifo_less(struct bpf_rb_node *a, const struct bpf_rb_node *b)
+{
+ struct pfq_rq_core *na, *nb;
+
+ na = container_of(a, struct pfq_rq_core, fifo_rb_node);
+ nb = container_of(b, struct pfq_rq_core, fifo_rb_node);
+ if (na->qid != nb->qid)
+ return na->qid < nb->qid;
+ if (na->fifo_deadline_ns != nb->fifo_deadline_ns)
+ return na->fifo_deadline_ns < nb->fifo_deadline_ns;
+ return na->sector < nb->sector;
+}
+
+/* Tag bit in request->elv.priv[1] after dispatch (see pfq_rq_bind_qid). */
+#define PFQ_RQ_QID_TAG 0x80000000U
+
+/*
+ * UFQ owns elv.priv[0]. PFQ uses priv[1] for a non-owning core address while
+ * queued, then a tagged qid after dispatch. BPF cannot write the slot
+ * directly, so stores use the KF_SPINLOCK_SAFE kfunc.
+ */
+static __always_inline void pfq_rq_set_core(struct request *rq,
+ struct pfq_rq_core *core)
+{
+ __u64 addr = 0;
+
+ /*
+ * Pass address as scalar: do not hand a ref-tracked MEM_ALLOC pointer
+ * into priv (non-owning cache only).
+ */
+ bpf_probe_read_kernel(&addr, sizeof(addr), &core);
+ bpf_request_set_elv_priv1(rq, addr);
+}
+
+static __always_inline void pfq_rq_clear_node(struct request *rq)
+{
+ bpf_request_set_elv_priv1(rq, 0);
+}
+
+/*
+ * Keep the dispatched queue identity for finish_req after the wrapper is
+ * dropped. pfq_rq_bound_qid() checks both the tag and the encoded qid range.
+ */
+static __always_inline void pfq_rq_bind_qid(struct request *rq, u32 qid)
+{
+ bpf_request_set_elv_priv1(rq, (u64)(qid | PFQ_RQ_QID_TAG));
+}
+
+/* Must run before dd->lock: bpf_core_read() calls a probe-read helper. */
+static __always_inline u32 pfq_rq_bound_qid(struct request *rq)
+{
+ u64 tag = 0;
+
+ /* A C cast would retain the verifier pointer type. */
+ if (bpf_core_read(&tag, sizeof(tag), &rq->elv.priv[1]))
+ return PFQ_QID_NONE;
+
+ /* Only accept encoded qids, never a live pfq_rq_core pointer. */
+ if (!(tag & PFQ_RQ_QID_TAG) ||
+ (tag & ~(u64)PFQ_RQ_QID_TAG) >= PFQ_TOTAL_QUEUES)
+ return PFQ_QID_NONE;
+ return (u32)tag & ~PFQ_RQ_QID_TAG;
+}
+
+/*
+ * Graph nodes are opaque and there is no parent-access kfunc. Find sector
+ * neighbors by searching from the root using the available child kfuncs.
+ */
+
+/* Earliest FIFO deadline within qid. Caller holds dd->lock. */
+static struct pfq_rq_core *
+pfq_fifo_first_qid(struct bpf_rb_root *root, u32 qid)
+{
+ struct pfq_rq_core *c, *best = NULL;
+ struct bpf_rb_node *node;
+ int i;
+
+ if (!root)
+ return NULL;
+
+ node = bpf_rbtree_root(root);
+ for (i = 0; node && i < PFQ_RB_HEIGHT_MAX; i++) {
+ c = container_of(node, struct pfq_rq_core, fifo_rb_node);
+ if (c->qid < qid) {
+ node = bpf_rbtree_right(root, node);
+ } else {
+ best = c;
+ node = bpf_rbtree_left(root, node);
+ }
+ }
+ if (!best || best->qid != qid)
+ return NULL;
+ return best;
+}
+
+/* Nearest sector >= @sector within this qid; tree ordered by (qid, sector). */
+static struct pfq_rq_core *pfq_rq_tree_next(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ u64 sector)
+{
+ struct bpf_rb_root *root = pfq_rq_tree(dd, queue->qid);
+ struct pfq_rq_core *core, *best = NULL;
+ struct bpf_rb_node *node;
+ int i;
+
+ if (!root)
+ return NULL;
+ node = bpf_rbtree_root(root);
+ for (i = 0; node && i < PFQ_RB_HEIGHT_MAX; i++) {
+ core = container_of(node, struct pfq_rq_core, rb_node);
+ if (core->qid < queue->qid ||
+ (core->qid == queue->qid && core->sector < sector)) {
+ node = bpf_rbtree_right(root, node);
+ } else {
+ if (core->qid == queue->qid)
+ best = core;
+ node = bpf_rbtree_left(root, node);
+ }
+ }
+ return best;
+}
+
+/* Nearest sector < @sector within this qid. */
+static struct pfq_rq_core *pfq_rq_tree_prev(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ u64 sector)
+{
+ struct bpf_rb_root *root = pfq_rq_tree(dd, queue->qid);
+ struct pfq_rq_core *core, *best = NULL;
+ struct bpf_rb_node *node;
+ int i;
+
+ if (!root)
+ return NULL;
+ node = bpf_rbtree_root(root);
+ for (i = 0; node && i < PFQ_RB_HEIGHT_MAX; i++) {
+ core = container_of(node, struct pfq_rq_core, rb_node);
+ if (core->qid > queue->qid ||
+ (core->qid == queue->qid && core->sector >= sector)) {
+ node = bpf_rbtree_left(root, node);
+ } else {
+ if (core->qid == queue->qid)
+ best = core;
+ node = bpf_rbtree_right(root, node);
+ }
+ }
+ return best;
+}
+
+/*
+ * BPF v3 cannot directly encode acquire loads or release stores. A no-op CAS
+ * reads the state with a full barrier (stronger than acquire); release XCHG
+ * publishes preceding queue initialization or failure cleanup. Use these for
+ * all accesses to a published disk header, including accesses under dd->lock.
+ */
+static __always_inline int pfq_init_state_load(struct pfq_disk_data *dd)
+{
+ return __sync_val_compare_and_swap(&dd->init_state,
+ PFQ_UNINITIALIZED,
+ PFQ_UNINITIALIZED);
+}
+
+static __always_inline void pfq_init_state_store(struct pfq_disk_data *dd,
+ int state)
+{
+ (void)__atomic_exchange_n(&dd->init_state, state, __ATOMIC_RELEASE);
+}
+
+/**
+ * pfq_disk_lookup - Look up an initialized disk
+ * @disk_id: request_queue::id used as the map key
+ *
+ * Acquire initialization before exposing a ready disk to callbacks.
+ * UFQ must quiesce callbacks before exit_sched: the state check does not pin
+ * the disk or its queue entries against teardown.
+ *
+ * Context: Caller must not hold a BPF spin lock.
+ *
+ * Return: Map-value pointer, or NULL if the disk is absent or not
+ * initialized.
+ */
+static __always_inline struct pfq_disk_data *pfq_disk_lookup(int disk_id)
+{
+ struct pfq_disk_data *dd;
+
+ dd = bpf_map_lookup_elem(&pfq_map, &disk_id);
+ if (!dd || pfq_init_state_load(dd) != PFQ_INITIALIZED)
+ return NULL;
+ return dd;
+}
+
+/* Raw map lookup (init/exit only); ignores init_state. */
+static __always_inline struct pfq_disk_data *pfq_disk_lookup_raw(int disk_id)
+{
+ return bpf_map_lookup_elem(&pfq_map, &disk_id);
+}
+
+/*
+ * Reset special BTF fields with empty compound literals; a bulk memset
+ * would generate accesses the verifier rejects.
+ *
+ * Keep this inlinable: noinline + compound-literal temps can leave a stack
+ * address in R0 at subprog exit (verifier: cannot return stack pointer).
+ */
+static void pfq_disk_scratch_reset(struct pfq_disk_data *d)
+{
+ u32 i;
+
+ d->lock = (struct bpf_spin_lock){};
+ d->dispatch = (struct bpf_list_head){};
+ d->service_tree = (struct bpf_rb_root){};
+ for (i = 0; i < PFQ_RQ_TREES; i++) {
+ d->rq_trees[i] = (struct bpf_rb_root){};
+ d->fifo_trees[i] = (struct bpf_rb_root){};
+ }
+ d->disk_id = 0;
+ d->vtime = 0;
+ d->min_vtime = 0;
+ d->last_sector = 0;
+ d->nr_active = 0;
+ d->nr_queued_total = 0;
+ d->in_service_qid = PFQ_QID_NONE;
+ d->in_service_dispatched = 0;
+ d->waiting_idle = false;
+ d->idle_queue_qid = PFQ_QID_NONE;
+ d->idle_deadline_ns = 0;
+ d->fifo_expire_ns[PFQ_READ] = 0;
+ d->fifo_expire_ns[PFQ_WRITE] = 0;
+ d->fifo_expire_ns[PFQ_SYNC_WRITE] = 0;
+ d->max_batch[PFQ_READ] = 0;
+ d->max_batch[PFQ_WRITE] = 0;
+ d->max_batch[PFQ_SYNC_WRITE] = 0;
+ d->weight_base = 0;
+ d->idle_delay_min_ms = 0;
+ d->idle_delay_max_ms = 0;
+ d->init_state = PFQ_UNINITIALIZED;
+}
+
+/*
+ * Publish an UNINITIALIZED disk header. BPF_NOEXIST prevents replacement of
+ * an existing header; callers then claim initialization under that dd->lock.
+ * Keep the header on init failure so concurrent callers share the same lock
+ * and state until exit_sched.
+ */
+static int pfq_disk_create(int disk_id)
+{
+ struct pfq_disk_data *scratch;
+ u32 zero = 0;
+
+ scratch = bpf_map_lookup_elem(&pfq_disk_scratch, &zero);
+ if (!scratch)
+ return -ENOMEM;
+
+ pfq_disk_scratch_reset(scratch);
+ scratch->disk_id = disk_id;
+ return bpf_map_update_elem(&pfq_map, &disk_id, scratch, BPF_NOEXIST);
+}
+
+/*
+ * Only the INITIALIZING owner may create or unwind queue entries. Return 1
+ * when this attempt created the stats entry, 0 when it reused an existing one.
+ */
+static int pfq_disk_create_queues(int disk_id)
+{
+ struct pfq_stats empty_stats = {};
+ struct pfq_queue_data *queue;
+ bool stats_created = false;
+ int ret;
+ u32 i;
+
+ ret = bpf_map_update_elem(&pfq_stats_map, &disk_id, &empty_stats,
+ BPF_NOEXIST);
+ if (ret == -EEXIST) {
+ if (!pfq_stats_lookup(disk_id))
+ return -ENOENT;
+ } else if (ret) {
+ return ret;
+ } else {
+ stats_created = true;
+ }
+
+ for (i = 0; i < PFQ_TOTAL_QUEUES; i++) {
+ queue = bpf_obj_new(typeof(*queue));
+ if (!queue) {
+ ret = -ENOMEM;
+ goto err_queues;
+ }
+
+ ret = pfq_queue_map_insert(disk_id, i, queue);
+ if (ret) {
+ bpf_obj_drop(queue);
+ goto err_queues;
+ }
+ }
+
+ return stats_created ? 1 : 0;
+
+err_queues:
+ pfq_queues_destroy(disk_id);
+ if (stats_created)
+ bpf_map_delete_elem(&pfq_stats_map, &disk_id);
+ return ret;
+}
+
+/*
+ * Keep map-value and queue-object loads in separate unoptimized subprograms.
+ * Merging them can make the verifier reject one instruction used with
+ * different pointer types.
+ */
+static __attribute__((noinline, optnone)) void
+pfq_min_vtime_from_dd(struct pfq_disk_data *dd)
+{
+ dd->min_vtime = dd->vtime;
+}
+
+static __attribute__((noinline, optnone)) void
+pfq_min_vtime_from_queue(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ dd->min_vtime = queue->vtime;
+}
+
+static __noinline void pfq_refresh_min_vtime(struct pfq_disk_data *dd)
+{
+ struct pfq_queue_data *queue;
+ struct bpf_rb_node *node;
+
+ node = bpf_rbtree_first(&dd->service_tree);
+ if (!node) {
+ pfq_min_vtime_from_dd(dd);
+ return;
+ }
+ queue = container_of(node, struct pfq_queue_data, svc_rb_node);
+ pfq_min_vtime_from_queue(dd, queue);
+}
+
+/**
+ * pfq_insert_queue_to_service_tree - Insert a queue at its current vtime
+ * @dd: disk whose service tree receives the queue
+ * @queue: queue for which the caller retains an owning reference
+ *
+ * Acquire a second owning reference for the service tree without consuming
+ * the caller's @queue reference. bpf_rbtree_add() consumes tree_q on both
+ * success and failure, so do not drop tree_q again on insertion failure.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: true if inserted, or false if reference acquisition or insertion
+ * fails.
+ */
+static bool pfq_insert_queue_to_service_tree(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ struct pfq_queue_data *tree_q;
+ int ret;
+
+ tree_q = bpf_refcount_acquire(queue);
+ if (!tree_q)
+ return false;
+
+ ret = bpf_rbtree_add(&dd->service_tree, &tree_q->svc_rb_node, svc_less);
+ if (ret)
+ return false;
+ pfq_refresh_min_vtime(dd);
+ return true;
+}
+
+static void pfq_deactivate_queue(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ if (!queue->in_tree)
+ return;
+ queue->in_tree = false;
+ if (dd->nr_active > 0)
+ dd->nr_active--;
+}
+
+/*
+ * Detach and deactivate atomically. Return the tree's owning ref for the
+ * caller to drop outside dd->lock; never unlock inside this helper.
+ */
+static struct pfq_queue_data *
+pfq_remove_queue_from_service_tree_locked(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ struct bpf_rb_node *removed;
+
+ if (!queue->in_tree)
+ return NULL;
+ removed = bpf_rbtree_remove(&dd->service_tree, &queue->svc_rb_node);
+ if (!removed)
+ return NULL;
+
+ pfq_deactivate_queue(dd, queue);
+ pfq_refresh_min_vtime(dd);
+ return container_of(removed, struct pfq_queue_data, svc_rb_node);
+}
+
+/* Retain the tree reference and active count while updating its position. */
+static void pfq_requeue_in_service_tree(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ struct pfq_queue_data *owned;
+ struct bpf_rb_node *removed;
+ int ret;
+
+ if (!queue->in_tree)
+ return;
+
+ removed = bpf_rbtree_remove(&dd->service_tree, &queue->svc_rb_node);
+ if (!removed)
+ goto deactivate;
+
+ owned = container_of(removed, struct pfq_queue_data, svc_rb_node);
+ ret = bpf_rbtree_add(&dd->service_tree, &owned->svc_rb_node, svc_less);
+ if (ret)
+ goto deactivate;
+
+ pfq_refresh_min_vtime(dd);
+ return;
+
+deactivate:
+ pfq_deactivate_queue(dd, queue);
+}
+
+/**
+ * pfq_activate_queue - Activate a queue at the virtual-time floor
+ * @dd: per-disk scheduler state
+ * @queue: queue receiving its first queued request
+ *
+ * Publish tree membership and active accounting together under dd->lock.
+ * A newly active queue starts no earlier than the disk's virtual-time floor.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: true if already active or successfully activated, otherwise
+ * false.
+ */
+static bool pfq_activate_queue(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ if (queue->in_tree)
+ return true;
+
+ if (queue->vtime < dd->vtime)
+ queue->vtime = dd->vtime;
+
+ if (!pfq_insert_queue_to_service_tree(dd, queue))
+ return false;
+
+ queue->in_tree = true;
+ dd->nr_active++;
+ return true;
+}
+
+/**
+ * pfq_queue_snap - Snapshot queue scheduling fields
+ * @queue: logical queue being operated on
+ * @nq: output request count
+ * @in_tree: output service-tree membership flag
+ *
+ * Caller holds dd->lock when reading queue scheduling fields.
+ *
+ * Context: Caller holds the disk lock.
+ */
+static __always_inline void pfq_queue_snap(struct pfq_queue_data *queue,
+ u32 *nq, bool *in_tree)
+{
+ *nq = queue->nr_queued;
+ *in_tree = queue->in_tree;
+}
+
+/**
+ * pfq_queue_nq_get - Read the queued request count
+ * @queue: queue whose pending requests are counted
+ *
+ * Context: Caller holds the disk lock.
+ *
+ * Return: Number of requests queued in the logical queue.
+ */
+static __always_inline u32 pfq_queue_nq_get(struct pfq_queue_data *queue)
+{
+ bool in_tree;
+ u32 nq;
+
+ pfq_queue_snap(queue, &nq, &in_tree);
+ return nq;
+}
+
+/* Clear SPECIAL+READ idle-hold state without removing queue from tree. */
+static void pfq_clear_idle_hold(struct pfq_disk_data *dd)
+{
+ dd->waiting_idle = false;
+ WRITE_ONCE(dd->idle_queue_qid, PFQ_QID_NONE);
+ dd->idle_deadline_ns = 0;
+}
+
+/*
+ * Clang may merge adjacent u32 stores into u64; verifier rejects that as
+ * access beyond BTF member bounds (e.g. nr_active).
+ */
+static __always_inline void pfq_compiler_barrier(void)
+{
+ asm volatile("" ::: "memory");
+}
+
+static void pfq_disk_reset_counters(struct pfq_disk_data *dd)
+{
+ dd->nr_active = 0;
+ pfq_compiler_barrier();
+ dd->nr_queued_total = 0;
+ pfq_compiler_barrier();
+ WRITE_ONCE(dd->in_service_qid, PFQ_QID_NONE);
+ pfq_compiler_barrier();
+ dd->in_service_dispatched = 0;
+}
+
+/* Cancel the reservation when its queue receives new work. */
+static void pfq_stop_idle_hold(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ if (dd->waiting_idle && dd->idle_queue_qid == queue->qid)
+ pfq_clear_idle_hold(dd);
+}
+
+/**
+ * pfq_expire_idle_hold - Expire a SPECIAL read reservation
+ * @dd: per-disk scheduler state
+ * @now: boot time sampled before taking the disk lock
+ * @cleanup: prefetched queue and output detached service reference
+ *
+ * Caller holds dd->lock and has prefetched cleanup->queue. Revalidate the
+ * queue and current deadline; a stale or missing prefetch defers expiry.
+ * cleanup->service must be NULL on entry. Detached and prefetched references
+ * are released after the caller's final unlock.
+ *
+ * Context: Caller holds @dd->lock.
+ */
+static void pfq_expire_idle_hold(struct pfq_disk_data *dd, u64 now,
+ struct pfq_idle_cleanup *cleanup)
+{
+ struct pfq_queue_data *queue = cleanup->queue;
+ u32 qid;
+
+ if (!queue || !dd->waiting_idle ||
+ dd->idle_queue_qid != queue->qid || now < dd->idle_deadline_ns)
+ return;
+
+ qid = queue->qid;
+ pfq_clear_idle_hold(dd);
+ if (dd->in_service_qid == qid)
+ WRITE_ONCE(dd->in_service_qid, PFQ_QID_NONE);
+ if (!queue->nr_queued)
+ cleanup->service =
+ pfq_remove_queue_from_service_tree_locked(dd, queue);
+}
+
+/**
+ * pfq_choose_core - Choose between cached dispatch candidates
+ * @c1: first candidate, possibly NULL
+ * @c2: second candidate, possibly NULL
+ * @last_sector: end sector of the last fair-queue dispatch
+ *
+ * Prefer metadata, then smaller seek distance, then the lower sector.
+ *
+ * Context: Caller holds the disk lock.
+ *
+ * Return: Preferred non-owning candidate, or NULL if both candidates are
+ * NULL.
+ */
+static __always_inline struct pfq_rq_core *
+pfq_choose_core(struct pfq_rq_core *c1, struct pfq_rq_core *c2,
+ u64 last_sector)
+{
+ u64 d1, d2, s1, s2;
+
+ if (!c1 || c1 == c2)
+ return c2;
+ if (!c2)
+ return c1;
+
+ if (c1->is_meta && !c2->is_meta)
+ return c1;
+ if (!c1->is_meta && c2->is_meta)
+ return c2;
+
+ s1 = c1->sector;
+ s2 = c2->sector;
+ d1 = pfq_distance(s1, last_sector);
+ d2 = pfq_distance(s2, last_sector);
+ if (d1 < d2)
+ return c1;
+ if (d2 < d1)
+ return c2;
+ return s1 <= s2 ? c1 : c2;
+}
+
+/**
+ * pfq_check_fifo_impl - Check the earliest FIFO deadline
+ * @dd: per-disk scheduler state
+ * @queue: queue whose FIFO deadline is checked
+ * @last: node to reject if it is the earliest node, possibly NULL
+ * @now: boot time used to check the deadline
+ *
+ * Return the queue's earliest node if expired and different from last.
+ * The result is non-owning and valid only while dd->lock is held.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: Non-owning expired node, or NULL if absent, unexpired or equal
+ * to @last.
+ */
+static __noinline struct pfq_rq_core *
+pfq_check_fifo_impl(struct pfq_disk_data *dd, struct pfq_queue_data *queue,
+ struct pfq_rq_core *last, u64 now)
+{
+ struct bpf_rb_root *fifo;
+ struct pfq_rq_core *core;
+
+ fifo = pfq_fifo_tree(dd, queue->qid);
+ if (!fifo)
+ return NULL;
+
+ core = pfq_fifo_first_qid(fifo, queue->qid);
+ if (!core || core == last || now < core->fifo_deadline_ns)
+ return NULL;
+
+ return core;
+}
+
+/**
+ * pfq_check_fifo - Check for an expired FIFO request
+ * @dd: per-disk scheduler state
+ * @queue: queue whose FIFO deadline is checked
+ * @last: node to reject if it is the earliest node, possibly NULL
+ * @now: boot time used to check the deadline
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: Non-owning expired node, or NULL if absent, unexpired or equal
+ * to @last.
+ */
+static struct pfq_rq_core *pfq_check_fifo(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_rq_core *last, u64 now)
+{
+ return pfq_check_fifo_impl(dd, queue, last, now);
+}
+
+/* FIFO expiry wins; otherwise compare candidates on either side of the head. */
+static struct pfq_rq_core *pfq_update_next_rq(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ u64 now)
+{
+ struct pfq_rq_core *next_rn, *prev_rn, *fifo_rn;
+
+ fifo_rn = pfq_check_fifo(dd, queue, NULL, now);
+ if (fifo_rn)
+ return fifo_rn;
+
+ next_rn = pfq_rq_tree_next(dd, queue, dd->last_sector);
+ prev_rn = pfq_rq_tree_prev(dd, queue, dd->last_sector);
+ return pfq_choose_core(next_rn, prev_rn, dd->last_sector);
+}
+
+/**
+ * pfq_reposition_rq_in_tree - Reposition a request after a front merge
+ * @dd: per-disk scheduler state
+ * @queue: queue containing the merged request
+ * @rn: non-owning wrapper whose cached sector has been updated
+ *
+ * Reposition after a front merge; caller has updated rn->sector and holds
+ * dd->lock. Preserve the FIFO deadline and FIFO-tree reference. If re-add
+ * fails, the FIFO tree keeps rn alive and dequeue must handle that alone.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: true if reinserted, otherwise false.
+ */
+static bool pfq_reposition_rq_in_tree(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_rq_core *rn)
+{
+ struct bpf_rb_root *rq_root = pfq_rq_tree(dd, queue->qid);
+ struct pfq_rq_core *tree_core;
+
+ if (!rq_root)
+ return false;
+
+ tree_core = pfq_rq_tree_remove_take(rq_root, &rn->rb_node);
+ if (!tree_core)
+ return false;
+
+ if (bpf_rbtree_add(rq_root, &tree_core->rb_node, rq_less))
+ return false;
+ return true;
+}
+
+/*
+ * Caller holds dd->lock. A NULL queue adjusts only the disk total.
+ */
+static __always_inline void pfq_account_inc(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ if (queue)
+ queue->nr_queued++;
+ dd->nr_queued_total++;
+}
+
+static __always_inline void pfq_account_dec(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue)
+{
+ if (queue)
+ queue->nr_queued--;
+ dd->nr_queued_total--;
+}
+
+/* Owning insert references retained until the callback's unlocked exit. */
+struct pfq_insert_cleanup {
+ struct pfq_dispatch_node *dn;
+ struct pfq_queue_data *queue;
+ struct pfq_rq_core *list_ref;
+ struct pfq_rq_core *extra;
+ struct pfq_rq_core *core;
+ struct request *stale;
+ struct request *req;
+ struct pfq_stats *stats;
+};
+
+/*
+ * Release each owning reference separately, even when both refer to the
+ * same wrapper. list_n is the FIFO-tree reference, despite its name.
+ */
+static __noinline void pfq_drop_rq_container_refs(struct pfq_rq_core *list_n,
+ struct pfq_rq_core *tree_n)
+{
+ if (list_n)
+ bpf_obj_drop(list_n);
+ if (tree_n)
+ bpf_obj_drop(tree_n);
+}
+
+/* Caller has completed all state updates and released dd->lock. */
+static void pfq_insert_cleanup(struct pfq_insert_cleanup *cleanup)
+{
+ if (cleanup->req)
+ bpf_request_release(cleanup->req);
+ if (cleanup->stale)
+ bpf_request_release(cleanup->stale);
+ pfq_drop_rq_container_refs(cleanup->list_ref, cleanup->core);
+ if (cleanup->extra)
+ bpf_obj_drop(cleanup->extra);
+ if (cleanup->dn)
+ bpf_obj_drop(cleanup->dn);
+ if (cleanup->queue)
+ bpf_obj_drop(cleanup->queue);
+}
+
+/**
+ * pfq_finish_front_bio_merge - Update the front-merge survivor index
+ * @dd: per-disk scheduler state
+ * @queue: queue containing the surviving request
+ * @rn: non-owning wrapper of the surviving request
+ * @new_sector: request start sector after the bio merge
+ *
+ * Update the survivor's cached sector before repositioning under dd->lock.
+ *
+ * Context: Caller holds @dd->lock.
+ */
+static void pfq_finish_front_bio_merge(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_rq_core *rn,
+ u64 new_sector)
+{
+ rn->sector = new_sector;
+ (void)pfq_reposition_rq_in_tree(dd, queue, rn);
+}
+
+/**
+ * pfq_dequeue_rq_core_locked - Detach a request from both indexes
+ * @dd: per-disk scheduler state
+ * @queue: queue containing the request
+ * @rn: non-NULL, non-owning request wrapper
+ * @cleanup: output tree references to release after unlocking
+ *
+ * Detach rn and transfer its request reference without unlocking dd->lock.
+ * A failed reposition may have left only the FIFO-tree reference. Transfer
+ * each removed reference to cleanup for release after unlocking.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: Owning request reference, or NULL if no request can be taken.
+ */
+static __noinline struct request *
+pfq_dequeue_rq_core_locked(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_rq_core *rn,
+ struct pfq_dispatch_cleanup *cleanup)
+{
+ struct bpf_rb_root *rq_root = pfq_rq_tree(dd, queue->qid);
+ struct bpf_rb_root *fifo_root = pfq_fifo_tree(dd, queue->qid);
+ struct pfq_rq_core *list_n, *tree_n;
+ bool unlinked = false;
+ struct request *rq;
+
+ tree_n = pfq_rq_tree_remove_take(rq_root, &rn->rb_node);
+ if (tree_n)
+ list_n = pfq_fifo_tree_remove_take(fifo_root,
+ &tree_n->fifo_rb_node);
+ else
+ list_n = pfq_fifo_tree_remove_take(fifo_root, &rn->fifo_rb_node);
+
+ if (tree_n) {
+ rq = bpf_kptr_xchg(&tree_n->req, NULL);
+ unlinked = true;
+ } else if (list_n) {
+ rq = bpf_kptr_xchg(&list_n->req, NULL);
+ unlinked = true;
+ } else {
+ rq = NULL;
+ }
+
+ if (unlinked)
+ pfq_account_dec(dd, queue);
+ cleanup->rq_node = tree_n;
+ cleanup->fifo_node = list_n;
+ return rq;
+}
+
+/**
+ * pfq_undo_insert_activate_fail - Roll back a failed queue activation
+ * @dd: per-disk scheduler state
+ * @queue: queue whose activation failed
+ * @held: owning wrapper reference transferred to cleanup
+ * @rq: request whose non-owning private cache must be cleared
+ * @cleanup: output references to release at the callback exit
+ *
+ * Roll back failed activation under dd->lock. Consume held and transfer
+ * all owning references to cleanup for release at the callback's exit.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: -ENOMEM.
+ */
+static __noinline int
+pfq_undo_insert_activate_fail(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_rq_core *held,
+ struct request *rq,
+ struct pfq_insert_cleanup *cleanup)
+{
+ struct bpf_rb_root *rq_root = pfq_rq_tree(dd, queue->qid);
+ struct bpf_rb_root *fifo_root = pfq_fifo_tree(dd, queue->qid);
+ struct pfq_rq_core *list_n, *tree_n;
+ struct request *old;
+ bool unlinked = false;
+
+ pfq_rq_clear_node(rq);
+
+ /*
+ * Match pfq_dequeue_rq_core_locked control flow so clang does not emit
+ * pointer |= pointer for (tree_n || list_n).
+ */
+ tree_n = pfq_rq_tree_remove_take(rq_root, &held->rb_node);
+ if (tree_n)
+ list_n = pfq_fifo_tree_remove_take(fifo_root,
+ &tree_n->fifo_rb_node);
+ else
+ list_n = pfq_fifo_tree_remove_take(fifo_root,
+ &held->fifo_rb_node);
+
+ if (tree_n) {
+ old = bpf_kptr_xchg(&tree_n->req, NULL);
+ unlinked = true;
+ } else if (list_n) {
+ old = bpf_kptr_xchg(&list_n->req, NULL);
+ unlinked = true;
+ } else {
+ old = bpf_kptr_xchg(&held->req, NULL);
+ }
+
+ if (unlinked)
+ pfq_account_dec(dd, queue);
+ cleanup->req = old;
+ cleanup->list_ref = list_n;
+ cleanup->core = tree_n;
+ cleanup->extra = held;
+ return -ENOMEM;
+}
+
+/**
+ * pfq_insert_core_locked - Insert a request into both indexes
+ * @dd: per-disk scheduler state
+ * @queue: target logical queue
+ * @core: owning wrapper reference consumed by insertion or cleanup
+ * @rq: incoming request with its non-owning private cache set
+ * @cleanup: output temporary or detached references for cleanup
+ *
+ * Each request index consumes one owning reference. Keep dd->lock held
+ * through insertion and activation; return temporary or detached references
+ * through cleanup on either success or failure.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: 0 on success, or a negative errno after preparing deferred
+ * cleanup.
+ */
+static __noinline int
+pfq_insert_core_locked(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_rq_core *core,
+ struct request *rq,
+ struct pfq_insert_cleanup *cleanup)
+{
+ struct pfq_rq_core *list_ref, *list_n, *tree_n, *held = NULL;
+ struct bpf_rb_root *rq_root = pfq_rq_tree(dd, queue->qid);
+ struct bpf_rb_root *fifo_root = pfq_fifo_tree(dd, queue->qid);
+ struct request *old;
+
+ list_ref = bpf_refcount_acquire(core);
+ if (!list_ref) {
+ old = bpf_kptr_xchg(&core->req, NULL);
+ pfq_rq_clear_node(rq);
+ cleanup->req = old;
+ cleanup->core = core;
+ return -ENOMEM;
+ }
+
+ if (!rq_root || bpf_rbtree_add(rq_root, &core->rb_node, rq_less)) {
+ old = bpf_kptr_xchg(&list_ref->req, NULL);
+ pfq_rq_clear_node(rq);
+ cleanup->req = old;
+ /* Failed rbtree_add consumes core; a missing root does not. */
+ if (!rq_root)
+ cleanup->core = core;
+ cleanup->list_ref = list_ref;
+ return -EINVAL;
+ }
+ if (!fifo_root ||
+ bpf_rbtree_add(fifo_root, &list_ref->fifo_rb_node, fifo_less)) {
+ pfq_rq_clear_node(rq);
+ tree_n = pfq_rq_tree_remove_take(rq_root, &core->rb_node);
+ old = tree_n ? bpf_kptr_xchg(&tree_n->req, NULL) : NULL;
+ cleanup->req = old;
+ cleanup->core = tree_n;
+ /* A failed add consumes list_ref; a missing root does not. */
+ if (!fifo_root)
+ cleanup->list_ref = list_ref;
+ return -EINVAL;
+ }
+
+ /*
+ * Keep an owning reference for activation rollback or final cleanup.
+ * Successful insertion retains both tree references independently.
+ */
+ held = bpf_refcount_acquire(core);
+ if (!held) {
+ pfq_rq_clear_node(rq);
+ tree_n = pfq_rq_tree_remove_take(rq_root, &core->rb_node);
+ list_n = tree_n ?
+ pfq_fifo_tree_remove_take(fifo_root,
+ &tree_n->fifo_rb_node) :
+ NULL;
+ old = tree_n ? bpf_kptr_xchg(&tree_n->req, NULL) : NULL;
+ cleanup->req = old;
+ cleanup->list_ref = list_n;
+ cleanup->core = tree_n;
+ return -ENOMEM;
+ }
+
+ pfq_account_inc(dd, queue);
+
+ if (queue->in_tree) {
+ pfq_stop_idle_hold(dd, queue);
+ goto success_defer_held;
+ }
+ if (pfq_activate_queue(dd, queue))
+ goto success_defer_held;
+
+ return pfq_undo_insert_activate_fail(dd, queue, held, rq, cleanup);
+
+success_defer_held:
+ cleanup->extra = held;
+ return 0;
+}
+
+/**
+ * pfq_update_queue_vtime - Charge virtual time for a dispatch
+ * @dd: disk supplying the virtual-time floor
+ * @queue: queue whose virtual time is charged
+ * @sectors: number of sectors dispatched
+ *
+ * Context: Caller holds @dd->lock.
+ */
+static void pfq_update_queue_vtime(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ u32 sectors)
+{
+ u64 delta;
+
+ if (queue->vtime < dd->vtime)
+ queue->vtime = dd->vtime;
+ delta = ((u64)sectors << PFQ_VTIME_SHIFT) / queue->weight;
+ /* Integer truncation must not give small requests free service. */
+ if (sectors && !delta)
+ delta = 1;
+ queue->vtime += delta;
+}
+
+/* Begin a new in-service batch on queue @qid. */
+static void pfq_set_in_service(struct pfq_disk_data *dd, u32 qid)
+{
+ WRITE_ONCE(dd->in_service_qid, qid);
+ dd->in_service_dispatched = 0;
+}
+
+/*
+ * Keep an empty SPECIAL+READ queue selected until its deadline.
+ * @random and @now are sampled before dd->lock; never drop the lock between
+ * checking idle eligibility and publishing the hold.
+ */
+static void pfq_start_idle_hold(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue, u64 now, u32 random,
+ struct pfq_stats *stats)
+{
+ u64 range, extra = 0;
+
+ if (!queue->in_tree || queue->nr_queued ||
+ dd->in_service_qid != queue->qid || dd->waiting_idle)
+ return;
+
+ /* Use u64 so the inclusive range cannot wrap for UINT_MAX tunables. */
+ if (dd->idle_delay_max_ms >= dd->idle_delay_min_ms) {
+ range = (u64)dd->idle_delay_max_ms - dd->idle_delay_min_ms + 1;
+ extra = random % range;
+ }
+
+ WRITE_ONCE(dd->idle_queue_qid, queue->qid);
+ dd->waiting_idle = true;
+ dd->idle_deadline_ns = now +
+ ((u64)dd->idle_delay_min_ms + extra) * 1000000ULL;
+ pfq_stat_add(stats, PFQ_STAT_IDLE_HOLD, 1);
+}
+
+/*
+ * Return an owning queue reference under dd->lock. Revalidate the prefetched
+ * current queue and use the latest batch state, including a new batch on
+ * the same qid. A stale prefetch returns NULL without dequeuing. Transfer
+ * removed service references to cleanup for release after unlocking.
+ */
+static struct pfq_queue_data *
+pfq_select_next_queue_locked(struct pfq_disk_data *dd,
+ struct pfq_queue_data *current,
+ u64 now, u32 random,
+ struct pfq_dispatch_cleanup *cleanup)
+{
+ struct pfq_queue_data *queue;
+ struct bpf_rb_node *node;
+
+ if (current) {
+ if (dd->in_service_qid != current->qid)
+ return NULL;
+ } else if (pfq_qid_valid(dd->in_service_qid)) {
+ return NULL;
+ }
+ if (!dd->nr_active)
+ return NULL;
+
+ if (current) {
+ if (current->nr_queued > 0) {
+ if (dd->in_service_dispatched <
+ pfq_max_batch(dd, current->io_type))
+ return bpf_refcount_acquire(current);
+ if (current->in_tree)
+ pfq_requeue_in_service_tree(dd, current);
+ } else {
+ if (current->pq_class == PFQ_SPECIAL_CLASS &&
+ current->io_type == PFQ_READ)
+ pfq_start_idle_hold(dd, current, now, random,
+ cleanup->stats);
+ if (dd->waiting_idle &&
+ dd->idle_queue_qid == current->qid)
+ return bpf_refcount_acquire(current);
+ cleanup->old_service =
+ pfq_remove_queue_from_service_tree_locked(dd, current);
+ }
+ WRITE_ONCE(dd->in_service_qid, PFQ_QID_NONE);
+ }
+
+ node = bpf_rbtree_first(&dd->service_tree);
+ if (!node)
+ return NULL;
+ queue = container_of(node, struct pfq_queue_data, svc_rb_node);
+ pfq_set_in_service(dd, queue->qid);
+ return bpf_refcount_acquire(queue);
+}
+
+/* Only the INITIALIZING owner may initialize queue fields. */
+static int pfq_init_queue_fields(struct pfq_disk_data *dd, u32 idx)
+{
+ struct pfq_queue_data *queue = pfq_queue_scalar(dd, idx);
+
+ if (!queue)
+ return -ENOENT;
+
+ queue->qid = idx;
+ queue->pq_class = PFQ_GET_CLASS(idx);
+ queue->priority = PFQ_GET_PRIO(idx);
+ queue->io_type = PFQ_GET_TYPE(idx);
+ queue->nr_queued = 0;
+ queue->vtime = 0;
+ queue->weight = pfq_calc_weight(dd, queue);
+ queue->in_tree = false;
+ return 0;
+}
+
+/**
+ * pfq_disk_init - Reset disk scheduling state
+ * @dd: unpublished disk state with tunables already applied
+ *
+ * Reset scheduling state under dd->lock before publishing initialized queues.
+ * Tunables have already been applied; event counters are managed separately.
+ *
+ * Context: Caller holds @dd->lock.
+ */
+static void pfq_disk_init(struct pfq_disk_data *dd)
+{
+ dd->vtime = 0;
+ dd->min_vtime = 0;
+ dd->last_sector = 0;
+ pfq_disk_reset_counters(dd);
+ dd->waiting_idle = false;
+ WRITE_ONCE(dd->idle_queue_qid, PFQ_QID_NONE);
+ dd->idle_deadline_ns = 0;
+ dd->fifo_expire_ns[PFQ_READ] = PFQ_READ_EXPIRE_NS;
+ dd->fifo_expire_ns[PFQ_SYNC_WRITE] = PFQ_SYNC_WRITE_EXPIRE_NS;
+ dd->fifo_expire_ns[PFQ_WRITE] = PFQ_WRITE_EXPIRE_NS;
+}
+
+/**
+ * pfq_do_init_sched - Initialize disk state and logical queues
+ * @q: request queue being initialized
+ *
+ * Initialize at elevator setup or lazily at insert after BPF attach.
+ *
+ * Claim UNINITIALIZED -> INITIALIZING under dd->lock before creating queues.
+ * The owner stays INITIALIZING across unlocks, including failure cleanup.
+ * Concurrent callers return -EAGAIN while initialization or cleanup is in
+ * progress. Publish INITIALIZED on success or UNINITIALIZED after cleanup.
+ *
+ * Context: Caller must not hold a BPF spin lock.
+ *
+ * Return: 0 if ready, -EAGAIN if initialization is in progress, or another
+ * negative errno on failure.
+ */
+static int pfq_do_init_sched(struct request_queue *q)
+{
+ struct pfq_tunables tun = {}, *map_tun;
+ u32 tun_key = PFQ_TUNABLE_KEY, i;
+ int id = q->id, ret, state;
+ struct pfq_disk_data *dd;
+ bool stats_created;
+
+ dd = pfq_disk_lookup_raw(id);
+ if (!dd) {
+ ret = pfq_disk_create(id);
+ if (ret && ret != -EEXIST)
+ return ret;
+ dd = pfq_disk_lookup_raw(id);
+ if (!dd)
+ return -ENOENT;
+ }
+
+ bpf_spin_lock(&dd->lock);
+ state = pfq_init_state_load(dd);
+ if (state != PFQ_UNINITIALIZED) {
+ ret = state == PFQ_INITIALIZED ? 0 : -EAGAIN;
+ bpf_spin_unlock(&dd->lock);
+ return ret;
+ }
+ pfq_init_state_store(dd, PFQ_INITIALIZING);
+ bpf_spin_unlock(&dd->lock);
+
+ ret = pfq_disk_create_queues(id);
+ if (ret < 0)
+ goto reset_state;
+ stats_created = ret > 0;
+
+ map_tun = bpf_map_lookup_elem(&pfq_tunables, &tun_key);
+ if (map_tun)
+ tun = *map_tun;
+
+ bpf_spin_lock(&dd->lock);
+ pfq_apply_tunables(dd, &tun);
+ pfq_disk_init(dd);
+ for (i = 0; i < PFQ_TOTAL_QUEUES; i++) {
+ ret = pfq_init_queue_fields(dd, i);
+ if (ret)
+ goto err_queues;
+ }
+ pfq_init_state_store(dd, PFQ_INITIALIZED);
+ bpf_spin_unlock(&dd->lock);
+
+ return 0;
+
+err_queues:
+ bpf_spin_unlock(&dd->lock);
+ pfq_queues_destroy(id);
+ if (stats_created)
+ bpf_map_delete_elem(&pfq_stats_map, &id);
+reset_state:
+ /* Complete cleanup before another caller can claim initialization. */
+ pfq_init_state_store(dd, PFQ_UNINITIALIZED);
+ return ret;
+}
+
+/*
+ * pfq_init_sched - Initialize a disk using the attached policy
+ * @q: request queue switching to UFQ
+ *
+ * Return: 0 on success, or a negative errno.
+ */
+int BPF_STRUCT_OPS(pfq_init_sched, struct request_queue *q)
+{
+ return pfq_do_init_sched(q);
+}
+
+/* UFQ must quiesce this disk's callbacks before removing its state. */
+int BPF_STRUCT_OPS(pfq_exit_sched, struct request_queue *q)
+{
+ struct pfq_disk_data *dd;
+ int id = q->id;
+
+ dd = pfq_disk_lookup_raw(id);
+ if (dd) {
+ pfq_init_state_store(dd, PFQ_UNINITIALIZED);
+ pfq_queues_destroy(id);
+ }
+
+ bpf_map_delete_elem(&pfq_stats_map, &id);
+ /* The value destructor releases graph nodes and request kptrs. */
+ bpf_map_delete_elem(&pfq_map, &id);
+ return 0;
+}
+
+/*
+ * pfq_has_req - Report requests awaiting dispatch
+ * @q: request queue being queried
+ * @rqs_count: UFQ pending request count used if disk state is unavailable
+ *
+ * Report queued work, excluding empty idle reservations. If disk state is
+ * unavailable, fall back to UFQ's count of requests awaiting dispatch.
+ * Idle expiry is checked here; UFQ does not arrange an idle-expiry timer.
+ *
+ * Return: true if requests await dispatch, otherwise false.
+ */
+bool BPF_STRUCT_OPS(pfq_has_req, struct request_queue *q, int rqs_count)
+{
+ struct pfq_idle_cleanup idle_cleanup = {};
+ struct pfq_disk_data *dd;
+ bool has = false;
+ int id = q->id;
+ u64 now;
+
+ dd = pfq_disk_lookup(id);
+ if (!dd)
+ return rqs_count > 0;
+
+ now = pfq_now_ns();
+ idle_cleanup.queue = pfq_queue_acquire(id, READ_ONCE(dd->idle_queue_qid));
+ bpf_spin_lock(&dd->lock);
+ pfq_expire_idle_hold(dd, now, &idle_cleanup);
+ has = !bpf_list_empty(&dd->dispatch) || dd->nr_queued_total > 0;
+ bpf_spin_unlock(&dd->lock);
+ pfq_idle_cleanup(&idle_cleanup);
+ return has;
+}
+
+/**
+ * pfq_insert_normal - Insert a request into a fair queue
+ * @dd: per-disk scheduler state
+ * @rq: incoming request
+ * @qid: queue index resolved before taking the disk lock
+ * @now: boot time sampled before taking the disk lock
+ * @cleanup: prefetched statistics and output references for cleanup
+ *
+ * Enter and leave unlocked. Prepare references before taking dd->lock, then
+ * hold it through tree insertion, activation and accounting. Return all
+ * temporary or detached references in cleanup for the callback's exit.
+ *
+ * Context: Caller must not hold @dd->lock.
+ *
+ * Return: 0 on success, or a negative errno.
+ */
+static int pfq_insert_normal(struct pfq_disk_data *dd, struct request *rq,
+ u32 qid, u64 now, struct pfq_insert_cleanup *cleanup)
+{
+ struct pfq_queue_data *queue;
+ struct pfq_rq_core *core;
+ struct request *acquired;
+ int link_ret;
+
+ queue = pfq_queue_acquire(dd->disk_id, qid);
+ if (!queue)
+ return -EINVAL;
+ cleanup->queue = queue;
+
+ core = bpf_obj_new(typeof(*core));
+ if (!core)
+ return -ENOMEM;
+
+ acquired = bpf_request_acquire(rq);
+ if (!acquired) {
+ cleanup->core = core;
+ return -EPERM;
+ }
+
+ /* The probe-read helper requires an unlocked context. */
+ pfq_rq_set_core(rq, core);
+
+ /* jiffies helpers forbidden under bpf_spin_lock. */
+ unsigned long jnow = (unsigned long)bpf_jiffies64();
+ unsigned long ft = (unsigned long)rq->fifo_time;
+ u64 fifo_deadline_ns;
+
+ if (ft > jnow)
+ fifo_deadline_ns = pfq_jiffies_to_deadline_ns(ft, now);
+ else
+ /* dd->fifo_expire_ns read unlocked; stable after disk_init. */
+ fifo_deadline_ns = now + pfq_fifo_expire_ns(dd, queue->io_type);
+
+ bpf_spin_lock(&dd->lock);
+
+ cleanup->stale = bpf_kptr_xchg(&core->req, acquired);
+
+ core->qid = queue->qid;
+ core->sector = rq->__sector;
+ core->is_meta = !!(rq->cmd_flags & PFQ_REQ_META);
+ core->fifo_deadline_ns = fifo_deadline_ns;
+
+ link_ret = pfq_insert_core_locked(dd, queue, core, rq, cleanup);
+ if (!link_ret) {
+ pfq_stat_add(cleanup->stats, PFQ_STAT_INSERT_CNT, 1);
+ pfq_stat_add(cleanup->stats, PFQ_STAT_INSERT_SIZE, rq->__data_len);
+ }
+ bpf_spin_unlock(&dd->lock);
+ return link_ret;
+}
+
+/*
+ * pfq_insert_req - Insert a request using the UFQ policy
+ * @q: request queue receiving the request
+ * @rq: incoming request
+ * @flags: insertion flags; BLK_MQ_INSERT_AT_HEAD selects urgent dispatch
+ *
+ * Head inserts use the urgent dispatch list; other inserts use fair queues.
+ * Retain idle-expiration references until the final callback cleanup.
+ *
+ * Return: 0 on success, or a negative errno for UFQ insert-error handling.
+ */
+int BPF_STRUCT_OPS(pfq_insert_req, struct request_queue *q,
+ struct request *rq, blk_insert_t flags)
+{
+ struct pfq_idle_cleanup idle_cleanup = {};
+ struct pfq_insert_cleanup cleanup = {};
+ struct pfq_dispatch_node *dn;
+ struct pfq_disk_data *dd;
+ struct request *acquired;
+ struct pfq_stats *stats;
+ int id = q->id, ret = 0;
+ bool interactive;
+ u32 qid = 0;
+ u64 now;
+
+ dd = pfq_disk_lookup(id);
+ if (!dd) {
+ /* Initialize after attach or retry an earlier failure. */
+ ret = pfq_do_init_sched(q);
+ if (ret)
+ return ret;
+ dd = pfq_disk_lookup(id);
+ if (!dd)
+ return -EINVAL;
+ }
+
+ interactive = pfq_current_is_interactive();
+ if (!(flags & BLK_MQ_INSERT_AT_HEAD))
+ qid = pfq_rq_qid(rq, interactive);
+ now = pfq_now_ns();
+ stats = pfq_stats_lookup(dd->disk_id);
+ cleanup.stats = stats;
+ idle_cleanup.queue = pfq_queue_acquire(id, READ_ONCE(dd->idle_queue_qid));
+ bpf_spin_lock(&dd->lock);
+ pfq_expire_idle_hold(dd, now, &idle_cleanup);
+ bpf_spin_unlock(&dd->lock);
+
+ if (flags & BLK_MQ_INSERT_AT_HEAD) {
+ dn = bpf_obj_new(typeof(*dn));
+ if (!dn) {
+ ret = -ENOMEM;
+ goto out;
+ }
+ acquired = bpf_request_acquire(rq);
+ if (!acquired) {
+ cleanup.dn = dn;
+ ret = -EPERM;
+ goto out;
+ }
+ bpf_spin_lock(&dd->lock);
+ cleanup.stale = bpf_kptr_xchg(&dn->req, acquired);
+ ret = pfq_dispatch_push_dn(&dd->dispatch, dn);
+ if (!ret) {
+ pfq_stat_add(stats, PFQ_STAT_AT_HEAD_CNT, 1);
+ pfq_stat_add(stats, PFQ_STAT_AT_HEAD_SIZE, rq->__data_len);
+ }
+ bpf_spin_unlock(&dd->lock);
+ } else {
+ ret = pfq_insert_normal(dd, rq, qid, now, &cleanup);
+ if (!ret && interactive)
+ pfq_stat_add(stats, PFQ_STAT_INTERACTIVE, 1);
+ }
+out:
+ if (ret)
+ pfq_stat_add(stats, PFQ_STAT_INSERT_ERR, 1);
+ pfq_idle_cleanup(&idle_cleanup);
+ pfq_insert_cleanup(&cleanup);
+ return ret;
+}
+
+/*
+ * Commit post-dispatch service state without unlocking. Return any removed
+ * service-tree ref for deferred release by the dispatch caller.
+ */
+static __always_inline struct pfq_queue_data *
+pfq_reposition_after_dispatch_locked(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ bool force_requeue,
+ u64 now, u32 random,
+ struct pfq_stats *stats)
+{
+ struct pfq_queue_data *removed;
+
+ if (queue->nr_queued > 0) {
+ if (force_requeue)
+ pfq_requeue_in_service_tree(dd, queue);
+ return NULL;
+ }
+ if (queue->pq_class == PFQ_SPECIAL_CLASS && queue->io_type == PFQ_READ &&
+ dd->in_service_qid == queue->qid && queue->in_tree) {
+ pfq_requeue_in_service_tree(dd, queue);
+ pfq_start_idle_hold(dd, queue, now, random, stats);
+ if (dd->waiting_idle && dd->idle_queue_qid == queue->qid)
+ return NULL;
+ }
+ removed = pfq_remove_queue_from_service_tree_locked(dd, queue);
+ if (dd->in_service_qid == queue->qid)
+ WRITE_ONCE(dd->in_service_qid, PFQ_QID_NONE);
+ return removed;
+}
+
+/**
+ * pfq_dispatch_rq_from_queue_locked - Dequeue and charge a selected request
+ * @dd: per-disk scheduler state
+ * @queue: selected logical queue
+ * @now: boot time used for FIFO expiry
+ * @cleanup: output detached tree references for deferred release
+ *
+ * Pick, dequeue and charge virtual time with dd->lock held throughout.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: Owning request reference, or NULL if no request can be
+ * dispatched.
+ */
+static struct request *
+pfq_dispatch_rq_from_queue_locked(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue, u64 now,
+ struct pfq_dispatch_cleanup *cleanup)
+{
+ struct pfq_rq_core *rn;
+ struct request *rq;
+ u32 sectors;
+
+ rn = pfq_update_next_rq(dd, queue, now);
+ if (!rn)
+ return NULL;
+
+ rq = pfq_dequeue_rq_core_locked(dd, queue, rn, cleanup);
+ if (!rq)
+ return NULL;
+
+ pfq_rq_bind_qid(rq, queue->qid);
+
+ sectors = rq->__data_len >> SECTOR_SHIFT;
+ dd->last_sector = rq->__sector + sectors;
+ pfq_update_queue_vtime(dd, queue, sectors);
+
+ return rq;
+}
+
+/*
+ * pfq_dispatch_req - Dispatch the next UFQ request
+ * @q: request queue being dispatched
+ *
+ * Prefetch the current queue and acquire newly selected queues from the
+ * service tree. Idle expiry, selection, dequeue and accounting share one
+ * lock hold; release detached and temporary references after unlocking.
+ *
+ * Return: Owning request reference, or NULL if no request is selected.
+ */
+struct request *BPF_STRUCT_OPS(pfq_dispatch_req, struct request_queue *q)
+{
+ struct pfq_queue_data *current, *queue = NULL;
+ struct pfq_idle_cleanup idle_cleanup = {};
+ struct pfq_dispatch_cleanup cleanup = {};
+ struct pfq_dispatch_node *dn = NULL;
+ bool at_head = false, more_batch;
+ struct request *rq = NULL;
+ struct pfq_disk_data *dd;
+ struct bpf_list_node *ln;
+ struct pfq_stats *stats;
+ int id = q->id;
+ u32 random, qid;
+ u64 now;
+
+ dd = pfq_disk_lookup(id);
+ if (!dd)
+ return NULL;
+
+ now = pfq_now_ns();
+ random = bpf_get_prandom_u32();
+ stats = pfq_stats_lookup(dd->disk_id);
+ qid = READ_ONCE(dd->in_service_qid);
+ current = pfq_queue_acquire(id, qid);
+ /* The idle hold belongs to current; share its prefetched reference. */
+ idle_cleanup.queue = current;
+
+ bpf_spin_lock(&dd->lock);
+ pfq_expire_idle_hold(dd, now, &idle_cleanup);
+ ln = bpf_list_pop_front(&dd->dispatch);
+ if (ln) {
+ dn = container_of(ln, struct pfq_dispatch_node, node);
+ rq = bpf_kptr_xchg(&dn->req, NULL);
+ at_head = true;
+ goto out;
+ }
+
+ cleanup.stats = stats;
+ /* Expiration may have cleared the current slot. */
+ queue = pfq_select_next_queue_locked(dd,
+ dd->in_service_qid == PFQ_QID_NONE ? NULL : current,
+ now, random, &cleanup);
+ if (!queue)
+ goto out;
+
+ rq = pfq_dispatch_rq_from_queue_locked(dd, queue, now, &cleanup);
+ if (!rq)
+ goto out;
+
+ dd->in_service_dispatched++;
+ /* Keep the queue in service until its batch budget is exhausted. */
+ more_batch = pfq_queue_nq_get(queue) > 0 &&
+ dd->in_service_dispatched <
+ pfq_max_batch(dd, queue->io_type);
+ cleanup.empty_service =
+ pfq_reposition_after_dispatch_locked(dd, queue, !more_batch,
+ now, random, stats);
+
+ /* Advance global vtime floor to lagging active queues. */
+ if (dd->nr_active > 0 && dd->vtime != dd->min_vtime)
+ dd->vtime = dd->min_vtime;
+
+out:
+ if (rq) {
+ if (at_head) {
+ pfq_stat_add(stats, PFQ_STAT_DISPATCH_AT_HEAD_CNT, 1);
+ pfq_stat_add(stats, PFQ_STAT_DISPATCH_AT_HEAD_SIZE,
+ rq->__data_len);
+ } else {
+ pfq_stat_add(stats, PFQ_STAT_DISPATCH_CNT, 1);
+ pfq_stat_add(stats, PFQ_STAT_DISPATCH_SIZE, rq->__data_len);
+ }
+ }
+ bpf_spin_unlock(&dd->lock);
+ pfq_idle_cleanup(&idle_cleanup);
+ pfq_drop_rq_container_refs(cleanup.fifo_node, cleanup.rq_node);
+ if (cleanup.old_service)
+ bpf_obj_drop(cleanup.old_service);
+ if (cleanup.empty_service)
+ bpf_obj_drop(cleanup.empty_service);
+ if (queue)
+ bpf_obj_drop(queue);
+ if (dn)
+ bpf_obj_drop(dn);
+ return rq;
+}
+
+/*
+ * pfq_finish_req - Account completion and consider idle hold
+ * @rq: completed request carrying its dispatched queue identity
+ *
+ * Account completions for initialized disks. The qid bound at dispatch
+ * identifies whether an empty SPECIAL+READ queue still owns the service
+ * slot and may start an idle hold.
+ */
+void BPF_STRUCT_OPS(pfq_finish_req, struct request *rq)
+{
+ struct pfq_queue_data *queue = NULL;
+ struct pfq_disk_data *dd;
+ struct pfq_stats *stats;
+ u32 rq_qid, random;
+ u64 now;
+ int id;
+
+ if (!rq || !rq->q)
+ return;
+
+ id = rq->q->id;
+ dd = pfq_disk_lookup(id);
+ if (!dd)
+ return;
+
+ rq_qid = pfq_rq_bound_qid(rq);
+ now = pfq_now_ns();
+ random = bpf_get_prandom_u32();
+ stats = pfq_stats_lookup(dd->disk_id);
+ queue = pfq_queue_acquire(id, rq_qid);
+ bpf_spin_lock(&dd->lock);
+
+ pfq_stat_add(stats, PFQ_STAT_FINISH_CNT, 1);
+ pfq_stat_add(stats, PFQ_STAT_FINISH_SIZE,
+ (u64)rq->stats_sectors << SECTOR_SHIFT);
+
+ /* The service slot may have changed while acquiring the bound queue. */
+ if (!queue || rq_qid != dd->in_service_qid)
+ goto out;
+
+ {
+ bool may_idle = false;
+
+ if (queue->pq_class == PFQ_SPECIAL_CLASS &&
+ queue->io_type == PFQ_READ &&
+ queue->in_tree && !queue->nr_queued)
+ may_idle = true;
+ if (may_idle && !dd->waiting_idle)
+ pfq_start_idle_hold(dd, queue, now, random, stats);
+ }
+
+out:
+ pfq_rq_clear_node(rq);
+ bpf_spin_unlock(&dd->lock);
+ if (queue)
+ bpf_obj_drop(queue);
+}
+
+/* Group arguments to stay within the BPF subprogram limit of five. */
+struct pfq_adj_find_ctl {
+ sector_t probe_start;
+ sector_t probe_end;
+ struct pfq_rq_core **out_rn;
+ struct request **out_stale;
+ struct request **out_req;
+};
+
+/**
+ * pfq_rq_tree_find_adjacent - Find a queued sector-adjacent request
+ * @dd: per-disk scheduler state
+ * @queue: logical queue to search
+ * @ctl: probe sector range and output wrapper/request references
+ *
+ * Find an adjacent request within queue's qid with dd->lock held. Return
+ * its owning request reference in ctl->out_req and a non-owning wrapper in
+ * ctl->out_rn. The caller must restore the request or commit the merge
+ * before unlocking. Unexpected displaced references go to ctl->out_stale
+ * for release after unlocking.
+ *
+ * Context: Caller holds @dd->lock.
+ *
+ * Return: ELEVATOR_FRONT_MERGE, ELEVATOR_BACK_MERGE or ELEVATOR_NO_MERGE.
+ */
+static enum elv_merge
+pfq_rq_tree_find_adjacent(struct pfq_disk_data *dd,
+ struct pfq_queue_data *queue,
+ struct pfq_adj_find_ctl *ctl)
+{
+ struct bpf_rb_root *rq_root = pfq_rq_tree(dd, queue->qid);
+ sector_t probe_start, probe_end, cand_start, cand_end;
+ enum elv_merge mt = ELEVATOR_NO_MERGE;
+ struct request *req, *stale = NULL;
+ struct bpf_rb_node *node;
+ struct pfq_rq_core *rn;
+ int count = 0;
+
+ if (!ctl)
+ return ELEVATOR_NO_MERGE;
+
+ probe_start = ctl->probe_start;
+ probe_end = ctl->probe_end;
+ *ctl->out_rn = NULL;
+ *ctl->out_req = NULL;
+ *ctl->out_stale = NULL;
+
+ if (!rq_root)
+ return ELEVATOR_NO_MERGE;
+
+ node = bpf_rbtree_root(rq_root);
+ while (node && count++ < PFQ_LOOP_MAX) {
+ rn = container_of(node, struct pfq_rq_core, rb_node);
+ if (rn->qid < queue->qid) {
+ node = bpf_rbtree_right(rq_root, node);
+ continue;
+ }
+ if (rn->qid > queue->qid) {
+ node = bpf_rbtree_left(rq_root, node);
+ continue;
+ }
+
+ req = bpf_kptr_xchg(&rn->req, NULL);
+ if (!req)
+ break;
+
+ cand_start = req->__sector;
+ cand_end = cand_start + (req->__data_len >> SECTOR_SHIFT);
+
+ if (probe_end < cand_start) {
+ node = bpf_rbtree_left(rq_root, node);
+ } else if (probe_start > cand_end) {
+ node = bpf_rbtree_right(rq_root, node);
+ } else if (cand_start == probe_end) {
+ *ctl->out_rn = rn;
+ *ctl->out_req = req;
+ mt = ELEVATOR_FRONT_MERGE;
+ break;
+ } else if (cand_end == probe_start) {
+ *ctl->out_rn = rn;
+ *ctl->out_req = req;
+ mt = ELEVATOR_BACK_MERGE;
+ break;
+ } else {
+ stale = bpf_kptr_xchg(&rn->req, req);
+ break;
+ }
+
+ stale = bpf_kptr_xchg(&rn->req, req);
+ if (stale)
+ break;
+ }
+
+ if (stale) {
+ *ctl->out_stale = stale;
+ return ELEVATOR_NO_MERGE;
+ }
+ return mt;
+}
+
+/*
+ * Reject clearly unmergeable requests before unlinking a candidate.
+ * The kernel performs the remaining merge checks after the callback returns.
+ */
+static __always_inline bool pfq_rq_mergeable_one(struct request *rq)
+{
+ u32 op = rq->cmd_flags & REQ_OP_MASK;
+
+ if (op == PFQ_REQ_OP_DRV_IN || op == PFQ_REQ_OP_DRV_OUT)
+ return false;
+ if (op == PFQ_REQ_OP_FLUSH || op == PFQ_REQ_OP_WRITE_ZEROES ||
+ op == PFQ_REQ_OP_ZONE_APPEND)
+ return false;
+ if (rq->cmd_flags & PFQ_REQ_NOMERGE_FLAGS)
+ return false;
+ if (rq->rq_flags & PFQ_RQF_NOMERGE_FLAGS)
+ return false;
+ return true;
+}
+
+/**
+ * pfq_rq_merge_attrs_ok - Precheck request merge attributes
+ * @a: first merge candidate
+ * @b: second merge candidate
+ *
+ * Check a subset of attempt_merge() attributes before unlinking. The caller
+ * checks mq_ctx and mq_hctx; the kernel still validates the full merge.
+ *
+ * Context: Caller holds the disk lock.
+ *
+ * Return: true if the checked attributes match, otherwise false.
+ */
+static __always_inline bool pfq_rq_merge_attrs_ok(struct request *a,
+ struct request *b)
+{
+ u32 opa, opb;
+
+ if (!pfq_rq_mergeable_one(a) || !pfq_rq_mergeable_one(b))
+ return false;
+
+ opa = a->cmd_flags & REQ_OP_MASK;
+ opb = b->cmd_flags & REQ_OP_MASK;
+ if (opa != opb)
+ return false;
+ /* data direction: LSB of op (same as rq_data_dir). */
+ if ((opa & 1) != (opb & 1))
+ return false;
+ if (!a->bio || !b->bio)
+ return false;
+ if (a->bio->bi_write_hint != b->bio->bi_write_hint)
+ return false;
+ if (a->bio->bi_ioprio != b->bio->bi_ioprio)
+ return false;
+ return true;
+}
+
+/*
+ * pfq_merge_req - Detach an adjacent request for merging
+ * @q: request queue containing the candidates
+ * @rq: incoming request to merge
+ * @type: output merge type; ELEVATOR_NO_MERGE when no candidate is
+ * returned
+ *
+ * Return an owning reference to an adjacent candidate and set the merge type.
+ *
+ * Unlink the candidate from PFQ; UFQ attempts the merge and reinserts the
+ * candidate if it fails. For a back merge, incoming rq is absorbed into the
+ * candidate; for a front merge, the candidate is absorbed into rq.
+ *
+ * Return: Owning reference to the detached candidate, or NULL.
+ */
+struct request *BPF_STRUCT_OPS(pfq_merge_req, struct request_queue *q,
+ struct request *rq, int *type)
+{
+ struct pfq_rq_core *tree_core, *list_n, *rn;
+ struct request *stale = NULL, *targ = NULL;
+ enum elv_merge mt = ELEVATOR_NO_MERGE;
+ struct pfq_queue_data *queue;
+ sector_t rq_start, rq_end;
+ struct pfq_disk_data *dd;
+ struct pfq_stats *stats;
+ bool interactive;
+ u32 qid;
+
+ *type = ELEVATOR_NO_MERGE;
+ if (!pfq_rq_mergeable_one(rq))
+ return NULL;
+
+ dd = pfq_disk_lookup(q->id);
+ if (!dd)
+ return NULL;
+
+ interactive = pfq_current_is_interactive();
+ qid = pfq_rq_qid(rq, interactive);
+ rq_start = rq->__sector;
+ rq_end = rq_start + (rq->__data_len >> SECTOR_SHIFT);
+
+ stats = pfq_stats_lookup(dd->disk_id);
+ queue = pfq_queue_acquire(dd->disk_id, qid);
+ if (!queue)
+ return NULL;
+ bpf_spin_lock(&dd->lock);
+ {
+ struct pfq_adj_find_ctl ctl = {
+ .probe_start = rq_start,
+ .probe_end = rq_end,
+ .out_rn = &rn,
+ .out_req = &targ,
+ .out_stale = &stale,
+ };
+
+ mt = pfq_rq_tree_find_adjacent(dd, queue, &ctl);
+ }
+ if (mt == ELEVATOR_NO_MERGE)
+ goto out;
+
+ if (rq->mq_ctx != targ->mq_ctx || rq->mq_hctx != targ->mq_hctx ||
+ !pfq_rq_merge_attrs_ok(rq, targ)) {
+ stale = pfq_core_put_req_stash(rn, targ);
+ targ = NULL;
+ goto out;
+ }
+
+ tree_core = pfq_rq_tree_remove_take(pfq_rq_tree(dd, queue->qid),
+ &rn->rb_node);
+ list_n = tree_core ?
+ pfq_fifo_tree_remove_take(pfq_fifo_tree(dd, queue->qid),
+ &tree_core->fifo_rb_node) :
+ NULL;
+
+ /* fifo_time inheritance is done by ufq-iosched.c after merge_req. */
+ pfq_account_dec(dd, queue);
+ pfq_rq_clear_node(targ);
+ pfq_stat_add(stats, PFQ_STAT_RQMERGE_CNT, 1);
+ pfq_stat_add(stats, PFQ_STAT_RQMERGE_SIZE, targ->__data_len);
+ bpf_spin_unlock(&dd->lock);
+ bpf_obj_drop(queue);
+ pfq_drop_rq_container_refs(list_n, tree_core);
+ *type = mt;
+ return targ;
+
+out:
+ bpf_spin_unlock(&dd->lock);
+ if (stale)
+ bpf_request_release(stale);
+ bpf_obj_drop(queue);
+ return NULL;
+}
+
+/*
+ * pfq_merge_bio - Merge a bio into a queued request
+ * @q: request queue receiving the bio
+ * @bio: incoming bio
+ * @nr_segs: segment count passed to the merge kfunc
+ * @merged: set true when the bio is absorbed
+ *
+ * Merge within the bio's logical queue, including the current task-name
+ * classification. Set merged on success and keep the survivor queued;
+ * this callback does not coalesce neighboring requests.
+ *
+ * The merge kfunc uses spin_trylock() on ctx->lock to avoid blocking on a
+ * nested lock while dd->lock is held.
+ *
+ * Return: NULL; the surviving request stays queued.
+ */
+struct request *BPF_STRUCT_OPS(pfq_merge_bio, struct request_queue *q,
+ struct bio *bio, unsigned int nr_segs,
+ bool *merged)
+{
+ struct request *cand = NULL, *stale = NULL;
+ struct pfq_queue_data *queue;
+ struct pfq_disk_data *dd;
+ struct pfq_stats *stats;
+ struct pfq_rq_core *rn;
+ sector_t start, end;
+ enum elv_merge mt;
+ bool interactive;
+ u32 qid;
+
+ if (!merged)
+ return NULL;
+
+ dd = pfq_disk_lookup(q->id);
+ if (!dd)
+ return NULL;
+
+ interactive = pfq_current_is_interactive();
+ qid = pfq_bio_qid(bio, interactive);
+ start = bio->bi_iter.bi_sector;
+ end = start + (bio->bi_iter.bi_size >> SECTOR_SHIFT);
+
+ stats = pfq_stats_lookup(dd->disk_id);
+ queue = pfq_queue_acquire(dd->disk_id, qid);
+ if (!queue)
+ return NULL;
+ bpf_spin_lock(&dd->lock);
+ {
+ struct pfq_adj_find_ctl ctl = {
+ .probe_start = start,
+ .probe_end = end,
+ .out_rn = &rn,
+ .out_req = &cand,
+ .out_stale = &stale,
+ };
+
+ mt = pfq_rq_tree_find_adjacent(dd, queue, &ctl);
+ }
+ if (mt == ELEVATOR_NO_MERGE)
+ goto out;
+
+ if (!bpf_request_bio_try_merge(cand, bio, nr_segs)) {
+ stale = pfq_core_put_req_stash(rn, cand);
+ goto out;
+ }
+ *merged = true;
+ stale = pfq_core_put_req_stash(rn, cand);
+ if (mt == ELEVATOR_FRONT_MERGE)
+ pfq_finish_front_bio_merge(dd, queue, rn,
+ bio->bi_iter.bi_sector);
+ pfq_stat_add(stats, PFQ_STAT_BIOMERGE_CNT, 1);
+ pfq_stat_add(stats, PFQ_STAT_BIOMERGE_SIZE, bio->bi_iter.bi_size);
+
+out:
+ bpf_spin_unlock(&dd->lock);
+ if (stale)
+ bpf_request_release(stale);
+ bpf_obj_drop(queue);
+ return NULL;
+}
+
+UFQ_OPS_DEFINE(pfq_ops,
+ .init_sched = (void *)pfq_init_sched,
+ .exit_sched = (void *)pfq_exit_sched,
+ .insert_req = (void *)pfq_insert_req,
+ .dispatch_req = (void *)pfq_dispatch_req,
+ .has_req = (void *)pfq_has_req,
+ .finish_req = (void *)pfq_finish_req,
+ .merge_req = (void *)pfq_merge_req,
+ .merge_bio = (void *)pfq_merge_bio,
+ .name = "pfq_ebpf");
diff --git a/tools/ufq_iosched/pfq.c b/tools/ufq_iosched/pfq.c
new file mode 100644
index 000000000000..78cce0d47abf
--- /dev/null
+++ b/tools/ufq_iosched/pfq.c
@@ -0,0 +1,329 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Copyright (c) 2026 KylinSoft Corporation.
+ * Copyright (c) 2026 Kaitao Cheng <chengkaitao@kylinos.cn>
+ * Copyright (c) 2026 Li Youhong <liyouhong@kylinos.cn>
+ *
+ * Userspace loader for the PFQ eBPF scheduler (UFQ struct_ops backend).
+ */
+#include <stdio.h>
+#include <stdlib.h>
+#include <string.h>
+#include <unistd.h>
+#include <signal.h>
+#include <stdarg.h>
+#include <libgen.h>
+#include <bpf/bpf.h>
+#include <ufq/common.h>
+#include <ufq/pfq.bpf.h>
+#include <ufq/pfq_stat.h>
+#include <ufq/pfq_tunable.h>
+#include <ufq/pfq_disk.h>
+#include "pfq.bpf.skel.h"
+
+const char help_fmt[] =
+"PFQ eBPF scheduler for the UFQ iosched framework.\n"
+"\n"
+"Usage: %s [-v] [-d] [-t SEC] [-i NAME]... [-h]\n"
+"\n"
+" -v Print version\n"
+" -d Print libbpf debug messages\n"
+" -t SEC Stats print interval in seconds (default: 3)\n"
+" -i NAME Prioritize an exact task comm (case-sensitive, 1-15 bytes)\n"
+" Repeat for up to 64 names; default: no interactive names\n"
+" -h Display this help and exit\n";
+
+#define PFQ_VERSION "0.1.0"
+#define TIME_INTERVAL 3
+
+static bool verbose;
+static volatile sig_atomic_t exit_req;
+static unsigned int time_interval = TIME_INTERVAL;
+static __u64 old_stats[PFQ_STAT_MAX];
+static struct pfq_comm_key interactive_comms[PFQ_INTERACTIVE_MAX];
+static unsigned int nr_interactive_comms;
+
+static int libbpf_print_fn(enum libbpf_print_level level, const char *format,
+ va_list args)
+{
+ if (level == LIBBPF_DEBUG && !verbose)
+ return 0;
+ return vfprintf(stderr, format, args);
+}
+
+static void sigint_handler(int sig)
+{
+ (void)sig;
+ exit_req = 1;
+}
+
+static int add_interactive_comm(const char *name)
+{
+ size_t len = strlen(name);
+ unsigned int i;
+
+ if (!len || len >= PFQ_COMM_LEN) {
+ fprintf(stderr, "pfq: -i requires a task name of 1-%d bytes\n",
+ PFQ_COMM_LEN - 1);
+ return -1;
+ }
+ for (i = 0; i < nr_interactive_comms; i++) {
+ if (!strcmp(interactive_comms[i].name, name))
+ return 0;
+ }
+ if (nr_interactive_comms == PFQ_INTERACTIVE_MAX) {
+ fprintf(stderr, "pfq: at most %d distinct task names are supported\n",
+ PFQ_INTERACTIVE_MAX);
+ return -1;
+ }
+
+ /* Static storage keeps the rest of each fixed-size map key zeroed. */
+ memcpy(interactive_comms[nr_interactive_comms++].name, name, len);
+ return 0;
+}
+
+static int init_interactive_comms(struct pfq *skel)
+{
+ int fd = bpf_map__fd(skel->maps.pfq_interactive_comms);
+ __u8 enabled = 1;
+ unsigned int i;
+
+ for (i = 0; i < nr_interactive_comms; i++) {
+ if (bpf_map_update_elem(fd, &interactive_comms[i],
+ &enabled, BPF_ANY) < 0) {
+ fprintf(stderr, "pfq: cannot configure task '%s': %s\n",
+ interactive_comms[i].name, strerror(errno));
+ return -1;
+ }
+ }
+ return 0;
+}
+
+static void init_tunables(struct pfq *skel)
+{
+ int fd = bpf_map__fd(skel->maps.pfq_tunables);
+ __u32 key = PFQ_TUNABLE_KEY;
+ struct pfq_tunables tun = {
+ .weight_base = PFQ_WEIGHT_BASE_DEFAULT,
+ .idle_delay_min_ms = PFQ_IDLE_DELAY_MIN_MS_DEFAULT,
+ .idle_delay_max_ms = PFQ_IDLE_DELAY_MAX_MS_DEFAULT,
+ .batch_read = PFQ_BATCH_READ_DEFAULT,
+ .batch_sync_write = PFQ_BATCH_SYNC_WRITE_DEFAULT,
+ .batch_write = PFQ_BATCH_WRITE_DEFAULT,
+ };
+
+ if (bpf_map_update_elem(fd, &key, &tun, BPF_ANY) < 0)
+ fprintf(stderr, "pfq: failed to seed tunables map\n");
+}
+
+/* Sum CPUs within each disk, then sum the independent disk totals. */
+static int read_stats(struct pfq *skel, __u64 *stats)
+{
+ __u64 total[PFQ_STAT_MAX] = {}, disk[PFQ_STAT_MAX];
+ __s32 key, next_key, seen[PFQ_DISK_MAP_MAX];
+ int nr_cpus = libbpf_num_possible_cpus();
+ int fd = bpf_map__fd(skel->maps.pfq_stats_map);
+ int cpu, idx, ret = -1, nr_seen = 0;
+ struct pfq_stats *values;
+
+ if (nr_cpus <= 0)
+ return -1;
+ values = calloc(nr_cpus, sizeof(*values));
+ if (!values)
+ return -1;
+
+ while (!bpf_map_get_next_key(fd, nr_seen ? &key : NULL, &next_key)) {
+ /*
+ * A deletion can restart iteration. Skip this sample to avoid
+ * double-counting a disk or looping during disk churn.
+ */
+ if (nr_seen == PFQ_DISK_MAP_MAX)
+ goto out;
+ for (idx = 0; idx < nr_seen; idx++)
+ if (seen[idx] == next_key)
+ goto out;
+ seen[nr_seen++] = next_key;
+ key = next_key;
+ if (bpf_map_lookup_elem(fd, &key, values)) {
+ if (errno == ENOENT)
+ continue;
+ goto out;
+ }
+
+ memset(disk, 0, sizeof(disk));
+ for (cpu = 0; cpu < nr_cpus; cpu++)
+ for (idx = 0; idx < PFQ_STAT_MAX; idx++)
+ disk[idx] += values[cpu].counters[idx];
+ /*
+ * Saturate each disk separately: a skewed sample on one disk
+ * must not deduct inserts from another disk.
+ */
+ disk[PFQ_STAT_INSERT_CNT] =
+ disk[PFQ_STAT_INSERT_CNT] >= disk[PFQ_STAT_RQMERGE_CNT] ?
+ disk[PFQ_STAT_INSERT_CNT] - disk[PFQ_STAT_RQMERGE_CNT] : 0;
+ disk[PFQ_STAT_INSERT_SIZE] =
+ disk[PFQ_STAT_INSERT_SIZE] >= disk[PFQ_STAT_RQMERGE_SIZE] ?
+ disk[PFQ_STAT_INSERT_SIZE] - disk[PFQ_STAT_RQMERGE_SIZE] : 0;
+ for (idx = 0; idx < PFQ_STAT_MAX; idx++)
+ total[idx] += disk[idx];
+ }
+ if (errno == ENOENT) {
+ memcpy(stats, total, sizeof(total));
+ ret = 0;
+ }
+out:
+ free(values);
+ return ret;
+}
+
+/* Removing/reinitializing a disk can reduce the aggregate counters. */
+static __u64 stat_delta(__u64 current, __u64 previous)
+{
+ return current >= previous ? current - previous : 0;
+}
+
+/* Sum live per-disk queue state from pfq_map. */
+static void read_disk_state(struct pfq *skel, __u32 *nr_queued_total,
+ __u32 *nr_active)
+{
+ int map_fd = bpf_map__fd(skel->maps.pfq_map);
+ __u32 vsize = bpf_map__value_size(skel->maps.pfq_map);
+ __u32 queued = 0, active = 0;
+ __s32 key = 0, next_key;
+ bool first = true;
+ void *buf;
+
+ if (vsize < sizeof(struct pfq_disk_data_user)) {
+ fprintf(stderr, "pfq: pfq_map value smaller than expected\n");
+ return;
+ }
+
+ buf = malloc(vsize);
+ if (!buf)
+ return;
+
+ while (!bpf_map_get_next_key(map_fd, first ? NULL : &key, &next_key)) {
+ struct pfq_disk_data_user *dd = buf;
+
+ if (!bpf_map_lookup_elem(map_fd, &next_key, buf)) {
+ queued += dd->nr_queued_total;
+ active += dd->nr_active;
+ }
+ key = next_key;
+ first = false;
+ }
+ free(buf);
+ if (nr_queued_total)
+ *nr_queued_total = queued;
+ if (nr_active)
+ *nr_active = active;
+}
+
+static void print_stats(struct pfq *skel)
+{
+ __u32 nr_queued_total = 0, nr_active = 0;
+ __u64 stats[PFQ_STAT_MAX];
+
+ if (read_stats(skel, stats)) {
+ fprintf(stderr, "pfq: cannot read per-CPU statistics\n");
+ return;
+ }
+ read_disk_state(skel, &nr_queued_total, &nr_active);
+
+ printf("bps:%lluk iops:%llu\n",
+ stat_delta(stats[PFQ_STAT_FINISH_SIZE],
+ old_stats[PFQ_STAT_FINISH_SIZE]) / 1024 / time_interval,
+ stat_delta(stats[PFQ_STAT_FINISH_CNT],
+ old_stats[PFQ_STAT_FINISH_CNT]) / time_interval);
+ printf("(insert: cnt=%llu size=%llu err=%llu) (interactive: %llu)\n",
+ stats[PFQ_STAT_INSERT_CNT], stats[PFQ_STAT_INSERT_SIZE],
+ stats[PFQ_STAT_INSERT_ERR],
+ stats[PFQ_STAT_INTERACTIVE]);
+ printf("(at_head: cnt=%llu size=%llu)\n",
+ stats[PFQ_STAT_AT_HEAD_CNT], stats[PFQ_STAT_AT_HEAD_SIZE]);
+ printf("(rqmerge: cnt=%llu size=%llu) (biomerge: cnt=%llu size=%llu)\n",
+ stats[PFQ_STAT_RQMERGE_CNT], stats[PFQ_STAT_RQMERGE_SIZE],
+ stats[PFQ_STAT_BIOMERGE_CNT], stats[PFQ_STAT_BIOMERGE_SIZE]);
+ printf("(dispatch: cnt=%llu size=%llu) (at_head_disp: cnt=%llu size=%llu)\n",
+ stats[PFQ_STAT_DISPATCH_CNT], stats[PFQ_STAT_DISPATCH_SIZE],
+ stats[PFQ_STAT_DISPATCH_AT_HEAD_CNT],
+ stats[PFQ_STAT_DISPATCH_AT_HEAD_SIZE]);
+ printf("(finish: cnt=%llu size=%llu)\n",
+ stats[PFQ_STAT_FINISH_CNT], stats[PFQ_STAT_FINISH_SIZE]);
+ printf("(nr_queued_total: %u) (nr_active: %u) (idle_hold: %llu)\n",
+ nr_queued_total, nr_active, stats[PFQ_STAT_IDLE_HOLD]);
+ memcpy(old_stats, stats, sizeof(old_stats));
+}
+
+int main(int argc, char **argv)
+{
+ struct bpf_link *link;
+ struct pfq *skel;
+ int opt;
+
+ libbpf_set_print(libbpf_print_fn);
+ signal(SIGINT, sigint_handler);
+ signal(SIGTERM, sigint_handler);
+
+ skel = UFQ_OPS_OPEN(pfq_ops, pfq);
+
+ while ((opt = getopt(argc, argv, "vdht:i:")) != -1) {
+ switch (opt) {
+ case 'v':
+ printf("pfq version: %s\n", PFQ_VERSION);
+ pfq__destroy(skel);
+ return 0;
+ case 'd':
+ verbose = true;
+ break;
+ case 'i':
+ if (add_interactive_comm(optarg)) {
+ pfq__destroy(skel);
+ return 1;
+ }
+ break;
+ case 't': {
+ char *end;
+ unsigned long v = strtoul(optarg, &end, 10);
+
+ if (*end || v == 0) {
+ fprintf(stderr, "pfq: invalid -t value\n");
+ pfq__destroy(skel);
+ return 1;
+ }
+ time_interval = v;
+ break;
+ }
+ default:
+ fprintf(stderr, help_fmt, basename(argv[0]));
+ pfq__destroy(skel);
+ return opt != 'h';
+ }
+ }
+
+ if (optind != argc) {
+ fprintf(stderr, "pfq: unexpected argument '%s' (use -i NAME)\n",
+ argv[optind]);
+ pfq__destroy(skel);
+ return 1;
+ }
+
+ UFQ_OPS_LOAD(skel, pfq_ops, pfq);
+ if (init_interactive_comms(skel)) {
+ pfq__destroy(skel);
+ return 1;
+ }
+ init_tunables(skel);
+ link = UFQ_OPS_ATTACH(skel, pfq_ops, pfq);
+
+ printf("PFQ eBPF scheduler v%s attached (struct_ops)\n", PFQ_VERSION);
+
+ while (!exit_req) {
+ sleep(time_interval);
+ print_stats(skel);
+ }
+
+ bpf_link__destroy(link);
+ pfq__destroy(skel);
+ return 0;
+}
--
2.53.0
^ permalink raw reply [flat|nested] 6+ messages in thread* Re: [RFC v3 0/3] block: Introduce a BPF-based I/O scheduler
2026-10-03 4:27 [RFC v3 0/3] block: Introduce a BPF-based I/O scheduler Kaitao Cheng
` (2 preceding siblings ...)
2026-10-03 4:27 ` [RFC v3 3/3] tools/ufq_iosched: add PFQ eBPF " Kaitao Cheng
@ 2026-10-03 9:36 ` Alexei Starovoitov
2026-10-03 11:15 ` Kaitao Cheng
3 siblings, 1 reply; 6+ messages in thread
From: Alexei Starovoitov @ 2026-10-03 9:36 UTC (permalink / raw)
To: Kaitao Cheng, Jens Axboe; +Cc: linux-block, linux-kernel, bpf, Kaitao Cheng
On Sat, Oct 03, 2026 at 12:27 PM Kaitao Cheng <kaitao.cheng@linux.dev> wrote:
> PFQ gives us a concrete policy to explore how well the UFQ interface
> supports more involved scheduling decisions and to guide further work on
> the framework. It has not yet been used in production, and further testing
> and workload evaluation are needed.
Third version in six months and still not a single number.
v2 got replies from the bots only.
Without a solid use case there is no point in polishing this.
> In particular, I would appreciate suggestions on the boundary between the
> UFQ framework and BPF policies, the struct_ops interface, and request
> ownership and fallback handling.
One global ufq_ops for all disks is not the best shape.
I'd do it like bpf_qdisc.
The ownership is the bigger problem.
The request sits in ctx->rq_lists and in a bpf map at the same time
and the kernel relies on the prog to keep the two in sync.
The prog holds rq->ref. rq holds q_usage_counter until
__blk_mq_free_request(). One request that the prog didn't return
from dispatch_req or left in a map and blk_mq_freeze_queue() waits
forever.
sched_ext has a watchdog that kicks the bpf scheduler out.
Something like that is necessary here too.
pw-bot: cr
^ permalink raw reply [flat|nested] 6+ messages in thread* Re: [RFC v3 0/3] block: Introduce a BPF-based I/O scheduler
2026-10-03 9:36 ` [RFC v3 0/3] block: Introduce a BPF-based " Alexei Starovoitov
@ 2026-10-03 11:15 ` Kaitao Cheng
0 siblings, 0 replies; 6+ messages in thread
From: Kaitao Cheng @ 2026-10-03 11:15 UTC (permalink / raw)
To: Alexei Starovoitov, Jens Axboe
Cc: linux-block, linux-kernel, bpf, Kaitao Cheng
在 2026/10/3 17:36, Alexei Starovoitov 写道:
> On Sat, Oct 03, 2026 at 12:27 PM Kaitao Cheng <kaitao.cheng@linux.dev> wrote:
>> PFQ gives us a concrete policy to explore how well the UFQ interface
>> supports more involved scheduling decisions and to guide further work on
>> the framework. It has not yet been used in production, and further testing
>> and workload evaluation are needed.
>
> Third version in six months and still not a single number.
> v2 got replies from the bots only.
> Without a solid use case there is no point in polishing this.
Thank you very much for your review. This is a fairly large patch series,
and I have gone through several iterations because I was concerned that
my approach might not align with the community's expectations. I wanted
to get feedback early so that any architectural issues could be corrected
promptly. I am not seeking inclusion at this stage; I posted the series
for discussion. When I started developing this project, I wasn't entirely
sure it would work either, so I have been writing and testing as I go, haha!
I have actually run some fio tests locally. Initially, moving the I/O
scheduling policy into BPF caused a substantial performance regression.
After three iterations, however, the current implementation can match
the performance of native schedulers such as mq-deadline.
It also performs comparably to the none scheduler on slower SSDs, but
there is still a substantial performance gap on fast NVMe devices. My
analysis suggests that this is because the none scheduler can frequently
bypass ordering in the ctx queues and issue requests directly to the
driver (see blk_mq_try_issue_directly). I'll try to optimize it further.
I am also exploring and testing real-world use cases, which is why I
developed PFQ in the third patch. Next, I plan to evaluate PFQ with
foreground/background applications and with co-located containerized
online services and batch workloads. I may be able to share the results
in the next iteration.
>> In particular, I would appreciate suggestions on the boundary between the
>> UFQ framework and BPF policies, the struct_ops interface, and request
>> ownership and fallback handling.
>
> One global ufq_ops for all disks is not the best shape.
> I'd do it like bpf_qdisc.
Good suggestion. I'll give it a try. Thanks!
> The ownership is the bigger problem.
> The request sits in ctx->rq_lists and in a bpf map at the same time
> and the kernel relies on the prog to keep the two in sync.
> The prog holds rq->ref. rq holds q_usage_counter until
> __blk_mq_free_request(). One request that the prog didn't return
> from dispatch_req or left in a map and blk_mq_freeze_queue() waits
> forever.
> sched_ext has a watchdog that kicks the bpf scheduler out.
> Something like that is necessary here too.
When the BPF scheduler leaks an rq reference, system I/O does indeed stall.
However, because the requests are also kept in ctx->rq_lists, everything
returns to normal once the BPF scheduler exits. A watchdog that kicks the
BPF scheduler out is necessary, and it's also something I was planning to
implement. Thank you very much!
--
Thanks
Kaitao Cheng
^ permalink raw reply [flat|nested] 6+ messages in thread