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

ACK前停机的零食售货机

恢复在ACK前停机的零食售货机

在 TT Lab 中继续学习

目标

实现持久游标、原子处理和重放范围,并通过真实的 TCP 重新连接来防止重复处理。

为什么重要

连接恢复与业务恢复是不同的。通过亲自检查在保存与确认之间崩溃的接收端,来学习二者的边界。建议已掌握 Python 函数、异常、async/await 和 SQL 基础,并学过前面的实时订阅者课程。不需要安装任何东西,也不需要互联网。这是 75 分钟的实验,所以请在默认 60 分钟的会话中用“+时间”按钮延长。会话结束时所有文件都会消失,需要的代码请另外保存。

步骤

  1. 不把残缺的最后一行当作命令接收——实现 parse_line(line)。只接收 bytes,把 64 字节以内的 ASCII EVENT 序号 变化量 LF 作为 (seq, delta) 元组返回。seq 为 0–2147483647,delta 为 -1000–1000。除数字 0 之外的前导 0、+ 号、-0、CRLF、缺少换行、多行以及非 ASCII,均为 ValueError。
  2. 打开重启后仍能记住的存储——open_store(path) 返回 sqlite3 连接。仅在不存在时创建 checkpoint(id INTEGER PRIMARY KEY CHECK(id=1), last INTEGER NOT NULL, total INTEGER NOT NULL) 和 ledger(seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL) 两张表,并仅在不存在时插入 checkpoint 的初始行 (1,-1,0)。state(con) 是 (last,total) 元组。为了能使用显式 SQL 事务,请以 isolation_level=None、busy timeout=1 秒建立连接。初始化失败时关闭连接并传递错误。
  3. 合计与游标要么一起保存,要么一起撤销——声明 Exception 的子类 Gap 和 Conflict,并实现 apply_event(con,seq,delta,fault=None)。只允许除 bool 之外、处于第 1 步范围内的 int,否则为 ValueError。seq=last+1 时,把 ledger 插入 → total 更新 → 若有 fault 则 fault() → last 更新放在一个事务中提交,并返回 True。seq<=last 时,若 ledger 中的 delta 相同则返回 False,内容不同则为 Conflict。更大的序号是 Gap。失败时回滚所有改动并传递原始错误,不要遗留事务。
  4. 不把被截断的记录说成恢复成功——声明 Exception 的子类 ResyncRequired,并实现 replay(events,last,limit)。events 是由 1–128 个 (seq,delta) 元组组成的 list,所有序号连续且升序,值为第 1 步范围内的 int。last 是 -1–2147483647 的 int,limit 是 1–16 的 int。bool、类型与范围错误、空日志、last 大于最后一个序号,均为 ValueError。last 小于首个序号减 1 时为 ResyncRequired,其余情况则把 seq>last 的条目最多 limit 个以新 list 返回。先检查全部输入,且不修改输入。
  5. 只有提交之后才生成业务 ACK——consume_line(con,line,fault=None) 调用 parse_line 和 apply_event,成功后返回 b"ACK 序号\n"。请把 fault 传给 apply_event。内容相同的重复也要 ACK,但解析失败、Gap、冲突、业务失败要传递错误,不发送 ACK。
  6. 即使出错和取消也要关闭连接——实现 async consume_session(reader,writer,con,timeout,before_ack=None)。timeout 是不含 bool 的正有限 int/float,否则为 ValueError。每次 readline 等待都用 timeout 限制,遇到 EOF 则返回 state。用 consume_line 提交一行之后,如果有 before_ack,就调用 before_ack(seq),执行 writer.write(ack) 以及受 timeout 限制的 drain。传递错误与取消,无论正常、失败还是取消,都按 close → 受 timeout 限制的 wait_closed 来回收。关闭等待时的 TimeoutError、ConnectionError 要调用 writer.transport.abort,且不要掩盖原有错误。不关闭 DB 连接。对于无效 timeout,关闭等待使用 1 秒。
  7. 终止接收端,并在同一个 DB 上接续——实现 async run_client(host,port,path,timeout,before_ack=None)。先验证 timeout,然后依次执行 open_store(path)、受 timeout 限制的 asyncio.open_connection(host,port,limit=64)。把 DB 中的 last 写成 b"RESUME last\n",经受 timeout 限制的 drain 之后交给 consume_session,并返回其返回值。所有路径都要关闭 DB 连接,在交给 consume_session 之前就失败时,要亲自回收套接字。传递错误与取消。检查器会创建临时 DB 和 loopback 服务器,在提交 1 与 ACK 1 之间终止第一个进程,然后用新进程重新连接。不得把重发的 1 重复计入。

参考

先执行 mkdir -p /root/resume,然后在 /root/resume/client.py 中实现。检查器也会运行此前的步骤,并准备和回收临时 DB、loopback 端口和接收进程。学习者无需一直开着服务器或 DB 文件。不要修改提交的文件。每一步的运行限制为 15 秒,并不是限制学习时间。最后的检查会真正终止进程,所以不要假定清理用的 finally 会被执行。本实验并没有实现面向生产的 TLS、认证、自动退避、快照恢复和分布式 exactly-once。

不把残缺的最后一行当作命令接收

实现 parse_line(line)。只接收 bytes,把 64 字节以内的 ASCII EVENT 序号 变化量 LF 作为 (seq, delta) 元组返回。seq 为 0–2147483647,delta 为 -1000–1000。除数字 0 之外的前导 0、+ 号、-0、CRLF、缺少换行、多行以及非 ASCII,均为 ValueError。

在用 strip 去掉最后的换行之前,先确认这是一个完整的帧。

打开重启后仍能记住的存储

open_store(path) 返回 sqlite3 连接。仅在不存在时创建 checkpoint(id INTEGER PRIMARY KEY CHECK(id=1), last INTEGER NOT NULL, total INTEGER NOT NULL) 和 ledger(seq INTEGER PRIMARY KEY, delta INTEGER NOT NULL) 两张表,并仅在不存在时插入 checkpoint 的初始行 (1,-1,0)。state(con) 是 (last,total) 元组。为了能使用显式 SQL 事务,请以 isolation_level=None、busy timeout=1 秒建立连接。初始化失败时关闭连接并传递错误。

如果重置已保存的游标,就不是重连,而是从头开始重复处理。

合计与游标要么一起保存,要么一起撤销

声明 Exception 的子类 Gap 和 Conflict,并实现 apply_event(con,seq,delta,fault=None)。只允许除 bool 之外、处于第 1 步范围内的 int,否则为 ValueError。seq=last+1 时,把 ledger 插入 → total 更新 → 若有 fault 则 fault() → last 更新放在一个事务中提交,并返回 True。seq<=last 时,若 ledger 中的 delta 相同则返回 False,内容不同则为 Conflict。更大的序号是 Gap。失败时回滚所有改动并传递原始错误,不要遗留事务。

无论先提交游标,还是先单独提交合计,一旦出现故障,两者都会不一致。

不把被截断的记录说成恢复成功

声明 Exception 的子类 ResyncRequired,并实现 replay(events,last,limit)。events 是由 1–128 个 (seq,delta) 元组组成的 list,所有序号连续且升序,值为第 1 步范围内的 int。last 是 -1–2147483647 的 int,limit 是 1–16 的 int。bool、类型与范围错误、空日志、last 大于最后一个序号,均为 ValueError。last 小于首个序号减 1 时为 ResyncRequired,其余情况则把 seq>last 的条目最多 limit 个以新 list 返回。先检查全部输入,且不修改输入。

保留起点 10 与游标 9 是衔接的,但游标 8 缺少 9。

只有提交之后才生成业务 ACK

consume_line(con,line,fault=None) 调用 parse_line 和 apply_event,成功后返回 b"ACK 序号\n"。请把 fault 传给 apply_event。内容相同的重复也要 ACK,但解析失败、Gap、冲突、业务失败要传递错误,不发送 ACK。

ACK 之后,从另一个 DB 连接也应能看到处理结果。

即使出错和取消也要关闭连接

实现 async consume_session(reader,writer,con,timeout,before_ack=None)。timeout 是不含 bool 的正有限 int/float,否则为 ValueError。每次 readline 等待都用 timeout 限制,遇到 EOF 则返回 state。用 consume_line 提交一行之后,如果有 before_ack,就调用 before_ack(seq),执行 writer.write(ack) 以及受 timeout 限制的 drain。传递错误与取消,无论正常、失败还是取消,都按 close → 受 timeout 限制的 wait_closed 来回收。关闭等待时的 TimeoutError、ConnectionError 要调用 writer.transport.abort,且不要掩盖原有错误。不关闭 DB 连接。对于无效 timeout,关闭等待使用 1 秒。

ACK 之前的钩子会重现“保存已完成,但发布方并不知道”的故障窗口。

终止接收端,并在同一个 DB 上接续

实现 async run_client(host,port,path,timeout,before_ack=None)。先验证 timeout,然后依次执行 open_store(path)、受 timeout 限制的 asyncio.open_connection(host,port,limit=64)。把 DB 中的 last 写成 b"RESUME last\n",经受 timeout 限制的 drain 之后交给 consume_session,并返回其返回值。所有路径都要关闭 DB 连接,在交给 consume_session 之前就失败时,要亲自回收套接字。传递错误与取消。检查器会创建临时 DB 和 loopback 服务器,在提交 1 与 ACK 1 之间终止第一个进程,然后用新进程重新连接。不得把重发的 1 重复计入。

如果总是发送 last 为 -1,或者每次都清空存储,在新进程的真实请求中就会暴露出来。