Lely CANopen spscring 单生产者单消费者环形队列完整机制学习笔记

在这里插入图片描述

@[toc]


1. 先明确边界:spscring 不是数据缓冲区

SPSC 是 Single-Producer Single-Consumer,即单生产者、单消费者。

Lely 的 spscring 与常见“结构体内部直接包含数组”的环形队列不同。它只管理:

  • 哪些索引可以由生产者写入;
  • 哪些索引可以由消费者读取;
  • 生产者何时把写入结果发布给消费者;
  • 消费者何时把已读槽位归还给生产者;
  • 队列由空变为可读或由满变为可写时,是否需要触发回调。

实际数据存储由调用者单独提供:

1
2
struct spscring ring;
struct io_can_frame *rxbuf;

两者的关系是:

1
2
3
4
5
spscring
= 索引所有权与可见性控制器

rxbuf
= 真正保存 CAN 帧的内存数组

因此,调用 spscring_p_alloc() 只会得到一个可写索引,不会返回数据指针,也不会复制任何数据。

1
2
3
4
5
6
7
flowchart LR
Prod[Producer] -->|p_alloc obtains index| Ring[spscring index controller]
Prod -->|writes payload| Buf[User data buffer]
Prod -->|p_commit publishes index| Ring
Ring -->|c_alloc obtains index| Cons[Consumer]
Buf -->|reads payload| Cons
Cons -->|c_commit releases index| Ring

这是一种低层、通用的设计。相同的 spscring 可以控制:

  • CAN 帧数组;
  • 音频采样数组;
  • 日志记录数组;
  • DMA 描述符数组;
  • 共享内存块;
  • 文件中的逻辑槽位。

2. 它在 Lely Linux CAN 接收链路中的位置

在 Linux CanChannel 中,spscring 只用于控制用户态接收队列:

1
2
3
io_can_chan_impl
├─ struct spscring rxring
└─ struct io_can_frame *rxbuf

初始化时:

1
2
spscring_init(&impl->rxring, rxlen);
impl->rxbuf = calloc(rxlen, sizeof(struct io_can_frame));

实际职责划分如下:

组件 角色 主要操作
rxbuf_task 逻辑生产者 从 SocketCAN recvmsg(),写 rxbuf[i],执行 p_commit()
read_task 逻辑消费者 rxbuf[i] 取帧,完成异步读请求,执行 c_commit()
同步 read() 逻辑消费者的另一个入口 同样从 rxring 取帧,但由 c_mtx 与异步消费者串行化
c_mtx 消费者侧串行化 把多个消费入口约束成“一个逻辑消费者”
io_can_chan_impl_c_signal() 可读通知桥梁 队列产生数据后,重新投递 read_task

因此它与上一篇 IO2 笔记的连接关系是:

1
2
3
4
5
6
7
8
9
10
11
flowchart TB
Sock[SocketCAN fd] --> RxTask[rxbuf_task]
RxTask -->|p_alloc| Ring[rxring]
RxTask -->|recvmsg writes| Buf[rxbuf array]
RxTask -->|p_commit| Ring
Ring -->|consumer signal| Sig[io_can_chan_impl_c_signal]
Sig -->|post| ReadTask[read_task]
ReadTask -->|c_alloc| Ring
Buf -->|copy frame| ReadTask
ReadTask -->|c_commit| Ring
ReadTask --> Done[User read completion task]

SocketCAN 负责把帧交给进程,spscring 负责在 rxbuf_task 与读任务之间安全转移槽位所有权。


3. 公共 API 按生产者和消费者严格分组

3.1 通用接口

API 作用
spscring_init() 初始化环形队列的索引状态
spscring_size() 返回逻辑容量

3.2 生产者接口

API 作用
spscring_p_capacity() 查询总可写槽位数,允许跨越数组末尾
spscring_p_capacity_no_wrap() 查询当前连续可写槽位数,不跨越数组末尾
spscring_p_alloc() 申请可写索引范围,允许逻辑环绕
spscring_p_alloc_no_wrap() 申请连续可写索引范围,不允许环绕
spscring_p_commit() 发布已写槽位,使消费者可见
spscring_p_submit_wait() 可写空间不足时注册通知
spscring_p_abort_wait() 取消生产者等待

3.3 消费者接口

API 作用
spscring_c_capacity() 查询总可读槽位数,允许跨越数组末尾
spscring_c_capacity_no_wrap() 查询当前连续可读槽位数,不跨越数组末尾
spscring_c_alloc() 申请可读索引范围,允许逻辑环绕
spscring_c_alloc_no_wrap() 申请连续可读索引范围,不允许环绕
spscring_c_commit() 归还已读槽位,使生产者可重用
spscring_c_submit_wait() 可读数据不足时注册通知
spscring_c_abort_wait() 取消消费者等待

最重要的使用约束是:

1
2
所有 p_* 接口只能由一个逻辑生产者调用;
所有 c_* 接口只能由一个逻辑消费者调用。

这里的“逻辑”允许外部使用互斥锁把多个入口串行化,但一旦没有串行化,算法就不再是安全的 SPSC。


4. struct spscring 的内存布局

核心结构可以概括为:

1
2
3
4
5
6
7
8
9
spscring
├─ p: producer side
│ ├─ ctx: producer local context
│ ├─ pos: producer published position
│ └─ sig: producer wait registration
└─ c: consumer side
├─ ctx: consumer local context
├─ pos: consumer published position
└─ sig: consumer wait registration

4.1 本地上下文 spscring_ctx

1
2
3
4
5
6
struct spscring_ctx {
size_t size;
size_t base;
size_t pos;
size_t end;
};
字段 含义
size 环形队列容量 N
base 当前本地轮次对应的绝对基址
pos 当前物理索引,范围为 [0, N-1]
end 本地缓存的逻辑可用区间终点,可能大于 N

生产者只修改 p.ctx,消费者只修改 c.ctx。这些字段不是原子的,因为不会被对端直接访问。

4.2 发布位置 p.posc.pos

1
2
p.pos = 已经提交、可供消费者读取的生产者绝对位置
c.pos = 已经提交、可供生产者重用的消费者绝对位置

访问关系:

原子位置 谁写 谁读
p.pos 生产者 消费者
c.pos 消费者 生产者

这两个位置是生产者和消费者之间最核心的共享状态。

4.3 等待记录 spscring_sig

1
2
3
4
5
struct spscring_sig {
spscring_atomic_t size;
void (*func)(struct spscring *ring, void *arg);
void *arg;
};
等待记录 谁注册/取消 谁检查并触发
p.sig 生产者 消费者在 c_commit() 后检查
c.sig 消费者 生产者在 p_commit() 后检查

p.sig 的意思不是“生产者通知消费者”,而是“生产者正在等待可写空间”。

c.sig 的意思不是“消费者通知生产者”,而是“消费者正在等待可读数据”。


5. 为什么大量使用缓存行对齐与填充

结构体中的 ctxpossig 都使用:

1
_Alignas(LEVEL1_DCACHE_LINESIZE)

并通过 _pad 填充到一个一级数据缓存行。

其目的不是为了数据对齐美观,而是避免 false sharing,即伪共享。

5.1 没有隔离时的问题

假设生产者频繁修改 p.ctx.pos,消费者频繁读取 p.pos。如果它们位于同一个缓存行,即使双方访问的是不同变量,该缓存行仍会在 CPU 核之间不断失效和迁移。

5.2 当前布局的热点划分

1
2
3
4
生产者高频私有写:p.ctx
消费者高频私有写:c.ctx
跨线程发布位置:p.pos / c.pos
等待协议状态:p.sig / c.sig

每类热点单独占据缓存行后:

  • 本地 ctx 更新通常不会触发对端缓存失效;
  • 只有真正需要通信的 possig.size 才产生共享缓存流量;
  • 在 64 字节缓存行平台上,整个对象大约会占用六个缓存行,这是用空间换并发稳定性。

因此 spscring 不适合作为海量小对象嵌入到每个消息中,它更适合作为长期存在的队列控制对象。


6. 两套坐标:物理索引与绝对位置

普通环形队列通常只保存:

1
2
read_index
write_index

当两者相等时,必须额外区分“空”还是“满”,常见做法是浪费一个槽位。

spscring 不浪费槽位。它同时维护:

  • 物理索引:pos,用于访问数组;
  • 绝对位置:base + pos,用于区分不同轮次。

假设容量 N = 8

1
2
3
4
绝对位置 0  → 物理索引 0
绝对位置 7 → 物理索引 7
绝对位置 8 → 物理索引 0
绝对位置 13 → 物理索引 5

base 保存当前物理轮次的起点:

1
physical_index = absolute_position - base

pos 前进到 N 时:

1
2
3
base += N;
pos -= N;
end -= N;

于是物理索引回到数组开头,而绝对位置仍继续前进。

6.1 可用数据与可写空间公式

定义:

1
2
3
P = producer absolute position
C = consumer absolute position
N = ring size

则:

1
2
可读数量 = P - C
可写数量 = C + N - P

并始终满足:

1
可读数量 + 可写数量 = N

这就是它可以使用全部 N 个槽位而不需要浪费一个槽位的原因。


7. ctx.end 是本地缓存,不是第三个共享游标

ctx.end 容易被误解成全局结束位置。实际上它只是本地线程对“当前可用范围终点”的缓存。

生产者查询容量时:

1
2
3
4
size_t cpos = atomic_load(&ring->c.pos) + ctx->size;
cpos -= ctx->base;
ctx->end = cpos;
return ctx->end - ctx->pos;

消费者查询容量时:

1
2
3
4
size_t ppos = atomic_load(&ring->p.pos);
ppos -= ctx->base;
ctx->end = ppos;
return ctx->end - ctx->pos;

7.1 为什么允许 end > size

假设容量为 8,生产者物理位置在 6,而消费者已前进到下一轮的物理位置 2。

生产者可写槽位是:

1
6, 7, 0, 1

为了把这个跨界范围表示成一个线性区间,本地上下文可表示为:

1
2
3
pos = 6
end = 10
capacity = end - pos = 4

其中逻辑索引 8、9 分别映射到物理索引 0、1。

7.2 缓存的性能意义

p_alloc()c_alloc() 会先使用:

1
ctx.end - ctx.pos

只有请求数量超过本地缓存容量时,才重新原子读取对端发布位置。

因此,在批量操作期间,多次 alloc() 不一定每次都产生跨核原子读取。


8. alloc()commit() 是一个两阶段事务

8.1 alloc() 只授予本地访问权

生产者:

1
2
3
4
p_alloc
→ 返回可写起始索引
→ 调整请求数量到实际可用数量
→ 不修改共享 p.pos

消费者:

1
2
3
4
c_alloc
→ 返回可读起始索引
→ 调整请求数量到实际可用数量
→ 不修改共享 c.pos

因此在没有 commit() 的情况下,重复调用相同的 alloc() 是幂等的,返回的仍是同一段尚未提交范围。

8.2 commit() 才转移所有权

生产者完成数据写入后调用:

1
spscring_p_commit(ring, count);

这一步会:

  1. 推进生产者本地位置;
  2. 必要时处理物理环绕;
  3. release-store 新的 p.pos
  4. 检查消费者等待条件;
  5. 条件满足时同步调用消费者信号函数。

消费者完成数据读取后调用:

1
spscring_c_commit(ring, count);

它会把槽位归还给生产者,并检查生产者是否正在等待可写空间。

8.3 正确访问区间

1
2
生产者:p_alloc → 写用户缓冲区 → p_commit
消费者:c_alloc → 读用户缓冲区 → c_commit

禁止把数据访问放到事务边界之外:

1
2
错误:先 p_commit,再写缓冲区
错误:先 c_commit,再读取缓冲区

提交后,槽位所有权已经转移,对端可能立即访问或覆盖该槽位。


9. 生产者路径逐步解析

9.1 查询总可写容量

1
size_t capacity = spscring_p_capacity(ring);

其逻辑等价于:

1
free = consumer_position + ring_size - producer_position

初始化时:

1
2
3
P = 0
C = 0
free = 0 + N - 0 = N

所以空队列的生产者容量是全部 N 个槽位。

9.2 分配可写范围

1
2
size_t n = requested;
size_t i = spscring_p_alloc(ring, &n);

行为:

  • 返回当前生产者物理索引;
  • 若请求超过容量,把 n 截断到实际容量;
  • 不推进生产者位置;
  • 允许返回的逻辑范围跨越数组末尾。

使用允许环绕的 API 时,调用者必须自行取模:

1
2
for (size_t j = 0; j < n; j++)
buffer[(i + j) % ring_size] = input[j];

9.3 发布已写数据

1
spscring_p_commit(ring, n);

位置推进代码的核心含义是:

1
2
3
4
5
6
new local pos = old pos + n

if new local pos reaches ring_size:
advance base by one ring
wrap local pos to beginning
normalize cached end

然后发布:

1
p.pos = p.ctx.base + p.ctx.pos

返回值是提交后的下一个生产者物理索引。


10. 消费者路径与生产者完全对称

消费者容量:

1
readable = producer_position - consumer_position

初始化时为 0。

当生产者提交 3 个槽位后:

1
2
3
P = 3
C = 0
readable = 3

消费者典型流程:

1
2
3
4
5
6
7
size_t n = requested;
size_t i = spscring_c_alloc(ring, &n);

for (size_t j = 0; j < n; j++)
output[j] = buffer[(i + j) % ring_size];

spscring_c_commit(ring, n);

c_commit() 发布新的 c.pos 后,生产者才能重新覆盖这些槽位。


11. no_wrap 接口解决连续内存需求

11.1 总容量与连续容量不同

假设:

1
2
3
ring size = 8
producer pos = 6
总可写容量 = 4

总可写槽位为:

1
6, 7, 0, 1

但从索引 6 到数组末尾只有两个连续槽位。

因此:

1
2
p_capacity()         = 4
p_capacity_no_wrap() = 2

no_wrap 容量公式为:

1
min(end, size) - pos

11.2 什么时候应该使用 no_wrap

适合:

  • 单次 memcpy()
  • DMA 连续描述符;
  • 要求返回单一连续内存区间的 API;
  • 不希望调用方做取模循环的批量数据处理。

不适合:

  • 希望一次申请全部剩余总容量;
  • 可以自然处理两段数据的场景。

跨界批量访问通常要拆成两次:

1
2
3
4
第一次:处理数组尾部连续段
commit
第二次:处理数组头部连续段
commit

12. 数值运行示例:容量为 8

12.1 初始化

1
2
3
P = 0, C = 0
readable = 0
writable = 8

12.2 生产者写入 3 个元素

1
2
3
p_alloc(3)  → index 0, count 3
write → buffer[0..2]
p_commit(3) → P = 3

状态:

1
2
readable = 3
writable = 5

12.3 消费者读取 2 个元素

1
2
3
c_alloc(2)  → index 0, count 2
read → buffer[0..1]
c_commit(2) → C = 2

状态:

1
2
readable = 1
writable = 7

12.4 生产者再写 5 个元素

当前生产者物理索引为 3,数组尾部正好有 5 个槽位:

1
2
3
p_alloc_no_wrap(5) → index 3, count 5
write → buffer[3..7]
p_commit(5) → P = 8, physical pos wraps to 0

状态:

1
2
readable = P - C = 8 - 2 = 6
writable = 2

消费者可读物理顺序为:

1
2, 3, 4, 5, 6, 7

12.5 消费者全部读完

1
2
c_alloc_no_wrap(6) → index 2, count 6
c_commit(6) → C = 8, physical pos wraps to 0

最终:

1
2
3
P = 8, C = 8
readable = 0
writable = 8

虽然物理索引再次相等,但绝对位置表明这是“第二轮的空队列”,不会与初始状态或满队列混淆。


13. Acquire/Release 如何保证数据可见性

在正常原子构建中:

1
2
atomic_load_explicit(..., memory_order_acquire);
atomic_store_explicit(..., memory_order_release);

13.1 生产者到消费者

1
2
3
4
5
6
7
Producer                              Consumer
-------- --------
write buffer[i]
write buffer[i + 1]
release-store p.pos -------------> acquire-load p.pos
read buffer[i]
read buffer[i + 1]

当消费者通过 acquire-load 观察到新的 p.pos 时,生产者在 release-store 之前完成的缓冲区写入也对消费者可见。

13.2 消费者到生产者

1
2
3
4
5
Consumer                              Producer
-------- --------
read buffer[i]
release-store c.pos -------------> acquire-load c.pos
overwrite buffer[i]

生产者只有观察到消费者提交的新 c.pos 后,才会认为槽位可重用。

13.3 环形队列不会替调用者同步其他无关数据

保证只覆盖与槽位发布顺序相关的数据访问。若回调参数或其他共享对象还会被并发修改,仍需要它们自己的同步规则。


14. 等待机制不是阻塞队列,而是条件通知注册

spscring 不直接创建条件变量、信号量或线程等待对象。它提供的是:

1
2
如果当前条件不满足,登记一个回调;
当对端 commit 使条件满足时,调用该回调。

生产者可等待“至少有 size 个可写槽位”:

1
spscring_p_submit_wait(ring, size, producer_signal, arg);

消费者可等待“至少有 size 个可读槽位”:

1
spscring_c_submit_wait(ring, size, consumer_signal, arg);

返回值:

返回值 含义
0 条件已经满足,没有注册等待,也不会调用回调
1 等待已经注册;稍后条件满足时由对端 commit() 调用回调,或等待可能被取消

回调不是由独立线程自动执行,而是在使条件满足的对端 commit() 调用栈中同步执行。

这一点必须牢记:

1
2
3
4
5
consumer wait callback
由 producer 的 p_commit() 线程执行

producer wait callback
由 consumer 的 c_commit() 线程执行

15. 双重检查如何避免丢失唤醒

以消费者等待一个可读元素为例。

错误的简单实现可能是:

1
2
3
4
1. 检查队列为空
2. 生产者提交数据并检查“无人等待”
3. 消费者登记等待
4. 数据已经存在,但之后再也没有提交,消费者永久沉睡

Lely 使用“检查—发布等待—再次检查”的流程:

1
2
3
4
5
6
7
8
9
1. abort previous wait
2. 保存 func/arg
3. 第一次检查容量
4. 若已满足,直接返回 0
5. 原子发布 sig.size = requested
6. 第二次检查容量
7. 若并发期间条件已满足,尝试把 sig.size 从 requested 改回 0
8. CAS 成功:等待尚未被对端接管,返回 0
9. CAS 失败:对端已开始或完成信号流程,返回 1
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
sequenceDiagram
participant Cons as Consumer
participant Ring as spscring
participant Prod as Producer

Cons->>Ring: first capacity check
Ring-->>Cons: not enough data
Cons->>Ring: publish c.sig.size
Prod->>Ring: write payload and p_commit
Prod->>Ring: check c.sig and available data
Cons->>Ring: second capacity check

alt Consumer removes wait first
Cons->>Ring: CAS requested to 0
Ring-->>Cons: condition already satisfied, return 0
else Producer claims signal first
Prod->>Ring: CAS requested to SIZE_MAX
Prod->>Ring: copy callback and clear size
Prod-->>Cons: invoke callback, submit_wait returns 1
end

第二次检查与 CAS 竞争共同保证:

1
条件不会在“检查”和“登记”等待之间悄悄变为满足而丢失通知。

16. sig.size 实际是一个三态同步协议

sig.size 不只是“等待数量”,还承担状态机作用:

状态
0 没有等待,或等待已经取消/完成
1..N 已注册等待,数值是条件阈值
SIZE_MAX 对端已经接管通知,正在读取 funcarg

16.1 为什么需要 SIZE_MAX

funcarg 本身不是原子变量。

如果取消等待的一方在对端读取函数指针期间立即覆盖它们,会产生数据竞争。信号方先执行:

1
CAS sig.size: requested → SIZE_MAX

这相当于声明:

1
我已经进入读取 callback 元数据的临界窗口。

随后:

  1. 复制 funcarg 到局部变量;
  2. release-store sig.size = 0
  3. 调用局部副本中的函数。

等待方看到 SIZE_MAX 时不能直接取消,只能短暂自旋,直到信号方完成元数据读取。

16.2 为什么请求数量不会与 SIZE_MAX 冲突

等待数量会被限制到环形队列容量:

1
2
if (size > ctx->size)
size = ctx->size;

同时队列容量被限制为:

1
size <= SIZE_MAX / 2 + 1

因此正常请求阈值不会取到 SIZE_MAX,该值可以安全作为内部哨兵。


17. abort_wait() 的真实语义

取消等待时:

1
2
3
4
5
6
7
8
load sig.size

while size != 0:
if size == SIZE_MAX:
processor pause
reload
else:
CAS size → 0

返回值:

返回值 含义
1 成功取消了一个尚未触发的等待
0 当前没有已注册等待

17.1 取消与回调可能存在竞争

如果信号方已经把状态改成 SIZE_MAX,取消方不能阻止本次回调,因为信号方已经接管了函数指针和参数。

此时 abort_wait() 会等待信号方把状态恢复为 0,随后返回 0。

因此调用者不能假设:

1
只要调用 abort_wait(),回调就绝对不会发生。

更准确的语义是:

1
2
若等待仍处于可取消状态,则取消;
若通知已经开始,则等待其完成接管。

18. pc_signal()cp_signal() 的命名和判断公式

18.1 spscring_pc_signal()

在生产者 p_commit() 后调用:

1
producer commit → consumer signal

它检查 c.sig,判断消费者请求的可读数量是否已经满足。

源码使用:

1
readable = ring_size - producer_free_capacity

若:

1
requested_readable <= readable

则触发消费者回调。

18.2 spscring_cp_signal()

在消费者 c_commit() 后调用:

1
consumer commit → producer signal

它检查 p.sig,判断生产者请求的可写数量是否已经满足。

源码使用:

1
writable = ring_size - consumer_readable_capacity

若:

1
requested_writable <= writable

则触发生产者回调。


19. 回调函数必须保持轻量且避免重入死锁

信号函数在 commit() 内同步执行。因此回调会延长提交操作的执行时间。

推荐回调只做:

  • 设置原子标志;
  • 释放信号量;
  • 向 Executor 投递任务;
  • 写 eventfd;
  • 触发轻量通知。

不推荐:

  • 长时间阻塞;
  • 执行复杂协议状态机;
  • 进行大量内存分配;
  • 获取可能已由提交方持有的同一把锁;
  • 在不满足线程所有权的情况下直接调用对端专属 p_*c_* API。

Lely CanChannel 采用的是安全模式:

1
2
3
p_commit()
→ io_can_chan_impl_c_signal()
→ 只检查状态并 post read_task

而不是在生产者线程中直接消费 CAN 帧。

另一个关键细节是,rxbuf_task 在执行 p_commit() 时没有持有 impl->mtx,而回调 io_can_chan_impl_c_signal() 会获取该互斥锁。这样避免了提交线程在同步回调中重复获取同一把锁而死锁。


20. Lock-free 与 Wait-free 在这里分别意味着什么

20.1 Wait-free

当没有注册信号函数时,普通容量、分配和提交路径只包含有限数量的本地操作与原子读写,不需要等待另一个线程完成某个临界区。

这不表示:

1
2
生产者一定能分配到空间;
消费者一定能读到数据。

队列满或空时,操作会立即返回容量 0。这里的 wait-free 描述的是算法步骤有界,而不是业务条件永远满足。

20.2 Lock-free

等待注册、取消和通知使用 CAS 循环。单个线程可能短暂重试,但系统整体会持续取得进展。

当回调存在时,commit() 的进度还取决于回调:

1
2
回调是 lock-free → 提交路径仍可视为 lock-free
回调会阻塞 → 整体不再具备该性质

20.3 它不是多生产者多消费者队列

多个生产者同时修改 p.ctx,或多个消费者同时修改 c.ctx,不会被原子位置保护。

若业务确实存在多生产者或多消费者:

  • 使用外部互斥锁把同侧访问串行化;或
  • 选择专门的 MPSC、SPMC、MPMC 队列;
  • 不要仅因为 p.pos/c.pos 是原子的,就误认为整个结构支持 MPMC。

21. spscring_yield() 不是操作系统线程让步

等待协议中的自旋提示实现为:

1
2
3
Windows   → YieldProcessor()
x86/x64 → pause 指令
其他平台 → 可能为空操作

它主要用于降低自旋时的流水线和功耗代价,并不等价于:

1
sched_yield();

因此,在信号方长期阻塞或被高优先级线程永久抢占时,取消方可能持续占用 CPU 自旋。

正常设计要求 SIZE_MAX 状态窗口极短,只覆盖复制函数指针与参数的几条指令。


22. LELY_NO_ATOMICS 分支应如何理解

头文件根据构建环境选择:

1
2
3
C++                → std::atomic_size_t
C11 → atomic_size_t
LELY_NO_ATOMICS → size_t

正常多线程环境依赖 C/C++ 原子操作的 acquire/release 与 CAS 语义。

LELY_NO_ATOMICS 分支的普通位置加载和存储退化为直接内存访问,不提供标准 C11 跨线程内存序保证。Windows 的部分 CAS 路径使用 Interlocked,但这不能把所有普通访问自动变成完整、可移植的原子协议。

工程上应采用以下判断:

1
2
3
4
5
单线程或明确受平台约束的特殊构建
→ 可以评估 LELY_NO_ATOMICS

正常 Linux/RTOS 多线程生产者-消费者
→ 应启用编译器原子支持

不要仅因为程序“在当前 CPU 上看起来能运行”就忽略数据竞争与内存可见性问题。


23. 容量上限与 size_t 回绕

初始化断言:

1
assert(size <= SIZE_MAX / 2 + 1);

其作用包括:

  1. 保证生产者与消费者之间的最大逻辑距离不超过半个无符号计数空间;
  2. 让无符号位置回绕后,差值仍能在环容量约束下被唯一解释;
  3. 避免等待数量与内部 SIZE_MAX 哨兵冲突。

绝对位置不是永远不溢出的数学整数,而是按 size_t 模数回绕的逻辑计数器。算法依赖以下不变量:

1
2
0 <= P - C <= N
N <= SIZE_MAX / 2 + 1

在此约束下,即使 PC 穿过 SIZE_MAX,无符号减法仍能得到正确的环内距离。

23.1 容量不能为 0

当前 spscring_init() 没有显式断言 size > 0,但后续所有操作都要求:

1
assert(ctx->pos < ctx->size);

因此 size = 0 在实际使用中无效。工程接口应在上层明确保证:

1
ring size >= 1

24. Lely CanChannel 接收流程的完整时序

24.1 空队列上的异步读

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
sequenceDiagram
participant App as Application
participant Rd as read_task
participant Ring as rxring
participant Rx as rxbuf_task
participant Sock as SocketCAN
participant Sig as c_signal

App->>Rd: submit async read
Rd->>Ring: c_alloc 1
Ring-->>Rd: count 0
Rd->>Ring: c_submit_wait 1
Rd->>Rx: post rxbuf_task
Rx->>Ring: p_alloc 1
Ring-->>Rx: writable index
Rx->>Sock: recvmsg
Sock-->>Rx: CAN frame
Rx->>Rx: write rxbuf[index]
Rx->>Ring: p_commit 1
Ring->>Sig: invoke consumer signal
Sig->>Rd: post read_task
Rd->>Ring: c_alloc 1
Ring-->>Rd: readable index
Rd->>Rd: copy and decode rxbuf[index]
Rd->>Ring: c_commit 1
Rd-->>App: post completion task

24.2 为什么 c_mtx 仍然存在

spscring 只允许一个消费者,但 CanChannel 具有两个可能的消费入口:

1
2
同步 io_can_chan_impl_read()
异步 io_can_chan_impl_do_read()

两条路径都通过 c_mtx 保护 c_alloc/c_commit,把它们串行化为一个逻辑消费者。

因此:

1
2
spscring 的无锁性质
不等于整个 CanChannel 接收路径完全没有互斥锁

环形队列解决的是生产者与消费者之间的槽位转移;上层仍可能需要互斥锁维护更复杂的 API 并发语义。

24.3 MSG_CONFIRM 不进入普通接收环

rxbuf_task 收到 SocketCAN 自发帧确认时:

1
2
3
4
MSG_CONFIRM
→ 转为 can_msg
→ 匹配 confirm_queue
→ 不提交到 rxring

普通 CAN/CAN FD 帧才会执行 p_commit() 并成为用户读请求的数据来源。

24.4 环满与发送确认并存时的源码注意点

rxbuf_task 在存在写确认等待时,即使普通接收环没有空槽,也可能继续调用 recvmsg() 以寻找 MSG_CONFIRM

此时:

1
2
p_alloc() 可能返回 count = 0
frame 指向栈上临时对象

若读到的是普通帧而不是确认帧,该帧无法进入已经满的 rxring。调试高总线负载、启用 txwait 且消费端处理过慢的问题时,应把这一版本相关行为纳入丢帧分析。

根本治理方向仍然是:

  • 提高 rxlen
  • 避免读任务长期阻塞;
  • 缩短完成回调;
  • 监控应用层消费延迟;
  • 根据发送确认需求评估是否必须启用 txwait

25. 一个最小、正确的用户数据队列示例

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
34
35
36
37
38
39
40
41
42
43
44
45
46
47
#include <lely/util/spscring.h>
#include <stdbool.h>
#include <stddef.h>

#define SAMPLE_QUEUE_SIZE 8

struct sample {
unsigned int id;
unsigned int value;
};

struct sample_queue {
struct spscring ring;
struct sample data[SAMPLE_QUEUE_SIZE];
};

static void
sample_queue_init(struct sample_queue *queue)
{
spscring_init(&queue->ring, SAMPLE_QUEUE_SIZE);
}

static bool
sample_queue_push(struct sample_queue *queue, const struct sample *sample)
{
size_t n = 1;
size_t index = spscring_p_alloc_no_wrap(&queue->ring, &n);
if (!n)
return false;

queue->data[index] = *sample;
spscring_p_commit(&queue->ring, 1);
return true;
}

static bool
sample_queue_pop(struct sample_queue *queue, struct sample *sample)
{
size_t n = 1;
size_t index = spscring_c_alloc_no_wrap(&queue->ring, &n);
if (!n)
return false;

*sample = queue->data[index];
spscring_c_commit(&queue->ring, 1);
return true;
}

成立条件:

1
2
3
4
sample_queue_push() 只由一个生产者线程调用;
sample_queue_pop() 只由一个消费者线程调用;
对象初始化完成后再启动两个线程;
队列销毁前两个线程已经停止。

如果同侧有多个调用者,必须在该侧增加外部串行化。


26. 批量写入的两种正确方式

26.1 允许环绕,一次逻辑申请

1
2
3
4
5
6
7
size_t n = requested;
size_t index = spscring_p_alloc(&queue->ring, &n);

for (size_t i = 0; i < n; i++)
queue->data[(index + i) % SAMPLE_QUEUE_SIZE] = input[i];

spscring_p_commit(&queue->ring, n);

优点:一次取得总容量。
缺点:每个元素需要取模或拆分成两段。

26.2 不允许环绕,分两次连续处理

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
size_t remaining = requested;
size_t offset = 0;

while (remaining) {
size_t n = remaining;
size_t index = spscring_p_alloc_no_wrap(&queue->ring, &n);
if (!n)
break;

memcpy(&queue->data[index], &input[offset], n * sizeof(input[0]));
spscring_p_commit(&queue->ring, n);

offset += n;
remaining -= n;
}

优点:每段都是连续内存,适合 memcpy()
缺点:跨界时需要两次申请与提交。


27. 等待回调的推荐接法

下面的回调只负责通知,不直接消费:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
struct consumer_context {
/* executor, semaphore or event object */
};

static void
sample_queue_readable(struct spscring *ring, void *arg)
{
(void)ring;
struct consumer_context *ctx = arg;

/* Post a task or release a semaphore. Keep this callback short. */
consumer_schedule(ctx);
}

static bool
sample_queue_wait_readable(struct sample_queue *queue,
struct consumer_context *ctx)
{
return spscring_c_submit_wait(
&queue->ring, 1, &sample_queue_readable, ctx) != 0;
}

调用方必须处理两种返回路径:

1
2
返回 0:现在已经可读,调用方应直接继续处理,回调不会发生。
返回 1:等待已经登记,稍后由生产者提交路径触发回调。

错误模式是只等待回调而忽略返回 0,这会在条件已经满足时造成任务停滞。


28. 调试时建议观察的状态与不变量

28.1 核心不变量

对生产者上下文:

1
2
3
pos < size
pos <= end
end - pos <= size

消费者上下文相同。

全局逻辑关系:

1
0 <= producer_absolute - consumer_absolute <= size

28.2 关键观测点

位置 建议记录
p_alloc 返回后 请求数量、实际数量、起始索引、p.ctx.base/pos/end
p_commit 前后 提交数量、新 p.pos、是否发生环绕
c_alloc 返回后 请求数量、实际数量、起始索引、c.ctx.base/pos/end
c_commit 前后 提交数量、新 c.pos、是否发生环绕
submit_wait 阈值、第一次容量、第二次容量、返回值
signal sig.size 状态、CAS 是否成功、回调执行线程
abort 是否观察到 SIZE_MAX、自旋时间、最终返回值

28.3 常见故障定位

现象 优先检查
偶发读取旧数据 是否在 p_commit() 前完成数据写入;是否禁用了原子支持
偶发被覆盖 消费者是否在读取完成前调用了 c_commit()
队列容量异常 是否同侧存在多个未串行化调用者;提交数量是否超过分配数量
永久不再唤醒 是否忽略 submit_wait() 返回 0;是否错误取消等待
abort_wait() 长时间自旋 回调接管窗口是否被阻塞;提交线程是否被异常挂起
回调死锁 commit() 时是否持有回调还会获取的锁
CAN 接收丢帧 rxring 是否满;消费任务是否阻塞;发送确认路径是否迫使继续读取

29. 单元测试应该覆盖哪些边界

Lely 自带测试重点覆盖:

  • 空队列生产者容量等于 size
  • 提交一个、多个和全部槽位后的容量;
  • 请求数量大于容量时被截断;
  • no_wrap 在数组尾部的连续容量;
  • 生产者提交触发消费者信号;
  • 消费者提交触发生产者信号;
  • 条件已满足时 submit_wait() 返回 0;
  • 条件不满足时等待成功注册;
  • abort_wait() 后回调不再执行。

工程移植时还应增加:

  1. 多轮绝对位置回绕测试;
  2. 批量跨物理边界读写测试;
  3. 高并发长时间顺序校验;
  4. 生产者和消费者不同 CPU 核绑定测试;
  5. ThreadSanitizer 或等效数据竞争检测;
  6. 回调注册、触发和取消的竞争测试;
  7. 队列容量为 1 的边界测试;
  8. 上层线程停止与对象销毁竞争测试。

30. 设计优点与限制

30.1 优点

  • 使用完整 N 个槽位,不浪费一个空槽;
  • 生产者与消费者各自维护本地上下文,减少共享原子访问;
  • 缓存行隔离降低伪共享;
  • acquire/release 正确发布用户数据;
  • 支持总容量与连续容量两类查询;
  • alloc/commit 允许批量和增量提交;
  • 等待协议通过双重检查避免丢失唤醒;
  • 不绑定线程、条件变量或操作系统,可嵌入 Executor、RTOS 与裸机适配层。

30.2 限制

  • 只支持一个逻辑生产者和一个逻辑消费者;
  • 只管理索引,不管理数据内存和对象生命周期;
  • 回调在对端提交线程同步执行,使用不当会阻塞或死锁;
  • 取消等待存在短暂忙等;
  • 对象体积明显大于简单读写索引环;
  • LELY_NO_ATOMICS 不应被当作通用多线程实现;
  • 上层仍需负责 shutdown、线程退出和缓冲区生命周期。

31. 最终心智模型

可以把 spscring 理解为一个只有两把“所有权印章”的仓库:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
用户数据数组
= 仓库货架

p_alloc
= 生产者取得空货架的临时写权限

p_commit
= 生产者盖章,宣布货物已上架,可供消费者领取

c_alloc
= 消费者取得已上架货物的临时读权限

c_commit
= 消费者盖章,宣布货架已清空,可重新使用

p.pos / c.pos
= 双方发布给对端的所有权边界

p.ctx / c.ctx
= 双方自己的本地工作账本

p.sig / c.sig
= 空间或数据不足时登记的通知条件