跳到主要内容

订阅与计费生命周期

让订阅经历试用期、周期扣费、价格调整和取消,同时在等待期间不丢失状态。

产品需求

Flow 以客户和订阅计划作为输入。计划包含试用时长、计费周期时长、最大计费周期数和当前扣费金额。它必须:

  1. 存储客户并初始化计费计数器。
  2. 发送欢迎邮件,然后等待试用期 Timer 触发,再开始计费。
  3. 等待每个计费周期 Timer,扣除当期费用,并重复直到达到配置的周期数。
  4. 在试用或计费 Timer 仍在等待时接受扣费金额更新。持久化新金额,让下一次扣费使用它。
  5. 随时接受取消请求。发送取消邮件并完成整个 Flow,包括仍在等待的 sibling Step。
  6. 在最后一个计费周期后结束订阅,并通知客户。
  7. 在 Flow 仍处于活动状态时,通过 Describe RPC 返回已存储的计划。

Flow 等待两类持久条件:试用期和每个计费周期的 Timer,以及取消和金额更新的 Channel 消息。Flow 可以在任一等待中停留数天或数月;Worker 重启不会重新开始试用期,也不会丢失已被接受的消息。

潜在故障需要产品层面的决策:

  • 欢迎邮件、扣费或结束通知可能在外部提供方已经接受后失败。这些操作必须是幂等的,因为重试的 Step 可能再次调用提供方。
  • 金额更新必须恰好包含一个金额。无效或意外消息会让更新 Step 失败,而不是静默选择一个金额。
  • 取消可能与即将触发的计费 Timer 竞争。支付提供方和取消策略必须能够安全处理该竞争,包括退款或权益策略。
  • 取消或正常完成后再发布消息不能重新打开 Flow。客户端应将其视为终态,并在需要时启动新的订阅。

Flow 设计

这个交互式 definition graph 由可运行的 Python Flow 生成。

Definition graph

SubscriptionFlow

Valid

python · examples/python/dex_examples/products/subscription/subscription_flow.py

初始化 Step 扇出为三个长期运行的分支:

InitializeTrialChargeCurrentBill

InitializeCancel

InitializeUpdateChargeAmount

Trial 管理试用期 Timer,并把控制权交给计费循环。ChargeCurrentBill 在下一个 Timer 前递增持久化的计费周期 Attribute;计划结束时强制完成,否则扣除当期费用。CancelUpdateChargeAmount 与该路径并行运行,二者每次各等待一条 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 运行存储了订阅,扇出 TrialCancelUpdateChargeAmount,应用价格更新,然后通过取消完成。已完成节点是已完成的 Step execution;等待节点展示了取消关闭 Flow 时仍处于活动状态的持久条件。

打开图片可按原始分辨率查看。

已完成的 SubscriptionFlow Step execution graph