TT Lab
开始
学习 学习路径 课程

ACK前停机的零食售货机

连接恢复不等于业务恢复

在 TT Lab 中继续学习

一句话总结

断开的连接可以重新打开,但要决定是否重做已经完成的工作,就必须把业务结果和处理游标放在同一个事务里保存。

为什么需要它

太空站里的零食自动售货机在接收事件。EVENT 0 7 的意思是把库存增加 7 个。接收端刚保存完库存,电源就断了,没能发出 ACK。对发布方来说,“已处理但应答丢失”和“根本没能处理”看起来一模一样。仅凭 TCP 重连成功,无法区分这两种情况。需要这样一种设计:重发不确定的事件,同时让接收一侧能认出自己已经处理过的事件。

前一门课程中的 Cursor 是运行中对象的记忆。重新启动进程后,这份记忆也会消失。这次我们把序号和合计保存到 SQLite 文件中,并在 ACK 之前真正终止接收进程。然后观察第二个进程能否打开同一个文件并接续处理。这不是搭建支付服务器的实验,而是用一条流中有序的库存变化量来揭示恢复边界的实验。

工作原理

先把字节变成契约

练习协议由 ASCII 的 EVENT、空格、序号、空格、变化量和一个 LF 构成。例如 EVENT 0 7 后面跟一个换行就是一个事件。序号是从 0 开始的连续整数,变化量范围是 -1000 到 1000。不允许前导 0、+ 号、-0、CRLF,以及没有换行的最后一个片段。严格的语法是作者定下的,用来区分“已经转换成数字就行了”和“这是约定好的消息”。如果其他协议允许 CRLF,就必须遵循那个协议的契约。

这种消息不是 WebSocket 或 SSE。这两个标准各自有独立的帧和重连规则。这里是在 TCP 字节流之上叠加一个简单的行协议,把注意力集中在处理语义上。每行上限为 64 字节,open_connection 的 StreamReader limit 也设置为 64。读到错误消息后如果不加分辨地从下一行继续,可能会掩盖实际缺失的命令,所以要关闭连接。这里没有实现 HTTP 认证、TLS 和按用户划分的权限,因此不能原样暴露到互联网上。

把业务效果和游标放在同一处

checkpoint 表中唯一的一行 id=1 保存 last 和 total。尚未处理任何事件时,last=-1,total=0。ledger 以 seq 为主键,保存已处理的 delta。本实验不会删除 ledger,处理历史的长期保留和压缩策略需要另行设计。

apply_event 用 BEGIN IMMEDIATE 打开写事务并读取当前的 last。如果 seq 等于 last+1,就写入 ledger,更新 total,推进 last,然后 COMMIT。中途出错则 ROLLBACK,并把原始错误返回给调用方。实验中的 fault 钩子在修改 total 之后、修改 last 之前执行。检查器此时在同一连接中观察到部分改动并注入错误,随后在另一个连接中确认没有留下任何改动。它还会在这个钩子处立即终止一个独立进程,以检查 finally 不会被执行的情形。不会拿一句“捕获了异常”的输出来充当原子性。

如果相同的 seq 和 delta 再次到来,就不累加结果,直接返回 False。seq 相同但 delta 不同则是 Conflict。这是为了不把重用事件 ID 却修改内容的发布方悄悄当作重复事件接受。如果 seq 大于 last+1,则是 Gap。处理到 4 之后收到 6,若直接跳到 last=6,就会误以为稍后到达的 5 已经处理过。应当先停下,再恢复缺失的区间。

这里通过 sqlite3.connect 的 isolation_level=None 关闭自动 BEGIN,并用 SQL 显式写出 BEGIN、COMMIT 和 ROLLBACK。不要与 Python 3.12 的 autocommit=False 方式混用而造成嵌套 BEGIN。with con 这种语法并不会关闭连接本身,所以由所有者在 finally 中 close。open_store 不会用初始值覆盖已有的行。如果重启时把 last 初始化为 -1,写文件就失去了意义。

ACK 在提交之后,恢复从游标之后开始

consume_line 在解析和 apply_event 都成功后,返回 ACK 序号和 LF。内容相同的重复事件也要 ACK。这样既不会再次累加结果,又能让发布方结束重试。如果还没提交就先发出 ACK 然后崩溃,发布方可能会把事件删掉。反过来,提交之后 ACK 丢失的情况下可能产生重复,但可以借助处理历史把它过滤掉。

run_client 会打开存储,通过 TCP 连接,并发送 RESUME last 和 LF。重新连接通过新的 run_client 调用显式完成,不会暗藏无限自动重试循环。在最终检查中,第一个接收端保存了 0 和 1 之后,在 ACK 1 之前通过 os._exit 终止。第二个进程发送 RESUME 1,收到服务器故意重发的 1 和新事件 2。检查会把服务器观察到的请求与 ACK,以及通过独立连接读出的 ledger 和合计一并核对。这不是偷换运行中对象的模拟重启。

超出保留日志范围时,恢复方式会不同

replay 会从发布方一侧的有限日志中,返回序号大于游标的条目,最多 limit 个。如果当前保留起点是 10,那么 last=9 可以从 10 开始读取,但 last=8 缺少所需的 9,因此是 ResyncRequired。只给最新条目就说成功,是悄无声息的丢失。实际产品应当提供一致的快照,并从该快照的游标开始重新同步,或者给出明确的恢复失败。本实验并不实现快照的生成与安装本身。

空日志在没有保留信息的情况下,无法区分“一开始就是空的”和“全部被截断了”,所以在这个 API 中会抛出 ValueError。输入限制为 1–128 个连续条目,批次限制为 1–16 个。实际的日志服务还需要保留下限、上限之类的独立元数据。Redis XREAD 同样会读取指定 ID 之后的条目,但不能把本实验中的整数序号和异常理解为 Redis 的真实行为。

在现场相遇的样子

通知、协作界面和 AI 流式输出,即便连接能长时间保持,也会因为部署和网络变更而断开。重连与重放、处理游标、保留策略必须一起设计。连接数增多后,重试还需要指数退避、抖动和最大次数,以及按账户的资源限制。这次只测试了两次连接和一个很小的日志,因此不能把它当作互联网延迟或大规模吞吐量的证据。

同一个 SQLite 事务中合计与游标的原子性,并不能把外部邮件、支付 API 也一并绑定进来。在 DB 提交与 HTTP 请求之间同样会出现新的故障窗口。对这类系统,需要另外学习接收侧的幂等键、outbox 之类的设计以及对账流程。“有了 ACK 就能在整个分布式系统中恰好一次”这种说法是错误的。

实验文件在进程重启期间会保留,但 LabHub 会话结束时会消失。这次的进程故障检查并没有测试主机断电、磁盘损坏或备份恢复。SQLite 的持久性同样以对文件系统和同步行为的某些前提为基础。请把观察到的故障类型和未做保证的类型区分开来汇报。

下一项实验要做什么

按帧解析 → 存储 → 原子处理 → 判断保留范围 → 提交后 ACK → 回收连接 → 真实进程重新连接的顺序完成 client.py。timeout 是每个网络 await 的等待上限,并不是包含整个会话和 DB 操作在内的总期限。SQLite 调用是同步的,而这个小实验只有单个接收端。应避免在生产环境的事件循环中直接执行耗时较长的 DB 操作。

通过官方文档进一步阅读