buffered-transform.ts 2.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758
  1. import type { ValueOrPromise } from '@yume-chan/struct';
  2. import { BufferedReadableStream, BufferedReadableStreamEndedError } from './buffered.js';
  3. import { PushReadableStream, PushReadableStreamController } from './push-readable.js';
  4. import { ReadableStream, ReadableWritablePair, WritableStream } from './stream.js';
  5. // TODO: BufferedTransformStream: find better implementation
  6. export class BufferedTransformStream<T> implements ReadableWritablePair<T, Uint8Array> {
  7. private _readable: ReadableStream<T>;
  8. public get readable() { return this._readable; }
  9. private _writable: WritableStream<Uint8Array>;
  10. public get writable() { return this._writable; }
  11. constructor(transform: (stream: BufferedReadableStream) => ValueOrPromise<T>) {
  12. // Convert incoming chunks to a `BufferedReadableStream`
  13. let sourceStreamController!: PushReadableStreamController<Uint8Array>;
  14. const buffered = new BufferedReadableStream(new PushReadableStream<Uint8Array>(
  15. controller =>
  16. sourceStreamController = controller,
  17. ));
  18. this._readable = new ReadableStream<T>({
  19. async pull(controller) {
  20. try {
  21. const value = await transform(buffered);
  22. controller.enqueue(value);
  23. } catch (e) {
  24. // TODO: BufferedTransformStream: The semantic of stream ending is not clear
  25. // If the `transform` started but did not finish, it should really be an error?
  26. // But we can't detect that, unless there is a `peek` method on buffered stream.
  27. if (e instanceof BufferedReadableStreamEndedError) {
  28. controller.close();
  29. return;
  30. }
  31. throw e;
  32. }
  33. },
  34. cancel: (reason) => {
  35. // Propagate cancel to the source stream
  36. // So future writes will be rejected
  37. buffered.cancel(reason);
  38. }
  39. });
  40. this._writable = new WritableStream({
  41. async write(chunk) {
  42. await sourceStreamController.enqueue(chunk);
  43. },
  44. abort() {
  45. sourceStreamController.close();
  46. },
  47. close() {
  48. sourceStreamController.close();
  49. },
  50. });
  51. }
  52. }