文章目录
交易系统:UDP McastSocket 类 交易系统:UDP McastSocket 类 交易系统:UDP McastSocket 类
成员对象
struct McastSocket {
init socket_fd_ = -1;
// 应用层缓冲区
std::vector<char> outbound_data_;
size_t next_send_valid_index_ = 0;
std::vector<char> inbound_data_;
size_t next_rcv_valid_index_ = 0;
// 只有 MDC 需要用到这个回调函数吗?
std::function<void(McastSocket* s)> recv_callback_ = nullptr;
std::string time_str_;
Logger &logger;
}
在 McstSocket 当中的 logger 对象是引用类型的,而在一些其它对象当中的 logger 是非引用类型。McastSocket 借用所属组件的 logger,把日志写入该组件的日志文件。外部 logger 必须活得比 McastSocket 对象更久。
McastSocket::init()
负责创建并配置 multicast socket。我发现 UDP multicast 的 socket 的配置比 TCP socket 的配置简单很多,是不是?
auto McastSocket::init(const std::string &ip, const std::string &iface, int port, bool is_listening) {
const Common::SocketCfg socket_cfg{ip, iface, port, true, is_listening, false};
socket_fd_ = Common::createSocket(logger_, socket_cfg);
return socket_fd_;
}
join() 和 leave()
发送者向组播 IP 发送,加入该组的接收者接收数据。
McastSocket::join() 用于接收端加入组播组。而发送方是无需加入就可发送的。
由于 McastSocket::leave() 的实现实质上是关闭 socket,因此发送方和接收方都可以调用
bool McastSocket::join(const std::string& ip) {
return Common::join(socket_fd_, ip);
}
auto McastSocket::leave(const std::string&, int) {
// 这个方法的参数没写错,是合法的“省略参数名”语法,表示函数接收这些参数,但实现不使用它们
// 这里的 leave() 是直接关闭 socket,发送端也可调用。leave 的方法名其实不够准确
// TODO:重构代码,把 leave() 写成单纯从组播组中退出;为关闭 socket 写一个新方法
close(socket_fd_);
socket_fd_ = -1;
}
McastSocket::sendAndRecv()
负责把 socket 内核缓冲区当中当中的字节流复制到应用层接收缓冲区,并且把应用层发送缓冲区的字节流复制到 socket 内核缓冲区。
auto McastSocket::sendAndRecv() noexcept {
// 把内核缓冲区的信息拷贝到应用层缓冲区,使用 recv() 系统调用
const ssize_t n_rcv = recv(socket_fd_, inbound_data_.data() + next_rcv_valid_index_, McastBufferSize - next_rcv_valid_index_, MSG_DONTWAIT);
if (n_rcv > 0) {
next_rcv_valid_index_ += n_rcv;
logger_.log("%:% %() % read socket:% len:%\n", __FILE__, __LINE__, __FUNCTION__, Common::getCurrentTimeStr(&time_str_), socket_fd_, next_rcv_valid_index_);
// 让回调函数来处理应用层缓冲区的数据
recv_callback_(this);
}
// 把应用层缓冲区的字节移动到内核缓冲区
if (next_send_valid_index_ > 0) {
ssize_t 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 (n_rcv > 0);
}
}
这里使用的 recv() 也是系统调用。recvmsg() 与 recv() 的区别在于:
recv()直接传入缓冲区地址、长度和标志recvmsg()通过msghdr描述多个缓冲区,因此可以接收实践戳等辅助信息
在 OrderServer 当中,FIFOSequencer 依赖每个订单到达内核的时间,因此需要额外的缓冲区来存储时间戳这一辅助信息,因此 TCPSocket 使用 recvmsg() 系统调用。但是在 MarketDataPublisher 这边,推送市场更新的 McastSocket 不需要记录辅助信息,因此使用 recv() 系统调用即可。
McastSocket::send()
auto McastSocket::send(const void* data, size_t len) noexcept {
// const void* 表示指向未知类型数据的指针,不能通过它修改数据
// 因此 send() 可以接收不同结构体或字节数组的地址,配合 len 确定复制多少字节
// memcpy 把待发送数据追加到应用层缓冲区
// 参数:memcpy(目标地址, 源地址, 字节数); memcpy 不是系统调用,完全运行在用户态
memcpy(outbound_data_.data() + next_send_valid_index_, data, len);
next_send_valid_index_ += len;
ASSERT(next_send_valid_index_ < McastBufferSize, "Mcast socket buffer filled up and sendAndRecv() not called.");
}