Skip to content
This repository has been archived by the owner on Feb 12, 2024. It is now read-only.

Feat/floodsub #531

Closed
wants to merge 3 commits into from
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 64 additions & 0 deletions src/core/components/pubsub.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
'use strict'

const promisify = require('promisify-es6')
const Stream = require('stream')

const FloodSub = require('libp2p-floodsub')

const OFFLINE_ERROR = require('../utils').OFFLINE_ERROR

module.exports = function pubsub (self) {
let fsub

return {
start: promisify((libp2pNode, callback) => {
fsub = new FloodSub(libp2pNode)
callback(null)
}),

sub: promisify((topic, options, callback) => {
if (typeof options === 'function') {
callback = options
options = {}
}

if (!self.isOnline()) {
return callback(OFFLINE_ERROR)
}

var rs = new Stream()
rs.readable = true
rs._read = () => {}
rs.cancel = () => {
fsub.unsubscribe(topic)
}

fsub.on(topic, (data) => {
console.log('PUBSUB DATA:', data.toString())
rs.emit('data', {
data: data.toString(),
topicIDs: [topic]
// these fields are currently missing from message
// (but are present in messages from go-ipfs pubsub)
// from: bs58.encode(message.from),
// seqno: Base64.decode(message.seqno)
})
})

fsub.subscribe(topic)
callback(null, rs)
}),

pub: promisify((topic, data, callback) => {
if (!self.isOnline()) {
return callback(OFFLINE_ERROR)
}

const buf = Buffer.isBuffer(data) ? data : new Buffer(data)

fsub.publish(topic, buf)
callback(null)
})

}
}
2 changes: 2 additions & 0 deletions src/core/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ const swarm = require('./components/swarm')
const ping = require('./components/ping')
const files = require('./components/files')
const bitswap = require('./components/bitswap')
const pubsub = require('./components/pubsub')

exports = module.exports = IPFS

Expand Down Expand Up @@ -67,4 +68,5 @@ function IPFS (repoInstance) {
this.files = files(this)
this.bitswap = bitswap(this)
this.ping = ping(this)
this.pubsub = pubsub(this)
}