Connection 封装详解
Connection是 muduo 架构中负责管理一个已连接连接的核心组件。它管理一个已连接的 socket fd 及其对应事件(通过Channel),提供输入/输出缓冲区方便粘包处理,管理连接状态(连接中、已连接、半关闭、已关闭),提供回调机制让用户处理连接建立、消息到达、连接关闭等事件,支持非活跃连接的超时释放(通过定时器)和协议切换(Upgrade),并保证线程安全:所有关键操作都在所属的EventLoop线程中执行。
一、整体定位与设计思想
1. 功能概述
Connection 负责:
- 管理一个已连接的 socket fd,及其对应的事件(通过
Channel)。 - 提供输入/输出缓冲区,方便数据的粘包处理。
- 管理连接状态(连接中、已连接、半关闭、已关闭)。
- 提供回调机制,让用户处理连接建立、消息到达、连接关闭等事件。
- 支持非活跃连接的超时释放(通过定时器)。
- 支持协议切换(Upgrade),可在运行时改变回调函数和上下文。
- 保证线程安全:所有关键操作都在所属的
EventLoop线程中执行。
2. 与其他模块的关系
- 持有一个
Socket对象管理 fd。 - 持有一个
Channel对象,将 fd 的事件注册到EventLoop。 - 持有一个
EventLoop*指针,用于投递任务(RunInLoop、QueueInLoop)。 - 使用
Buffer作为输入/输出缓冲区。 - 使用
Any保存用户自定义上下文。 - 使用定时器(通过
EventLoop的TimerWheel)实现非活跃释放。
二、类声明与继承
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;用于类型别名。 - 状态枚举:定义了连接可能处于的四种状态。
PtrConnection:shared_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 粘包的关键。_context:Any类型,用于存储任意用户数据,例如协议状态。
回调函数成员
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();
}
- 当对端关闭连接(
EPOLLHUP或EPOLLRDHUP)时触发。 - 先处理输入缓冲区中剩余的数据(如果有),然后释放连接。
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 线程中执行)
这些函数通过 RunInLoop 或 QueueInLoop 被调度到正确的线程执行。
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。 - 从
EventLoop的Poller中移除该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));
}
- 保存参数,初始化
Socket和Channel对象。 - 设置
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 的公有方法,都通过 RunInLoop 或 QueueInLoop 投递到 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类有更精细的操作和边界控制。 - 支持设置高/低水位回调,避免缓冲区无限增长。
- 有
forceClose和shutdown区分。 - 定时器使用
TimerId管理,更安全。 - 使用
WeakCallback机制,在定时器回调中绑定weak_ptr,避免悬垂。 - 提供了
setTcpNoDelay、setKeepAlive等 socket 选项设置。 - 生命周期管理更严格,通过
Channel::tie绑定shared_ptr,防止回调期间对象被销毁。
- 使用
此处的 Connection 是简化版本,但体现了 muduo 的核心设计思想。
转载自 CSDN-专业IT技术社区



