For AI agents: the complete documentation index is available at https://a3s-lab.github.io/Flow/llms.txt, the full documentation bundle is available at https://a3s-lab.github.io/Flow/llms-full.txt, and this page is available as Markdown at https://a3s-lab.github.io/Flow/guide/runtime-contract.md.
  • 简体中文
  • v1.1.0
  • 运行时边界

    Flow 的运行时接口故意只有两个入口。run_workflow 决定下一步,run_step 执行可能产生副作用的工作。把这条线画对,进程重启后的重放才有确定结果。

    #[async_trait]
    pub trait FlowRuntime: Send + Sync {
        async fn run_workflow(
            &self,
            invocation: WorkflowInvocation,
        ) -> a3s_flow::Result<RuntimeCommand>;
    
        async fn run_step(
            &self,
            invocation: StepInvocation,
        ) -> a3s_flow::Result<JsonValue>;
    }

    工作流代码可以做什么

    WorkflowInvocation 提供初始输入、固定的 WorkflowSpec 和完整历史。invocation.context() 在这些不可变数据上提供读取与命令构造方法。

    工作流代码适合做以下事情。

    • 检查某个步骤、等待、信号、Hook 或子工作流是否已经完成
    • 读取已经提交的 JSON 输出,并解码成业务类型
    • 根据取消请求进入稳定的清理分支
    • 返回一个步骤、批次、等待、信号等待、Hook 或子工作流命令
    • 写入进度、继续新历史段,或到达一个终态

    工作流代码不应直接做以下事情。

    • 读取当前时间并用它临时决定截止时间
    • 生成随机 ID 后再把它当作步骤身份
    • 发起网络请求、写文件、发送消息或修改数据库
    • 读取会在两次重放之间变化的环境变量
    • 依靠进程内计数器判断执行位置

    需要时间时,由主机计算一个具体 UTC 截止时间,再把它放进会被持久化的命令。需要唯一身份时,由调用方或确定性输入派生稳定值。

    每次只返回一个命令

    RuntimeCommand 是工作流决定的完整边界。常用命令包括下面这些。

    命令提交后的行为
    ScheduleStep持久化步骤定义,执行步骤,再写入输出或失败
    ScheduleSteps原子记录一个步骤批次,随后并发推进
    WaitUntil记录 UTC 截止时间并挂起
    WaitForSignal声明一个命名信号等待并挂起
    CreateHook记录外部回调身份与元数据并挂起
    StartChildWorkflow先记录子运行身份,再创建和推进子运行
    ContinueAsNew关闭当前历史段并创建继任段
    CompleteFailCancelTimeout写入唯一的终态

    一个命令提交后,引擎重新投影历史并再次调用工作流。工作流不会在一个调用栈里连续执行任意数量的决定。

    稳定 ID 是重放契约

    步骤、等待、信号等待、Hook 和子工作流都使用父运行内稳定的 ID。相同 ID 再次出现时,Flow 会比较已经持久化的定义。

    下面的改变会触发 NonDeterministic,而不会悄悄修改旧运行。

    • 同一个 step_id 换了处理器名称、输入或重试策略
    • 同一个 wait_id 换了截止时间
    • 同一个信号等待换了消息名称
    • 同一个 hook_id 换了 token 或元数据
    • 同一个子流程 ID 换了定义、输入或取消策略

    稳定 ID 最好直接写在工作流分支中,或从持久输入和确定性索引生成。

    let steps = items
        .iter()
        .enumerate()
        .map(|(index, item)| {
            StepCommand::new(
                format!("reserve-{index:04}"),
                "reserve_item",
                serde_json::to_value(item)?,
            )
        })
        .collect::<Result<Vec<_>, serde_json::Error>>()?;
    
    Ok(RuntimeCommand::schedule_steps(steps))

    不要从并发完成顺序生成下一批 ID。完成顺序会随调度改变,输入顺序不会。

    步骤是至少一次交付

    步骤输出只有在 StepCompleted 提交后才对工作流可见。一个无法消除的窗口仍然存在。

    1. 外部服务已经完成请求。
    2. worker 在步骤输出写入 Flow 历史前退出。
    3. 恢复后,同一次步骤尝试再次交付。

    Flow 不会把这个窗口描述成恰好一次。步骤实现应把稳定幂等键传给外部服务。

    async fn run_step(
        &self,
        invocation: StepInvocation,
    ) -> a3s_flow::Result<serde_json::Value> {
        let key = format!("flow:{}:{}", invocation.run_id, invocation.step_id);
        let request = invocation.input_as::<CaptureRequest>()?;
        let receipt = self.payments.capture(request, &key).await?;
        Ok(serde_json::to_value(receipt)?)
    }

    如果外部系统不支持幂等键,主机需要一个业务侧去重记录,或准备能验证结果的补偿流程。

    使用类型化输入输出

    历史的线格式是 JSON。应用边界可以继续使用 serde 类型。

    #[derive(serde::Deserialize)]
    struct OrderInput {
        order_id: String,
    }
    
    #[derive(serde::Deserialize)]
    struct Reservation {
        reservation_id: String,
    }
    
    let order = ctx.input_as::<OrderInput>()?;
    let reservation = ctx.step_output_as::<Reservation>("reserve")?;

    解码失败会返回 FlowError::Serialization。不要用默认值吞掉旧历史与新类型之间的真实不兼容。

    版本变更怎样进入工作流

    WorkflowSpec.version 是应用定义版本,runtime_build_id 是可执行构建身份,两者用途不同。

    • 输入结构或工作流语义发生有意变化时,提升定义版本。
    • 部署每个可重放构建时,给它一个明确的运行时构建 ID。
    • 只改变新运行的分支时,使用不可变补丁标记,让旧历史继续走原分支。

    生产发布还需要保留旧构建路由。详见 Worker 与发布