CMVR-IOT-UI/src/plugins/websocket.js
lixiaolong fabfa05db5 fix(websocket): 解决WebSocket重连后频道订阅丢失问题
- 移除调试用的console.log语句
- 添加channelSubscriptions集合跟踪频道订阅状态
- 在WebSocket连接建立时恢复之前的频道订阅
- 添加refreshAfterSocketReconnect函数用于重连后刷新数据
- 在组件挂载时监听socket open事件并绑定重连刷新逻辑
- 在组件卸载时移除socket open事件监听器
2026-08-14 13:31:37 +08:00

230 lines
6.6 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

const MEDIA_HEADER_SIZE = 26;
const MEDIA_MAGIC = [0x43, 0x4d, 0x56, 0x46]; // CMVF
const MEDIA_MIME_TYPES = {
1: 'video/mp4; codecs="avc1.42C01E"'
};
class WebSocketManager {
constructor(url, userId) {
this.url = url;
this.userId = userId;
this.ws = null;
this.events = {};
this.reconnectTimer = null;
this.heartbeatTimer = null;
this.isManualClose = false;
this.messageQueue = [];
this.processingMessages = false;
this.maxQueueSize = 100;
this.status = false;
this.channelSubscriptions = new Set();
}
connect() {
if (!window.WebSocket) {
console.error('浏览器不支持 WebSocket');
return;
}
this.ws = new WebSocket(`${this.url}?userId=${this.userId}`, [], { perMessageDeflate: true });
this.ws.binaryType = 'arraybuffer'; // 确保接收 ArrayBuffer
this.ws.onopen = () => {
console.log('WebSocket 已连接');
this.status = true;
this.startHeartbeat();
this.processMessageQueue();
this.restoreChannelSubscriptions();
this.dispatch('open');
};
this.ws.onmessage = async (event) => {
const data = event.data instanceof Blob ? await event.data.arrayBuffer() : event.data;
if (data instanceof ArrayBuffer && this.isMediaPacket(data)) {
this.dispatchMediaPacket(data);
return;
}
if (this.messageQueue.length >= this.maxQueueSize) {
console.warn('WebSocket控制消息队列已满,丢弃最旧消息');
this.messageQueue.shift();
}
this.messageQueue.push(data);
this.processMessageQueue();
};
this.ws.onclose = () => {
console.log('WebSocket 断开');
this.status = false;
this.stopHeartbeat();
this.dispatch('close');
if (!this.isManualClose) {
this.reconnect();
}
};
this.ws.onerror = (err) => {
console.error('WebSocket 错误:', err);
};
}
isMediaPacket(buffer) {
if (buffer.byteLength < MEDIA_HEADER_SIZE) return false;
const bytes = new Uint8Array(buffer, 0, MEDIA_MAGIC.length);
return MEDIA_MAGIC.every((value, index) => bytes[index] === value);
}
dispatchMediaPacket(buffer) {
try {
const view = new DataView(buffer);
const version = view.getUint8(4);
const mediaType = view.getUint8(5);
const timestamp = view.getUint32(8) * 0x100000000 + view.getUint32(12);
const sequence = view.getUint32(16);
const channelLength = view.getUint16(20);
const payloadLength = view.getUint32(22);
const payloadOffset = MEDIA_HEADER_SIZE + channelLength;
if (version !== 1 || payloadOffset + payloadLength !== buffer.byteLength) {
throw new Error('媒体数据包头无效');
}
const channel = new TextDecoder().decode(new Uint8Array(buffer, MEDIA_HEADER_SIZE, channelLength));
const mimeType = MEDIA_MIME_TYPES[mediaType];
if (!mimeType) {
throw new Error(`不支持的媒体类型: ${mediaType}`);
}
this.dispatch(channel, {
media: true,
mimeType,
sequence,
timestamp,
data: buffer.slice(payloadOffset, payloadOffset + payloadLength)
});
} catch (error) {
console.error('WebSocket媒体包解析失败:', error);
}
}
async processMessageQueue() {
if (this.processingMessages || !this.ws) return;
this.processingMessages = true;
try {
while (this.messageQueue.length > 0) {
const data = this.messageQueue.shift();
let message;
if (typeof data === 'string') {
message = JSON.parse(data);
} else if (data instanceof ArrayBuffer) {
const decompressedStream = new Blob([data])
.stream()
.pipeThrough(new DecompressionStream('gzip'));
const decompressedData = await new Response(decompressedStream).arrayBuffer();
message = JSON.parse(new TextDecoder().decode(decompressedData));
} else {
throw new Error(`不支持的WebSocket数据类型: ${typeof data}`);
}
if (message.channel) {
this.dispatch(message.channel, message.payload);
} else if (message.type) {
this.dispatch(message.type, message.payload ?? message.content ?? message);
}
}
} catch (err) {
console.error('WebSocket消息解析失败:', err);
} finally {
this.processingMessages = false;
if (this.messageQueue.length > 0) {
queueMicrotask(() => this.processMessageQueue());
}
}
}
dispatch(eventName, payload) {
const callbacks = this.events[eventName];
if (callbacks) {
callbacks.slice().forEach(callback => callback(payload));
}
}
on(eventName, callback) {
if (!this.events[eventName]) {
this.events[eventName] = [];
}
this.events[eventName].push(callback);
}
off(eventName, callback) {
if (this.events[eventName]) {
const index = this.events[eventName].indexOf(callback);
if (index !== -1) {
this.events[eventName].splice(index, 1);
}
}
}
send(payload) {
if (payload?.type === 'channel_subscription' && payload.channel) {
if (payload.action === 'subscribe') {
this.channelSubscriptions.add(payload.channel);
} else if (payload.action === 'unsubscribe') {
this.channelSubscriptions.delete(payload.channel);
}
}
if (this.ws && this.ws.readyState === 1) {
this.ws.send(JSON.stringify(payload));
} else if (payload?.type !== 'channel_subscription') {
console.warn('WebSocket 未连接,消息发送失败');
}
}
restoreChannelSubscriptions() {
if (!this.ws || this.ws.readyState !== 1) return;
this.channelSubscriptions.forEach((channel) => {
this.ws.send(JSON.stringify({
type: 'channel_subscription',
action: 'subscribe',
channel,
}));
});
}
close() {
this.isManualClose = true;
this.status = false;
this.stopHeartbeat();
if (this.ws) {
this.ws.close();
}
}
reconnect() {
clearTimeout(this.reconnectTimer);
this.reconnectTimer = setTimeout(() => {
console.log('尝试重新连接 WebSocket...');
this.connect();
}, 5000);
}
startHeartbeat() {
this.heartbeatTimer = setInterval(() => {
this.send({ type: 'heartbeat' });
}, 30000);
}
stopHeartbeat() {
clearInterval(this.heartbeatTimer);
}
status() {
return this.status;
}
}
export default {
install(app, options) {
const { url, userId } = options;
const wsManager = new WebSocketManager(url, userId);
wsManager.connect();
app.config.globalProperties.$ws = wsManager;
app.provide('ws', wsManager);
}
};