MCPcopy Create free account
hub / github.com/bigskysoftware/_hyperscript / createStream

Function createStream

www/js/ext/eventsource.js:338–394  ·  view source on GitHub ↗
(response, runtime, context)

Source from the content-addressed store, hash-verified

336 return data;
337 }
338 function createStream(response, runtime, context) {
339 var element = context.me;
340 var reader = response.body.getReader();
341 var messages = [];
342 var waiting = null;
343 var done = false;
344 (async function() {
345 try {
346 for await (var msg of parseSSE(reader)) {
347 var eventType = msg.event || "message";
348 if (msg.event) {
349 runtime.triggerEvent(element, eventType, {
350 data: msg.data,
351 lastEventId: msg.id || ""
352 });
353 } else {
354 messages.push(msg.data);
355 if (waiting) {
356 waiting.resolve({ value: msg.data, done: false });
357 waiting = null;
358 }
359 }
360 }
361 } catch (err) {
362 runtime.triggerEvent(element, "stream-error", { error: err });
363 }
364 done = true;
365 if (waiting) {
366 waiting.resolve({ value: void 0, done: true });
367 waiting = null;
368 }
369 runtime.triggerEvent(element, "streamEnd", {});
370 })();
371 var stream = {
372 element,
373 [Symbol.asyncIterator]: function() {
374 var index = 0;
375 return {
376 next: function() {
377 if (index < messages.length) {
378 return Promise.resolve({ value: messages[index++], done: false });
379 }
380 if (done) {
381 return Promise.resolve({ value: void 0, done: true });
382 }
383 return new Promise(function(resolve) {
384 waiting = { resolve };
385 }).then(function(result) {
386 if (!result.done) index++;
387 return result;
388 });
389 }
390 };
391 }
392 };
393 return stream;
394 }
395 var streamConversion = function(response, runtime, context) {

Callers 1

streamConversionFunction · 0.70

Calls 3

parseSSEFunction · 0.70
triggerEventMethod · 0.45
resolveMethod · 0.45

Tested by

no test coverage detected