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

Function createStream

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

Source from the content-addressed store, hash-verified

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