3 ms·
I wrote a library a while back which implements this. It handles all the edge cases, provides the ability to measure message backpressure and allows you to clos
by socketcluster 3y ago
I wrote a library a while back which implements this. It handles all the edge cases, provides the ability to measure message backpressure and allows you to close or abruptly kill a stream:
https://github.com/SocketCluster/writable-consumable-stream?tab=readme-ov-file#writable-consumable-stream https://github.com/SocketCluster/writable-consumable-stream?...
It has full test coverage and is a core part SocketCluster which has almost full e2e test coverage.
To put it in perspective, the library itself is very lightweight with under 300 lines of code including dependencies (only 1 of them) and has 1.6K lines of test code...
It supports concurrent consumption/iteration by default (each consuming loop causes a new async iterator to be created implicitly in accordance with the AsyncIterator spec) but it also gives you the ability to manually create consumer streams using `stream.createConsumer()` for cases where you want to start buffering messages some time before consuming/iterating over them.