Lely canopenio_poll 机制原理与完整运行流程详解

在这里插入图片描述

@[toc]

1. 先给结论

Lely 的 io_poll 是一个基于 Linux epoll 的 I/O 事件分发器。它不解析 CANopen,也不主动读取 SocketCAN;它只负责:

  1. 把需要关注的文件描述符 fd 登记到 Linux epoll
  2. 阻塞等待这些 fd 出现可读、可写、异常或断开事件;
  3. 根据内核返回的 fd,在红黑树中找到对应的 io_poll_watch
  4. 调用 watch->func(watch, events)
  5. 因为所有登记都附加了 EPOLLONESHOT,每次通知后当前一次监听失效;
  6. 调用方若要继续接收后续事件,必须再次调用 io_poll_watch() 恢复下一次监听。

最重要的完整周期是:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
创建 Poll

初始化 io_poll_watch

io_poll_watch(events != 0)

首次登记:EPOLL_CTL_ADD

ev_poll_wait() → epoll_pwait()

内核返回 fd 和 EPOLL* 事件

红黑树按 fd 找到 io_poll_watch

将 watch->_events 清零,表示本次一次性监听已消费

调用 watch->func()

回调读取/写入 fd

仍需继续监听?
├─ 是:再次 io_poll_watch(events != 0) → EPOLL_CTL_MOD,恢复下一次监听
└─ 否:io_poll_watch(events = 0) → EPOLL_CTL_DEL,彻底取消登记
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
flowchart TD
A[创建 io_poll] --> B[初始化 io_poll_watch]
B --> C[io_poll_watch: events 非 0]
C --> D{红黑树中已有该 fd?}
D -- 否 --> E[EPOLL_CTL_ADD]
D -- 是且本次监听需更新/恢复 --> F[EPOLL_CTL_MOD]
E --> G[当前一次监听有效]
F --> G
G --> H[epoll_pwait 等待]
H --> I[内核返回 fd 和 EPOLL 事件]
I --> J[红黑树按 fd 找到 watch]
J --> K[_events=0, nwatch--]
K --> L[调用 watch->func]
L --> M{还要继续监听?}
M -- 是 --> C
M -- 否 --> N[io_poll_watch: events=0]
N --> O[EPOLL_CTL_DEL + 红黑树删除]
O --> P[销毁 Poll]

2. Poll 在 Lely CANopen 全栈中的位置:哪些模块使用它、分别做什么

本章基于 Lely CANopen 源码中的实际调用关系回答两个问题:

  1. 哪些模块直接调用 io_poll_watch(),把文件描述符交给 Poll 监听;
  2. 哪些 CANopen 组件虽然不直接调用 Poll,但其异步收发和定时机制最终依赖 Poll

这里分析的是新 I/O 库 liblely-io2 的 Linux/POSIX 实现,对应:

1
2
3
4
include/lely/io2/...
src/io2/...
src/ev/...
src/coapp/...

不是旧 I/O 库中的:

1
2
include/lely/io/poll.h
src/io/poll.c

两套库中都存在名为 poll 的接口,但对象模型和调用方式不同。本文后续所说的 Poll,均指本文件正在分析的 io2 实现:

1
2
3
include/lely/io2/posix/poll.h
include/lely/io2/posix/poll.hpp
src/io2/linux/poll.c

2.1 先给结论:直接把 fd 登记到 Poll 的核心组件只有三个

liblely-io2 的 Linux/POSIX 后端中,搜索 IO_POLL_WATCH_INITio_poll_watch() 可以看到三个核心设备实现直接持有 struct io_poll_watch

直接使用 Poll 的组件 源码文件 交给 Poll 的 fd 监听事件 主要作用
Linux CAN 通道 src/io2/linux/can_chan.c SocketCAN raw socket IO_EVENT_INIO_EVENT_OUT CAN 帧异步接收、发送阻塞后的重试、发送确认处理
Linux 系统定时器 src/io2/linux/timer.c timerfd IO_EVENT_IN 把定时器到期转换成异步 wait 完成事件
POSIX 信号集合 src/io2/posix/sigset.c self-pipe 的读端 IO_EVENT_IN 把 POSIX signal 安全地转换成事件循环任务

此外,还有一个非常关键但角色不同的组件:

组件 是否调用 io_poll_watch() 它对 Poll 做什么
ev::Loop / ev_loop 调用 ev_poll_wait() 驱动 Poll 等待事件;需要唤醒时调用 ev_poll_kill()

再往上的组件都属于间接依赖:

上层组件 是否直接操作 Poll 实际依赖路径
io::CanNet 使用 CanChannel 收发 CAN 帧,使用 Timer 驱动 CAN 网络时间和超时
canopen::Node 继承/组合 CanNet,把收到的 CAN 帧交给 CANopen 协议对象
BasicMaster / AsyncMaster 通过 Node → CanNet → CanChannel/Timer → Poll 间接使用
BasicSlave 同样通过 Node 和 I/O 层间接使用
DriverBase / BasicDriver / LoopDriver 消费 Master 分发的 CANopen 事件,不负责监听 SocketCAN fd

因此,不能把 Poll 理解成“CANopen 主站模块的一部分”。更准确的层次是:

1
2
3
4
Poll 是底层异步 I/O 事件源聚合器;
CanChannel、Timer、SignalSet 是它的直接设备客户端;
ev::Loop 是它的等待驱动者;
CanNet、Node、Master、Slave 和 CANopen Driver 是间接消费者。

2.2 全栈依赖关系图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
flowchart TB
APP[用户应用程序] --> CTX[io::Context]
APP --> POLL[io::Poll]
APP --> LOOP[ev::Loop]
APP --> TIMER[io::Timer]
APP --> CHAN[io::CanChannel]
APP --> SIGSET[io::SignalSet]

CTX --> POLL
POLL -- get_poll() --> LOOP
LOOP -- ev_poll_wait() --> POLL
LOOP -- 执行已投递任务 --> EXEC[Executor task queue]

POLL -- IO_EVENT_IN / OUT --> CHAN
CHAN --> CANFD[SocketCAN raw socket fd]

POLL -- IO_EVENT_IN --> TIMER
TIMER --> TFD[timerfd]

POLL -- IO_EVENT_IN --> SIGSET
SIGSET --> PIPE[self-pipe read fd]

CHAN -- 异步读写完成任务 --> EXEC
TIMER -- 定时等待完成任务 --> EXEC
SIGSET -- 信号处理任务 --> EXEC

CHAN --> CANNET[io::CanNet]
TIMER --> CANNET
CANNET --> PASSIVE[can_net_t 被动 CAN 网络核心]
PASSIVE --> NODE[canopen::Node]
NODE --> MASTER[BasicMaster / AsyncMaster]
NODE --> SLAVE[BasicSlave]
MASTER --> DRIVER[DriverBase / BasicDriver / LoopDriver]

图中最关键的边界是:

1
2
3
4
5
6
7
8
9
10
11
Linux fd 就绪通知

Poll

设备回调只负责投递 task

ev::Loop / Executor 执行 task

CanChannel、Timer、SignalSet 完成真正处理

CanNet 把 I/O 结果送入被动 CAN/CANopen 协议核心

Poll 的回调通常不直接执行完整协议逻辑,而是尽快把后续工作投递给 executor。这样可以缩短 epoll_pwait() 返回后的临界处理时间,并避免在 Poll 内部互斥锁上下文中执行复杂业务。

2.3 C++ 应用是怎样把这些组件连接起来的

典型 C++ 初始化关系如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
io::Context ctx;
io::Poll poll(ctx);
ev::Loop loop(poll.get_poll());
auto exec = loop.get_executor();

io::Timer timer(poll, exec, CLOCK_MONOTONIC);
io::CanController ctrl("can0");
io::CanChannel chan(poll, exec);
chan.open(ctrl);

canopen::AsyncMaster master(timer, chan, /* DCF 等参数 */);

io::SignalSet sigset(poll, exec);
loop.run();

每个构造参数都对应一条真实依赖:

构造关系 含义
Poll(ctx) Poll 注册到 I/O context,统一管理生命周期和 shutdown/fork 通知
Loop(poll.get_poll()) Event loop 获得与平台无关的 ev_poll_t 接口,用它阻塞等待 I/O
Timer(poll, exec, ...) Timer 把 timerfd 登记到 Poll;到期后将 wait task 投递到 executor
CanChannel(poll, exec) CAN channel 把 SocketCAN fd 的可读/可写事件交给 Poll
AsyncMaster(timer, chan, ...) Master 不直接拿 Poll;它通过 Timer 和 CanChannel 获得时间和 CAN I/O
SignalSet(poll, exec) SignalSet 把 self-pipe 读端登记到 Poll,统一处理退出信号
loop.run() Event loop 开始执行任务,并在任务队列为空时进入 ev_poll_wait()

换句话说,Poll 不是由 AsyncMaster 创建或拥有的。通常由应用在最外层创建,然后把同一个 Poll 传给所有需要异步 fd 事件的 I/O 组件。

2.4 直接用户一:Linux CAN 通道 can_chan.c

2.4.1 它持有哪些与 Poll 相关的成员

struct io_can_chan_impl 中的关键成员是:

1
2
3
4
5
6
7
8
9
struct io_can_chan_impl {
/* ... */
io_poll_t *poll;
struct io_poll_watch watch;
/* ... */
int fd;
int events;
/* ... */
};

它们分别表示:

成员 含义
poll 当前 CAN channel 使用的 Poll 实例
watch 该 SocketCAN fd 对应的回调和 Poll 内部索引节点
fd Linux SocketCAN raw socket
events CAN channel 当前希望继续关注的 Lely I/O 事件位

初始化时,CAN channel 保存 Poll,并给 watch 设置回调:

1
2
3
4
5
6
7
impl->poll = poll;

impl->watch = (struct io_poll_watch)IO_POLL_WATCH_INIT(
&io_can_chan_impl_watch_func);

impl->fd = -1;
impl->events = 0;

注意:创建 CanChannel 时并不一定立即登记 SocketCAN fd。只有 channel 打开/绑定具体 CAN 控制器,并且异步操作遇到需要等待的条件后,才会根据需要登记 IO_EVENT_INIO_EVENT_OUT

2.4.2 读路径:只有非阻塞读取返回 EAGAIN 时才监听 IO_EVENT_IN

CAN channel 的接收任务会先直接尝试从 SocketCAN fd 读取帧。只要 socket 中仍有数据,就继续读取并填充用户态接收队列。

read() 返回:

1
EAGAIN 或 EWOULDBLOCK

表示当前没有更多 CAN 帧可读。如果仍存在待完成的读请求,代码才执行:

1
2
3
4
5
6
int events = impl->events | IO_EVENT_IN;

if (!io_poll_watch(impl->poll, impl->fd, events, &impl->watch)) {
impl->events = events;
post_rxbuf = 0;
}

这里的动作是:

  1. 保留此前可能已经关注的事件位;
  2. 增加 IO_EVENT_IN
  3. 把 SocketCAN fd 登记或重新登记到 Poll;
  4. 不再反复投递接收任务,等待内核通知“现在又可以读了”。

所以 Poll 在读路径中的价值不是“替代 read()”,而是:

1
2
3
4
5
直接 read() 已经无法继续推进

用 Poll 睡眠等待 SocketCAN 重新可读

避免 CPU 忙轮询

2.4.3 写路径:非阻塞写返回 EAGAIN 时监听 IO_EVENT_OUT

发送任务同样会先直接尝试写 SocketCAN fd:

1
2
3
4
5
int result = io_can_fd_write_msg(fd, write->msg,
impl->poll ? 0 : LELY_IO_TX_TIMEOUT);

int errc = !result ? 0 : errno;
int wouldblock = errc == EAGAIN || errc == EWOULDBLOCK;

使用 Poll 时,写操作采用非阻塞模式。若发送缓冲区暂时无法接收新帧,待发送操作被重新放回队列,然后登记 IO_EVENT_OUT

1
2
3
4
5
6
int events = impl->events | IO_EVENT_OUT;

if (!io_poll_watch(impl->poll, impl->fd, events, &impl->watch)) {
impl->events = events;
post_write = 0;
}

等 Linux 报告 EPOLLOUT 后,Poll 转换为 IO_EVENT_OUT,再触发 CAN channel 的 watch 回调,由回调投递写任务重新尝试发送。

这里也不能把 IO_EVENT_OUT 理解成“这一帧必定已经发到 CAN 总线”。它只表示 socket 当前允许继续推进写入。若启用了 txwait,CAN channel 还会等待 SocketCAN 的发送确认消息。

2.4.4 同一个 SocketCAN fd 可以同时监听 IN 和 OUT

impl->events 是位掩码,因此可能出现:

1
IO_EVENT_IN | IO_EVENT_OUT

典型场景是:

  • 接收队列正等待新 CAN 帧;
  • 发送队列又因为 socket 暂时不可写而等待重试。

此时同一个 fd 只需要一个 io_poll_watch,但关注两个方向的事件。

2.4.5 CAN channel 的 Poll 回调做什么

回调入口:

1
2
3
4
5
6
7
static void
io_can_chan_impl_watch_func(struct io_poll_watch *watch, int events)
{
struct io_can_chan_impl *impl =
structof(watch, struct io_can_chan_impl, watch);
/* ... */
}

它首先处理 EPOLLONESHOT 带来的“本次通知已消费”问题:

1
2
3
4
5
6
7
8
9
if ((events & IO_EVENT_ERR) || impl->fd == -1 || impl->shutdown) {
impl->events = 0;
} else if ((impl->events &= ~events) != 0) {
if (io_poll_watch(impl->poll, impl->fd, impl->events,
&impl->watch) == -1) {
impl->events = 0;
events |= IO_EVENT_ERR;
}
}

这段代码的含义是:

  1. 如果发生错误、fd 已关闭或 channel 正在 shutdown,不再恢复监听;
  2. 否则,从 impl->events 中清除本次已经发生的事件位;
  3. 如果仍剩其他关注事件,立即重新调用 io_poll_watch()
  4. 这样可以继续等待“尚未发生的另一个方向事件”。

举例:

1
2
3
4
原来关注:IN | OUT
本次只返回:IN
清除后剩余:OUT
动作:立即重新登记 OUT

随后回调不直接在 Poll 上下文中完成全部收发,而是投递任务:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
if ((events & (IO_EVENT_IN | IO_EVENT_ERR)) && !impl->shutdown) {
post_rxbuf = !impl->rxbuf_posted;
impl->rxbuf_posted = 1;
}

if ((events & (IO_EVENT_OUT | IO_EVENT_ERR))
&& !sllist_empty(&impl->write_queue)
&& !impl->shutdown) {
post_write = !impl->write_posted;
impl->write_posted = 1;
}

if (post_rxbuf)
ev_exec_post(impl->rxbuf_task.exec, &impl->rxbuf_task);
if (post_write)
ev_exec_post(impl->write_task.exec, &impl->write_task);

因此它的职责可以概括为:

Poll 事件 CAN channel 回调动作 后续真正处理者
IO_EVENT_IN 投递 rxbuf_task 接收任务读取 SocketCAN、填充接收队列、完成异步读请求
IO_EVENT_OUT 投递 write_task 写任务重试发送队列中的 CAN 帧
IO_EVENT_ERR 同时唤醒读写处理路径 读写任务读取错误状态并完成/取消相关请求

2.4.6 CAN 接收的完整链路

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
sequenceDiagram
participant K as Linux SocketCAN
participant P as io::Poll
participant C as CanChannel watch callback
participant E as Executor / ev::Loop
participant R as rxbuf_task/read_task
participant N as io::CanNet
participant CN as can_net_t / CANopen objects

R->>K: 非阻塞 read()
K-->>R: EAGAIN
R->>P: io_poll_watch(fd, IO_EVENT_IN)
P->>K: epoll 监听 EPOLLIN|EPOLLONESHOT
K-->>P: fd 可读
P->>C: watch->>func(IO_EVENT_IN)
C->>E: ev_exec_post(rxbuf_task)
E->>R: 执行接收任务
R->>K: 读取并排空当前 CAN 帧
R->>N: 完成 io_can_chan_read
N->>CN: can_net_recv(msg)
R->>P: 若再次 EAGAIN 且仍需接收,重新监听 IN

2.4.7 CAN 发送的完整链路

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
sequenceDiagram
participant CN as CANopen protocol
participant N as io::CanNet
participant W as CanChannel write_task
participant K as Linux SocketCAN
participant P as io::Poll
participant C as CanChannel watch callback
participant E as Executor / ev::Loop

CN->>N: send callback(CAN frame)
N->>W: io_can_chan_submit_write()
W->>K: 非阻塞 write()
K-->>W: EAGAIN
W->>P: io_poll_watch(fd, IO_EVENT_OUT)
K-->>P: fd 可写
P->>C: watch->>func(IO_EVENT_OUT)
C->>E: ev_exec_post(write_task)
E->>W: 重试发送
W->>K: write(CAN frame)
K-->>W: 成功或发送确认
W->>N: 完成 write operation

2.4.8 关闭时做什么

CAN channel shutdown 或关闭 socket 时会执行:

1
2
3
impl->events = 0;
io_poll_watch(impl->poll, impl->fd, 0, &impl->watch);
close(impl->fd);

顺序很重要:先从 Poll 删除 fd,再关闭 fd,避免后续事件仍通过旧的 watch 进入回调。

2.5 直接用户二:Linux 系统定时器 timer.c

2.5.1 Timer 把 Linux timerfd 当作普通可读 fd

Linux timer 实现持有:

1
2
3
4
5
6
7
8
9
10
struct io_timer_impl {
/* ... */
io_poll_t *poll;
struct io_poll_watch watch;
int tfd;
struct ev_task wait_task;
struct sllist wait_queue;
int overrun;
/* ... */
};

打开定时器时:

1
2
3
4
5
6
7
8
impl->tfd = timerfd_create(
impl->clockid,
TFD_NONBLOCK | TFD_CLOEXEC);

if (io_poll_watch(impl->poll, impl->tfd,
IO_EVENT_IN, &impl->watch) == -1) {
/* error handling */
}

timerfd 到期后会变成“可读”。因此 Timer 不需要特殊的 Poll 类型,只需监听 IO_EVENT_IN

2.5.2 Poll 回调只负责投递 wait_task

Timer 的 watch 回调很短:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
static void
io_timer_impl_watch_func(struct io_poll_watch *watch, int events)
{
struct io_timer_impl *impl =
structof(watch, struct io_timer_impl, watch);
(void)events;

int post_wait = !impl->wait_posted && !impl->shutdown;
if (post_wait)
impl->wait_posted = 1;

if (post_wait)
ev_exec_post(impl->wait_task.exec, &impl->wait_task);
}

它没有在 Poll 回调中直接读取 timerfd,而是把任务交给 executor。

2.5.3 wait_task 读取计数并恢复下一次监听

timerfd 的一次 read() 返回的是累计到期次数,而不是普通字节流。任务会循环读取并累加 overrun

1
2
uintmax_t value = 0;
result = read(impl->tfd, &value, sizeof(value));

当数据被排空后,非阻塞 read() 返回 EAGAIN,代码设置:

1
events |= IO_EVENT_IN;

随后:

1
2
if (events && !impl->shutdown)
io_poll_watch(impl->poll, impl->tfd, events, &impl->watch);

这一步恢复下一次 EPOLLONESHOT 监听。

完整周期是:

1
2
3
4
5
6
7
8
9
10
11
12
13
timerfd_settime() 设置到期时间

timerfd 到期,可读

Poll 返回 IO_EVENT_IN

Timer watch callback 投递 wait_task

wait_task 读取到期计数

完成 wait_queue 中的异步等待操作

重新 io_poll_watch(IO_EVENT_IN)

2.5.4 Timer 与 CANopen 定时逻辑的关系

io::Timer 本身不知道 SDO、NMT、Heartbeat 或 PDO。它只提供通用异步时间等待。

真正把它接到 CANopen 时间系统的是 io::CanNet

1
2
3
4
5
io::CanNet
├─ 持有 io_timer_t *timer
├─ 基于 timer 创建 io_tqueue
├─ 等待 can_net_t 给出的“下一次到期时间”
└─ 到期时调用 can_net_set_time()

因此应区分两层:

认识的概念
src/io2/linux/timer.c timerfd、到期次数、异步 wait queue
io_can_net / can_net_t CAN 网络当前时间、下一个协议定时器、超时回调

Poll 只负责让第一层在到期时被唤醒;CANopen 协议栈根据自己的 timer heap 决定具体执行哪个协议定时动作。

2.6 直接用户三:POSIX 信号集合 sigset.c

2.6.1 Poll 并不是直接监听 POSIX signal

普通 signal 不是一个 fd,不能直接加入 epoll。Lely 使用 self-pipe:

1
2
3
4
5
6
7
signal handler
↓ 写 1 字节
非阻塞 pipe 写端

非阻塞 pipe 读端变为可读

Poll 返回 IO_EVENT_IN

关键成员:

1
2
3
4
5
6
7
8
9
struct io_sigset_impl {
/* ... */
io_poll_t *poll;
struct io_poll_watch watch;
int fds[2];
struct ev_task read_task;
struct ev_task wait_task;
/* ... */
};

2.6.2 创建 self-pipe 并监听读端

Linux 下:

1
2
3
4
pipe2(impl->fds, O_CLOEXEC | O_NONBLOCK);

io_poll_watch(impl->poll, impl->fds[0],
IO_EVENT_IN, &impl->watch);

其中:

fd 用途
fds[0] pipe 读端,交给 Poll 监听
fds[1] pipe 写端,由 signal handler 写入

2.6.3 signal handler 只做最小工作

信号处理函数记录 pending 状态并向 pipe 写一个字节:

1
2
3
4
5
6
7
static void
io_sigset_handler(int signo)
{
io_sigset_shared.sig[signo - 1].pending = 1;
io_sigset_shared.pending = 1;
io_sigset_kill(signo);
}

io_sigset_kill() 的核心动作是:

1
write(pipe_write_fd, "", 1);

这样做的原因是 signal handler 上下文限制很多,不能安全执行复杂的 C/C++ 逻辑、锁操作或用户回调。写 self-pipe 是一种把异步 signal 转换到普通事件循环上下文的方式。

2.6.4 Poll 回调和后续任务

pipe 可读时,watch 回调只投递读取任务:

1
2
3
4
5
int post_read = !impl->read_posted;
impl->read_posted = 1;

if (post_read)
ev_exec_post(impl->read_task.exec, &impl->read_task);

read_task 会:

  1. 循环读空 pipe;
  2. 调用 io_sigset_process_all() 处理 pending signal;
  3. 把对应 signal 交给等待队列;
  4. 当读到 EAGAIN 后重新登记 IO_EVENT_IN

恢复监听代码:

1
2
3
if (events && !impl->shutdown)
io_poll_watch(impl->poll, impl->fds[0],
events, &impl->watch);

完整链路:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
sequenceDiagram
participant OS as POSIX signal
participant H as signal handler
participant W as self-pipe write fd
participant R as self-pipe read fd
participant P as io::Poll
participant E as Executor / ev::Loop
participant S as SignalSet tasks
participant A as Application shutdown logic

OS->>H: SIGINT / SIGTERM
H->>W: write 1 byte + set pending flag
W-->>R: pipe becomes readable
R-->>P: EPOLLIN
P->>S: watch callback(IO_EVENT_IN)
S->>E: post read_task
E->>S: drain pipe and process pending signals
S->>A: complete async signal wait
S->>P: re-register IO_EVENT_IN

2.6.5 在 CANopen 程序中的作用

SignalSet 通常用于:

  • 捕获 SIGINTSIGTERM
  • 请求 Context shutdown;
  • 让 CAN channel、timer、CAN network 等注册在 context 中的服务依次停止;
  • 最终让 event loop 退出。

它不是 CANopen 协议通信对象,而是应用生命周期和进程退出控制组件。

2.7 ev::Loop:它不登记设备 fd,但负责真正驱动 Poll

Poll 只实现了 ev_poll_t 接口:

1
2
3
self()
wait(timeout)
kill(thread_token)

ev_loop 保存这个接口:

1
2
3
4
struct ev_loop {
ev_poll_t *poll;
/* task queue, executor, contexts... */
};

当任务队列中存在真实任务时,event loop 优先执行任务:

1
ev_exec_run(task->exec, task);

当没有任务可执行、并且允许当前线程进入 Poll 时,才调用:

1
2
3
int result = ev_poll_wait(
loop->poll,
empty ? -1 : 0);

含义是:

event loop 状态 Poll timeout
任务队列为空 -1,允许阻塞等待 I/O
已经还有任务 0,只做一次非阻塞 I/O 检查

当其他线程向 event loop 投递任务,需要把正在 epoll_pwait() 中睡眠的线程叫醒时,event loop 调用:

1
2
3
return ctx->polling
? ev_poll_kill(ctx->loop->poll, ctx->thr)
: 0;

所以二者分工为:

1
2
3
4
Poll:知道 Linux epoll 和 fd 事件;
Loop:知道任务队列、executor、何时阻塞、何时唤醒;
设备组件:知道 fd 就绪后应做什么;
CANopen 上层:知道收到 CAN 帧或定时器到期后应执行哪个协议动作。

2.8 io::CanNet:Poll 驱动 I/O 与被动 CANopen 核心之间的桥梁

io_can_net 不包含 io_poll_t *,也不调用 io_poll_watch()。它持有的是:

1
2
3
4
io_timer_t *timer;
io_tqueue_t *tq;
io_can_chan_t *chan;
can_net_t *net;

它把三个不同世界连接起来:

下层 中间桥梁 上层
Poll 驱动的 CanChannel io_can_net_read_func() / write_func() 被动 can_net_t
Poll 驱动的 Timer io_tqueue / wait_next_func() can_net_set_time() 和 timer heap
CANopen 发送请求 io_can_net_send_func() CanChannel 异步写操作

2.8.1 接收链路

io_can_net_read_func() 在异步 CAN read 完成后执行:

1
2
3
4
if (read->r.result == 1) {
io_can_net_set_time(net);
can_net_recv(net->net, &net->read_msg);
}

随后它再次提交读请求:

1
io_can_chan_submit_read(net->chan, &net->read);

完整链路为:

1
2
3
4
5
6
SocketCAN fd
→ Poll
→ CanChannel rx/read task
→ io_can_net_read_func()
→ can_net_recv()
→ CANopen NMT/SDO/PDO/SYNC/EMCY 等接收处理器

最后一层具体进入哪些通信对象,取决于 can_net_t 中注册的接收者和当前编译配置。

2.8.2 定时链路

初始化时,io_can_net 注册:

1
2
can_net_set_next_func(net->net,
&io_can_net_next_func, net);

当 CAN 网络核心的“最近到期时间”发生变化时,io_can_net_next_func() 把该绝对时间提交到 io_tqueue

Timer 到期后,io_can_net_wait_next_func() 执行:

1
io_can_net_set_time(net);

io_can_net_set_time() 最终调用:

1
can_net_set_time(io_can_net_get_net(net), &now);

这会让被动 CAN 网络核心检查 timer heap,并执行所有已经到期的协议定时器。

完整链路:

1
2
3
4
5
6
7
8
9
CANopen/CAN timer heap 更新最近到期时间
→ io_can_net_next_func()
→ io_tqueue_submit_wait()
→ Timer 设置 timerfd
→ Poll 等待 timerfd 可读
→ Timer wait_task
→ io_can_net_wait_next_func()
→ can_net_set_time()
→ 执行到期协议回调

2.8.3 发送链路

CAN/CANopen 核心需要发送帧时调用已注册的 send callback:

1
2
can_net_set_send_func(net->net,
&io_can_net_send_func, net);

io_can_net_send_func() 把 CAN 帧放入用户态发送环形队列,随后 io_can_net_do_write() 调用:

1
io_can_chan_submit_write(net->chan, &net->write);

若 SocketCAN 暂时不可写,则回到前面所述的 CanChannel IO_EVENT_OUT + Poll 重试机制。

完整链路:

1
2
3
4
5
6
7
8
CANopen 协议对象请求发送
→ can_net_t send callback
→ io_can_net_send_func()
→ io_can_net 发送队列
→ io_can_chan_submit_write()
→ SocketCAN write()
→ 若 EAGAIN:Poll 等待 IO_EVENT_OUT
→ write_task 重试

2.9 NodeMasterSlave 和 CANopen Driver 到底怎样依赖 Poll

2.9.1 canopen::Node

Node 构造函数首先构造 io::CanNet

1
2
3
4
5
6
7
8
9
10
Node::Node(ev_exec_t* exec,
io::TimerBase& timer,
io::CanChannelBase& chan,
__co_dev* dev,
uint8_t id)
: io::CanNet(exec, timer, chan, 0, 0),
Device(dev, id, this),
impl_(new Impl_(this, net(), Device::dev())) {
start();
}

因此 Node 的 I/O 入口已经由 CanNet 提供。start() 后,CanNet 持续提交异步 CAN read;收到帧后进入 can_net_recv(),再由 Node 内部 NMT、PDO、SYNC、TIME、EMCY 等对象接收。

Node 本身不需要知道 SocketCAN fd,也不调用 io_poll_watch()

2.9.2 BasicMaster / AsyncMaster

Master 构造函数继续把 timer 和 channel 交给 Node:

1
2
3
4
5
6
7
8
BasicMaster::BasicMaster(ev_exec_t* exec,
io::TimerBase& timer,
io::CanChannelBase& chan,
__co_dev* dev,
uint8_t id)
: Node(exec, timer, chan, dev, id) {
/* master-specific initialization */
}

因此 Master 对 Poll 的实际依赖链是:

1
2
3
4
5
BasicMaster / AsyncMaster
→ Node
→ CanNet
├─ CanChannel → Poll → SocketCAN fd
└─ Timer → Poll → timerfd

SDO 请求、NMT boot、heartbeat 监控等上层行为并不是直接向 Poll 登记事件,而是创建 CAN 操作或协议 timer,最终通过上述链路落到底层 I/O。

2.9.3 BasicSlave

Slave 与 Master 一样建立在 Node 和 CanNet 上。它接收 SDO、NMT、RPDO 等帧时,最底层仍由同一个 CanChannel → Poll 接收链路驱动;协议定时行为仍由 Timer → Poll 驱动。

2.9.4 CANopen “Driver” 不是 Linux fd 驱动

Lely coapp 中的:

1
2
3
4
DriverBase
BasicDriver
LoopDriver
FiberDriver

这里的 Driver 表示“远程 CANopen 节点的应用逻辑驱动/处理器”,不是 Linux 字符设备驱动,也不是负责 SocketCAN fd 的 I/O driver。

它们的主要职责是:

  • 接收 Master 分发的远程节点状态事件;
  • 执行从站配置;
  • 发起和等待 SDO 请求;
  • 处理 RPDO/TPDO、heartbeat、boot 等上层回调;
  • 在需要时使用自己的 strand、fiber 或局部 event loop 串行化业务逻辑。

它们不直接调用 io_poll_watch()。事件在到达 Driver 之前已经经过:

1
2
3
4
5
Poll
→ CanChannel / Timer task
→ CanNet
→ CANopen Node/Master
→ Driver callback

LoopDriver 虽然内部有一个 event loop 和线程,但它的 loop 主要用于执行该远程节点 Driver 的异步任务。它不会重新监听 CAN fd;CAN fd 仍归 Master 共享的 CanChannel 和主 event loop 管理。

2.10 三条最重要的完整运行流程

2.10.1 收到一帧 TPDO、SDO 响应或 heartbeat

1
2
3
4
5
6
7
8
9
10
11
1. CAN 控制器收到帧,SocketCAN raw socket 变为可读
2. epoll 返回 EPOLLIN
3. io_poll 转换为 IO_EVENT_IN
4. io_can_chan_impl_watch_func() 投递 rxbuf_task
5. rxbuf_task 读取 SocketCAN 帧并填入接收队列
6. CanChannel 完成 io_can_chan_read
7. io_can_net_read_func() 更新时间并调用 can_net_recv()
8. can_net_t 根据 CAN-ID 分发给已注册的 CAN/CANopen 接收对象
9. Node/Master 触发 PDO、SDO、NMT、heartbeat 等对应回调
10. coapp Driver 或用户回调在 executor 上运行
11. CanChannel 再次等待下一帧;读到 EAGAIN 后重新登记 IO_EVENT_IN

2.10.2 CANopen 协议定时器到期

1
2
3
4
5
6
7
8
9
10
11
1. CANopen 对象在 can_net_t timer heap 中更新下一到期时间
2. io_can_net_next_func() 把时间提交给 io_tqueue
3. Timer 设置 timerfd
4. timerfd 到期并变为可读
5. Poll 返回 IO_EVENT_IN
6. Timer watch callback 投递 wait_task
7. wait_task 读取 timerfd 到期计数
8. io_tqueue 完成对应 wait
9. io_can_net_wait_next_func() 调用 can_net_set_time()
10. can_net_t 执行已到期 timer callback
11. 若仍有下一到期时间,重新设置 timerfd 并恢复 Poll 监听

2.10.3 Ctrl+C 退出程序

1
2
3
4
5
6
7
8
9
10
11
1. 进程收到 SIGINT
2. SignalSet signal handler 设置 pending,并写 self-pipe
3. self-pipe 读端变为可读
4. Poll 返回 IO_EVENT_IN
5. SignalSet watch callback 投递 read_task
6. read_task 读空 pipe,确认具体 signal
7. 异步 signal wait 完成
8. 应用调用 Context shutdown / Loop stop
9. CanChannel、Timer、SignalSet、CanNet 取消等待和登记
10. ev_poll_kill() 唤醒仍在 epoll_pwait() 中的 event-loop 线程
11. loop.run() 返回,应用进行析构

2.11 为什么 Poll 可以被多个组件共享

同一个 Poll 可以同时维护多个 fd:

1
2
3
4
5
6
SocketCAN channel A fd
SocketCAN channel B fd
timerfd A
timerfd B
SignalSet self-pipe fd
其他可轮询 I/O fd

每个 fd 都有独立的 io_poll_watch,而 Poll 内部红黑树按 fd 建立映射。内核返回某个 fd 后,Poll 能找到其所属组件的 watch callback。

这也是红黑树存在的实际应用背景:不仅一个 CAN channel 会使用 Poll;一个应用可以共享同一个 Poll 给多个 CAN channel、多个 timer 和一个 SignalSet。随着节点或网络数量增加,Poll 中登记的 fd 数量也会增加。

但需要注意:

  • 一个具体 fd 在同一 Poll 中只能对应一个 watch;
  • 每个 CanChannelTimerSignalSet 对象都持有自己的 watch;
  • EPOLLONESHOT 使每个组件必须在自己的任务中明确恢复下一次监听;
  • 多个 CANopen Node 可以共享 Poll、Loop 和 Controller,但通常需要各自独立的 Timer 和 CanChannel,以保持异步操作状态相互独立。

2.12 直接与间接依赖总表

层次 模块/组件 直接调用 io_poll_watch() 被监听对象 Poll 事件后的动作
平台 I/O io_can_chan_impl SocketCAN fd 投递接收或写重试任务
平台 I/O io_timer_impl timerfd 投递定时 wait 任务
平台 I/O io_sigset_impl self-pipe read fd 投递 pipe 读取和 signal 分发任务
Event ev_loop 不拥有具体设备 fd 调用 ev_poll_wait(),执行 executor task,必要时 ev_poll_kill()
I/O 适配 io_can_net 通过 Timer/CanChannel 间接依赖 把 CAN 帧、发送请求和时间同步给 can_net_t
CAN 核心 can_net_t 无系统 fd 被动接收帧、维护接收者和 timer heap、请求发送
CANopen 应用 canopen::Node 无系统 fd 构造协议对象并启动 CanNet 收发
CANopen 应用 BasicMaster / AsyncMaster 无系统 fd 主站 NMT、SDO、节点管理等
CANopen 应用 BasicSlave 无系统 fd 从站协议处理
节点业务 DriverBase / BasicDriver 无系统 fd 处理远程节点事件和业务逻辑
节点业务 LoopDriver / FiberDriver 不重新监听 CAN fd 提供业务任务的同步/串行执行环境

2.13 阅读这些模块时的推荐源码顺序

要从 Poll 一直读到 CANopen 主站,推荐按下面顺序:

  1. include/lely/io2/posix/poll.h:理解 watch 的一次性回调契约;
  2. src/io2/linux/poll.c:理解 epoll、红黑树、wait 和 kill;
  3. src/ev/loop.c:理解谁调用 ev_poll_wait(),谁运行任务;
  4. src/io2/linux/can_chan.c:理解 SocketCAN 的异步读写如何登记 IN/OUT;
  5. src/io2/linux/timer.c:理解 timerfd 如何产生异步定时完成;
  6. src/io2/posix/sigset.c:理解 self-pipe 信号转换;
  7. src/io2/can_net.c:理解 CanChannel/Timer 如何接入被动 can_net_t
  8. src/coapp/node.cpp:理解 Node 如何基于 CanNet 启动 CANopen 协议对象;
  9. src/coapp/master.cpp:理解 Master 如何在 Node 之上增加远程节点管理;
  10. src/coapp/driver.cpploop_driver.cpp 或对应头文件:理解远程节点业务 Driver 如何消费上层事件。

按照这个顺序,不会把以下三个层次混在一起:

1
2
3
Linux fd 是否就绪
CAN 帧如何异步收发
CANopen 协议收到帧后做什么

3. 阅读源码前必须先懂的 Linux epoll

3.1 epoll 解决什么问题

进程可能同时管理很多文件描述符,例如:

  • SocketCAN socket;
  • TCP/UDP socket;
  • timerfd
  • eventfd
  • pipe;
  • 串口或其他可轮询设备。

如果每个 fd 都用一个线程阻塞读取,线程数量和切换成本会很高。epoll 允许一个线程统一等待多个 fd

1
2
3
4
5
6
7
8
fd 3: SocketCAN
fd 4: timerfd
fd 7: TCP socket
fd 9: eventfd
↓ 全部登记到同一个 epoll 实例
epoll_pwait()

内核只返回当前已经就绪的 fd

Linux epoll 的核心概念可以简化为两个集合:

内核集合 含义
interest list 调用方已经登记、希望监控的 fd 集合
ready list interest list 中当前已经发生事件、可以返回给调用方的 fd 集合

3.2 epfd 是什么

源码通过:

1
poll->epfd = epoll_create1(EPOLL_CLOEXEC);

创建一个 epoll 实例。返回值 epfd 本身也是一个文件描述符,但它代表的是“内核中的 epoll 管理对象”,不是被监控的 SocketCAN 或 timerfd。

EPOLL_CLOEXEC 表示执行 exec() 装载新程序时自动关闭该 epfd,避免文件描述符泄漏到新程序。

3.3 struct epoll_event 的两个关键字段

源码宏最终构造的是:

1
2
3
4
struct epoll_event event = {
.events = /* EPOLLIN、EPOLLOUT、EPOLLONESHOT 等位掩码 */,
.data.fd = fd
};

含义:

字段 用途
event.events 告诉内核要关注哪些事件,以及采用什么行为模式
event.data.fd 与该登记项绑定的用户数据;事件返回时内核原样带回来

Lely 只把 fd 放入 data.fd,没有把 watch* 直接放入 data.ptr。因此内核返回事件后,Lely 还需要用 fd 去红黑树查找对应的 io_poll_watch

3.4 epoll_ctl():管理 interest list

函数原型可理解为:

1
int epoll_ctl(int epfd, int op, int fd, struct epoll_event *event);

三个操作是理解 io_poll_watch() 的基础:

操作 含义 Lely 中的使用场景
EPOLL_CTL_ADD 把一个新 fd 加入 epoll interest list 第一次监听该 fd
EPOLL_CTL_MOD 修改该 fd 的事件掩码;对于 EPOLLONESHOT,同时恢复下一次监听 修改关注类型,或一次通知后继续监听
EPOLL_CTL_DEL 从 interest list 删除该 fd 不再监听该 fd

需要严格区分:

1
2
3
ADD:此前内核中没有该登记项
MOD:此前已有登记项,只更新或恢复
DEL:彻底删除登记项

3.5 epoll_pwait():等待事件

源码调用:

1
2
3
4
5
6
int nevents = epoll_pwait(
poll->epfd,
events,
LELY_IO_EPOLL_MAXEVENTS,
timeout,
&set);

参数含义:

参数 含义
poll->epfd 等待哪个 epoll 实例
events 输出数组,内核把已经发生的事件写入这里
LELY_IO_EPOLL_MAXEVENTS 本次数组最多容纳多少个事件
timeout 等待时间,单位为毫秒;-1 无限等待,0 立即返回
&set 等待期间临时使用的信号掩码

返回值:

返回值 含义
> 0 本次返回的事件数量
0 超时,没有事件
-1 系统调用失败,错误在 errno 中;被信号打断时通常为 EINTR

epoll_pwait()epoll_wait() 的主要区别是:它能在“进入等待”的同时原子地切换信号掩码,用于避免检查状态与进入阻塞之间的丢失唤醒竞态。

3.6 本实现是“电平触发 + 一次性通知”

宏没有加入 EPOLLET,所以 epoll 使用默认的 level-triggered,即电平触发模式;同时又加入 EPOLLONESHOT

组合后的实际行为:

1
2
3
4
5
只要 fd 仍满足可读/可写条件,电平触发本来可以持续报告
+
EPOLLONESHOT 规定每次恢复监听后最多只报告一次
=
收到一次通知后暂停该 fd,调用方处理完后显式决定是否继续

如果回调没有消除就绪条件,例如 socket 中仍有未读数据,那么再次 EPOLL_CTL_MOD 恢复监听后,可能立刻再次收到事件。这是正常的电平触发行为。


4. IO_EVENT_*EPOLL* 到底是什么关系

4.1 两套事件名称属于不同层

名称前缀 所属层 用途
IO_EVENT_* Lely 跨平台抽象层 上层设备代码使用,不直接依赖 Linux
EPOLL* Linux epoll API Linux 内核能识别的事件和控制标志

因此,IO_EVENT_IN 不是 EPOLLIN 的别名,而是 Lely 的通用“读侧事件”抽象。Linux 后端负责把它转换成 EPOLLIN 等内核标志。

4.2 事件转换宏逐项解释

源码:

1
2
3
4
5
6
7
8
#define EPOLL_EVENT_INIT(events, fd) \
{ \
(((events) & IO_EVENT_IN) ? (EPOLLIN | EPOLLRDHUP) : 0) \
| (((events) & IO_EVENT_PRI) ? EPOLLPRI : 0) \
| (((events) & IO_EVENT_OUT) ? EPOLLOUT : 0) \
| EPOLLONESHOT, \
{ .fd = (fd) } \
}

它构造一个 struct epoll_event,第一部分是 event.events,第二部分是 event.data.fd

4.2.1 IO_EVENT_IN

Lely 含义:调用方关注“读侧有事件”。通常意味着应该尝试 read()recv() 或设备对应的读取操作。

它在 Linux 后端被映射为:

1
EPOLLIN | EPOLLRDHUP

4.2.2 EPOLLIN

Linux 含义:关联对象当前可进行读取操作。

需要注意:

  • 它表示“应该尝试读取”,不是保证一定能读取到完整业务帧;
  • 多线程或其他代码也可能先一步读走数据,因此非阻塞读取仍应处理 EAGAIN
  • 对 pipe 或 stream socket,读侧 EOF/关闭也可能伴随读事件,最终应以 read() 返回值判断;
  • 对 SocketCAN,收到 CAN 帧后通常表现为 EPOLLIN

4.2.3 EPOLLRDHUP

Linux 含义:面向连接的流式 socket 对端关闭了写方向,或者关闭了连接。

Lely 把它也映射为 IO_EVENT_IN

1
2
if (events[i].events & (EPOLLIN | EPOLLRDHUP))
revents |= IO_EVENT_IN;

原因是读侧处理函数需要被唤醒,然后通过 read()/recv() 判断数据、EOF 或断开状态。

对于 SocketCAN 这类非 TCP 流式连接,EPOLLRDHUP 通常不是主要事件;但宏是通用 POSIX I/O 后端,所以统一加入。

4.2.4 IO_EVENT_PRI

Lely 含义:调用方关注“高优先级或异常条件事件”。

它被映射为:

1
EPOLLPRI

4.2.5 EPOLLPRI

Linux 含义:文件描述符出现 exceptional condition,即异常优先级条件。典型例子包括某些 socket 的带外数据或特定设备的紧急状态。

它不是普通输入数据:

1
2
普通可读数据 → EPOLLIN
异常/紧急条件 → EPOLLPRI

是否使用 IO_EVENT_PRI 取决于具体 fd 类型。仅从 poll.c 无法断定某个 CANopen 设备模块会不会登记该事件。

4.2.6 IO_EVENT_OUT

Lely 含义:调用方关注“写侧可以继续推进”。通常用于非阻塞 socket 发送缓冲区有空间、异步连接完成等情况。

它被映射为:

1
EPOLLOUT

4.2.7 EPOLLOUT

Linux 含义:关联对象当前可执行写操作。

它不保证一次 write() 就能写完全部数据,只说明至少可以尝试推进写操作。正确的非阻塞代码仍要处理:

  • 部分写入;
  • EAGAIN
  • 连接错误;
  • 待发送缓冲区剩余数据。

4.2.8 “所有登记强制附加 EPOLLONESHOT

宏最后无条件执行:

1
| EPOLLONESHOT

因此,只要 events != 0 并进入 ADD 或 MOD,生成的 Linux 事件掩码一定包含 EPOLLONESHOT,无论调用方登记的是:

1
2
3
4
5
IO_EVENT_IN
IO_EVENT_PRI
IO_EVENT_OUT
IO_EVENT_IN | IO_EVENT_OUT
或其他受支持组合

EPOLLONESHOT 的准确含义是:

  1. 内核返回该 fd 的一次事件;
  2. fd 的 epoll 登记项仍存在,但暂时不再继续报告事件;
  3. 调用方必须执行 epoll_ctl(..., EPOLL_CTL_MOD, ...) 才能恢复下一次通知。

本文把这个过程称为:

1
2
3
第一次 ADD/MOD 后:当前一次监听有效
事件返回后:本次监听已消费
再次 MOD 后:恢复下一次监听

4.2.9 为什么宏没显式加入 EPOLLERREPOLLHUP

Linux 会始终报告 EPOLLERREPOLLHUP,即使调用 epoll_ctl() 时没有把它们写入事件掩码。因此源码只在事件返回阶段进行转换:

1
2
3
4
if (events[i].events & EPOLLERR)
revents |= IO_EVENT_ERR;
if (events[i].events & EPOLLHUP)
revents |= IO_EVENT_HUP;

5. 核心对象以及本文使用的状态术语

5.1 struct io_poll

1
2
3
4
5
6
7
8
9
10
11
12
struct io_poll {
struct io_svc svc;
const struct ev_poll_vtbl *poll_vptr;
io_ctx_t *ctx;
int signo;
struct sigaction oact;
sigset_t oset;
int epfd;
pthread_mutex_t mtx;
struct rbtree tree;
size_t nwatch;
};
字段 含义
svc 注册到 io_ctx 的服务对象,用于接收 fork 通知
poll_vptr ev_poll_t 虚函数表入口,包含 self/wait/kill
ctx 所属 I/O context
signo 跨线程唤醒 epoll_pwait() 使用的信号
oact 保存原信号处理函数
oset 保存原信号掩码
epfd Linux epoll 实例文件描述符
mtx 保护内部红黑树、计数和停止状态的互斥锁
tree fd 为键保存 io_poll_watch 的红黑树
nwatch 当前仍处于“一次监听有效”状态的 watch 数量

5.2 struct io_poll_watch

1
2
3
4
5
6
struct io_poll_watch {
io_poll_watch_func_t *func;
int _fd;
struct rbnode _node;
int _events;
};
字段 含义
func 事件发生时调用的回调函数
_fd 该 watch 对应的文件描述符,同时作为红黑树键
_node 嵌入在 watch 内部的红黑树节点
_events Lely 认为当前有效的事件掩码;0 表示当前没有有效的一次监听

5.3 watch 的三种监听状态

状态 红黑树中是否存在 _events 内核状态 含义
未登记 通常为 0 epoll 中无该登记项 io_poll 不再管理该 fd
已登记且当前一次监听有效 非 0 ADD/MOD 后可返回一次事件 正在等待下一次通知
已登记,但本次监听已消费 0 EPOLLONESHOT 已暂停该项 关联关系还在,但必须 MOD 才能继续收到通知

第三种状态最容易误解:红黑树节点还在,不代表当前还能收到事件。它只是保留了 fd → watch 关联,便于后续通过 EPOLL_CTL_MOD 恢复监听,或者通过 EPOLL_CTL_DEL 彻底删除。


6. 为什么需要红黑树

6.1 先回答“是不是因为可能挂载很多事件”

是,但这只是原因之一。

epoll 本身适合管理大量文件描述符,Lely 选择红黑树也使用户态的 fd → watch 查找在监控对象较多时仍保持稳定的 O(log N) 复杂度。不过,从这三个文件不能推断实际运行时一定会有很多 fd;这是面向通用 I/O 库的可扩展设计。

6.2 红黑树在本实现中承担四个具体职责

6.2.1 由内核返回的 fd 找到回调对象

内核只带回:

1
int fd = events[i].data.fd;

Lely 随后执行:

1
struct rbnode *node = rbtree_find(&poll->tree, &fd);

找到节点后,再通过:

1
io_poll_watch_from_node(node)

恢复到包含它的 struct io_poll_watch

6.2.2 保证同一 fd 只能绑定一个 watch 对象

源码:

1
2
3
4
if (node && node != &watch->_node) {
errsv = EEXIST;
goto error;
}

如果树中已经有这个 fd,但对应的节点不是当前传入的 watch->_node,说明调用方试图用另一个 watch 接管同一个 fd,函数返回 EEXIST

6.2.3 过滤已经取消登记的陈旧事件

并发情况下,内核事件数组中可能已经包含某个 fd,但另一个线程随后取消了登记。事件分发时再次查树:

1
2
3
if (node) {
...
}

如果节点已经不存在,就忽略该事件,不会直接解引用一个可能已经失效的 watch*

这也是为什么实现选择 event.data.fd + 红黑树,而不是简单地把裸 watch* 放在 event.data.ptr 后直接调用。

6.2.4 fork 后遍历并重建 epoll 登记项

子进程收到 IO_FORK_CHILD 后,源码遍历红黑树,对仍有 _events 的 watch 重新执行 EPOLL_CTL_ADD

6.3 为什么不是数组、链表或哈希表

结构 优点 对本实现的限制
数组按 fd 直接索引 查找 O(1) fd 数值可能稀疏且很大,浪费空间
链表 实现简单 每次按 fd 查找 O(N)
哈希表 平均 O(1) 需要容量、扩容和哈希策略,最坏情况不稳定
红黑树 最坏 O(log N),支持有序遍历,适合稀疏整数键 实现比链表复杂

上述“为什么选红黑树”属于基于数据结构和调用方式的设计分析;源码没有注释直接说明作者的完整选型理由。


7. io_fd_cmp() 实际比较的是什么

源码:

1
2
3
4
5
6
7
8
9
10
static int
io_fd_cmp(const void *p1, const void *p2)
{
assert(p1);
int fd1 = *(const int *)p1;
assert(p2);
int fd2 = *(const int *)p2;

return (fd2 < fd1) - (fd1 < fd2);
}

7.1 p1p2 都指向一个 int fd

红黑树初始化时:

1
rbtree_init(&poll->tree, &io_fd_cmp);

插入节点时:

1
2
watch->_fd = fd;
rbnode_init(&watch->_node, &watch->_fd);

所以节点的键是 &watch->_fd

查找时:

1
rbtree_find(&poll->tree, &fd);

传入的查找键也是一个 int 的地址。

因此,io_fd_cmp() 比较的不是:

  • watch 地址;
  • rbnode 地址;
  • 事件掩码;
  • epoll 实例 epfd

它只比较两个整数文件描述符 fd1fd2

7.2 返回表达式如何工作

1
return (fd2 < fd1) - (fd1 < fd2);

C 中比较表达式结果为 01

条件 第一项 第二项 返回值 比较结果
fd1 < fd2 0 1 -1 第一个键更小
fd1 == fd2 0 0 0 两个键相等
fd1 > fd2 1 0 1 第一个键更大

例如:

1
2
3
fd1=3, fd2=5 → (5<3)-(3<5) = 0-1 = -1
fd1=5, fd2=3 → (3<5)-(5<3) = 1-0 = 1
fd1=4, fd2=4 → 0-0 = 0

这等价于普通升序比较,但避免写成:

1
return fd1 - fd2;

后者对通用整数键可能存在减法溢出风险。


8. 完整使用周期第 1 步:创建 io_poll

8.1 C++ 入口:lely::io::Poll::Poll()

poll.hpp

1
2
3
Poll(ContextBase& ctx, int signo = 0) : poll_(io_poll_create(ctx, signo)) {
if (!poll_) util::throw_errc("Poll");
}

C++ 类只是 RAII 包装:

  • 构造时调用 io_poll_create()
  • 析构时调用 io_poll_destroy()
  • watch() 转调 io_poll_watch()
  • get_poll() 返回通用 ev::Poll 接口。

实际 epoll 算法全部在 poll.c

8.2 io_poll_create()

1
2
3
io_poll_t *poll = io_poll_alloc();
...
io_poll_t *tmp = io_poll_init(poll, ctx, signo);

职责:

  1. malloc(sizeof(io_poll_t))
  2. 调用 io_poll_init()
  3. 初始化失败时释放内存并保留错误码。

8.3 io_poll_init()

关键步骤按执行顺序如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
设置 io_svc 虚函数表

保存 ctx

设置 ev_poll self/wait/kill 虚函数表

选择唤醒信号,signo=0 时使用 SIGUSR1

安装空信号处理函数 sig_ign

正常执行期间阻塞该唤醒信号

初始化 mutex

初始化红黑树,比较函数为 io_fd_cmp

nwatch=0

epoll_create1(EPOLL_CLOEXEC)

注册到 io_ctx

初始化完成后的关键状态:

1
2
3
poll->epfd   = 有效 epoll fd
poll->tree = 空
poll->nwatch = 0

9. 完整使用周期第 2 步:初始化 io_poll_watch

头文件提供:

1
2
3
4
#define IO_POLL_WATCH_INIT(func) \
{ \
(func), -1, RBNODE_INIT, 0 \
}

例如:

1
2
3
static void on_event(struct io_poll_watch *watch, int events);

struct io_poll_watch watch = IO_POLL_WATCH_INIT(&on_event);

初始值:

1
2
3
4
watch.func    = on_event
watch._fd = -1
watch._node = 未插入红黑树
watch._events = 0

调用方应该把 io_poll_watch 嵌入自己的设备对象中,以便在回调中通过 structof() 找回上层对象。该 watch 的内存必须一直存活到取消登记完成。


10. 完整使用周期第 3 步:调用 io_poll_watch() 登记 fd

这是整个文件最关键的函数。

10.1 函数入口和参数检查

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
int
io_poll_watch(io_poll_t *poll, int fd, int events,
struct io_poll_watch *watch)
{
assert(poll);
assert(watch);
int epfd = poll->epfd;

if (fd == -1 || fd == epfd) {
errno = EBADF;
return -1;
}

if (events < 0) {
errno = EINVAL;
return -1;
}
events &= IO_EVENT_MASK;

逐项解释:

  1. fd == -1:无效文件描述符;
  2. fd == epfd:不允许把当前 epoll 实例自身作为目标 fd;
  3. events < 0:事件掩码不合法;
  4. events &= IO_EVENT_MASK:只保留 Lely 支持的事件位,其他位被清除。IO_EVENT_MASK 的具体位值定义不在当前三个文件中。

10.2 加锁和红黑树查找

1
2
3
4
5
6
7
pthread_mutex_lock(&poll->mtx);

struct rbnode *node = rbtree_find(&poll->tree, &fd);
if (node && node != &watch->_node) {
errsv = EEXIST;
goto error;
}

此时:

  • node == NULL:该 fd 从未登记,或者已经彻底取消登记;
  • node == &watch->_node:该 fd 已经由当前 watch 管理;
  • node != NULL && node != &watch->_node:同一个 fd 已绑定另一个 watch,返回 EEXIST

10.3 核心代码原文

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
if (events) {
struct epoll_event event = EPOLL_EVENT_INIT(events, fd);
if (node && events != watch->_events) {
if (epoll_ctl(epfd, EPOLL_CTL_MOD, fd, &event) == -1) {
errsv = errno;
epoll_ctl(epfd, EPOLL_CTL_DEL, fd, NULL);
rbtree_remove(&poll->tree, node);
if (watch->_events)
poll->nwatch--;
watch->_events = 0;
goto error;
}
} else if (!node) {
if (epoll_ctl(epfd, EPOLL_CTL_ADD, fd, &event) == -1) {
errsv = errno;
goto error;
}
watch->_fd = fd;
rbnode_init(&watch->_node, &watch->_fd);
watch->_events = 0;
rbtree_insert(&poll->tree, &watch->_node);
}
if (!watch->_events)
poll->nwatch++;
watch->_events = events;
} else if (node) {
epoll_ctl(epfd, EPOLL_CTL_DEL, fd, NULL);
rbtree_remove(&poll->tree, node);
if (watch->_events)
poll->nwatch--;
watch->_events = 0;
}

10.4 先看总决策树

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
flowchart TD
A[进入 io_poll_watch] --> B[校验 fd 和 events]
B --> C[加锁并按 fd 查红黑树]
C --> D{同 fd 是否属于另一个 watch?}
D -- 是 --> X[返回 EEXIST]
D -- 否 --> E{events 是否非 0?}
E -- 否 --> F{node 是否存在?}
F -- 否 --> G[什么也不做,返回成功]
F -- 是 --> H[EPOLL_CTL_DEL]
H --> I[红黑树删除]
I --> J[必要时 nwatch--]
J --> K[_events=0]
E -- 是 --> L[构造 epoll_event]
L --> M{node 是否存在?}
M -- 否 --> N[EPOLL_CTL_ADD]
N --> O[保存 fd并初始化 rbnode]
O --> P[插入红黑树]
M -- 是 --> Q{events 与旧 _events 是否不同?}
Q -- 是 --> R[EPOLL_CTL_MOD]
Q -- 否 --> S[内核配置不变]
R --> T{MOD 是否成功?}
T -- 否 --> U[DEL + 树删除 + 状态清理 + 返回错误]
T -- 是 --> V[更新状态]
P --> V
S --> V
V --> W{旧 _events 是否为 0?}
W -- 是 --> Y[nwatch++]
W -- 否 --> Z[nwatch 不变]
Y --> AA[_events=events]
Z --> AA

11. io_poll_watch() 分支一:首次登记,执行 EPOLL_CTL_ADD

触发条件:

1
2
events != 0
node == NULL

对应源码:

1
2
3
4
5
6
7
8
9
10
11
12
13
} else if (!node) {
if (epoll_ctl(epfd, EPOLL_CTL_ADD, fd, &event) == -1) {
errsv = errno;
goto error;
}
watch->_fd = fd;
rbnode_init(&watch->_node, &watch->_fd);
watch->_events = 0;
rbtree_insert(&poll->tree, &watch->_node);
}
if (!watch->_events)
poll->nwatch++;
watch->_events = events;

按时间顺序解释:

11.1 构造内核事件描述

1
struct epoll_event event = EPOLL_EVENT_INIT(events, fd);

假设调用:

1
io_poll_watch(poll, fd, IO_EVENT_IN, watch);

则宏生成的事件近似为:

1
2
event.events = EPOLLIN | EPOLLRDHUP | EPOLLONESHOT;
event.data.fd = fd;

11.2 EPOLL_CTL_ADD

1
epoll_ctl(epfd, EPOLL_CTL_ADD, fd, &event);

内核把该 fd 加入 epfd 的 interest list。

失败时:

  • 不修改 watch->_fd
  • 不插入红黑树;
  • 不增加 nwatch
  • 返回 -1errno 保留为 epoll_ctl() 的错误。

11.3 建立用户态关联

ADD 成功后:

1
2
3
4
watch->_fd = fd;
rbnode_init(&watch->_node, &watch->_fd);
watch->_events = 0;
rbtree_insert(&poll->tree, &watch->_node);

作用:

1
2
内核:epfd 已经管理 fd
用户态:红黑树建立 fd → watch 映射

11.4 更新 nwatch_events

1
2
3
if (!watch->_events)
poll->nwatch++;
watch->_events = events;

因为首次登记时 _events 被设为 0,所以 nwatch++

完成后的状态:

1
2
3
4
红黑树:有节点
watch->_events:非 0
nwatch:增加 1
内核:当前一次监听有效,最多返回一次事件

12. io_poll_watch() 分支二:修改事件或恢复下一次监听,执行 EPOLL_CTL_MOD

触发条件:

1
2
3
events != 0
node != NULL
events != watch->_events

源码:

1
2
3
4
5
if (node && events != watch->_events) {
if (epoll_ctl(epfd, EPOLL_CTL_MOD, fd, &event) == -1) {
...
}
}

这个分支包含两种不同场景。

12.1 场景 A:真正修改关注事件

原来:

1
watch->_events = IO_EVENT_IN

现在调用:

1
events = IO_EVENT_OUT

因为不同,执行 MOD,把内核事件掩码从读侧关注改成写侧关注。

12.2 场景 B:一次通知后继续监听

事件处理前,io_poll_process() 会执行:

1
2
watch->_events = 0;
poll->nwatch--;

回调若再次调用:

1
io_poll_watch(poll, fd, IO_EVENT_IN, watch);

此时:

1
2
3
node 仍存在
watch->_events == 0
events == IO_EVENT_IN != 0

所以必然进入 EPOLL_CTL_MOD

这次 MOD 不只是“修改事件类型”,更重要的是让此前因 EPOLLONESHOT 暂停的内核登记项恢复下一次通知。

12.3 MOD 成功后的计数

1
2
3
if (!watch->_events)
poll->nwatch++;
watch->_events = events;
  • 如果原 _events == 0:说明本次监听已消费,现在恢复监听,所以 nwatch++
  • 如果原 _events != 0:只是从一种有效事件配置改成另一种,数量没有变化。

12.4 MOD 失败时为什么做彻底清理

源码:

1
2
3
4
5
6
7
8
9
if (epoll_ctl(epfd, EPOLL_CTL_MOD, fd, &event) == -1) {
errsv = errno;
epoll_ctl(epfd, EPOLL_CTL_DEL, fd, NULL);
rbtree_remove(&poll->tree, node);
if (watch->_events)
poll->nwatch--;
watch->_events = 0;
goto error;
}

执行顺序:

  1. 保存 MOD 失败原因;
  2. 尝试从内核 epoll 中删除该 fd
  3. 从红黑树删除关联;
  4. 若原来仍计入 nwatch,则减一;
  5. _events = 0
  6. 返回失败。

这是“失败后不保留半有效状态”的处理:MOD 失败后,代码不再假设内核状态仍可靠,而是把该 watch 从 io_poll 管理关系中移除。

需要注意:清理用的 EPOLL_CTL_DEL 返回值被忽略,最终 errno 仍报告最初的 MOD 失败原因。


13. io_poll_watch() 分支三:重复传入相同事件

触发条件:

1
2
3
events != 0
node != NULL
events == watch->_events

这时:

1
2
3
4
5
if (node && events != watch->_events) {
...
} else if (!node) {
...
}

两个分支都不执行,因此不会调用 epoll_ctl()

随后:

1
2
3
if (!watch->_events)
poll->nwatch++;
watch->_events = events;

由于 _events 原本与 events 相同且非 0,nwatch 不变。

实际效果:这是一个幂等调用,不重复修改内核登记。

注意:在“本次监听已消费”的状态下 _events == 0,而传入 events != 0,不可能进入此分支;它一定会进入 MOD 以恢复监听。


14. io_poll_watch() 分支四:events == 0,彻底取消登记

源码:

1
2
3
4
5
6
7
} else if (node) {
epoll_ctl(epfd, EPOLL_CTL_DEL, fd, NULL);
rbtree_remove(&poll->tree, node);
if (watch->_events)
poll->nwatch--;
watch->_events = 0;
}

调用形式:

1
io_poll_watch(poll, fd, 0, watch);

14.1 EPOLL_CTL_DEL

1
epoll_ctl(epfd, EPOLL_CTL_DEL, fd, NULL);

从内核 epoll interest list 中删除该 fd

对 DEL 操作,event 参数不需要内容,因此传 NULL

14.2 从红黑树删除

1
rbtree_remove(&poll->tree, node);

删除后,后续内核即使返回一个已经进入事件数组的陈旧 fd,分发阶段也无法在树中找到 watch,因此会忽略。

14.3 nwatch 为什么有条件递减

1
2
if (watch->_events)
poll->nwatch--;

两种情况:

取消前状态 _events 是否已计入 nwatch 是否递减
当前一次监听仍有效 非 0
本次一次性监听已经消费 0 否,事件分发时已经减过

14.4 DEL 失败的源码行为

这里源码没有检查 epoll_ctl(... DEL ...) 的返回值,而是继续删除用户态红黑树节点,并最终返回成功。

这是源码客观行为。它意味着取消登记时以“清理用户态关联”为优先;但如果要诊断内核 DEL 失败,这个函数不会把该错误返回给调用方。

14.5 events == 0 && node == NULL

函数不做任何操作并返回成功,相当于允许重复取消一个已经不存在的登记。


15. io_poll_watch() 的全部场景汇总

红黑树 node 传入 events 与旧 _events 关系 epoll 操作 结果
非 0 不适用 ADD 首次建立监听
有,属于当前 watch 非 0 不同 MOD 修改事件或恢复下一次监听
有,属于当前 watch 非 0 相同 幂等,保持现状
有,属于其他 watch 任意 不适用 EEXIST
0 不适用 DEL 彻底取消登记并删除树节点
0 不适用 重复取消,返回成功

errno 处理还应注意:函数入口先保存原 errno,成功时会恢复原值。因此调用方只能在返回 -1 时读取 errno,不能因为成功后 errno 非 0 就判断失败。


16. 完整使用周期第 4 步:io_poll_poll_wait() 等待事件

上层通过 ev_poll_wait() 最终进入:

1
static int io_poll_poll_wait(ev_poll_t *poll_, int timeout)

16.1 获取当前线程状态

1
2
void *thr_ = io_poll_poll_self(poll_);
struct io_poll_thrd *thr = (struct io_poll_thrd *)thr_;

io_poll_poll_self() 为每个线程维护一个 _Thread_local 对象:

1
2
3
4
struct io_poll_thrd {
int stopped;
pthread_t *thread;
};

作用:

  • 保存实际 pthread_t,供其他线程定向发送唤醒信号;
  • stopped 表示当前这次 wait() 已经被要求停止阻塞、进入非阻塞收尾。

16.2 调用 epoll_pwait() 前释放互斥锁

1
2
pthread_mutex_unlock(&poll->mtx);
int nevents = epoll_pwait(...);

不能持锁阻塞等待,否则其他线程无法:

  • 调用 io_poll_watch() ADD/MOD/DEL;
  • 调用 io_poll_poll_kill() 修改停止状态;
  • 更新内部元数据。

16.3 timeout 语义

timeout 行为
-1 无限等待,直到 I/O 或信号发生
0 非阻塞查询,立即返回
> 0 最多等待指定毫秒数

源码在 timeout == 0 时先设置:

1
thr->stopped = 1;

从而本次只做非阻塞事件收集,不进入后续阻塞循环。

16.4 EINTR 分支

1
2
3
4
5
6
if (nevents == -1 && errno == EINTR) {
sigaddset(&set, poll->signo);
pthread_mutex_lock(&poll->mtx);
thr->stopped = 1;
continue;
}

含义:

  1. 等待被信号中断;
  2. 本次 epoll_pwait() 没有返回 I/O 事件;
  3. 下一轮把唤醒信号加入临时阻塞集合;
  4. 设置 stopped=1,使下一轮 timeout=0
  5. 再进行一次非阻塞 I/O 检查。

这样可以避免大量唤醒信号持续打断 wait,导致真正已经就绪的 I/O 长期得不到处理。


17. 完整使用周期第 5 步:把 Linux 事件转换回 Lely 事件

内核返回 struct epoll_event events[] 后,源码逐项转换:

1
2
3
4
5
6
7
8
9
10
11
int revents = 0;
if (events[i].events & (EPOLLIN | EPOLLRDHUP))
revents |= IO_EVENT_IN;
if (events[i].events & EPOLLPRI)
revents |= IO_EVENT_PRI;
if (events[i].events & EPOLLOUT)
revents |= IO_EVENT_OUT;
if (events[i].events & EPOLLERR)
revents |= IO_EVENT_ERR;
if (events[i].events & EPOLLHUP)
revents |= IO_EVENT_HUP;

映射表:

Linux 返回事件 Lely 回调事件
EPOLLIN IO_EVENT_IN
EPOLLRDHUP IO_EVENT_IN
EPOLLPRI IO_EVENT_PRI
EPOLLOUT IO_EVENT_OUT
EPOLLERR IO_EVENT_ERR
EPOLLHUP IO_EVENT_HUP

多个事件可以同时存在。例如:

1
2
3
EPOLLIN | EPOLLHUP

IO_EVENT_IN | IO_EVENT_HUP

回调不能假设 events 只有一个位。


18. 完整使用周期第 6 步:按 fd 查红黑树并找到 watch

源码:

1
2
3
4
5
6
7
8
9
int fd = events[i].data.fd;
pthread_mutex_lock(&poll->mtx);
struct rbnode *node = rbtree_find(&poll->tree, &fd);
if (node) {
struct io_poll_watch *watch = io_poll_watch_from_node(node);
io_poll_process(poll, revents, watch);
n += n < INT_MAX;
}
pthread_mutex_unlock(&poll->mtx);

流程:

1
2
3
4
5
6
7
8
9
内核返回 event.data.fd

rbtree_find(tree, &fd)

取得 rbnode

structof(node, io_poll_watch, _node)

取得 watch

io_poll_watch_from_node()

1
return structof(node, struct io_poll_watch, _node);

因为 _node 是嵌入在 io_poll_watch 内部的成员,structof() 根据成员地址反推出外层结构体地址。

如果 node == NULL,事件被忽略。常见原因是该 fd 在事件进入内核返回数组后,被另一个线程取消登记并从红黑树移除。


19. 完整使用周期第 7 步:io_poll_process() 调用回调

源码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
static void
io_poll_process(io_poll_t *poll, int revents,
struct io_poll_watch *watch)
{
assert(poll);
assert(poll->nwatch);
assert(watch);
assert(watch->_events);

watch->_events = 0;
poll->nwatch--;

if (watch->func) {
pthread_mutex_unlock(&poll->mtx);
watch->func(watch, revents);
pthread_mutex_lock(&poll->mtx);
}
}

19.1 为什么先 _events = 0

这表示:

1
2
该 fd 的本次 EPOLLONESHOT 通知已经被消费
当前不再认为它能继续产生下一次回调

它与内核行为一致:内核已经因 EPOLLONESHOT 暂停该登记项。

19.2 为什么 nwatch--

nwatch 统计的是当前仍有一次有效监听的 watch 数量。事件已经返回并准备分发,本次一次性监听已经结束,所以减一。

19.3 为什么回调前释放锁

1
2
3
pthread_mutex_unlock(&poll->mtx);
watch->func(watch, revents);
pthread_mutex_lock(&poll->mtx);

回调可能需要:

  • 再次调用 io_poll_watch() 恢复监听;
  • 调用 io_poll_watch(..., 0, ...) 取消登记;
  • 登记其他 fd;
  • 读取或写入设备;
  • 投递任务到 executor。

如果带锁调用回调,回调再次进入 io_poll_watch() 会尝试锁同一把 mutex,造成自锁;复杂回调还会长时间阻塞其他 I/O 管理操作。


20. 完整使用周期第 8 步:回调中处理 fd,并决定是否继续

一个典型读事件回调的逻辑应是:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
static void
on_read(struct io_poll_watch *watch, int events)
{
struct my_device *dev = structof(watch, struct my_device, watch);

if (events & IO_EVENT_IN) {
// 调用 read()/recv(),消费当前可读条件。
}

if (events & (IO_EVENT_ERR | IO_EVENT_HUP)) {
// 查询具体错误,决定关闭还是继续。
}

if (dev->continue_reading) {
io_poll_watch(dev->poll, dev->fd, IO_EVENT_IN, &dev->watch);
} else {
io_poll_watch(dev->poll, dev->fd, 0, &dev->watch);
}
}

这段是示意代码,具体错误处理取决于设备类型。

20.1 继续监听

再次调用:

1
io_poll_watch(poll, fd, IO_EVENT_IN, watch);

此时红黑树节点仍在,但 _events == 0,因此进入 MOD 分支:

1
2
3
4
5
6
7
EPOLL_CTL_MOD

nwatch++

watch->_events = IO_EVENT_IN

当前下一次监听重新生效

20.2 不再监听

调用:

1
io_poll_watch(poll, fd, 0, watch);

进入 DEL 分支,彻底删除内核登记和红黑树节点。

20.3 不做任何调用

如果回调返回前既没有 MOD,也没有 DEL:

1
2
3
4
红黑树节点仍存在
watch->_events == 0
内核 EPOLLONESHOT 登记项暂停
后续不会继续收到该 fd 的事件

这通常是“只收到一次,后续再也没有”的直接原因。


21. 一个 SocketCAN 可读事件的完整时序示例

以下只描述 Poll 层,不代表完整 CANopen 协议调用链:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
sequenceDiagram
participant App as CAN 设备适配层
participant Poll as io_poll
participant Tree as 红黑树
participant Kernel as Linux epoll
participant CB as watch 回调

App->>Poll: io_poll_watch(can_fd, IO_EVENT_IN, watch)
Poll->>Kernel: EPOLL_CTL_ADD<br/>EPOLLIN|EPOLLRDHUP|EPOLLONESHOT
Poll->>Tree: 插入 can_fd → watch
Poll->>Poll: nwatch++, _events=IO_EVENT_IN

App->>Poll: ev_poll_wait(timeout)
Poll->>Kernel: epoll_pwait()
Kernel-->>Poll: data.fd=can_fd, events=EPOLLIN
Poll->>Tree: rbtree_find(can_fd)
Tree-->>Poll: watch
Poll->>Poll: _events=0, nwatch--
Poll->>CB: watch->>func(watch, IO_EVENT_IN)
CB->>CB: recv/read struct can_frame
CB->>Poll: io_poll_watch(can_fd, IO_EVENT_IN, watch)
Poll->>Kernel: EPOLL_CTL_MOD<br/>恢复下一次通知
Poll->>Poll: nwatch++, _events=IO_EVENT_IN

Poll 层只知道 can_fd 可读。CAN 帧解析、COB-ID 分类、NMT/SDO/PDO/Heartbeat 状态机均由更上层代码完成。


22. nwatch 的准确含义

nwatch 不是:

  • 红黑树节点数量;
  • 已经处理过的事件数量;
  • 当前 epoll ready list 数量;
  • 所有历史登记数量。

它表示:

当前在 Lely 状态中仍有一次有效监听、理论上可以产生下一次通知的 watch 数量。

状态变化:

操作 _events 变化 nwatch 变化
首次 ADD 0 → 非0 +1
有效状态下修改事件 非0 → 非0 0
事件分发 非0 → 0 -1
回调后 MOD 恢复监听 0 → 非0 +1
有效状态下 DEL 非0 → 0 -1
已消费状态下 DEL 0 → 0 0

因此可能出现:

1
红黑树节点数 > nwatch

原因是树中可能保留“本次监听已消费但尚未恢复或删除”的 watch。


23. io_poll_poll_wait() 为什么可能连续调用多次 epoll_pwait()

每轮完成后:

1
2
thr->stopped = 1;
stopped = nevents != LELY_IO_EPOLL_MAXEVENTS;

如果本次返回数量刚好等于数组容量:

1
nevents == LELY_IO_EPOLL_MAXEVENTS

说明内核 ready list 可能还有更多事件未能装入本次数组。代码会再次循环;因为 thr->stopped == 1,下一轮把 timeout 改为 0,只做非阻塞收集。

所以一次 wait() 的策略是:

1
2
3
最多阻塞一次
+
尽量排空当时已经就绪的事件

它不会每处理一批事件后再次按原始 timeout 阻塞。


24. 跨线程唤醒:self()kill() 和信号

这部分不是正常 I/O 周期的主线,但解释了 executor 有新任务时如何让 poll 线程退出阻塞。

24.1 io_poll_poll_self()

为当前线程返回一个线程局部令牌:

1
static _Thread_local struct io_poll_thrd thr = { 0, NULL };

第一次调用时保存:

1
2
thread = pthread_self();
thr.thread = &thread;

24.2 io_poll_poll_kill() 不是杀死线程

它的真实语义是:

请求目标线程结束当前这一次阻塞等待,完成一次非阻塞收尾后从 wait() 返回。

流程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
检查目标是否为当前线程

加锁读取 target->stopped

若尚未停止,设置 stopped=1

解锁

pthread_kill(target, poll->signo)

目标线程的 epoll_pwait() 返回 -1/EINTR

目标线程再做一次 timeout=0 的 I/O 检查

wait() 返回

24.3 为什么既要 stopped 又要信号

机制 单独使用的问题
只设置 stopped 目标线程若已经阻塞在内核中,不能及时看到变量变化
只发送信号 线程醒来后没有持久状态表明应该停止,可能再次阻塞
两者组合 状态表达停止意图,信号负责打断阻塞

25. fork 后恢复:io_poll_svc_notify_fork()

收到 IO_FORK_CHILD 时:

1
2
3
4
5
6
7
8
9
关闭旧 epfd

创建新 epfd

遍历红黑树

_events != 0:重新 EPOLL_CTL_ADD
_events == 0:从树中移除
ADD 失败:从树中移除并返回错误

源码中的判断:

1
2
3
4
5
if (events) {
... EPOLL_CTL_ADD ...
} else {
rbtree_remove(&poll->tree, node);
}

这再次证明 _events == 0 的节点只是暂存关联,不能视为当前有效监听。

注意:当前源码在 fork 恢复时删除节点,但没有同步调整 nwatch。本文只陈述代码行为,不对更高层 fork 调用约束作额外假设。


26. 完整使用周期第 9 步:销毁

26.1 C++ 析构

1
~Poll() { io_poll_destroy(*this); }

26.2 io_poll_destroy()

1
2
3
4
if (poll) {
io_poll_fini(poll);
io_poll_free(poll);
}

26.3 io_poll_fini()

执行顺序:

  1. io_ctx 移除服务;
  2. 关闭 epfd
  3. 销毁 mutex;
  4. 清除可能残留的已阻塞唤醒信号;
  5. 恢复原线程信号掩码;
  6. 恢复原信号处理函数。

调用方应在销毁 Poll 前,确保没有其他线程仍在使用它,也没有回调正在访问其内部对象。


27. 按函数理解整个文件

27.1 生命周期函数

函数 作用
io_poll_alloc() 分配 io_poll_t 内存
io_poll_free() 释放内存
io_poll_init() 初始化信号、mutex、红黑树、epoll 和 context 注册
io_poll_fini() 逆向清理初始化资源
io_poll_create() alloc + init
io_poll_destroy() fini + free
io_poll_open() 关闭旧 epfd 后创建新 epoll 实例
io_poll_close() 关闭 epfd,并先将成员设为 -1

27.2 对外访问函数

函数 作用
io_poll_get_ctx() 返回所属 io_ctx_t*
io_poll_get_poll() 返回通用 ev_poll_t* 接口入口
io_poll_watch() ADD/MOD/DEL 一个 fd,并维护红黑树和 _events/nwatch

27.3 ev_poll 后端函数

函数 作用
io_poll_poll_self() 获取当前 poll 线程令牌
io_poll_poll_wait() epoll_pwait() 等待、转换并分发事件
io_poll_poll_kill() 用状态和线程信号结束目标线程当前一次等待

27.4 内部辅助函数

函数 作用
io_poll_process() 标记本次监听已消费并调用回调
io_fd_cmp() 比较两个 int fd,供红黑树排序和查找
io_poll_watch_from_node() 由嵌入式 rbnode 找回 io_poll_watch
io_poll_from_svc() io_svc 成员找回 io_poll_t
io_poll_from_poll() poll_vptr 成员找回 io_poll_t
io_poll_svc_notify_fork() 子进程中重建 epoll 登记
sig_ign() 空信号处理函数,只用于打断等待

28. 最容易出现的理解和使用错误

28.1 把 _events == 0 理解为“已经彻底注销”

错误。只有红黑树节点被删除并执行 DEL,才是彻底取消登记。

1
2
_events == 0 + 节点仍在树中
= 本次一次性监听已消费,等待调用方恢复或删除

28.2 回调处理完后没有再次调用 io_poll_watch()

结果:只收到第一次事件,后续永远没有通知。

28.3 看到 EPOLLOUT 就假设所有数据都能一次写完

错误。仍需处理部分写和 EAGAIN

28.4 只处理 IO_EVENT_IN,忽略同时出现的 ERR/HUP

revents 是位组合。回调应逐位判断,不应使用互斥的单值比较。

28.5 在取消登记前释放 watch

红黑树节点嵌入在 io_poll_watch 对象中。对象释放后树中节点会变成悬空内存。

正确顺序:

1
2
3
4
5
6
7
8
9
停止产生新操作

io_poll_watch(..., 0, watch)

确认相关回调/任务结束

close(fd)

释放包含 watch 的对象

28.6 把 ev_poll_kill() 理解为线程终止

错误。它只结束当前一次 wait,不销毁线程或 Poll。


29. 推荐断点和观察变量

断点位置 重点观察
io_poll_watch():310 fdnodewatch 是否匹配
io_poll_watch():317 宏转换后的 event.eventsevent.data.fd
io_poll_watch():319 为什么进入 MOD:修改事件还是 _events==0
io_poll_watch():329 首次 ADD 的返回值和 errno
io_poll_watch():338 nwatch 是否需要增加
io_poll_watch():341 DEL 前后的树和计数
io_poll_poll_wait():455 timeoutneventserrno
io_poll_poll_wait():482 每个内核事件位和转换后的 revents
io_poll_poll_wait():498 fd 是否还能在红黑树找到
io_poll_process():602 回调前 _events 从非 0 变为 0
具体 watch->func 是否读取/写入并再次调用 io_poll_watch()
io_poll_poll_kill():540 stoppedpthread_kill() 返回值

推荐日志:

1
2
3
4
5
thread_id, epfd, fd,
requested_io_events, generated_epoll_events,
returned_epoll_events, callback_io_events,
watch_pointer, watch->_events,
nwatch, node_pointer, timeout, errno

30. 用一句话记住每个关键概念

概念 一句话
IO_EVENT_IN Lely 层请求读侧通知
EPOLLIN Linux 表示当前可以尝试读取
IO_EVENT_PRI Lely 层请求异常优先级通知
EPOLLPRI Linux 表示出现 exceptional condition
IO_EVENT_OUT Lely 层请求写侧通知
EPOLLOUT Linux 表示当前可以尝试推进写入
EPOLLONESHOT 每次 ADD/MOD 后最多通知一次,之后必须 MOD 才能继续
EPOLL_CTL_ADD 首次加入 epoll interest list
EPOLL_CTL_MOD 修改事件或恢复下一次一次性监听
EPOLL_CTL_DEL 从 epoll interest list 彻底删除
红黑树 用户态 fd → watch 索引、唯一性检查和陈旧事件过滤
io_fd_cmp() 比较两个整数 fd 的大小
_events != 0 当前有一次有效监听
_events == 0 且节点仍在 本次监听已消费,需恢复或删除
nwatch 当前有效一次监听的数量,不是树节点数

31. 源码位置索引

主题 位置
EPOLL_EVENT_INIT 事件映射 poll.c:44-52
io_poll 数据结构 poll.c:88-101
初始化和 epoll 创建 poll.c:131-205572-592
io_poll_watch() poll.c:286-357
fork 恢复 poll.c:359-402
线程令牌 self() poll.c:404-420
wait() poll.c:422-523
kill() poll.c:525-554
事件回调前状态更新 poll.c:594-614
io_fd_cmp() poll.c:616-625
watch 结构和 API 约定 poll.h:35-123
C++ RAII 包装 poll.hpp:36-103
Event loop 驱动 Poll src/ev/loop.cev_loop_ctx_wait_one()ev_loop_ctx_kill()
SocketCAN fd 直接使用 Poll src/io2/linux/can_chan.cio_can_chan_impl_watch_func()、读写 task
timerfd 直接使用 Poll src/io2/linux/timer.cio_timer_impl_open()io_timer_impl_wait_task_func()
self-pipe 直接使用 Poll src/io2/posix/sigset.cio_sigset_impl_open()io_sigset_impl_read_task_func()
I/O 到 CAN 网络核心的桥接 src/io2/can_net.c:read、write、next-time callback
CANopen Node 间接依赖链 src/coapp/node.cppNode::Node()
CANopen Master 间接依赖链 src/coapp/master.cppBasicMaster::BasicMaster()

32. 参考资料

32.1 本次分析的 Poll 源码

  • poll.c:Lely Linux epoll 后端实现。
  • poll.hio_poll_watch、回调语义及对外 C API。
  • poll.hpp:C++ RAII 包装。

32.2 Lely 全栈模块源码依据

本章的组件调用关系按 Lely liblely-io2 / liblely-coapp 源码分析:

  • src/ev/loop.c:event loop 调用 ev_poll_wait()ev_poll_kill()
  • src/io2/linux/can_chan.c:SocketCAN fd 的 IO_EVENT_IN/OUT 登记、收发任务和发送确认。
  • src/io2/linux/timer.ctimerfd 创建、可读事件、到期计数和 wait queue。
  • src/io2/posix/sigset.c:self-pipe、signal handler、pipe 可读事件和信号分发。
  • src/io2/can_net.c:CanChannel/Timer 与被动 can_net_t 之间的桥接。
  • src/coapp/node.cpp:Node 基于 io::CanNet 启动 CANopen 收发。
  • src/coapp/master.cpp:Master 在 Node 之上实现主站和远程节点管理。
  • src/coapp/loop_driver.cpp:CANopen 远程节点 Driver 的独立业务 event loop;不直接监听 CAN fd。
  • include/lely/io2/posix/poll.hpplinux/can.hppsys/timer.hppsys/sigset.hpp:对应 C++ 包装及构造依赖。

源码层级以 Lely 官方 lely-core 仓库为准;章节中的短代码片段使用同源公开镜像提交 e2ca7cb9e719e3371b88c3441d7f66e54a1a408d 校核路径和调用关系。

32.3 Linux API 依据