(response: Response)
| 447 | } |
| 448 | |
| 449 | async function* parseServerSentEvents(response: Response) { |
| 450 | const body = response.body |
| 451 | if (!body) { |
| 452 | throw new APIConnectionError({ message: 'Missing streaming response body.' }) |
| 453 | } |
| 454 | const reader = body.getReader() |
| 455 | const decoder = new TextDecoder() |
| 456 | let buffer = '' |
| 457 | while (true) { |
| 458 | const { done, value } = await reader.read() |
| 459 | if (done) break |
| 460 | buffer += decoder.decode(value, { stream: true }) |
| 461 | while (true) { |
| 462 | const separatorIndex = buffer.indexOf('\n\n') |
| 463 | if (separatorIndex === -1) break |
| 464 | const rawEvent = buffer.slice(0, separatorIndex) |
| 465 | buffer = buffer.slice(separatorIndex + 2) |
| 466 | const data = rawEvent |
| 467 | .split(/\r?\n/) |
| 468 | .filter(line => line.startsWith('data:')) |
| 469 | .map(line => line.slice(5).trimStart()) |
| 470 | .join('\n') |
| 471 | if (!data) continue |
| 472 | if (data === '[DONE]') return |
| 473 | yield JSON.parse(data) as OpenAIChatCompletionChunk |
| 474 | } |
| 475 | } |
| 476 | const trailing = buffer.trim() |
| 477 | if (!trailing) return |
| 478 | const trailingData = trailing |
| 479 | .split(/\r?\n/) |
| 480 | .filter(line => line.startsWith('data:')) |
| 481 | .map(line => line.slice(5).trimStart()) |
| 482 | .join('\n') |
| 483 | if (trailingData && trailingData !== '[DONE]') { |
| 484 | yield JSON.parse(trailingData) as OpenAIChatCompletionChunk |
| 485 | } |
| 486 | } |
| 487 | |
| 488 | function createOpenAIEventStream( |
| 489 | response: Response, |
no test coverage detected