The Reactive Extensions for JavaScript provides integration points to the core Node.js libraries.
Deprecated in favor of Rx.Observable.fromCallback in rx.async.js.
Converts a callback function to an observable sequence.
func(Function): Callback function[scheduler = Rx.Scheduler.timeout](Scheduler): Scheduler used to execute the callback.[context](Any): The context to execute the callback.
(Function): Function, when called with arguments, creates an Observable sequence from the callback.
var fs = require('fs');
var Rx = require('Rx');
// Wrap exists
var exists = Rx.Node.fromCallback(fs.exists);
// Call exists
var source = exists('/etc/passwd');
var observer = Rx.Observer.create(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}
);
var subscription = source.subscribe(observer);
// => Next: true
// => Completed- rx.node.js
Deprecated in favor of Rx.Observable.fromNodeCallback in rx.async.js.
Converts a Node.js callback style function to an observable sequence. This must be in function (err, ...) format.
func(Function): Callback function which must be in function (err, ...) format.[scheduler = Rx.Scheduler.timeout](Scheduler): Scheduler used to execute the callback.[context](Any): The context to execute the callback.
(Function): An function which when applied, returns an observable sequence with the callback arguments as an array.
var fs = require('fs');
var Rx = require('Rx');
var source = Rx.Node.fromNodeCallback(fs.stat)('file.txt');
var observer = Rx.Observer.create(
function (x) {
var stat = x[0];
console.log('Next: ' + stat.isFile());
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}
);
var subscription = source.subscribe(observer);
// => Next: true
// => Completed- rx.node.js
Handles an event from the given EventEmitter as an observable sequence.
eventEmitter(EventEmitter): The EventEmitter to subscribe to the given event.eventName(String): The event name to subscribe.
(Observable): An observable sequence generated from the named event from the given EventEmitter.
var EventEmitter = require('events').EventEmitter;
var Rx = require('Rx');
var emitter = new EventEmitter();
var source = Rx.Node.fromEvent(emitter, 'data');
var observer = Rx.Observer.create(
function (x) {
console.log('Next: ' + x[0]);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}
);
var subscription = source.subscribe(observer);
emitter.emit('data', 'foo');
// => Next: foo- rx.node.js
Converts the given observable sequence to an event emitter with the given event name. The errors are handled on the 'error' event and completion on the 'end' event.
observable(Obsesrvable): The observable sequence to convert to an EventEmitter.eventName(String): The event name to subscribe.
(EventEmitter): An EventEmitter which emits the given eventName for each onNext call in addition to 'error' and 'end' events.
var Rx = require('Rx');
var source = Rx.Observable.return(42);
var emitter = Rx.Node.toEventEmitter(source, 'data');
emitter.on('data', function (data) {
console.log('Data: ' + data);
});
emitter.on('end', function () {
console.log('End');
});
// Ensure to call publish to fire events from the observable
emitter.publish();
// => Data: 42
// => End- rx.node.js
Converts a flowing stream to an Observable sequence.
stream(Stream): A stream to convert to a observable sequence.
(Observable): An observable sequence which fires on each 'data' event as well as handling 'error' and 'end' events.
var Rx = require('rx');
var subscription = Rx.Node.fromStream(process.stdin)
.subscribe(function (x) { console.log(x); });
// => r<Buffer 72>
// => x<Buffer 78>- rx.node.js
Writes an observable sequence to a stream.
observable(Obsesrvable): Observable sequence to write to a stream.stream(Stream): The stream to write to.[encoding](String): The encoding of the item to write.
(Disposable): The subscription handle.
var Rx = require('Rx');
var source = Rx.Observable.range(0, 5);
var subscription = Rx.Node.writeToStream(source, process.stdout, 'utf8');
// => 012- rx.node.js