249 lines
9.5 KiB
Markdown
249 lines
9.5 KiB
Markdown
# libtask 模块分析
|
||
|
||
**日期**: 2026-06-15(更新)
|
||
**基于源码**: `src/public/libtask/src/myTask.c`(1个文件,913行)
|
||
|
||
---
|
||
|
||
## 1. 模块定位
|
||
|
||
`libtask` 是 RTU 的**线程间同步原语库**,位于公共库层。它提供三个独立的子系统:事件标志、消息队列、定时器。全部基于 POSIX 标准 API(pthread、timer_create),零第三方依赖,纯 C 实现。
|
||
|
||
RTU 的 9 个应用线程(app_sys、app_cmd、app_comm_channel 等)全部通过 libtask 的原语进行同步和通信。
|
||
|
||
## 2. 三个子系统
|
||
|
||
### 2.1 事件标志(Task Event)
|
||
|
||
**用途**: 线程间事件通知,支持多事件同时等待。
|
||
|
||
**核心结构**:
|
||
```c
|
||
typedef struct {
|
||
char name[64];
|
||
uint32_t events; // 当前挂起的事件位掩码
|
||
pthread_mutex_t mutex;
|
||
pthread_cond_t cond;
|
||
int ref_count;
|
||
} stru_task_event;
|
||
```
|
||
|
||
**API**:
|
||
|
||
| 函数 | 说明 |
|
||
|------|------|
|
||
| `task_event_create(name)` | 创建事件对象,ref_count=1 |
|
||
| `task_event_send(event, bits)` | 设置事件位(OR 操作),signal 唤醒等待者 |
|
||
| `task_event_recv(event, set, opt, timeout_ms, &recved)` | 等待事件。opt: AND/OR + CLEAR |
|
||
| `task_event_clear(event, bits)` | 手动清除事件位 |
|
||
| `task_event_query(event)` | **新增** — 不等待,直接返回当前挂起的事件位 |
|
||
| `task_event_destroy(event)` | 引用计数减 1,归零则释放 |
|
||
|
||
**等待逻辑**:
|
||
```
|
||
task_event_recv(set, opt, timeout):
|
||
lock mutex
|
||
while 事件未满足:
|
||
if timeout == 0: → 返回 -1 (非阻塞)
|
||
if timeout == FOREVER: → pthread_cond_wait (永久阻塞)
|
||
else: → pthread_cond_timedwait (超时返回 -1)
|
||
if 满足:
|
||
OR 模式: 检查 (events & mask) == mask
|
||
AND 模式: 检查 events 包含所有 bits
|
||
if CLEAR flag: events &= ~mask (自动清除已捕获的事件)
|
||
unlock mutex
|
||
```
|
||
|
||
**关键设计**: 通过 `TASK_EVENT_FLAG_CLEAR` 实现边沿触发语义 —— 唤醒后自动清除,防止重复处理。
|
||
|
||
---
|
||
|
||
### 2.2 消息队列(Task Message Queue)
|
||
|
||
**用途**: 线程间数据传递,环形缓冲区 + 长度前缀编码。
|
||
|
||
**核心结构**:
|
||
```c
|
||
typedef struct {
|
||
char name[64];
|
||
uint8_t *buffer; // 环形缓冲区(msg_size × max_msgs)
|
||
uint32_t size; // 单条消息的最大大小
|
||
uint32_t max_msgs; // 队列容量
|
||
uint32_t msg_count; // 当前消息数
|
||
uint32_t head; // 读指针
|
||
uint32_t tail; // 写指针
|
||
pthread_mutex_t mutex;
|
||
pthread_cond_t cond;
|
||
int ref_count;
|
||
} stru_task_msg_queue;
|
||
```
|
||
|
||
**消息编码格式**:
|
||
|
||
每条消息在缓冲区中的格式为:
|
||
```
|
||
[4字节 msg_len] [msg_len 字节 payload]
|
||
```
|
||
|
||
写入时先写长度再写数据,读时先读长度再读数据。这允许可变长度消息。
|
||
|
||
**API**:
|
||
|
||
| 函数 | 说明 |
|
||
|------|------|
|
||
| `task_msg_queue_create(name, msg_size, msg_num)` | 创建队列,分配 `msg_size × msg_num` 字节缓冲区 |
|
||
| `task_msg_queue_send(queue, msg, size)` | 非阻塞写入。队列满时返回 -1 |
|
||
| `task_msg_queue_send_timeout(queue, msg, size, timeout_ms)` | **新增** — 写入。队列满时阻塞等待(支持 FOREVER / 定时 / 0=非阻塞) |
|
||
| `task_msg_queue_recv(queue, msg, size, timeout_ms)` | 阻塞读取(支持 FOREVER / 定时 / 0=try) |
|
||
| `task_msg_queue_try_recv(queue, msg, size)` | 非阻塞读取(内部调用 recv(0)) |
|
||
| `task_msg_queue_get_count(queue)` | 返回当前消息数 |
|
||
| `task_msg_queue_space(queue)` | 返回剩余空间 |
|
||
| `task_msg_queue_destroy(queue)` | 引用计数减 1,归零释放 |
|
||
|
||
**关键行为**:
|
||
- `send` 满时**直接返回 -1**,不阻塞,不覆盖旧数据。生产者需自行处理
|
||
- `recv` 队列空时 `pthread_cond_wait` 阻塞等待
|
||
- 头部 4 字节长度前缀确保可变长度消息的正确读取
|
||
|
||
---
|
||
|
||
### 2.3 定时器(Task Timer)
|
||
|
||
**用途**: 一次性或周期性定时任务。
|
||
|
||
**核心结构**:
|
||
```c
|
||
typedef struct {
|
||
char name[64];
|
||
timer_t timerid; // POSIX timer ID
|
||
timer_func_cb fun; // 回调函数
|
||
void *arg; // 回调参数
|
||
uint32_t timeout_ms; // 超时毫秒
|
||
int flags; // TASK_TIMER_FLAG_PERIODIC
|
||
int active; // 是否活跃
|
||
pthread_mutex_t mutex;
|
||
int ref_count;
|
||
} stru_task_timer;
|
||
```
|
||
|
||
**定时器机制** ✅ **已优化 (2026-06-15)**: 使用 Linux **timerfd** + **epoll** + **单例管理器线程** 替代原 `SIGEV_THREAD`
|
||
|
||
```
|
||
task_timer_create → timerfd_create(CLOCK_REALTIME, TFD_NONBLOCK) 创建 fd
|
||
→ task_timer_start → timerfd_settime 装备 → epoll_ctl(ADD) 注册
|
||
→ epoll_wait 就绪 → read(fd) 清除过期计数 → fun(arg) 调用回调
|
||
→ 周期模式: timerfd 自动重触发
|
||
→ 单次模式: active=0,回调不再执行
|
||
```
|
||
|
||
**API**:
|
||
|
||
| 函数 | 说明 |
|
||
|------|------|
|
||
| `task_timer_create(name, cb, arg, timeout_ms, flags)` | 创建定时器(未启动) |
|
||
| `task_timer_start(timer)` | 启动定时器 |
|
||
| `task_timer_stop(timer)` | 停止定时器(设置 its={0,0}) |
|
||
| `task_timer_restart(timer, new_timeout)` | 修改超时并重启 |
|
||
| `task_timer_is_active(timer)` | 查询是否活跃 |
|
||
| `task_timer_destroy(timer)` | 引用计数减 1,归零删除 timer + 释放 |
|
||
|
||
---
|
||
|
||
## 3. 引用计数机制
|
||
|
||
三个子系统都使用引用计数(`ref_count`)进行生命周期管理:
|
||
|
||
```c
|
||
// 创建
|
||
p->ref_count = 1;
|
||
|
||
// 获取 → ref_count++
|
||
// 不需要显式 get API,通过外部指针共享隐式增加
|
||
|
||
// 销毁
|
||
lock mutex
|
||
p->ref_count--;
|
||
if (p->ref_count > 0) { unlock; return 0; } // 仍有使用者,不释放
|
||
unlock mutex
|
||
free(p); // 最后一个使用者,释放
|
||
```
|
||
|
||
**特点**: 多线程共享同一个事件/队列/定时器对象时,最后一个使用者释放。但没有显式的 `retain`/`release` API,ref_count 的实际增减依赖外部代码手动管理。
|
||
|
||
---
|
||
|
||
## 4. 辅助函数
|
||
|
||
```c
|
||
void task_sleep_ms(uint32_t ms):
|
||
// 使用 nanosleep 实现毫秒级休眠,支持 EINTR 中断后自动恢复
|
||
struct timespec ts = {ms/1000, (ms%1000)*1000000};
|
||
while(nanosleep(&ts, &ts) == -1 && errno == EINTR);
|
||
```
|
||
|
||
---
|
||
|
||
## 5. 在 RTU 中的应用
|
||
|
||
RTU 使用 libtask 实现线程调度框架 `myTask.h`(位于 `release/inc/`),定义:
|
||
|
||
```c
|
||
// 在 myTask.h 中定义(不在 libtask 目录内)
|
||
#define EV_TIMER1 (1 << 0) // 10ms
|
||
#define EV_TIMER2 (1 << 1) // 100ms
|
||
#define EV_TIMER3 (1 << 2) // 1000ms
|
||
#define EV_TIMER4 (1 << 3)
|
||
|
||
#define TASK_EVENT_WAIT_FOREVER ~0
|
||
#define TASK_EVENT_FLAG_OR 0
|
||
#define TASK_EVENT_FLAG_AND 1
|
||
#define TASK_EVENT_FLAG_CLEAR 2
|
||
|
||
#define TASK_TIMER_FLAG_PERIODIC 1
|
||
#define TASK_TIMER_FLAG_ONE_SHOT 0
|
||
|
||
typedef void *(*task_thread_func)(void *arg);
|
||
```
|
||
|
||
每个应用线程的主循环模式:
|
||
```cpp
|
||
while (1) {
|
||
task_event_recv(p_event, EV_TIMER1 | EV_TIMER2 | EV_TIMER3,
|
||
TASK_EVENT_FLAG_OR | TASK_EVENT_FLAG_CLEAR,
|
||
TASK_EVENT_WAIT_FOREVER, &event);
|
||
if (event & EV_TIMER1) { /* 10ms 任务 */ }
|
||
if (event & EV_TIMER2) { /* 100ms 任务 */ }
|
||
if (event & EV_TIMER3) { /* 1000ms 任务 */ }
|
||
}
|
||
```
|
||
|
||
---
|
||
|
||
## 6. 优点
|
||
|
||
| 优点 | 说明 |
|
||
|------|------|
|
||
| **纯 POSIX + Linux 原生** | 事件/消息队列用 pthread,定时器用 timerfd+epoll,无第三方库依赖,交叉编译无障碍 |
|
||
| **单线程管理所有定时器** | ✅ **已优化 (2026-06-15)** — 一个持久化 epoll 线程管理全部 27 个定时器,消除原 SIGEV_THREAD 每秒 ~1000 次线程创建/销毁 |
|
||
| **stop/destroy 零 CPU 同步** | ✅ **已优化 (2026-06-15)** — condvar 替代自旋,`while(running) cond_wait` 阻塞等待无 CPU 开销 |
|
||
| **完备的超时支持** | 事件和消息队列都支持三种模式:永久等待 / 指定超时 / 非阻塞立即返回 |
|
||
| **双条件事件** | 支持 AND(全部触发)和 OR(任一触发)两种等待模式,灵活匹配不同场景 |
|
||
| **边沿触发** | `TASK_EVENT_FLAG_CLEAR` 自动在等待返回时清除事件位,防止重复处理 |
|
||
| **可变长度消息** | 消息队列通过 4 字节长度前缀编码,支持不同大小的消息共用同一队列 |
|
||
| **完备的队列查询** | `get_count` 和 `space` 提供队列状态查询,生产者可据此做背压决策 |
|
||
| **引用计数生命周期** | 多线程共享对象时安全释放,不会出现 use-after-free |
|
||
| **代码紧凑** | 全部功能在 913 行 C 代码中,易于审计和理解 |
|
||
|
||
## 7. 缺点
|
||
|
||
| 缺点 | 严重度 | 说明 |
|
||
|------|--------|------|
|
||
| **定时器 destroy use-after-free** | 高 | ✅ **已修复 (2026-06-15)**: `task_timer_destroy` 中原 `free(p)` 后在 LOG 中访问 `p->name`,`free` 移到 LOG 之后 |
|
||
| **消息队列 destroy use-after-free** | 高 | ✅ **已修复 (2026-06-15)**: `task_msg_queue_destroy` 同上 |
|
||
| **事件 destroy use-after-free** | 高 | ✅ **已修复 (2026-06-15)**: `task_event_destroy` 同上 |
|
||
| **事件 recv AND 位测试陷阱** | 中 | ✅ **已修复 (2026-06-15)**: `opt == TASK_EVENT_FLAG_AND` 改为 `opt & TASK_EVENT_FLAG_AND`,避免 AND\|CLEAR 组合被误识别 |
|
||
| **事件 send 用 signal 非 broadcast** | 低 | ✅ **已修复 (2026-06-15)**: `pthread_cond_signal` 改为 `pthread_cond_broadcast`,多等待者场景更稳健 |
|
||
| **无事件优先级** | 低 | `task_event_recv` 按位掩码匹配,所有事件平等。多事件同时触发无优先级排序 |
|
||
| **ref_count 缺乏原子性** | 低 | 引用计数通过 mutex 保护无 CAS,极高频场景下瓶颈(本场景不适用) |
|
||
| **头文件不在模块内** | 低 | `myTask.h` 位于 `release/inc/`,新开发者可能找不到 |
|