跳到主内容

v0.6.0 版本说明

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完整生命周期
消息传递延迟 P991.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)
    }
}

状态转换图:

当前状态事件新状态
CreatedInitInitialized
InitializedStartRunning
RunningSuspendSuspended
SuspendedResumeRunning
任意StopStopped
任意FailFailed

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() 方法
  • 依赖:引入 tokioasync-traitchrono
  • 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-
消息传递(同节点)P991.8 ms5.5 万/秒
状态机转换18 ns-
1024 Agent 内存占用48 MB-

消息传递延迟分布:

百分位延迟
P50420 μs
P901.1 ms
P991.8 ms
P99.93.2 ms

贡献者

贡献者角色提交数
@eneros-foundation架构师38
@actor-model-expertActor 模型专家35
@async-rustacean异步 Rust 工程师41
@ros2-contributorROS 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(()) }
}