| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248 |
- // 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();
- }
- };
- }
|