For AI agents: the complete documentation index is available at https://a3s-lab.github.io/Flow/v1.0.0/llms.txt, the full documentation bundle is available at https://a3s-lab.github.io/Flow/v1.0.0/llms-full.txt, and this page is available as Markdown at https://a3s-lab.github.io/Flow/v1.0.0/concepts/durable-primitives.md.
  • 简体中文
  • v1.0.0
  • 持久原语

    Flow 的持久原语都遵循同一条规则。工作流先返回一个带稳定身份的命令,引擎提交事件,随后重放代码并从历史读取结果。运行时代码不需要维护隐藏游标。

    身份先于执行

    原语稳定身份参数漂移时的处理
    步骤step_id名称、输入或重试策略变化时拒绝重放
    定时等待wait_idUTC 截止时间变化时拒绝重放
    信号等待wait_id信号名变化时拒绝重放
    Hookhook_id令牌或元数据变化时拒绝重放
    信号投递signal_id名称或载荷变化时返回冲突
    子工作流child_id定义、输入或取消策略变化时拒绝重放
    进度progress_id数量、消息或详情变化时返回冲突
    子操作引用reference_id操作身份或元数据变化时返回冲突

    ID 应描述业务阶段,不要使用数组下标之外的临时位置、随机值或当前时间。对集合进行批处理时,先稳定排序,再生成带固定宽度序号的 ID。

    步骤

    步骤是外部副作用边界。run_workflow() 决定是否安排步骤,run_step() 执行网络、数据库、文件或工具调用。

    if let Some(output) = ctx.step_output("charge") {
        return Ok(ctx.complete(output.clone()));
    }
    
    Ok(ctx.schedule_step(
        "charge",
        "chargeCard",
        serde_json::json!({
            "paymentId": "pay-8842",
            "idempotencyKey": format!("{}:charge", ctx.run_id()),
        }),
    ))

    步骤状态和输出保存在快照中。成功事件提交后,重放只读取结果。外部调用完成但提交前崩溃时,步骤可能再次执行。

    原子批次

    schedule_steps() 一次声明多个互相独立的步骤。引擎先校验整批 ID 唯一且参数稳定,再推进成员。结果顺序采用请求顺序,不采用完成时间。

    let steps = items
        .into_iter()
        .enumerate()
        .map(|(index, item)| {
            ctx.step(
                format!("item-{index:04}"),
                "processItem",
                serde_json::json!({ "item": item }),
            )
        })
        .collect();
    
    Ok(ctx.schedule_steps(steps))

    批次不能表达步骤之间的依赖。后一阶段需要前一阶段输出时,等第一批全部进入历史,再由下一次重放安排第二批。

    重试

    重试策略是步骤命令的一部分。支持不重试、固定延迟和带上限的指数退避。

    use a3s_flow::RetryPolicy;
    use std::time::Duration;
    
    let retry = RetryPolicy::exponential(
        8,
        Duration::from_secs(1),
        Duration::from_secs(30),
    );
    
    Ok(ctx.schedule_step_with_retry(
        "reserve",
        "reserveInventory",
        input,
        retry,
    ))

    指数策略使用不可变的运行、步骤和尝试身份推导确定性抖动。重启不会改变已选择的下一次 UTC 时间。默认在耗尽后终止运行,continue_workflow_on_failure() 允许工作流读取 step_failed() 并进入补偿或降级分支。

    定时等待

    wait_until() 保存绝对 UTC 时间,不保存进程内计时器。

    if ctx.wait_completed("payment-window") {
        return Ok(ctx.complete(serde_json::json!({ "ready": true })));
    }
    
    Ok(ctx.wait_until("payment-window", resume_at))

    到期只表示运行可以恢复。调度器可能重复派发任务,恢复操作会检查等待状态并安全收敛。

    信号与 Hook

    信号是运行 ID 寻址的命名消息,Hook 是令牌寻址的外部回调。两者都会暂停运行并在载荷提交后重放。详细的声明、查重和撤回规则见信号与 Hook

    进度

    进度适合控制面查询,不能用来驱动工作流分支。工作流可以返回 ctx.record_progress(),宿主也可以调用 engine.record_progress()

    use a3s_flow::WorkflowProgress;
    
    let progress = WorkflowProgress::new("import-page-0042", 42)
        .with_total(100)
        .with_message("page committed");
    
    engine.record_progress(&run_id, progress).await?;

    相同 progress_id 和相同内容可以重投。更新百分比时要使用新的身份,例如页码或单调递增的业务序号。

    子操作引用

    外部系统自己管理的任务不需要伪装成第一类子工作流。用 ChildOperationReference 保存稳定关联,再通过步骤、信号或 Hook 管理它的生命周期。

    use a3s_flow::ChildOperationReference;
    
    let child = ChildOperationReference::new(
        "video-render",
        "render-job",
        "render-8821",
    )
    .with_metadata(serde_json::json!({ "region": "cn-east" }));
    
    Ok(ctx.link_child_operation(child))

    这种引用只负责持久关联,不会自动传播取消。由 Flow 管理生命周期的子运行应使用第一类子工作流

    续段

    continue_as_new() 用新输入创建后继运行并关闭当前历史段。适合轮询、批量导入和长期周期任务。

    if cursor < total {
        return Ok(ctx.continue_as_new(serde_json::json!({
            "cursor": cursor + 1,
            "total": total,
        })));
    }

    每个段有独立事件流,后继段继承完整工作流定义。continuation_chain() 可以从根运行查询整条链。续段用于控制重放长度,不负责修改旧历史。

    组合原则

    一条常见生产路径可以按下面的顺序组合。

    1. 用步骤创建外部请求。
    2. 用 Hook 或信号等待结果。
    3. 用定时等待实现业务截止时间。
    4. 用清理步骤处理取消。
    5. 用进度和观察者提供运维视图。
    6. 历史达到预定长度后续段。

    每一步都有独立稳定身份,恢复点就会清楚。多个职责塞进一个步骤会让重试、审计和补偿边界一起变模糊。