MConnection

MConnection 是一种多路复用连接,它在单个 TCP 连接之上支持多个彼此独立的流,并为这些流提供不同的服务质量保证。 每个流都称为一个 Channel,并且每个 Channel 都有一个全局唯一的 byte id。 每个 Channel 还具有一个相对优先级,用于决定该 Channel 相比其他 Channel 的服务质量。 每个 Channel 的 byte id 和相对优先级都会在连接初始化时完成配置。 MConnection 支持三种数据包类型:
  • Ping
  • Pong
  • Msg

Ping 和 Pong

ping 和 pong 消息分别通过向连接写入单个字节来表示,即 0x1 和 0x2。 当我们在 pingTimeout 时间内没有在某个 MConnection 上收到任何消息时,就会发送一条 ping 消息。 当 MConnection 收到 ping 时,只有在没有其他消息需要发送,且对端没有向我们发送过多 ping(TODO)的情况下,才会回复 pong。 如果在发送 ping 后,经过足够时间仍未收到 pong 或消息,则会断开该对端连接。

Msg

各个通道中的消息会被拆分为更小的 msgPacket 以便进行多路复用。
type msgPacket struct {
 ChannelID byte
 EOF       byte // 1 means message ends here.
 Bytes     []byte
}
msgPacket 使用 Proto3 进行序列化。 接收到的一组连续数据包中的 Bytes 会被依次追加拼接, 直到收到一个 EOF=1 的数据包;此时会返回完整的序列化消息, 并交由对应通道的 onReceive 函数处理。

多路复用

消息由单个 sendRoutine 发送。该例程会循环执行一个 select 语句,并最终发送 ping、pong 或一批数据消息。这一批数据消息可能包含来自多个通道的消息。 待发送的消息字节会排入各自所属通道的队列中,每个通道同一时间只保存一条尚未发送的消息。 每次为一个批次选择消息时,都会从“最近已发送字节数与通道优先级之比”最低的通道中取出一条消息。

发送消息

发送消息有两种方法:
func (m MConnection) Send(chID byte, msg interface{}) bool {}
func (m MConnection) TrySend(chID byte, msg interface{}) bool {}
Send(chID, msg) 是一个阻塞调用,会一直等待,直到 msg 被成功加入给定 id 字节 chID 对应通道的队列中。消息 msg 会使用 protobuf marshal 进行序列化。 TrySend(chID, msg) 是一个非阻塞调用;如果队列未满,就会将消息 msg 加入给定 id 字节 chID 对应的通道中,否则会立即返回 false。 每个 Peer 也会暴露 Send() 和 TrySend()。

Peer

每个对端都有一个 MConnection 实例,并且还包含其他信息,例如该连接是否为出站连接、连接关闭后是否应重新创建、节点的各类身份信息,以及供 reactors 使用的其他更高层级线程安全数据。

Switch/Reactor

Switch 负责处理对端连接,并暴露一个 API,用于在 Reactors 上接收传入消息。 每个 Reactor 负责处理一个或多个 Channels 的传入消息。 因此,尽管发送出站消息通常是在 peer 上执行,传入消息则是在 reactor 上接收。
// Declare a MyReactor reactor that handles messages on MyChannelID.
type MyReactor struct{}

func (reactor MyReactor) GetChannels() []*ChannelDescriptor {
    return []*ChannelDescriptor{ChannelDescriptor{ID:MyChannelID, Priority: 1}}
}

func (reactor MyReactor) Receive(chID byte, peer *Peer, msgBytes []byte) {
    r, n, err := bytes.NewBuffer(msgBytes), new(int64), new(error)
    msgString := ReadString(r, n, err)
    fmt.Println(msgString)
}

// Other Reactor methods omitted for brevity
...

switch := NewSwitch([]Reactor{MyReactor{}})

...

// Send a random message to all outbound connections
for _, peer := range switch.Peers().List() {
    if peer.IsOutbound() {
        peer.Send(MyChannelID, "Here's a random message")
    }
}

MConnection

MConnection is a multiplex connection that supports multiple independent streams with distinct quality of service guarantees atop a single TCP connection. Each stream is known as a Channel and each Channel has a globally unique byte id. Each Channel also has a relative priority that determines the quality of service of the Channel compared to other Channels. The byte id and the relative priorities of each Channel are configured upon initialization of the connection. The MConnection supports three packet types:
  • Ping
  • Pong
  • Msg

Ping and Pong

The ping and pong messages consist of writing a single byte to the connection; 0x1 and 0x2, respectively. When we haven’t received any messages on an MConnection in time pingTimeout, we send a ping message. When a ping is received on the MConnection, a pong is sent in response only if there are no other messages to send and the peer has not sent us too many pings (TODO). If a pong or message is not received in sufficient time after a ping, the peer is disconnected from.

Msg

Messages in channels are chopped into smaller msgPackets for multiplexing.
type msgPacket struct {
 ChannelID byte
 EOF       byte // 1 means message ends here.
 Bytes     []byte
}
The msgPacket is serialized using Proto3. The received Bytes of a sequential set of packets are appended together until a packet with EOF=1 is received, then the complete serialized message is returned for processing by the onReceive function of the corresponding channel.

Multiplexing

Messages are sent from a single sendRoutine, which loops over a select statement and results in the sending of a ping, a pong, or a batch of data messages. The batch of data messages may include messages from multiple channels. Message bytes are queued for sending in their respective channel, with each channel holding one unsent message at a time. Messages are chosen for a batch one at a time from the channel with the lowest ratio of recently sent bytes to channel priority.

Sending Messages

There are two methods for sending messages:
func (m MConnection) Send(chID byte, msg interface{}) bool {}
func (m MConnection) TrySend(chID byte, msg interface{}) bool {}
Send(chID, msg) is a blocking call that waits until msg is successfully queued for the channel with the given id byte chID. The message msg is serialized using protobuf marshalling. TrySend(chID, msg) is a nonblocking call that queues the message msg in the channel with the given id byte chID if the queue is not full; otherwise it returns false immediately. Send() and TrySend() are also exposed for each Peer.

Peer

Each peer has one MConnection instance, and includes other information such as whether the connection was outbound, whether the connection should be recreated if it closes, various identity information about the node, and other higher level thread-safe data used by the reactors.

Switch/Reactor

The Switch handles peer connections and exposes an API to receive incoming messages on Reactors. Each Reactor is responsible for handling incoming messages of one or more Channels. So while sending outgoing messages is typically performed on the peer, incoming messages are received on the reactor.
// Declare a MyReactor reactor that handles messages on MyChannelID.
type MyReactor struct{}

func (reactor MyReactor) GetChannels() []*ChannelDescriptor {
    return []*ChannelDescriptor{ChannelDescriptor{ID:MyChannelID, Priority: 1}}
}

func (reactor MyReactor) Receive(chID byte, peer *Peer, msgBytes []byte) {
    r, n, err := bytes.NewBuffer(msgBytes), new(int64), new(error)
    msgString := ReadString(r, n, err)
    fmt.Println(msgString)
}

// Other Reactor methods omitted for brevity
...

switch := NewSwitch([]Reactor{MyReactor{}})

...

// Send a random message to all outbound connections
for _, peer := range switch.Peers().List() {
    if peer.IsOutbound() {
        peer.Send(MyChannelID, "Here's a random message")
    }
}