MCPcopy Create free account
hub / github.com/cloudflare/computer / transformModuleEvents

Function transformModuleEvents

packages/computer/src/runtime/runtime.ts:338–411  ·  view source on GitHub ↗
(
  source: ReadableStream<WorkspaceRuntimeEvent>,
  encoding: E | undefined,
)

Source from the content-addressed store, hash-verified

336}
337
338function transformModuleEvents<E extends ExecEncoding>(
339 source: ReadableStream<WorkspaceRuntimeEvent>,
340 encoding: E | undefined,
341): ReadableStream<WorkspaceRuntimeEvent<E>> {
342 if (encoding !== "utf8") {
343 return source as ReadableStream<WorkspaceRuntimeEvent<E>>;
344 }
345 const stdoutDecoder = new TextDecoder();
346 const stderrDecoder = new TextDecoder();
347 let stdoutMeta: { id: string; seq: number } | undefined;
348 let stderrMeta: { id: string; seq: number } | undefined;
349 let lastSeq = 0;
350 const enqueue = (
351 controller: TransformStreamDefaultController<WorkspaceRuntimeEvent<E>>,
352 event: WorkspaceRuntimeEvent<E>,
353 ) => {
354 lastSeq = event.seq;
355 controller.enqueue(event);
356 };
357 const flushPending = (
358 controller: TransformStreamDefaultController<WorkspaceRuntimeEvent<E>>,
359 beforeSeq?: number,
360 ) => {
361 const pending: WorkspaceRuntimeEvent<E>[] = [];
362 const stdout = stdoutDecoder.decode();
363 const stderr = stderrDecoder.decode();
364 if (stdout && stdoutMeta) {
365 pending.push({
366 ...stdoutMeta,
367 name: "stdout",
368 value: stdout,
369 } as WorkspaceRuntimeEvent<E>);
370 }
371 if (stderr && stderrMeta) {
372 pending.push({
373 ...stderrMeta,
374 name: "stderr",
375 value: stderr,
376 } as WorkspaceRuntimeEvent<E>);
377 }
378 pending.sort((a, b) => a.seq - b.seq);
379 const span = beforeSeq !== undefined && beforeSeq > lastSeq ? beforeSeq - lastSeq : 1;
380 let index = 0;
381 for (const event of pending) {
382 index += 1;
383 enqueue(controller, {
384 ...event,
385 seq: lastSeq + (span * index) / (pending.length + 1),
386 } as WorkspaceRuntimeEvent<E>);
387 }
388 stdoutMeta = undefined;
389 stderrMeta = undefined;
390 };
391 return source.pipeThrough(
392 new TransformStream<WorkspaceRuntimeEvent, WorkspaceRuntimeEvent<E>>({
393 transform(event, controller) {
394 if (event.name === "stdout" || event.name === "stderr") {
395 if (event.name === "stdout") stdoutMeta = { id: event.id, seq: event.seq };

Callers 1

pullFunction · 0.85

Calls 1

pipeThroughMethod · 0.65

Tested by

no test coverage detected