组件必须实现 p2p.Reactor 接口,才能使用 p2p 层提供的通信服务。 该接口目前也是 reactor 的主要文档来源。 本文档的目标是说明 p2p 通信层在与 reactor 交互时的行为。 因此,虽然 Reactor 接口声明了会被调用的方法,并规定了 p2p 层对 reactor 的期望, 本文档关注的是 reactor 实现应当从 p2p 层预期到的时序行为。 也就是说,这些函数可能会以什么顺序被调用。 本规范配套提供了 reactor.qnt 文件, 这是一个使用Quint编写的、更完整的 reactor 运行模型。Quint 是一种可执行规范语言。 Reactor 接口中声明的方法在 Quint 中被建模为 pure def 方法, 并给出了一些关于这些方法应如何实现的示例。 p2p 层通过调用这些接口方法与 reactor 交互时的行为, 则以状态迁移的形式建模,在 Quint 术语中称为 action。

概览

下面的_语法_是对 p2p 层调用 reactor 的预期调用序列的简化表示。 请注意,该语法描述的是与_单个 reactor_相关的事件,而 p2p 层支持同时运行多个 reactor。 如需更详细地了解 p2p 层到各个 reactor 的调用序列, 请参阅配套的 Quint 模型。 语法适合用于概览 reactor 的运行方式, 但它在可表达的行为上也存在一些限制。 例如,下面的语法只表示了_单个 peer_的管理, 即具有某个给定 ID 的 peer,它可以与节点多次连接、断开并重连。 p2p 层以及每个 reactor 都应能够并行处理多个不同的 peer。 这意味着下面语法中的多个非终结符 peer-management 出现时, 都可以彼此独立并并行“运行”,每一个都对应并产生与不同 peer 相关的事件:
start           = registration on-start *peer-management on-stop
registration    = get-channels set-switch

; Refers to a single peer, a reactor must support multiple concurrent peers
peer-management = init-peer start-peer stop-peer
start-peer      = [*receive] (connected-peer / start-error)
connected-peer  = add-peer *receive
stop-peer       = [peer-error] remove-peer

; Service interface
on-start        = %s"OnStart()"
on-stop         = %s"OnStop()"
; Reactor interface
get-channels    = %s"GetChannels()"
set-switch      = %s"SetSwitch(*Switch)"
init-peer       = %s"InitPeer(Peer)"
add-peer        = %s"AddPeer(Peer)"
remove-peer     = %s"RemovePeer(Peer, reason)"
receive         = %s"Receive(Envelope)"

; Errors, for reference
start-error     = %s"log(Error starting peer)"
peer-error      = %s"log(Stopping peer for error)"
该语法使用大小写敏感的扩展巴科斯-瑙尔范式(ABNF, 定义见 IETF RFC 7405)。 它借鉴了用于说明 CometBFT 与 ABCI++ 应用交互的语法, 可在这里查看。

注册

要成为一个 reactor,组件首先必须实现 Reactor 接口, 然后通过 Switch.AddReactor(name string, reactor Reactor) 方法 将该实现注册到 p2p 层,并为该 reactor 提供一个全局唯一的 name。 注册必须发生在节点整体启动之前,尤其必须发生在 p2p 层启动之前。 换句话说,当前不支持在运行中的节点上注册 reactor: reactor 必须作为节点初始化配置的一部分完成注册。
registration    = get-channels set-switch
p2p 层会通过 GetChannels() 方法,从 reactor 获取其负责的 channel 列表。 此后,reactor 实现应当预期接收到 p2p 层在这些已声明 channel 上收到的所有消息。 第二个方法 SetSwitch(Switch) 完成了 reactor 与 p2p 层之间的握手。 Switch 是 p2p 层的核心组件,负责与 peer 建立连接以及路由消息。 Switch 实例为所有已注册的 reactor 提供了一组方法, 这些方法记录在配套文档 API for Reactors 中。

服务接口

reactor 必须实现 Service 接口, 特别是启动方法 OnStart() 和关闭方法 OnStop():
start           = registration on-start *peer-management on-stop
作为节点启动过程的一部分,所有已注册的 reactor 都会由 p2p 层启动。 当节点关闭时,所有已注册的 reactor 也都会由 p2p 层停止。 需要注意的是,Service 接口规范规定一个 service 只能启动和停止一次。 因此,在被 p2p 层启动之前,或者在被停止之后,reactor 都不应预期发生任何交互。

Peer 管理

reactor 运行的核心,是与 peer 交互,或者更准确地说, 与运行在连接到该节点的 peer 上、并使用相同 channel 的对应 reactor 交互。 下面的语法片段表示 reactor 与单个 peer 的交互:
; Refers to a single peer, a reactor must support multiple concurrent peers
peer-management = init-peer start-peer stop-peer
当 p2p 层与某个 Peer 建立连接时,会通过 InitPeer(Peer) 方法通知所有已注册的 reactor。 调用该方法时,Peer 尚未启动, 也就是说,与该 peer 发送和接收消息相关的例程还没有运行。 这个方法应用于初始化与新 peer 相关的状态或数据, 但不应在此阶段与其交互。 下一步是启动与这个新 Peer 的通信例程。 如后文所述,这个过程可能成功,也可能失败。 无论如何,最终该 peer 都会被停止,从而结束对这个 Peer 实例的管理。

启动 Peer

当对每个已注册 reactor 都调用完 InitPeer(Peer) 之后,p2p 层会启动该 peer 的通信例程,并将这个 Peer 加入已连接 peer 集合。 如果这两个步骤都无错误完成,就会调用 reactor 的 AddPeer(Peer):
start-peer      = [*receive] (connected-peer / start-error)
connected-peer  = add-peer *receive
如果发生错误,会记录一条日志,说明 p2p 层启动该 peer 失败。 这并不是常见场景,通常只会在与行为异常或响应缓慢的 peer 交互时出现。 一个实际示例记录在这个issue中。 如何处理 AddPeer(Peer) 事件由 reactor 自行定义。 典型行为是启动一些例程,在满足某些条件或发生某些事件时, 使用提供的 Peer 实例向新加入的 peer 发送消息。 配套文档 API for Reactors 说明了 Peer 实例提供的方法,这些方法从 peer 被添加到 reactor 时起即可使用。

停止 Peer

当 p2p 层与某个 Peer 断开连接时,会通过 RemovePeer(Peer, reason) 方法通知所有已注册的 reactor:
stop-peer       = [peer-error] remove-peer
该方法会在 p2p 层停止 peer 的发送和接收例程之后被调用。 根据停止该 peer 的 reason 不同,可能会产生不同的日志消息。 当一个 peer 从所有 reactor 中移除后,该 Peer 实例也会从已连接 peer 集合中移除。 这使得同一个 peer 可以重新连接,并针对新的连接再次调用 InitPeer(Peer)。 一旦某个 Peer 被移除,reactor 就不应再收到来自该 peer 的任何消息, 并且也绝不能再尝试向这个已移除的 peer 发送消息。 这通常意味着要停止那些由对应 Add(Peer) 方法启动的例程。

接收消息

reactor 的主要职责,是处理其向 p2p 层注册的各个 channel 上收到的入站消息。 从某个 Peer 接收消息的_前置条件_是 p2p 层此前已经调用过 InitPeer(Peer)。 这意味着 reactor 必须能够在调用 AddPeer(Peer)_之前_接收来自某个 Peer 的消息。 之所以会这样,是因为 peer 的发送和接收例程会更早启动, 并且当 p2p 层将该 peer 添加到每个已注册 reactor 时,这些例程应该已经在运行。
start-peer      = [*receive] (connected-peer / start-error)
connected-peer  = add-peer *receive
不过,更常见的情况仍然是在调用 AddPeer(Peer) 之后才开始接收来自该 peer 的消息。 在该 peer 被停止并调用 RemovePeer(Peer) 之前, 都可能接收到任意数量的消息。 当已连接 peer 在 reactor 已注册的任意 channel 上发送来消息时, p2p 层会通过 Receive(Envelope) 方法将消息投递给 reactor。 该消息会被封装为一个 Envelope,其中包含:
  • ChannelID:消息所属的 channel
  • Src:源 Peer 句柄,即消息的接收来源
  • Message:消息的实际负载,使用 protocol buffers 反序列化后得到
关于 Receive 方法的实现,有两个重要注意点:
  1. 并发性:实现应考虑 Receive 方法可能被并发调用, 并携带来自不同 peer 的消息,因为与不同 peer 的交互是相互独立的,消息也可以并行接收。
  2. 非阻塞:Receive 方法的实现预期不应阻塞, 因为它是由接收例程直接调用的。 换句话说,只要 Receive 尚未返回,同一发送方的其他消息就不会被投递到任何 reactor。

A component has to implement the p2p.Reactor interface in order to use communication services provided by the p2p layer. This interface is currently the main source of documentation for a reactor. The goal of this document is to specify the behaviour of the p2p communication layer when interacting with a reactor. So while the Reactor interface declares the methods invoked and determines what the p2p layer expects from a reactor, this documentation focuses on the temporal behaviour that a reactor implementation should expect from the p2p layer. (That is, in which orders the functions may be called) This specification is accompanied by the reactor.qnt file, a more comprehensive model of the reactor’s operation written in Quint, an executable specification language. The methods declared in the Reactor interface are modeled in Quint, in the form of pure def methods, providing some examples of how they should be implemented. The behaviour of the p2p layer when interacting with a reactor, by invoking the interface methods, is modeled in the form of state transitions, or actions in the Quint nomenclature.

Overview

The following grammar is a simplified representation of the expected sequence of calls from the p2p layer to a reactor. Note that the grammar represents events referring to a single reactor, while the p2p layer supports the execution of multiple reactors. For a more detailed representation of the sequence of calls from the p2p layer to reactors, please refer to the companion Quint model. While useful to provide an overview of the operation of a reactor, grammars have some limitations in terms of the behaviour they can express. For instance, the following grammar only represents the management of a single peer, namely of a peer with a given ID which can connect, disconnect, and reconnect multiple times to the node. The p2p layer and every reactor should be able to handle multiple distinct peers in parallel. This means that multiple occurrences of non-terminal peer-management of the grammar below can “run” independently and in parallel, each one referring and producing events associated to a different peer:
start           = registration on-start *peer-management on-stop
registration    = get-channels set-switch

; Refers to a single peer, a reactor must support multiple concurrent peers
peer-management = init-peer start-peer stop-peer
start-peer      = [*receive] (connected-peer / start-error)
connected-peer  = add-peer *receive
stop-peer       = [peer-error] remove-peer

; Service interface
on-start        = %s"OnStart()"
on-stop         = %s"OnStop()"
; Reactor interface
get-channels    = %s"GetChannels()"
set-switch      = %s"SetSwitch(*Switch)"
init-peer       = %s"InitPeer(Peer)"
add-peer        = %s"AddPeer(Peer)"
remove-peer     = %s"RemovePeer(Peer, reason)"
receive         = %s"Receive(Envelope)"

; Errors, for reference
start-error     = %s"log(Error starting peer)"
peer-error      = %s"log(Stopping peer for error)"
The grammar is written in case-sensitive Augmented Backus–Naur form (ABNF, specified in IETF RFC 7405). It is inspired on the grammar produced to specify the interaction of CometBFT with an ABCI++ application, available here.

Registration

To become a reactor, a component has first to implement the Reactor interface, then to register the implementation with the p2p layer, using the Switch.AddReactor(name string, reactor Reactor) method, with a global unique name for the reactor. The registration must happen before the node, in general, and the p2p layer, in particular, are started. In other words, there is no support for registering a reactor on a running node: reactors must be registered as part of the setup of a node.
registration    = get-channels set-switch
The p2p layer retrieves from the reactor a list of channels the reactor is responsible for, using the GetChannels() method. The reactor implementation should thereafter expect the delivery of every message received by the p2p layer in the informed channels. The second method SetSwitch(Switch) concludes the handshake between the reactor and the p2p layer. The Switch is the main component of the p2p layer, being responsible for establishing connections with peers and routing messages. The Switch instance provides a number of methods for all registered reactors, documented in the companion API for Reactors document.

Service interface

A reactor must implement the Service interface, in particular, a startup OnStart() and a shutdown OnStop() methods:
start           = registration on-start *peer-management on-stop
As part of the startup of a node, all registered reactors are started by the p2p layer. And when the node is shut down, all registered reactors are stopped by the p2p layer. Observe that the Service interface specification establishes that a service can be started and stopped only once. So before being started or once stopped by the p2p layer, the reactor should not expect any interaction.

Peer management

The core of a reactor’s operation is the interaction with peers or, more precisely, with companion reactors operating on the same channels in peers connected to the node. The grammar extract below represents the interaction of the reactor with a single peer:
; Refers to a single peer, a reactor must support multiple concurrent peers
peer-management = init-peer start-peer stop-peer
The p2p layer informs all registered reactors when it establishes a connection with a Peer, using the InitPeer(Peer) method. When this method is invoked, the Peer has not yet been started, namely the routines for sending messages to and receiving messages from the peer are not running. This method should be used to initialize state or data related to the new peer, but not to interact with it. The next step is to start the communication routines with the new Peer. As detailed in the following, this procedure may or may not succeed. In any case, the peer is eventually stopped, which concludes the management of that Peer instance.

Start peer

Once InitPeer(Peer) is invoked for every registered reactor, the p2p layer starts the peer’s communication routines and adds the Peer to the set of connected peers. If both steps are concluded without errors, the reactor’s AddPeer(Peer) is invoked:
start-peer      = [*receive] (connected-peer / start-error)
connected-peer  = add-peer *receive
In case of errors, a message is logged informing that the p2p layer failed to start the peer. This is not a common scenario and it is only expected to happen when interacting with a misbehaving or slow peer. A practical example is reported on this issue. It is up to the reactor to define how to process the AddPeer(Peer) event. The typical behavior is to start routines that, given some conditions or events, send messages to the added peer, using the provided Peer instance. The companion API for Reactors documents the methods provided by Peer instances, available from when they are added to the reactors.

Stop Peer

The p2p layer informs all registered reactors when it disconnects from a Peer, using the RemovePeer(Peer, reason) method:
stop-peer       = [peer-error] remove-peer
This method is invoked after the p2p layer has stopped peer’s send and receive routines. Depending of the reason for which the peer was stopped, different log messages can be produced. After removing a peer from all reactors, the Peer instance is also removed from the set of connected peers. This enables the same peer to reconnect and InitPeer(Peer) to be invoked for the new connection. From the removal of a Peer , the reactor should not receive any further message from the peer and must not try sending messages to the removed peer. This usually means stopping the routines that were started by the companion Add(Peer) method.

Receive messages

The main duty of a reactor is to handle incoming messages on the channels it has registered with the p2p layer. The pre-condition for receiving a message from a Peer is that the p2p layer has previously invoked InitPeer(Peer). This means that the reactor must be able to receive a message from a Peer before AddPeer(Peer) is invoked. This happens because the peer’s send and receive routines are started before, and should be already running when the p2p layer adds the peer to every registered reactor.
start-peer      = [*receive] (connected-peer / start-error)
connected-peer  = add-peer *receive
The most common scenario, however, is to start receiving messages from a peer after AddPeer(Peer) is invoked. An arbitrary number of messages can be received, until the peer is stopped and RemovePeer(Peer) is invoked. When a message is received from a connected peer on any of the channels registered by the reactor, the p2p layer will deliver the message to the reactor via the Receive(Envelope) method. The message is packed into an Envelope that contains:
  • ChannelID: the channel the message belongs to
  • Src: the source Peer handler, from which the message was received
  • Message: the actual message’s payload, unmarshalled using protocol buffers
Two important observations regarding the implementation of the Receive method:
  1. Concurrency: the implementation should consider concurrent invocations of the Receive method carrying messages from different peers, as the interaction with different peers is independent and messages can be received in parallel.
  2. Non-blocking: the implementation of the Receive method is expected not to block, as it is invoked directly by the receive routines. In other words, while Receive does not return, other messages from the same sender are not delivered to any reactor.