sse.js 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248
  1. // sse.js —— 跨端 SSE 连接
  2. // H5 / 浏览器:原生 EventSource
  3. // App(5+):plus.net.XMLHttpRequest 流式读取 responseText 增量
  4. // (注意:App 端不支持 uni.request 的 enableChunked / onChunkReceived)
  5. // 小程序:uni.request 的 enableChunked 分块接收
  6. // ---- UTF-8 解码器:支持跨分片的多字节字符(中文、emoji),避免截断乱码 ----
  7. // 仅用于小程序 onChunkReceived(拿到的是 ArrayBuffer 字节流)
  8. function createUtf8Decoder() {
  9. let pending = [];
  10. return function decode(bytes) {
  11. const all = pending.length ? pending.concat(Array.from(bytes)) : Array.from(bytes);
  12. let out = '';
  13. let i = 0;
  14. const len = all.length;
  15. while (i < len) {
  16. const c = all[i];
  17. if (c < 0x80) {
  18. out += String.fromCharCode(c);
  19. i++;
  20. } else if (c >> 5 === 0x06) {
  21. // 2 字节
  22. if (i + 1 < len) {
  23. out += String.fromCharCode(((c & 0x1f) << 6) | (all[i + 1] & 0x3f));
  24. i += 2;
  25. } else { pending = all.slice(i); return out; }
  26. } else if (c >> 4 === 0x0e) {
  27. // 3 字节
  28. if (i + 2 < len) {
  29. out += String.fromCharCode(((c & 0x0f) << 12) | ((all[i + 1] & 0x3f) << 6) | (all[i + 2] & 0x3f));
  30. i += 3;
  31. } else { pending = all.slice(i); return out; }
  32. } else if (c >> 3 === 0x1e) {
  33. // 4 字节(代理对)
  34. if (i + 3 < len) {
  35. let cp = ((c & 0x07) << 18) | ((all[i + 1] & 0x3f) << 12) | ((all[i + 2] & 0x3f) << 6) | (all[i + 3] & 0x3f);
  36. cp -= 0x10000;
  37. out += String.fromCharCode(0xd800 + (cp >> 10)) + String.fromCharCode(0xdc00 + (cp & 0x3ff));
  38. i += 4;
  39. } else { pending = all.slice(i); return out; }
  40. } else {
  41. i++; // 非法字节,跳过
  42. }
  43. }
  44. pending = [];
  45. return out;
  46. };
  47. }
  48. // ---- SSE 协议解析:以空行(\n\n)分隔事件,提取 data: 字段 ----
  49. function createSseParser(onData) {
  50. let buffer = '';
  51. return function feed(text) {
  52. buffer += text.replace(/\r\n/g, '\n').replace(/\r/g, '\n');
  53. let idx;
  54. while ((idx = buffer.indexOf('\n\n')) !== -1) {
  55. const raw = buffer.slice(0, idx);
  56. buffer = buffer.slice(idx + 2);
  57. const lines = raw.split('\n');
  58. let dataStr = '';
  59. for (const line of lines) {
  60. if (line.indexOf('data:') === 0) {
  61. dataStr += line.slice(5).replace(/^ /, '');
  62. }
  63. }
  64. if (dataStr !== '') {
  65. onData(dataStr);
  66. }
  67. }
  68. };
  69. }
  70. // 把已转义的 \n 还原(服务端常把换行转义成字面量 "\\n")
  71. function unescapeData(dataStr) {
  72. return dataStr.replace(/\\n/g, '\n');
  73. }
  74. // H5:原生 EventSource(浏览器自带 SSE 解析,直接拿 data)
  75. function connectWithEventSource(url, handlers) {
  76. const eventSource = new EventSource(url);
  77. eventSource.addEventListener('message', (e) => {
  78. try {
  79. handlers._gotData = true;
  80. handlers.onMessage(unescapeData(e.data));
  81. } catch (err) {
  82. console.error('Error parsing SSE data:', err);
  83. }
  84. });
  85. eventSource.addEventListener('error', (e) => {
  86. // EventSource 在服务端正常关闭连接时也会触发 error,区分一下:
  87. // 已经收到过数据视为正常结束,没收到过才当作真正的连接失败
  88. if (handlers._gotData) {
  89. if (typeof handlers.onClose === 'function') handlers.onClose();
  90. } else {
  91. handlers.onError(e);
  92. }
  93. });
  94. return {
  95. close() {
  96. try { eventSource.close(); } catch (e) {}
  97. }
  98. };
  99. }
  100. // App(5+):plus.net.XMLHttpRequest 流式
  101. // 通过 onreadystatechange 在 readyState=3/4 时读取 responseText 增量,
  102. // 某些机型不实时推流时,会在 readyState=4 一次性吐出完整内容(降级但可用)。
  103. function connectWithPlusXhr(url, handlers) {
  104. const feed = createSseParser((dataStr) => {
  105. try {
  106. handlers._gotData = true;
  107. handlers.onMessage(unescapeData(dataStr));
  108. } catch (err) {
  109. console.error('Error parsing SSE data:', err);
  110. }
  111. });
  112. let xhr = null;
  113. let processed = 0;
  114. let closed = false;
  115. try {
  116. xhr = new plus.net.XMLHttpRequest();
  117. } catch (e) {
  118. handlers.onError(new Error('当前 App 环境不支持 plus.net.XMLHttpRequest'));
  119. return { close() {} };
  120. }
  121. xhr.open('GET', url);
  122. try { xhr.setRequestHeader('Accept', 'text/event-stream'); } catch (e) {}
  123. xhr.onreadystatechange = () => {
  124. try {
  125. // readyState 3(LOADING) / 4(DONE) 都读取增量
  126. if (xhr.readyState === 3 || xhr.readyState === 4) {
  127. const full = xhr.responseText || '';
  128. if (full.length > processed) {
  129. const chunk = full.substring(processed);
  130. processed = full.length;
  131. feed(chunk);
  132. }
  133. }
  134. if (xhr.readyState === 4) {
  135. if (typeof handlers.onClose === 'function') handlers.onClose();
  136. }
  137. } catch (err) {
  138. console.error('plus XHR onreadystatechange error:', err);
  139. }
  140. };
  141. xhr.onerror = (e) => {
  142. if (!closed) handlers.onError(e);
  143. };
  144. xhr.ontimeout = () => {
  145. if (!closed) handlers.onError(new Error('请求超时'));
  146. };
  147. try {
  148. xhr.send();
  149. } catch (e) {
  150. handlers.onError(e);
  151. }
  152. return {
  153. close() {
  154. closed = true;
  155. try { xhr && xhr.abort(); } catch (e) {}
  156. }
  157. };
  158. }
  159. // 小程序:uni.request 分块流式接收
  160. function connectWithUniRequest(url, handlers) {
  161. const decoder = createUtf8Decoder();
  162. const feed = createSseParser((dataStr) => {
  163. try {
  164. handlers._gotData = true;
  165. handlers.onMessage(unescapeData(dataStr));
  166. } catch (err) {
  167. console.error('Error parsing SSE data:', err);
  168. }
  169. });
  170. let requestTask = null;
  171. let closed = false;
  172. requestTask = uni.request({
  173. url,
  174. method: 'GET',
  175. enableChunked: true,
  176. responseType: 'text',
  177. success: () => {
  178. if (typeof handlers.onClose === 'function') handlers.onClose();
  179. },
  180. fail: (err) => {
  181. if (!closed) handlers.onError(err);
  182. }
  183. });
  184. if (requestTask && typeof requestTask.onChunkReceived === 'function') {
  185. requestTask.onChunkReceived((res) => {
  186. const bytes = new Uint8Array(res.data);
  187. const text = decoder(bytes);
  188. if (text) feed(text);
  189. });
  190. } else {
  191. handlers.onError(new Error('当前环境不支持流式接收(enableChunked / onChunkReceived)'));
  192. }
  193. return {
  194. close() {
  195. closed = true;
  196. try { requestTask && requestTask.abort && requestTask.abort(); } catch (e) {}
  197. }
  198. };
  199. }
  200. // 统一入口:按环境选择连接方式
  201. export function createSSEConnection(url, options = {}) {
  202. const handlers = {
  203. onMessage: () => {},
  204. onError: () => {},
  205. onClose: () => {},
  206. _gotData: false
  207. };
  208. let conn;
  209. if (typeof EventSource !== 'undefined') {
  210. conn = connectWithEventSource(url, handlers); // H5
  211. } else if (typeof plus !== 'undefined' && plus.net && plus.net.XMLHttpRequest) {
  212. conn = connectWithPlusXhr(url, handlers); // App(5+)
  213. } else {
  214. conn = connectWithUniRequest(url, handlers); // 小程序
  215. }
  216. return {
  217. onMessage(callback) {
  218. handlers.onMessage = callback || (() => {});
  219. },
  220. onError(callback) {
  221. handlers.onError = callback || (() => {});
  222. },
  223. onClose(callback) {
  224. handlers.onClose = callback || (() => {});
  225. },
  226. close() {
  227. conn.close();
  228. }
  229. };
  230. }