| 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) { |