sse.js 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149
  1. // sse.js —— 跨端 SSE 连接
  2. // H5 / 浏览器:直接用原生 EventSource
  3. // 小程序(无 EventSource):用 uni.request 的 enableChunked 接收分块,手动解析 SSE 协议
  4. // ---- UTF-8 解码器:支持跨分片的多字节字符(中文、emoji),避免截断乱码 ----
  5. function createUtf8Decoder() {
  6. let pending = [];
  7. return function decode(bytes) {
  8. const all = pending.length ? pending.concat(Array.from(bytes)) : Array.from(bytes);
  9. let out = '';
  10. let i = 0;
  11. const len = all.length;
  12. while (i < len) {
  13. const c = all[i];
  14. if (c < 0x80) {
  15. out += String.fromCharCode(c);
  16. i++;
  17. } else if (c >> 5 === 0x06) {
  18. // 2 字节
  19. if (i + 1 < len) {
  20. out += String.fromCharCode(((c & 0x1f) << 6) | (all[i + 1] & 0x3f));
  21. i += 2;
  22. } else { pending = all.slice(i); return out; }
  23. } else if (c >> 4 === 0x0e) {
  24. // 3 字节
  25. if (i + 2 < len) {
  26. out += String.fromCharCode(((c & 0x0f) << 12) | ((all[i + 1] & 0x3f) << 6) | (all[i + 2] & 0x3f));
  27. i += 3;
  28. } else { pending = all.slice(i); return out; }
  29. } else if (c >> 3 === 0x1e) {
  30. // 4 字节(代理对)
  31. if (i + 3 < len) {
  32. let cp = ((c & 0x07) << 18) | ((all[i + 1] & 0x3f) << 12) | ((all[i + 2] & 0x3f) << 6) | (all[i + 3] & 0x3f);
  33. cp -= 0x10000;
  34. out += String.fromCharCode(0xd800 + (cp >> 10)) + String.fromCharCode(0xdc00 + (cp & 0x3ff));
  35. i += 4;
  36. } else { pending = all.slice(i); return out; }
  37. } else {
  38. i++; // 非法字节,跳过
  39. }
  40. }
  41. pending = [];
  42. return out;
  43. };
  44. }
  45. // ---- SSE 协议解析:以空行(\n\n)分隔事件,提取 data: 字段 ----
  46. function createSseParser(onData) {
  47. let buffer = '';
  48. return function feed(text) {
  49. buffer += text.replace(/\r\n/g, '\n').replace(/\r/g, '\n');
  50. let idx;
  51. while ((idx = buffer.indexOf('\n\n')) !== -1) {
  52. const raw = buffer.slice(0, idx);
  53. buffer = buffer.slice(idx + 2);
  54. const lines = raw.split('\n');
  55. let dataStr = '';
  56. for (const line of lines) {
  57. if (line.indexOf('data:') === 0) {
  58. dataStr += line.slice(5).replace(/^ /, '');
  59. }
  60. }
  61. if (dataStr !== '') {
  62. onData(dataStr);
  63. }
  64. }
  65. };
  66. }
  67. // H5:原生 EventSource
  68. function connectWithEventSource(url, handlers) {
  69. const eventSource = new EventSource(url);
  70. eventSource.addEventListener('message', (e) => {
  71. try {
  72. handlers.onMessage(e.data.replace(/\\n/g, '\n'));
  73. } catch (err) {
  74. console.error('Error parsing SSE data:', err);
  75. }
  76. });
  77. eventSource.addEventListener('error', (e) => {
  78. handlers.onError(e);
  79. });
  80. return {
  81. close() {
  82. try { eventSource.close(); } catch (e) {}
  83. }
  84. };
  85. }
  86. // 小程序:uni.request 分块流式接收
  87. function connectWithUniRequest(url, handlers) {
  88. const decoder = createUtf8Decoder();
  89. const feed = createSseParser((dataStr) => {
  90. try {
  91. handlers.onMessage(dataStr.replace(/\\n/g, '\n'));
  92. } catch (err) {
  93. console.error('Error parsing SSE data:', err);
  94. }
  95. });
  96. let requestTask = null;
  97. let closed = false;
  98. requestTask = uni.request({
  99. url,
  100. method: 'GET',
  101. enableChunked: true,
  102. responseType: 'text',
  103. success: () => {},
  104. fail: (err) => {
  105. if (!closed) handlers.onError(err);
  106. }
  107. });
  108. if (requestTask && typeof requestTask.onChunkReceived === 'function') {
  109. requestTask.onChunkReceived((res) => {
  110. const bytes = new Uint8Array(res.data);
  111. const text = decoder(bytes);
  112. if (text) feed(text);
  113. });
  114. } else {
  115. handlers.onError(new Error('当前环境不支持流式接收(enableChunked / onChunkReceived)'));
  116. }
  117. return {
  118. close() {
  119. closed = true;
  120. try { requestTask && requestTask.abort && requestTask.abort(); } catch (e) {}
  121. }
  122. };
  123. }
  124. export function createSSEConnection(url, options = {}) {
  125. const handlers = { onMessage: () => {}, onError: () => {} };
  126. const conn = (typeof EventSource !== 'undefined')
  127. ? connectWithEventSource(url, handlers)
  128. : connectWithUniRequest(url, handlers);
  129. return {
  130. onMessage(callback) {
  131. handlers.onMessage = callback || (() => {});
  132. },
  133. onError(callback) {
  134. handlers.onError = callback || (() => {});
  135. },
  136. close() {
  137. conn.close();
  138. }
  139. };
  140. }