Deal DSL
把持久化交易定义成数据,再用同一个解释器处理任意可销售的商品。
产品需求
这个流程表示一位买家针对一件商品发起的一次交易。商品可以是数据集、实体商品、订阅、服务或其他可销售内容。流程必须:
- 接收一份定义,其中包含流程 ID、商品 ID 和名称、初始状态、初始状态数据、状态、动作、等待条件和状态转换。
- 在执行前校验定义。所有被引用的状态和动作都必须存在,标识符必须有效,每个条件名称只能对应一个等待点。
- 只运行一次 InitializeDeal,保存不可变的定义快照、买家和商品身份、初始状态数据以及第一个待进入状态。
- 在进入状态前运行 WaitForDealCondition。如果状态声明了前置条件,Step 会等待该条件对应的键控 Channel 消息;否则立即继续。
- 把条件消息合并进 DealStateData,再按声明顺序执行状态动作。每个动作都有独立的 ExecuteDealAction execution 和持久化检查点。
- 动作结束后运行 EvaluateDealTransition。它可以等待另一个外部条件,根据状态数据按顺序匹配等值分支,并选择匹配状态或兜底状态。
- 重复解释器循环,直到进入没有后续转换的状态。Flow 最终用累计的状态数据完成。
- 保证运行中的 execution 不受后续定义修改影响。Go 应用把可编辑定义保存在 PostgreSQL 中,并在启动新 Flow 时传入快照。
Flow 会等待前置条件消息、状态转换消息,以及 Step 失败后的重试间隔。等待期间不会占用 Worker。
潜在失败包括:
- 定义缺少商品、引用不存在的初始或目标状态、重复使用条件名称、标识符无效或动作未注册时,校验会拒绝定义。
- 条件消息可能格式错误、指向不存在或已经停止的 Flow,或者使用定义中未暴露的条件名称。
- 一个键控等待必须且只能得到一条消息。结果不唯一时,Step 会失败,不会合并有歧义的输入。
- 支付、退款、样品交付或商品交付可能超时或报错。Dex 只重试失败的 Step,不会重跑之前已经完成的动作。
- 外部动作可能已经成功,但响应在返回途中丢失。生产适配器应使用 Flow ID 和 Step execution 标识生成幂等键。
- 状态数据可能缺少判断字段。此时状态转换会走显式兜底分支。
- 定义可能形成无法自然结束的业务循环。Flow 会保持持久化运行,但产品层校验或超时策略需要限制非预期循环。
Flow 设计
这个交互式 definition graph 由可运行的 Python Flow 生成。
这个 Flow 是一个小型持久化解释器。定义决定业务路径,四种 Step 提供执行机制:
DealStart(定义快照、买家、商品)
│
▼
InitializeDeal
│
▼
WaitForDealCondition ◄──────────────┐
等待键控 Channel │
│ │
▼ │
ExecuteDealAction │
每个动作一个 execution │
│ │
▼ │
EvaluateDealTransition │
可选等待 + 分支 / 兜底 │
│ │
有下一状态?── 是 ─────────────┘
│ 否
▼
用累计状态数据完成
设计使用了以下模式:
- 解释器模式让可执行 Flow 保持稳定,由流程定义提供状态、等待、动作和状态转换。
- 有限状态机显式记录当前状态和每一次转换,便于检查和调试。
- 不可变 execution 快照保证目录中的定义修改不会改变运行中的交易。
- 键控 Channel 会合把每个命名外部条件映射到独立的持久化收件箱,无需为每个条件声明新的 Channel 类型。
- 单动作检查点隔离重试并精确记录失败动作,多个有序动作不会被折叠成一个不透明 Step。
- 带兜底的守卫转换按顺序比较共享状态数据中的值,并始终提供明确的默认路径。
- 黑板状态让条件消息和动作把小型键值更新汇总到持久化的 DealStateData map。
- Entity store 为每个 Go execution 提供稳定的 Flow ID,并把流程、商品、买家、当前状态和等待条件保存成可搜索 Attribute。
各语言的通用 SDK 示例为每个状态使用一个动作列表。Go 产品示例在同一解释器循环上扩展了前置动作、后置动作、PostgreSQL 定义目录、搜索 API 和浏览器 UI。
核心实现
每个标签都包含 Flow signature、持久化 schema 和四个 Step 的完整路由。所有代码都来自链接中的可运行示例。
class DealDSLFlow(Flow[DealStart]):
definition = Attribute("DealDefinition", DealDefinition)
state_data = Attribute("DealStateData", dict[str, str])
condition_messages = ChannelMap("DealConditionMessages", dict[str, str])
def get_steps(self) -> StepList[DealStart]:
return StepList.start_step(self.initialize).other_steps(
self.wait_for_condition,
self.execute_action_step,
self.evaluate_transition,
)
def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(
self.definition,
self.state_data,
self.process_id,
self.item_id,
self.buyer_id,
self.current_state,
self.pending_condition,
self.condition_messages,
)
class WaitForDealCondition(Step[StateStepInput]):
def wait_for(self, context: Context, input: StateStepInput) -> Wait:
state = self.flow.definition.get(context).state(input.state_name)
if state.pre_condition is None:
return Wait.skip_immediately()
return Wait.until(self.flow.condition_messages.for_one(state.pre_condition.name))
def execute(self, context: Context, input: StateStepInput) -> StepDecision:
state = self.flow.definition.get(context).state(input.state_name)
if state.pre_condition is not None:
self.flow.merge_condition(context, state.pre_condition.name)
if state.actions:
return go_to(ExecuteDealAction, ActionStepInput(state.name, 0))
return go_to(EvaluateDealTransition, input)
class ExecuteDealAction(Step[ActionStepInput]):
def execute(self, context: Context, input: ActionStepInput) -> StepDecision:
state = self.flow.definition.get(context).state(input.state_name)
self.flow.execute_action(context, state.actions[input.action_index])
next_index = input.action_index + 1
if next_index < len(state.actions):
return go_to(ExecuteDealAction, ActionStepInput(state.name, next_index))
return go_to(EvaluateDealTransition, StateStepInput(state.name))
class EvaluateDealTransition(Step[StateStepInput]):
def execute(self, context: Context, input: StateStepInput) -> StepDecision:
transition = self.flow.definition.get(context).state(input.state_name).transition
if transition is None:
return graceful_complete(self.flow.state_data.get(context))
state_data = self.flow.state_data.get(context)
next_state = transition.else_state
for deal_case in transition.cases:
if state_data.get(transition.key) == deal_case.equals:
next_state = deal_case.go_to_state
break
return go_to(WaitForDealCondition, StateStepInput(next_state))
例子: examples/python/dex_examples/products/deal_dsl/deal_dsl_flow.py
Demo
这次完成的 Go 运行把 Premium research package 商品卖给了 buyer-demo。Flow 等待 buyer-confirmation,向买家收款,交付商品,并以 itemDeliveryStatus: delivered 完成。每个动作在图中都有独立且已完成的 ExecuteDealAction execution。