WebSocket 事件流
EnerOS 通过 WebSocket 提供双向实时通信通道,用于 Agent 状态推送、电网事件广播、远程指令下发等场景。底层由 eneros-eventbus 产生事件,经 eneros-gateway 鉴权后推送至订阅客户端。WebSocket 协议基于 HTTP 升级握手,连接建立后可保持长连接,相比 REST 反复建连,可显著降低延迟与资源消耗。
端点信息
| 项目 | 取值 |
|---|
| URL | ws://localhost:8080/ws |
| 生产 URL | wss://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 文本帧(不支持二进制帧),统一字段约定如下:
| 字段 | 类型 | 必填 | 说明 |
|---|
type | string | 是 | 消息类型,见下方枚举 |
id | string | 否 | 消息 ID,用于 ACK 关联 |
timestamp | string | 否 | 消息时间戳(RFC 3339) |
payload | object | 否 | 消息体,结构随 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"]
}
| 字段 | 类型 | 必填 | 默认值 | 说明 |
|---|
token | string | 是 | — | Bearer Token |
client_id | string | 否 | 自动生成 | 客户端标识 |
client_version | string | 否 | "" | 客户端版本 |
capabilities | array | 否 | [] | 客户端能力声明 |
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"
}
| 字段 | 类型 | 必填 | 默认值 | 说明 |
|---|
topics | array | 是 | — | 主题列表,支持通配符 * |
ack | bool | 否 | false | 是否要求事件 ACK |
last_event_id | string | 否 | — | 断点续传,从此 ID 之后开始补发 |
subscribed 帧
{
"type": "subscribed",
"topics": ["network.net_8f3a2b", "agent.*", "alarm.critical"],
"pending_replay": 12
}
| 字段 | 类型 | 必填 | 说明 |
|---|
topics | array | 是 | 已订阅主题列表 |
pending_replay | int | 是 | 待补发的事件数量 |
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"}
}
}
| 字段 | 类型 | 必填 | 说明 |
|---|
id | string | 是 | 事件唯一 ID,用于 ACK 与断点续传 |
topic | string | 是 | 事件所属主题 |
timestamp | string | 是 | 事件时间戳 |
payload.kind | string | 是 | 事件类型枚举 |
payload | object | 是 | 事件具体内容 |
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
}
| 字段 | 类型 | 必填 | 默认值 | 说明 |
|---|
id | string | 是 | — | 指令 ID,用于关联结果 |
agent_id | string | 是 | — | 目标 Agent |
command | string | 是 | — | 指令名称 |
params | object | 否 | {} | 指令参数 |
priority | string | 否 | normal | low/normal/high/emergency |
timeout_ms | int | 否 | 5000 | 超时时间 |
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"
}
| 字段 | 类型 | 必填 | 说明 |
|---|
id | string | 是 | 关联的指令 ID |
status | string | 是 | success / failed / timeout |
result | object | 否 | 执行结果 |
duration_ms | int | 是 | 执行耗时 |
error | object | 否 | 失败时的错误信息 |
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.critical | alarm.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 关键字段 |
|---|
AgentStatusChanged | Agent 启停/切换状态 | agent_id, old_status, new_status |
AgentAction | Agent 执行动作 | agent_id, command, params |
AgentError | Agent 内部错误 | agent_id, error |
AgentHeartbeat | Agent 心跳 | 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_FAILED | token 无效或过期 | 重新获取 token |
AUTH_TIMEOUT | 5 秒内未发送 auth 帧 | 立即发送鉴权帧 |
TOPIC_DENIED | 无订阅权限 | 联系管理员调整策略 |
TOPIC_LIMIT_EXCEEDED | 超过单连接最大订阅数 | 取消部分订阅 |
RATE_LIMITED | 发送频率超限 | 降低下发频率 |
INVALID_MESSAGE | 消息格式错误 | 校验 JSON 结构 |
MESSAGE_TOO_LARGE | 消息超过 256KB | 拆分消息 |
AGENT_NOT_FOUND | 指令目标 Agent 不存在 | 核对 agent_id |
AGENT_BUSY | Agent 忙,无法处理指令 | 排队或稍后重试 |
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 小时(用于断点续传) |
相关文档