Sohara

Sohara 重新设计与多阶段增量规划

状态:已采纳(accepted,含 challenge 评审修订) 目标:把 tiger(事件驱动模块框架)与 rec-core(流式数据处理框架)的概念合并为一套统一的单机自动化数据/流程工作流框架,并给出从最小可验证 MVP 到完善自动化框架的增量路线。

本文是对 docs/design/draft/yaml-defiition.md 中草稿与开放问题的正式回答与落地。

评审修订结论(challenge)

  1. 单机优先,不做分布式 —— 多节点能力不在本路线图,未来以更高层抽象(外部队列/调度器)承载。
  2. AI agent 为未来扩展,不在本范围实现
  3. S0–S4 采用无 schema JSON 单路径 —— Schema/DataType/DataFrame 不进入核心,作为 S5 可选增强。 地基决策见 §2.3(数据模型)、§2.4(投递/控制流/错误/并发)、§2.5(投递语义)。

目录

  1. 现状与概念映射
  2. 重新设计
  3. 多阶段增量路线图
  4. 对草案开放问题的回答索引

1. 现状与概念映射

1.1 三个项目的定位

项目 语言 定位 核心范式
tiger TypeScript 极简事件驱动服务器(webhook/cron/queue) 有状态模块 + 协议解析器
rec-core Java 数据记录文件的校验与转换框架 函数式流式管线(Source→Tee→Target)
sohara(现状) Rust 轻量事件驱动数据处理框架(早期) Source/Sink/Transform/Pipeline + Record(JSON)

三者本质都在回答同一组问题:数据从哪来(源/触发)、如何被加工(变换)、到哪去(输出)、如何编排(流程)、如何记住状态(状态/持久化)、如何被扩展(插件)。tiger 侧重「事件驱动的流程编排」,rec-core 侧重「数据流的类型化加工」,二者互补而非互斥,可以合并。

1.2 tiger 的核心概念(已读源码归纳)

1.3 rec-core 的核心概念(已读源码归纳)

1.4 概念映射(合并结论)

tiger rec-core Sohara(统一后)
Module Source / Tee / Target Step(统一节点,按 kind 区分角色)
protocol resolver(define/notified) 源/目标工厂 ComponentRegistry + Trigger/Source/Sink trait
notify(target, param) 管线边(tee/.to() Flow 图的 Edge
Module state stateful tee / ExecutionContext StepState + ExecutionContext(checkpoint/resume)
PersistenceProvider + 任务队列 cache 回放 + .retry 文件 StateStore(状态/历史/队列/检查点)
queue 插件(进程内总线) ReactiveTee(推式链) 事件总线(有界 channel/队列)
无 schema(消息即 JSON) Schema / DataType 无 schema JSON 单路径(S0–S4);Schema 作为 S5 可选增强
serve 常驻 run 脚本(一次性) 双模式:sohara serve / sohara run

明确收敛的决策

  1. tiger 的 distributed 不在本路线图实现(单机优先);仅保留「部署维度」扩展点,未来以更高层抽象(外部队列/调度器)承载,不把多节点协议污染进单机模型。
  2. rec-core 的 DataFrame/Schema/DataType 不进入核心路径:S0–S4 只做无 schema JSON 单路径,Schema 作为 S5 可选增强(见 §2.3)。
  3. 两套 JS 脚本(Rhino / QuickJS)统一为 QuickJS(沿用 Cargo.toml 已声明的 quick-js),只保留一套脚本桥;MVP 采用同步宿主调用(学 rec 的 Rhino 全同步),异步桥接为后续可选增强(见 quickjs-api.md §7)。
  4. tiger 的「协议字符串 protocol:path」保留其命名空间思想kind:type),但不复用其字符串解析约定;Sohara 使用结构化的 kind + type 字段,便于 YAML 与校验。
  5. rec-core 的 ReactiveTee 推式语义由「运行时调度器 + 有界 channel 事件总线」统一承载,不再单独成为一种用户可见的管线种类。

2. 重新设计

2.1 统一概念模型

2.2 crate 布局(在现有 workspace 上增量演进)

现有 workspace 仅含 sohara-coreCargo.toml 已声明 tokio / axum / cron / sqlx / quick-js / crossbeam-channel / tokio-stream / chrono / uuid / dashmap 等依赖,恰好预示了下述各 crate 的分工,将在对应阶段接入。

sohara/
├── sohara-core        # 数据模型(Record[serde_json]) + Step/Source/Transform/Sink trait + Error + Registry
├── sohara-runtime     # Flow 图、调度器(顺序/并发/批量)、ExecutionContext、StateStore trait、事件总线、优雅停机
├── sohara-config      # YAML/serde schema、校验、加载、imports
├── sohara-builtins    # 内置步骤:filter/map/aggregate/merge/assert、file/log/vec 等
├── sohara-io          # 数据格式与连接器:csv/json/jsonl/parquet、db(sqlx)、http 客户端
├── sohara-triggers    # http(axum)、cron、queue/事件总线
├── sohara-js          # QuickJS 脚本桥 + script 步骤
├── sohara-cli         # `sohara init` / `run <flow.yaml>` / `serve <flow.yaml>`
├── sohara-persistence # StateStore 实现:memory/rocksdb/sqlite(单机)
└── sohara-server      # 运行历史/指标/管理 API(单机)

sohara-agent 与分布式(多节点)能力为未来扩展,不在本路线图

crate 与阶段对应

crate 引入阶段
sohara-core S0(在现有基础上收敛数据模型与 trait)
sohara-configsohara-builtinssohara-cli S1
sohara-runtime S2(图与调度)
sohara-triggers S3
sohara-persistence S4
sohara-iosohara-js S5
sohara-server S6

2.3 数据模型

地基决策(定稿):S0–S4 采用无 schema JSON 单路径;rec-core 的 Schema/DataType/DataFrame 不进入核心,作为 S5 可选增强(届时以「Schema 可选字段 + 校验层」方式引入,不替代、不破坏 JSON 路径与既有 API)。

Record(单路径 JSON)

沿用现状签名,payload 保持 JSON 值,id/timestamp/metadata 不变——从而 Record::from_json(...)record.set(field, serde_json::Value) 与 README Quick Start 全部无需重写:

pub struct Record {
    pub id: String,
    pub timestamp: DateTime<Utc>,
    pub payload: serde_json::Value,   // S0–S4 唯一路径
    pub metadata: HashMap<String, String>,
}

TransformOutcome

修复现有 transform.rs 中「用 Error::Transform 表示被过滤」的语义缺陷,并为错误显式建模:

pub enum TransformOutcome {
    Pass(Record),            // 继续流动
    Filtered,                // 被过滤,停止(计入 filtered)
    Expand(Vec<Record>),     // 一对多(split / flat-map)
    Fail(Error),             // 步骤失败(计入 errors,按 on_error 策略处理)
}

2.4 执行模型

2.5 状态与恢复(长流程)

决策:单机优先,不做分布式;以下 StateStore 均为单机实现(memory/rocksdb/sqlite)。

统一 tiger 的持久化队列语义与 rec 的 checkpoint 语义,形成一套单机 at-least-once + 幂等去重模型:

StateStore trait:
  step_state    // 步骤累加状态(tiger module state / rec stateful)
  run_history   // 运行历史(tiger monitor recordRun)
  job_queue     // claim/ack/fail/requeue(崩溃恢复)
  cron_schedule // cron 下次运行持久化
  checkpoint    // 源位置 + 步骤状态 + 未决任务

2.6 插件 / 组件注册表

ComponentRegistry:
  register(kind, type, factory)
  resolve(kind, type) -> Step

2.7 YAML schema(分阶段演进)

完整权威 schema(顶层字段、全部 step 类型与 config、edges、表达式、imports、triggers、checkpoint、校验规则、版本演进)见 yaml-workflow-schema.md

S1(线性形态,最小可用)

name: example
version: "1"
steps:
  - { id: in,    kind: source,    type: file,      format: csv,   path: data.csv, columns: [name, age] }
  - { id: adult, kind: transform, type: filter,    where: "age > 18" }
  - { id: out,   kind: sink,      type: file,      format: jsonl, path: out.jsonl }
edges: [[in, adult], [adult, out]]

省略 edges 时,仅当有且只有一个 source/trigger 才允许线性串联;否则必须显式写 edges(见 schema §10)。

S2(图 + 控制流)

steps:
  - { id: fanout, kind: control, type: parallel, branches: [a, b] }
  - { id: branch, kind: control, type: switch, cases: [{ when: "amount > 1000", to: big }], default: small }
  - { id: loop,   kind: control, type: foreach, over: "$.items", as: item }
  - { id: batch,  kind: transform, type: batch, size: 100, within: 5s }

S3(触发器)

triggers:
  - { id: webhook, type: http, method: POST, path: /webhook }
  - { id: tick,    type: cron, expression: "*/5 * * * * *" }
  - { id: bus,     type: queue, topic: hello }
steps:
  - { id: handle, kind: transform, type: map, ... }
  - { id: sink,   kind: sink, type: log }
edges: [[webhook, handle], [tick, handle], [bus, handle], [handle, sink]]

S4(状态与恢复)

checkpoint: { every: 1000 }      # 计数 checkpoint(顶层)
steps:
  - { id: count, kind: transform, type: map, state: { count: 0 },
      on_error: retry, retry: { max: 3, backoff: 1s } }   # 步骤级状态 + 错误策略
  - { id: approve, kind: control, type: approve, config: { title: "请审批", owners: [alice] } }

S5(扩展与复用)

imports: [common-steps.yaml]     # YAML 片段复用(答复「yaml 导入」)
steps:
  - { id: enrich, kind: transform, type: script, script: enrich.js }   # QuickJS
  - { id: from_db, kind: source, type: db, query: "SELECT ..." }
  - { id: to_parquet, kind: sink, type: parquet, path: out.parquet }

2.8 表达式语言


3. 多阶段增量路线图

原则:每个阶段可独立交付、可独立验证,阶段间单向依赖 S(n) 依赖 S(n-1)。验收统一以「cargo test(含集成测试)+ 一个可复现 example + 文档同步」为准。AI agent 与分布式明确不在本路线图,S6 之后为完善与扩展点预留。

S0 — 核心 MVP:JSON Record + 线性管线

S1 — 声明式 + CLI(run 模式)

S2 — 流程 DAG + 并发

S3 — 触发器 + serve 模式

S4 — 持久化 + 恢复 + 人工在环

S5 — 连接器 + 脚本

S6 — 可观测 + 管理(单机)

S7 — 完善 + 打包(扩展点预留)


4. 对草案开放问题的回答索引

yaml-defiition.md 的问题 本文回答位置
trigger / source / transform / sink 这一套 §2.1(Step.kind 统一);§2.7 YAML
sequential / concurrent / cascaded / batch 这一套 §2.4(调度语义:顺序/并发/级联/批量);§2.7 S2;路线图 S2
Rescenario 场景(capture、transform、assertion/expect) assert 归入 transform(S0 提供 AssertTransform);「场景」即一个带 assert 的 Flow
sohara serve / sohara run §2.4;S1(run)、S3(serve)
yaml 文件互相导入共享 step §2.7 S5(imports);S5 实施
执行可记录、可恢复、流程可持久化(长时流程) §2.5(StateStore + checkpoint);S4 实施
条件定义(分支/循环/复杂业务流) §2.7 S2(switch/foreach/loop);S2 实施
human-in-the-loop(approve/退回) §2.5(wait/approve);S4 实施

附:术语对照(中 / 英)

中文 英文
流程 Flow
步骤 / 节点 Step / Node
Edge
源 / 触发器 Source / Trigger
变换 Transform
输出 / 汇聚 Sink / Target
记录 Record
模式 / 表结构 Schema
状态 State
执行上下文 ExecutionContext
检查点 / 恢复 Checkpoint / Resume
组件注册表 ComponentRegistry
状态存储 StateStore
人工在环 Human-in-the-loop