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); } };