// sse.js —— 跨端 SSE 连接 // H5 / 浏览器:直接用原生 EventSource // 小程序(无 EventSource):用 uni.request 的 enableChunked 接收分块,手动解析 SSE 协议 // ---- UTF-8 解码器:支持跨分片的多字节字符(中文、emoji),避免截断乱码 ---- function createUtf8Decoder() { let pending = []; return function decode(bytes) { const all = pending.length ? pending.concat(Array.from(bytes)) : Array.from(bytes); let out = ''; let i = 0; const len = all.length; while (i < len) { const c = all[i]; if (c < 0x80) { out += String.fromCharCode(c); i++; } else if (c >> 5 === 0x06) { // 2 字节 if (i + 1 < len) { out += String.fromCharCode(((c & 0x1f) << 6) | (all[i + 1] & 0x3f)); i += 2; } else { pending = all.slice(i); return out; } } else if (c >> 4 === 0x0e) { // 3 字节 if (i + 2 < len) { out += String.fromCharCode(((c & 0x0f) << 12) | ((all[i + 1] & 0x3f) << 6) | (all[i + 2] & 0x3f)); i += 3; } else { pending = all.slice(i); return out; } } else if (c >> 3 === 0x1e) { // 4 字节(代理对) if (i + 3 < len) { let cp = ((c & 0x07) << 18) | ((all[i + 1] & 0x3f) << 12) | ((all[i + 2] & 0x3f) << 6) | (all[i + 3] & 0x3f); cp -= 0x10000; out += String.fromCharCode(0xd800 + (cp >> 10)) + String.fromCharCode(0xdc00 + (cp & 0x3ff)); i += 4; } else { pending = all.slice(i); return out; } } else { i++; // 非法字节,跳过 } } pending = []; return out; }; } // ---- SSE 协议解析:以空行(\n\n)分隔事件,提取 data: 字段 ---- function createSseParser(onData) { let buffer = ''; return function feed(text) { buffer += text.replace(/\r\n/g, '\n').replace(/\r/g, '\n'); let idx; while ((idx = buffer.indexOf('\n\n')) !== -1) { const raw = buffer.slice(0, idx); buffer = buffer.slice(idx + 2); const lines = raw.split('\n'); let dataStr = ''; for (const line of lines) { if (line.indexOf('data:') === 0) { dataStr += line.slice(5).replace(/^ /, ''); } } if (dataStr !== '') { onData(dataStr); } } }; } // H5:原生 EventSource function connectWithEventSource(url, handlers) { const eventSource = new EventSource(url); eventSource.addEventListener('message', (e) => { try { handlers.onMessage(e.data.replace(/\\n/g, '\n')); } catch (err) { console.error('Error parsing SSE data:', err); } }); eventSource.addEventListener('error', (e) => { handlers.onError(e); }); return { close() { try { eventSource.close(); } catch (e) {} } }; } // 小程序:uni.request 分块流式接收 function connectWithUniRequest(url, handlers) { const decoder = createUtf8Decoder(); const feed = createSseParser((dataStr) => { try { handlers.onMessage(dataStr.replace(/\\n/g, '\n')); } catch (err) { console.error('Error parsing SSE data:', err); } }); let requestTask = null; let closed = false; requestTask = uni.request({ url, method: 'GET', enableChunked: true, responseType: 'text', success: () => {}, fail: (err) => { if (!closed) handlers.onError(err); } }); if (requestTask && typeof requestTask.onChunkReceived === 'function') { requestTask.onChunkReceived((res) => { const bytes = new Uint8Array(res.data); const text = decoder(bytes); if (text) feed(text); }); } else { handlers.onError(new Error('当前环境不支持流式接收(enableChunked / onChunkReceived)')); } return { close() { closed = true; try { requestTask && requestTask.abort && requestTask.abort(); } catch (e) {} } }; } export function createSSEConnection(url, options = {}) { const handlers = { onMessage: () => {}, onError: () => {} }; const conn = (typeof EventSource !== 'undefined') ? connectWithEventSource(url, handlers) : connectWithUniRequest(url, handlers); return { onMessage(callback) { handlers.onMessage = callback || (() => {}); }, onError(callback) { handlers.onError = callback || (() => {}); }, close() { conn.close(); } }; }