文章目录
交易系统:TCPSocket 类 交易系统:TCPSocket 类 交易系统:TCPSocket 类
这个类封装了最多的网络相关系统调用。
TCPSocket 的成员对象
struct TCPSocket {
int socket_fd_ = -1;
std::vector<char> outbound_data_; // 是 TCPSocket 自己维护的应用层发送缓冲区,不是 Linux 内核的 socket 缓冲区。
size_t next_send_valid_index_ = 0; // 表示缓冲区中待发送的有效字节数,也就是下一次追加数据的位置
std::vector<char> inbound_data_;
size_t next_rcv_valid_index_ = 0;
// sockaddr_in 类型结构体用来保存 IPv4 地址
struct sockaddr_in socket_attrib_{};
// 该 socket 事件的回调函数
std::function<void(TCPSocket* s, Nanos rx_time)> recv_callback_ = nullptr; // 这是什么玩意?这是 functor 吗?哦好像是,因为具体的接收到消息时调用的回调函数应该要注册,而不是硬编码写死了。我在想不同的 socket 的回调函数有何不同?服务器侧的 socket 的回调函数与客户端侧的 socket 回调函数有何不同?
// 每个 socket 只能注册一个回调函数是吧!
std::string time_str_;
Logger& logger_;
}
下面逐个逐个来解释:
std::vector<char> outbound_data_; // 应用层缓冲区
size_t next_send_valid_index_ = 0; // 缓冲区中待发送的字节数,即下一次追加数据的位置
需要把 socket 作为应用层概念与作为内核概念分开理解:作为应用层的 socket 就是我们在这里定义的 TCPSocket,负责管理应用层缓冲区并调用系统 API;而作为内核概念的 socket 负责管理 TCP 可靠传输(按字节流发送、分段、排序、确认、重传、流量/拥塞控制)。两个 socket 概念都有缓冲区:TCPSocket 管理 “应用层缓冲区”,内核 socket 管理内核缓冲区。
两种缓冲区的职责不同,outbound_data_ 这个应用层缓冲区相当于提供了一层封装:
- 程序需要调用
send()把其中的数据转移到 socket 的内核发送缓冲区。 send()是 Linux 系统调用,它只负责把用户态缓冲区中的字节交给该 socket 的内核发送缓冲区,然后就返回了。- 后续数据从内核发送缓冲区推进到网卡就不由
send()负责了。
TCPSocket::outbound_data_ 作为应用层缓冲区,则在内核缓冲区之上增加了封装:
- 把多个消息连续拼接
- 记录有效数据长度。
next_send_valid_index_表示缓冲区中待发送的有效字节数,也就是下一次追加数据的位置 - 在非阻塞发送时保留未发送部分
- 将业务层的
send(data, len)转成网络层可执行
std::vector<char> inbound_data_;
size_t next_rcv_valid_index_ = 0;
inbound_data_ 是 TCPSocket 管理的应用层缓冲区。仅当 inbound_data 上面有数据的时候(而不是 outbound_data_ 有数据的时候)调用 recv_callback_ 让注册的回调函数处理请求。对于 OrderServer 的连接 socket,就是使用 OrderServer::recvCallback() 来处理。
std::function<void(TCPSocket *s, Nanos rx_time)> recv_callback_ = nullptr;
在当前的实现下,每个 socket 只有一个回调函数。但是实际上一个 socket 可以有多个回调。允许不同的回调用于不同的事件。总结一下实现多种回调的方式:
TCPServer根据 epoll 事件分支,直接调用内部的onReadable(),onWritable()等方法TCPSocket持有多个回调,也就是指向不同回调函数的函数指针TCPSocket持有一个统一回调,回调内部再写不同的回调分支
struct sockaddr_in socket_attrib_{};
sockaddr_in 是 Linux 头文件 <netinet/in.h> 里定义的结构体,用来表示 IPv4 地址和端口信息。
这个对象具体用于保存什么地址和端口号,要视乎当前 socket 的定位而异:
- 当前 socket 是监听 socket 的时候:
sock_attrib_保存服务器本地 IP 和端口号,用于 bind 和 listen - 当前 socket 是连接 socket 的时候:
sock_attrib_保存对端(客户端)的 IP 和端口号 - 当前 socket 是客户端 socket 的时候:
sock_attrib_保存客户端自己的 IP 和端口号
TCPSocket 的构造函数
buffer 在这里初始化。在之后,buffer 都不会有 clear 的操作了,只会被覆写
TCPSocket::TCPSocket() {
inbound_data_.resize(TCPBufferSize);
outbound_data_.resize(TCPBufferSize);
}
TCPSocket::connect()
这个方法只有监听 socket 和客户端 socket 需要调用。
- 对于监听 socket:该方法创建并配置本地监听端点,使其绑定接口和端口并进入监听模式
- 对于客户端 socket:该方法创建并配置客户端端点,使其连接指定的远端 IP 和端口
auto TCPSocket::connect(const std::string& ip, const std::string& iface, int port, bool is_listening) {
// 配置:is_udp = false; is_listening 表示是监听 socket;true 表示开启时间戳
const SocketCfg socket_cfg{ip, iface, port, false, is_listening, true};
// createSocket() 完成底层工作,例如创建 socket, 设置非阻塞模式等
socket_fd_ = createSocket(logger_, socket_cfg);
// 填充当前 socket 对象的 socket_attrib_ 结构体。该结构体存储 IPv4 IP 和端口信息
socket_attrib_.sin_addr.s_addr = INADDR_ANY;
socket_attrib_.sin_port = htons(port);
socket_attrib_.sin_family = AF_INET;
return socket_fd_;
}
先看参数:
auto TCPSocket::connect(
const std::string& ip, // 目标 IP:客户端模式下为要交易所IP;监听模式下为空
const std::string& iface, // 本机网络接口名称:指定 socket 使用哪块网卡
int port, // 端口:客户端模式下为远端服务端口;监听模式下为本地监听端口
bool is_listening) // 是否为监听 socket
要理解 TCPSocket::connect() 被调用的链条,需要把监听 socket 与另外两种 socket 区分开来
当前 socket 是监听 socket 的时候:
OrderServer::start()- 进而调用
TCPServer::listen(iface, port) - 进而调用
listener_socket_.connect("", iface, port, true) - 进而调用
createSocket() - 进而调用
::socket()、::bind()和::listen()
学到这里,我发现本书的监听 socket、连接 socket 与客户端 socket 共享一个 TCPSocket 的定义,导致在许多方法内经常要区分当前 socket 是什么 socket。我觉得更合理的方式是分为两种 socket 定义:ListeningSocket 和 ConnectedSocket,前者用于监听 socket,后者用于连接 socket 和客户端 socket。
// TODO:重构本书的 socket 实现,分为两种 socket 定义。
TCPSocket::sendAndRecv()
// TODO:原书使用了 sendAndRecv() 方法把 send 和 recv 的职责放在同一个函数。也许为了清晰,更适合把这两个职责拆开
这个函数先处理 recv 的职责,把该 socket 在内核缓冲区内的数据,读入应用层缓冲区。然后处理 send 的职责。调用链条是:OrderServer::run() 依次调用 TCPServer::poll() 和 TCPServer::sendAndRecv() 进而调用 TCPSocket::sendAndRecv()。效果就是每当 OrderServer 的 poll 唤醒后,对每一个连接和监听 socket 调用 TCPSocket::sendAndRecv。
auto TCPSocket::sendAndRecv() noexcept {
char ctrl[CMSG_SPACE(sizeof(struct timeval))];
auto cmsg = reinterpret_cast<struct cmsgdr*>(&ctrl);
iovec iov{inbound_data_.data() + next_rcv_valid_index_, TCPBufferSize - next_rcv_valid_index_};
msghdr msg{&socket_attrib_, sizeof(socket_attrib_), &iov, 1, ctrl, sizeof(ctrl), 0};
const auto read_size = recvmsg(socket_fd, &msg, MSG_DONTWAIT);
if (read_size > 0) {
next_rcv_valid_index_ += read_size;
Nanos kernel_time = 0;
timeval time_kernel;
if (cmsg->cmsg_level == SOL_SOCKET && cmsg->cmsg_type == SCM_TIMESTAMP && cmsg->cmsg_len == CMSG_LEN(sizeof(time_kernel))) { // 这个判断意义何在
memcpy(&time_kernel, CMSG_DATA(cmsg), sizeof(time_kernel));
kernel_time = time_kernel.tv_sec * NANO_TO_SECS_ + time_kernel.tv_usec * NANOS_TO_MICROS;
}
const auto user_time = getCurrentNanos();
logger_.log("%:% %() % read socket:% len:% utime:% ktime:% diff:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket_fd_, next_rcv_valid_index_, user_time, kernel_time, (user_time - kernel_time));
// 终于调用回调函数了。
// 上层(比如 OrderServer 或 MDP)为 TCPSocket 对象注册回调函数
// 但是仍然由 TCPSocket 对象在 sendAndRecv() 函数中调用回调函数
recv_callback_(this, kernel_time);
}
// 这里突然就从接收信息转为发送信息了
if (next_send_valid_index_ > 0) {
// 这是非阻塞的
const auto n = ::send(socket_fd_, outbound_data_.data(), next_send_valid_index_, MSG_DONTWAIT | MSG_NOSIGNAL);
}
logger_.log("%:% %() % send socket:% len:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket_fd_, n);
next_send_valid_index_ = 0;
return (read_size > 0);
}
char ctrl[CMSG_SPACE(sizeof(struct timeval))];
ctrl 是给 recvmsg() 准备的 辅助 数据缓冲区,用于接收 socket 附带的控制信息。在这个函数当中,就是接收时间戳信息的缓冲区。
auto cmsg = reinterpret_cast<struct cmsghdr *>(&ctrl);
cmsg 是 control message 的简写。这里把 ctrl 这个辅助缓冲区解释成 cmsghdr * 类型。因为 ctrl 只是原始字节数组,而随后调用的 recvmsg() 这个系统调用写入的辅助数据有标准格式 cmsghdr,其实是 control message header。
struct cmsghdr {
size_t cmsg_len;
int cmsg_level;
int cmsg_type;
}
把 ctrl 解释成 cmsghdr* 才能按字段读取 cmsg_len、cmsg_level、cmsg_type。
iovec iov{inbound_data_.data() + next_rcv_valid_index_, TCPBufferSize - next_rcv_valid_index_};
msghdr msg{&socket_attrib_, sizeof(socket_attrib_), &iov, 1, ctrl, sizeof(ctrl), 0};
把 TCPSocket 类的应用层缓冲区 inbound_data_ 转成 准备接受缓冲区 iov。注意 iov 是写入消息正文的,而 ctrl 是写入时间戳这样的辅助信息的。
iovec 是 I/O vector 的简写。该类型是 Linux 定义的结构体类型,用来描述一段内存缓冲区。因此上述创建 iov 时的两个参数表示把新收到的字节写入 inbound_data_ 尾部未被使用的空间
struct iovec {
void *iov_base;
size_t iov_len;
};
然后
msghdr msg{&socket_attrib_, sizeof(socket_attrib_), &iov, 1, ctrl, sizeof(ctrl), 0};
这一步使用辅助数据缓冲区 ctrl 与正文缓冲区描述符 iov 构造 msg 对象/之后就可以使用 msg
msghdr msg{
&socket_attrib_, // msg_name:发送方地址写入位置
sizeof(socket_attrib_), // msg_namelen:地址结构大小
&iov, // 数据缓冲区数组
1, // 缓冲区数量
ctrl, // 辅助数据缓冲区数组
sizeof(ctrl), // 辅助数据缓冲区大小
0 // 由内核返回的附加标志
};
const auto read_size = recvmsg(socket_fd_, &msg, MSG_DONTWAIT);
if (read_size > 0) {
// 记录收到的字节数
next_rcv_valid_index_ += read_size;
Nanos kernel_time = 0;
timeval time_kernel;
// 检查辅助数据是否包含内核时间戳
if (cmsg->cmsg_level == SOL_SOCKET &&
cmsg->cmsg_type == SCM_TIMESTAMP &&
cmsg->cmsg_len == CMSG_LEN(sizeof(time_kernel))) {
// 提取并转换时间戳
memcpy(&time_kernel, CMSG_DATA(cmsg), sizeof(time_kernel));
kernel_time = time_kernel.tv_sec * NANOS_TO_SECS + time_kernel.tv_usec * NANOS_TO_MICROS; // convert timestamp to nanoseconds.
}
这里终于读入消息了。recvmsg() 也是系统调用。recvmsg() 接收 &msg 后,根据这些描述信息,把内核 socket 缓冲区当中的正文写入 iov 指向的 inbound_data_,把辅助数据写入 ctrl。read_size > 0 表示读到了数据;== 0 表示对端关闭连接;< 0 表示没有数据或发生错误。
ctrl 保存的信息当中,只有内核时间戳被读取: memcpy(&time_kernel, CMSG_DATA(cmsg), sizeof(time_kernel))。而 ctrl 是 sendAndRecv() 的局部遍历,因此其它数据被丢弃了。
recv_callback_(this, kernel_time);
调用回调函数来处理读入到应用层缓冲区 TCPServer::inbound_data_ 的消息。注意这里在调用这个回调函数的时候并没有传入信息,而是传入 this 指针。后续在回调函数内部通过该指针访问 inbound_data_ 缓冲区,然后按照协议大小完整解析数据(对于 OrderServer 来说,就是按照 OMClientRequest 的大小来解析数据。)
if (next_send_valid_index_ > 0) {
const auto n = ::send(socket_fd_, outbound_data_.data(), next_send_valid_index_, MSG_DONTWAIT | MSG_NOSIGNAL);
logger_.log("%:% %() % send socket:% len:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket_fd_, n);
next_send_valid_index_ = 0;
return (read_size > 0);
}
然后轮到发送了。发送其实像一个顺手做的事情,也就是在接收内核接收缓冲区数据的同时,顺手把应用层发送缓冲区的数据 ::send() 到内核发送缓冲区。最后的返回值是本轮接收数据是否成功,也可也看出本方法的重点是接收。
TCPSocket::send()
这个 send() 与 ::send() 系统调用不同,是吧新信息加入 TCPSocket::outbound_data_ 应用层缓冲区当中,而不是把应用层缓冲区的数据转移到内核缓冲区。
auto TCPSocket::send(const void* data, size_t len) noexcept {
memcpy(outbound_data_.data() + next_send_valid_index_, data, len);
next_send_valid_index_ += len;
}