剖析 muduo 网络库 TcpConnection 类的实现机制
在 muduo 网络库中,TcpConnection 是最核心的类之一。它封装了一个已建立的 TCP 连接,负责处理数据的收发、连接状态的管理以及事件回调的执行。它是 TcpServer 与客户端通信的桥梁,内部维护了输入输出缓冲区(Buffer),以应对非阻塞网络 I/O 中的数据积压问题。
TcpConnection 的核心结构
TcpConnection 继承自 std::enable_shared_from_this,这是为了确保在回调函数执行期间,对象不会被意外销毁。以下是其核心成员和接口的精简定义:
class TcpConnection : noncopyable, public std::enable_shared_from_this<TcpConnection> {
public:
TcpConnection(EventLoop* loop, const string& name, int fd, const InetAddress& local, const InetAddress& peer);
void send(const string& data);
void shutdown();
void forceClose();
void setConnectionCallback(const ConnectionCallback& cb) { connectionCb_ = cb; }
void setMessageCallback(const MessageCallback& cb) { messageCb_ = cb; }
void setWriteCompleteCallback(const WriteCompleteCallback& cb) { writeCompleteCb_ = cb; }
void connectEstablished();
void connectDestroyed();
private:
enum StateE { kDisconnected, kConnecting, kConnected, kDisconnecting };
void handleRead(Timestamp receiveTime);
void handleWrite();
void handleClose();
void handleError();
void sendInLoop(const void* data, size_t len);
void shutdownInLoop();
EventLoop* loop_;
StateE state_;
std::unique_ptr<Socket> socket_;
std::unique_ptr<Channel> channel_;
const InetAddress localAddr_;
const InetAddress peerAddr_;
Buffer inputBuf_; // 接收缓冲区
Buffer outputBuf_; // 发送缓冲区
// ... 回调函数及其他成员
};
初始化过程
在构造函数中,TcpConnection 会将其绑定的文件描述符(fd)注册到 Channel 中,并设置读、写、关闭及错误处理的回调。这些回调最终会由 Poller 触发并在 EventLoop 中执行。
TcpConnection::TcpConnection(EventLoop* loop, ...)
: loop_(loop),
state_(kConnecting),
socket_(new Socket(sockfd)),
channel_(new Channel(loop, sockfd)) {
channel_->setReadCallback(std::bind(&TcpConnection::handleRead, this, _1));
channel_->setWriteCallback(std::bind(&TcpConnection::handleWrite, this));
channel_->setCloseCallback(std::bind(&TcpConnection::handleClose, this));
channel_->setErrorCallback(std::bind(&TcpConnection::handleError, this));
socket_->setKeepAlive(true);
}
数据发送逻辑
send 接口是线程安全的。如果调用者不在当前 I/O 线程,它会将发送任务转发到连接所属的 EventLoop 线程中执行。核心逻辑由 sendInLoop 实现:
void TcpConnection::sendInLoop(const void* data, size_t len) {
loop_->assertInLoopThread();
ssize_t bytes_sent = 0;
size_t remaining = len;
bool error_occurred = false;
if (state_ == kDisconnected) return;
// 如果当前没有正在写入,且输出缓冲区为空,尝试直接写内核内核空间
if (!channel_->isWriting() && outputBuf_.readableBytes() == 0) {
bytes_sent = sockets::write(channel_->fd(), data, len);
if (bytes_sent >= 0) {
remaining = len - bytes_sent;
if (remaining == 0 && writeCompleteCb_) {
loop_->queueInLoop(std::bind(writeCompleteCb_, shared_from_this()));
}
} else {
bytes_sent = 0;
if (errno != EWOULDBLOCK) {
if (errno == EPIPE || errno == ECONNRESET) error_occurred = true;
}
}
}
// 若数据未完全发完,则存入应用层缓冲区,并注册可写事件
if (!error_occurred && remaining > 0) {
size_t current_buffer_size = outputBuf_.readableBytes();
// 高水位回调触发逻辑(防止内存积压过多)
if (current_buffer_size + remaining >= highWaterMark_ && current_buffer_size < highWaterMark_) {
loop_->queueInLoop(std::bind(highWaterMarkCb_, shared_from_this(), current_buffer_size + remaining));
}
outputBuf_.append(static_cast<const char*>(data) + bytes_sent, remaining);
if (!channel_->isWriting()) {
channel_->enableWriting();
}
}
}
事件处理:读取与写入
当底层 Poller 通知 fd 可读或可写时,对应的 handleRead 或 handleWrite 会被调用。
handleRead: 从内核读取数据到 inputBuf_。如果 read 返回 0,说明对端已关闭连接,此时调用 handleClose。如果有数据读取成功,则触发 messageCb_ 供用户处理。
handleWrite: 当内核缓冲区有空间时,将 outputBuf_ 中的内容继续写入。一旦 outputBuf_ 为空,必须立即停止监听可写事件,否则会造成 Busy Loop(因为非阻塞模式下,内核缓冲区通常是可写的)。
void TcpConnection::handleWrite() {
loop_->assertInLoopThread();
if (channel_->isWriting()) {
ssize_t n = sockets::write(channel_->fd(), outputBuf_.peek(), outputBuf_.readableBytes());
if (n > 0) {
outputBuf_.retrieve(n);
if (outputBuf_.readableBytes() == 0) {
channel_->disableWriting(); // 停止关注可写事件
if (writeCompleteCb_) {
loop_->queueInLoop(std::bind(writeCompleteCb_, shared_from_this()));
}
if (state_ == kDisconnecting) shutdownInLoop();
}
}
}
}
连接的关闭
muduo 倾向于优雅地关闭连接(Half-close)。shutdown 调用会关闭 TCP 的写端(shutdown(fd, SHUT_WR)),但此时连接仍可读取数据,直到对端也关闭连接并触发读 0 事件。
void TcpConnection::shutdownInLoop() {
loop_->assertInLoopThread();
if (!channel_->isWriting()) {
socket_->shutdownWrite();
}
}
这种设计确保了 outputBuf_ 中的剩余数据能被安全发送,同时不丢失对端可能正在传输的数据。