-
Notifications
You must be signed in to change notification settings - Fork 1
/
channel.js
82 lines (59 loc) · 1.95 KB
/
channel.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
'use strict';
module.exports = Channel;
function Channel() {
return {
buffer: [],
consumers: [],
closed: false
};
}
const END = Channel.END = Symbol('END');
const KEEP_OPEN = Channel.KEEP_OPEN = Symbol('KEEP_OPEN');
const CLOSE_BOTH = Channel.CLOSE_BOTH = Symbol('CLOSE_BOTH');
Channel.put = (channel, value) => {
if (channel.closed) return channel;
if (value === END) channel.closed = true;
channel.buffer.push(value);
return runTick(channel);
};
Channel.take = (channel, callback) => {
if (channel.closed && channel.buffer.length === 0) return channel;
if (isFunction(callback)) channel.consumers.push(callback);
return runTick(channel);
};
Channel.close = (channel) => {
return Channel.put(channel, END);
};
Channel.pipe = (input, output, keepOpen = KEEP_OPEN, transform = identity) => {
const consume = (value) => {
if (!(keepOpen === KEEP_OPEN && value === END))
Channel.put(output, transform(value));
if (value !== END) Channel.take(input, consume);
};
Channel.take(input, consume);
return output;
};
Channel.demux = (channels, output, keepOpen, transform) => {
channels.forEach((channel) => Channel.pipe(channel, output, keepOpen, transform));
return output;
};
Channel.mux = (input, channels, keepOpen, transform = identity) => {
const consume = (value) => {
if (!(keepOpen === KEEP_OPEN && value === END))
channels.forEach((channel) => Channel.put(channel, transform(value)));
if (value !== END) Channel.take(input, consume);
};
Channel.take(input, consume);
return channels;
};
const runTick = (channel) => {
if (!(channel.buffer.length !== 0 && channel.consumers.length !== 0)) {
return channel;
}
const message = channel.buffer.shift();
const consumer = channel.consumers.shift();
setImmediate(consumer.bind(null, message));
return runTick(channel);
};
const identity = (value) => value;
const isFunction = (fn) => ({}).toString.call(fn) === '[object Function]';