iotspool 原理:append-only 日志、CRC 截断恢复与退避调度

声明: iotspool 是一个开源项目,作者为 Vanderhell。本文是阅读该项目源码和文档后整理的学习笔记,用于理解嵌入式持久化消息队列的工程实现方式。本文作者不是该项目的开发者,未参与该项目的任何代码贡献。 文中所有工程细节均来自对开源代码的分析,不代表本文作者的设计决策。

项目仓库:github.com/Vanderhell/iotspool

上一篇讲了 iotspool 怎么用。本文拆解内部三件事:记录怎么编码、掉电怎么恢复、Broker 不可达时怎么退避。源码体量不大,核心逻辑都在 src/spool.csrc/record.c 两个文件里。

记录格式:ENQ 与 ACK 两种定长头

磁盘日志只有两种记录类型,所有整数小端序:

ENQ 与 ACK 记录格式

ENQ 是变长记录:头 20 字节 + topic + payload + 可选 SHA-256 + 4 字节 CRC。ACK 固定 12 字节(4 字节头 + 4 字节 msg_id + 4 字节 CRC)。

追加写与 CRC 截断恢复

核心设计选择是只追加、不原地修改

入队时,把 ENQ 记录拼到 scratch 缓冲区、算 CRC,调用 store.append() 写到文件尾,再调 store.sync()(POSIX 上就是 fsync)。sync 返回成功才算入队完成;如果在 append 过程中断电,这条记录只写了一半,下次启动会被识别为脏尾。

发送成功后不删除 ENQ,直接在文件尾再追加一条 ACK 记录指向同一个 msg_id。整个写入路径始终是 append-only,实现简单且掉电安全。

恢复过程(iotspool_recover)是顺序扫一遍整个 store 文件:

  1. 读 superblock,校验 magic/version/generation。
  2. 循环读下一条记录:根据 type 判断是 ENQ 还是 ACK,读固定头,算出整个记录长度(含 payload 和 CRC)。
  3. 若剩余字节不够读完整记录,说明写到一半断电了,truncate_to 当前位置,停止。
  4. 读完整记录后校验 CRC32,CRC 不对同样截断到当前位置,停止。
  5. CRC 正确的 ENQ 放进 RAM 索引,标记为 pending;ACK 把对应 msg_id 从索引里移除。
  6. 扫到文件正常结束,恢复完成。

截断到"最后一条完整记录的尾部"是关键:这保证了日志的一致性视图永远只包含已经完整 sync 过的记录,半截的永远不会被 replay。ACK 丢了最坏情况是消息重发一次(at-least-once),不会重复投递同一条多次,因为 ACK 也是 append 的。

追加写日志与 compact 前后对比

RAM 环形索引

所有 ENQ 记录的磁盘位置不会常驻内存。RAM 里只保存一个 iotspool_entry_t 环形数组:每条 entry 包含 msg_idgenerationrecord_offsetrecord_lentopic_lenpayload_lenqosretain

这就是 STM32 上开 64 条 pending 也只占几百字节 RAM 的原因:真正的数据只在要发送的那一刻才读进 scratch。

compact:重写消除 ACK 碎片

一直 append 会让 store 文件越滚越大:已经被 ACK 的 ENQ 和对应的 ACK 记录都成了历史垃圾。当 store 大小达到 max_store_bytes 的某个阈值时,库会触发 iotspool_compact()

  1. 在 scratch 或工作缓冲区里重建一份新的 superblock。
  2. 遍历当前 RAM 索引里所有 pending 的 ENQ,用 read_at 从原文件读出完整记录,直接追加到新镜像。
  3. 新镜像里只包含未被 ACK 的 ENQ,没有 ACK 记录,没有已确认的历史消息。
  4. 调用 store.replace() 原子替换整个文件(POSIX 上通常是写临时文件 + rename)。
  5. generation 加 1,索引里的 entry 重新指向新文件中的 offset。

compact 完成后,store 大小回到"所有 pending 消息的体积之和",不会无限增长。drop_oldest_on_full=true 时,入队发现空间不够会先尝试丢最老的 pending 再 compact。

退避算法:Full Jitter

发送失败时(iotspool_on_publish_fail 或新 API 的 iotspool_publish_failed),库使用 AWS 架构博客推广的 Full Jitter 指数退避:

sleep = random_between(0, cap)
cap   = min(cap * 2, max_retry_ms)

初始 cap 是 min_retry_ms,每次失败翻倍直到 max_retry_ms 封顶。每次成功 ACK 后 cap 重置回 min_retry_ms

加随机抖动是为了避免经典的"雪崩":一组设备同时断线、同时等到固定退避窗口到期、同时重连把 Broker 打死。Full Jitter 让重试时间均匀分布在 [0, cap],打散峰值。

iotspool_peek_ready(now_ms) 会比较当前时间和 retry_deadline_ms,没到点就返回 IOTSPOOL_ENOTFOUND,调用方的主循环可以去做别的事,不用自己算等待时间。

存储后端 vtable

库本身不做任何系统调用,所有 I/O 走 iotspool_store_t 的 6 个回调:

回调 职责
append(ctx, data, len) 把数据写到当前文件尾
read_at(ctx, off, out, cap, *out_len) 从指定偏移读数据
sync(ctx) 落盘(fsync 或 Flash 擦写确认)
size_bytes(ctx) 返回当前 store 大小
truncate_to(ctx, new_size) 恢复时截掉脏尾
replace(ctx, data, len) compact 时原子替换整个文件

POSIX 后端 store_posix.c 把这 6 个回调映射到 write/pread/fsync/fstat/ftruncate/tmpfile+rename。ESP-IDF 直接复用 POSIX 后端走 VFS。裸机上只要把这几个对接到你自己的 Flash 驱动、LittleFS、或 FAL 分区即可,核心代码不用改一行。

并发模型

库默认假设调用方自己序列化所有公开 API 调用(通常是单任务主循环 + 事件队列的嵌入式模型)。如果网络任务和采样任务会并发访问 spool,在 cfg 里填 lock/unlock 两个回调(FreeRTOS 里就是 xSemaphoreTakeRecursive / xSemaphoreGiveRecursive),核心在改内部状态和调 store 回调前后自动加锁。

锁回调必须不可重入。store 回调内部不能再调任何 iotspool API,否则死锁。ISR 里调用必须让锁回调返回 IOTSPOOL_EBUSY,或者只在 ISR 里发事件、真正 enqueue 放到任务里做。

源码阅读入口

文件极少,读的顺序推荐:

  1. include/iotspool.h:公开 API、错误码、配置结构、entry/inflight 结构。
  2. src/record.c:记录编解码、CRC32、可选 SHA-256。
  3. src/spool.c:init/recover/enqueue/peek/ack/compact 主流程。
  4. src/backoff.c:Full Jitter 实现,几十行。
  5. src/store_posix.c:POSIX 后端,作为移植参考。
  6. tests/test_main.c:覆盖掉电截断、recover、backoff 等关键路径的单测。

设计上的三件事值得记下来: