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/reference/api.md.
  • 简体中文
  • v1.0.0
  • 公共 API

    这一页按使用场景整理 Flow 1.0 的公开 Rust 接口。完整签名、字段和 Rustdoc 以 docs.rs 为准。

    功能开关

    功能默认提供内容
    native-tsNativeTsRuntime 与原生编译调用
    sqliteSqliteEventStore、SQL 迁移与历史保留
    postgresPostgresEventStore、独立迁移入口与历史保留
    bootBootFlowTaskManager 与任务策略
    a3s-event提交后事件总线 Sink

    最小 Rust 主机可以关闭默认功能。

    a3s-flow = { version = "=1.0.0", default-features = false }

    引擎构建

    API用途
    FlowEngine::in_memory(runtime)用内存历史创建测试引擎
    FlowEngine::new(store, runtime)用明确存储创建引擎
    FlowEngine::builder(runtime)配置存储、观察者、版本准入和安全上限
    FlowEngineBuilder::with_store()替换事件存储
    with_observer()接收提交后的事件
    with_runtime_build_compatibility()声明当前和可重放版本
    with_max_replay_iterations()限制一次驱动的重放循环
    with_max_continue_as_new_hops()限制一次驱动跨越的续段数量
    with_max_child_workflow_depth()限制父子运行嵌套深度

    FlowEngine 可克隆,内部存储和运行时通过 Arc 共享。公开运行时、存储、队列和观察者 Trait 都要求 Send + Sync

    创建与驱动运行

    API返回说明
    start(spec, input)String生成运行 ID 并驱动到终态或暂停
    start_with_id(run_id, spec, input)String以调用方 ID 幂等创建并驱动
    drive(run_id)WorkflowRunSnapshot从当前历史继续运行
    continuation_chain(run_id)Vec<WorkflowRunSnapshot>从任一段读取整条续段链

    start_with_id() 的重投只有在定义与输入完全一致时成功。定义名称、版本、运行时、入口、运行版本、补丁标记、允许信号或输入变化都会返回 RunConflict

    检查接口

    API用途
    snapshot(run_id)投影单条运行
    history(run_id)读取完整事件包络
    list_run_ids()稳定顺序列出运行 ID
    list_snapshots()投影全部运行
    run_summary()汇总状态数量
    list_open_suspensions(now)列出等待、重试和活动 Hook
    next_wakeup(now)找到下一条定时暂停
    list_active_hooks()列出活动回调入口

    大规模生产查询应依赖 SQL 索引投影,并在控制面增加分页。list_snapshots() 会读取很多历史,不适合无限增长的数据集热路径。

    外部输入与控制

    API语义
    send_signal(run_id, WorkflowSignal)幂等提交命名消息并驱动运行
    resume_hook(run_id, hook_id, payload)以稳定身份恢复 Hook,支持重投
    resume_hook_by_token(token, payload)解析活动令牌并恢复一次
    dispose_hook(run_id, hook_id)以稳定身份撤回 Hook
    dispose_hook_by_token(token)按活动令牌撤回 Hook
    request_cancellation(run_id, request)请求可恢复清理并重放
    force_cancel(run_id, reason)不运行清理,立即形成取消终态
    terminate_for_timeout()写入带截止时间的超时终态
    terminate_for_host_shutdown()写入明确不可恢复的主机关闭终态
    record_progress()幂等保存宿主进度
    link_child_operation()幂等保存外部子操作关联

    按令牌查询只覆盖活动 Hook。可靠重投要保存第一次解析得到的运行 ID 与 Hook ID。

    FlowRuntime

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

    WorkflowInvocation 提供运行 ID、工作流定义、初始输入和完整历史。context() 返回 WorkflowContext,集中提供投影查询与命令构造。

    StepInvocation 提供运行 ID、步骤 ID、步骤名、输入和历史。外部副作用只放在 run_step()

    WorkflowContext 查询

    方法读取内容
    run_id()input()input_as()spec()运行权威与初始输入
    history()原始事件包络
    step_output()step_output_as()成功步骤输出
    step_completed()step_failed()步骤终态
    wait_completed()定时等待结果
    signal_payload()signal_payload_as()已匹配信号载荷
    hook_payload()hook_payload_as()hook_disposed()Hook 结果
    cancellation_request()清理式取消请求
    child_workflow_run_id()child_workflow_outcome()第一类子运行
    progress()child_operation()控制面进度与外部关联
    has_patch_marker()不可变代码分支选择

    类型化读取通过 Serde 反序列化,失败返回 FlowError,不会产生 panic。

    WorkflowContext 命令

    方法形成的决定
    complete(output)成功终态
    fail(error)失败终态
    cancel()完成已经请求的取消
    timeout(deadline, reason)超时终态
    schedule_step()schedule_step_with_retry()单个持久步骤
    step()step_with_retry()构造批次步骤
    schedule_steps()原子声明步骤批次
    wait_until()持久 UTC 定时等待
    wait_for_signal()命名信号等待
    create_hook()create_hook_with_metadata()外部回调入口
    start_child_workflow()单个第一类子运行
    child_workflow()start_child_workflows()构造并启动子运行批次
    continue_as_new()关闭当前历史并创建后继段
    record_progress()保存工作流进度
    link_child_operation()保存外部子操作关联

    一个 run_workflow() 调用只返回一个 RuntimeCommand

    定义和策略类型

    类型角色
    WorkflowSpec工作流名称、版本、运行时入口、运行版本、补丁与信号声明
    RuntimeSpecRuntimeKindRust 嵌入式或原生 TypeScript 入口
    RuntimeBuildId具体可执行代码身份
    RuntimeBuildCompatibilityWorker 可接受版本集合
    WorkflowPatchId有界、不可变、回放安全的分支标记
    RetryPolicyRetryBackoff尝试次数、延迟和退避
    StepFailureAction重试耗尽后终止运行或返回工作流
    StepCommand批次中的步骤定义
    ChildWorkflowCommand批次中的子运行定义
    ChildWorkflowCancellationPolicy父取消时请求取消或放弃关联
    CancellationRequest持久停止请求与理由
    WorkflowSignal调用方拥有 ID 的命名消息
    HookMetadataHookCallbackRouteHook 审计和宿主路由元数据

    快照与终态

    主要只读投影包括 WorkflowRunSnapshotStepSnapshotWaitSnapshotHookSnapshotSignalWaitSnapshotChildWorkflowSnapshotScheduledWakeup

    WorkflowRunStatus 表达当前生命周期,WorkflowTerminalOutcome 表达带字段的最终结果。公开枚举为 #[non_exhaustive],匹配时保留兜底分支。

    存储

    类型功能要求
    InMemoryEventStore
    LocalFileEventStore
    SqliteEventStoresqlite
    PostgresEventStorepostgres
    FlowEventStore自定义存储 Trait
    FlowHistoryRetentionPolicySQL 历史保留条件
    FlowHistoryHold持久审计保留标记
    FlowHistoryTombstone删除后的最小校验记录

    PostgreSQL 生产迁移入口是 migrate_postgres_flow(),服务进程使用验证构造器。

    调度与任务

    类型角色
    FlowScheduler扫描到期等待和重试,派发按运行任务
    FlowSchedulerTick本轮到期项与派发数量
    FlowTask驱动、等待、Hook、信号和定时恢复载荷
    FlowTaskDispatcher任务派发 Trait
    FlowTaskQueue租约、确认和心跳 Trait
    FlowWorker处理嵌入式队列租约
    RuntimeBuildTaskRouter按持久运行版本精确路由
    BootFlowTaskManagerBoot 处理器与派发器
    BootFlowTaskPolicy任务重试、超时、清理和去重策略

    FlowWorker 适合嵌入式队列。生产宿主优先使用任务管理组件统一处理生命周期。

    观察

    FlowEventObserver 在事件提交后接收包络。内置实现包括 NoopFlowEventObserverInMemoryFlowEventObserverFanoutFlowEventObserver 和本地 JSONL Sink。

    A3sFlowEvent 提供低基数的事件投影,safe_metric_labels() 只返回适合指标标签的字段。观察失败不能反向改变已经提交的工作流历史。

    工作流图

    类型用途
    WorkflowDsl完整 YAML 或 JSON 文档、版本与扩展字段
    WorkflowDag节点、边、计划和语义摘要
    WorkflowDagNodeWorkflowDagEdge程序化构图
    WorkflowDagPlan顶层与容器作用域的确定性顺序
    WorkflowDslCompatibility当前、旧版警告或需确认的版本分类
    WorkflowDslError文档大小、解析、结构与版本错误

    错误处理

    crate 统一返回 a3s_flow::Result<T>,错误类型是 FlowError。生产代码通常需要单独处理以下类别。

    错误操作含义
    RunNotFound返回未找到或检查路由目标
    RunConflict幂等身份被不同权威复用
    RunTerminal外部输入到达已结束运行
    EventConflict并发写入获胜者已产生,可重新读取
    NonDeterministic代码决定与历史不一致,停止自动重试
    SignalConflictHookConflict外部幂等身份载荷漂移
    RuntimeBuildUnavailable当前 Worker 不具备重放权威
    RuntimeBuildRouteNotFound调度器缺少目标版本派发器
    ReplayLimitExceeded重放或恢复循环超过安全上限
    StoreRuntime后端或运行时错误,按上下文分类重试

    不要把所有错误无条件重试。非确定性、身份冲突和版本准入错误需要修复代码、路由或调用方身份。