Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 741a606d00 | |||
| 7537d7fa37 |
@@ -1,4 +1,4 @@
|
||||
# stream-to-push-stream
|
||||
# @push-stream-std/stream-to-push-stream
|
||||
|
||||
Convert native Node.js streams to push-stream source, sink, and duplex interfaces.
|
||||
|
||||
@@ -41,8 +41,7 @@ npm config set -- //hub.kl1.tenere.ai/api/packages/push-stream-std/npm/:_authTok
|
||||
npm install @push-stream-std/stream-to-push-stream
|
||||
```
|
||||
|
||||
## Description
|
||||
|
||||
## Why
|
||||
Utilities for bridging between the Node.js `stream` module (Readable, Writable, Duplex) and the push-stream protocol. PushSource wraps a Readable, listening to `data`/`end`/`error` events and implementing `pipe()`, `resume()`, and `abort()`. PushSink wraps a Writable, translating `write()` calls to `writable.write()` with `drain`-based backpressure and `end()` to `writable.end()`. PushDuplex wraps a Duplex stream as a combined source/sink for bidirectional push-stream pipelines.
|
||||
|
||||
## API
|
||||
@@ -78,8 +77,5 @@ Tests cover readable-to-source data flow and end handling, writable-to-sink writ
|
||||
|
||||
## Requirements
|
||||
|
||||
Node.js 18 or later. ESM only.
|
||||
Node.js 22 or later. ESM only.
|
||||
|
||||
## License
|
||||
|
||||
MIT
|
||||
|
||||
@@ -1,5 +1,12 @@
|
||||
import { Readable, Writable, Duplex } from 'stream';
|
||||
|
||||
/**
|
||||
* stream-to-push-stream - Convert Node.js streams to push-streams
|
||||
*
|
||||
* Bridges Node.js Readable/Writable/Duplex streams into the push-stream protocol
|
||||
* so they can be composed with `@push-stream-std/*` operators.
|
||||
*/
|
||||
|
||||
class PushSource {
|
||||
constructor(readable) {
|
||||
this.readable = readable;
|
||||
@@ -223,6 +230,11 @@ class PushDuplex {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap a Node.js Readable stream as a push-stream source.
|
||||
* @param {import('stream').Readable} readableStream
|
||||
* @returns {object} A push-stream source
|
||||
*/
|
||||
export function source(readableStream) {
|
||||
if (!(readableStream instanceof Readable)) {
|
||||
throw new Error('Expected a Node.js Readable stream');
|
||||
@@ -230,6 +242,11 @@ export function source(readableStream) {
|
||||
return new PushSource(readableStream);
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap a Node.js Writable stream as a push-stream sink.
|
||||
* @param {import('stream').Writable} writableStream
|
||||
* @returns {object} A push-stream sink
|
||||
*/
|
||||
export function sink(writableStream) {
|
||||
if (!(writableStream instanceof Writable)) {
|
||||
throw new Error('Expected a Node.js Writable stream');
|
||||
@@ -237,6 +254,11 @@ export function sink(writableStream) {
|
||||
return new PushSink(writableStream);
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap a Node.js Duplex stream as a push-stream duplex (`{ source, sink }`).
|
||||
* @param {import('stream').Duplex} duplexStream
|
||||
* @returns {{source: object, sink: object}} Push-stream duplex
|
||||
*/
|
||||
export function duplex(duplexStream) {
|
||||
if (!(duplexStream instanceof Duplex)) {
|
||||
throw new Error('Expected a Node.js Duplex stream');
|
||||
|
||||
+5
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@push-stream-std/stream-to-push-stream",
|
||||
"version": "1.0.1",
|
||||
"version": "1.0.3",
|
||||
"description": "Convert Node.js streams to push-streams",
|
||||
"type": "module",
|
||||
"main": "index.js",
|
||||
@@ -29,5 +29,9 @@
|
||||
],
|
||||
"publishConfig": {
|
||||
"registry": "https://hub.kl1.tenere.ai/api/packages/push-stream-std/npm/"
|
||||
},
|
||||
"repository": {
|
||||
"type": "git",
|
||||
"url": "https://hub.kl1.tenere.ai/push-stream-std/stream-to-push-stream.git"
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user