用户入驻流程
引导新用户完成账号创建、邮箱验证和两个必需的设置任务。
产品需求
流程接收用户名、邮箱地址、名字和姓氏。它必须:
- 保存注册表单,将用户标记为等待验证,并发送验证邮件。
- 等待用户验证邮箱。示例中每 24 秒发送一次提醒,直到收到验证。
- 要求已验证的用户完成任务 1,然后等待任务完成。任务未完成时持续发送任务 1 提醒。
- 只有任务 1 完成后,才要求用户完成任务 2。任务未完成时持续发送任务 2 提醒。
- 任务 2 完成后,将入驻状态标记为完成,发送欢迎邮件,并完成 Flow。
- 通过 RPC 暴露邮箱验证和两个任务完成操作,让应用可以在 Flow 等待时推进流程。
每个等待 Step 都通过 AnyOf 同时等待 Timer 和 Channel。Channel 消息会立即推进用户;Timer 会发送提醒,并将同一个 Step 移回自身。Worker 在任何等待期间重启,都不会丢失当前阶段、提醒计划或用户操作。
潜在故障和例外结果包括:
- 重复用户名不能使用同一个 Flow ID 启动第二个 Flow。启动端点会报告该入驻流程已经存在。
- 邮箱验证或任务可能始终没有完成。Flow 会持久等待并发送提醒,直到流程完成、被取消或达到配置的 Flow 超时。
- 任务完成操作可能乱序到达。对应 RPC 会返回该任务当前不在等待状态,并且不会发布 Channel 消息。
- Step method 或 RPC 可能在返回前失败。Dex 不会 commit 该 method 的任何 Attribute 写入或 Channel publish,因此重试会从之前已 commit 的 durable state 开始。
- 邮件发送可能发生瞬时故障。邮件适配器应向上返回错误,让 Dex 重试该 Step。
- 邮件供应商可能已经接受邮件,但 Worker 尚未记录成功。生产环境中的邮件调用需要幂等键或等效的去重机制。
Go、Java、Python、TypeScript 和 Rust 示例都实现并通过集成测试覆盖完整的邮箱验证、任务 1、任务 2 和完成路径。
Flow 设计
这个交互式 definition graph 由可运行的 Python Flow 生成。
该 Flow 是一个带可中断等待的顺序状态机:
Submit → VerifyEmail → AccomplishTask1 → AccomplishTask2
它组合了以下设计模式:
- 持久提醒在 Timer 触发后把等待 Step 移回自身。
- 可中断等待使用 AnyOf,让 Channel 消息可以在提醒 Timer 之前唤醒 Step。
- 响应式更新使用 RPC 校验当前阶段,并在 Flow 等待期间发布对应的 Channel 消息。
- 顺序状态机明确规定邮箱验证和两个任务的顺序,后面的任务不能绕过前面的任务。
- 持久进度状态由每个目标阶段的 WaitFor 方法记录。Status Attribute 会在持久等待安装时更新,因此 RPC 和运维人员看到的就是 Flow 当前实际等待的阶段。
- Method 原子提交让每次 WaitFor、Execute 和 RPC 调用都有独立的 commit 边界。Dex Server 只在 method 成功后,才会一起 commit 其中所有 Attribute 写入和 Channel publish,因此失败的 method 不会暴露部分 durable state 或消息。
- 幂等副作用让重试的 Step 可以安全地重新调用具备外部去重能力的邮件服务。
邮箱验证、任务 1 和任务 2 分别使用独立的 Channel。这样每个外部动作只会解除一个等待,也能直接拒绝无效状态转换。
核心实现
每个标签页首先展示 Flow signature、Step 拓扑和 persistence schema。第二段代码展示一个有代表性的 Timer 与 Channel 组合等待。链接指向完整的可运行示例,其中包含所有 RPC、HTTP controller 和端到端集成测试。
class UserOnboardingFlow(Flow[SignupForm]):
form = Attribute("Form", SignupForm)
status = Attribute("Status", str)
verify_email = Channel[None]("VerifyEmail", type(None))
task_1_completed = Channel[None]("Task1Completed", type(None))
task_2_completed = Channel[None]("Task2Completed", type(None))
def __init__(self, service: MyDependencyService) -> None:
self.service = service
self.verify_step = VerifyEmail(service, self.form, self.verify_email, self.status)
self.task_1_step = AccomplishTask1(
service, self.form, self.status, self.task_1_completed
)
self.task_2_step = AccomplishTask2(
service, self.form, self.status, self.task_2_completed
)
self.submit = Submit(service, self.form, self.status, self.verify_step)
def get_steps(self) -> StepList[SignupForm]:
return StepList.start_step(self.submit).other_steps(
self.verify_step,
self.task_1_step,
self.task_2_step,
)
def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(
self.form,
self.status,
self.verify_email,
self.task_1_completed,
self.task_2_completed,
)
class AccomplishTask1(Step[None]):
def wait_for(self, context: Context, input: None) -> Wait:
self.status.set(context, "waiting_for_task_1")
return Wait.any_of(
Timer.by_duration(timedelta(seconds=24)),
self.task_1_completed.for_one(),
)
def execute(self, context: Context, input: None) -> StepDecision:
signup_form = self.form.get(context)
if self.task_1_completed.results(context):
self.service.send_email(
signup_form.email,
"complete onboarding task 2",
"task 2 is ready",
)
return go_to(AccomplishTask2, None)
self.service.send_email(
signup_form.email,
"task 1 reminder",
"please complete onboarding task 1",
)
return go_to(AccomplishTask1, None)
class UserOnboardingFlow(Flow[SignupForm]):
@rpc
def accomplish_task_1(self, context: Context) -> RPCResult[str]:
if self.status.get(context) != "waiting_for_task_1":
return RPCResult("task 1 is not waiting")
self.task_1_completed.publish(context, None)
return RPCResult("task 1 accomplished")
例子: examples/python/dex_examples/products/signup/user_signup_flow.py
演示
下面是一次完整的 Go 本地运行:提交用户、验证邮箱、完成任务 1,再完成任务 2。图中每个业务阶段都是独立的已完成 Step;每个等待 Step 都展示了用户操作到达前处于活动状态的 Timer 和 Channel。
打开无损 PNG 可以查看完整分辨率的执行图。