CMVR-IOT-UI/src/plugins/websocket.js
2025-08-21 14:11:58 +08:00

183 lines
5.3 KiB
JavaScript

import { fa } from 'element-plus/es/locale/index.mjs';
import { throttle } from 'lodash';
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.maxQueueSize = 1;
this.status = false;
}
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.ws.onmessage = (event) => {
if (this.messageQueue.length < this.maxQueueSize) {
this.messageQueue.push(event);
this.processMessageQueue();
} else {
console.warn('消息队列已满,丢弃消息');
}
};
this.ws.onclose = () => {
console.log('WebSocket 断开');
this.status = false;
this.stopHeartbeat();
if (!this.isManualClose) {
this.reconnect();
}
};
this.ws.onerror = (err) => {
console.error('WebSocket 错误:', err);
};
}
async processMessageQueue() {
if (!this.ws || this.messageQueue.length === 0) return;
const event = this.messageQueue.shift();
try {
// 处理消息数据
console.log('收到消息数据类型:', event.data.constructor.name);
let compressedData;
// 处理 Blob 或 ArrayBuffer
if (event.data instanceof Blob) {
console.log('处理 Blob 数据,大小:', event.data.size);
compressedData = new Uint8Array(await event.data.arrayBuffer());
} else if (event.data instanceof ArrayBuffer) {
console.log('处理 ArrayBuffer 数据,大小:', event.data.byteLength);
compressedData = new Uint8Array(event.data);
} else {
throw new Error('不支持的数据类型: ' + event.data.constructor.name);
}
// 使用 DecompressionStream 解压 GZIP 数据
const decompressionStream = new DecompressionStream('gzip');
const writer = decompressionStream.writable.getWriter();
writer.write(compressedData);
writer.close();
const decompressedStream = decompressionStream.readable;
const reader = decompressedStream.getReader();
let decompressedData = new Uint8Array(0);
while (true) {
const { done, value } = await reader.read();
if (done) break;
const newData = new Uint8Array(decompressedData.length + value.length);
newData.set(decompressedData);
newData.set(value, decompressedData.length);
decompressedData = newData;
}
// 将解压后的数据转换为字符串
const textDecoder = new TextDecoder();
const decompressedString = textDecoder.decode(decompressedData);
console.log('解压后的数据:', JSON.parse(decompressedString));
// 解析 JSON
const { type, channel, payload,timestamp } = JSON.parse(decompressedString);
console.log('后端到前端延迟:', ((new Date()).getTime() - timestamp));
if (channel && this.events[channel]) {
this.events[channel].forEach(callback => callback(payload));
} else if (type && this.events[type]) {
this.events[type].forEach(callback => callback(payload));
}
setTimeout(() => this.processMessageQueue(), 0);
} catch (err) {
console.error('WebSocket 消息解析失败:', err);
console.error('错误堆栈:', err.stack);
console.error('原始数据类型:', event.data.constructor.name);
console.error('原始数据大小:', event.data instanceof Blob ? event.data.size : event.data.byteLength || '未知');
}
}
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 (this.ws && this.ws.readyState === 1) {
this.ws.send(JSON.stringify(payload));
} else {
console.warn('WebSocket 未连接,消息发送失败');
}
}
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);
}
};