Node clone readable stream

Мне необходимо клонировать Readable stream и записать его содержимое в два Writeble stream Потом посередине поставлю Transform и сделаю resize картинки в два формата с помощью только одной операции чтения данных.

Для старта нашел небольшой пакет : https://github.com/levansuper/readable-stream-clone/blob/master/readable-stream-clone.js

И переписал его на современный лад (Node 14 LTS): https://github.com/Ulibka68/s3-compress-image-node/blob/master/src/cloneRead.ts

Данный код работает, но проблема заключается в том что я не понимаю почему он работает :)

Помогите пожалуйста понять - почему данный код работает?

Основная идея достаточно проста - сначала создается основной экземпляр stream Readable потом создается два дополнительных экземпляра класса stream Readable . Дополнительные экземпляры подписываются на событие on ('data') основного потока и с помощью метода push кладут в свой внутренний буфер данные.

Дальше создаются два экземпляра stream Writeble которые присоединяются через pipe и запускается процесс копирования.

Что мне непонятно - после того как запустился первый процесс readClone1.pipe(writeStream1); (если я правильно понимаю) - основной поток должен быть полностью прочитан, положен по частям - по 3 байта во внутренний буфер readStream и readClone1 . Первый чанк также попадет в readClone2.

Первый чанк запишется в writeStream1 (по идее первый чанк во writeStream2 не пишется - или пишется ? )

После того как первый чанк записался в writeStream1 - очишается внутренний буфер чтения readStream и readClone1 и открывается возможность чтения второго чанка с основного потока.

При этом если первый чанк не записался во writeStream2 - то тогда внутренний буфер в readClone2 не освободился и чтение из него не возможно.

Т.е. грубо говоря по моему мнению в файле 'text2.txt' я должен увидеть только первые 3 байта - однако данный код прекрасно отрабатывает и файлы 'text.txt','text1.txt' и 'text2.txt' - полностью идентичны.

import { Readable, ReadableOptions } from 'stream';

class ReadableStreamClone extends Readable {
  constructor(readableStream: Readable, options?: ReadableOptions) {
    super(options);

    readableStream.on('data', (chunk) => {
      this.push(chunk);
    });

    readableStream.on('end', () => {
      this.push(null);
    });

    readableStream.on('error', (err) => {
      this.emit('error', err);
    });

    readableStream.on('error', (err) => {
      this.emit('error', err);
    });
  }

  // eslint-disable-next-line @typescript-eslint/no-empty-function
  _read() {}
}

module.exports = ReadableStreamClone;

import * as fs from 'fs';

const readStream = fs.createReadStream('text.txt', { highWaterMark: 3 });

const readClone1 = new ReadableStreamClone(readStream, { highWaterMark: 3 });
const readClone2 = new ReadableStreamClone(readStream, { highWaterMark: 3 });

const writeStream1 = fs.createWriteStream('text1.txt');
const writeStream2 = fs.createWriteStream('text2.txt');

readClone1.pipe(writeStream1);
readClone2.pipe(writeStream2);


Ответы (0 шт):