Implemented writable streams interface

This commit is contained in:
birme 2020-06-22 13:45:43 +02:00
parent ec9854efe7
commit ef94263b8a
6 changed files with 63 additions and 6 deletions

View file

@ -32,7 +32,7 @@ if (fd) {
```
class SRT {
createSocket(): socket:Number
createSocket(sender?:Boolean): socket:Number
bind(socket:Number, address:String, port:Number): result:Number
listen(socket:Number, backlog:Number): result:Number
connect(socket:Number, host:String, port:Number): result:Number
@ -66,6 +66,21 @@ srt.connect(readStream => {
});
```
### Writable Stream
Example of a writable stream
```
const fs = require('fs');
const source = fs.createReadStream('./input');
const { SRTWriteStream } = require('@eyevinn/srt');
const srt = new SRTWriteStream('127.0.0.1', 1234);
srt.connect(writeStream => {
source.pipe(writeStream);
});
```
## [Contributing](CONTRIBUTING.md)
In addition to contributing code, you can help to triage issues. This can include reproducing bug reports, or asking for vital information such as version numbers or reproduction instructions.

8
examples/writable.js Normal file
View file

@ -0,0 +1,8 @@
const fs = require('fs');
const source = fs.createReadStream(process.argv[2], { highWaterMark: 1316 });
const { SRTWriteStream } = require('../index.js');
const srt = new SRTWriteStream('127.0.0.1', 1234);
srt.connect(writeStream => {
source.pipe(writeStream);
});

View file

@ -1,9 +1,10 @@
const LIB = require('./build/Release/node_srt.node');
const Server = require('./src/server.js');
const { SRTReadStream } = require('./src/stream.js');
const { SRTReadStream, SRTWriteStream } = require('./src/stream.js');
module.exports = {
SRT: LIB.SRT,
Server: Server,
SRTReadStream
SRTReadStream,
SRTWriteStream
}

View file

@ -44,11 +44,20 @@ Napi::Value NodeSRT::CreateSocket(const Napi::CallbackInfo& info) {
Napi::Env env = info.Env();
Napi::HandleScope scope(env);
Napi::Boolean isSender = Napi::Boolean::New(env, false);
if (info.Length() > 0) {
isSender = info[0].As<Napi::Boolean>();
}
SRTSOCKET socket = srt_socket(AF_INET, SOCK_DGRAM, 0);
if (socket == SRT_ERROR) {
Napi::Error::New(env, srt_getlasterror_str()).ThrowAsJavaScriptException();
return Napi::Number::New(env, SRT_ERROR);
}
if (isSender) {
int yes = 1;
srt_setsockflag(socket, SRTO_SENDER, &yes, sizeof(yes));
}
return Napi::Number::New(env, socket);
}

View file

@ -16,4 +16,4 @@ class NodeSRT : public Napi::ObjectWrap<NodeSRT> {
Napi::Value Close(const Napi::CallbackInfo& info);
Napi::Value Read(const Napi::CallbackInfo& info);
Napi::Value Write(const Napi::CallbackInfo& info);
};
};

View file

@ -1,4 +1,4 @@
const { Readable } = require('stream');
const { Readable, Writable } = require('stream');
const LIB = require('../build/Release/node_srt.node');
const debug = require('debug')('srt-stream');
@ -56,6 +56,30 @@ class SRTReadStream extends Readable {
}
}
class SRTWriteStream extends Writable {
constructor(address, port, opts) {
super();
this.srt = new LIB.SRT();
this.socket = this.srt.createSocket();
this.address = address;
this.port = port;
}
connect(cb) {
this.srt.connect(this.socket, this.address, this.port);
this.fd = this.socket;
if (this.fd) {
cb(this);
}
}
_write(chunk, encoding, callback) {
this.srt.write(this.fd, chunk);
callback();
}
}
module.exports = {
SRTReadStream
SRTReadStream,
SRTWriteStream
}