打开导航 打开导航 Open menu
量化交易系统开发(C++)/ 期权策略与风控 / 行业研究分析
量化交易系统開發(C++)/ 期权策略与风控 / 行业研究分析
C++ trading systems / options strategy & risk controls / equity research
adrian@adrianxv.cn
文章目录

交易系统: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_ 作为应用层缓冲区,则在内核缓冲区之上增加了封装:

  1. 把多个消息连续拼接
  2. 记录有效数据长度。next_send_valid_index_ 表示缓冲区中待发送的有效字节数,也就是下一次追加数据的位置
  3. 在非阻塞发送时保留未发送部分
  4. 将业务层的 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 可以有多个回调。允许不同的回调用于不同的事件。总结一下实现多种回调的方式:

  1. TCPServer 根据 epoll 事件分支,直接调用内部的 onReadable(), onWritable() 等方法
  2. TCPSocket 持有多个回调,也就是指向不同回调函数的函数指针
  3. TCPSocket 持有一个统一回调,回调内部再写不同的回调分支
struct sockaddr_in socket_attrib_{};

sockaddr_in 是 Linux 头文件 <netinet/in.h> 里定义的结构体,用来表示 IPv4 地址和端口信息。

这个对象具体用于保存什么地址和端口号,要视乎当前 socket 的定位而异:

  1. 当前 socket 是监听 socket 的时候:sock_attrib_ 保存服务器本地 IP 和端口号,用于 bind 和 listen
  2. 当前 socket 是连接 socket 的时候:sock_attrib_ 保存对端(客户端)的 IP 和端口号
  3. 当前 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 的时候:

  1. OrderServer::start()
  2. 进而调用 TCPServer::listen(iface, port)
  3. 进而调用 listener_socket_.connect("", iface, port, true)
  4. 进而调用 createSocket()
  5. 进而调用 ::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_lencmsg_levelcmsg_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))。而 ctrlsendAndRecv() 的局部遍历,因此其它数据被丢弃了。

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;
}