hehelm头像
关注

仿muduo库实现高并发服务器—Connection

Connection 封装详解

Connection 是 muduo 架构中负责管理一个已连接连接的核心组件。它管理一个已连接的 socket fd 及其对应事件(通过 Channel),提供输入/输出缓冲区方便粘包处理,管理连接状态(连接中、已连接、半关闭、已关闭),提供回调机制让用户处理连接建立、消息到达、连接关闭等事件,支持非活跃连接的超时释放(通过定时器)和协议切换(Upgrade),并保证线程安全:所有关键操作都在所属的 EventLoop 线程中执行。

一、整体定位与设计思想

1. 功能概述

Connection 负责:

  • 管理一个已连接的 socket fd,及其对应的事件(通过 Channel)。
  • 提供输入/输出缓冲区,方便数据的粘包处理。
  • 管理连接状态(连接中、已连接、半关闭、已关闭)。
  • 提供回调机制,让用户处理连接建立、消息到达、连接关闭等事件。
  • 支持非活跃连接的超时释放(通过定时器)。
  • 支持协议切换(Upgrade),可在运行时改变回调函数和上下文。
  • 保证线程安全:所有关键操作都在所属的 EventLoop 线程中执行。

2. 与其他模块的关系

  • 持有一个 Socket 对象管理 fd。
  • 持有一个 Channel 对象,将 fd 的事件注册到 EventLoop
  • 持有一个 EventLoop* 指针,用于投递任务(RunInLoopQueueInLoop)。
  • 使用 Buffer 作为输入/输出缓冲区。
  • 使用 Any 保存用户自定义上下文。
  • 使用定时器(通过 EventLoopTimerWheel)实现非活跃释放。

二、类声明与继承

class Connection;
typedef enum { DISCONNECTED, CONNECTING, CONNECTED, DISCONNECTING } ConnStatu;
using PtrConnection = std::shared_ptr<Connection>;

class Connection : public std::enable_shared_from_this<Connection> {
  • 前向声明class Connection; 用于类型别名。
  • 状态枚举:定义了连接可能处于的四种状态。
  • PtrConnectionshared_ptr<Connection> 的别名,方便使用。
  • 继承 std::enable_shared_from_this<Connection>:这是关键点。它允许在类的成员函数中通过 shared_from_this() 获取指向自身的 shared_ptr。因为连接对象通常由 shared_ptr 管理,在回调中传递 shared_ptr 可以保证对象在回调执行期间不被释放。但前提是对象必须由 shared_ptr 管理,且不能用于栈对象。

三、私有成员变量

uint64_t _conn_id;              // 连接唯一ID
int _sockfd;                    // 关联的 socket fd
bool _enable_inactive_release;  // 是否启用非活跃释放
EventLoop *_loop;               // 所属的事件循环
ConnStatu _statu;               // 连接状态
Socket _socket;                 // Socket 管理对象
Channel _channel;               // 事件管理
Buffer _in_buffer;              // 输入缓冲区
Buffer _out_buffer;             // 输出缓冲区
Any _context;                   // 用户自定义上下文
  • _conn_id:唯一标识连接,通常由服务器生成,也用作定时器 ID。
  • _sockfd:保存 fd,方便访问。
  • _enable_inactive_release:控制是否启用超时释放,默认 false。
  • _loop:指向所属的 EventLoop,所有操作必须在该线程执行。
  • _statu:当前状态,在状态转换中起到重要作用。
  • _socket:封装了 fd,提供底层 socket 操作(如读写、关闭)。
  • _channel:负责事件注册和回调,在构造函数中设置回调函数。
  • _in_buffer / _out_buffer:输入输出缓冲区,是处理 TCP 粘包的关键。
  • _contextAny 类型,用于存储任意用户数据,例如协议状态。

回调函数成员

using ConnectedCallback = std::function<void(const PtrConnection&)>;
using MessageCallback = std::function<void(const PtrConnection&, Buffer *)>;
using ClosedCallback = std::function<void(const PtrConnection&)>;
using AnyEventCallback = std::function<void(const PtrConnection&)>;

ConnectedCallback _connected_callback;
MessageCallback _message_callback;
ClosedCallback _closed_callback;
AnyEventCallback _event_callback;
ClosedCallback _server_closed_callback;
  • 这些回调由用户(或服务器框架)设置。
  • _connected_callback:连接建立完成时调用。
  • _message_callback:接收到数据时调用,传入输入缓冲区。
  • _closed_callback:连接关闭时调用(用户层面)。
  • _event_callback:任意事件发生时调用(如读、写、关闭等),可用于刷新活跃度。
  • _server_closed_callback:由服务器框架设置,用于在连接关闭时从服务器的连接表中移除该连接。

四、Channel 事件回调处理方法(私有)

这些函数是 Channel 的回调,当 fd 发生相应事件时由 EventLoop 调用。

1. HandleRead

void HandleRead() {
    char buf[65536];
    ssize_t ret = _socket.NonBlockRecv(buf, 65535);
    if (ret < 0) {
        return ShutdownInLoop();
    }
    _in_buffer.WriteAndPush(buf, ret);
    if (_in_buffer.ReadAbleSize() > 0) {
        return _message_callback(shared_from_this(), &_in_buffer);
    }
}
  • 使用非阻塞读取,一次性尽量多读(最多 65535 字节)。
  • 如果读取出错(ret < 0),则调用 ShutdownInLoop() 进入关闭流程。
  • 成功读取后,将数据追加到输入缓冲区。
  • 如果缓冲区有可读数据,则调用消息回调,让用户处理。注意:传入的 Buffer* 是输入缓冲区的指针,用户在处理时可以直接从中读取数据,但不要长期保存该指针,因为缓冲区内容可能后续被修改或清空。
  • 使用 shared_from_this() 获取自身的 shared_ptr,保证回调期间对象存活。

注意ret == 0 表示没有数据可读(非阻塞下正常),但此函数没有显式处理,因为 ret == 0 时不会进入 if (ret < 0),继续执行 WriteAndPush(buf, 0) 实际不会写入数据,然后判断缓冲区是否有数据,若无则不调用回调。这是合理的。

2. HandleWrite

void HandleWrite() {
    ssize_t ret = _socket.NonBlockSend(_out_buffer.ReadPosition(), _out_buffer.ReadAbleSize());
    if (ret < 0) {
        if (_in_buffer.ReadAbleSize() > 0) {
            _message_callback(shared_from_this(), &_in_buffer);
        }
        return Release();
    }
    _out_buffer.MoveReadOffset(ret);
    if (_out_buffer.ReadAbleSize() == 0) {
        _channel.DisableWrite();
        if (_statu == DISCONNECTING) {
            return Release();
        }
    }
    return;
}
  • 发送输出缓冲区中待发送的数据。
  • 如果发送出错(ret < 0),先处理可能未处理的输入数据(调用消息回调),然后直接释放连接(Release())。
  • 发送成功,移动输出缓冲区的读偏移。
  • 如果输出缓冲区已空,关闭写事件监控(因为不再有数据要写)。如果此时连接状态是 DISCONNECTING(半关闭,等待发送完数据后关闭),则执行释放。

3. HandleClose

void HandleClose() {
    if (_in_buffer.ReadAbleSize() > 0) {
        _message_callback(shared_from_this(), &_in_buffer);
    }
    return Release();
}
  • 当对端关闭连接(EPOLLHUPEPOLLRDHUP)时触发。
  • 先处理输入缓冲区中剩余的数据(如果有),然后释放连接。

4. HandleError

void HandleError() {
    return HandleClose();
}
  • 出错时按关闭处理。

5. HandleEvent

void HandleEvent() {
    if (_enable_inactive_release == true) {
        _loop->TimerRefresh(_conn_id);
    }
    if (_event_callback) {
        _event_callback(shared_from_this());
    }
}
  • 这是任意事件触发时都会调用的函数(因为 Channel::HandleEvent 最后会调用 _event_callback)。
  • 如果启用了非活跃释放,则刷新定时器,推迟超时时间。
  • 然后调用用户设置的 _event_callback(如果有)。

五、内部操作函数(在 EventLoop 线程中执行)

这些函数通过 RunInLoopQueueInLoop 被调度到正确的线程执行。

1. EstablishedInLoop

void EstablishedInLoop() {
    assert(_statu == CONNECTING);
    _statu = CONNECTED;
    _channel.EnableRead();
    if (_connected_callback) _connected_callback(shared_from_this());
}
  • 在连接建立后调用(通常由服务器在接受新连接后调用 Established() 投递)。
  • 断言当前状态必须是 CONNECTING,确保逻辑正确。
  • 修改状态为 CONNECTED
  • 启动读事件监控,开始接收数据。
  • 调用连接建立回调,通知用户。

2. ReleaseInLoop

void ReleaseInLoop() {
    _statu = DISCONNECTED;
    _channel.Remove();
    _socket.Close();
    if (_loop->HasTimer(_conn_id)) CancelInactiveReleaseInLoop();
    if (_closed_callback) _closed_callback(shared_from_this());
    if (_server_closed_callback) _server_closed_callback(shared_from_this());
}
  • 真正的资源释放函数。
  • 将状态设为 DISCONNECTED
  • EventLoopPoller 中移除该 Channel
  • 关闭 socket。
  • 如果存在非活跃释放定时器,取消它。
  • 调用用户关闭回调。
  • 调用服务器内部关闭回调(从连接表中移除)。
  • 注意执行顺序:先调用用户回调,再调用服务器回调,以保证在服务器移除连接之前,用户还有机会处理。最后服务器移除连接后,shared_ptr 引用计数可能降为 0,对象析构。

3. SendInLoop

void SendInLoop(Buffer &buf) {
    if (_statu == DISCONNECTED) return;
    _out_buffer.WriteBufferAndPush(buf);
    if (_channel.WriteAble() == false) {
        _channel.EnableWrite();
    }
}
  • 将数据追加到输出缓冲区。
  • 如果当前没有启动写事件监控,则启动写事件。当 socket 可写时,HandleWrite 会被调用发送数据。

4. ShutdownInLoop

void ShutdownInLoop() {
    _statu = DISCONNECTING;
    if (_in_buffer.ReadAbleSize() > 0) {
        if (_message_callback) _message_callback(shared_from_this(), &_in_buffer);
    }
    if (_out_buffer.ReadAbleSize() > 0) {
        if (_channel.WriteAble() == false) {
            _channel.EnableWrite();
        }
    }
    if (_out_buffer.ReadAbleSize() == 0) {
        Release();
    }
}
  • 半关闭操作:不再接收新数据,但可能还有数据需要发送。
  • 状态设为 DISCONNECTING
  • 处理输入缓冲区剩余数据(调用消息回调)。
  • 如果输出缓冲区还有数据待发送,确保写事件监控开启,等待发送完成。
  • 如果输出缓冲区为空,立即释放。

5. EnableInactiveReleaseInLoop

void EnableInactiveReleaseInLoop(int sec) {
    _enable_inactive_release = true;
    if (_loop->HasTimer(_conn_id)) {
        return _loop->TimerRefresh(_conn_id);
    }
    _loop->TimerAdd(_conn_id, sec, std::bind(&Connection::Release, this));
}
  • 启用非活跃超时释放。
  • 如果定时器已存在,刷新它;否则添加新定时器,超时后调用 Release()
  • 定时器 ID 使用连接 ID(_conn_id),确保唯一。
  • 注意:Release() 是公有接口,会调用 QueueInLoop 投递 ReleaseInLoop。由于定时器超时可能发生在 EventLoop 线程中,这里使用 QueueInLoop 而非 RunInLoop,但 Release() 内部已经处理了线程调度,所以没问题。

6. CancelInactiveReleaseInLoop

void CancelInactiveReleaseInLoop() {
    _enable_inactive_release = false;
    if (_loop->HasTimer(_conn_id)) { 
        _loop->TimerCancel(_conn_id); 
    }
}
  • 取消非活跃释放,删除定时器。

7. UpgradeInLoop

void UpgradeInLoop(const Any &context, 
                   const ConnectedCallback &conn, 
                   const MessageCallback &msg, 
                   const ClosedCallback &closed, 
                   const AnyEventCallback &event) {
    _context = context;
    _connected_callback = conn;
    _message_callback = msg;
    _closed_callback = closed;
    _event_callback = event;
}
  • 用于协议切换:更新上下文和所有回调函数。
  • 该函数必须在 EventLoop 线程中执行,因此外部调用 Upgrade 时先断言在循环线程,然后通过 RunInLoop 投递(虽然 RunInLoop 在循环线程内会直接执行,但为了统一,仍使用投递)。

六、公有接口

构造函数

Connection(EventLoop *loop, uint64_t conn_id, int sockfd)
    : _conn_id(conn_id), _sockfd(sockfd),
      _enable_inactive_release(false), _loop(loop), _statu(CONNECTING),
      _socket(_sockfd), _channel(loop, _sockfd) {
    _channel.SetCloseCallback(std::bind(&Connection::HandleClose, this));
    _channel.SetEventCallback(std::bind(&Connection::HandleEvent, this));
    _channel.SetReadCallback(std::bind(&Connection::HandleRead, this));
    _channel.SetWriteCallback(std::bind(&Connection::HandleWrite, this));
    _channel.SetErrorCallback(std::bind(&Connection::HandleError, this));
}
  • 保存参数,初始化 SocketChannel 对象。
  • 设置 Channel 的各个事件回调,绑定到 Connection 的成员函数。
  • 注意:此时连接状态为 CONNECTING,还没有启动读监控,等待 Established() 调用后才正式进入 CONNECTED 状态并开始读写。

析构函数

~Connection() { DBG_LOG("RELEASE CONNECTION:%p", this); }
  • 仅打印日志,实际资源释放已在 ReleaseInLoop 中完成。

访问器

int Fd() { return _sockfd; }
int Id() { return _conn_id; }
bool Connected() { return (_statu == CONNECTED); }
void SetContext(const Any &context) { _context = context; }
Any *GetContext() { return &_context; }
  • 提供 fd、ID、状态查询以及上下文设置/获取。
  • 注意SetContext 未做线程安全处理,应确保在 EventLoop 线程或连接建立前调用。

设置回调

void SetConnectedCallback(const ConnectedCallback&cb) { _connected_callback = cb; }
void SetMessageCallback(const MessageCallback&cb) { _message_callback = cb; }
void SetClosedCallback(const ClosedCallback&cb) { _closed_callback = cb; }
void SetAnyEventCallback(const AnyEventCallback&cb) { _event_callback = cb; }
void SetSrvClosedCallback(const ClosedCallback&cb) { _server_closed_callback = cb; }
  • 这些设置通常也在连接建立前进行,或在 EventLoop 线程内通过 RunInLoop 设置,以保证线程安全。

核心操作

Established
void Established() {
    _loop->RunInLoop(std::bind(&Connection::EstablishedInLoop, this));
}
  • EstablishedInLoop 投递到 EventLoop 线程执行。由于可能从其他线程调用,使用 RunInLoop 保证在正确的线程执行。
Send
void Send(const char *data, size_t len) {
    Buffer buf;
    buf.WriteAndPush(data, len);
    _loop->RunInLoop(std::bind(&Connection::SendInLoop, this, std::move(buf)));
}
  • 用户发送数据接口。
  • 将数据先拷贝到临时 Buffer,然后通过 RunInLoop 投递到 EventLoop 线程。
  • 使用 std::move 将临时 Buffer 的所有权转移到任务中,避免额外拷贝。
  • 安全考虑:如果直接传指针,可能因任务延迟执行导致指针悬垂,因此这里先拷贝一份是必要的。
Shutdown
void Shutdown() {
    _loop->RunInLoop(std::bind(&Connection::ShutdownInLoop, this));
}
  • 优雅关闭接口:停止接收,但会发送完缓冲区内数据。
Release
void Release() {
    _loop->QueueInLoop(std::bind(&Connection::ReleaseInLoop, this));
}
  • 强制关闭接口:立即关闭连接,但内部仍通过 QueueInLoop 投递,保证在 EventLoop 线程中执行释放。
  • 使用 QueueInLoop 而非 RunInLoop,因为调用者可能已经在 EventLoop 线程中,但为了统一,且释放操作不需要立即执行,使用 QueueInLoop 可以避免在调用点阻塞。
EnableInactiveRelease / CancelInactiveRelease
void EnableInactiveRelease(int sec) {
    _loop->RunInLoop(std::bind(&Connection::EnableInactiveReleaseInLoop, this, sec));
}
void CancelInactiveRelease() {
    _loop->RunInLoop(std::bind(&Connection::CancelInactiveReleaseInLoop, this));
}
  • 启用/取消非活跃超时释放。
Upgrade
void Upgrade(const Any &context, const ConnectedCallback &conn, const MessageCallback &msg, 
             const ClosedCallback &closed, const AnyEventCallback &event) {
    _loop->AssertInLoop();
    _loop->RunInLoop(std::bind(&Connection::UpgradeInLoop, this, context, conn, msg, closed, event));
}
  • 协议切换接口。要求调用者必须在 EventLoop 线程中(使用 AssertInLoop 断言),然后投递执行实际的更新。

七、关键设计点分析

1. 线程安全与任务投递

所有可能修改连接状态或操作 socket 的公有方法,都通过 RunInLoopQueueInLoop 投递到 EventLoop 线程执行。这保证了在单线程内串行处理,避免数据竞争。

  • RunInLoop:如果当前已在 EventLoop 线程,则直接执行;否则投递并唤醒。
  • QueueInLoop:总是投递到任务队列,不立即执行(即使当前在循环线程,也会排队)。
  • 这种设计是 muduo 的核心思想:每个连接属于一个线程,所有操作都在该线程完成

2. 生命周期管理

  • 使用 shared_ptr 管理 Connection 对象,通过 shared_from_this() 在回调中传递自身,避免在回调期间对象被销毁。
  • 释放流程:Release() -> QueueInLoop(ReleaseInLoop) -> ReleaseInLoop 中移除事件、关闭 socket、调用回调、最后服务器回调会从管理容器中移除 shared_ptr。当所有 shared_ptr 释放后,对象析构。
  • 注意:在 ReleaseInLoop 中调用用户回调后,不应再访问成员变量(因为可能被释放),所以将服务器回调放在最后,让服务器有机会从容器中移除自身引用。

3. 缓冲区设计

  • _in_buffer 用于缓存读到的数据,防止粘包;用户可从中读取数据,处理时可能需要读取多次。
  • _out_buffer 用于缓存待发送的数据,当 socket 不可写时,数据暂存于此,等待可写事件触发后发送。
  • 使用 Buffer 类(未给出,但应具备读写偏移、可读大小等操作)来管理数据,简化了操作。

4. 非活跃释放机制

  • 通过定时器实现:当连接在一定时间内无任何事件(读、写等)时,自动关闭。
  • 每次事件触发(HandleEvent)会刷新定时器,延迟超时时间。
  • 定时器 ID 与连接 ID 相同,避免冲突。
  • 该功能可选,默认关闭。

5. 协议切换 (Upgrade)

  • 允许在连接建立后动态改变回调函数和上下文,常用于切换应用层协议(如从 HTTP 升级到 WebSocket)。
  • 要求必须在 EventLoop 线程中调用,以保证更新立即生效,避免新事件使用旧回调。

八、与 muduo 原版对比

  • muduo 的 TcpConnection 更完善
    • 使用 Buffer 类有更精细的操作和边界控制。
    • 支持设置高/低水位回调,避免缓冲区无限增长。
    • forceCloseshutdown 区分。
    • 定时器使用 TimerId 管理,更安全。
    • 使用 WeakCallback 机制,在定时器回调中绑定 weak_ptr,避免悬垂。
    • 提供了 setTcpNoDelaysetKeepAlive 等 socket 选项设置。
    • 生命周期管理更严格,通过 Channel::tie 绑定 shared_ptr,防止回调期间对象被销毁。

此处的 Connection 是简化版本,但体现了 muduo 的核心设计思想。

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/hehelm/article/details/163999956

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--