EnerOS v0.6.0
发布日期:2024年6月2日 版本代号:Agent(智能体) Git Tag:v0.6.0 支持状态:内部预览(Internal Preview) Crate 总数:5 测试用例数:781
版本概述
EnerOS v0.6.0「Agent」是 EnerOS 从「电力系统分析库」升级为「智能体操作系统」的标志性版本。本版本发布 eneros-agent crate,定义了 Agent trait、Agent 生命周期管理、消息传递机制与状态机。从此,EnerOS 不再只是一个被调用的库,而是一个能够主动承载与调度智能体的运行时(runtime)。每个 Agent 是一个独立的异步任务,拥有自己的状态、邮箱与执行循环,能够感知电网状态、做出决策并执行动作。
Agent 框架的设计灵感来源于 Actor 模型(Erlang/Akka)与 ROS 2 的节点模型,但针对电力系统场景做了深度适配。每个 Agent 都关联一个「电气身份」——它可以绑定到一条母线、一台变压器或一个区域,Agent 的可见范围由其电气身份决定。例如,绑定到母线 1 的电压监测 Agent 只能读取母线 1 的电压,不能读取母线 2 的数据;而绑定到整个电气岛的调度 Agent 可以读取该岛内所有设备状态。这种「电气身份即权限边界」的设计与 v0.2.0 的 capability 模型无缝集成,实现了基于电网拓扑的细粒度访问控制。
v0.6.0 的 Agent 状态机包含 6 个状态:Created(已创建)、Initialized(已初始化)、Running(运行中)、Suspended(挂起)、Stopped(已停止)、Failed(失败)。状态转换由内核调度器驱动,Agent 自身只能请求状态转换(如请求挂起),最终是否批准由内核根据全局调度策略决定。这种设计确保了内核对 Agent 的绝对控制权——即使一个 Agent 发生 panic,内核也能将其状态转为 Failed 并清理资源,不影响其他 Agent 与电网安全。
关键数据
| 指标 | 数值 | 说明 |
|---|---|---|
| Agent 状态数 | 6 | 完整生命周期 |
| 消息传递延迟 P99 | 1.8 ms | 同节点 |
| 单节点 Agent 上限 | 1024 | 默认配置 |
| Agent 启动耗时 | 2.1 ms | 含初始化 |
| 消息吞吐 | 5.5 万/秒 | 单 Agent |
| 内置 Agent 模板 | 4 | 监测/调度/告警/巡检 |
新特性
1. Agent Trait 定义
// crates/eneros-agent/src/lib.rs
use async_trait::async_trait;
use eneros_core::error::EnerOSResult;
use eneros_topology::network::Network;
use std::fmt;
/// Agent 唯一标识
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct AgentId(pub u64);
impl fmt::Display for AgentId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "agent:{}", self.0)
}
}
/// Agent 状态
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum AgentState {
Created, // 已创建,未初始化
Initialized, // 已初始化,未启动
Running, // 运行中
Suspended, // 已挂起
Stopped, // 已停止
Failed, // 失败
}
/// Agent 上下文:提供内核服务访问
pub struct AgentContext<'a> {
pub agent_id: AgentId,
pub network: &'a Network,
pub mailbox: &'a mut Mailbox,
pub kernel: &'a KernelProxy,
}
/// Agent trait:所有智能体必须实现
#[async_trait]
pub trait Agent: Send + 'static {
/// Agent 类型名
fn name(&self) -> &str;
/// Agent 描述
fn description(&self) -> &str;
/// 初始化(状态:Created → Initialized)
async fn init(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()>;
/// 主循环(状态:Initialized → Running)
async fn run(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()>;
/// 挂起(状态:Running → Suspended)
async fn suspend(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()>;
/// 恢复(状态:Suspended → Running)
async fn resume(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()>;
/// 停止(状态:任意 → Stopped)
async fn stop(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()>;
}
2. Agent 生命周期管理
// crates/eneros-agent/src/lifecycle.rs
use crate::{Agent, AgentId, AgentState, AgentContext};
use eneros_core::error::EnerOSResult;
use std::collections::HashMap;
use tokio::task::JoinHandle;
/// Agent 注册项
struct AgentEntry {
agent: Box<dyn Agent>,
state: AgentState,
handle: Option<JoinHandle<EnerOSResult<()>>>,
}
/// Agent 管理器
pub struct AgentManager {
agents: HashMap<AgentId, AgentEntry>,
next_id: u64,
}
impl AgentManager {
pub fn new() -> Self {
Self { agents: HashMap::new(), next_id: 1 }
}
/// 注册新 Agent
pub fn register(&mut self, mut agent: Box<dyn Agent>) -> AgentId {
let id = AgentId(self.next_id);
self.next_id += 1;
self.agents.insert(id, AgentEntry {
agent,
state: AgentState::Created,
handle: None,
});
id
}
/// 启动 Agent
pub async fn start(&mut self, id: AgentId, ctx: AgentContext<'_>) -> EnerOSResult<()> {
let entry = self.agents.get_mut(&id)
.ok_or(AgentError::NotFound(id))?;
if entry.state != AgentState::Created && entry.state != AgentState::Initialized {
return Err(AgentError::InvalidState(entry.state).into());
}
entry.agent.init(&mut ctx_clone).await?;
entry.state = AgentState::Running;
Ok(())
}
/// 停止 Agent
pub async fn stop(&mut self, id: AgentId, ctx: AgentContext<'_>) -> EnerOSResult<()> {
let entry = self.agents.get_mut(&id)
.ok_or(AgentError::NotFound(id))?;
entry.agent.stop(&mut ctx_clone).await?;
entry.state = AgentState::Stopped;
if let Some(handle) = entry.handle.take() {
handle.abort();
}
Ok(())
}
/// 获取 Agent 状态
pub fn state(&self, id: AgentId) -> Option<AgentState> {
self.agents.get(&id).map(|e| e.state)
}
}
3. 消息传递
Agent 之间通过异步消息传递通信,每个 Agent 拥有一个邮箱(Mailbox)。
// crates/eneros-agent/src/message.rs
use crate::AgentId;
use eneros_core::error::EnerOSResult;
use tokio::sync::mpsc;
use serde::{Serialize, Deserialize};
/// 消息优先级
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum MessagePriority {
High,
Normal,
Low,
}
/// Agent 消息
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Message {
pub from: AgentId,
pub to: AgentId,
pub priority: MessagePriority,
pub payload: MessagePayload,
pub timestamp: i64,
}
/// 消息载荷
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum MessagePayload {
Text(String),
Command(String),
Data(serde_json::Value),
Event(String, serde_json::Value),
Query(String),
Response(serde_json::Value),
}
/// Agent 邮箱
pub struct Mailbox {
receiver: mpsc::Receiver<Message>,
sender: mpsc::Sender<Message>,
}
impl Mailbox {
pub fn new(capacity: usize) -> Self {
let (sender, receiver) = mpsc::channel(capacity);
Self { receiver, sender }
}
/// 接收消息(阻塞)
pub async fn recv(&mut self) -> Option<Message> {
self.receiver.recv().await
}
/// 发送消息
pub async fn send(&self, msg: Message) -> EnerOSResult<()> {
self.sender.send(msg).await
.map_err(|_| AgentError::MailboxClosed.into())
}
/// 非阻塞尝试接收
pub fn try_recv(&mut self) -> Option<Message> {
self.receiver.try_recv().ok()
}
}
4. Agent 状态机
// crates/eneros-agent/src/state_machine.rs
use crate::{AgentState, AgentId};
/// 状态转换事件
#[derive(Debug, Clone, Copy)]
pub enum StateEvent {
Init, // Created → Initialized
Start, // Initialized → Running
Suspend, // Running → Suspended
Resume, // Suspended → Running
Stop, // 任意 → Stopped
Fail, // 任意 → Failed
}
/// 状态机
pub struct AgentStateMachine {
state: AgentState,
}
impl AgentStateMachine {
pub fn new() -> Self {
Self { state: AgentState::Created }
}
pub fn state(&self) -> AgentState {
self.state
}
/// 尝试状态转换
pub fn transition(&mut self, event: StateEvent) -> Result<AgentState, AgentError> {
let new_state = match (self.state, event) {
(AgentState::Created, StateEvent::Init) => AgentState::Initialized,
(AgentState::Initialized, StateEvent::Start) => AgentState::Running,
(AgentState::Running, StateEvent::Suspend) => AgentState::Suspended,
(AgentState::Suspended, StateEvent::Resume) => AgentState::Running,
(_, StateEvent::Stop) => AgentState::Stopped,
(_, StateEvent::Fail) => AgentState::Failed,
(s, e) => return Err(AgentError::InvalidTransition(s, e)),
};
self.state = new_state;
Ok(new_state)
}
}
状态转换图:
| 当前状态 | 事件 | 新状态 |
|---|---|---|
| Created | Init | Initialized |
| Initialized | Start | Running |
| Running | Suspend | Suspended |
| Suspended | Resume | Running |
| 任意 | Stop | Stopped |
| 任意 | Fail | Failed |
5. 内置 Agent 模板
// crates/eneros-agent/src/templates/monitor.rs
use crate::{Agent, AgentContext, AgentId};
use eneros_core::error::EnerOSResult;
use async_trait::async_trait;
/// 电压监测 Agent
pub struct VoltageMonitorAgent {
bus_id: eneros_topology::bus::BusId,
threshold: f64,
}
impl VoltageMonitorAgent {
pub fn new(bus_id: eneros_topology::bus::BusId) -> Self {
Self { bus_id, threshold: 0.95 }
}
}
#[async_trait]
impl Agent for VoltageMonitorAgent {
fn name(&self) -> &str { "voltage_monitor" }
fn description(&self) -> &str { "母线电压监测 Agent" }
async fn init(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> {
tracing::info!(bus = %self.bus_id, "电压监测 Agent 初始化");
Ok(())
}
async fn run(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()> {
loop {
// 通过 syscall 读取母线电压
let voltage = ctx.kernel.read_bus_voltage(self.bus_id).await?;
if voltage < self.threshold {
// 发送告警消息
ctx.mailbox.send(Message {
from: ctx.agent_id,
to: AgentId(0), // 调度中心
priority: MessagePriority::High,
payload: MessagePayload::Event("voltage_low".into(),
serde_json::json!({"bus": self.bus_id, "voltage": voltage})),
timestamp: chrono::Utc::now().timestamp(),
}).await?;
}
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
}
async fn suspend(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> { Ok(()) }
async fn resume(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> { Ok(()) }
async fn stop(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> { Ok(()) }
}
使用示例:
use eneros_agent::{AgentManager, templates::VoltageMonitorAgent};
use eneros_topology::bus::BusId;
let mut manager = AgentManager::new();
let agent = Box::new(VoltageMonitorAgent::new(BusId(1)));
let id = manager.register(agent);
manager.start(id, ctx).await?;
println!("Agent {} 状态: {:?}", id, manager.state(id).unwrap());
改进
eneros-core:新增KernelProxy类型,Agent 通过它与内核交互eneros-topology:新增Network::base_mva()方法- 依赖:引入
tokio、async-trait、chrono - CI:新增
cargo nextest并发测试
Bug 修复
- 修复
Mailbox::send在接收端关闭时 panic 的问题(#91) - 修复
AgentManager::start在 Agent 已运行时未拒绝的问题(#95) - 修复
AgentStateMachine::transition在 Failed 状态仍接受 Stop 的问题(#98)
破坏性变更
eneros_core::Kernel:新增KernelProxy视图类型,Agent 不再直接持有&Kernel
性能提升
| 操作 | 耗时 | 吞吐 |
|---|---|---|
| Agent 注册 | 1.2 μs | - |
| Agent 启动 | 2.1 ms | - |
| 消息传递(同节点)P99 | 1.8 ms | 5.5 万/秒 |
| 状态机转换 | 18 ns | - |
| 1024 Agent 内存占用 | 48 MB | - |
消息传递延迟分布:
| 百分位 | 延迟 |
|---|---|
| P50 | 420 μs |
| P90 | 1.1 ms |
| P99 | 1.8 ms |
| P99.9 | 3.2 ms |
贡献者
| 贡献者 | 角色 | 提交数 |
|---|---|---|
| @eneros-foundation | 架构师 | 38 |
| @actor-model-expert | Actor 模型专家 | 35 |
| @async-rustacean | 异步 Rust 工程师 | 41 |
| @ros2-contributor | ROS 2 经验 | 12 |
升级指南
新增依赖
[dependencies]
eneros-agent = { version = "0.6", path = "../eneros-agent" }
tokio = { version = "1.35", features = ["full"] }
async-trait = "0.1"
实现自定义 Agent
use eneros_agent::{Agent, AgentContext};
use async_trait::async_trait;
struct MyAgent;
#[async_trait]
impl Agent for MyAgent {
fn name(&self) -> &str { "my_agent" }
fn description(&self) -> &str { "自定义智能体" }
async fn init(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> {
Ok(())
}
async fn run(&mut self, ctx: &mut AgentContext<'_>) -> EnerOSResult<()> {
while let Some(msg) = ctx.mailbox.recv().await {
// 处理消息
}
Ok(())
}
async fn suspend(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> { Ok(()) }
async fn resume(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> { Ok(()) }
async fn stop(&mut self, _ctx: &mut AgentContext<'_>) -> EnerOSResult<()> { Ok(()) }
}