wuminyu头像
关注

Iinux的双CQE通知以安全回收Buffer原理剖析


前言

本文旨在记录近期研读Java源码的学习心得与疑难问题。由于个人理解水平有限,文中内容难免存在疏漏,恳请读者不吝指正。

双CQE通知以安全回收Buffer原理剖析

1. 架构原理:双 CQE 生命周期与无锁状态机

在 Linux 6.0+ 的 IORING_OP_SEND_ZC 架构中,内核将传统套接字零拷贝复杂同步逻辑重构为基于 io_uring CQE (Completion Queue Event) 的异步状态机。

              ┌─────────────────────────────────────────────────────────┐
              │                   应用程序提交 SQE                      │
              │         io_uring_prep_send_zc(sqe, fd, buf...)          │
              └────────────────────────────┬────────────────────────────┘
                                           │
                                           ▼
                              【状态: BUF_STATE_SUBMITTED】
                              (物理内存被 GUP 钉扎,严禁改写)
                                           │
                   ┌───────────────────────┴───────────────────────┐
                   │                                               │
                   ▼ (内核成功入队)                                ▼ (内核直接拒绝/出错)
        【产生第一 CQE (Syscall)】                     【仅产生单一终结 CQE】
   (cqe->flags & IORING_CQE_F_MORE = True)            (cqe->flags & MORE = False)
   (cqe->res = 发送字节数或错误码)                    (cqe->res = 错误码)
                   │                                               │
                   ▼                                               │
      【状态: BUF_STATE_WAIT_NOTIF】                               │
   (应用层获知发送字节数,但 DMA 未完成)                            │
                   │                                               │
                   ▼ (网卡 DMA 完成, SKB 释放)                      │
        【产生第二 CQE (Notification)】                            │
   (cqe->flags & IORING_CQE_F_MORE = False)                        │
   (cqe->res = 0)                                                  │
                   │                                               │
                   └───────────────────────┬───────────────────────┘
                                           │
                                           ▼
                                 【状态: BUF_STATE_FREE】
                             (无锁回收,可安全写改或 free)

双 CQE 的精确判定契约

  1. 第一 CQE(传输提交结果)
  • cqe->res:表示推入 TCP 发送队列的字节数(若小于 0 则表示系统调用错误,如 -EAGAIN 或 -ECONNRESET)。
  • cqe->flags:当且仅当内核成功接管零拷贝内存并预计后续会有 DMA 释放通知时,内核会置位 IORING_CQE_F_MORE。
  1. 第二 CQE(DMA 完成与内存释放通知)
  • cqe->res:固定为 0(或特定的通知标记)。
  • cqe->flags:**清零 IORING_CQE_F_MORE**。收到该 CQE 代表底层 SKB 引用计数降低为 0,物理页已完全脱离网络栈。

2. 基于 liburing 的完整 C 语言示例代码

以下代码演示了基于 liburing 的 IORING_OP_SEND_ZC 零拷贝发送,包含非阻塞 Socket 对建立、Buffer 内存池设计、无锁双 CQE 解析及异常情况下的安全回收逻辑。

#define _GNU_SOURCE
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <errno.h>
#include <stdbool.h>
#include <fcntl.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <netinet/tcp.h>
#include <liburing.h>

#define RING_ENTRIES     64
#define BUF_SIZE         (64 * 1024) // 64KB 大包,确保正向零拷贝收益
#define POOL_CAPACITY    8

// ============================================================================
// 1. 数据结构与状态机定义
// ============================================================================

typedef enum {
    BUF_STATE_FREE = 0,         // 空闲:应用层拥有完整控制权,可写改或释放
    BUF_STATE_SUBMITTED,        // 已提交:内核托管中,严禁修改/释放
    BUF_STATE_WAIT_NOTIF        // 已收到第一 CQE:发送字节数已确认,等待 DMA 完成通知
} buf_state_t;

// 缓冲区上下文节点 (携带自描述元数据)
typedef struct {
    uint32_t        buf_id;     // 缓冲区唯一标识
    buf_state_t     state;      // 当前状态机位置
    size_t          len;        // 拟发送数据长度
    ssize_t         bytes_sent; // 实际完成发送的字节数
    char            *data;      // 实际数据内存指针 (按页对齐)
} tx_buffer_t;

// 内存池对象 (无锁单线程环形回收)
typedef struct {
    tx_buffer_t     buffers[POOL_CAPACITY];
    tx_buffer_t*    free_stack[POOL_CAPACITY];
    int             top;
} buffer_pool_t;

// ============================================================================
// 2. 内存池初始化与管理逻辑
// ============================================================================

static buffer_pool_t* init_buffer_pool(void) {
    buffer_pool_t *pool = calloc(1, sizeof(buffer_pool_t));
    if (!pool) return NULL;

    pool->top = -1;
    for (int i = 0; i < POOL_CAPACITY; i++) {
        pool->buffers[i].buf_id = i + 1000;
        pool->buffers[i].state = BUF_STATE_FREE;
        pool->buffers[i].len = BUF_SIZE;
        
        // 按照 4096 字节对齐分配物理页,提升 GUP (Get User Pages) 效率
        if (posix_memalign((void**)&pool->buffers[i].data, 4096, BUF_SIZE) != 0) {
            perror("posix_memalign failed");
            exit(EXIT_FAILURE);
        }
        // 压入空闲栈
        pool->free_stack[++pool->top] = &pool->buffers[i];
    }
    return pool;
}

static tx_buffer_t* alloc_buffer(buffer_pool_t *pool) {
    if (pool->top < 0) return NULL; // 池已耗尽
    tx_buffer_t *buf = pool->free_stack[pool->top--];
    buf->state = BUF_STATE_FREE;
    buf->bytes_sent = 0;
    return buf;
}

static void free_buffer(buffer_pool_t *pool, tx_buffer_t *buf) {
    // 使用 C11/GCC 内存屏障保证状态写入前 Buffer 数据改写已完成
    __atomic_store_n(&buf->state, BUF_STATE_FREE, __ATOMIC_RELEASE);
    pool->free_stack[++pool->top] = buf;
}

// 设置非阻塞 Socket
static int set_nonblocking(int fd) {
    int flags = fcntl(fd, F_GETFL, 0);
    if (flags < 0) return -1;
    return fcntl(fd, F_SETFL, flags | O_NONBLOCK);
}

// ============================================================================
// 3. 核心:双 CQE 无锁解析与回收逻辑
// ============================================================================

static void process_cqe_events(struct io_uring *ring, buffer_pool_t *pool) {
    struct io_uring_cqe *cqe;
    unsigned head;
    unsigned count = 0;

    // 【无锁批量轮询】:直接读取共享内存 CQ 环,零系统调用开销
    io_uring_for_each_cqe(ring, head, cqe) {
        count++;
        // 从 user_data 还原缓冲区描述符指针
        tx_buffer_t *buf = (tx_buffer_t *)(uintptr_t)io_uring_cqe_get_data64(cqe);
        if (!buf) {
            continue; // 非 zero-copy 操作或无句柄事件
        }

        bool has_more = (cqe->flags & IORING_CQE_F_MORE) != 0;

        if (has_more) {
            // ----------------------------------------------------------------
            // 情况 A:收到第一 CQE (提交结果通知)
            // ----------------------------------------------------------------
            if (cqe->res >= 0) {
                buf->bytes_sent = cqe->res;
                buf->state = BUF_STATE_WAIT_NOTIF;
                printf("[CQE-1 成功] Buffer ID: %u | 传输字节: %d | 标志: IORING_CQE_F_MORE | 状态 -> BUF_STATE_WAIT_NOTIF\n",
                       buf->buf_id, cqe->res);
                printf("          └─► ⚠️  安全屏障:DMA 传输中,物理页 %p 严禁写入/回收!\n", (void*)buf->data);
            } else {
                // 第一 CQE 即返回错误 (如 -EAGAIN / -EPIPE)
                fprintf(stderr, "[CQE-1 错误] Buffer ID: %u | 错误码: %d (%s)\n",
                        buf->buf_id, cqe->res, strerror(-cqe->res));
                buf->bytes_sent = cqe->res;
                // 注意:即使报错,若带有 F_MORE,仍必须等待第二 CQE 才能回收!
                buf->state = BUF_STATE_WAIT_NOTIF;
            }
        } else {
            // ----------------------------------------------------------------
            // 情况 B:收到终结 CQE (Notification 或者无 MORE 的单 CQE)
            // ----------------------------------------------------------------
            if (buf->state == BUF_STATE_WAIT_NOTIF) {
                printf("[CQE-2 完成] Buffer ID: %u | DMA 释放完成 | 状态 -> BUF_STATE_FREE\n", buf->buf_id);
                printf("          └─► ✅ 内存屏障解除:物理页 %p 已从内核网络栈脱离,安全的归还内存池!\n", (void*)buf->data);
            } else if (buf->state == BUF_STATE_SUBMITTED) {
                // 内核由于严重错误未置位 F_MORE 直接终结请求 (未产生第一 CQE)
                printf("[CQE-单终结] Buffer ID: %u | 直接终结 (res=%d) | 状态 -> BUF_STATE_FREE\n",
                       buf->buf_id, cqe->res);
            }

            // 【安全回收点】:无锁归还内存池
            free_buffer(pool, buf);
        }
    }

    // 批量更新 CQ 环头指针,通知内核事件已消费
    if (count > 0) {
        io_uring_cq_advance(ring, count);
    }
}

// ============================================================================
// 4. 主程序流程
// ============================================================================

int main(void) {
    int fds[2];
    struct io_uring ring;
    buffer_pool_t *pool;

    // 1. 创建 Unix 域套接字对用于测试
    if (socketpair(AF_UNIX, SOCK_STREAM, 0, fds) < 0) {
        perror("socketpair failed");
        exit(EXIT_FAILURE);
    }
    set_nonblocking(fds[0]);
    set_nonblocking(fds[1]);

    // 放大套接字发送/接收缓冲区,避免因 Socket 满了退化
    int sndbuf_size = 1024 * 1024;
    setsockopt(fds[0], SOL_SOCKET, SO_SNDBUF, &sndbuf_size, sizeof(sndbuf_size));
    setsockopt(fds[1], SOL_SOCKET, SO_RCVBUF, &sndbuf_size, sizeof(sndbuf_size));

    // 2. 初始化 io_uring
    if (io_uring_queue_init(RING_ENTRIES, &ring, 0) < 0) {
        perror("io_uring_queue_init failed");
        exit(EXIT_FAILURE);
    }

    // 3. 初始化内存池
    pool = init_buffer_pool();

    // 4. 从池中分配 Buffer 并填入数据
    tx_buffer_t *buf = alloc_buffer(pool);
    if (!buf) {
        fprintf(stderr, "Buffer 分配失败\n");
        exit(EXIT_FAILURE);
    }
    memset(buf->data, 'Z', BUF_SIZE);

    // 5. 准备并提交零拷贝发送请求 (IORING_OP_SEND_ZC)
    struct io_uring_sqe *sqe = io_uring_get_sqe(&ring);
    if (!sqe) {
        fprintf(stderr, "SQE 获取失败\n");
        exit(EXIT_FAILURE);
    }

    // 使用 liburing 准备 Send Zero-Copy 算子
    io_uring_prep_send_zc(sqe, fds[0], buf->data, buf->len, 0, 0);
    // 将 Buffer 句柄指针直接绑定到 user_data
    io_uring_sqe_set_data64(sqe, (uintptr_t)buf);
    
    // 标记状态为已提交
    buf->state = BUF_STATE_SUBMITTED;

    printf("==> 提交 IORING_OP_SEND_ZC 请求 (Buffer ID: %u, Size: %zu 字节)...\n", buf->buf_id, buf->len);
    io_uring_submit(&ring);

    // 6. 循环等待并无锁解析双 CQE
    int notif_count = 0;
    while (notif_count < 2) {
        // 等待至少 1 个 CQE 到达
        io_uring_wait_cqe(&ring, NULL);
        
        // 消费解析当前 CQ 环中的所有事件
        process_cqe_events(&ring, pool);

        // 如果 Buffer 已经被安全回收归还,说明双 CQE 流程已走完
        if (pool->top >= 0 && pool->free_stack[pool->top] == buf) {
            notif_count = 2; // 退出循环
        }
    }

    // 清理资源
    io_uring_queue_exit(&ring);
    close(fds[0]);
    close(fds[1]);
    printf("==> 测试完成,Buffer 安全回收成功。\n");
    return 0;
}


3. 代码关键节点与源码级机制深度拆解

3.1 内存元数据与 user_data 的无锁绑定机制

在提交 SQE 时,程序通过以下调用:

io_uring_sqe_set_data64(sqe, (uintptr_t)buf);

将应用层的 tx_buffer_t 结构体指针赋值给 sqe->user_data。

  • 内核透传原理:内核在处理 IORING_OP_SEND_ZC 时,分配的通知上下文 struct io_notif_data 会继承该 user_data 值。
  • 双 CQE 共享关联:无论是第一 CQE(发送结果)还是第二 CQE(DMA 释放通知),内核在向共享 CQ 环压入 struct io_uring_cqe 时,**均会原封不动地复制相同的 user_data**。这使得应用层在处理任何一个 CQE 时,均可通过 O ( 1 ) O(1) O(1) 的指针转换恢复上下文:
tx_buffer_t *buf = (tx_buffer_t *)(uintptr_t)io_uring_cqe_get_data64(cqe);

3.2 IORING_CQE_F_MORE 标志位的无锁状态切换

在 process_cqe_events 函数中,核心判定在于逻辑分流:

bool has_more = (cqe->flags & IORING_CQE_F_MORE) != 0;

状态机转化表
触发事件cqe->res 值IORING_CQE_F_MOREBuffer 转换后状态内存回收动作
CQE 1: 传输提交成功 > 0 > 0 >0 (发送字节数)Set (1)BUF_STATE_WAIT_NOTIF禁止回收 (网卡 DMA 尚在读取)
CQE 1: 传输提交失败 < 0 < 0 <0 (错误码)Set (1)BUF_STATE_WAIT_NOTIF禁止回收 (仍需等第二 CQE 解锁)
CQE 2: DMA 释放通知 = 0 = 0 =0Cleared (0)BUF_STATE_FREE安全回收 (归还无锁内存池)
异常: 单一终结 CQE < 0 < 0 <0 (拒绝请求)Cleared (0)BUF_STATE_FREE安全回收 (直接归还内存池)

3.3 无锁(Lock-Free)内存安全屏障设计

在传统多线程网络库中,内存回收往往需要争抢 pthread_mutex。而在 io_uring 模式下,结合单线程 Event-Loop,可通过轻量级 C11/GCC 内存屏障保障绝对的写屏障顺序:

static void free_buffer(buffer_pool_t *pool, tx_buffer_t *buf) {
    // __ATOMIC_RELEASE 保证在 state 更改为 FREE 之前,所有对于 buf->data 的读写已全部完成
    __atomic_store_n(&buf->state, BUF_STATE_FREE, __ATOMIC_RELEASE);
    pool->free_stack[++pool->top] = buf;
}

  • 写改防护:只要 buf->state 处于 BUF_STATE_SUBMITTED 或 BUF_STATE_WAIT_NOTIF,应用层任何工作线程试图获取或修改 buf->data 都将被拒绝。
  • CPU 乱序重排防护:__ATOMIC_RELEASE 阻止编译器和 CPU 将后续归还内存池的操作重排到 DMA 完成之前,从而在硬件层面杜绝了 Use-After-Free 或 Data Corruption 风险。

4. 生产级场景下的异常处理与边缘边界

在实际高性能网络编程中,应用层还需要处理以下三种复杂边缘场景:

4.1 短发送(Short Sends)与多次提交

由于 Socket 发送缓冲区满了,IORING_OP_SEND_ZC 的第一 CQE 返回的 cqe->res 可能是部分字节数(例如请求 64   KB 64\,\text{KB} 64KB,实际仅发送 16   KB 16\,\text{KB} 16KB)。

  • 处理规则:
  1. 第一 CQE 返回 res = 16384 且带 IORING_CQE_F_MORE。
  2. 必须记录 buf->bytes_sent = 16384。
  3. 注意:该 16   KB 16\,\text{KB} 16KB 对应的物理页仍被内核 Pin 住,直到第二 CQE 达到前,整块 64   KB 64\,\text{KB} 64KB 内存依然不可写改。
  4. 剩余 48   KB 48\,\text{KB} 48KB 如果需要再次发送,必须使用新的偏移量分配新 SQE 提交,不能直接覆盖原 Buffer!

4.2 网络连接异常断开(EPIPE / ECONNRESET)

当客户端强制断开连接时,提交零拷贝发送可能会直接报错。

  • 处理规则:
    内核网络栈可能依然会触发双 CQE 流程:
  • CQE 1:cqe->res = -EPIPE,flags 包含 IORING_CQE_F_MORE。
  • CQE 2:cqe->res = 0,flags 不带 IORING_CQE_F_MORE。
    解析逻辑必须确保:哪怕 CQE 1 已经返回了负数错误码,只要 IORING_CQE_F_MORE 被置位,就绝不能提前释放内存,必须死守第二 CQE 的到来,否则将引发内核网卡 DMA 访问已释放内存的 Kernel Panic!

4.3 零拷贝退化为 Copy 场景

当发送数据包小于内核阈值(或遭遇 Copy-On-Write 页面)时,内核内部可能会将零拷贝退化为常规内核拷贝。
此时,内核仍然会遵循契约投递双 CQE(或直接在第一 CQE 后快速投递第二 CQE)。应用层无需关注内核内部是否真的使用了零拷贝,只需统一按照 IORING_CQE_F_MORE 标志位判定即可,使得应用层与内核底层实现彻底解耦。

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/wuminyu/article/details/166634161

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--