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

ACK前停机的零食售货机

零食订单已保存,仓库却不知道

在 TT Lab 中继续学习

目标

分别以原子方式保存业务与发送意图、接收 ID 与业务效果,并在真实的 HTTP 重发和进程终止之后完成恢复。

为什么重要

如果仅仅因为没有响应就断定业务没有被执行,就会造成重复反映。反过来,如果在发送之前就标记为完成,就会丢失尚未投递的业务。用 Python 和 SQLite 把两个边界分开来实现,并确认在清理代码都不会执行的进程终止之后,仍然留存的状态。

预计需要 85 分钟。比默认会话更长,所以请在到期之前点击“+时间”按钮延长。会话结束后文件会消失,需要的代码请另行保存。开始前应已了解 Python 函数、异常、SQLite 事务,以及前面模块中的 ACK 恢复。

数据契约

所有提交的代码都是 /root/outbox/worker.py。下面这些函数收到的 con 是角色正确的已打开连接,调用前没有事务。除 open_store 之外的函数不会关闭借来的连接。无论成功还是失败,都不留下未结束的事务。source 和 sink 必须用各自独立的文件来测试,并且不连接生产 DB。

source schema:

CREATE TABLE orders (id TEXT PRIMARY KEY, qty INTEGER NOT NULL);
CREATE TABLE outbox (seq INTEGER PRIMARY KEY AUTOINCREMENT,
    id TEXT NOT NULL UNIQUE, qty INTEGER NOT NULL,
    sent INTEGER NOT NULL CHECK(sent IN (0,1)));

sink schema:

CREATE TABLE inbox (id TEXT PRIMARY KEY, qty INTEGER NOT NULL);
CREATE TABLE stock (id INTEGER PRIMARY KEY CHECK(id=1), total INTEGER NOT NULL);

stock 的初始行是 (1,0)。每个角色的初始化都在一个事务中完成,并保留已有内容。检查器会创建并清理临时 DB、loopback HTTP 服务器和子进程,所以不要把 DB 路径或固定端口写进代码。

步骤

  1. 制定发送 ID 与数量的契约——在 worker.py 中实现 Exception 的子类 Conflict、DeliveryError 以及 validate_event(event)。event 是恰好只有 id、qty 两个键的 dict。id 是由 ASCII 字母、数字、下划线、连字符组成的 1–64 个字符,qty 是不含 bool 的 int 1–1000。无效时为 ValueError,有效时不修改输入,并返回新的 dict 副本。
  2. 初始化两个不同的存储——open_store(path, role) 返回 source 或 sink 角色的 sqlite3.Connection。以 isolation_level=None、busy timeout 1 秒打开,并仅在不存在时创建下面的 schema。sink 的 stock 初始行 (1,0) 也仅在不存在时插入。初始化在一个事务中完成,失败时关闭连接并传递错误。不要用初始值覆盖已有的记录。其他 role 在连接之前就是 ValueError。
  3. 同时留下订单与发送意图——enqueue(con,event,fault=None) 在 validate_event 之后,于 BEGIN IMMEDIATE 中处理。新 ID 按 orders 插入 → fault('after-order') → outbox(sent=0) 插入 → fault('after-outbox') → COMMIT → fault('after-commit') 的顺序进行,并返回 True。fault 只在存在时才调用。已有 ID 且 qty 相同则不做修改并返回 False,qty 不同则为 Conflict。提交之前的错误全部回滚并传递原始错误,提交之后的错误不会删除已确定的记录。
  4. 按顺序读取尚未完成的批次——pending(con,limit) 按 seq 升序返回 sent=0 的 outbox,最多 limit 个。每个元素都是只有 id、qty 的新 dict。limit 是不含 bool 的 int 1–16,其他值为 ValueError。查询不会修改 DB 或输入。
  5. 把重复接收的效果限制为一次——receive(con,event,fault=None) 在 sink 连接上运行。验证之后,在 BEGIN IMMEDIATE 中按新 ID 的 inbox 插入 → fault('after-inbox') → 把 qty 加到 stock.total → fault('after-total') → COMMIT → fault('after-commit') 的顺序进行,并返回 True。相同 ID、相同 qty 返回 False,相同 ID、不同 qty 为 Conflict。不同 ID、相同 qty 是独立业务。提交之前的错误整体回滚,提交之后的错误保留已确定的状态,并传递原始错误。
  6. 只对已确认的条目留下完成标记——mark_sent(con,event) 只把与已验证的 id、qty 一致的 source outbox 的 sent 改为 1。不存在的 ID 为 ValueError,qty 不同则为 Conflict,失败时不改变 DB 状态。新标记返回 True,已经完成则返回 False。不删除 orders 和 outbox 的行,也不留下事务。
  7. ACK 之后再标记,并在第一次失败时停止——dispatch(con,send,limit=8,fault=None) 只按顺序处理 pending 的有限批次。把 event 副本交给 send,在收到严格的 True ACK 之后,调用 fault('after-delivery') 和 mark_sent。False、None、1、字符串为 DeliveryError,send 或 fault 的异常按原样传递。在第一次失败时停止,并保留此前已完成的部分。返回值是本次标记为完成的个数。发送期间不占用写事务,即使回调修改了参数,也只把最初选定的条目标记为完成。
  8. 借助 HTTP 与真实的进程重启来恢复——实现 send_http(url,event,timeout=1)。在 validate_event 之后发送 HTTP POST JSON,只有同时满足以下条件才返回 True:状态为 200,响应是不超过 1024 字节、恰好含 id 与 accepted 两个键的 JSON dict,id 与请求相同,且 accepted is True。无效的 ACK 为 DeliveryError,HTTP 和连接错误则传递出去,并关闭响应。timeout 是不含 bool、不超过 5 秒的正有限 int/float,否则为 ValueError。URL 只允许 http://127.0.0.1:포트/events만 (其中端口为占位符)这种形式,且必须是明确指定的端口 1–65535,不得含有用户信息、查询和 fragment。不使用自动代理和重定向。最终检查会在接收提交之后真正终止发布方,并确认新进程重发相同 ID 后库存只增加一次。

参考

制定发送 ID 与数量的契约

在 worker.py 中实现 Exception 的子类 Conflict、DeliveryError 以及 validate_event(event)。event 是恰好只有 id、qty 两个键的 dict。id 是由 ASCII 字母、数字、下划线、连字符组成的 1–64 个字符,qty 是不含 bool 的 int 1–1000。无效时为 ValueError,有效时不修改输入,并返回新的 dict 副本。

int(True) 等于 1 这一事实,与数量是否有效的契约是两回事。正则表达式要检查整个字符串。

初始化两个不同的存储

open_store(path, role) 返回 source 或 sink 角色的 sqlite3.Connection。以 isolation_level=None、busy timeout 1 秒打开,并仅在不存在时创建下面的 schema。sink 的 stock 初始行 (1,0) 也仅在不存在时插入。初始化在一个事务中完成,失败时关闭连接并传递错误。不要用初始值覆盖已有的记录。其他 role 在连接之前就是 ValueError。

CREATE TABLE IF NOT EXISTS 与 INSERT OR IGNORE 的作用不同。成功时不要留下未结束的事务。

同时留下订单与发送意图

enqueue(con,event,fault=None) 在 validate_event 之后,于 BEGIN IMMEDIATE 中处理。新 ID 按 orders 插入 → fault('after-order') → outbox(sent=0) 插入 → fault('after-outbox') → COMMIT → fault('after-commit') 的顺序进行,并返回 True。fault 只在存在时才调用。已有 ID 且 qty 相同则不做修改并返回 False,qty 不同则为 Conflict。提交之前的错误全部回滚并传递原始错误,提交之后的错误不会删除已确定的记录。

如果在两个 INSERT 之间放入 COMMIT,进程终止时就只会留下订单。必须保留 fault 点,才能测试中间状态。

按顺序读取尚未完成的批次

pending(con,limit) 按 seq 升序返回 sent=0 的 outbox,最多 limit 个。每个元素都是只有 id、qty 的新 dict。limit 是不含 bool 的 int 1–16,其他值为 ValueError。查询不会修改 DB 或输入。

ID z 可能比 a 更早受理。不要用业务 ID 的字典序,而要使用保存下来的插入顺序。

把重复接收的效果限制为一次

receive(con,event,fault=None) 在 sink 连接上运行。验证之后,在 BEGIN IMMEDIATE 中按新 ID 的 inbox 插入 → fault('after-inbox') → 把 qty 加到 stock.total → fault('after-total') → COMMIT → fault('after-commit') 的顺序进行,并返回 True。相同 ID、相同 qty 返回 False,相同 ID、不同 qty 为 Conflict。不同 ID、相同 qty 是独立业务。提交之前的错误整体回滚,提交之后的错误保留已确定的状态,并传递原始错误。

不要把接收记录与业务效果分开提交。False 不是失败,而是本次没有应用新效果的有效重复。

只对已确认的条目留下完成标记

mark_sent(con,event) 只把与已验证的 id、qty 一致的 source outbox 的 sent 改为 1。不存在的 ID 为 ValueError,qty 不同则为 Conflict,失败时不改变 DB 状态。新标记返回 True,已经完成则返回 False。不删除 orders 和 outbox 的行,也不留下事务。

保留日后要对照的业务与意图。不要在没有 WHERE id 条件的情况下把整个批次都标记为完成。

ACK 之后再标记,并在第一次失败时停止

dispatch(con,send,limit=8,fault=None) 只按顺序处理 pending 的有限批次。把 event 副本交给 send,在收到严格的 True ACK 之后,调用 fault('after-delivery') 和 mark_sent。False、None、1、字符串为 DeliveryError,send 或 fault 的异常按原样传递。在第一次失败时停止,并保留此前已完成的部分。返回值是本次标记为完成的个数。发送期间不占用写事务,即使回调修改了参数,也只把最初选定的条目标记为完成。

如果先标记为完成,在发送之前终止就会丢失未完成的条目。不要把 receive 的 False 与 HTTP ACK 失败混为一谈。

借助 HTTP 与真实的进程重启来恢复

实现 send_http(url,event,timeout=1)。在 validate_event 之后发送 HTTP POST JSON,只有同时满足以下条件才返回 True:状态为 200,响应是不超过 1024 字节、恰好含 id 与 accepted 两个键的 JSON dict,id 与请求相同,且 accepted is True。无效的 ACK 为 DeliveryError,HTTP 和连接错误则传递出去,并关闭响应。timeout 是不含 bool、不超过 5 秒的正有限 int/float,否则为 ValueError。URL 只允许 http://127.0.0.1:포트/events만 (其中端口为占位符)这种形式,且必须是明确指定的端口 1–65535,不得含有用户信息、查询和 fragment。不使用自动代理和重定向。最终检查会在接收提交之后真正终止发布方,并确认新进程重发相同 ID 后库存只增加一次。

请查看 urllib.request 的 ProxyHandler({}) 和 HTTPRedirectHandler。即使请求到达两次,仓库的效果也必须只有一次。