返回博客

源码精读:ActionEngine — 10 种执行器的调度器模式

ActionEngineService 是 coomia-dip Agent Runtime 层(Agent Runtime Layer)的核心调度引擎,负责将决策结果转化为可执行操作。它通过 gRPC 协议暴露 ExecuteAction、BatchExecuteActions、CancelAction 等 12 个 RPC 方法,支持规则驱动(RuleBasedExecution)和函数驱动(FunctionExecution)两种执行模式,外加通知、Webhook、本体编辑、函数调用四种副作用(SideEffect)。本文将深入剖析其 Proto 定义中的 oneof logic 分发策略、ExecutionState 八态状态机、批量执行的 fail-fast/parallel 控制、回滚策略选择器、以及 ActionTemplate 的实例化与参数约束验证流程。

Coomia发布于 2025年12月9日13 分钟阅读
分享本文Twitter / X

源码精读:ActionEngine — 10 种执行器的调度器模式

系列:S9 源码精读 · 第 13 篇 | 难度:高级 | 阅读时间:25 分钟

#TL;DR

ActionEngineService 是 coomia-dip Agent Runtime 层(Agent Runtime Layer)的核心调度引擎,负责将决策结果转化为可执行操作。它通过 gRPC 协议暴露 ExecuteActionBatchExecuteActionsCancelAction 等 12 个 RPC 方法,支持规则驱动(RuleBasedExecution)和函数驱动(FunctionExecution)两种执行模式,外加通知、Webhook、本体编辑、函数调用四种副作用(SideEffect)。本文将深入剖析其 Proto 定义中的 oneof logic 分发策略、ExecutionState 八态状态机、批量执行的 fail-fast/parallel 控制、回滚策略选择器、以及 ActionTemplate 的实例化与参数约束验证流程。

#目录

  1. 整体架构:从决策到行动的桥梁
  2. Proto 合约:ActionDefinition 的 oneof logic 设计
  3. ExecutionState 八态状态机
  4. ExecuteAction:单条执行的完整流程
  5. BatchExecuteActions:批量调度与并行控制
  6. RuleBasedExecution:声明式本体编辑
  7. FunctionExecution:函数驱动的行动
  8. SideEffect 四种副作用的触发机制
  9. ActionTemplate:模板实例化与参数约束
  10. 回滚策略:RollbackConfig 的四种模式
  11. Key Takeaways

#1. 整体架构:从决策到行动的桥梁

ActionEngine 在 coomia-dip 的八层架构中属于 Agent Runtime Layer(Agent Runtime),其核心职责是:接收来自 Reasoning & Decision Layer 决策引擎的决策结果,将其转化为具体的系统操作。代码结构如下:

Code
intelligence-Layer/src/agent_runtime_plane/
├── action/
│   ├── audit_emitter.py          # 审计事件发射器
│   └── ...
├── api/generated/plane_e/
│   ├── action_engine_pb2.py      # Proto 生成代码
│   └── action_engine_pb2_grpc.py # gRPC Stub
proto/plane_e/
└── action_engine.proto           # 867 行核心合约
python-sdk/ontology_sdk/
└── grpc_client/grpc_action_client.py  # SDK 客户端封装

设计哲学:ActionEngine 不直接实现具体操作逻辑,而是作为调度器(Dispatcher)根据 ActionDefinition.logiconeof 类型,将执行分发到不同的执行器。这种策略模式使新增执行器类型时只需扩展 proto 定义和对应的执行器实现,零修改核心调度逻辑。

#2. Proto 合约:ActionDefinition 的 oneof logic 设计

ActionEngine 最核心的设计决策体现在 ActionDefinition 消息的 oneof logic 字段上:

PROTOBUF
message ActionDefinition {
  string action_id = 1;
  string action_name = 2;

  // 两种执行模式二选一
  oneof logic {
    RuleBasedExecution rule_based = 3;
    FunctionExecution function = 4;
  }

  // 副作用
  repeated SideEffectExecution side_effects = 5;

  RetryPolicy retry_policy = 6;
  TimeoutConfig timeout = 7;
  RollbackConfig rollback_config = 8;
  string description = 9;
  repeated string tags = 10;
  bool enabled = 11;

  // 向后兼容
  ActionType legacy_type = 50;
  google.protobuf.Struct legacy_config = 51;
}

这里有几个关键的设计决策值得注意:

oneof 而非 enum + config:早期版本使用 ActionType 枚举加通用 Struct config 的方式。v1.5 重构为 oneof logic,每种执行模式有自己的强类型消息。这带来了编译时类型安全——你无法将 Webhook 配置传给函数执行器。

向后兼容保留:注意 legacy_typelegacy_config 字段编号从 50 开始,与正常字段保持距离,避免未来扩展时的编号冲突。枚举 ActionType 虽然标记为 DEPRECATED,但保留其定义确保旧客户端不会反序列化失败。

副作用与主逻辑分离side_effectsrepeated 而非 oneof 的一部分,意味着一个 Action 可以有零到多个副作用,且副作用与主逻辑独立运行。

#3. ExecutionState 八态状态机

ActionEngine 定义了八种执行状态,构成了清晰的状态机:

PROTOBUF
enum ExecutionState {
  EXECUTION_STATE_UNSPECIFIED = 0;
  EXECUTION_STATE_PENDING = 1;
  EXECUTION_STATE_RUNNING = 2;
  EXECUTION_STATE_COMPLETED = 3;
  EXECUTION_STATE_FAILED = 4;
  EXECUTION_STATE_CANCELLED = 5;
  EXECUTION_STATE_ROLLBACK_IN_PROGRESS = 6;
  EXECUTION_STATE_ROLLBACK_COMPLETED = 7;
  EXECUTION_STATE_WAITING_APPROVAL = 8;
}

合法的状态转换路径:

Code
PENDING → RUNNING → COMPLETED
                  → FAILED → ROLLBACK_IN_PROGRESS → ROLLBACK_COMPLETED
PENDING → WAITING_APPROVAL → RUNNING
RUNNING → CANCELLED
ROLLBACK_IN_PROGRESS → FAILED (回滚失败)

WAITING_APPROVAL 状态是 Human-in-the-Loop 模式的核心——当 Action 需要人工审批时,执行引擎会暂停在此状态,等待 Reasoning & Decision Layer 的 ApprovalWorkflowService 发出审批通过信号后才继续执行。

SDK 客户端中对状态的封装同样值得关注:

Python
class ExecutionState(str, Enum):
    UNSPECIFIED = "UNSPECIFIED"
    PENDING = "PENDING"
    RUNNING = "RUNNING"
    COMPLETED = "COMPLETED"
    FAILED = "FAILED"
    CANCELLED = "CANCELLED"
    ROLLBACK_IN_PROGRESS = "ROLLBACK_IN_PROGRESS"
    ROLLBACK_COMPLETED = "ROLLBACK_COMPLETED"
    WAITING_APPROVAL = "WAITING_APPROVAL"

使用 str, Enum 双继承使得状态值既可以作为枚举比较,也可以直接序列化为 JSON 字符串。

#4. ExecuteAction:单条执行的完整流程

单条 Action 执行的请求包含五个核心字段:

PROTOBUF
message ExecuteActionRequest {
  com.onto.common.v1.RequestContext context = 1;
  string decision_id = 2;          // 关联决策 ID
  ActionDefinition action = 3;
  google.protobuf.Struct input = 4;
  ExecuteOptions options = 5;
  string template_id = 6;          // 从模板加载
}

ExecuteOptions 中的 dry_run 字段允许在不实际执行的情况下验证 Action 的可行性:

PROTOBUF
message ExecuteOptions {
  bool async_execution = 1;
  bool dry_run = 2;
  string idempotency_key = 3;
  int32 priority = 4;              // 0-100
}

幂等性设计idempotency_key 允许客户端在网络重试时避免重复执行。执行引擎会检查此 key 是否已有对应的执行记录,如果有则直接返回已有结果。

优先级队列priority 字段(0-100)在异步执行模式下生效,高优先级的 Action 会被优先调度。

响应中的 ActionStatus 包含了丰富的执行信息:

PROTOBUF
message ActionStatus {
  string execution_id = 1;
  ExecutionState state = 2;
  google.protobuf.Timestamp started_at = 3;
  google.protobuf.Timestamp completed_at = 4;
  int32 retry_count = 5;
  string current_step = 6;
  double progress = 7;             // 0.0 to 1.0
  string error_code = 8;
  string error_message = 9;
}

progress 字段(0.0 到 1.0)为前端 UI 提供实时进度反馈,特别是在批量本体编辑操作中,每完成一条规则执行进度就会推进。

#5. BatchExecuteActions:批量调度与并行控制

批量执行是 ActionEngine 的高频使用场景,其控制选项设计得很精细:

PROTOBUF
message BatchExecuteOptions {
  bool fail_fast = 1;              // 第一个失败即停止
  int32 max_parallel = 2;          // 最大并行数(0 = 串行)
  bool async_execution = 3;
}

fail_fast 语义:当 fail_fast = true 时,第一个 Action 执行失败会立即取消所有未开始的 Action。已经在运行中的 Action 不会被强制取消(需要通过 CancelAction 显式取消),但不会再启动新的 Action。

max_parallel 控制max_parallel = 0 表示严格串行执行,保证操作的顺序性。正整数值限制了并发执行的最大数量,防止在大批量操作时压垮下游服务。

批量执行的每个 item 都有独立的 item_id 用于关联:

PROTOBUF
message ActionExecutionItem {
  string item_id = 1;              // 客户端提供的关联 ID
  ActionDefinition action = 2;
  google.protobuf.Struct input = 3;
  ExecuteOptions options = 4;
}

响应中包含完整的汇总信息:

PROTOBUF
message BatchExecutionSummary {
  int32 total = 1;
  int32 succeeded = 2;
  int32 failed = 3;
  int32 cancelled = 4;
  double total_duration_ms = 5;
}

#6. RuleBasedExecution:声明式本体编辑

RuleBasedExecution 对标 Palantir Foundry 的 Rules-based Action,允许用户通过声明式规则描述本体编辑操作:

PROTOBUF
message RuleBasedExecution {
  repeated OntologyEditExecution rules = 1;
}

message OntologyEditExecution {
  string rule_id = 1;
  string operation_type = 2;       // create_object, modify_object, delete_object,
                                   // create_link, delete_link, create_or_modify
  google.protobuf.Struct parameters = 3;
  string condition = 4;            // 执行条件(可选)
  int32 order = 5;                 // 执行顺序
}

六种操作类型涵盖了本体实例的完整 CRUD:

  1. create_object — 创建新的本体实例
  2. modify_object — 修改已有实例的属性
  3. delete_object — 删除实例(软删除)
  4. create_link — 建立实例间的关系
  5. delete_link — 删除实例间的关系
  6. create_or_modify — 幂等的 upsert 操作

条件执行condition 字段允许规则根据运行时上下文条件性地执行。例如 "input.quantity > 0" 只在数量大于零时执行创建操作。

有序执行order 字段确保多条规则按指定顺序执行。这对于先创建主对象、再创建关系这类有依赖的操作序列至关重要。

#7. FunctionExecution:函数驱动的行动

FunctionExecution 将 Action 的执行逻辑委托给 Reasoning & Decision Layer 的 FunctionRuntime:

PROTOBUF
message FunctionExecution {
  string function_id = 1;          // 引用 Reasoning & Decision Layer FunctionRuntime 中注册的函数
  string function_version = 2;     // 可选
  google.protobuf.Struct parameters = 3;
}

这种设计实现了 Action 定义与执行逻辑的完全解耦。Action 只关心"调用哪个函数、传什么参数",具体的函数实现由 Reasoning & Decision Layer 管理。function_version 字段允许锁定特定版本的函数,防止函数升级导致行为变化。

#8. SideEffect 四种副作用的触发机制

SideEffect 是 v1.5 引入的重要概念,将原来混在 ActionType 中的通知和 Webhook 提取为独立的副作用系统:

PROTOBUF
enum SideEffectType {
  SIDE_EFFECT_TYPE_UNSPECIFIED = 0;
  SIDE_EFFECT_TYPE_NOTIFICATION = 1;
  SIDE_EFFECT_TYPE_WEBHOOK = 2;
  SIDE_EFFECT_TYPE_ONTOLOGY_EDIT = 3;
  SIDE_EFFECT_TYPE_FUNCTION_CALL = 4;
}

enum SideEffectTrigger {
  SIDE_EFFECT_TRIGGER_UNSPECIFIED = 0;
  SIDE_EFFECT_TRIGGER_ON_SUCCESS = 1;
  SIDE_EFFECT_TRIGGER_ON_FAILURE = 2;
  SIDE_EFFECT_TRIGGER_ALWAYS = 3;
}

触发时机控制是 SideEffect 设计的精髓:

  • ON_SUCCESS:主逻辑成功后触发,典型场景是发送通知邮件
  • ON_FAILURE:主逻辑失败后触发,典型场景是告警 Webhook
  • ALWAYS:无论成败都触发,典型场景是审计日志

Webhook 的执行模式(FEAT-009)进一步区分了同步和异步语义:

PROTOBUF
enum WebhookExecutionMode {
  WEBHOOK_EXECUTION_MODE_UNSPECIFIED = 0;
  WEBHOOK_EXECUTION_MODE_WRITEBACK = 1;    // Pre-commit: 同步调用,失败回滚
  WEBHOOK_EXECUTION_MODE_SIDE_EFFECT = 2;  // Post-commit: 异步调用,失败仅记录
}

WRITEBACK 模式下 Webhook 的返回值可以通过 WebhookOutputMapping 写回到操作上下文中供后续规则使用,实现了外部系统的数据回写能力。

#9. ActionTemplate:模板实例化与参数约束

ActionTemplate 是 ActionEngine 的高级特性,允许将常用的 Action 配置保存为可复用的模板:

PROTOBUF
message ActionTemplate {
  string template_id = 1;
  string name = 2;
  string display_name = 3;

  oneof logic {
    RuleBasedExecution rule_based_logic = 4;
    FunctionExecution function_logic = 5;
  }

  TemplateStatus status = 7;
  string version = 8;
  repeated TemplateParameter parameters = 10;

  // 提交条件与副作用
  SubmissionCriteria submission_criteria = 23;
  repeated SideEffect side_effects = 24;
}

参数约束系统(Sprint 4)支持五种约束类型:

PROTOBUF
enum ConstraintType {
  CONSTRAINT_TYPE_RANGE = 1;       // 数值范围
  CONSTRAINT_TYPE_LENGTH = 2;      // 字符串长度
  CONSTRAINT_TYPE_PATTERN = 3;     // 正则匹配
  CONSTRAINT_TYPE_ENUM = 4;        // 枚举值
  CONSTRAINT_TYPE_CUSTOM = 5;      // 自定义验证函数
}

SubmissionCriteria 将约束和规则组合在一起:

PROTOBUF
message SubmissionCriteria {
  string validation_function_id = 1;         // 外部验证函数
  repeated ParameterConstraint constraints = 2;
  repeated SubmissionRule rules = 3;
}

SubmissionRule 支持跨参数的布尔表达式验证,并区分 ERROR 和 WARNING 两个严重级别——WARNING 级别的验证失败会告警但不阻止提交。

模板的生命周期通过 TemplateStatus 枚举管理:

Code
DRAFT → ENABLED ↔ DISABLED → DEPRECATED → ARCHIVED

#10. 回滚策略:RollbackConfig 的四种模式

当 Action 执行失败时,RollbackConfig 决定如何恢复:

PROTOBUF
enum RollbackStrategy {
  ROLLBACK_STRATEGY_NONE = 1;              // 不回滚
  ROLLBACK_STRATEGY_COMPENSATE = 2;        // 执行补偿 Action
  ROLLBACK_STRATEGY_RESTORE_SNAPSHOT = 3;  // 恢复快照
  ROLLBACK_STRATEGY_MANUAL = 4;            // 人工介入
}

message RollbackConfig {
  bool enabled = 1;
  RollbackStrategy strategy = 2;
  bool auto_rollback_on_failure = 3;
  int32 rollback_timeout_ms = 4;
  string compensating_action_id = 5;
}

COMPENSATE 模式是最常用的策略——通过 compensating_action_id 指向一个专门的补偿 Action,实现 Saga 模式的最终一致性。

RESTORE_SNAPSHOT 依赖 Data Layer 的 Nessie 分支能力,在执行前创建数据快照,失败时回滚到快照版本。

auto_rollback_on_failure 控制是否自动触发回滚,设为 false 时需要人工通过管理界面确认回滚。

#11. Key Takeaways

  1. oneof 替代 enum+configActionDefinition.logiconeof 设计提供了编译时类型安全,是 Proto3 中常被低估的特性。
  2. 主逻辑与副作用分离:SideEffect 的独立性使得通知、Webhook、审计等关注点完全解耦。
  3. 八态状态机:从 PENDING 到 ROLLBACK_COMPLETED 的完整状态机覆盖了企业级场景的所有路径。
  4. 批量执行的精细控制:fail_fast + max_parallel 的组合提供了安全性和性能之间的灵活权衡。
  5. 模板系统的约束验证:ConstraintType + SubmissionRule 的双层验证确保了参数在提交前的完整校验。
  6. 四种回滚策略:从简单的不回滚到 Saga 补偿模式,覆盖了不同的一致性需求。

#下一篇

S9-14:FunctionRuntime — 多语言沙箱的统一接口,我们将深入 Reasoning & Decision Layer 的函数运行时,了解 Python/TypeScript/Groovy 三种语言的沙箱隔离机制。

Tags: #coomia-dip #source-code-reading #action-engine #dispatcher-pattern #grpc #Layer-e #agent-runtime