This repository has been archived by the owner on Jun 27, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 16
/
index.js
102 lines (90 loc) · 2.5 KB
/
index.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
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
'use strict'
const debugName = 'libp2p:floodsub'
// @ts-ignore time-cache does not export types
const TimeCache = require('time-cache')
const toString = require('uint8arrays/to-string')
const BaseProtocol = require('libp2p-interfaces/src/pubsub')
const { utils } = require('libp2p-interfaces/src/pubsub')
const { multicodec } = require('./config')
/**
* @typedef {import('libp2p-interfaces/src/pubsub').InMessage} InMessage
*/
/**
* FloodSub (aka dumbsub is an implementation of pubsub focused on
* delivering an API for Publish/Subscribe, but with no CastTree Forming
* (it just floods the network).
*/
class FloodSub extends BaseProtocol {
/**
* @param {import('libp2p')} libp2p - instance of libp2p
* @param {Object} [options]
* @param {boolean} [options.emitSelf] - if publish should emit to self, if subscribed, defaults to false
* @class
*/
constructor (libp2p, options = {}) {
super({
debugName: debugName,
multicodecs: multicodec,
libp2p,
canRelayMessage: true,
...options
})
/**
* Cache of seen messages
*
* @type {TimeCache}
*/
this.seenCache = new TimeCache()
}
/**
* Process incoming message
* Extends base implementation to check router cache.
*
* @override
* @param {InMessage} message - The message to process
* @returns {Promise<void>}
*/
async _processRpcMessage (message) {
// Check if I've seen the message, if yes, ignore
const seqno = await this.getMsgId(message)
const msgIdStr = toString(seqno, 'base64')
if (this.seenCache.has(msgIdStr)) {
return
}
this.seenCache.put(msgIdStr)
await super._processRpcMessage(message)
}
/**
* Publish message created. Forward it to the peers.
*
* @override
* @param {InMessage} message
* @returns {Promise<void>}
*/
_publish (message) {
this._forwardMessage(message)
return Promise.resolve()
}
/**
* Forward message to peers.
*
* @param {InMessage} message
* @returns {void}
*/
_forwardMessage (message) {
message.topicIDs.forEach((topic) => {
const peers = this.topics.get(topic)
if (!peers) {
return
}
peers.forEach((id) => {
this.log('publish msgs on topics', message.topicIDs, id)
if (id !== this.peerId.toB58String() && id !== message.receivedFrom) {
this._sendRpc(id, { msgs: [utils.normalizeOutRpcMessage(message)] })
}
})
})
}
}
module.exports = FloodSub
module.exports.multicodec = multicodec