@effection/channel
Advanced tools
Comparing version 2.0.0-preview.4 to 2.0.0-preview.5
# Changelog | ||
## 2.0.0-preview.5 | ||
### Minor Changes | ||
- 9cf6053: Increase default value of max subscribers on channel | ||
- 0b24415: Add WritableStream interface and implement it for channels | ||
### Patch Changes | ||
- Updated dependencies [0b24415] | ||
- Updated dependencies [22e5230] | ||
- Updated dependencies [70c358f] | ||
- Updated dependencies [3983202] | ||
- Updated dependencies [2c2749d] | ||
- @effection/subscription@2.0.0-preview.5 | ||
- @effection/core@2.0.0-preview.5 | ||
## 2.0.0-preview.4 | ||
@@ -4,0 +21,0 @@ |
@@ -16,5 +16,7 @@ 'use strict'; | ||
bus.setMaxListeners(options.maxSubscribers); | ||
} else { | ||
bus.setMaxListeners(100000); | ||
} | ||
var subscribable = subscription.createStream(function (publish) { | ||
var stream = subscription.createStream(function (publish) { | ||
return function* (task) { | ||
@@ -35,17 +37,22 @@ var subscription = events.on(bus, 'event').subscribe(task); | ||
}); | ||
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({ | ||
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); | ||
send: send, | ||
close: close, | ||
stream: stream | ||
}, stream); | ||
} | ||
@@ -52,0 +59,0 @@ |
@@ -1,2 +0,2 @@ | ||
"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)}; | ||
"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;i.setMaxListeners(r.maxSubscribers?r.maxSubscribers:1e5);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({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]})},stream:s},s)}; | ||
//# sourceMappingURL=channel.cjs.production.min.js.map |
@@ -1,10 +0,11 @@ | ||
import { Stream } from '@effection/subscription'; | ||
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']; | ||
export declare type ChannelOptions = { | ||
maxSubscribers?: number; | ||
}; | ||
export interface Channel<T, TClose = undefined> extends Stream<T, TClose> { | ||
send(message: T): void; | ||
close(...args: TClose extends undefined ? [] : [TClose]): void; | ||
export interface Channel<T, TClose = undefined> extends WritableStream<T, T, TClose> { | ||
close: Close<TClose>; | ||
stream: Stream<T, TClose>; | ||
} | ||
export declare function createChannel<T, TClose = undefined>(options?: ChannelOptions): Channel<T, TClose>; |
@@ -14,5 +14,7 @@ import { createStream } from '@effection/subscription'; | ||
bus.setMaxListeners(options.maxSubscribers); | ||
} else { | ||
bus.setMaxListeners(100000); | ||
} | ||
var subscribable = createStream(function (publish) { | ||
var stream = createStream(function (publish) { | ||
return function* (task) { | ||
@@ -33,17 +35,22 @@ var subscription = on(bus, 'event').subscribe(task); | ||
}); | ||
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({ | ||
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); | ||
send: send, | ||
close: close, | ||
stream: stream | ||
}, stream); | ||
} | ||
@@ -50,0 +57,0 @@ |
{ | ||
"name": "@effection/channel", | ||
"version": "2.0.0-preview.4", | ||
"version": "2.0.0-preview.5", | ||
"description": "MPMC Channel implementation for effection", | ||
@@ -23,3 +23,3 @@ "main": "dist/index.js", | ||
"devDependencies": { | ||
"@effection/mocha": "2.0.0-preview.3", | ||
"@effection/mocha": "2.0.0-preview.4", | ||
"@frontside/tsconfig": "^0.0.1", | ||
@@ -38,6 +38,6 @@ "@types/node": "^12.7.11", | ||
"dependencies": { | ||
"@effection/core": "^2.0.0-preview.4", | ||
"@effection/core": "^2.0.0-preview.5", | ||
"@effection/events": "^2.0.0-preview.4", | ||
"@effection/subscription": "^2.0.0-preview.4" | ||
"@effection/subscription": "^2.0.0-preview.5" | ||
} | ||
} |
@@ -1,5 +0,9 @@ | ||
import { createStream, Stream } from '@effection/subscription'; | ||
import { createStream, WritableStream, Writable, 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 = { | ||
@@ -9,5 +13,4 @@ maxSubscribers?: number; | ||
export interface Channel<T, TClose = undefined> extends Stream<T, TClose> { | ||
send(message: T): void; | ||
close(...args: TClose extends undefined ? [] : [TClose]): void; | ||
export interface Channel<T, TClose = undefined> extends WritableStream<T, T, TClose> { | ||
close: Close<TClose>; | ||
stream: Stream<T, TClose>; | ||
@@ -21,5 +24,7 @@ } | ||
bus.setMaxListeners(options.maxSubscribers); | ||
} else { | ||
bus.setMaxListeners(100000); | ||
} | ||
let subscribable = createStream<T, TClose>((publish) => function*(task) { | ||
let stream = createStream<T, TClose>((publish) => function*(task) { | ||
let subscription = on(bus, 'event').subscribe(task); | ||
@@ -36,13 +41,11 @@ while(true) { | ||
return Object.assign({ | ||
stream: subscribable, | ||
let send: Send<T> = (message: T) => { | ||
bus.emit('event', { done: false, value: message }); | ||
}; | ||
send(message: T) { | ||
bus.emit('event', { done: false, value: message }); | ||
}, | ||
let close: Close<TClose> = (...args) => { | ||
bus.emit('event', { done: true, value: args[0] }); | ||
}; | ||
close(...args: TClose extends undefined ? [] : [TClose]) { | ||
bus.emit('event', { done: true, value: args[0] }); | ||
} | ||
}, subscribable); | ||
return Object.assign({ send, close, stream }, stream); | ||
} |
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
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
18859
156