p2p

package
v0.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 22, 2025 License: LGPL-3.0 Imports: 23 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func LoadPeers

func LoadPeers(h host.Host, peers []*peer.AddrInfo)

LoadPeers clears out peerstore and loads new peers into it

func NewHost

func NewHost(privKey crypto.PrivKey, networkTopology *topology.NetworkTopology, cg *ConnectionGate, port uint16) (host.Host, error)

NewHost creates new host.Host from private key and relayer configuration

func ReadStream

func ReadStream(r *bufio.Reader) ([]byte, error)

ReadStream reads data from the given stream

func WriteStream

func WriteStream(msg []byte, w *bufio.Writer) error

WriteStream writes the message to stream

Types

type ConnectionGate

type ConnectionGate struct {
	// contains filtered or unexported fields
}

ConnectionGate implements libp2p ConnectionGater to prevent inbound and outbound requests to peers not specified in topology

func NewConnectionGate

func NewConnectionGate(topology *topology.NetworkTopology) *ConnectionGate

func (*ConnectionGate) InterceptAccept

func (cg *ConnectionGate) InterceptAccept(network.ConnMultiaddrs) (allow bool)

func (*ConnectionGate) InterceptAddrDial

func (cg *ConnectionGate) InterceptAddrDial(peer.ID, ma.Multiaddr) (allow bool)

func (*ConnectionGate) InterceptPeerDial

func (cg *ConnectionGate) InterceptPeerDial(p peer.ID) (allow bool)

func (*ConnectionGate) InterceptSecured

func (cg *ConnectionGate) InterceptSecured(nd network.Direction, p peer.ID, cm network.ConnMultiaddrs) (allow bool)

func (*ConnectionGate) InterceptUpgraded

func (cg *ConnectionGate) InterceptUpgraded(network.Conn) (allow bool, reason control.DisconnectReason)

func (*ConnectionGate) SetTopology

func (cg *ConnectionGate) SetTopology(topology *topology.NetworkTopology)

type Libp2pCommunication

type Libp2pCommunication struct {
	SessionSubscriptionManager
	// contains filtered or unexported fields
}

func NewCommunication

func NewCommunication(h host.Host, protocolID protocol.ID) Libp2pCommunication

func (Libp2pCommunication) Broadcast

func (c Libp2pCommunication) Broadcast(
	peers peer.IDSlice,
	msg []byte,
	msgType comm.MessageType,
	sessionID string,
) error

func (Libp2pCommunication) CloseSession

func (c Libp2pCommunication) CloseSession(sessionID string)

func (Libp2pCommunication) ProcessMessagesFromStream

func (c Libp2pCommunication) ProcessMessagesFromStream(s network.Stream)

func (Libp2pCommunication) StreamHandlerFunc

func (c Libp2pCommunication) StreamHandlerFunc(s network.Stream)

func (Libp2pCommunication) Subscribe

func (c Libp2pCommunication) Subscribe(
	sessionID string,
	msgType comm.MessageType,
	channel chan *comm.WrappedMessage,
) comm.SubscriptionID

func (Libp2pCommunication) UnSubscribe

func (c Libp2pCommunication) UnSubscribe(
	subID comm.SubscriptionID,
)

type SessionSubscriptionManager

type SessionSubscriptionManager struct {
	// contains filtered or unexported fields
}

SessionSubscriptionManager manages channel subscriptions by comm.SessionID

func NewSessionSubscriptionManager

func NewSessionSubscriptionManager() SessionSubscriptionManager

func (*SessionSubscriptionManager) GetSubscribers

func (ms *SessionSubscriptionManager) GetSubscribers(
	sessionID string,
	msgType comm.MessageType,
) []chan *comm.WrappedMessage

func (*SessionSubscriptionManager) SubscribeTo

func (ms *SessionSubscriptionManager) SubscribeTo(
	sessionID string, msgType comm.MessageType, channel chan *comm.WrappedMessage,
) comm.SubscriptionID

func (*SessionSubscriptionManager) UnSubscribeFrom

func (ms *SessionSubscriptionManager) UnSubscribeFrom(
	subscriptionID comm.SubscriptionID,
)

type StreamManager

type StreamManager struct {
	// contains filtered or unexported fields
}

StreamManager manages instances of network.Stream

Each stream is mapped to a specific session, by sessionID

func NewStreamManager

func NewStreamManager() *StreamManager

NewStreamManager creates new StreamManager

func (*StreamManager) AddStream

func (sm *StreamManager) AddStream(sessionID string, peerID peer.ID, stream network.Stream)

AddStream saves and maps provided stream to sessionID

func (*StreamManager) ReleaseStreams

func (sm *StreamManager) ReleaseStreams(sessionID string)

ReleaseStream removes reference on streams mapped to provided sessionID and closes them

func (*StreamManager) Stream

func (sm *StreamManager) Stream(sessionID string, peerID peer.ID) (network.Stream, error)

Stream fetches stream by peer and session ID

Directories

Path Synopsis
mock
conn
Package mock_network is a generated GoMock package.
Package mock_network is a generated GoMock package.
host
Package mock_host is a generated GoMock package.
Package mock_host is a generated GoMock package.
stream
Package mock_network is a generated GoMock package.
Package mock_network is a generated GoMock package.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL