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/operations/persistence.md.
  • 简体中文
  • v1.1.0
  • 存储与迁移

    事件存储是 Flow 的状态权威。步骤输出、等待截止时间、回调载荷、取消请求和终态都要先进入事件历史,运行才能继续。队列、指标和观察日志可以重建,事件历史不能靠它们恢复。

    选择存储

    存储适用范围多进程写入生产建议
    InMemoryEventStore单元测试、临时演示进程退出即丢失
    LocalFileEventStore单进程嵌入式、开发环境备份整个 JSONL 根目录
    SqliteEventStore单节点持久服务不建议多个所有者一个写入服务拥有数据库文件
    PostgresEventStore多进程 Worker、共享控制面迁移角色与服务角色分离

    默认功能只启用原生 TypeScript 适配器。SQL 存储需要明确打开功能。

    [dependencies]
    a3s-flow = { version = "=1.1.0", features = ["sqlite"] }

    PostgreSQL 改用 features = ["postgres"]。只用 Rust 且不需要原生 TypeScript 时,可以同时设置 default-features = false

    内存和本地 JSONL

    use a3s_flow::{FlowEngine, LocalFileEventStore};
    use std::sync::Arc;
    
    let store = Arc::new(LocalFileEventStore::new(".a3s/flow/history"));
    let engine = FlowEngine::new(store, runtime);

    本地存储按运行保存追加式 JSONL。写入会检查尾部记录,完整但缺少换行的最后一条记录可以恢复,损坏的中间记录会失败关闭。它不提供跨进程协调,不要让两个服务实例同时拥有同一目录。

    备份时要一起保存历史目录和宿主用于外部幂等的数据库恢复点。只复制一半会让工作流历史与业务副作用失去对应关系。

    SQLite

    use a3s_flow::{FlowEngine, SqliteEventStore};
    use std::sync::Arc;
    
    let store = Arc::new(
        SqliteEventStore::connect("sqlite://.a3s/flow/flow.db").await?,
    );
    let engine = FlowEngine::new(store, runtime);

    connect() 打开数据库并运行规范化、带校验和的迁移。迁移与事件追加都在事务中执行。生产环境只保留一个 SQLite 所有者,启动新版本前先停旧进程,并验证文件级备份可以恢复。

    SQLite 维护活动 Hook 和定时唤醒投影索引。索引用于查询加速,原始事件仍是权威。不要手工更新投影表或触发器。

    PostgreSQL

    开发环境可以直接使用 PostgresEventStore::connect(),它会迁移并打开存储。生产环境建议把 DDL 权限与服务权限分开。

    use a3s_flow::{migrate_postgres_flow, PostgresEventStore};
    
    // Run once with a dedicated migration executor.
    let report = migrate_postgres_flow(&migration_executor).await?;
    println!("applied={}", report.applied.len());
    
    // Serving workers verify checksums without holding DDL authority.
    let store = PostgresEventStore::from_executor_verified(
        serving_executor,
    )
    .await?;

    迁移通过数据库咨询锁串行执行。connect_verified()from_executor_verified() 只验证所需迁移及校验和,缺少迁移时拒绝服务。部署顺序应固定为迁移任务成功、服务实例通过验证、再恢复新运行与调度。

    数据库账户至少要分成两类。

    • 迁移账户可以创建表、索引和触发器,并写迁移账本。
    • 服务账户可以读写 Flow 表,不能修改结构和迁移记录。

    自定义 FlowEventStore

    自定义后端必须实现四个核心方法。

    #[async_trait::async_trait]
    trait FlowEventStore {
        async fn append(&self, run_id: &str, event: FlowEvent)
            -> Result<FlowEventEnvelope>;
        async fn append_if_sequence(
            &self,
            run_id: &str,
            expected_sequence: u64,
            event: FlowEvent,
        ) -> Result<FlowEventEnvelope>;
        async fn list(&self, run_id: &str)
            -> Result<Vec<FlowEventEnvelope>>;
        async fn list_run_ids(&self) -> Result<Vec<String>>;
    }

    append_if_sequence() 必须以原子方式检查最新序号并追加,竞争失败返回 FlowError::EventConflictlist() 必须提供从 1 开始、无缺口、按序排列的完整事件。

    活动 Hook 和定时唤醒方法有基于全量重放的默认实现。数据量大时应实现索引投影,并用同一事务随事件写入更新。投影错误不能改变原始事件。

    前向迁移

    Flow 迁移只向前执行并固定校验和。已经在任何环境应用的迁移文件不能修改。修复结构时新增迁移。

    一次生产迁移要保留以下证据。

    1. 升级前数据库恢复点与恢复演练结果。
    2. 候选二进制版本和 Cargo 锁定依赖。
    3. 迁移前后的 a3s_orm_migrations 内容。
    4. 代表性运行的历史长度、最后序号和快照状态。
    5. 一个真实等待或中断步骤的恢复结果。

    迁移提交后不能只回滚二进制。旧二进制不认识新迁移 ID,会拒绝打开数据库。需要回滚时,先停止全部写入者,再恢复升级前数据库与匹配的持久文件。

    历史保留

    SQLite 和 PostgreSQL 支持完整历史删除。Flow 不压缩事件流的一部分,也不改写旧事件。

    use a3s_flow::FlowHistoryRetentionPolicy;
    use chrono::{Duration, Utc};
    
    let policy = FlowHistoryRetentionPolicy::new(
        Utc::now() - Duration::days(90),
    );
    let report = store.prune_terminal_history(policy).await?;
    
    println!("deleted={:?}", report.deleted_run_ids);

    运行只有同时满足以下条件才会删除。

    • 已形成终态且终态时间早于截止点。
    • 没有持久审计保留标记。
    • 所属续段链和父子组件中的所有历史都在本次扫描中符合条件。
    • 没有缺失的链接目标。

    删除前会写入墓碑,记录终态序号、事件 ID、事件键和整条历史的 SHA-256。墓碑阻止相同运行 ID 被重新创建。

    审计保留标记

    store
        .hold_history(
            "invoice-2026-0001",
            "legal-case-8821",
            "payment dispute",
        )
        .await?;
    
    let holds = store.history_holds("invoice-2026-0001").await?;

    相同运行 ID、标记 ID 和理由可以重投。用相同标记 ID 改理由会返回冲突。解除标记使用 release_history_hold(),之后还要等下一次保留扫描。

    日常校验

    • 定期恢复备份,并驱动一条中断运行到终态。
    • 监控事件追加冲突率、查询耗时、调度索引延迟和活动 Hook 数量。
    • 对迁移账本、Flow 表和投影结构做漂移检查。
    • 保留删除报告与墓碑,核对审计标记审批流程。
    • 长期循环优先使用 continue_as_new() 控制重放长度,再按完整组件执行保留。