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/signals-and-hooks.md.
  • 简体中文
  • v1.1.0
  • 信号与 Hook

    工作流经常要停下来等一件暂时没有结果的事。付款状态可能几分钟后才更新,人工审批也可能隔天返回。Flow 会把等待本身写进历史,当前进程可以退出,之后由任何兼容的 Worker 接着处理。

    信号和 Hook 都能恢复工作流,但它们解决的入口不同。

    入口适合处理外部寻址方式工作流读取方式
    命名信号订单状态、设备事件、业务消息运行 ID 加信号名signal_payload()
    Hook审批、Webhook、一次性回调公开令牌或运行 ID 加 Hook IDhook_payload()

    两者都采用至少一次投递。调用方必须保留稳定的消息身份,并把重复调用当成正常恢复过程。

    声明并等待信号

    信号名先写进 WorkflowSpec。这一步把运行允许接收的消息类型固定到历史中,拼错名称的投递会在写入前被拒绝。

    use a3s_flow::WorkflowSpec;
    
    const APPROVAL_SIGNAL: &str = "invoice.approved";
    
    let spec = WorkflowSpec::rust_embedded(
        "billing.invoice-approval",
        "1",
        "billing",
        "main",
    )
    .with_signal(APPROVAL_SIGNAL);

    工作流用稳定的等待 ID 创建等待。信号已经先到也没有关系,Flow 会把最早一条尚未消费、名称匹配的消息配给这次等待。

    use serde::Deserialize;
    use serde_json::json;
    
    #[derive(Deserialize)]
    struct Approval {
        reviewer: String,
    }
    
    let ctx = invocation.context();
    let Some(approval) = ctx.signal_payload_as::<Approval>("approval")? else {
        return Ok(ctx.wait_for_signal("approval", APPROVAL_SIGNAL));
    };
    
    Ok(ctx.complete(json!({
        "status": "approved",
        "reviewer": approval.reviewer,
    })))

    approval 是工作流内部的等待 ID,invoice.approved 是对外的消息契约。前者在一次运行中保持稳定,后者可以被多个连续等待复用。

    投递信号

    调用方负责生成 signal_id。应优先使用消息总线事件 ID、审批决定 ID 或业务操作 ID,不要在每次 HTTP 重试时生成新值。

    use a3s_flow::WorkflowSignal;
    use serde_json::json;
    
    let snapshot = engine
        .send_signal(
            "invoice-2026-0001",
            WorkflowSignal::new(
                "approval-decision-2026-0001",
                "invoice.approved",
                json!({ "reviewer": "finance@example.com" }),
            ),
        )
        .await?;

    相同运行 ID 与 signal_id 再次投递时,名称和载荷完全一致就会返回已经形成的结果。名称或载荷不同会得到 SignalConflict。运行经过 continue_as_new 后,查重范围会沿整条续段链展开,因此调用方仍然重用根运行 ID。

    消息先于等待到达时,消息会留在历史里。工作流创建同名等待后,Flow 按接收顺序消费最早一条。每条消息只会完成一个等待。

    创建外部 Hook

    Hook 适合外部系统只持有回调令牌的情况。令牌由宿主管理,要具备足够熵,也不能出现在日志、错误详情或普通查询参数中。

    use a3s_flow::{HookCallbackRoute, HookMetadata};
    
    let metadata = HookMetadata::human_approval("invoice:inv-0001")
        .with_callback_route(HookCallbackRoute::post(
            "/callbacks/flow/hooks/{token}",
        ))
        .with_label("tenant", "north")
        .with_data("invoiceId", "inv-0001");
    
    return Ok(ctx.create_hook_with_metadata(
        "approval",
        "public-random-token",
        metadata,
    )?);

    HookMetadata 只保存路由和审计信息。HTTP 服务、鉴权、限流、签名校验和令牌交换仍由宿主负责。回调路由模板也不会自动创建接口。

    当外部入口只有令牌时,可以恢复活动 Hook。

    let (run_id, hook_id) = engine
        .resume_hook_by_token(
            "public-random-token",
            serde_json::json!({
                "approved": true,
                "reviewer": "finance@example.com",
            }),
        )
        .await?;

    工作流重放后通过 ctx.hook_payload("approval") 读取已经提交的载荷。载荷入库成功之前,工作流看不到这次回调。

    正确处理回调重投

    按令牌查询只覆盖活动 Hook。第一次回调成功后,令牌不再出现在活动索引里。可靠回调消费者应在第一次解析令牌后保存稳定的运行 ID 与 Hook ID,后续重投调用 resume_hook()

    engine
        .resume_hook(
            &run_id,
            &hook_id,
            serde_json::json!({ "approved": true }),
        )
        .await?;

    已收到相同载荷时,这个调用可以安全重复。载荷不同、Hook 已撤回或 Hook 已被取消时会明确报冲突。不要把这些冲突转换成成功响应,它们通常表示上游复用了错误的幂等身份。

    撤回尚未完成的 Hook

    审批请求过期或外部任务关闭时,用 dispose_hook_by_token()dispose_hook() 撤回入口。工作流通过 ctx.hook_disposed("approval") 进入稳定的撤回分支。

    if ctx.hook_disposed("approval") {
        return Ok(ctx.complete(serde_json::json!({
            "status": "withdrawn",
        })));
    }

    重复撤回同一个 Hook 是幂等的。已经接收载荷的 Hook 不能再撤回,撤回后的令牌也不能接受迟到回调。

    上线前检查

    • WorkflowSpec 明确列出每个允许的信号名。
    • 信号 ID 来自业务消息身份,HTTP 重试不会生成新 ID。
    • 等待 ID、Hook ID 和令牌在重放时保持不变。
    • 回调处理器先完成签名校验和授权,再调用 Flow。
    • 可靠消费者保存运行 ID 与 Hook ID,用于确认结果不明时的重投。
    • 监控活动 Hook 数量、最早创建时间、冲突数和无法解析的令牌。

    完整程序可运行 cargo run --example workflow_signalscargo run --example hook_approvalcargo run --example hook_disposal