Async API server/streams impl & functional tests + Allow to use local SRT source repo instead of remote (#9)

* replace stream module by improved version of readable/writable impl

* rm server.js

* async api improvments:
- better tracing of calls from worker back and forth
- fix transferrable handling to avoid copying buffers for r/w
- optional debug logs
- completed jsdocs annotations
- add dispose method
- add setLogLevel method (analoguous to added binding)

* node-srt C bindings:
- add SetLogLevel to get libSRT log output if desired
- add OK static member
- add #define EPOLL_EVENTS_NUM_MAX 1024
- improve error string thrown in Read (add that it comes from srt_recvmsg)
- improve error string thrown in Write (add that it comes from srt_sendmsg2)
- misc isofunctional improvements (var names) and comments

* add SRT logging related JS-side helper

* rewrite flat TypeScript decl files without "module" keyword

* add ts enum decl for all libSRT enums

* async-worker: enable using transferrable for zero-copy
+ allow better debugging (like in api/dispatcher side)
+ misc improvements on code quality

* add async-helpers: various functions to help dealing with transferrables
+ tracing calls to native bindings in debug output

* add async read/write modes functions + async-reader-writer class
- these will allow for performing high-level r/w operations conveniently
at optimum throughput for larger pieces of payload i.e list of packets.

* add srt-server and srt-connection (can manage multiple clients),
- based on async-api
- can be used with reader/writer (i.e the underlying modes)

* srt-server/connection typings

* async srt spec: add dispose method usage (but commented out as crashing atm)

* async srt spec: rm redundant checks on SRT static members (they are done
in other spec already)

* promises api spec: formal fixes

* stream spec: add dummy test

* package.json:
- put gyp toolchain in runtime deps (since the build happens on install)
- add JEST test runner
- shorten check-tsc script
- rebuild script: check & use all CPU cores available
- run rebuild actually on install, not preinstall (fixes deps not being there)
- remove preinstall and thus "npm install git-clone" in the package scripts

* update package lock

* update typings index not to need triple-slashs anymore

* in srt.ts example: check for read return value type

* build-srt-sdk script:
- allow to use any local libSRT code repo
- when using make: use all amount of cores available for build
- isolate better code running on different platforms

* update package main index with new things

* add enum typings index

* add jest config

* add "use strict" on async-srt-await example

* add integration/smoke testing for client-to-server one-way burst write

* readme: add note on build prerequisites

* readme: add infos on new components SRTServer/Connection & AsyncReaderWriter
This commit is contained in:
Stephan Hesse 2020-10-21 13:09:52 +02:00 • committed by GitHub
parent bf8745b66f
commit bf60795889
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
35 changed files with 6548 additions and 758 deletions

View file

@ -13,6 +13,21 @@ npm install --save @eyevinn/srt
Installing from NPM downloads and builds SRT SDK and NodeJS addon for your operating system × architecture.
## Prerequisities
Please refer to build instructions for your OS of the [SRT project](https://github.com/Haivision/srt#requirements).
This is providing a NodeJS binding layer that at its build-time assumes that plain libSRT already *can* be built on your system.
We will merely pull in an SRT codebase from a GIT repo here (that may be local or remote), and attempt to compile it (using the specific toolchain and commands invoked for each OS). See `scripts/build-srt-sdk.js`. Then linking the result of it into the compiled NodeJS add-on we provide here via the Gyp tool.
Everything we need on the NodeJS side of things (N-API, Gyp) gets installed via NPM,
as you install this package.
However, it is not provided with this any prerequisits of building libSRT itself (we only try to invoke the toolchain correctly, whichever it is, on your OS).
As you install all of the prerequisites, or even build the library already in your environment, this package should work also likewise.
## Example
```
@ -93,6 +108,21 @@ or with promises:
});
```
### High-performance read/write use-cases & server/multi-connection implementation
In order to perform on certain use-cases where larger chunks of data split into packets
would need to be sent/received in bursts, we provide "modes" and a high-level class called `AsyncReaderWriter`, which can be initialized an existing `AsyncSRT` instance, i.e a worker thread, and a given SRT socket identifier. It therefore allows to plug-into any concept of a
connection on either side. A client connection instance can therefore just be a socket created
with a successful connection state via the SRT API, and then using the reader-writer for
any sort of transmission upon it.
Also, we provide a class to allow building a server that can accept multiple incoming connections(`SRTServer` and its friend `SRTConnection`). The latter server-side connection object has the method `SRTConnection#getReaderWriter()` to allow using the reader-writer
in order to communicate with the respective client.
The best example to see all this in action at once is taking a look at the respective integration test(s).
These components also all have JSdoc annotations that should help with their usage.
### Readable Stream
A custom readable stream API is also available, example (in listener mode):

View file

@ -1,3 +1,5 @@
"use strict";
const { AsyncSRT } = require('../src/async.js');
const asyncSrt = new AsyncSRT();

View file

@ -1,4 +1,4 @@
import {SRT} from '../index';
import {SRT, SRTReadReturn} from '../index';
const srt = new SRT();
const socket = srt.createSocket();
@ -26,6 +26,8 @@ const fhandle = srt.accept(socket);
if (fhandle) {
console.log("Client connected");
const chunk = srt.read(fhandle, 1316);
console.log("Read chunk: " + chunk.length);
const chunk: SRTReadReturn = srt.read(fhandle, 1316);
if (chunk instanceof Uint8Array) {
console.log("Read chunk: " + chunk.length);
}
}

2
index-enums.ts Normal file
View file

@ -0,0 +1,2 @@
export * from "./src/srt-api-enums";
export * from "./src/async-api-enums";

13
index.d.ts vendored
View file

@ -1,7 +1,8 @@
/* eslint-disable @typescript-eslint/triple-slash-reference */
/// <reference path="./types/srt-api.d.ts" />
/// <reference path="./types/srt-stream.d.ts" />
/// <reference path="./types/srt-server.d.ts" />
/// <reference path="./types/srt-api-async.d.ts" />
import { SRTLoggingLevel } from "./src/srt-api-enums";
export * from "srt";
export * from "./types/srt-api";
export * from "./types/srt-api-async";
export * from "./types/srt-server";
export * from "./types/srt-stream";
export function setSRTLoggingLevel(level: SRTLoggingLevel);

View file

@ -1,12 +1,15 @@
const { SRT } = require('./build/Release/node_srt.node');
const Server = require('./src/server.js');
const { SRTReadStream, SRTWriteStream } = require('./src/stream.js');
const { AsyncSRT } = require('./src/async');
const { SRTReadStream } = require('./src/srt-stream-readable.js');
const { SRTWriteStream } = require('./src/srt-stream-writable.js');
const { SRTServer } = require('./src/srt-server');
const { setSRTLoggingLevel } = require('./src/logging');
module.exports = {
SRT,
Server,
AsyncSRT,
SRTServer,
SRTReadStream,
SRTWriteStream,
AsyncSRT,
setSRTLoggingLevel
};

View file

@ -0,0 +1,175 @@
const { SRT, AsyncSRT, SRTServer } = require('../index');
const {
writeChunksWithYieldingLoop,
writeChunksWithExplicitScheduling
} = require('../src/async-write-modes');
const {sliceBufferToChunks, copyChunksIntoBuffer} = require('../src/tools')
const fs = require("fs");
const path = require("path");
const {performance} = require("perf_hooks");
const now = performance.now;
const testFiles = [
"data/SpringBlenderOpenMovie.mp4.ts"
]
jest && jest.setTimeout(5000)
describe("AsyncSRT to SRTServer one-way transmission", () => {
it("should transmit data written (yielding-loop)", async done => {
transmitClientToServerLoopback(9000, done, false);
});
it("should transmit data written (explicit-scheduling)", async done => {
transmitClientToServerLoopback(9001, done, true);
});
});
async function transmitClientToServerLoopback(localServerPort, done, useExplicitScheduling) {
const fileReadStartTime = now();
const sourceDataBuf = fs.readFileSync(path.resolve(__dirname, testFiles[0]))
const fileReadTimeDiffMs = now() - fileReadStartTime;
const localServerBindIface = '127.0.0.1';
const chunkMaxSize = 1024;
const numChunks = 8 * 1024;
// type size NodeJS-internal readable grabs in binary streams
const readBufSize = 1024 * 1024;
const bytesShouldSendTotal
= Math.min(numChunks * chunkMaxSize, sourceDataBuf.byteLength);
const clientWritesPerTick = 128;
console.log(`Read ${sourceDataBuf.byteLength} bytes from file into buffer in ${fileReadTimeDiffMs.toFixed(3)} ms`);
const packetDataSlicingStartTime = now();
const chunks = sliceBufferToChunks(sourceDataBuf, chunkMaxSize, bytesShouldSendTotal);
const packetDataSlicingTimeD = now() - packetDataSlicingStartTime;
console.log('Pre-slicing packet data took millis:', packetDataSlicingTimeD);
// we need two instances of task-runners here,
// because otherwise awaiting server accept
// result would deadlock
// client connection tasks
const asyncSrtServer = new SRTServer(localServerPort);
asyncSrtServer.on('connection', (connection) => {
onClientConnected(connection);
});
const asyncSrtClient = new AsyncSRT();
const [clientSideSocket] = await Promise.all([
asyncSrtClient.createSocket(), // we could also use the server-runner here. doesnt matter.
asyncSrtServer.create().then(s => s.open())
]);
console.log('Got socket handles (client/server):',
clientSideSocket, '/',
asyncSrtServer.socket);
clientWriteToConnection();
let clientWriteStartTime;
let clientWriteDoneTime;
let bytesSentCount = 0;
async function clientWriteToConnection() {
let result = await asyncSrtClient.connect(clientSideSocket,
localServerBindIface, localServerPort);
if (result === SRT.ERROR) {
throw new Error('client connect failed');
}
console.log('connect result:', result)
clientWriteStartTime = now();
if (useExplicitScheduling) {
writeChunksWithExplicitScheduling(asyncSrtClient,
clientSideSocket, chunks, onWrite, clientWritesPerTick);
} else {
writeChunksWithYieldingLoop(asyncSrtClient,
clientSideSocket, chunks, onWrite, clientWritesPerTick);
}
function onWrite(byteLength) {
bytesSentCount += byteLength;
if(bytesSentCount >= bytesShouldSendTotal) {
console.log('done writing, took millis:',
now() - clientWriteStartTime);
clientWriteDoneTime = now();
}
}
}
function onClientConnected(connection) {
console.log('Got new connection:', connection.fd)
let bytesRead = 0;
let firstByteReadTime;
const serverConnectionAcceptTime = now();
connection.on('data', async () => {
if (!connection.gotFirstData) {
onClientData();
}
});
const reader = connection.getReaderWriter();
async function onClientData() {
const chunks = await reader.readChunks(
bytesShouldSendTotal,
readBufSize,
(readBuf) => {
if (!firstByteReadTime) {
firstByteReadTime = now();
}
//console.log('Read buffer of size:', readBuf.byteLength)
bytesRead += readBuf.byteLength;
}, (errRes) => {
console.log('Error reading, got result:', errRes);
});
const readDoneTime = now();
const readTimeDiffMs = readDoneTime - serverConnectionAcceptTime;
const readBandwidthEstimKbps = (8 * (bytesShouldSendTotal / readTimeDiffMs))
console.log('Done reading stream, took millis:', readTimeDiffMs, 'for kbytes:~',
(bytesSentCount / 1000), 'of', (bytesShouldSendTotal / 1000));
console.log('Estimated read-bandwidth (kb/s):', readBandwidthEstimKbps.toFixed(3))
console.log('First-byte-write-to-read latency millis:',
firstByteReadTime - clientWriteStartTime)
console.log('End-to-end transfer latency millis:', readDoneTime - clientWriteStartTime)
console.log('Client-side writing took millis:',
clientWriteDoneTime - clientWriteStartTime);
expect(bytesSentCount).toEqual(bytesShouldSendTotal);
const receivedBuffer = copyChunksIntoBuffer(chunks);
expect(receivedBuffer.byteLength).toEqual(bytesSentCount);
/*
for (let i = 0; i < receivedBuffer.byteLength; i++) {
expect(sourceDataBuf.readInt8(i)).toEqual(receivedBuffer.readInt8(i));
}
*/
done();
}
}
}

200
jest.config.js Normal file
View file

@ -0,0 +1,200 @@
// For a detailed explanation regarding each configuration property, visit:
// https://jestjs.io/docs/en/configuration.html
module.exports = {
// All imported modules in your tests should be mocked automatically
// automock: false,
// Stop running tests after `n` failures
// bail: 0,
// The directory where Jest should store its cached dependency information
// cacheDirectory: "/private/var/folders/j9/xvw659rx2gx3b_4gmk6rx2nw0000gn/T/jest_dx",
// Automatically clear mock calls and instances between every test
clearMocks: true,
// Indicates whether the coverage information should be collected while executing the test
// collectCoverage: false,
// An array of glob patterns indicating a set of files for which coverage information should be collected
// collectCoverageFrom: undefined,
// The directory where Jest should output its coverage files
coverageDirectory: "coverage",
// An array of regexp pattern strings used to skip coverage collection
// coveragePathIgnorePatterns: [
// "/node_modules/"
// ],
// Indicates which provider should be used to instrument code for coverage
coverageProvider: "v8",
// A list of reporter names that Jest uses when writing coverage reports
// coverageReporters: [
// "json",
// "text",
// "lcov",
// "clover"
// ],
// An object that configures minimum threshold enforcement for coverage results
// coverageThreshold: undefined,
// A path to a custom dependency extractor
// dependencyExtractor: undefined,
// Make calling deprecated APIs throw helpful error messages
// errorOnDeprecated: false,
// Force coverage collection from ignored files using an array of glob patterns
// forceCoverageMatch: [],
// A path to a module which exports an async function that is triggered once before all test suites
// globalSetup: undefined,
// A path to a module which exports an async function that is triggered once after all test suites
// globalTeardown: undefined,
// A set of global variables that need to be available in all test environments
// globals: {},
// The maximum amount of workers used to run your tests. Can be specified as % or a number. E.g. maxWorkers: 10% will use 10% of your CPU amount + 1 as the maximum worker number. maxWorkers: 2 will use a maximum of 2 workers.
// maxWorkers: "50%",
// An array of directory names to be searched recursively up from the requiring module's location
// moduleDirectories: [
// "node_modules"
// ],
// An array of file extensions your modules use
// moduleFileExtensions: [
// "js",
// "json",
// "jsx",
// "ts",
// "tsx",
// "node"
// ],
// A map from regular expressions to module names or to arrays of module names that allow to stub out resources with a single module
// moduleNameMapper: {},
// An array of regexp pattern strings, matched against all module paths before considered 'visible' to the module loader
// modulePathIgnorePatterns: [],
// Activates notifications for test results
// notify: false,
// An enum that specifies notification mode. Requires { notify: true }
// notifyMode: "failure-change",
// A preset that is used as a base for Jest's configuration
// preset: undefined,
// Run tests from one or more projects
// projects: undefined,
// Use this configuration option to add custom reporters to Jest
// reporters: undefined,
// Automatically reset mock state between every test
// resetMocks: false,
// Reset the module registry before running each individual test
// resetModules: false,
// A path to a custom resolver
// resolver: undefined,
// Automatically restore mock state between every test
// restoreMocks: false,
// The root directory that Jest should scan for tests and modules within
// rootDir: undefined,
// A list of paths to directories that Jest should use to search for files in
// roots: [
// "<rootDir>"
// ],
// Allows you to use a custom runner instead of Jest's default test runner
// runner: "jest-runner",
// The paths to modules that run some code to configure or set up the testing environment before each test
// setupFiles: [],
// A list of paths to modules that run some code to configure or set up the testing framework before each test
// setupFilesAfterEnv: [],
// The number of seconds after which a test is considered as slow and reported as such in the results.
// slowTestThreshold: 5,
// A list of paths to snapshot serializer modules Jest should use for snapshot testing
// snapshotSerializers: [],
// The test environment that will be used for testing
testEnvironment: "node",
preset: 'ts-jest',
// Options that will be passed to the testEnvironment
// testEnvironmentOptions: {},
// Adds a location field to test results
// testLocationInResults: false,
// The glob patterns Jest uses to detect test files
testMatch: [
"/**/*_spec.js",
"/**/*_test.js",
],
// An array of regexp pattern strings that are matched against all test paths, matched tests are skipped
testPathIgnorePatterns: [
"/build/",
"/deps/",
"/examples/",
"/lib/",
"/scripts/",
"/types/",
"/node_modules/"
],
// The regexp pattern or array of patterns that Jest uses to detect test files
// testRegex: [],
// This option allows the use of a custom results processor
// testResultsProcessor: undefined,
// This option allows use of a custom test runner
// testRunner: "jasmine2",
// This option sets the URL for the jsdom environment. It is reflected in properties such as location.href
// testURL: "http://localhost",
// Setting this value to "fake" allows the use of fake timers for functions such as "setTimeout"
// timers: "real",
// A map from regular expressions to paths to transformers
// transform: undefined,
// An array of regexp pattern strings that are matched against all source file paths, matched files will skip transformation
// transformIgnorePatterns: [
// "/node_modules/",
// "\\.pnp\\.[^\\/]+$"
// ],
// An array of regexp pattern strings that are matched against all modules before the module loader will automatically return a mock for them
// unmockedModulePathPatterns: undefined,
// Indicates whether each individual test should be reported during the run
// verbose: undefined,
// An array of regexp patterns that are matched against all source file paths before re-running tests in watch mode
// watchPathIgnorePatterns: [],
// Whether to use watchman for file crawling
// watchman: true,
};

4616
package-lock.json generated

File diff suppressed because it is too large Load diff

View file

@ -4,13 +4,14 @@
"description": "Nodejs bindings for Secure Reliable Transport SDK",
"main": "index.js",
"scripts": {
"preinstall": "npm install git-clone && node scripts/build-srt-sdk.js",
"install": "node-gyp rebuild -j 8",
"rebuild": "node-gyp rebuild",
"install": "npm run build-srt && npm run rebuild",
"build-srt": "node scripts/build-srt-sdk.js",
"rebuild": "node-gyp rebuild -j $(echo \"console.log(require('os').cpus().length)\" | node)",
"clean": "node-gyp clean",
"test": "$(npm bin)/jasmine",
"test": "jasmine",
"test-jest": "jest --runInBand --detectOpenHandles",
"lint": "eslint . --ext .js --ext .ts",
"check-tsc": "./node_modules/.bin/tsc examples/srt.ts --outDir ./tsc-lib",
"check-tsc": "tsc examples/srt.ts --outDir ./tsc-lib",
"postversion": "git push && git push --tags"
},
"repository": {
@ -26,20 +27,24 @@
],
"license": "MIT",
"devDependencies": {
"@types/jest": "^26.0.10",
"@types/node": "^14.0.23",
"@typescript-eslint/eslint-plugin": "^3.6.1",
"@typescript-eslint/parser": "^3.6.1",
"eslint": "^7.4.0",
"eslint-plugin-jasmine": "^4.1.1",
"jasmine": "^3.5.0",
"node-gyp": "^7.0.0",
"jest": "^26.4.1",
"ts-jest": "^26.2.0",
"ts-node": "^8.10.2",
"typescript": "^3.9.6"
},
"dependencies": {
"debug": "^4.1.1",
"del": "^5.1.0",
"git-clone": "^0.1.0",
"node-addon-api": "^3.0.0"
"node-addon-api": "^3.0.0",
"node-gyp": "^7.0.0"
},
"bugs": {
"url": "https://github.com/Eyevinn/node-srt/issues"

183
scripts/build-srt-sdk.js Normal file → Executable file
View file

@ -1,94 +1,123 @@
#!/usr/bin/env node
"use strict";
const path = require('path');
const fs = require('fs');
const process = require('process');
const clone = require('git-clone');
const del = require('del');
const { spawnSync } = require('child_process');
const os = require('os');
const env = process.env;
const SRT_REPO = "https://github.com/Haivision/srt.git";
const SRT_VERSION = "v1.4.1";
const SRT_REPO = env.NODE_SRT_REPO || "https://github.com/Haivision/srt.git";
const SRT_CHECKOUT = "v1.4.1";
const srtRepoPath = env.NODE_SRT_LOCAL_REPO ? `file://${path.join(__dirname, env.NODE_SRT_LOCAL_REPO)}` : SRT_REPO;
const srtCheckout = env.NODE_SRT_CHECKOUT || SRT_CHECKOUT;
const depsPath = path.join(__dirname, '../', 'deps');
const srtSourcePath = path.join(depsPath, 'srt');
const buildDir = path.join(depsPath, 'build');
const buildDir = path.join(depsPath, 'build'); // FIXME: name this srt-build (in case other deps come up)
const numCpus = os.cpus().length; // NOTE: not the actual physical cores amount btw, see https://www.npmjs.com/package/physical-cpu-count
if (!fs.existsSync(depsPath)) {
console.log('Creating dir:', depsPath)
fs.mkdirSync(depsPath);
}
if (!fs.existsSync(buildDir)) {
console.log(`Cloning ${SRT_REPO}:${SRT_VERSION}`);
clone(SRT_REPO, srtSourcePath, { checkout: SRT_VERSION }, () => {
if (process.platform === "win32") {
process.env.SRT_ROOT = srtSourcePath;
fs.mkdirSync(buildDir);
console.log("Building OpenSSL");
const openssl = spawnSync('vcpkg', [ 'install', 'openssl', '--triplet', `${process.arch}-windows` ], { cwd: process.env.VCPKG_ROOT, shell: true } );
if (openssl.stdout)
console.log(openssl.stdout.toString());
if (openssl.status) {
console.log(openssl.stderr.toString());
process.exit(openssl.status);
}
if (!fs.existsSync(srtSourcePath)) {
console.log(`Cloning ${srtRepoPath}#${srtCheckout}`);
clone(srtRepoPath, srtSourcePath, { checkout: srtCheckout }, (err) => {
console.log("Building pthreads");
const pthreads = spawnSync('vcpkg', [ 'install', 'pthreads', '--triplet', `${process.arch}-windows` ], { cwd: process.env.VCPKG_ROOT, shell: true } );
if (pthreads.stdout)
console.log(pthreads.stdout.toString());
if (pthreads.status) {
console.log(pthreads.stderr.toString());
process.exit(pthreads.status);
}
console.log("Integrate vcpkg build system");
const integrate = spawnSync('vcpkg', [ 'integrate', 'install' ], { cwd: process.env.VCPKG_ROOT, shell: true } );
if (integrate.stdout)
console.log(integrate.stdout.toString());
if (integrate.status) {
console.log(integrate.stderr.toString());
process.exit(integrate.status);
}
console.log("Running cmake generator");
const generator = spawnSync('cmake', [ srtSourcePath, '-DCMAKE_BUILD_TYPE=Release', '-G"Visual Studio 16 2019"', '-A', process.arch, '-DCMAKE_TOOLCHAIN_FILE="%VCPKG_ROOT%\\scripts\\buildsystems\\vcpkg.cmake' ], { cwd: buildDir, shell: true } );
if (generator.stdout)
console.log(generator.stdout.toString());
if (generator.status) {
console.log(generator.stderr.toString());
process.exit(generator.status);
}
console.log("Running cmake build");
const build = spawnSync('cmake', [ '--build', buildDir, '--config', 'Release' ], { cwd: buildDir, shell: true } );
if (build.stdout)
console.log(build.stdout.toString());
if (build.status) {
console.log(build.stderr.toString());
process.exit(build.status);
}
} else {
console.log("Running ./configure");
const configure = spawnSync('./configure', [ '--prefix', buildDir ], { cwd: srtSourcePath, shell: true } );
console.log(configure.stdout.toString());
if (configure.status) {
console.log(configure.stderr.toString());
process.exit(configure.status);
}
console.log("Running make");
const make = spawnSync('make', [], { cwd: srtSourcePath, shell: true });
console.log(make.stdout.toString());
if (make.status) {
console.log(make.stderr.toString());
process.exit(make.status);
}
console.log("Running make install");
const install = spawnSync('make', [ 'install' ], { cwd: srtSourcePath, shell: true });
console.log(install.stdout.toString());
if (install.status) {
console.log(install.stderr.toString());
process.exit(install.status);
}
if (err) {
console.error(err.message);
if (fs.existsSync(srtSourcePath)) del.sync(srtSourcePath);
process.exit(1);
}
build();
});
}
} else {
build();
}
function build() {
console.log('Building SRT SDK and prerequisites')
if (process.platform === "win32") {
buildWin32();
} else {
buildNx();
}
}
function buildWin32() {
process.env.SRT_ROOT = srtSourcePath;
fs.mkdirSync(buildDir);
console.log("Building OpenSSL");
const openssl = spawnSync('vcpkg', [ 'install', 'openssl', '--triplet', `${process.arch}-windows` ], { cwd: process.env.VCPKG_ROOT, shell: true } );
if (openssl.stdout)
console.log(openssl.stdout.toString());
if (openssl.status) {
console.log(openssl.stderr.toString());
process.exit(openssl.status);
}
console.log("Building pthreads");
const pthreads = spawnSync('vcpkg', [ 'install', 'pthreads', '--triplet', `${process.arch}-windows` ], { cwd: process.env.VCPKG_ROOT, shell: true } );
if (pthreads.stdout)
console.log(pthreads.stdout.toString());
if (pthreads.status) {
console.log(pthreads.stderr.toString());
process.exit(pthreads.status);
}
console.log("Integrate vcpkg build system");
const integrate = spawnSync('vcpkg', [ 'integrate', 'install' ], { cwd: process.env.VCPKG_ROOT, shell: true } );
if (integrate.stdout)
console.log(integrate.stdout.toString());
if (integrate.status) {
console.log(integrate.stderr.toString());
process.exit(integrate.status);
}
console.log("Running cmake generator");
const generator = spawnSync('cmake', [ srtSourcePath, '-DCMAKE_BUILD_TYPE=Release', '-G"Visual Studio 16 2019"', '-A', process.arch, '-DCMAKE_TOOLCHAIN_FILE="%VCPKG_ROOT%\\scripts\\buildsystems\\vcpkg.cmake' ], { cwd: buildDir, shell: true } );
if (generator.stdout)
console.log(generator.stdout.toString());
if (generator.status) {
console.log(generator.stderr.toString());
process.exit(generator.status);
}
console.log("Running CMake build");
const build = spawnSync('cmake', [ '--build', buildDir, '--config', 'Release' ], { cwd: buildDir, shell: true } );
if (build.stdout)
console.log(build.stdout.toString());
if (build.status) {
console.log(build.stderr.toString());
process.exit(build.status);
}
}
function buildNx() {
console.log("Running ./configure");
const configure = spawnSync('./configure', [ '--prefix', buildDir ], { cwd: srtSourcePath, shell: true, stdio: 'inherit' } );
if (configure.status) {
process.exit(configure.status);
}
console.log("Running make with threads:", numCpus);
const make = spawnSync('make', [`-j${numCpus}`], { cwd: srtSourcePath, shell: true, stdio: 'inherit' });
if (make.status) {
process.exit(make.status);
}
console.log("Running make install");
const install = spawnSync('make', ['install'], { cwd: srtSourcePath, shell: true, stdio: 'inherit' });
if (install.status) {
process.exit(install.status);
}
}

View file

@ -1,11 +1,13 @@
const { SRT, AsyncSRT } = require('../index.js');
describe("Async SRT Library with async/await", () => {
describe("Async SRT API with async/await", () => {
it("can create an SRT socket", async () => {
const asyncSrt = new AsyncSRT();
const socket = await asyncSrt.createSocket(false);
expect(socket).not.toEqual(SRT.ERROR);
//return await asyncSrt.dispose();
});
it("can create an SRT socket for sending data", async () => {
@ -13,5 +15,7 @@ describe("Async SRT Library with async/await", () => {
const socket = await asyncSrt.createSocket(true);
expect(socket).not.toEqual(SRT.ERROR);
//return await asyncSrt.dispose();
});
});
});

View file

@ -1,10 +1,10 @@
const { SRT, AsyncSRT } = require('../index.js');
describe("Async SRT Library with promises", () => {
describe("Async SRT API with promises", () => {
it("can create an SRT socket", done => {
const asyncSrt = new AsyncSRT();
asyncSrt.createSocket(false)
.then(socket => {
.then(socket => {
expect(socket).not.toEqual(SRT.ERROR);
done();
}).catch(done.fail);
@ -13,9 +13,9 @@ describe("Async SRT Library with promises", () => {
it("can create an SRT socket for sending data", done => {
const asyncSrt = new AsyncSRT();
asyncSrt.createSocket(true)
.then(socket => {
.then(socket => {
expect(socket).not.toEqual(SRT.ERROR);
done();
}).catch(done.fail);
});
});
});

View file

@ -1,15 +1,6 @@
const { SRT, AsyncSRT } = require('../index.js');
describe("Async SRT Library", () => {
it("exposes constants", () => {
expect(SRT.ERROR).toEqual(-1);
expect(SRT.INVALID_SOCK).toEqual(-1);
});
it("exposes socket options", () => {
expect(SRT.SRTO_UDP_SNDBUF).toEqual(8);
expect(SRT.SRTO_RCVLATENCY).toEqual(43);
});
describe("Async SRT API with callbacks", () => {
it("can create an SRT socket", done => {
const asyncSrt = new AsyncSRT();

View file

@ -1,3 +1,9 @@
const fs = require('fs');
const dest = fs.createWriteStream('/dev/null');
const { SRTReadStream } = require('../index.js');
describe("SRTReadStream", () => {
it('can be constructed without throwing an exception', () => {
new SRTReadStream();
})
});

13
src/async-api-enums.ts Normal file
View file

@ -0,0 +1,13 @@
export enum SRTServerEvent {
CREATED = "created",
CONNECTION = "connection",
DISCONNECTION = "disconnection",
DISPOSED = "disposed"
}
export enum SRTConnectionEvent {
DATA = "data",
CLOSING = "closing",
CLOSED = "closed"
}

45
src/async-helpers.js Normal file
View file

@ -0,0 +1,45 @@
function argsToString(args) {
const list = args
.map(argsItemToString).join(', ');
return `[${list}]`;
}
function isBufferOrTypedArray(elem) {
return elem.buffer
&& elem.buffer instanceof ArrayBuffer;
}
function argsItemToString(elem) {
if (isBufferOrTypedArray(elem)) {
return `${elem.constructor.name}<bytes=${elem.byteLength}>`
} else {
return elem;
}
}
function traceCallToString(method, args) {
return `SRT.${method}(...${argsToString(args)});`
}
/**
* @see https://nodejs.org/api/worker_threads.html#worker_threads_worker_threads
* @see https://developer.mozilla.org/en-US/docs/Web/API/Transferable
* @param {any[]} args Used parameter list to extract Transferrables from
* @returns {ArrayBuffer[]} List of transferrable objects owned by items of a parameter list
*/
function extractTransferListFromParams(args) {
const transferList = args.reduce((accu, item, index) => {
if (isBufferOrTypedArray(item)) {
accu.push(item.buffer);
}
return accu;
}, []);
return transferList;
}
module.exports = {
argsToString,
traceCallToString,
isBufferOrTypedArray,
extractTransferListFromParams
};

42
src/async-read-modes.js Normal file
View file

@ -0,0 +1,42 @@
const READ_BUF_SIZE = 16 * 1024;
/**
* Will read at least max number of bytes from SRT socket in async loop.
*
* Returns Promise of array of buffers.
*
* @param {AsyncSRT} asyncSrt
* @param {number} socketFd
* @param {number} minBytesRead
* @param {Function} onRead
* @param {Function} onError
* @returns {Promise<Uint8Array[]>}
*/
async function readChunks(asyncSrt, socketFd, minBytesRead, readBufSize = READ_BUF_SIZE,
onRead = null, onError = null) {
let bytesRead = 0;
const chunks = [];
while (bytesRead < minBytesRead) {
const readReturn = await asyncSrt.read(socketFd, readBufSize);
if (readReturn instanceof Uint8Array) {
const readBuf = readReturn;
bytesRead += readBuf.byteLength;
if (onRead) {
onRead(readBuf);
}
chunks.push(readBuf);
} else if (result === SRT.ERROR || result === null) {
if (onError) {
onError(result);
}
} else {
throw new Error('Got unexpected read-result')
}
}
return chunks;
}
module.exports = {
READ_BUF_SIZE,
readChunks
}

View file

@ -0,0 +1,70 @@
const EventEmitter = require("events");
const {
writeChunksWithYieldingLoop
} = require('../src/async-write-modes');
const {
READ_BUF_SIZE,
readChunks
} = require('../src/async-read-modes');
const DEFAULT_MTU_SIZE = 1316; // (for writes) should be the maximum on all IP networks cases
const DEFAULT_WRITES_PER_TICK = 128; // tbi
const DEFAULT_READ_BUFFER = READ_BUF_SIZE; // typical stream buffer size read in Node-JS internals
class AsyncReaderWriter {
constructor(asyncSrt, socketFd) {
this._asyncSrt = asyncSrt;
this._fd = socketFd;
}
/**
*
* @param {Uint8Array | Buffer} buffer
* @param {number} writesPerTick
* @param {number} mtuSize
* @param {Function} onWrite
* @returns {Promise<void>}
*/
async writeChunks(buffer,
writesPerTick = DEFAULT_WRITES_PER_TICK,
mtuSize = DEFAULT_MTU_SIZE,
onWrite = null) {
const chunks = sliceBufferToChunks(buffer, mtuSize,
buffer.byteLength, 0);
return writeChunksWithYieldingLoop(this._asyncSrt, this._fd, chunks,
onWrite, writesPerTick);
}
/**
* Will read at least a number of bytes from SRT socket in async loop.
*
* Returns Promise on array of buffers.
*
* The amount read (sum of bytes of array of buffers returned)
* may differ (exceed min bytes) by less than one MTU size.
*
* @param {number} minBytesRead
* @param {Function} onRead
* @param {Function} onError
* @param {number} readBufSize
* @returns {Promise<Uint8Array[]>}
*/
async readChunks(minBytesRead = DEFAULT_MTU_SIZE,
readBufSize = DEFAULT_READ_BUFFER,
onRead = null,
onError = null) {
return readChunks(this._asyncSrt, this._fd, minBytesRead, readBufSize, onRead, onError)
}
}
module.exports = {
AsyncReaderWriter,
DEFAULT_MTU_SIZE,
DEFAULT_WRITES_PER_TICK,
DEFAULT_READ_BUFFER
}

View file

@ -1,41 +1,79 @@
const {
Worker, isMainThread, parentPort, workerData
isMainThread, parentPort
} = require('worker_threads');
const debug = require('debug')('srt-async-worker');
const { SRT } = require('../build/Release/node_srt.node');
const { argsToString, traceCallToString, extractTransferListFromParams } = require('./async-helpers');
const DEBUG = false;
const DRY_RUN = false;
if (isMainThread) {
throw new Error("Worker module can not load on main thread");
}
(function run() {
const libSRT = new SRT();
try {
run()
} catch(err) {
console.error('AsyncSRT task-runner internal exception:', err);
}
function run() {
DEBUG && debug('AsyncSRT: Launching task-runner');
const srtNapiObjw = new SRT();
DEBUG && debug('AsyncSRT: SRT native object-wrap created');
parentPort.on('close', () => {
DEBUG && debug('AsyncSRT: Closing task-runner');
})
parentPort.on('message', (data) => {
if (!data.method) {
throw new Error('Worker message needs `method` property');
}
/*
if (!data.workId) {
throw new Error('Worker message needs `workId` property');
}
*/
let result = libSRT[data.method].apply(libSRT, data.args);
// TODO: see if we can do this using SharedArrayBuffer for example,
// or just leveraging Transferable objects capabilities ... ?
// FIXME: Performance ... ?
if (result instanceof Buffer) {
const buf = Buffer.allocUnsafe(result.length);
result.copy(buf);
result = buf;
if (data.args.some((arg) => arg === undefined)) {
const err = new Error(
`Ignoring call: Can't have any arguments be undefined: ${argsToString(data.args)}`);
parentPort.postMessage(err);
return;
}
DEBUG && debug('Received call:', traceCallToString(data.method, data.args));
let result = 0;
if (!DRY_RUN) {
try {
result = srtNapiObjw[data.method].apply(srtNapiObjw, data.args);
} catch(err) {
console.error(
`Exception thrown by native binding call "${traceCallToString(data.method, data.args)}":`,
err);
parentPort.postMessage({err, call: data});
return;
}
}
const transferList = extractTransferListFromParams([result]);
parentPort.postMessage({
// workId: data.workId,
timestamp: data.timestamp,
result
});
}, transferList);
});
})();
}

204
src/async-write-modes.js Normal file
View file

@ -0,0 +1,204 @@
const { SRT } = require('../build/Release/node_srt.node');
/**
* @module async-write-modes
*
* @author Stephan Hesse <stephan@emliri.com>
* @copyright EMLIRI, Stephan Hesse (c) 2020
*
* Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
*
*/
/**
*
* @function
*
* Description:
*
* This function allows to dispatch in high-performing mode a number
* of sequential `write()` calls to an `AsyncSRT` instance for a given `socketFd`.
*
* Specifications:
*
* - It allows for burst-writes to achieve throughput peaks as necessary in latency-critical
* applications.
*
* - The function is passed an array of data-buffer slices
* (that need to be individually "neuterable" or referencing a shared memory,
* see `SharedArrayBuffer` in JavaScript language documentations),
* which may be referred to packets.
*
* - The chunks size can not exceed the system or network
* specific MTU which the underlying SRT write binding will be able to accept.
*
* - It operates fully non-blocking. The main-thread is only used for queuing write-calls,
* and the maximum occupation time per iteration for queuing calls is parameterizable
* (see `writesPerTick`). If you figure that main-thread slots are too long when calling this
* function, consider adapting that parameter (which will cause scheduling overhead however
* to how this function runs).
*
* - When using `ArrayBuffer`, i.e not using `SharedArrayBuffer` as a memory implementation:
* It expects an array of buffers (chunks), which are referencing independent non
* overlapping memory buffers, which is critical here as there will be no copy performed,
* but each underlying buffer that is passed to the write-method will get neutered
* by the worker-thread it is passed to.
*
* Performance:
*
* This note below is here mainly to explain how
* this mode could potentially result in a different runtime-internal execution
* (and performance) then another implementation, the "explict scheduling" one.
*
* What we call a "Yielding Loop" will attempt to use any task-slot available by the runtime
* for queing writes (i.e transferring payload ownership to a worker),
* while ensuring a limited number of calls per main-loop (that is corresponding to how the tick()
* function here is invoked).
*
* By it's usage of the async/await paradigm, it causes the runtime
* to internally use generators as promise-executor and implicitely involves a `yield`.
*
* "yielding" means, stepping in and out of the async function context
* at a given point while maintaining context, while stepping back
* in "asap" when the awaited promise gets resolved - here inside the while loop.
* In order to allow for that context preservation and stack resumage (via async/await),
* this will result in so-called generator-functions internally
* with the JavaScript runtime. Now, depending how all that is implemented,
* there can be an overhead to doing that, in comparison with different,
* more hands-on "explicit" approaches to scheduling write calls,
* that are not using generators.
*
* @param {AsyncSRT} asyncSrt
* @param {number} socketFd
* @param {Array<Uint8Array>} chunks
* @param {Function} onWrite
* @param {number} writesPerTick
* @returns {Promise<void>}
*
*/
async function writeChunksWithYieldingLoop(asyncSrt, socketFd, chunks,
onWrite = null, writesPerTick = 1) {
let chunkIndex = 0;
let chunkWrittenIdx = 0;
while(chunkIndex < chunks.length) await tick();
function tick() {
const writeResultPromises = [];
for (let i = 0; i < writesPerTick; i++) {
if(chunkIndex >= chunks.length) {
break;
}
const chunkBuf = chunks[chunkIndex++];
const whenWritten = asyncSrt.write(socketFd, chunkBuf);
writeResultPromises.push(whenWritten);
whenWritten.then((writeRes) => {
if (writeRes === SRT.ERROR) {
throw new Error('AsyncSRT.write() failed');
}
if (onWrite) {
onWrite(writeRes, chunkWrittenIdx);
}
chunkWrittenIdx++
});
}
return Promise.all(writeResultPromises);
}
}
/**
*
* This is (almost) isofunctional to the yielding-loop (when intervalMs = 0),
* but implemented without using the runtimes generator-function support. Instead,
* we explicitely schedule all write calls using plain-old main-loop timers.
* When we do that setting "immediate tasks" by using a zero-timeout value,
* it will result in the exact same behavior from the perspective of the runtime,
* to use any available slot potentially (as the "yielding loop" mode).
*
* However the yielding-loop "awaits" until all writes that had been
* dispatched are resolved, while here we just keep pressuring the event queue.
*
* TODO: Implement rescheduling based on write-resolution (optional).
*
* The clear advantage of this is the explict nature of scheduling,
* allow pace of calls to be throttled by the set interval value.
* This can be very useful when we can afford to write data not asap,
* but in a given minimal rate instead of bursts and avoid peaking CPU with that.
*
* When using zero as interval timeout, in principle this should perform
* almost exactly like the generator-based mode (using any available
* slot for running a tick()).
*
* But the explicit scheduling has in theory also the advantage to have
* less runtime overhead (as it does not use any generator-function based features).
* However, this may depend on how tasks get prioritized in the end also,
* or how well the runtime is optimized and implemented for one or the other.
* Maybe an await-resolution gets more attention to appear on the main-loop
* than the scheduled interval, if it has immediate timeout (i.e 0), or the overhead
* caused by yielding generators is neglectable.
*
* @param {AsyncSRT} asyncSrt
* @param {number} socketFd
* @param {Array<Uint8Array>} chunks
* @param {Function} onWrite
* @param {number} writesPerTick
* @param {number} intervalMs
*/
function writeChunksWithExplicitScheduling(asyncSrt, socketFd, chunks,
onWrite = null, writesPerTick = 1, intervalMs = 0) {
let chunkIndex = 0;
let chunkWrittenIdx = 0;
// schedule tick-interval
const writeTimer = setInterval(tick, intervalMs);
// run once directly in current stack
tick();
function tick() {
for (let i = 0; i < writesPerTick; i++) {
if(chunkIndex >= chunks.length) {
clearInterval(writeTimer);
break;
}
const chunkBuf = chunks[chunkIndex++];
asyncSrt.write(socketFd, chunkBuf)
.then((writeRes) => {
if (writeRes === SRT.ERROR) {
throw new Error('AsyncSRT.write() failed');
}
if (onWrite) {
onWrite(writeRes, chunkWrittenIdx);
}
chunkWrittenIdx++
});
}
}
}
module.exports = {
writeChunksWithYieldingLoop,
writeChunksWithExplicitScheduling
}

View file

@ -1,12 +1,16 @@
const {
Worker, isMainThread, parentPort, workerData
} = require('worker_threads');
const { Worker } = require('worker_threads');
const path = require('path');
const {performance} = require("perf_hooks");
const debug = require('debug')('srt-async');
const { traceCallToString, extractTransferListFromParams } = require('./async-helpers');
const { SRT } = require('../build/Release/node_srt.node');
const DEFAULT_PROMISE_TIMEOUT_MS = 3000;
const DEBUG = false;
/*
const WORK_ID_GEN_MOD = 0xFFF;
*/
@ -20,6 +24,9 @@ class AsyncSRT {
static TimeoutMs = DEFAULT_PROMISE_TIMEOUT_MS;
constructor() {
DEBUG && debug('Creating task-runner worker instance');
this._worker = new Worker(path.resolve(__dirname, './async-worker.js'));
this._worker.on('message', this._onWorkerMessage.bind(this));
/*
@ -29,15 +36,41 @@ class AsyncSRT {
this._workCbQueue = [];
}
/**
* @returns {Promise<number>} Resolves to exit code of Worker
*/
dispose() {
const worker = this._worker;
this._worker = null;
if (this._workCbQueue.length !== 0) {
console.warn(`AsyncSRT: flushing callback-queue with ${this._workCbQueue.length} remaining jobs awaiting.`);
this._workCbQueue.length = 0;
}
return worker.terminate();
}
/**
* @private
* @param {*} data
* @param {object} data
*/
_onWorkerMessage(data) {
const resolveTime = performance.now();
const {timestamp, result, workId} = data;
// not sure if there can still be message event
// after calling terminate
// but let's guard from that state anyway.
if (this._worker === null) return;
const resolveTime = performance.now();
const callback = this._workCbQueue.shift();
if (data.err) {
console.error('AsyncSRT: Error from task-runner:', data.err.message,
'\n Binding call:', traceCallToString(data.call.method, data.call.args),
//'\n Stacktrace:', data.err.stack
);
return;
}
const {timestamp, result, workId} = data;
callback(result);
}
@ -63,8 +96,12 @@ class AsyncSRT {
this._workCbMap.set(workId, callback);
*/
DEBUG && debug('Sending call:', traceCallToString(method, args));
const transferList = extractTransferListFromParams(args);
this._workCbQueue.push(callback);
this._worker.postMessage({method, args, /*workId,*/ timestamp});
this._worker.postMessage({method, args, /*workId,*/ timestamp}, transferList);
}
/**
@ -72,8 +109,15 @@ class AsyncSRT {
* @param {string} method
* @param {Array<any>} args optional
* @param {Function} callback optional
* @param {boolean} useTimeout
* @param {number} timeoutMs
*/
_createAsyncWorkPromise(method, args = [], callback = null, useTimeout = true, timeoutMs = AsyncSRT.TimeoutMs) {
_createAsyncWorkPromise(method,
args = [],
callback = null,
useTimeout = false,
timeoutMs = AsyncSRT.TimeoutMs) {
return new Promise((resolve, reject) => {
let timeout;
let rejected = false;
@ -91,7 +135,7 @@ class AsyncSRT {
};
if (useTimeout) {
timeout = setTimeout(() => {
reject(new Error('Timeout exceeded while awaiting result from worker running native-addon module functions'));
reject(new Error(`Timeout exceeded (${timeoutMs} ms) while awaiting method result: ${traceCallToString(method, args)}`));
rejected = true;
}, timeoutMs);
}
@ -101,18 +145,17 @@ class AsyncSRT {
/**
*
* @param {boolean} sender
* @returns SRTSOCKET identifier (integer value) or -1 (SRT_ERROR)
* @param {boolean} sender default: false. only needed to specify if local/remote SRT ver < 1.3 or no other HSv5 support
*/
createSocket(sender, callback) {
createSocket(sender = false, callback) {
return this._createAsyncWorkPromise("createSocket", [sender], callback);
}
/**
*
* @param socket
* @param address
* @param port
* @param {number} socket
* @param {string} address
* @param {number} port
*/
bind(socket, address, port, callback) {
return this._createAsyncWorkPromise("bind", [socket, address, port], callback);
@ -120,8 +163,8 @@ class AsyncSRT {
/**
*
* @param socket
* @param backlog
* @param {number} socket
* @param {number} backlog
*/
listen(socket, backlog, callback) {
return this._createAsyncWorkPromise("listen", [socket, backlog], callback);
@ -129,9 +172,9 @@ class AsyncSRT {
/**
*
* @param socket
* @param host
* @param port
* @param {number} socket
* @param {string} host
* @param {number} port
*/
connect(socket, host, port, callback) {
return this._createAsyncWorkPromise("connect", [socket, host, port], callback);
@ -139,8 +182,7 @@ class AsyncSRT {
/**
*
* @param socket
* @returns File descriptor of incoming connection pipe
* @param {number} socket
*/
accept(socket, callback, useTimeout = false, timeoutMs = AsyncSRT.TimeoutMs) {
return this._createAsyncWorkPromise("accept", [socket], callback, useTimeout, timeoutMs);
@ -148,7 +190,7 @@ class AsyncSRT {
/**
*
* @param socket
* @param {number} socket
*/
close(socket, callback) {
return this._createAsyncWorkPromise("close", [socket], callback);
@ -156,9 +198,9 @@ class AsyncSRT {
/**
*
* @param socket
* @param chunkSize
* @returns {Promise<Buffer>}
* @param {number} socket
* @param {number} chunkSize
* @returns {Promise<Buffer | SRTResult.SRT_ERROR | null>}
*/
read(socket, chunkSize, callback) {
return this._createAsyncWorkPromise("read", [socket, chunkSize], callback);
@ -166,24 +208,45 @@ class AsyncSRT {
/**
*
* @param socket
* @param {Buffer} chunk
* Pass a packet buffer to write to the socket.
*
* The size of the buffer must not exceed the SRT payload MTU
* (usually 1316 bytes).
*
* Otherwise the call will resolve to SRT_ERROR.
*
* A system-specific socket-message error message may show in logs as enabled
* where the error is thrown (on the binding call to the native SRT API),
* and in the async API internals as it gets propagated back from the task-runner).
*
* Note that any underlying data buffer passed in
* will be *neutered* by our worker thread and
* therefore become unusable (i.e go to detached state, `byteLengh === 0`)
* for the calling thread of this method.
* When consuming from a larger piece of data,
* chunks written will need to be slice copies of the source buffer.
*
* For a usage example, check the performance & smoke testbench.
*
* @param {number} socket Socket identifier to write to
* @param {Buffer | Uint8Array} chunk The underlying `buffer` (ArrayBufferLike) will get "neutered" by creating the async task. Pass in or use a copy respectively if concurrent data usage is intended.
*/
write(socket, chunk, callback) {
// TODO: see if we can do this using SharedArrayBuffer for example,
// or just leveraging Transferable objects capabilities ... ?
// FIXME: Performance ... ?
const buf = Buffer.allocUnsafe(chunk.length);
chunk.copy(buf);
chunk = buf;
return this._createAsyncWorkPromise("write", [socket, chunk], callback);
const byteLength = chunk.byteLength;
DEBUG && debug(`write ${byteLength} to socket:`, socket)
return this._createAsyncWorkPromise("write", [socket, chunk], callback)
.then((result) => {
if (result !== SRT.ERROR) {
return byteLength;
}
});
}
/**
*
* @param socket
* @param option
* @param value
* @param {number} socket
* @param {number} option
* @param {number} value
*/
setSockOpt(socket, option, value, callback) {
return this._createAsyncWorkPromise("setSockOpt", [socket, option, value], callback);
@ -191,8 +254,8 @@ class AsyncSRT {
/**
*
* @param socket
* @param option
* @param {number} socket
* @param {number} option
*/
getSockOpt(socket, option, callback) {
return this._createAsyncWorkPromise("getSockOpt", [socket, option], callback);
@ -200,14 +263,14 @@ class AsyncSRT {
/**
*
* @param socket
* @param {number} socket
*/
getSockState(socket, callback) {
return this._createAsyncWorkPromise("getSockState", [socket], callback);
}
/**
* @returns epid
* @returns {number} epid
*/
epollCreate(callback) {
return this._createAsyncWorkPromise("epollCreate", [], callback);
@ -215,9 +278,9 @@ class AsyncSRT {
/**
*
* @param epid
* @param socket
* @param events
* @param {number} epid
* @param {number} socket
* @param {number} events
*/
epollAddUsock(epid, socket, events, callback) {
return this._createAsyncWorkPromise("epollAddUsock", [epid, socket, events], callback);
@ -225,12 +288,21 @@ class AsyncSRT {
/**
*
* @param epid
* @param msTimeOut
* @param {number} epid
* @param {number} msTimeOut
*/
epollUWait(epid, msTimeOut, callback) {
return this._createAsyncWorkPromise("epollUWait", [epid, msTimeOut], callback);
}
/**
*
* @param {number | SRTLoggingLevel} logLevel
* @returns {Promise<SRTResult>}
*/
setLogLevel(logLevel, callback) {
return this._createAsyncWorkPromise("setLogLevel", [logLevel], callback);
}
}
module.exports = {AsyncSRT};

18
src/logging.js Normal file
View file

@ -0,0 +1,18 @@
const {SRT} = require('../build/Release/node_srt.node');
let srt = null;
/**
*
* @param {number | SRTLoggingLevel} level
*/
function setSRTLoggingLevel(level) {
if (!srt) {
srt = new SRT();
}
srt.setLogLevel(level)
}
module.exports = {
setSRTLoggingLevel
}

View file

@ -8,6 +8,10 @@
#include "node-srt.h"
#include "srt-enums.h"
using namespace std;
#define EPOLL_EVENTS_NUM_MAX 1024
Napi::FunctionReference NodeSRT::constructor;
Napi::Object NodeSRT::Init(Napi::Env env, Napi::Object exports) {
@ -28,13 +32,15 @@ Napi::Object NodeSRT::Init(Napi::Env env, Napi::Object exports) {
InstanceMethod("epollCreate", &NodeSRT::EpollCreate),
InstanceMethod("epollAddUsock", &NodeSRT::EpollAddUsock),
InstanceMethod("epollUWait", &NodeSRT::EpollUWait),
InstanceMethod("setLogLevel", &NodeSRT::SetLogLevel),
StaticValue("OK", Napi::Number::New(env, 0)),
StaticValue("ERROR", Napi::Number::New(env, SRT_ERROR)),
StaticValue("INVALID_SOCK", Napi::Number::New(env, SRT_INVALID_SOCK)),
// Socket options
SOCKET_OPTIONS,
// Socket status
SOCKET_STATUS,
@ -56,10 +62,12 @@ NodeSRT::NodeSRT(const Napi::CallbackInfo& info) : Napi::ObjectWrap<NodeSRT>(inf
Napi::Env env = info.Env();
Napi::HandleScope scope(env);
// Q: should we avoid to call this repeatedly (with potentially several ObjectWrap instances created) ?
srt_startup();
}
NodeSRT::~NodeSRT() {
srt_cleanup();
}
@ -69,9 +77,10 @@ Napi::Value NodeSRT::CreateSocket(const Napi::CallbackInfo& info) {
Napi::Boolean isSender = Napi::Boolean::New(env, false);
if (info.Length() > 0) {
// FIXME: throws exception when `arg[0] === undefined`
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();
@ -88,7 +97,7 @@ Napi::Value NodeSRT::Bind(const Napi::CallbackInfo& info) {
Napi::Env env = info.Env();
Napi::HandleScope scope(env);
Napi::Number socketValue = info[0].As<Napi::Number>();
Napi::Number socketId = info[0].As<Napi::Number>();
Napi::String address = info[1].As<Napi::String>();
Napi::Number port = info[2].As<Napi::Number>();
@ -102,9 +111,9 @@ Napi::Value NodeSRT::Bind(const Napi::CallbackInfo& info) {
return Napi::Number::New(env, result);
}
result = srt_bind(socketValue, (struct sockaddr *)&addr, sizeof(addr));
result = srt_bind(socketId, (struct sockaddr *)&addr, sizeof(addr));
if (result == SRT_ERROR) {
srt_close(socketValue);
srt_close(socketId);
Napi::Error::New(env, srt_getlasterror_str()).ThrowAsJavaScriptException();
return Napi::Number::New(env, SRT_ERROR);
}
@ -192,16 +201,20 @@ Napi::Value NodeSRT::Read(const Napi::CallbackInfo& info) {
Napi::Number socketValue = info[0].As<Napi::Number>();
Napi::Number chunkSize = info[1].As<Napi::Number>();
// Q: why not converting to `int` directly here?
size_t bufferSize = uint32_t(chunkSize);
uint8_t *buffer = (uint8_t *)malloc(bufferSize);
memset(buffer, 0, bufferSize);
int nb = srt_recvmsg(socketValue, (char *)buffer, (int)bufferSize);
if (nb == SRT_ERROR) {
Napi::Error::New(env, srt_getlasterror_str()).ThrowAsJavaScriptException();
string err(string("srt_recvmsg: ")
+ string(srt_getlasterror_str()));
Napi::Error::New(env, err).ThrowAsJavaScriptException();
return Napi::Number::New(env, SRT_ERROR);
}
// Q: why not using char as data/template type?
Napi::Value nbuff = Napi::Buffer<uint8_t>::Copy(env, buffer, nb);
free(buffer);
@ -213,11 +226,15 @@ Napi::Value NodeSRT::Write(const Napi::CallbackInfo& info) {
Napi::HandleScope scope(env);
Napi::Number socketValue = info[0].As<Napi::Number>();
// Q: why not using char as data/template type?
Napi::Buffer<uint8_t> chunk = info[1].As<Napi::Buffer<uint8_t>>();
int result = srt_sendmsg2(socketValue, (const char *)chunk.Data(), chunk.Length(), nullptr);
if (result == SRT_ERROR) {
Napi::Error::New(env, srt_getlasterror_str()).ThrowAsJavaScriptException();
string err(string("srt_sendmsg2: ")
+ string(srt_getlasterror_str()));
Napi::Error::New(env, err).ThrowAsJavaScriptException();
return Napi::Number::New(env, SRT_ERROR);
}
return Napi::Number::New(env, result);
@ -337,7 +354,7 @@ Napi::Value NodeSRT::GetSockOpt(const Napi::CallbackInfo& info) {
Napi::Error::New(env, "SOCKOPT not implemented yet").ThrowAsJavaScriptException();
break;
}
if (result == SRT_ERROR) {
Napi::Error::New(env, srt_getlasterror_str()).ThrowAsJavaScriptException();
return empty;
@ -385,11 +402,11 @@ Napi::Value NodeSRT::EpollAddUsock(const Napi::CallbackInfo& info) {
Napi::Value NodeSRT::EpollUWait(const Napi::CallbackInfo& info) {
Napi::Env env = info.Env();
Napi::HandleScope scope(env);
Napi::Number epidValue = info[0].As<Napi::Number>();
Napi::Number msTimeOut = info[1].As<Napi::Number>();
const int fdsSetSize = 100;
const int fdsSetSize = EPOLL_EVENTS_NUM_MAX;
SRT_EPOLL_EVENT fdsSet[fdsSetSize];
int n = srt_epoll_uwait(epidValue, fdsSet, fdsSetSize, msTimeOut);
Napi::Array events = Napi::Array::New(env, n);
@ -402,3 +419,19 @@ Napi::Value NodeSRT::EpollUWait(const Napi::CallbackInfo& info) {
return events;
}
Napi::Value NodeSRT::SetLogLevel(const Napi::CallbackInfo& info) {
Napi::Env env = info.Env();
Napi::HandleScope scope(env);
int logLevel = info[0].As<Napi::Number>().Int32Value();
int result;
if (logLevel >= 0 && logLevel <= 7) {
srt_setloglevel(logLevel);
result = 0;
} else {
result = SRT_ERROR;
}
return Napi::Number::New(env, result);
}

View file

@ -23,4 +23,6 @@ class NodeSRT : public Napi::ObjectWrap<NodeSRT> {
Napi::Value EpollCreate(const Napi::CallbackInfo& info);
Napi::Value EpollAddUsock(const Napi::CallbackInfo& info);
Napi::Value EpollUWait(const Napi::CallbackInfo& info);
Napi::Value SetLogLevel(const Napi::CallbackInfo& info);
};

View file

@ -1,30 +0,0 @@
const { SRT } = require('../build/Release/node_srt.node');
const EventEmitter = require("events");
const debug = require('debug')('srt-server');
const libSRT = new SRT();
class SRTServer extends EventEmitter {
constructor() {
super();
this.iface = null;
this.port = null;
this.socket = libSRT.createSocket();
}
listen({ address, port }) {
const iface = address || "0.0.0.0";
libSRT.bind(this.socket, iface, port);
libSRT.listen(this.socket, 2);
this.iface = address;
this.port = port;
this.emit("listening", iface, port);
while (true) {
const fhandle = libSRT.accept(this.socket);
debug("Client connection accepted");
this.emit("accepted", fhandle);
}
}
}
module.exports = SRTServer;

108
src/srt-api-enums.ts Normal file
View file

@ -0,0 +1,108 @@
export enum SRTSockOpt {
SRTO_MSS = 0, // the Maximum Transfer Unit
SRTO_SNDSYN = 1, // if sending is blocking
SRTO_RCVSYN = 2, // if receiving is blocking
SRTO_ISN = 3, // Initial Sequence Number (valid only after srt_connect or srt_accept-ed sockets)
SRTO_FC = 4, // Flight flag size (window size)
SRTO_SNDBUF = 5, // maximum buffer in sending queue
SRTO_RCVBUF = 6, // UDT receiving buffer size
SRTO_LINGER = 7, // waiting for unsent data when closing
SRTO_UDP_SNDBUF = 8, // UDP sending buffer size
SRTO_UDP_RCVBUF = 9, // UDP receiving buffer size
// XXX Free space for 2 options
// after deprecated ones are removed
SRTO_RENDEZVOUS = 12, // rendezvous connection mode
SRTO_SNDTIMEO = 13, // send() timeout
SRTO_RCVTIMEO = 14, // recv() timeout
SRTO_REUSEADDR = 15, // reuse an existing port or create a new one
SRTO_MAXBW = 16, // maximum bandwidth (bytes per second) that the connection can use
SRTO_STATE = 17, // current socket state, see UDTSTATUS, read only
SRTO_EVENT = 18, // current available events associated with the socket
SRTO_SNDDATA = 19, // size of data in the sending buffer
SRTO_RCVDATA = 20, // size of data available for recv
SRTO_SENDER = 21, // Sender mode (independent of conn mode), for encryption, tsbpd handshake.
SRTO_TSBPDMODE = 22, // Enable/Disable TsbPd. Enable -> Tx set origin timestamp, Rx deliver packet at origin time + delay
SRTO_LATENCY = 23, // NOT RECOMMENDED. SET: to both SRTO_RCVLATENCY and SRTO_PEERLATENCY. GET: same as SRTO_RCVLATENCY.
SRTO_TSBPDDELAY = 23, // DEPRECATED. ALIAS: SRTO_LATENCY
SRTO_INPUTBW = 24, // Estimated input stream rate.
SRTO_OHEADBW, // MaxBW ceiling based on % over input stream rate. Applies when UDT_MAXBW=0 (auto).
SRTO_PASSPHRASE = 26, // Crypto PBKDF2 Passphrase size[0,10..64] 0:disable crypto
SRTO_PBKEYLEN, // Crypto key len in bytes {16,24,32} Default: 16 (128-bit)
SRTO_KMSTATE, // Key Material exchange status (UDT_SRTKmState)
SRTO_IPTTL = 29, // IP Time To Live (passthru for system sockopt IPPROTO_IP/IP_TTL)
SRTO_IPTOS, // IP Type of Service (passthru for system sockopt IPPROTO_IP/IP_TOS)
SRTO_TLPKTDROP = 31, // Enable receiver pkt drop
SRTO_SNDDROPDELAY = 32, // Extra delay towards latency for sender TLPKTDROP decision (-1 to off)
SRTO_NAKREPORT = 33, // Enable receiver to send periodic NAK reports
SRTO_VERSION = 34, // Local SRT Version
SRTO_PEERVERSION, // Peer SRT Version (from SRT Handshake)
SRTO_CONNTIMEO = 36, // Connect timeout in msec. Ccaller default: 3000, rendezvous (x 10)
// deprecated: SRTO_TWOWAYDATA, SRTO_SNDPBKEYLEN, SRTO_RCVPBKEYLEN (@c below)
_DEPRECATED_SRTO_SNDPBKEYLEN = 38, // (needed to use inside the code without generating -Wswitch)
//
SRTO_SNDKMSTATE = 40, // (GET) the current state of the encryption at the peer side
SRTO_RCVKMSTATE, // (GET) the current state of the encryption at the agent side
SRTO_LOSSMAXTTL, // Maximum possible packet reorder tolerance (number of packets to receive after loss to send lossreport)
SRTO_RCVLATENCY, // TsbPd receiver delay (mSec) to absorb burst of missed packet retransmission
SRTO_PEERLATENCY, // Minimum value of the TsbPd receiver delay (mSec) for the opposite side (peer)
SRTO_MINVERSION, // Minimum SRT version needed for the peer (peers with less version will get connection reject)
SRTO_STREAMID, // A string set to a socket and passed to the listener's accepted socket
SRTO_CONGESTION, // Congestion controller type selection
SRTO_MESSAGEAPI, // In File mode, use message API (portions of data with boundaries)
SRTO_PAYLOADSIZE, // Maximum payload size sent in one UDP packet (0 if unlimited)
SRTO_TRANSTYPE = 50, // Transmission type (set of options required for given transmission type)
SRTO_KMREFRESHRATE, // After sending how many packets the encryption key should be flipped to the new key
SRTO_KMPREANNOUNCE, // How many packets before key flip the new key is annnounced and after key flip the old one decommissioned
SRTO_ENFORCEDENCRYPTION, // Connection to be rejected or quickly broken when one side encryption set or bad password
SRTO_IPV6ONLY, // IPV6_V6ONLY mode
SRTO_PEERIDLETIMEO, // Peer-idle timeout (max time of silence heard from peer) in [ms]
// (some space left)
SRTO_PACKETFILTER = 60 // Add and configure a packet filter
}
export enum SRTEpollOpt
{
SRT_EPOLL_OPT_NONE = 0x0, // fallback
// this values are defined same as linux epoll.h
// so that if system values are used by mistake, they should have the same effect
SRT_EPOLL_IN = 0x1,
SRT_EPOLL_OUT = 0x4,
SRT_EPOLL_ERR = 0x8,
SRT_EPOLL_ET = 0x80000000 // (2147483648) C: 1u << 31
};
export enum SRTSockStatus {
SRTS_INIT = 1,
SRTS_OPENED,
SRTS_LISTENING,
SRTS_CONNECTING,
SRTS_CONNECTED,
SRTS_BROKEN,
SRTS_CLOSING,
SRTS_CLOSED,
SRTS_NONEXIST
}
export enum SRTResult {
SRT_ERROR = -1,
SRT_OK = 0
}
/**
* Enum values as taken from native SRT logging API declarations
*/
export enum SRTLoggingLevel {
FATAL = 2,
// Fatal vs. Error: with Error, you can still continue.
ERROR = 3,
// Error vs. Warning: Warning isn't considered a problem for the library.
WARNING = 4,
// Warning vs. Note: Note means something unusual, but completely correct behavior.
NOTE = 5,
// Note vs. Debug: Debug may occur even multiple times in a millisecond.
// (Well, worth noting that Error and Warning potentially also can).
DEBUG = 7
}

322
src/srt-server.js Normal file
View file

@ -0,0 +1,322 @@
const { AsyncSRT } = require('./async');
const { AsyncReaderWriter } = require('./async-reader-writer');
const { SRT } = require('../build/Release/node_srt.node');
const EventEmitter = require("events");
const debug = require('debug')('srt-server');
const DEBUG = false;
const EPOLL_PERIOD_MS_DEFAULT = 0;
const EPOLLUWAIT_TIMEOUT_MS = 0;
const SOCKET_LISTEN_BACKLOG = 128;
/**
* @emits data
* @emits closing
* @emits closed
*/
class SRTConnection extends EventEmitter {
/**
*
* @param {AsyncSRT} asyncSrt
* @param {number} fd
*/
constructor(asyncSrt, fd) {
super();
this._asyncSrt = asyncSrt;
this._fd = fd;
this._gotFirstData = false;
}
/**
* @returns {number}
*/
get fd() {
return this._fd;
}
/**
* Will be false until *after* emit of first `data` event.
* After that will be true.
*/
get gotFirstData() {
return this._gotFirstData;
}
/**
* @returns {AsyncReaderWriter}
*/
getReaderWriter() {
return new AsyncReaderWriter(this._asyncSrt, this.fd);
}
/**
*
* @param {number} bytes
* @returns {Promise<Buffer | SRTResult.SRT_ERROR | null>}
*/
async read(bytes) {
return await this._asyncSrt.read(this.fd, bytes);
}
/**
*
* Pass a packet buffer to write to the connection.
*
* The size of the buffer must not exceed the SRT payload MTU
* (usually 1316 bytes).
*
* Otherwise the call will resolve to SRT_ERROR.
*
* A system-specific socket-message error message may show in logs as enabled
* where the error is thrown (on the binding call to the native SRT API),
* and in the async API internals as it gets propagated back from the task-runner).
*
* Note that any underlying data buffer passed in
* will be *neutered* by our worker thread and
* therefore become unusable (i.e go to detached state, `byteLengh === 0`)
* for the calling thread of this method.
* When consuming from a larger piece of data,
* chunks written will need to be slice copies of the source buffer.
*
* @param {Buffer | Uint8Array} chunk
*/
async write(chunk) {
return await this._asyncSrt.write(this.fd, chunk);
}
/**
* @returns {Promise<SRTResult | null>}
*/
async close() {
if (this.isClosed()) return null;
const asyncSrt = this._asyncSrt;
this._asyncSrt = null;
this.emit('closing');
const result = await asyncSrt.close(this.fd);
this.emit('closed', result);
this.off();
return result;
}
isClosed() {
return ! this._asyncSrt;
}
onData() {
this.emit('data');
if (!this.gotFirstData) {
this._gotFirstData = true;
}
}
}
/**
* @emits created
* @emits opened
* @emits connection
* @emits disconnection
* @emits disposed
*/
class SRTServer extends EventEmitter {
/**
*
* @param {number} port socket port number
* @param {string} address optional, default: '0.0.0.0'
* @param {number} epollPeriodMs optional, default: EPOLL_PERIOD_MS_DEFAULT
* @returns {Promise<SRTServer>}
*/
static create(port, address, epollPeriodMs) {
return new SRTServer(port, address, epollPeriodMs).create();
}
/**
*
* @param {number} port socket port number
* @param {string} address optional, default: '0.0.0.0'
* @param {number} epollPeriodMs optional, default: EPOLL_PERIOD_MS_DEFAULT
*/
constructor(port, address = '0.0.0.0', epollPeriodMs = EPOLL_PERIOD_MS_DEFAULT) {
super();
if (!Number.isInteger(port) || port <= 0 || port > 65535)
throw new Error('Need a valid port number but got: ' + port);
this.port = port;
this.address = address;
this.epollPeriodMs = epollPeriodMs;
this.socket = null;
this.epid = null;
this._pollEventsTimer = null;
this._asyncSrt = new AsyncSRT();
this._connectionMap = {};
}
async dispose() {
clearTimeout(this._pollEventsTimer);
await this._asyncSrt.close(this.socket);
this.socket = null;
const res = await this._asyncSrt.dispose();
this._asyncSrt = null;
this.emit('disposed');
return res;
}
/**
* Call this before `open`.
* Call `setSocketFlags` after this.
*
* @return {Promise<SRTServer>}
*/
async create() {
this.socket = await this._asyncSrt.createSocket();
this.emit('created');
return this;
}
/**
* Call this after `create`.
* Call `setSocketFlags` before calling this.
*
* @return {Promise<SRTServer>}
*/
async open() {
let result;
result = await this._asyncSrt.bind(this.socket, this.address, this.port);
if (result === SRT.ERROR) {
throw new Error('SRT.bind() failed');
}
result = await this._asyncSrt.listen(this.socket, SOCKET_LISTEN_BACKLOG);
if (result === SRT.ERROR) {
throw new Error('SRT.listen() failed');
}
result = await this._asyncSrt.epollCreate();
if (result === SRT.ERROR) {
throw new Error('SRT.epollCreate() failed');
}
this.epid = result;
this.emit('opened');
// we should await the epoll subscribe result before continuing
// since it is useless to poll events otherwise
// and we should also yield from the stack at this point
// since the `opened` event handlers above may do whatever
await this._asyncSrt.epollAddUsock(this.epid, this.socket, SRT.EPOLL_IN | SRT.EPOLL_ERR);
this._pollEvents();
return this;
}
/**
*
* @param {SRTSockOpt[]} opts
* @param {SRTSockOptValue[]} values
* @returns {Promise<SRTResult[]>}
*/
async setSocketFlags(opts, values) {
if (opts.length !== values.length)
throw new Error('opts and values must have same length');
const promises = [];
opts.forEach((opt, index) => {
const p = this._asyncSrt.setSockOpt(this.socket, opt, values[index]);
promises.push(p);
})
return Promise.all(promises);
}
/**
*
* @param {number} fd
* @returns {SRTConnection | null}
*/
getConnectionByHandle(fd) {
return this._connectionMap[fd] || null;
}
/**
* @returns {Array<SRTConnection>}
*/
getAllConnections() {
return Array.from(Object.values(this._connectionMap));
}
/**
* @private
* @param {SRTEpollEvent} event
*/
async _handleEvent(event) {
const status = await this._asyncSrt.getSockState(event.socket);
// our local listener socket
if (event.socket === this.socket) {
if (status === SRT.SRTS_LISTENING) {
const fd = await this._asyncSrt.accept(this.socket);
// no need to await the epoll subscribe result before continuing
this._asyncSrt.epollAddUsock(this.epid, fd, SRT.EPOLL_IN | SRT.EPOLL_ERR);
debug("Accepted client connection with file-descriptor:", fd);
// create new client connection handle
// and emit accept event
const connection = new SRTConnection(this._asyncSrt, fd);
connection.on('closing', () => {
// remove handle
delete this._connectionMap[fd];
});
this._connectionMap[fd] = connection;
this.emit('connection', connection);
}
// a client socket / fd
// check if broken or closed
} else if (status === SRT.SRTS_BROKEN
|| status === SRT.SRTS_NONEXIST
|| status === SRT.SRTS_CLOSED) {
const fd = event.socket;
debug("Client disconnected on fd:", fd);
if (this._connectionMap[fd]) {
await this._connectionMap[fd].close();
this.emit('disconnection', fd);
}
// not broken, just new data
} else {
const fd = event.socket;
DEBUG && debug("Got data from connection on fd:", fd);
const connection = this.getConnectionByHandle(fd);
if (!connection) {
console.warn("Got event for fd not in connections map:", fd);
return;
}
connection.onData();
}
}
/**
* @private
*/
async _pollEvents() {
const events = await this._asyncSrt.epollUWait(this.epid, EPOLLUWAIT_TIMEOUT_MS);
events.forEach((event) => {
this._handleEvent(event);
});
// clearing in case we get called multiple times
// when already timer scheduled
// will be no-op if timer-id invalid or old
clearTimeout(this._pollEventsTimer);
this._pollEventsTimer
= setTimeout(this._pollEvents.bind(this), this.epollPeriodMs)
}
}
module.exports = {
SRTConnection,
SRTServer
};

View file

@ -1,6 +1,6 @@
const { Readable, Writable } = require('stream');
const LIB = require('../build/Release/node_srt.node');
const debug = require('debug')('srt-stream');
const { Readable } = require('stream');
const { SRT } = require('../build/Release/node_srt.node');
const debug = require('debug')('srt-read-stream');
const CONNECTION_ACCEPT_POLLING_INTERVAL_MS = 50;
const READ_WAIT_INTERVAL_MS = 50;
@ -24,10 +24,30 @@ class SRTReadStream extends Readable {
// Q: not better if port (mandatory) is before, and address is optional (default to "0.0.0.0")?
constructor(address, port, opts) {
super();
this.srt = new LIB.SRT();
/**
* @member {SRT}
*/
this.srt = new SRT();
/**
* @member {number}
*/
this.socket = this.srt.createSocket();
/**
* @member {string}
*/
this.address = address;
/**
* @member {number}
*/
this.port = port;
/**
* @member {number | null}
*/
this.fd = null;
this._eventPollInterval = null;
@ -39,28 +59,33 @@ class SRTReadStream extends Readable {
* @param {Function} onData Passes this stream instance as first arg to callback
*/
listen(onData) {
if (this.fd !== null) {
throw new Error('listen() called but stream file-descriptor already initialized');
}
this.srt.bind(this.socket, this.address, this.port);
this.srt.listen(this.socket, SOCKET_LISTEN_BACKLOG);
const epid = this.srt.epollCreate();
this.srt.epollAddUsock(epid, this.socket, LIB.SRT.EPOLL_IN | LIB.SRT.EPOLL_ERR);
this.srt.epollAddUsock(epid, this.socket, SRT.EPOLL_IN | SRT.EPOLL_ERR);
const interval = this._eventPollInterval = setInterval(() => {
const events = this.srt.epollUWait(epid, EPOLLUWAIT_TIMEOUT_MS);
events.forEach(event => {
const status = this.srt.getSockState(event.socket);
if (status === LIB.SRT.SRTS_BROKEN || status === LIB.SRT.SRTS_NONEXIST || status === LIB.SRT.SRTS_CLOSED) {
if (status === SRT.SRTS_BROKEN || status === SRT.SRTS_NONEXIST || status === SRT.SRTS_CLOSED) {
debug("Client disconnected with socket:", event.socket);
this.srt.close(event.socket);
this.push(null);
this.emit('end');
} else if (event.socket === this.socket) {
const fhandle = this.srt.accept(this.socket);
debug("New client connected with socket:", this.socket, "and fhandle:", fhandle);
this.srt.epollAddUsock(epid, fhandle, LIB.SRT.EPOLL_IN | LIB.SRT.EPOLL_ERR);
debug("Accepted client connection with file-descriptor:", fhandle);
this.srt.epollAddUsock(epid, fhandle, SRT.EPOLL_IN | SRT.EPOLL_ERR);
this.emit('readable');
} else {
debug("Data from client on fd:", event.socket);
debug("Got data from connection on fd:", event.socket);
this.fd = event.socket;
clearInterval(interval);
onData(this);
@ -70,11 +95,20 @@ class SRTReadStream extends Readable {
}, CONNECTION_ACCEPT_POLLING_INTERVAL_MS);
}
connect(cb) {
/**
*
* @param {Function} onConnect
*/
connect(onConnect) {
if (this.fd !== null) {
throw new Error('connect() called but stream file-descriptor already initialized');
}
this.srt.connect(this.socket, this.address, this.port);
this.fd = this.socket;
if (this.fd) {
cb(this);
onConnect(this);
}
}
@ -148,48 +182,10 @@ class SRTReadStream extends Readable {
this.srt.close(this.socket);
this.fd = null;
this._clearScheduledRead();
this.emit('close');
}
}
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);
}
}
close() {
this.srt.close(this.socket);
this.fd = null;
}
_write(chunk, encoding, callback) {
debug(`Writing chunk ${chunk.length}`);
if (this.fd) {
this.srt.write(this.fd, chunk);
callback();
} else {
callback(new Error("Socket was closed"));
}
}
_destroy(err, callback) {
this.close();
if (cb) cb(err);
}
}
module.exports = {
SRTReadStream,
SRTWriteStream
SRTReadStream
};

View file

@ -0,0 +1,44 @@
const { Writable } = require('stream');
const { SRT } = require('../build/Release/node_srt.node');
const debug = require('debug')('srt-write-stream');
class SRTWriteStream extends Writable {
constructor(address, port, opts) {
super();
this.srt = new 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);
}
}
close() {
this.srt.close(this.socket);
this.fd = null;
}
_write(chunk, encoding, callback) {
debug(`Writing chunk ${chunk.length}`);
if (this.fd) {
this.srt.write(this.fd, chunk);
callback();
} else {
callback(new Error("Socket was closed"));
}
}
_destroy(err, callback) {
this.close();
}
}
module.exports = {
SRTWriteStream
};

60
src/tools.js Normal file
View file

@ -0,0 +1,60 @@
/**
*
* @param {Buffer | Uint8Array} srcData
* @param {number} chunkMaxSize
* @param {number} byteLength
* @param {number} initialOffset
* @returns {Array<Uint8Array>}
*/
function sliceBufferToChunks(srcData, chunkMaxSize,
byteLength = srcData.byteLength, initialOffset = 0) {
const chunks = [];
let relativeOffset = 0;
for (let offset = initialOffset; relativeOffset < byteLength; offset += chunkMaxSize) {
relativeOffset = offset - initialOffset;
const size = Math.min(chunkMaxSize, byteLength - relativeOffset);
const chunkBuf
= Uint8Array.prototype
.slice.call(srcData, offset, offset + size);
chunks.push(chunkBuf);
}
return chunks;
}
/**
*
* @param {Array<Uint8Array>} chunks Input chunks
* @returns {number}
*/
function getChunksTotalByteLength(chunks) {
return chunks.reduce((sumBytes, chunk) => (sumBytes + chunk.byteLength), 0)
}
/**
*
* @param {Array<Uint8Array>} chunks Input chunks
* @param {Buffer} targetBuffer Optional, must have sufficient size
* @returns {Buffer} Passed buffer or newly allocated
*/
function copyChunksIntoBuffer(chunks, targetBuffer = null) {
if (!targetBuffer) {
const totalSize = getChunksTotalByteLength(chunks);
targetBuffer = Buffer.alloc(totalSize);
}
let offset = 0;
for (let i = 0; i < chunks.length; i++) {
if (offset >= targetBuffer.length) {
throw new Error('Target buffer to merge chunks in is too small');
}
Buffer.from(chunks[i]).copy(targetBuffer, offset)
offset += chunks[i].byteLength;
}
return targetBuffer;
}
module.exports = {
getChunksTotalByteLength,
copyChunksIntoBuffer,
sliceBufferToChunks
}

View file

@ -1,107 +1,115 @@
declare module "srt" {
type AsyncSRTCallback<T> = (result: T) => void;
import { SRTLoggingLevel, SRTResult, SRTSockOpt, SRTSockStatus } from "../src/srt-api-enums";
class AsyncSRT {
import {SRTReadReturn, SRTFileDescriptor, SRTEpollEvent, SRTSockOptValue} from "./srt-api"
static TimeoutMs: number;
export type AsyncSRTCallback<T> = (result: T) => void;
/**
*
* @param sender
* @returns SRTSOCKET identifier (integer value)
*/
createSocket(sender: boolean, callback?: AsyncSRTCallback<number>): Promise<number>
export class AsyncSRT {
/**
*
* @param socket
* @param address
* @param port
*/
bind(socket: number, address: string, port: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
static TimeoutMs: number;
/**
*
* @param socket
* @param backlog
*/
listen(socket: number, backlog: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param sender
* @returns SRTSOCKET identifier (integer value)
*/
createSocket(sender: boolean, callback?: AsyncSRTCallback<number>): Promise<number>
/**
*
* @param socket
* @param host
* @param port
*/
connect(socket: number, host: string, port: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
* @param address
* @param port
*/
bind(socket: number, address: string, port: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
* @returns File descriptor of incoming connection pipe
*/
accept(socket: number, callback?: AsyncSRTCallback<SRTFileDescriptor>): Promise<SRTFileDescriptor>
/**
*
* @param socket
* @param backlog
*/
listen(socket: number, backlog: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
*/
close(socket: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
* @param host
* @param port
*/
connect(socket: number, host: string, port: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
* @param chunkSize
*/
read(socket: number, chunkSize: number, callback?: AsyncSRTCallback<SRTReadReturn>): Promise<SRTReadReturn>
/**
*
* @param socket
* @returns File descriptor of incoming connection pipe
*/
accept(socket: number, callback?: AsyncSRTCallback<SRTFileDescriptor>): Promise<SRTFileDescriptor>
/**
*
* @param socket
* @param chunk
*/
write(socket: number, chunk: Buffer, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
*/
close(socket: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
* @param option
* @param value
*/
setSockOpt(socket: number, option: SRTSockOpt, value: SRTSockOptValue, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
* @param chunkSize
*/
read(socket: number, chunkSize: number, callback?: AsyncSRTCallback<SRTReadReturn>): Promise<SRTReadReturn>
/**
*
* @param socket
* @param option
*/
getSockOpt(socket: number, option: SRTSockOpt, callback?: AsyncSRTCallback<SRTSockOptValue>): Promise<SRTSockOptValue>
/**
*
* @param socket
* @param chunk
*/
write(socket: number, chunk: Buffer, callback?: AsyncSRTCallback<SRTResult>): Promise<number | SRTResult.SRT_ERROR>
/**
*
* @param socket
*/
getSockState(socket: number, callback?: AsyncSRTCallback<SRTSockStatus>): Promise<SRTSockStatus>
/**
*
* @param socket
* @param option
* @param value
*/
setSockOpt(socket: number, option: SRTSockOpt, value: SRTSockOptValue, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
* @returns epid
*/
epollCreate(callback?: AsyncSRTCallback<number>): Promise<number>
/**
*
* @param socket
* @param option
*/
getSockOpt(socket: number, option: SRTSockOpt, callback?: AsyncSRTCallback<SRTSockOptValue>): Promise<SRTSockOptValue>
/**
*
* @param epid
* @param socket
* @param events
*/
epollAddUsock(epid: number, socket: number, events: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param socket
*/
getSockState(socket: number, callback?: AsyncSRTCallback<SRTSockStatus>): Promise<SRTSockStatus>
/**
*
* @param epid
* @param msTimeOut
*/
epollUWait(epid: number, msTimeOut: number, callback?: AsyncSRTCallback<SRTEpollEvent[]>): Promise<SRTEpollEvent[]>
}
/**
* @returns epid
*/
epollCreate(callback?: AsyncSRTCallback<number>): Promise<number>
/**
*
* @param epid
* @param socket
* @param events
*/
epollAddUsock(epid: number, socket: number, events: number, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
/**
*
* @param epid
* @param msTimeOut
*/
epollUWait(epid: number, msTimeOut: number, callback?: AsyncSRTCallback<SRTEpollEvent[]>): Promise<SRTEpollEvent[]>
/**
*
* @param logLevel
*/
setLogLevel(logLevel: SRTLoggingLevel, callback?: AsyncSRTCallback<SRTResult>): Promise<SRTResult>
}

314
types/srt-api.d.ts vendored
View file

@ -1,191 +1,127 @@
declare module "srt" {
class SRT {
/**
*
* @param sender
* @returns SRTSOCKET identifier (integer value)
*/
createSocket(sender?: boolean): number
import { SRTLoggingLevel, SRTResult, SRTSockOpt, SRTSockStatus } from "../src/srt-api-enums";
/**
*
* @param socket
* @param address
* @param port
*/
bind(socket: number, address: string, port: number): SRTResult
/**
*
* @param socket
* @param backlog
*/
listen(socket: number, backlog: number): SRTResult
/**
*
* @param socket
* @param host
* @param port
*/
connect(socket: number, host: string, port: number): SRTResult
/**
*
* @param socket
* @returns File descriptor of incoming connection pipe
*/
accept(socket: number): SRTFileDescriptor
/**
*
* @param socket
*/
close(socket: number): SRTResult
/**
*
* @param socket
* @param chunkSize
*/
read(socket: number, chunkSize: number): SRTReadReturn
/**
*
* @param socket
* @param chunk
*/
write(socket: number, chunk: Buffer): SRTResult
/**
*
* @param socket
* @param option
* @param value
*/
setSockOpt(socket: number, option: SRTSockOpt, value: SRTSockOptValue): SRTResult
/**
*
* @param socket
* @param option
*/
getSockOpt(socket: number, option: SRTSockOpt): SRTSockOptValue
/**
*
* @param socket
*/
getSockState(socket: number): SRTSockStatus
/**
* @returns epid
*/
epollCreate(): number
/**
*
* @param epid
* @param socket
* @param events
*/
epollAddUsock(epid: number, socket: number, events: number): SRTResult
/**
*
* @param epid
* @param msTimeOut
*/
epollUWait(epid: number, msTimeOut: number): SRTEpollEvent[]
}
interface SRTEpollEvent {
socket: SRTFileDescriptor
events: number
}
type SRTReadReturn = Buffer | SRTResult.SRT_ERROR | null
type SRTFileDescriptor = number;
type SRTSockOptValue = boolean | number | string
enum SRTResult {
SRT_ERROR = -1,
SRT_OK = 0
}
enum SRTSockOpt {
SRTO_MSS = 0, // the Maximum Transfer Unit
SRTO_SNDSYN = 1, // if sending is blocking
SRTO_RCVSYN = 2, // if receiving is blocking
SRTO_ISN = 3, // Initial Sequence Number (valid only after srt_connect or srt_accept-ed sockets)
SRTO_FC = 4, // Flight flag size (window size)
SRTO_SNDBUF = 5, // maximum buffer in sending queue
SRTO_RCVBUF = 6, // UDT receiving buffer size
SRTO_LINGER = 7, // waiting for unsent data when closing
SRTO_UDP_SNDBUF = 8, // UDP sending buffer size
SRTO_UDP_RCVBUF = 9, // UDP receiving buffer size
// XXX Free space for 2 options
// after deprecated ones are removed
SRTO_RENDEZVOUS = 12, // rendezvous connection mode
SRTO_SNDTIMEO = 13, // send() timeout
SRTO_RCVTIMEO = 14, // recv() timeout
SRTO_REUSEADDR = 15, // reuse an existing port or create a new one
SRTO_MAXBW = 16, // maximum bandwidth (bytes per second) that the connection can use
SRTO_STATE = 17, // current socket state, see UDTSTATUS, read only
SRTO_EVENT = 18, // current available events associated with the socket
SRTO_SNDDATA = 19, // size of data in the sending buffer
SRTO_RCVDATA = 20, // size of data available for recv
SRTO_SENDER = 21, // Sender mode (independent of conn mode), for encryption, tsbpd handshake.
SRTO_TSBPDMODE = 22, // Enable/Disable TsbPd. Enable -> Tx set origin timestamp, Rx deliver packet at origin time + delay
SRTO_LATENCY = 23, // NOT RECOMMENDED. SET: to both SRTO_RCVLATENCY and SRTO_PEERLATENCY. GET: same as SRTO_RCVLATENCY.
SRTO_TSBPDDELAY = 23, // DEPRECATED. ALIAS: SRTO_LATENCY
SRTO_INPUTBW = 24, // Estimated input stream rate.
SRTO_OHEADBW, // MaxBW ceiling based on % over input stream rate. Applies when UDT_MAXBW=0 (auto).
SRTO_PASSPHRASE = 26, // Crypto PBKDF2 Passphrase size[0,10..64] 0:disable crypto
SRTO_PBKEYLEN, // Crypto key len in bytes {16,24,32} Default: 16 (128-bit)
SRTO_KMSTATE, // Key Material exchange status (UDT_SRTKmState)
SRTO_IPTTL = 29, // IP Time To Live (passthru for system sockopt IPPROTO_IP/IP_TTL)
SRTO_IPTOS, // IP Type of Service (passthru for system sockopt IPPROTO_IP/IP_TOS)
SRTO_TLPKTDROP = 31, // Enable receiver pkt drop
SRTO_SNDDROPDELAY = 32, // Extra delay towards latency for sender TLPKTDROP decision (-1 to off)
SRTO_NAKREPORT = 33, // Enable receiver to send periodic NAK reports
SRTO_VERSION = 34, // Local SRT Version
SRTO_PEERVERSION, // Peer SRT Version (from SRT Handshake)
SRTO_CONNTIMEO = 36, // Connect timeout in msec. Ccaller default: 3000, rendezvous (x 10)
// deprecated: SRTO_TWOWAYDATA, SRTO_SNDPBKEYLEN, SRTO_RCVPBKEYLEN (@c below)
_DEPRECATED_SRTO_SNDPBKEYLEN = 38, // (needed to use inside the code without generating -Wswitch)
//
SRTO_SNDKMSTATE = 40, // (GET) the current state of the encryption at the peer side
SRTO_RCVKMSTATE, // (GET) the current state of the encryption at the agent side
SRTO_LOSSMAXTTL, // Maximum possible packet reorder tolerance (number of packets to receive after loss to send lossreport)
SRTO_RCVLATENCY, // TsbPd receiver delay (mSec) to absorb burst of missed packet retransmission
SRTO_PEERLATENCY, // Minimum value of the TsbPd receiver delay (mSec) for the opposite side (peer)
SRTO_MINVERSION, // Minimum SRT version needed for the peer (peers with less version will get connection reject)
SRTO_STREAMID, // A string set to a socket and passed to the listener's accepted socket
SRTO_CONGESTION, // Congestion controller type selection
SRTO_MESSAGEAPI, // In File mode, use message API (portions of data with boundaries)
SRTO_PAYLOADSIZE, // Maximum payload size sent in one UDP packet (0 if unlimited)
SRTO_TRANSTYPE = 50, // Transmission type (set of options required for given transmission type)
SRTO_KMREFRESHRATE, // After sending how many packets the encryption key should be flipped to the new key
SRTO_KMPREANNOUNCE, // How many packets before key flip the new key is annnounced and after key flip the old one decommissioned
SRTO_ENFORCEDENCRYPTION, // Connection to be rejected or quickly broken when one side encryption set or bad password
SRTO_IPV6ONLY, // IPV6_V6ONLY mode
SRTO_PEERIDLETIMEO, // Peer-idle timeout (max time of silence heard from peer) in [ms]
// (some space left)
SRTO_PACKETFILTER = 60 // Add and configure a packet filter
}
enum SRTSockStatus {
SRTS_INIT = 1,
SRTS_OPENED,
SRTS_LISTENING,
SRTS_CONNECTING,
SRTS_CONNECTED,
SRTS_BROKEN,
SRTS_CLOSING,
SRTS_CLOSED,
SRTS_NONEXIST
}
export interface SRTEpollEvent {
socket: SRTFileDescriptor
events: number
}
export type SRTReadReturn = Uint8Array | null | SRTResult.SRT_ERROR;
export type SRTFileDescriptor = number;
export type SRTSockOptValue = boolean | number | string
export class SRT {
static OK: SRTResult.SRT_OK;
static ERROR: SRTResult.SRT_ERROR;
static INVALID_SOCK: SRTResult.SRT_ERROR;
// TODO: add SOCKET_OPTIONS, SOCKET_STATUS enums
// and EPOLL_OPTS
/**
*
* @param sender
* @returns SRTSOCKET identifier (integer value)
*/
createSocket(sender?: boolean): number
/**
*
* @param socket
* @param address
* @param port
*/
bind(socket: number, address: string, port: number): SRTResult
/**
*
* @param socket
* @param backlog
*/
listen(socket: number, backlog: number): SRTResult
/**
*
* @param socket
* @param host
* @param port
*/
connect(socket: number, host: string, port: number): SRTResult
/**
*
* @param socket
* @returns File descriptor of incoming connection pipe
*/
accept(socket: number): SRTFileDescriptor
/**
*
* @param socket
*/
close(socket: number): SRTResult
/**
*
* @param socket
* @param chunkSize
*/
read(socket: number, chunkSize: number): SRTReadReturn
/**
*
* @param socket
* @param chunk
*/
write(socket: number, chunk: Buffer): SRTResult
/**
*
* @param socket
* @param option
* @param value
*/
setSockOpt(socket: number, option: SRTSockOpt, value: SRTSockOptValue): SRTResult
/**
*
* @param socket
* @param option
*/
getSockOpt(socket: number, option: SRTSockOpt): SRTSockOptValue
/**
*
* @param socket
*/
getSockState(socket: number): SRTSockStatus
/**
* @returns epid
*/
epollCreate(): number
/**
*
* @param epid
* @param socket
* @param events
*/
epollAddUsock(epid: number, socket: number, events: number): SRTResult
/**
*
* @param epid
* @param msTimeOut
*/
epollUWait(epid: number, msTimeOut: number): SRTEpollEvent[]
/**
*
* @param logLevel Or 0 - 7 integer (not all values present in enum)
*/
setLogLevel(logLevel: SRTLoggingLevel): SRTResult;
}

68
types/srt-server.d.ts vendored
View file

@ -1,22 +1,58 @@
/// <reference types="node" />
declare module "srt" {
import {EventEmitter} from 'events';
import { SRTResult, SRTSockOpt } from '../src/srt-api-enums';
import { SRTSockOptValue } from './srt-api';
import { AsyncSRT } from './srt-api-async';
import {EventEmitter} from 'events';
export class AsyncReaderWriter {
constructor(asyncSrt: AsyncSRT, socketFd: number);
interface SRTServerBindOpts {
/**
* default: "0.0.0.0"
*/
address?: string
port: number
}
type SRTServerEvent = "listening"; /*| "foobar" */
class SRTServer extends EventEmitter /*<SRTServerEvent>*/ {
listen(opts: SRTServerBindOpts): void
}
writeChunks(buffer: Uint8Array | Buffer, mtuSize: number,
writesPerTick: number): Promise<void>;
readChunks(minBytesRead: number,
readBufSize: number,
onRead: (buf: Uint8Array) => void,
onError: (readResult: (SRTResult.SRT_ERROR | null)) => void): Promise<Uint8Array[]>;
}
export class SRTConnection extends EventEmitter {
readonly fd: number;
readonly gotFirstData: boolean;
read(): Promise<Uint8Array | SRTResult.SRT_ERROR | null>;
write(chunk: Buffer | Uint8Array): Promise<SRTResult>;
close(): Promise<SRTResult | null>;
isClosed(): boolean;
onData(): void;
getReaderWriter(): AsyncReaderWriter;
}
export class SRTServer extends EventEmitter /*<SRTServerEvent>*/ {
static create(port: number, address?: string,
epollPeriodMs?: number): Promise<SRTServer>;
port: number;
address: string;
epollPeriodMs: number;
socket: number;
epid: number;
constructor(port: number, address?: string, epollPeriodMs?: number);
create(): Promise<SRTServer>;
open(): Promise<SRTServer>;
dispose(): Promise<SRTResult>;
setSocketFlags(opts: SRTSockOpt[], values: SRTSockOptValue[]): Promise<SRTResult[]>;
getConnectionByHandle(fd: number);
getAllConnections(): SRTConnection[];
}

97
types/srt-stream.d.ts vendored
View file

@ -1,53 +1,52 @@
/// <reference types="node" />
import { Writable, Readable } from "stream";
import { SRT, SRTFileDescriptor } from "./srt-api";
declare module "srt" {
interface SRTConnectionState {
readonly srt: SRT;
readonly socket: number;
readonly address: string;
readonly port: number;
}
interface SRTCallerState extends SRTConnectionState {
readonly fd: SRTFileDescriptor | null;
connect(callback: (state: SRTCallerState) => void);
close();
}
interface SRTListenerState extends SRTConnectionState {
listen(callback: (state: SRTListenerState) => void);
}
class SRTReadStream extends Readable implements SRTCallerState, SRTListenerState {
readonly srt: SRT;
readonly fd: SRTFileDescriptor | null;
readonly socket: number;
readonly address: string;
readonly port: number;
readonly readTimer: number | null;
constructor(address: string, port: number, opts?: unknown);
connect(callback: (state: SRTCallerState) => void);
close();
listen(callback: (state: SRTListenerState) => void);
}
class SRTWriteStream extends Writable implements SRTCallerState {
readonly srt: SRT;
readonly socket: number;
readonly address: string;
readonly port: number;
readonly fd: SRTFileDescriptor | null;
constructor(address: string, port: number, opts?: unknown);
connect(callback: (state: SRTCallerState) => void);
close();
}
interface SRTConnectionState {
readonly srt: SRT;
readonly socket: number;
readonly address: string;
readonly port: number;
}
interface SRTCallerState extends SRTConnectionState {
readonly fd: SRTFileDescriptor | null;
connect(callback: (state: SRTCallerState) => void);
close();
}
export interface SRTListenerState extends SRTConnectionState {
listen(callback: (state: SRTListenerState) => void);
}
export class SRTReadStream extends Readable implements SRTCallerState, SRTListenerState {
readonly srt: SRT;
readonly fd: SRTFileDescriptor | null;
readonly socket: number;
readonly address: string;
readonly port: number;
readonly readTimer: number | null;
constructor(address: string, port: number, opts?: unknown);
connect(callback: (state: SRTCallerState) => void);
close();
listen(callback: (state: SRTListenerState) => void);
}
export class SRTWriteStream extends Writable implements SRTCallerState {
readonly srt: SRT;
readonly socket: number;
readonly address: string;
readonly port: number;
readonly fd: SRTFileDescriptor | null;
constructor(address: string, port: number, opts?: unknown);
connect(callback: (state: SRTCallerState) => void);
close();
}