-
Notifications
You must be signed in to change notification settings - Fork 0
/
app.js
31 lines (25 loc) · 959 Bytes
/
app.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
const Observable = require("rxjs/Rx").Observable;
const mqtt = require("mqtt").connect("mqtt://localhost");
function isPrime(value) {
for (let i = 2; i < value; i++) if (value % i === 0) return false;
return value > 1;
}
// once connected
mqtt.on("connect", () => {
console.info("connected");
mqtt.subscribe("#");
mqtt.publish("node/connection", "je suis connecté");
// publisher
Observable.from(["a", "b", "cd", "efg", "h", "ijk", "lm"])
.filter(ev => ev.length < 3)
.delay(1000)
.subscribe(ev => mqtt.publish("node/fromArray", ev));
// publisher
Observable.interval(100)
.filter(isPrime)
.bufferTime(10000)
.subscribe(ev => mqtt.publish("node/prime", ev.toString()));
});
// consumer
Observable.fromEvent(mqtt, "message", (topic, message) => ({ topic, message })) // we're converting the useful arguments of the mqtt event into a single object
.subscribe(ev => console.log(`${ev.topic} > ${ev.message}`));