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