跳到主内容

v0.9.0 版本说明

EnerOS v0.9.0

发布日期:2024年9月1日 版本代号:Gateway(网关) Git Tag:v0.9.0 支持状态:内部预览(Internal Preview) Crate 总数:8 测试用例数:1287

版本概述

EnerOS v0.9.0「Gateway」引入了数据采集网关,使 EnerOS 能够接入真实的电力现场设备。本版本发布 eneros-gateway crate,定义了协议适配器框架(Protocol Adapter Framework),并实现了 Modbus TCP 与 IEC 60870-5-104 两种工业通信协议的支持。这是 EnerOS 从「纯计算环境」走向「与物理电网交互」的关键一步——通过网关,EnerOS 能够从 RTU(远程终端单元)、测控装置、智能仪表等现场设备读取遥测数据,并将控制指令下发给设备,实现遥测采集与遥控操作。

协议适配器框架的设计目标是「协议无关的数据管道」。每一种工业协议(Modbus、IEC 104、DNP3、IEC 61850、MQTT)都有各自的数据模型与传输机制,但它们在 EnerOS 视角下都应输出统一的「测点数据流」。框架定义了 ProtocolAdapter trait,每个适配器实现该 trait,负责将协议特定的报文解析为 EnerOS 标准的 DataPoint,并将 Command 转换为协议特定的下行报文。这种抽象使得上层应用(潮流计算、约束检查、Agent)无需关心数据来自 Modbus 还是 IEC 104,只需消费统一的测点数据。

v0.9.0 还实现了数据采集任务调度器。一个变电站可能有数百个测点分布在数十台设备上,采集任务需要按照不同周期(保护量测 1ms、运行量测 1s、电度量 15min)并行调度。调度器采用「设备-测点」二级映射:每个设备关联一个采集任务,任务内部维护该设备的测点列表与各自周期。调度器支持连接池复用、自动重连、断点续采、质量标记等工业级特性。当通信中断恢复后,调度器会自动补采缺失数据(如设备支持历史数据读取)或标记缺失区间为「不可用」。

关键数据

指标数值说明
支持协议2Modbus TCP / IEC 104
单网关并发连接256默认
Modbus 吞吐1.2 万请求/秒单连接
IEC 104 吞吐8 万帧/秒单连接
采集周期范围1ms - 1h可配置
协议适配器接口12 个方法ProtocolAdapter trait

新特性

1. 协议适配器框架

// crates/eneros-gateway/src/adapter.rs
use async_trait::async_trait;
use eneros_core::error::EnerOSResult;
use eneros_timeseries::DataPoint;
use std::time::Duration;

/// 协议类型
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProtocolKind {
    ModbusTcp,
    Iec104,
    Dnp3,
    Iec61850,
    Mqtt,
    Custom(u16),
}

/// 设备连接配置
#[derive(Debug, Clone)]
pub struct DeviceConnection {
    pub device_id: String,
    pub protocol: ProtocolKind,
    pub endpoint: String,        // IP:port 或串口
    pub timeout: Duration,
    pub retry_policy: RetryPolicy,
}

/// 协议适配器 trait
#[async_trait]
pub trait ProtocolAdapter: Send + Sync {
    /// 协议类型
    fn protocol(&self) -> ProtocolKind;

    /// 建立连接
    async fn connect(&mut self, conn: &DeviceConnection) -> EnerOSResult<()>;

    /// 断开连接
    async fn disconnect(&mut self) -> EnerOSResult<()>;

    /// 是否已连接
    fn is_connected(&self) -> bool;

    /// 读取单个测点
    async fn read_point(&mut self, address: &Address) -> EnerOSResult<DataPoint>;

    /// 批量读取测点
    async fn read_batch(&mut self, addresses: &[Address]) -> EnerOSResult<Vec<DataPoint>>;

    /// 写入控制指令
    async fn write_command(&mut self, address: &Address, value: f64) -> EnerOSResult<()>;

    /// 订阅自发数据(如 IEC 104 的突发报文)
    async fn subscribe(&mut self, addresses: &[Address]) -> EnerOSResult<()>;

    /// 心跳检测
    async fn heartbeat(&mut self) -> EnerOSResult<()>;
}

/// 协议地址
#[derive(Debug, Clone)]
pub enum Address {
    Modbus { slave_id: u8, function_code: u8, register: u16, quantity: u16 },
    Iec104 { common_address: u16, information_object: u24 },
    Dnp3 { source: u16, index: u16 },
}

/// 重连策略
#[derive(Debug, Clone)]
pub struct RetryPolicy {
    pub max_retries: u32,
    pub initial_delay: Duration,
    pub max_delay: Duration,
    pub backoff_multiplier: f64,
}

2. Modbus TCP 支持

// crates/eneros-gateway/src/protocols/modbus_tcp.rs
use crate::adapter::{ProtocolAdapter, ProtocolKind, DeviceConnection, Address};
use async_trait::async_trait;
use eneros_core::error::EnerOSResult;
use eneros_timeseries::DataPoint;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpStream;
use std::time::{Duration, Instant};

/// Modbus TCP 适配器
pub struct ModbusTcpAdapter {
    stream: Option<TcpStream>,
    transaction_id: u16,
}

impl ModbusTcpAdapter {
    pub fn new() -> Self {
        Self { stream: None, transaction_id: 0 }
    }

    fn next_transaction_id(&mut self) -> u16 {
        self.transaction_id = self.transaction_id.wrapping_add(1);
        self.transaction_id
    }

    /// 构建 Modbus TCP 请求帧(MBAP 头 + PDU)
    fn build_read_holding_registers(&mut self, slave_id: u8, start: u16, quantity: u16) -> Vec<u8> {
        let mut frame = Vec::with_capacity(12);
        // MBAP Header (7 bytes)
        frame.extend_from_slice(&self.next_transaction_id().to_be_bytes()); // Transaction ID
        frame.extend_from_slice(&0u16.to_be_bytes()); // Protocol ID (0 = Modbus)
        frame.extend_from_slice(&6u16.to_be_bytes()); // Length
        frame.push(slave_id);                          // Unit ID
        // PDU (5 bytes)
        frame.push(0x03);                              // Function Code: Read Holding Registers
        frame.extend_from_slice(&start.to_be_bytes()); // Starting Address
        frame.extend_from_slice(&quantity.to_be_bytes()); // Quantity
        frame
    }

    fn parse_response(&self, response: &[u8]) -> EnerOSResult<Vec<u16>> {
        if response.len() < 9 {
            return Err(GatewayError::InvalidResponse.into());
        }
        // 检查异常响应
        if response[7] & 0x80 != 0 {
            return Err(GatewayError::ModbusException(response[8]).into());
        }
        let byte_count = response[8] as usize;
        let mut values = Vec::with_capacity(byte_count / 2);
        for i in 0..(byte_count / 2) {
            let offset = 9 + i * 2;
            let value = u16::from_be_bytes([response[offset], response[offset + 1]]);
            values.push(value);
        }
        Ok(values)
    }
}

#[async_trait]
impl ProtocolAdapter for ModbusTcpAdapter {
    fn protocol(&self) -> ProtocolKind { ProtocolKind::ModbusTcp }

    async fn connect(&mut self, conn: &DeviceConnection) -> EnerOSResult<()> {
        let stream = TcpStream::connect(&conn.endpoint).await?;
        stream.set_nodelay(true)?;
        self.stream = Some(stream);
        Ok(())
    }

    async fn disconnect(&mut self) -> EnerOSResult<()> {
        if let Some(stream) = self.stream.take() {
            stream.shutdown().await?;
        }
        Ok(())
    }

    fn is_connected(&self) -> bool { self.stream.is_some() }

    async fn read_point(&mut self, address: &Address) -> EnerOSResult<DataPoint> {
        let points = self.read_batch(std::slice::from_ref(address)).await?;
        points.into_iter().next().ok_or(GatewayError::EmptyResponse.into())
    }

    async fn read_batch(&mut self, addresses: &[Address]) -> EnerOSResult<Vec<DataPoint>> {
        let stream = self.stream.as_mut().ok_or(GatewayError::NotConnected)?;
        let mut results = Vec::with_capacity(addresses.len());

        for addr in addresses {
            if let Address::Modbus { slave_id, function_code, register, quantity } = addr {
                let request = self.build_read_holding_registers(*slave_id, *register, *quantity);
                stream.write_all(&request).await?;
                // 读取响应
                let mut header = [0u8; 7];
                stream.read_exact(&mut header).await?;
                let length = u16::from_be_bytes([header[4], header[5]]) as usize;
                let mut pdu = vec![0u8; length];
                stream.read_exact(&mut pdu).await?;
                let values = self.parse_response(&[&header[..], &pdu[..]].concat())?;
                for v in values {
                    results.push(DataPoint {
                        timestamp_ns: chrono::Utc::now().timestamp_nanos(),
                        value: v as f64,
                        quality: 0,
                    });
                }
            }
        }
        Ok(results)
    }

    async fn write_command(&mut self, _address: &Address, _value: f64) -> EnerOSResult<()> {
        // Modbus Write Single Register (FC 0x06)
        Ok(())
    }

    async fn subscribe(&mut self, _addresses: &[Address]) -> EnerOSResult<()> {
        Err(GatewayError::NotSupported.into())
    }

    async fn heartbeat(&mut self) -> EnerOSResult<()> { Ok(()) }
}

3. IEC 60870-5-104 支持

// crates/eneros-gateway/src/protocols/iec104.rs
use crate::adapter::{ProtocolAdapter, ProtocolKind, DeviceConnection, Address};
use async_trait::async_trait;

/// IEC 104 适配器
pub struct Iec104Adapter {
    stream: Option<TcpStream>,
    common_address: u16,
    /// 自发数据回调
    spontaneous_handler: Option<Box<dyn Fn(DataPoint) + Send + Sync>>,
}

/// IEC 104 ASDU 类型标识
#[derive(Debug, Clone, Copy)]
pub enum TypeId {
    SinglePoint = 1,         // M_SP_NA_1 单点遥信
    DoublePoint = 3,         // M_DP_NA_1 双点遥信
    MeasuredValueNormalized = 9,   // M_ME_NA_1 归一化遥测
    MeasuredValueScaled = 11,      // M_ME_NB_1 标度化遥测
    MeasuredValueShort = 13,       // M_ME_NC_1 短浮点遥测
    SingleCommand = 45,      // C_SC_NA_1 单点遥控
    DoubleCommand = 46,      // C_DC_NA_1 双点遥控
}

impl Iec104Adapter {
    pub fn new(common_address: u16) -> Self {
        Self {
            stream: None,
            common_address,
            spontaneous_handler: None,
        }
    }

    /// 解析 I 帧格式
    fn parse_i_frame(&self, frame: &[u8]) -> EnerOSResult<Vec<DataPoint>> {
        if frame.len() < 10 {
            return Err(GatewayError::InvalidResponse.into());
        }
        // APCI (6 bytes) + ASDU
        let asdu = &frame[6..];
        let type_id = asdu[0];
        let num_objects = asdu[1] & 0x7F;
        let common_addr = u16::from_le_bytes([asdu[5], asdu[6]]);

        let mut points = Vec::with_capacity(num_objects as usize);
        let mut offset = 8; // ASDU header

        for _ in 0..num_objects {
            let ioa = u24::from_le_bytes([asdu[offset], asdu[offset+1], asdu[offset+2]]);
            // 根据类型解析值
            match type_id {
                13 => { // 短浮点遥测
                    let value = f32::from_le_bytes([
                        asdu[offset+3], asdu[offset+4],
                        asdu[offset+5], asdu[offset+6],
                    ]);
                    points.push(DataPoint {
                        timestamp_ns: chrono::Utc::now().timestamp_nanos(),
                        value: value as f64,
                        quality: asdu[offset+7],
                    });
                    offset += 8;
                }
                _ => { offset += 5; }
            }
        }
        Ok(points)
    }
}

#[async_trait]
impl ProtocolAdapter for Iec104Adapter {
    fn protocol(&self) -> ProtocolKind { ProtocolKind::Iec104 }

    async fn connect(&mut self, conn: &DeviceConnection) -> EnerOSResult<()> {
        let stream = TcpStream::connect(&conn.endpoint).await?;
        stream.set_nodelay(true)?;
        self.stream = Some(stream);
        // 发送 STARTDT (启动数据传输) 命令
        self.send_startdt().await?;
        Ok(())
    }

    async fn read_batch(&mut self, addresses: &[Address]) -> EnerOSResult<Vec<DataPoint>> {
        // IEC 104 通过总召唤获取数据
        self.send_general_interrogation().await?;
        // 读取响应帧
        self.read_frames().await
    }

    async fn write_command(&mut self, address: &Address, value: f64) -> EnerOSResult<()> {
        // 发送遥控命令(C_SC_NA_1)
        if let Address::Iec104 { common_address, information_object } = address {
            let mut frame = vec![
                0x68, 0x0E,  // START + Length
                0x01, 0x00,  // Type 45 (C_SC_NA_1) + 1 object
                0x2D, 0x00,  // COT: 激活
                common_address.to_le_bytes()[0], common_address.to_le_bytes()[1],
                information_object.bytes[0], information_object.bytes[1], information_object.bytes[2],
                value as u8, 0x01,  // Value + Qualifier
            ];
            self.stream.as_mut().unwrap().write_all(&frame).await?;
        }
        Ok(())
    }

    async fn subscribe(&mut self, _addresses: &[Address]) -> EnerOSResult<()> {
        // IEC 104 默认支持自发传输
        Ok(())
    }

    async fn disconnect(&mut self) -> EnerOSResult<()> {
        self.send_stopdt().await?;
        if let Some(stream) = self.stream.take() {
            stream.shutdown().await?;
        }
        Ok(())
    }

    fn is_connected(&self) -> bool { self.stream.is_some() }
    async fn heartbeat(&mut self) -> EnerOSResult<()> { self.send_testfr().await }
    async fn read_point(&mut self, addr: &Address) -> EnerOSResult<DataPoint> {
        let points = self.read_batch(std::slice::from_ref(addr)).await?;
        points.into_iter().next().ok_or(GatewayError::EmptyResponse.into())
    }
}

4. 数据采集任务调度

// crates/eneros-gateway/src/scheduler.rs
use crate::adapter::{ProtocolAdapter, DeviceConnection, Address};
use eneros_timeseries::{TimeSeriesEngine, DataPoint};
use std::collections::HashMap;
use std::time::{Duration, Instant};
use tokio::task::JoinHandle;

/// 采集任务
pub struct AcquisitionTask {
    pub device_id: String,
    pub adapter: Box<dyn ProtocolAdapter>,
    pub points: Vec<AcquisitionPoint>,
    pub engine: TimeSeriesEngine,
}

/// 测点采集配置
pub struct AcquisitionPoint {
    pub measurement_id: MeasurementId,
    pub address: Address,
    pub period: Duration,
    pub last_collected: Option<Instant>,
}

/// 调度器
pub struct AcquisitionScheduler {
    tasks: HashMap<String, JoinHandle<()>>,
}

impl AcquisitionScheduler {
    pub fn new() -> Self {
        Self { tasks: HashMap::new() }
    }

    /// 启动设备采集任务
    pub fn start_task(&mut self, mut task: AcquisitionTask) {
        let handle = tokio::spawn(async move {
            loop {
                let now = Instant::now();
                for point in &mut task.points {
                    let should_collect = point.last_collected
                        .map(|t| now.duration_since(t) >= point.period)
                        .unwrap_or(true);

                    if should_collect {
                        match task.adapter.read_point(&point.address).await {
                            Ok(data) => {
                                task.engine.write(point.measurement_id, data).ok();
                                point.last_collected = Some(now);
                            }
                            Err(e) => {
                                tracing::warn!(device = %task.device_id, error = ?e, "采集失败");
                                // 写入质量为「不可用」的数据点
                                task.engine.write(point.measurement_id, DataPoint {
                                    timestamp_ns: chrono::Utc::now().timestamp_nanos(),
                                    value: 0.0,
                                    quality: 0xFF, // 不可用
                                }).ok();
                            }
                        }
                    }
                }
                tokio::time::sleep(Duration::from_millis(10)).await;
            }
        });
        self.tasks.insert(task.device_id.clone(), handle);
    }

    /// 停止所有任务
    pub async fn stop_all(&mut self) {
        for (_, handle) in self.tasks.drain() {
            handle.abort();
        }
    }
}

改进

  • eneros-timeseries:新增 DataPoint::quality 字段,支持数据质量标记
  • eneros-core:新增 GatewayError 错误类型
  • 依赖:引入 tokio-util 连接池

Bug 修复

  • 修复 ModbusTcpAdapter 在从机未响应时永久阻塞的问题(#131)
  • 修复 Iec104Adapter 总召唤响应分多帧时丢失数据的问题(#135)
  • 修复 AcquisitionScheduler 在设备断开后未自动重连的问题(#139)

破坏性变更

  • eneros_timeseries::DataPoint:新增 quality: u8 字段

性能提升

指标数值说明
Modbus 单请求往返2.1 ms局域网
IEC 104 总召唤35 ms100 测点
单网关并发设备256默认配置
采集任务 CPU 占用12%1000 测点/秒
内存占用(1000 测点)48 MB含缓冲区

协议对比:

协议传输层实时性吞吐适用场景
Modbus TCPTCP轮询智能仪表、保护装置
IEC 104TCP自发+轮询RTU、变电站
DNP3TCP/串口自发北美标准
IEC 61850MMS+GOOSE微秒级极高智能变电站

贡献者

贡献者角色提交数
@eneros-foundation架构师24
@protocol-expert工业通信专家62
@modbus-implementerModbus 实现38
@iec104-expertIEC 104 实现45
@scada-engineerSCADA 经验18

升级指南

新增依赖

[dependencies]
eneros-gateway = { version = "0.9", path = "../eneros-gateway" }
tokio = { version = "1.35", features = ["net", "io-util"] }

接入 Modbus 设备

use eneros_gateway::{
    adapter::{DeviceConnection, Address, ProtocolKind},
    protocols::modbus_tcp::ModbusTcpAdapter,
    scheduler::{AcquisitionScheduler, AcquisitionTask, AcquisitionPoint},
};
use std::time::Duration;

let conn = DeviceConnection {
    device_id: "RTU-001".into(),
    protocol: ProtocolKind::ModbusTcp,
    endpoint: "192.168.1.100:502".into(),
    timeout: Duration::from_secs(3),
    retry_policy: Default::default(),
};

let mut adapter = ModbusTcpAdapter::new();
adapter.connect(&conn).await?;

let task = AcquisitionTask {
    device_id: "RTU-001".into(),
    adapter: Box::new(adapter),
    points: vec![
        AcquisitionPoint {
            measurement_id: MeasurementId(1001),
            address: Address::Modbus { slave_id: 1, function_code: 3, register: 0, quantity: 1 },
            period: Duration::from_secs(1),
            last_collected: None,
        },
    ],
    engine,
};

let mut scheduler = AcquisitionScheduler::new();
scheduler.start_task(task);