订阅与计费生命周期
让订阅经历试用期、周期扣费、价格调整和取消,同时在等待期间不丢失状态。
产品需求
Flow 以客户和订阅计划作为输入。计划包含试用时长、计费周期时长、最大计费周期数和当前扣费金额。它必须:
- 存储客户并初始化计费计数器。
- 发送欢迎邮件,然后等待试用期 Timer 触发,再开始计费。
- 等待每个计费周期 Timer,扣除当期费用,并重复直到达到配置的周期数。
- 在试用或计费 Timer 仍在等待时接受扣费金额更新。持久化新金额,让下一次扣费使用它。
- 随时接受取消请求。发送取消邮件并完成整个 Flow,包括仍在等待的 sibling Step。
- 在最后一个计费周期后结束订阅,并通知客户。
- 在 Flow 仍处于活动状态时,通过 Describe RPC 返回已存储的计划。
Flow 等待两类持久条件:试用期和每个计费周期的 Timer,以及取消和金额更新的 Channel 消息。Flow 可以在任一等待中停留数天或数月;Worker 重启不会重新开始试用期,也不会丢失已被接受的消息。
潜在故障需要产品层面的决策:
- 欢迎邮件、扣费或结束通知可能在外部提供方已经接受后失败。这些操作必须是幂等的,因为重试的 Step 可能再次调用提供方。
- 金额更新必须恰好包含一个金额。无效或意外消息会让更新 Step 失败,而不是静默选择一个金额。
- 取消可能与即将触发的计费 Timer 竞争。支付提供方和取消策略必须能够安全处理该竞争,包括退款或权益策略。
- 取消或正常完成后再发布消息不能重新打开 Flow。客户端应将其视为终态,并在需要时启动新的订阅。
Flow 设计
这个交互式 definition graph 由可运行的 Python Flow 生成。
初始化 Step 扇出为三个长期运行的分支:
Initialize → Trial → ChargeCurrentBill ↻
Initialize → Cancel
Initialize → UpdateChargeAmount ↻
Trial 管理试用期 Timer,并把控制权交给计费循环。ChargeCurrentBill 在下一个 Timer 前递增持久化的计费周期 Attribute;计划结束时强制完成,否则扣除当期费用。Cancel 和 UpdateChargeAmount 与该路径并行运行,二者每次各等待一条 Channel 消息。
它组合了以下设计模式:
- 扇出并发 从一个初始化 Step 独立启动计费、取消和价格控制职责。
- 持久 Timer 循环 用自迁移表达周期计费计划,而不是依赖常驻内存的进程。
- 事件驱动控制面 用 Channel 处理客户发起的取消和运营发起的价格调整。
- 持久状态 将客户和计费计数器保存在 Attribute 中,使每个分支读取当前计划,而不是保留可变的 Worker 内存。
- 强制完成 让取消和订阅结束成为停止 sibling 等待的终态决定。
Rust 示例拥有同样的持久状态、计费循环和控制分支。它的公开金额更新和取消 RPC 会向内部 Channel 发布命令;其他 SDK 示例直接向 Channel 发布。
核心实现
每个标签页都包含可运行的 Flow 实现:签名、Step 拓扑、持久化 schema 和读取 RPC。链接的示例还包含请求模型、依赖服务、HTTP controller 和端到端测试。
class SubscriptionFlow(Flow[Customer]):
billing_period_number = Attribute("billing-period-number", int)
customer_details = Attribute("customer", Customer)
cancel_subscription = Channel[None]("cancel-subscription", type(None))
update_charge_amount = Channel("update-charge-amount", int)
def __init__(self, service: MyDependencyService) -> None:
self.service = service
self.charge_current_bill = ChargeCurrentBill(
service,
self.customer_details,
self.billing_period_number,
)
self.trial = Trial(
service,
self.customer_details,
self.billing_period_number,
self.charge_current_bill,
)
self.cancel = Cancel(
service,
self.customer_details,
self.cancel_subscription,
)
self.update_charge_amount_step = UpdateChargeAmount(
self.customer_details,
self.update_charge_amount,
)
self.initialize = Initialize(
self.customer_details,
self.trial,
self.cancel,
self.update_charge_amount_step,
)
def get_steps(self) -> StepList[Customer]:
return StepList.start_step(self.initialize).other_steps(
self.trial,
self.charge_current_bill,
self.cancel,
self.update_charge_amount_step,
)
def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(
self.billing_period_number,
self.customer_details,
self.cancel_subscription,
self.update_charge_amount,
)
@rpc
def describe(self, context: Context) -> RPCResult[Subscription]:
return RPCResult(self.customer_details.get(context).subscription)
例子: examples/python/dex_examples/products/subscription/subscription_flow.py
演示
下方完成的 Go 运行存储了订阅,扇出 Trial、Cancel 和 UpdateChargeAmount,应用价格更新,然后通过取消完成。已完成节点是已完成的 Step execution;等待节点展示了取消关闭 Flow 时仍处于活动状态的持久条件。
打开图片可按原始分辨率查看。
