// sse.js —— 跨端 SSE 连接 // H5 / 浏览器:原生 EventSource // App(5+):plus.net.XMLHttpRequest 流式读取 responseText 增量 // (注意:App 端不支持 uni.request 的 enableChunked / onChunkReceived) // 小程序:uni.request 的 enableChunked 分块接收 // ---- UTF-8 解码器:支持跨分片的多字节字符(中文、emoji),避免截断乱码 ---- // 仅用于小程序 onChunkReceived(拿到的是 ArrayBuffer 字节流) 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); } } }; } // 把已转义的 \n 还原(服务端常把换行转义成字面量 "\\n") function unescapeData(dataStr) { return dataStr.replace(/\\n/g, '\n'); } // H5:原生 EventSource(浏览器自带 SSE 解析,直接拿 data) function connectWithEventSource(url, handlers) { const eventSource = new EventSource(url); eventSource.addEventListener('message', (e) => { try { handlers._gotData = true; handlers.onMessage(unescapeData(e.data)); } catch (err) { console.error('Error parsing SSE data:', err); } }); eventSource.addEventListener('error', (e) => { // EventSource 在服务端正常关闭连接时也会触发 error,区分一下: // 已经收到过数据视为正常结束,没收到过才当作真正的连接失败 if (handlers._gotData) { if (typeof handlers.onClose === 'function') handlers.onClose(); } else { handlers.onError(e); } }); return { close() { try { eventSource.close(); } catch (e) {} } }; } // App(5+):plus.net.XMLHttpRequest 流式 // 通过 onreadystatechange 在 readyState=3/4 时读取 responseText 增量, // 某些机型不实时推流时,会在 readyState=4 一次性吐出完整内容(降级但可用)。 function connectWithPlusXhr(url, handlers) { const feed = createSseParser((dataStr) => { try { handlers._gotData = true; handlers.onMessage(unescapeData(dataStr)); } catch (err) { console.error('Error parsing SSE data:', err); } }); let xhr = null; let processed = 0; let closed = false; try { xhr = new plus.net.XMLHttpRequest(); } catch (e) { handlers.onError(new Error('当前 App 环境不支持 plus.net.XMLHttpRequest')); return { close() {} }; } xhr.open('GET', url); try { xhr.setRequestHeader('Accept', 'text/event-stream'); } catch (e) {} xhr.onreadystatechange = () => { try { // readyState 3(LOADING) / 4(DONE) 都读取增量 if (xhr.readyState === 3 || xhr.readyState === 4) { const full = xhr.responseText || ''; if (full.length > processed) { const chunk = full.substring(processed); processed = full.length; feed(chunk); } } if (xhr.readyState === 4) { if (typeof handlers.onClose === 'function') handlers.onClose(); } } catch (err) { console.error('plus XHR onreadystatechange error:', err); } }; xhr.onerror = (e) => { if (!closed) handlers.onError(e); }; xhr.ontimeout = () => { if (!closed) handlers.onError(new Error('请求超时')); }; try { xhr.send(); } catch (e) { handlers.onError(e); } return { close() { closed = true; try { xhr && xhr.abort(); } catch (e) {} } }; } // 小程序:uni.request 分块流式接收 function connectWithUniRequest(url, handlers) { const decoder = createUtf8Decoder(); const feed = createSseParser((dataStr) => { try { handlers._gotData = true; handlers.onMessage(unescapeData(dataStr)); } 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: () => { if (typeof handlers.onClose === 'function') handlers.onClose(); }, 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: () => {}, onClose: () => {}, _gotData: false }; let conn; if (typeof EventSource !== 'undefined') { conn = connectWithEventSource(url, handlers); // H5 } else if (typeof plus !== 'undefined' && plus.net && plus.net.XMLHttpRequest) { conn = connectWithPlusXhr(url, handlers); // App(5+) } else { conn = connectWithUniRequest(url, handlers); // 小程序 } return { onMessage(callback) { handlers.onMessage = callback || (() => {}); }, onError(callback) { handlers.onError = callback || (() => {}); }, onClose(callback) { handlers.onClose = callback || (() => {}); }, close() { conn.close(); } }; }