@effection/channel
Advanced tools
Comparing version 2.0.0-preview.3-2fc4f8e to 2.0.0-preview.3-2febf01
@@ -16,7 +16,5 @@ 'use strict'; | ||
bus.setMaxListeners(options.maxSubscribers); | ||
} else { | ||
bus.setMaxListeners(100000); | ||
} | ||
var stream = subscription.createStream(function (publish) { | ||
var subscribable = subscription.createStream(function (publish) { | ||
return function* (task) { | ||
@@ -27,4 +25,3 @@ var subscription = events.on(bus, 'event').subscribe(task); | ||
var _yield$subscription$n = yield subscription.next(), | ||
_yield$subscription$n2 = _yield$subscription$n.value, | ||
next = _yield$subscription$n2[0]; | ||
next = _yield$subscription$n.value; | ||
@@ -39,50 +36,20 @@ if (next.done) { | ||
}); | ||
var send = function send(message) { | ||
bus.emit('event', { | ||
done: false, | ||
value: message | ||
}); | ||
}; | ||
var close = function close() { | ||
bus.emit('event', { | ||
done: true, | ||
value: arguments.length <= 0 ? undefined : arguments[0] | ||
}); | ||
}; | ||
return Object.assign({ | ||
send: send, | ||
close: close, | ||
stream: stream | ||
}, stream); | ||
stream: subscribable, | ||
send: function send(message) { | ||
bus.emit('event', { | ||
done: false, | ||
value: message | ||
}); | ||
}, | ||
close: function close() { | ||
bus.emit('event', { | ||
done: true, | ||
value: arguments.length <= 0 ? undefined : arguments[0] | ||
}); | ||
} | ||
}, subscribable); | ||
} | ||
function createDuplexChannel(options) { | ||
if (options === void 0) { | ||
options = {}; | ||
} | ||
var leftChannel = createChannel(options); | ||
var rightChannel = createChannel(options); | ||
var close = function close() { | ||
leftChannel.close.apply(leftChannel, arguments); | ||
rightChannel.close.apply(rightChannel, arguments); | ||
}; | ||
return [Object.assign({ | ||
send: rightChannel.send, | ||
stream: leftChannel.stream, | ||
close: close | ||
}, leftChannel.stream), Object.assign({ | ||
send: leftChannel.send, | ||
stream: rightChannel.stream, | ||
close: close | ||
}, rightChannel.stream)]; | ||
} | ||
exports.createChannel = createChannel; | ||
exports.createDuplexChannel = createDuplexChannel; | ||
//# sourceMappingURL=channel.cjs.development.js.map |
@@ -1,2 +0,2 @@ | ||
"use strict";var e=require("@effection/subscription"),n=require("@effection/events"),t=require("events");function r(r){void 0===r&&(r={});var s=new t.EventEmitter;s.setMaxListeners(r.maxSubscribers?r.maxSubscribers:1e5);var a=e.createStream((function(e){return function*(t){for(var r=n.on(s,"event").subscribe(t);;){var a=(yield r.next()).value[0];if(a.done)return a.value;e(a.value)}}}));return Object.assign({send:function(e){s.emit("event",{done:!1,value:e})},close:function(){s.emit("event",{done:!0,value:arguments.length<=0?void 0:arguments[0]})},stream:a},a)}exports.createChannel=r,exports.createDuplexChannel=function(e){void 0===e&&(e={});var n=r(e),t=r(e),s=function(){n.close.apply(n,arguments),t.close.apply(t,arguments)};return[Object.assign({send:t.send,stream:n.stream,close:s},n.stream),Object.assign({send:n.send,stream:t.stream,close:s},t.stream)]}; | ||
"use strict";var e=require("@effection/subscription"),n=require("@effection/events"),t=require("events");exports.createChannel=function(r){void 0===r&&(r={});var i=new t.EventEmitter;r.maxSubscribers&&i.setMaxListeners(r.maxSubscribers);var s=e.createStream((function(e){return function*(t){for(var r=n.on(i,"event").subscribe(t);;){var s=(yield r.next()).value;if(s.done)return s.value;e(s.value)}}}));return Object.assign({stream:s,send:function(e){i.emit("event",{done:!1,value:e})},close:function(){i.emit("event",{done:!0,value:arguments.length<=0?void 0:arguments[0]})}},s)}; | ||
//# sourceMappingURL=channel.cjs.production.min.js.map |
@@ -1,11 +0,10 @@ | ||
import { WritableStream, Writable, Stream } from '@effection/subscription'; | ||
export declare type Close<T> = (...args: T extends undefined ? [] : [T]) => void; | ||
export declare type Send<T> = Writable<T>['send']; | ||
import { Stream } from '@effection/subscription'; | ||
export declare type ChannelOptions = { | ||
maxSubscribers?: number; | ||
}; | ||
export interface Channel<T, TClose = undefined> extends WritableStream<T, T, TClose> { | ||
close: Close<TClose>; | ||
export interface Channel<T, TClose = undefined> extends Stream<T, TClose> { | ||
send(message: T): void; | ||
close(...args: TClose extends undefined ? [] : [TClose]): void; | ||
stream: Stream<T, TClose>; | ||
} | ||
export declare function createChannel<T, TClose = undefined>(options?: ChannelOptions): Channel<T, TClose>; |
@@ -14,7 +14,5 @@ import { createStream } from '@effection/subscription'; | ||
bus.setMaxListeners(options.maxSubscribers); | ||
} else { | ||
bus.setMaxListeners(100000); | ||
} | ||
var stream = createStream(function (publish) { | ||
var subscribable = createStream(function (publish) { | ||
return function* (task) { | ||
@@ -25,4 +23,3 @@ var subscription = on(bus, 'event').subscribe(task); | ||
var _yield$subscription$n = yield subscription.next(), | ||
_yield$subscription$n2 = _yield$subscription$n.value, | ||
next = _yield$subscription$n2[0]; | ||
next = _yield$subscription$n.value; | ||
@@ -37,49 +34,20 @@ if (next.done) { | ||
}); | ||
var send = function send(message) { | ||
bus.emit('event', { | ||
done: false, | ||
value: message | ||
}); | ||
}; | ||
var close = function close() { | ||
bus.emit('event', { | ||
done: true, | ||
value: arguments.length <= 0 ? undefined : arguments[0] | ||
}); | ||
}; | ||
return Object.assign({ | ||
send: send, | ||
close: close, | ||
stream: stream | ||
}, stream); | ||
stream: subscribable, | ||
send: function send(message) { | ||
bus.emit('event', { | ||
done: false, | ||
value: message | ||
}); | ||
}, | ||
close: function close() { | ||
bus.emit('event', { | ||
done: true, | ||
value: arguments.length <= 0 ? undefined : arguments[0] | ||
}); | ||
} | ||
}, subscribable); | ||
} | ||
function createDuplexChannel(options) { | ||
if (options === void 0) { | ||
options = {}; | ||
} | ||
var leftChannel = createChannel(options); | ||
var rightChannel = createChannel(options); | ||
var close = function close() { | ||
leftChannel.close.apply(leftChannel, arguments); | ||
rightChannel.close.apply(rightChannel, arguments); | ||
}; | ||
return [Object.assign({ | ||
send: rightChannel.send, | ||
stream: leftChannel.stream, | ||
close: close | ||
}, leftChannel.stream), Object.assign({ | ||
send: leftChannel.send, | ||
stream: rightChannel.stream, | ||
close: close | ||
}, rightChannel.stream)]; | ||
} | ||
export { createChannel, createDuplexChannel }; | ||
export { createChannel }; | ||
//# sourceMappingURL=channel.esm.js.map |
export * from './channel'; | ||
export * from './duplex-channel'; |
{ | ||
"name": "@effection/channel", | ||
"version": "2.0.0-preview.3-2fc4f8e", | ||
"version": "2.0.0-preview.3-2febf01", | ||
"description": "MPMC Channel implementation for effection", | ||
@@ -23,7 +23,7 @@ "main": "dist/index.js", | ||
"devDependencies": { | ||
"@effection/mocha": "2.0.0-preview.2", | ||
"@frontside/tsconfig": "^0.0.1", | ||
"@types/node": "^12.7.11", | ||
"@effection/mocha": "2.0.0-preview.2", | ||
"expect": "^25.4.0", | ||
"mocha": "^7.2.0", | ||
"mocha": "^8.3.1", | ||
"ts-node": "^8.9.0", | ||
@@ -30,0 +30,0 @@ "tsdx": "0.13.2", |
@@ -1,9 +0,5 @@ | ||
import { createStream, WritableStream, Writable, Stream } from '@effection/subscription'; | ||
import { createStream, Stream } from '@effection/subscription'; | ||
import { on } from '@effection/events'; | ||
import { EventEmitter } from 'events'; | ||
export type Close<T> = (...args: T extends undefined ? [] : [T]) => void; | ||
export type Send<T> = Writable<T>['send']; | ||
export type ChannelOptions = { | ||
@@ -13,4 +9,5 @@ maxSubscribers?: number; | ||
export interface Channel<T, TClose = undefined> extends WritableStream<T, T, TClose> { | ||
close: Close<TClose>; | ||
export interface Channel<T, TClose = undefined> extends Stream<T, TClose> { | ||
send(message: T): void; | ||
close(...args: TClose extends undefined ? [] : [TClose]): void; | ||
stream: Stream<T, TClose>; | ||
@@ -24,10 +21,8 @@ } | ||
bus.setMaxListeners(options.maxSubscribers); | ||
} else { | ||
bus.setMaxListeners(100000); | ||
} | ||
let stream = createStream<T, TClose>((publish) => function*(task) { | ||
let subscribable = createStream<T, TClose>((publish) => function*(task) { | ||
let subscription = on(bus, 'event').subscribe(task); | ||
while(true) { | ||
let { value: [next] } = yield subscription.next(); | ||
let { value: next } = yield subscription.next(); | ||
if(next.done) { | ||
@@ -41,11 +36,13 @@ return next.value; | ||
let send: Send<T> = (message: T) => { | ||
bus.emit('event', { done: false, value: message }); | ||
}; | ||
return Object.assign({ | ||
stream: subscribable, | ||
let close: Close<TClose> = (...args) => { | ||
bus.emit('event', { done: true, value: args[0] }); | ||
}; | ||
send(message: T) { | ||
bus.emit('event', { done: false, value: message }); | ||
}, | ||
return Object.assign({ send, close, stream }, stream); | ||
close(...args: TClose extends undefined ? [] : [TClose]) { | ||
bus.emit('event', { done: true, value: args[0] }); | ||
} | ||
}, subscribable); | ||
} |
export * from './channel'; | ||
export * from './duplex-channel'; |
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is not supported yet
License Policy Violation
LicenseThis package is not allowed per your license policy. Review the package's license to ensure compliance.
Found 1 instance in 1 package
Minified code
QualityThis package contains minified code. This may be harmless in some cases where minified code is included in packaged libraries, however packages on npm should not minify code.
Found 1 instance in 1 package
License Policy Violation
LicenseThis package is not allowed per your license policy. Review the package's license to ensure compliance.
Found 1 instance in 1 package
Minified code
QualityThis package contains minified code. This may be harmless in some cases where minified code is included in packaged libraries, however packages on npm should not minify code.
Found 1 instance in 1 package
17215
14
146