()
| 386 | }; |
| 387 | |
| 388 | const flushBuffer = () => { |
| 389 | if (!buffer.trim()) { |
| 390 | buffer = ''; |
| 391 | return; |
| 392 | } |
| 393 | |
| 394 | if (isSSE) { |
| 395 | const events = parseSSEEvents(buffer + '\n\n'); |
| 396 | for (const event of events) { |
| 397 | emitMessage(event.event, event.data, event.id); |
| 398 | } |
| 399 | } else if (isNDJSON) { |
| 400 | const lines = buffer.split(/\r?\n/).filter(line => line.trim()); |
| 401 | for (const line of lines) { |
| 402 | emitMessage('message', line); |
| 403 | } |
| 404 | } |
| 405 | buffer = ''; |
| 406 | }; |
| 407 | |
| 408 | const shouldObserveStream = (stream) => { |
| 409 | if (!activeStream) { |
no test coverage detected