源码精读: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 的实例化与参数约束验证流程。
源码精读:ActionEngine — 10 种执行器的调度器模式
“系列:S9 源码精读 · 第 13 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
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 的实例化与参数约束验证流程。
#目录
- 整体架构:从决策到行动的桥梁
- Proto 合约:ActionDefinition 的 oneof logic 设计
- ExecutionState 八态状态机
- ExecuteAction:单条执行的完整流程
- BatchExecuteActions:批量调度与并行控制
- RuleBasedExecution:声明式本体编辑
- FunctionExecution:函数驱动的行动
- SideEffect 四种副作用的触发机制
- ActionTemplate:模板实例化与参数约束
- 回滚策略:RollbackConfig 的四种模式
- Key Takeaways
#1. 整体架构:从决策到行动的桥梁
ActionEngine 在 coomia-dip 的八层架构中属于 Agent Runtime Layer(Agent Runtime),其核心职责是:接收来自 Reasoning & Decision Layer 决策引擎的决策结果,将其转化为具体的系统操作。代码结构如下:
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.logic 的 oneof 类型,将执行分发到不同的执行器。这种策略模式使新增执行器类型时只需扩展 proto 定义和对应的执行器实现,零修改核心调度逻辑。
#2. Proto 合约:ActionDefinition 的 oneof logic 设计
ActionEngine 最核心的设计决策体现在 ActionDefinition 消息的 oneof logic 字段上:
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_type 和 legacy_config 字段编号从 50 开始,与正常字段保持距离,避免未来扩展时的编号冲突。枚举 ActionType 虽然标记为 DEPRECATED,但保留其定义确保旧客户端不会反序列化失败。
副作用与主逻辑分离:side_effects 是 repeated 而非 oneof 的一部分,意味着一个 Action 可以有零到多个副作用,且副作用与主逻辑独立运行。
#3. ExecutionState 八态状态机
ActionEngine 定义了八种执行状态,构成了清晰的状态机:
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;
}
合法的状态转换路径:
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 客户端中对状态的封装同样值得关注:
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 执行的请求包含五个核心字段:
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 的可行性:
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 包含了丰富的执行信息:
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 的高频使用场景,其控制选项设计得很精细:
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 用于关联:
message ActionExecutionItem {
string item_id = 1; // 客户端提供的关联 ID
ActionDefinition action = 2;
google.protobuf.Struct input = 3;
ExecuteOptions options = 4;
}
响应中包含完整的汇总信息:
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,允许用户通过声明式规则描述本体编辑操作:
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:
create_object— 创建新的本体实例modify_object— 修改已有实例的属性delete_object— 删除实例(软删除)create_link— 建立实例间的关系delete_link— 删除实例间的关系create_or_modify— 幂等的 upsert 操作
条件执行:condition 字段允许规则根据运行时上下文条件性地执行。例如 "input.quantity > 0" 只在数量大于零时执行创建操作。
有序执行:order 字段确保多条规则按指定顺序执行。这对于先创建主对象、再创建关系这类有依赖的操作序列至关重要。
#7. FunctionExecution:函数驱动的行动
FunctionExecution 将 Action 的执行逻辑委托给 Reasoning & Decision Layer 的 FunctionRuntime:
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 提取为独立的副作用系统:
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:主逻辑失败后触发,典型场景是告警 WebhookALWAYS:无论成败都触发,典型场景是审计日志
Webhook 的执行模式(FEAT-009)进一步区分了同步和异步语义:
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 配置保存为可复用的模板:
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)支持五种约束类型:
enum ConstraintType {
CONSTRAINT_TYPE_RANGE = 1; // 数值范围
CONSTRAINT_TYPE_LENGTH = 2; // 字符串长度
CONSTRAINT_TYPE_PATTERN = 3; // 正则匹配
CONSTRAINT_TYPE_ENUM = 4; // 枚举值
CONSTRAINT_TYPE_CUSTOM = 5; // 自定义验证函数
}
SubmissionCriteria 将约束和规则组合在一起:
message SubmissionCriteria {
string validation_function_id = 1; // 外部验证函数
repeated ParameterConstraint constraints = 2;
repeated SubmissionRule rules = 3;
}
SubmissionRule 支持跨参数的布尔表达式验证,并区分 ERROR 和 WARNING 两个严重级别——WARNING 级别的验证失败会告警但不阻止提交。
模板的生命周期通过 TemplateStatus 枚举管理:
DRAFT → ENABLED ↔ DISABLED → DEPRECATED → ARCHIVED
#10. 回滚策略:RollbackConfig 的四种模式
当 Action 执行失败时,RollbackConfig 决定如何恢复:
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
- oneof 替代 enum+config:
ActionDefinition.logic的oneof设计提供了编译时类型安全,是 Proto3 中常被低估的特性。 - 主逻辑与副作用分离:SideEffect 的独立性使得通知、Webhook、审计等关注点完全解耦。
- 八态状态机:从 PENDING 到 ROLLBACK_COMPLETED 的完整状态机覆盖了企业级场景的所有路径。
- 批量执行的精细控制:fail_fast + max_parallel 的组合提供了安全性和性能之间的灵活权衡。
- 模板系统的约束验证:ConstraintType + SubmissionRule 的双层验证确保了参数在提交前的完整校验。
- 四种回滚策略:从简单的不回滚到 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