跳到主内容

WebSocket

API 参考

WebSocket 事件流

EnerOS 通过 WebSocket 提供双向实时通信通道,用于 Agent 状态推送、电网事件广播、远程指令下发等场景。底层由 eneros-eventbus 产生事件,经 eneros-gateway 鉴权后推送至订阅客户端。WebSocket 协议基于 HTTP 升级握手,连接建立后可保持长连接,相比 REST 反复建连,可显著降低延迟与资源消耗。

端点信息

项目取值
URLws://localhost:8080/ws
生产 URLwss://api.example.com/ws
子协议eneros.v1(必须在握手时通过 Sec-WebSocket-Protocol 指定)
心跳周期30 秒
心跳超时60 秒(未响应 pong 则断开)
单连接最大订阅50 个主题
单租户最大连接20 个并发
消息大小上限256 KB
重连补发重连后自动补发最近 100 条事件
鉴权方式握手后第一条消息发送 auth

连接建立

握手流程

客户端                                          服务端
  │                                               │
  │── HTTP Upgrade + Sec-WebSocket-Protocol ────→ │
  │                                               │
  │←── 101 Switching Protocols ──────────────────│
  │                                               │
  │── auth 帧(含 token) ──────────────────────→ │
  │                                               │
  │←── auth_ok 帧 ───────────────────────────────│
  │                                               │
  │── subscribe 帧(订阅主题) ─────────────────→ │
  │                                               │
  │←── subscribed 帧(订阅确认) ────────────────│
  │                                               │
  │←── event 帧(事件推送) ─────────────────────│
  │                                               │
  │── ping 帧(每 30s) ────────────────────────→ │
  │                                               │
  │←── pong 帧 ───────────────────────────────────│

JavaScript 握手示例

const ws = new WebSocket('ws://localhost:8080/ws', ['eneros.v1']);

ws.onopen = () => {
  console.log('连接已建立,发送鉴权');
  ws.send(JSON.stringify({
    type: 'auth',
    token: '<bearer-token>',
    client_id: 'dashboard-001',
    client_version: '1.0.0',
  }));
};

ws.onmessage = (event) => {
  const msg = JSON.parse(event.data);
  switch (msg.type) {
    case 'auth_ok':
      console.log('鉴权成功,开始订阅');
      ws.send(JSON.stringify({
        type: 'subscribe',
        topics: ['network.net_8f3a2b', 'agent.*', 'alarm.critical'],
        ack: true,
      }));
      break;
    case 'event':
      console.log('收到事件:', msg.payload);
      break;
    case 'error':
      console.error('错误:', msg.payload);
      break;
  }
};

ws.onerror = (e) => console.error('连接错误', e);
ws.onclose = (e) => console.log('连接关闭:', e.code, e.reason);

消息格式

所有消息为 JSON 文本帧(不支持二进制帧),统一字段约定如下:

字段类型必填说明
typestring消息类型,见下方枚举
idstring消息 ID,用于 ACK 关联
timestampstring消息时间戳(RFC 3339)
payloadobject消息体,结构随 type 变化

消息类型枚举

方向类型说明
客户端 → 服务端auth鉴权
客户端 → 服务端subscribe订阅主题
客户端 → 服务端unsubscribe取消订阅
客户端 → 服务端command下发 Agent 指令
客户端 → 服务端ping心跳探测
客户端 → 服务端ack确认收到事件
服务端 → 客户端auth_ok鉴权成功
服务端 → 客户端auth_failed鉴权失败
服务端 → 客户端subscribed订阅确认
服务端 → 客户端unsubscribed取消订阅确认
服务端 → 客户端event事件推送
服务端 → 客户端command_result指令执行结果
服务端 → 客户端pong心跳响应
服务端 → 客户端error错误通知
服务端 → 客户端reconnect通知客户端重连(服务端即将关闭)

auth 帧

{
  "type": "auth",
  "token": "<bearer-token>",
  "client_id": "dashboard-001",
  "client_version": "1.0.0",
  "capabilities": ["ack", "batch"]
}
字段类型必填默认值说明
tokenstringBearer Token
client_idstring自动生成客户端标识
client_versionstring""客户端版本
capabilitiesarray[]客户端能力声明

auth_ok 帧

{
  "type": "auth_ok",
  "session_id": "sess_20260706_001",
  "tenant": "default",
  "expires_at": "2026-07-06T09:30:00Z",
  "server_time": "2026-07-06T08:30:00Z"
}

subscribe 帧

{
  "type": "subscribe",
  "topics": ["network.net_8f3a2b", "agent.*", "alarm.critical"],
  "ack": true,
  "last_event_id": "evt_20260706_001"
}
字段类型必填默认值说明
topicsarray主题列表,支持通配符 *
ackboolfalse是否要求事件 ACK
last_event_idstring断点续传,从此 ID 之后开始补发

subscribed 帧

{
  "type": "subscribed",
  "topics": ["network.net_8f3a2b", "agent.*", "alarm.critical"],
  "pending_replay": 12
}
字段类型必填说明
topicsarray已订阅主题列表
pending_replayint待补发的事件数量

event 帧

{
  "type": "event",
  "id": "evt_a1b2c3",
  "topic": "network.net_8f3a2b",
  "timestamp": "2026-07-06T08:30:00.123Z",
  "payload": {
    "kind": "TopologyChanged",
    "network_id": "net_8f3a2b",
    "changes": [{"branch": "1-2", "status": "open"}],
    "actor": {"type": "agent", "id": "agent_dispatch_01"}
  }
}
字段类型必填说明
idstring事件唯一 ID,用于 ACK 与断点续传
topicstring事件所属主题
timestampstring事件时间戳
payload.kindstring事件类型枚举
payloadobject事件具体内容

command 帧

{
  "type": "command",
  "id": "cmd_001",
  "agent_id": "agent_dispatch_01",
  "command": "set_generation",
  "params": {"bus": 2, "mw": 50.0},
  "priority": "normal",
  "timeout_ms": 5000
}
字段类型必填默认值说明
idstring指令 ID,用于关联结果
agent_idstring目标 Agent
commandstring指令名称
paramsobject{}指令参数
prioritystringnormallow/normal/high/emergency
timeout_msint5000超时时间

command_result 帧

{
  "type": "command_result",
  "id": "cmd_001",
  "status": "success",
  "result": {"bus": 2, "mw_set": 50.0, "audit_id": "aud_001"},
  "duration_ms": 12,
  "timestamp": "2026-07-06T08:30:00Z"
}
字段类型必填说明
idstring关联的指令 ID
statusstringsuccess / failed / timeout
resultobject执行结果
duration_msint执行耗时
errorobject失败时的错误信息

error 帧

{
  "type": "error",
  "code": "TOPIC_DENIED",
  "message": "No permission to subscribe: audit.*",
  "details": {"topic": "audit.*"},
  "timestamp": "2026-07-06T08:30:00Z"
}

事件主题列表

主题前缀说明通配符示例
network.{id}电网拓扑/状态变更支持 network.*network.net_8f3a2b
agent.{id}Agent 状态变更支持 agent.*agent.dispatch_01
alarm.{severity}告警按级别支持 alarm.* / alarm.criticalalarm.critical
powerflow.{id}潮流计算结果支持 powerflow.*powerflow.net_8f3a2b
device.{id}设备上下线支持 device.*device.dev_001
constraint.*约束越限事件仅通配constraint.*
audit.*审计事件(需审计权限)仅通配audit.*
task.{id}任务状态变更支持 task.*task.pf_20260706_001

通配符规则:

  • * 匹配单层,如 agent.* 匹配 agent.dispatch_01 但不匹配 agent.dispatch_01.sub
  • # 匹配多层(仅 MQTT 兼容模式)

事件类型枚举

NetworkEvent

事件 kind触发条件payload 关键字段
TopologyChanged母线/支路拓扑变更changes, network_id
NetworkCreated电网创建network_id, name
NetworkUpdated电网属性更新network_id, fields
NetworkDeleted电网删除network_id
NetworkLocked网络被锁定(计算中)network_id, lock_holder
NetworkUnlocked网络解锁network_id

PowerFlowEvent

事件 kind触发条件payload 关键字段
PowerFlowStarted潮流计算开始task_id, network_id, method
PowerFlowIteration每次迭代task_id, iteration, mismatch
PowerFlowCompleted潮流计算收敛task_id, iterations, duration_ms
PowerFlowDiverged潮流不收敛task_id, reason

AgentEvent

事件 kind触发条件payload 关键字段
AgentStatusChangedAgent 启停/切换状态agent_id, old_status, new_status
AgentActionAgent 执行动作agent_id, command, params
AgentErrorAgent 内部错误agent_id, error
AgentHeartbeatAgent 心跳agent_id, metrics

AlarmEvent

事件 kind触发条件payload 关键字段
AlarmRaised告警产生alarm_id, severity, message
AlarmAcknowledged告警被确认alarm_id, acknowledged_by
AlarmCleared告警清除alarm_id, reason

DeviceEvent

事件 kind触发条件payload 关键字段
DeviceOnline设备上线device_id, metadata
DeviceOffline设备离线device_id, reason
DeviceMeasurement设备测量上报device_id, metric, value

ConstraintEvent

事件 kind触发条件payload 关键字段
ConstraintViolation触发安全约束越限element, quantity, value, limit
ConstraintRestored约束恢复正常element, quantity

心跳机制

服务端每 30 秒发送 WebSocket 协议级 ping 帧,客户端需在 60 秒内回复 pong 帧。若客户端使用应用层心跳(如浏览器无法控制协议帧),可发送 ping 消息:

{"type": "ping", "timestamp": "2026-07-06T08:30:00Z"}

服务端响应:

{"type": "pong", "timestamp": "2026-07-06T08:30:00.001Z"}

重连机制

客户端断开后应实现指数退避重连:

class EnerOSWebSocket {
  constructor(url, token) {
    this.url = url;
    this.token = token;
    this.reconnectAttempts = 0;
    this.maxReconnectAttempts = 10;
    this.lastEventId = null;
    this.subscriptions = new Set();
    this.connect();
  }

  connect() {
    this.ws = new WebSocket(this.url, ['eneros.v1']);
    this.ws.onopen = () => this.onOpen();
    this.ws.onmessage = (e) => this.onMessage(e);
    this.ws.onclose = (e) => this.onClose(e);
    this.ws.onerror = (e) => console.error('WS error', e);
  }

  onOpen() {
    this.reconnectAttempts = 0;
    this.ws.send(JSON.stringify({
      type: 'auth',
      token: this.token,
    }));
  }

  onMessage(event) {
    const msg = JSON.parse(event.data);
    if (msg.type === 'auth_ok') {
      if (this.subscriptions.size > 0) {
        this.ws.send(JSON.stringify({
          type: 'subscribe',
          topics: Array.from(this.subscriptions),
          last_event_id: this.lastEventId,
        }));
      }
    } else if (msg.type === 'event') {
      this.lastEventId = msg.id;
      this.handleEvent(msg);
    }
  }

  onClose(event) {
    if (this.reconnectAttempts >= this.maxReconnectAttempts) {
      console.error('达到最大重连次数,放弃重连');
      return;
    }
    const delay = Math.min(1000 * Math.pow(2, this.reconnectAttempts), 30000);
    this.reconnectAttempts++;
    console.log(`${delay}ms 后重连(第 ${this.reconnectAttempts} 次)`);
    setTimeout(() => this.connect(), delay);
  }

  handleEvent(msg) {
    console.log('事件:', msg);
  }

  subscribe(topic) {
    this.subscriptions.add(topic);
    if (this.ws.readyState === WebSocket.OPEN) {
      this.ws.send(JSON.stringify({type: 'subscribe', topics: [topic]}));
    }
  }

  unsubscribe(topic) {
    this.subscriptions.delete(topic);
    if (this.ws.readyState === WebSocket.OPEN) {
      this.ws.send(JSON.stringify({type: 'unsubscribe', topics: [topic]}));
    }
  }
}

Rust 客户端示例

use tokio_tungstenite::{connect_async, tungstenite::Message};
use futures_util::{SinkExt, StreamExt};
use serde_json::json;

#[tokio::main]
async fn main() -> anyhow::Result<()> {
    let (ws_stream, _) = connect_async("ws://localhost:8080/ws").await?;
    let (mut write, mut read) = ws_stream.split();

    // 鉴权
    let auth_msg = json!({
        "type": "auth",
        "token": "<bearer-token>",
        "client_id": "rust-client-01",
    });
    write.send(Message::Text(auth_msg.to_string())).await?;

    // 等待 auth_ok
    if let Some(Ok(Message::Text(text))) = read.next().await {
        let resp: serde_json::Value = serde_json::from_str(&text)?;
        if resp["type"] == "auth_ok" {
            println!("鉴权成功");
        }
    }

    // 订阅
    let sub_msg = json!({
        "type": "subscribe",
        "topics": ["network.net_8f3a2b", "alarm.critical"],
    });
    write.send(Message::Text(sub_msg.to_string())).await?;

    // 接收事件循环
    while let Some(Ok(msg)) = read.next().await {
        match msg {
            Message::Text(text) => {
                let evt: serde_json::Value = serde_json::from_str(&text)?;
                if evt["type"] == "event" {
                    println!(
                        "[{}] {}: {:?}",
                        evt["timestamp"], evt["topic"], evt["payload"]
                    );
                }
            }
            Message::Ping(data) => {
                write.send(Message::Pong(data)).await?;
            }
            _ => {}
        }
    }

    Ok(())
}

Python 客户端示例

import asyncio
import json
import websockets

async def main():
    uri = "ws://localhost:8080/ws"
    async with websockets.connect(uri, subprotocols=["eneros.v1"]) as ws:
        # 鉴权
        await ws.send(json.dumps({
            "type": "auth",
            "token": "<bearer-token>",
            "client_id": "python-client-01",
        }))

        # 等待 auth_ok
        response = json.loads(await ws.recv())
        if response["type"] != "auth_ok":
            print("鉴权失败:", response)
            return

        print("鉴权成功")

        # 订阅
        await ws.send(json.dumps({
            "type": "subscribe",
            "topics": ["network.net_8f3a2b", "alarm.critical"],
        }))

        # 事件循环
        async for message in ws:
            evt = json.loads(message)
            if evt["type"] == "event":
                print(f"[{evt['timestamp']}] {evt['topic']}: {evt['payload']}")
            elif evt["type"] == "ping":
                await ws.send(json.dumps({"type": "pong"}))

asyncio.run(main())

ACK 与可靠传输

启用 ACK 后,客户端需对每个事件发送 ack 帧:

{
  "type": "ack",
  "event_id": "evt_a1b2c3"
}

若 5 秒内未收到 ACK,服务端将重发该事件(最多 3 次)。ACK 模式适用于关键告警与控制指令场景。

错误码

错误码含义处理建议
AUTH_FAILEDtoken 无效或过期重新获取 token
AUTH_TIMEOUT5 秒内未发送 auth 帧立即发送鉴权帧
TOPIC_DENIED无订阅权限联系管理员调整策略
TOPIC_LIMIT_EXCEEDED超过单连接最大订阅数取消部分订阅
RATE_LIMITED发送频率超限降低下发频率
INVALID_MESSAGE消息格式错误校验 JSON 结构
MESSAGE_TOO_LARGE消息超过 256KB拆分消息
AGENT_NOT_FOUND指令目标 Agent 不存在核对 agent_id
AGENT_BUSYAgent 忙,无法处理指令排队或稍后重试
COMMAND_TIMEOUT指令执行超时增大 timeout 或检查 Agent
SESSION_EXPIRED会话过期重新建立连接
SERVER_SHUTDOWN服务端即将关闭立即重连至其他节点

关闭码

WebSocket 关闭码遵循 RFC 6455,并扩展以下自定义关闭码:

关闭码含义处理建议
1000正常关闭
1006异常关闭(无关闭帧)检查网络
1011服务端内部错误重连
4001鉴权失败重新鉴权
4002心跳超时检查网络
4003订阅数超限取消部分订阅
4004限流退避后重连
4005服务端维护等待 5 分钟后重连

性能与限制

维度限制
单连接最大订阅50 个主题
单连接消息速率100 msg/s(入站)
单连接事件速率1000 msg/s(出站)
消息大小上限256 KB
单租户最大连接20
连接空闲超时5 分钟(无任何消息)
事件保留时长24 小时(用于断点续传)

相关文档