mirror of
https://github.com/fluencelabs/js-libp2p
synced 2025-05-28 18:01:19 +00:00
Co-authored-by: Alan Shaw <alan.shaw@protocol.ai> Co-authored-by: Alan Shaw <alan@tableflip.io> Co-authored-by: Arnaud <arnaud.valensi@gmail.com> Co-authored-by: David Dias <daviddias.p@gmail.com> Co-authored-by: David Dias <mail@daviddias.me> Co-authored-by: Dmitriy Ryajov <dryajov@gmail.com> Co-authored-by: Francisco Baio Dias <xicombd@gmail.com> Co-authored-by: Friedel Ziegelmayer <dignifiedquire@gmail.com> Co-authored-by: Haad <haadcode@users.noreply.github.com> Co-authored-by: Hugo Dias <mail@hugodias.me> Co-authored-by: Hugo Dias <hugomrdias@gmail.com> Co-authored-by: Jacob Heun <jacobheun@gmail.com> Co-authored-by: Kevin Kwok <antimatter15@gmail.com> Co-authored-by: Kobi Gurkan <kobigurk@gmail.com> Co-authored-by: Maciej Krüger <mkg20001@gmail.com> Co-authored-by: Matteo Collina <matteo.collina@gmail.com> Co-authored-by: Michael Fakhry <fakhrimichael@live.com> Co-authored-by: Oli Evans <oli@tableflip.io> Co-authored-by: Pau Ramon Revilla <masylum@gmail.com> Co-authored-by: Pedro Teixeira <i@pgte.me> Co-authored-by: Pius Nyakoojo <piusnyakoojo@gmail.com> Co-authored-by: Richard Littauer <richard.littauer@gmail.com> Co-authored-by: Sid Harder <sideharder@gmail.com> Co-authored-by: Vasco Santos <vasco.santos@ua.pt> Co-authored-by: harrshasri <35241544+harrshasri@users.noreply.github.com> Co-authored-by: kumavis <kumavis@users.noreply.github.com> Co-authored-by: ᴠɪᴄᴛᴏʀ ʙᴊᴇʟᴋʜᴏʟᴍ <victorbjelkholm@gmail.com>
49 lines
1.3 KiB
JavaScript
49 lines
1.3 KiB
JavaScript
'use strict'
|
|
|
|
const map = require('pull-stream/throughs/map')
|
|
const EventEmitter = require('events')
|
|
|
|
/**
|
|
* Takes a Switch and returns an Observer that can be used in conjunction with
|
|
* observe-connection.js. The returned Observer comes with `incoming` and
|
|
* `outgoing` properties that can be used in pull streams to emit all metadata
|
|
* for messages that pass through a Connection.
|
|
*
|
|
* @param {Switch} swtch
|
|
* @returns {EventEmitter}
|
|
*/
|
|
module.exports = (swtch) => {
|
|
const observer = Object.assign(new EventEmitter(), {
|
|
incoming: observe('in'),
|
|
outgoing: observe('out')
|
|
})
|
|
|
|
swtch.on('peer-mux-established', (peerInfo) => {
|
|
observer.emit('peer:connected', peerInfo.id.toB58String())
|
|
})
|
|
|
|
swtch.on('peer-mux-closed', (peerInfo) => {
|
|
observer.emit('peer:closed', peerInfo.id.toB58String())
|
|
})
|
|
|
|
return observer
|
|
|
|
function observe (direction) {
|
|
return (transport, protocol, peerInfo) => {
|
|
return map((buffer) => {
|
|
willObserve(peerInfo, transport, protocol, direction, buffer.length)
|
|
return buffer
|
|
})
|
|
}
|
|
}
|
|
|
|
function willObserve (peerInfo, transport, protocol, direction, bufferLength) {
|
|
peerInfo.then((_peerInfo) => {
|
|
if (_peerInfo) {
|
|
const peerId = _peerInfo.id.toB58String()
|
|
observer.emit('message', peerId, transport, protocol, direction, bufferLength)
|
|
}
|
|
})
|
|
}
|
|
}
|