本文档描述了 p2p 层向协议层提供的 API,即向已注册的 reactor 提供的 API。 该 API 由两个接口构成:一个是由 Switch 实例提供的接口,另一些是由多个 Peer 实例提供的接口,每个已连接的对等节点对应一个 Peer 实例。Switch 实例会作为 reactor 的注册流程的一部分提供给每个 reactor。每当与某个对等节点建立新连接时,多个 Peer 实例会提供给每个已注册的 reactor。
注意 出于实际原因,接口被拆分为 Switch 和 Peer 两部分,相关讨论可在知识库仓库中找到更详细的说明。

Switch API

Switch 是 p2p 层实现的核心组件。它管理节点中运行的所有 reactor,并跟踪与各个对等节点的连接。 下表总结了标准 reactor 与 Switch 的交互:
Switch API methodconsensusblock syncstate syncmempoolevidencePEX
Peers() IPeerSetxxx
NumPeers() (int, int, int)xx
Broadcast(Envelope) chan boolxxx
MarkPeerAsGood(Peer)x
StopPeerForError(Peer, interface{})xxxxxx
StopPeerGracefully(Peer)x
Reactor(string) Reactorx
上述列表并不完整,因为它未包含 PEX reactor 调用的所有 Switch 方法。PEX reactor 是一个特殊组件,应被视为 p2p 层的一部分。本文档不涵盖作为连接管理器的 PEX reactor 的运作方式。

对等节点状态

Switch API 中的前两个方法允许 reactor 查询 p2p 层的状态:已连接对等节点的集合。

    func (sw *Switch) Peers() IPeerSet
Peers() 方法返回当前已连接对等节点的集合。 返回的 IPeerSet 是该集合的一个不可变且并发安全的副本。 请注意,此方法返回的 Peer 处理器此前已经通过 InitPeer(Peer) 方法添加到 reactor 中,但尚未通过 RemovePeer(Peer) 方法移除。 因此,原则上 reactor 应当已经掌握这些信息。

    func (sw *Switch) NumPeers() (outbound, inbound, dialing int)
NumPeers() 方法返回当前已连接对等节点的数量,并区分 outbound 和 inbound 对等节点。 outbound 对等节点是节点主动拨号连接的对等节点,而 inbound 对等节点则是节点接受其连接请求的对等节点。 第三个字段 dialing 表示节点当前正在尝试连接、但尚未连接成功的对等节点数量。
注意 NumPeers() 返回的第三个字段,即处于 dialing 状态的对等节点数量,并不是协议层应当关心的信息。 实际上,除可视为 p2p 层实现一部分的 PEX reactor 外,没有任何标准 reactor 真正使用这项信息;在后续重构该接口时,这项信息可以被移除。

广播

出于历史原因或向后兼容的考虑,switch 提供了一个向所有已连接对等节点发送消息的方法:

    func (sw *Switch) Broadcast(e Envelope) chan bool
Broadcast() 方法是非阻塞的,并返回一个布尔值通道。 对于每个已连接的 Peer,它都会启动一个后台线程,使用 Peer.Send() 方法向该对等节点发送消息 (该方法是阻塞的,详见发送方法)。 每次单播发送操作的结果(成功或失败)都会写入返回的通道中,待所有操作完成后,该通道会被关闭。
注意
  • 当前 Switch.Broadcast(Envelope) 方法的_实现_效率不高,因为提供消息的编组是在 Peer.Send(Envelope) 辅助方法中完成的,也就是说,每个已连接对等节点都会执行一次。
  • 使用该方法的标准 reactor 都不会处理广播方法的返回值。其中一个原因是,无法将返回通道中的每个布尔输出与具体的某个对等节点对应起来。

对对等节点进行评估

p2p 层依赖已注册的 reactor 来评估对等节点的“质量”。 reactor 可以调用以下方法,通知 p2p 层某个对等节点表现出了“良好”行为。 该信息会记录到节点的地址簿中,并影响对等节点交换(PEX)协议的运作,因为节点发现过程会对“良好”对等节点有所偏向:

    func (sw *Switch) MarkPeerAsGood(peer Peer)
目前,由共识 reactor 负责对对等节点进行评估。 在当前逻辑中,只要共识协议从某个对等节点收集到 votesToContributeToBecomeGoodPeer = 10000 的整数倍有效投票,或 blocksToContributeToBecomeGoodPeer = 10000 的整数倍有效区块分片,就会将该对等节点标记为 good。 所谓“有效”,是指共识实现认为这些消息是合法的,并且节点是在预期接收此类信息的时机收到这些消息的,因此不包括重复消息或延迟收到的消息。
注意 当前 switch 并未提供将对等节点标记为坏节点的方法。 实际上,当前版本 p2p 层中的对等节点质量管理并没有真正实现。 这个话题正在知识库仓库中讨论。

停止对等节点

reactor 可以指示 p2p 层断开与某个对等节点的连接。 按照 p2p 层的术语,reactor 是在请求停止某个对等节点。 该对等节点的发送和接收例程实际上会被停止,从而中断与该对等节点的通信。 随后,该 Peer 会通过 RemovePeer(Peer) 方法从每个已注册的 reactor 中移除,并从已连接对等节点集合中删除。

    func (sw *Switch) StopPeerForError(peer Peer, reason interface{})
所有标准 reactor 都会在发生错误时使用上述方法断开与某个对等节点的连接。 这里的错误,是指在处理从某个 Peer 收到的消息时发生的错误。 生成的 error 会作为 reason 参数传给该方法。 StopPeerForError() 方法有一个重要的注意事项:如果要停止的对等节点被配置为_持久对等节点_,switch 将尝试重新连接到这个相同的对等节点。 当该方法由 p2p 层的其他组件调用时(例如发生通信错误的情况),这种行为是合理的;但当它由 reactor 调用时,这种行为就没有意义了。
注意 关于这个主题的更完整讨论,可参见知识库仓库。
func (sw *Switch) StopPeerGracefully(peer Peer) 第二个方法会指示 switch 在没有特定原因的情况下断开与某个对等节点的连接。 该方法仅由运行在_seed mode_ 下的节点的 PEX reactor 使用,因为 seed 节点会在与对等节点交换完地址后断开连接。

Reactor 表

switch 会跟踪所有已注册的 reactor,并按唯一的 reactor 名称建立索引。 因此,reactor 可以通过 switch 根据 name 访问另一个 Reactor:

    func (sw *Switch) Reactor(name string) Reactor
该方法目前仅被 Block Sync reactor 用来访问 Consensus reactor 的实现,并调用其导出的 SwitchToConsensus() 方法。 尽管这种方式可用,但不建议采用这种 reactor 间交互方式,并且应尽量避免,因为它违背了 reactor 彼此独立的假设。

Peer API

Peer 接口表示一个已连接的对等节点。 Peer 实例封装了一个多路复用连接,该连接实现了与对等节点之间的实际通信(发送和接收消息)。 当与某个对等节点建立连接时,Switch 会将相应的 Peer 实例提供给所有已注册的 reactor。 从这一时刻起,reactor 就可以使用这个新 Peer 实例的方法。 下表总结了标准 reactor 与已连接对等节点的交互,以及它们使用的 Peer 方法:
Peer API methodconsensusblock syncstate syncmempoolevidencePEX
ID() IDxxxxxx
IsRunning() boolxxx
Quit() <-chan struct{}xx
Get(string) interface{}xxx
Set(string, interface{})x
Send(Envelope) boolxxxxxx
TrySend(Envelope) boolxx
上述列表并不完整,因为它未包含 PEX reactor 调用的所有 Peer 方法。PEX reactor 是一个特殊组件,应被视为 p2p 层的一部分。本文档不涵盖作为连接管理器的 PEX reactor 的运作方式。

标识

p2p 网络中的节点都配置了唯一的密码学密钥对。 该密钥对的公钥部分会在与对等节点建立连接时、作为认证握手的一部分进行校验,并构成该对等节点的 ID:

    func (p Peer) ID() p2p.ID
请注意,节点每次连接到某个对等节点时(例如断开后重新连接),都会通过 InitPeer(Peer) 方法向 reactor 提供一个新的(不同的)Peer 处理器。 实际上,Peer 处理器关联的是与某个对等节点的_连接_,而不是网络中的实际_节点_。 要跟踪实际的对等节点,应使用上述方法提供的唯一对等节点 p2p.ID。

Peer 状态

交换机在通过 AddPeer(Peer) 方法将对等节点添加到每个已注册的 reactor 之前,会先启动该对等节点的发送和接收例程。 随后,reactor 通常会使用收到的 Peer 处理器,启动与这个新连接的对等节点交互的例程。 对于这些例程,检查该对等节点是否仍然保持连接,以及其发送和接收例程是否仍在运行,会很有用:

    func (p Peer) IsRunning() bool
    func (p Peer) Quit() <-chan struct{}
上述两个方法以两种不同方式提供了关于 Peer 实例状态的相同信息。 它们都定义在 Service 接口中。 IsRunning() 方法是同步的,返回该对等节点是否已启动且尚未停止。 Quit() 方法返回一个通道,当对等节点停止时该通道会被关闭;它是一种异步状态查询。

键值存储

每个 Peer 实例都提供了一个同步的键值存储,用于在 reactor 之间共享特定于该对等节点的状态:

    func (p Peer) Get(key string) interface{}
    func (p Peer) Set(key string, data interface{})
这个键值存储可以看作是在 reactor 之间交换对等节点状态的一种异步机制。 在该机制当前的使用场景中,Consensus reactor 会为每个已连接的对等节点在键值存储中填充一个 PeerState 实例。 与某个对等节点交互的 Consensus reactor 例程会读取并更新这份共享的对等节点状态。 而 Evidence 和 Mempool reactor 则会周期性查询每个对等节点的键值存储,尤其是获取该对等节点上报的最新高度。 这些由 Consensus reactor 产出的信息,会影响这两个 reactor 与其对等节点之间的交互。
说明 关于这个键值存储如何用于在 reactor 之间共享状态的更多细节,可参见 knowledge-base 仓库。

发送方法

最后,Peer 实例允许某个 reactor 向该对等节点上运行的对应 reactor 发送消息。 当交换机向已注册的 reactor 提供 Peer 实例时,这正是其最终目的。 有两种发送消息的方法:

    func (p Peer) Send(e Envelope) bool
    func (p Peer) TrySend(e Envelope) bool
这两个消息发送方法接收一个 Envelope,其内容应按如下方式设置:
  • ChannelID:消息应通过的通道,它决定了将处理该消息的 reactor;
  • Src:该字段表示传入消息的来源,对于传出消息无关紧要;
  • Message:消息的实际负载,使用 protocol buffers 进行编组。
这两个消息发送方法会尝试将消息(e.Payload)加入到该对等节点目标通道(e.ChannelID)的发送队列中。 对等节点支持的每个已注册通道都有一个发送队列,并且每个发送队列都有容量限制。 每个通道发送队列的容量由 reactor 通过相应的 ChannelDescriptor 进行配置。 这两个消息发送方法会返回是否成功将编组后的消息加入该通道的发送队列。 这些方法返回 false 的最常见原因是该通道的发送队列已满。 返回 false 的其他原因还包括:对等节点已停止、提供了未注册的通道 ID,或在编组消息负载时发生错误。 这两个消息发送方法的区别在于它们在何时返回 false。 Send() 方法是一个_阻塞_方法;如果由于通道发送队列仍然已满,导致消息在 10 秒的_超时_后仍无法入队,它就会返回 false。 TrySend() 方法是一个_非阻塞_方法;当通道发送队列已满时,它会_立即_返回 false。
This document describes the API provided by the p2p layer to the protocol layer, namely to the registered reactors. This API consists of two interfaces: the one provided by the Switch instance, and the ones provided by multiple Peer instances, one per connected peer. The Switch instance is provided to every reactor as part of the reactor’s registration procedure. The multiple Peer instances are provided to every registered reactor whenever a new connection with a peer is established.
Note The practical reasons that lead to the interface to be provided in two parts, Switch and Peer instances are discussed in more datail in the knowledge-base repository.

Switch API

The Switch is the central component of the p2p layer implementation. It manages all the reactors running in a node and keeps track of the connections with peers. The table below summarizes the interaction of the standard reactors with the Switch:
Switch API methodconsensusblock syncstate syncmempoolevidencePEX
Peers() IPeerSetxxx
NumPeers() (int, int, int)xx
Broadcast(Envelope) chan boolxxx
MarkPeerAsGood(Peer)x
StopPeerForError(Peer, interface{})xxxxxx
StopPeerGracefully(Peer)x
Reactor(string) Reactorx
The above list is not exhaustive as it does not include all the Switch methods invoked by the PEX reactor, a special component that should be considered part of the p2p layer. This document does not cover the operation of the PEX reactor as a connection manager.

Peers State

The first two methods in the switch API allow reactors to query the state of the p2p layer: the set of connected peers.

    func (sw *Switch) Peers() IPeerSet
The Peers() method returns the current set of connected peers. The returned IPeerSet is an immutable concurrency-safe copy of this set. Observe that the Peer handlers returned by this method were previously added to the reactor via the InitPeer(Peer) method, but not yet removed via the RemovePeer(Peer) method. Thus, a priori, reactors should already have this information.

    func (sw *Switch) NumPeers() (outbound, inbound, dialing int)
The NumPeers() method returns the current number of connected peers, distinguished between outbound and inbound peers. An outbound peer is a peer the node has dialed to, while an inbound peer is a peer the node has accepted a connection from. The third field dialing reports the number of peers to which the node is currently attempting to connect, so not (yet) connected peers.
Note The third field returned by NumPeers(), the number of peers in dialing state, is not an information that should regard the protocol layer. In fact, with the exception of the PEX reactor, which can be considered part of the p2p layer implementation, no standard reactor actually uses this information, that could be removed when this interface is refactored.

Broadcast

The switch provides, mostly for historical or retro-compatibility reasons, a method for sending a message to all connected peers:

    func (sw *Switch) Broadcast(e Envelope) chan bool
The Broadcast() method is not blocking and returns a channel of booleans. For every connected Peer, it starts a background thread for sending the message to that peer, using the Peer.Send() method (which is blocking, as detailed in Send Methods). The result of each unicast send operation (success or failure) is added to the returned channel, which is closed when all operations are completed.
Note
  • The current implementation of the Switch.Broadcast(Envelope) method is not efficient, as the marshalling of the provided message is performed as part of the Peer.Send(Envelope) helper method, that is, once per connected peer.
  • The return value of the broadcast method is not considered by any of the standard reactors that employ the method. One of the reasons is that is is not possible to associate each of the boolean outputs added to the returned channel to a peer.

Vetting Peers

The p2p layer relies on the registered reactors to gauge the quality of peers. The following method can be invoked by a reactor to inform the p2p layer that a peer has presented a “good” behaviour. This information is registered in the node’s address book and influences the operation of the Peer Exchange (PEX) protocol, as node discovery adopts a bias towards “good” peers:

    func (sw *Switch) MarkPeerAsGood(peer Peer)
At the moment, it is up to the consensus reactor to vet a peer. In the current logic, a peer is marked as good whenever the consensus protocol collects a multiple of votesToContributeToBecomeGoodPeer = 10000 useful votes or blocksToContributeToBecomeGoodPeer = 10000 useful block parts from that peer. By “useful”, the consensus implementation considers messages that are valid and that are received by the node when the node is expected for such information, which excludes duplicated or late received messages.
Note The switch doesn’t currently provide a method to mark a peer as a bad peer. In fact, the peer quality management is really implemented in the current version of the p2p layer. This topic is being discussed in the knowledge-base repository.

Stopping Peers

Reactors can instruct the p2p layer to disconnect from a peer. Using the p2p layer’s nomenclature, the reactor requests a peer to be stopped. The peer’s send and receive routines are in fact stopped, interrupting the communication with the peer. The Peer is then removed from every registered reactor, using the RemovePeer(Peer) method, and from the set of connected peers.

    func (sw *Switch) StopPeerForError(peer Peer, reason interface{})
All the standard reactors employ the above method for disconnecting from a peer in case of errors. These are errors that occur when processing a message received from a Peer. The produced error is provided to the method as the reason. The StopPeerForError() method has an important caveat: if the peer to be stopped is configured as a persistent peer, the switch will attempt reconnecting to that same peer. While this behaviour makes sense when the method is invoked by other components of the p2p layer (e.g., in the case of communication errors), it does not make sense when it is invoked by a reactor.
Note A more comprehensive discussion regarding this topic can be found on the knowledge-base repository.
func (sw *Switch) StopPeerGracefully(peer Peer) The second method instructs the switch to disconnect from a peer for no particular reason. This method is only adopted by the PEX reactor of a node operating in seed mode, as seed nodes disconnect from a peer after exchanging peer addresses with it.

Reactors Table

The switch keeps track of all registered reactors, indexed by unique reactor names. A reactor can therefore use the switch to access another Reactor from its name:

    func (sw *Switch) Reactor(name string) Reactor
This method is currently only used by the Block Sync reactor to access the Consensus reactor implementation, from which it uses the exported SwitchToConsensus() method. While available, this inter-reactor interaction approach is discouraged and should be avoided, as it violates the assumption that reactors are independent.

Peer API

The Peer interface represents a connected peer. A Peer instance encapsulates a multiplex connection that implements the actual communication (sending and receiving messages) with a peer. When a connection is established with a peer, the Switch provides the corresponding Peer instance to all registered reactors. From this point, reactors can use the methods of the new Peer instance. The table below summarizes the interaction of the standard reactors with connected peers, with the Peer methods used by them:
Peer API methodconsensusblock syncstate syncmempoolevidencePEX
ID() IDxxxxxx
IsRunning() boolxxx
Quit() <-chan struct{}xx
Get(string) interface{}xxx
Set(string, interface{})x
Send(Envelope) boolxxxxxx
TrySend(Envelope) boolxx
The above list is not exhaustive as it does not include all the Peer methods invoked by the PEX reactor, a special component that should be considered part of the p2p layer. This document does not cover the operation of the PEX reactor as a connection manager.

Identification

Nodes in the p2p network are configured with a unique cryptographic key pair. The public part of this key pair is verified when establishing a connection with the peer, as part of the authentication handshake, and constitutes the peer’s ID:

    func (p Peer) ID() p2p.ID
Observe that each time the node connects to a peer (e.g., after disconnecting from it), a new (distinct) Peer handler is provided to the reactors via InitPeer(Peer) method. In fact, the Peer handler is associated to a connection with a peer, not to the actual node in the network. To keep track of actual peers, the unique peer p2p.ID provided by the above method should be employed.

Peer state

The switch starts the peer’s send and receive routines before adding the peer to every registered reactor using the AddPeer(Peer) method. The reactors then usually start routines to interact with the new connected peer using the received Peer handler. For these routines it is useful to check whether the peer is still connected and its send and receive routines are still running:

    func (p Peer) IsRunning() bool
    func (p Peer) Quit() <-chan struct{}
The above two methods provide the same information about the state of a Peer instance in two different ways. Both of them are defined in the Service interface. The IsRunning() method is synchronous and returns whether the peer has been started and has not been stopped. The Quit() method returns a channel that is closed when the peer is stopped; it is an asynchronous state query.

Key-value store

Each Peer instance provides a synchronized key-value store that allows sharing peer-specific state between reactors:

    func (p Peer) Get(key string) interface{}
    func (p Peer) Set(key string, data interface{})
This key-value store can be seen as an asynchronous mechanism to exchange the state of a peer between reactors. In the current use-case of this mechanism, the Consensus reactor populates the key-value store with a PeerState instance for each connected peer. The Consensus reactor routines interacting with a peer read and update the shared peer state. The Evidence and Mempool reactors, in their turn, periodically query the key-value store of each peer for retrieving, in particular, the last height reported by the peer. This information, produced by the Consensus reactor, influences the interaction of these two reactors with their peers.
Note More details of how this key-value store is used to share state between reactors can be found on the knowledge-base repository.

Send methods

Finally, a Peer instance allows a reactor to send messages to companion reactors running at that peer. This is ultimately the goal of the switch when it provides Peer instances to the registered reactors. There are two methods for sending messages:

    func (p Peer) Send(e Envelope) bool
    func (p Peer) TrySend(e Envelope) bool
The two message-sending methods receive an Envelope, whose content should be set as follows:
  • ChannelID: the channel the message should be sent through, which defines the reactor that will process the message;
  • Src: this field represents the source of an incoming message, which is irrelevant for outgoing messages;
  • Message: the actual message’s payload, which is marshalled using protocol buffers.
The two message-sending methods attempt to add the message (e.Payload) to the send queue of the peer’s destination channel (e.ChannelID). There is a send queue for each registered channel supported by the peer, and each send queue has a capacity. The capacity of the send queues for each channel are configured by reactors via the corresponding ChannelDescriptor. The two message-sending methods return whether it was possible to enqueue the marshalled message to the channel’s send queue. The most common reason for these methods to return false is the channel’s send queue being full. Further reasons for returning false are: the peer being stopped, providing a non-registered channel ID, or errors when marshalling the message’s payload. The difference between the two message-sending methods is when they return false. The Send() method is a blocking method, it returns false if the message could not be enqueued, because the channel’s send queue is still full, after a 10-second timeout. The TrySend() method is a non-blocking method, it immediately returns false when the channel’s send queue is full.