(response, runtime, context)
| 367 | // ======================================== |
| 368 | |
| 369 | function createStream(response, runtime, context) { |
| 370 | var element = context.me; |
| 371 | var reader = response.body.getReader(); |
| 372 | var messages = []; |
| 373 | var waiting = null; |
| 374 | var done = false; |
| 375 | |
| 376 | // Read SSE events in the background, dispatch named events on the |
| 377 | // element, and buffer unnamed messages for async iteration. |
| 378 | (async function () { |
| 379 | try { |
| 380 | for await (var msg of parseSSE(reader)) { |
| 381 | var eventType = msg.event || 'message'; |
| 382 | if (msg.event) { |
| 383 | runtime.triggerEvent(element, eventType, { |
| 384 | data: msg.data, |
| 385 | lastEventId: msg.id || '' |
| 386 | }); |
| 387 | } else { |
| 388 | messages.push(msg.data); |
| 389 | if (waiting) { |
| 390 | waiting.resolve({ value: msg.data, done: false }); |
| 391 | waiting = null; |
| 392 | } |
| 393 | } |
| 394 | } |
| 395 | } catch (err) { |
| 396 | runtime.triggerEvent(element, 'stream-error', { error: err }); |
| 397 | } |
| 398 | done = true; |
| 399 | if (waiting) { |
| 400 | waiting.resolve({ value: undefined, done: true }); |
| 401 | waiting = null; |
| 402 | } |
| 403 | runtime.triggerEvent(element, 'streamEnd', {}); |
| 404 | })(); |
| 405 | |
| 406 | var stream = { |
| 407 | element: element, |
| 408 | [Symbol.asyncIterator]: function () { |
| 409 | var index = 0; |
| 410 | return { |
| 411 | next: function () { |
| 412 | if (index < messages.length) { |
| 413 | return Promise.resolve({ value: messages[index++], done: false }); |
| 414 | } |
| 415 | if (done) { |
| 416 | return Promise.resolve({ value: undefined, done: true }); |
| 417 | } |
| 418 | return new Promise(function (resolve) { |
| 419 | waiting = { resolve: resolve }; |
| 420 | }).then(function (result) { |
| 421 | if (!result.done) index++; |
| 422 | return result; |
| 423 | }); |
| 424 | } |
| 425 | }; |
| 426 | } |
no test coverage detected