所谓的 eventfd 其实就是一种事件通知的机制,就像 socket 有 socketfd,文件有文件的 fd 一样,eventfd 就是文件描述符的类型之一,专门用于事件通知的!
其本质就是在内核维护一个计数器,当我们创建了一个 eventfd 的时候就相当于创建了一个内核计数器。
每当我们使用 write 等接口向 eventfd 中写入一个数值的时候,计数器就会累加该数值,该数值表示事件通知的次数,即一共有多少个事件触发了!而当我们使用 read 等接口读取 eventfd 中的数值的时候,读取到的数值就是事件通知或者触发的次数,并且读取完之后,内核会自动将计数器清零!
而创建 eventfd 的接口如下所示:
#include <sys/eventfd.h>
int eventfd(unsigned int initval, int flags);- 功能:
- 创建一个文件描述符用于事件通知。
- 参数:
initval:计数器的初始值flags:- 对该文件描述符的属性标志,比如下面:
EFD_CLOEXEC:表示禁止进程复制EFD_NONBLOCK:启动非阻塞等待
- 对该文件描述符的属性标志,比如下面:
- 返回值:
- 成功返回一个文件描述符,失败返回
-1,并且设置错误码。
- 成功返回一个文件描述符,失败返回
需要注意的是,在对该 eventfd 进行读取或者写入的时候,数据大小必须是一个 8 字节数据的大小,即 uint64_t 类型!
并且一般 eventfd 是用来监听可读事件的,因为 eventfd 是一直可写的,所以一直都是有可写事件的,所以 eventfd 监听可写事件是没有意义的,就和文件描述符一样(文件描述符是一直可写和可读的,所以监听可读可写事件都没意义!这个也是为什么类似 ext4 这种文件不实现 poll 接口的原因!)
下面给出一个简单的使用样例:
#include <iostream>
#include <sys/eventfd.h>
#include <unistd.h>
using namespace std;
int main()
{
// 创建一个eventfd
int efd = eventfd(0, EFD_CLOEXEC | EFD_NONBLOCK);
if(efd < 0)
{
perror("eventfd error");
return -1;
}
// 对eventfd进行写两次,每次写1,所以此时eventfd中为2
uint64_t val = 1;
write(efd, &val, sizeof(uint64_t));
write(efd, &val, sizeof(uint64_t));
// 此时再对eventfd进行读
uint64_t ret = 0;
read(efd, &ret, sizeof(uint64_t));
cout << ret << endl;
close(efd);
return 0;
}
// 执行结果
[liren@VM-8-7-centos example]$ g++ -o eventfd eventfd.cpp
[liren@VM-8-7-centos example]$ ./eventfd
2Ⅱ. EventLoop描述符事件监控总模块的设计
该模块是对 Poller 模块、TimerQueue 模块、Socket 模块的一个整体封装,实现对所有描述符的事件监控。所以 EventLoop 模块必然是一个对象对应一个线程,在该线程执行函数中唯一要做的就是运行对应 EventLoop 的启动函数。
此外,EventLoop 模块为了保证整个服务器的线程安全问题,因此要求使用者对于 Connection 模块的所有操作一定要在其对应的 EventLoop 线程内完成,即每一个 Connection 连接对象都会绑定到一个 EventLoop 线程上,不能在其他线程中进行,避免线程安全问题!
比如组件使用者使用 Connection 模块发送数据,以及关闭连接等,这种操作都涉及到了线程安全问题,因为有可能这个描述符在多个线程中都触发了事件,需要多个线程都去处理的话,就会出现线程安全问题。所以要统一放在一个 EventLoop 线程内完成。所以说对于连接的所有操作,都需要放到 EventLoop 线程中执行!
至于如何保证对于一个连接操作都是在同一个线程中执行的话,就需要在每个 EventLoop 线程中配置一个任务队列,然后将该连接操作都包装一下,添加到其对应的 EventLoop 线程中的任务队列,而 EventLoop 线程要做的就是将任务队列中任务拿出来处理!
这里就引入了另一个问题,就是有可能 EventLoop 在监听事件就绪的时候会阻塞,导致其没办法去取出任务队列中的任务执行,导致效率低下,因此需要一个能够唤醒事件监控阻塞的事件通知机制,也就是上面提及到的 eventfd!
并且 EventLoop 模块保证自己内部所监控的所有描述符都必须是活跃连接,而非活跃连接就要及时释放避免资源浪费,所以就需要一个时间轮定时器 TimerQueue 来辅助管理释放这些非活跃连接!

下面是该模块内部需要包含成员变量如下:
-
⼀个
Poller对象:用于进行描述符的IO事件监控操作。 -
⼀个
TimerQueue对象:用于进行定时任务的管理,实现对非活跃连接的释放、刷新等操作。 -
⼀个
eventfd文件描述符:用于唤醒事件通知阻塞的事件通知机制,防止当前EventLoop线程因为等待IO事件就绪而阻塞,导致任务队列中的任务无法得到执行的问题! -
⼀个
PendingTask任务队列:使用者对Connection模块进行的所有操作,都要加入到任务队列中,由EventLoop模块进行管理,并在EventLoop模块对应的线程中执行。这样子做的原因是一个EventLoop对象中是有多个Connection连接对象的,所以需要让任务同步的执行。- 这里在实现的时候使用的是数组,这样子在拿出任务队列中的任务时候,可以用另一个数组来置换得到这些任务,然后直接执行即可,而不需要像队列那样子每次拿一个任务去执行的时候还得多次加锁保护这个出队列的过程!
-
一把互斥锁:保护任务队列操作的线程安全。
-
一个线程
ID:判断当前用户的某个操作所处线程是否与其EventLoop线程是对应的。当事件就绪需要处理的时候,如果执行的该操作就在该线程中,则不需要将操作压入任务队列,而是直接执行即可;如果执行的该操作不在该线程中的话,才需要将该操作加入任务队列!
该模块需要具备的功能设计如下:
- 添加连接操作任务到任务队列中的接口
- 定时任务的添加、删除、刷新
- 监控时间的添加、删除、修改
Ⅲ. 定时管理模块整合
之前我们在前置知识中讲到 timerfd 以及自主实现了一个时间轮 timerwheel,如果说我们想实现一个完整的定时管理模块的话,需要将这两者整合起来,因为它们各司其职,充当不同的功能,如下所示:
timerfd:实现内核每隔一段时间,给进程发送一次超时事件(通过触发可读事件)timerwheel:实现程序开始执行之后,可以执行不同时期的非活跃连接或者定时任务。
所以要实现一个完整的秒级定时器,就需要将这两者结合起来!让 timerfd 设置每秒钟触发一次定时事件,当定时事件触发时,则运行 timerwheel 中的 run_timer() 函数,执行一下所有的过期定时任务!
下面直接拿了之前对时间轮类的实现代码,然后加入 timerfd,然后用 Channel 类以及智能指针将其管理起来,然后在构造函数中对其设置可读事件的监控,当触发可读事件的时候,即超时之后,则进行定时任务的删除,所以就有了 timer_read() 函数来处理!
此外还有一个细节,就是添加、刷新、移除定时任务这三个接口,因为涉及到定时器中的 _table 成员也就是哈希表的操作,并且定时器有可能在多线程中进行,因此需要考虑线程安全问题,如果不想加锁,那就把对定期的所有操作,都放到一个线程中进行,所以就统一将它们这些操作添加到对应的 EventLoop 对象中的任务队列中,实现这个目的!
此外还需要注意的是 _loop 成员要在 _timer_channel 之前声明,因为 c++ 中构造函数的初始化顺序和声明次序是一样的!
using func_t = std::function<void()>; // 超时任务的函数类型,由使用者传入
using remove_t = std::function<void()>; // 用于释放weak_ptr的函数类型,由TimerWheel传入
// 定时任务类,封装一个定时任务
class TimerTask
{
private:
uint64_t _id; // 当前超时任务类的ID
uint32_t _timeout; // 超时时间
func_t _task; // 超时任务
remove_t _remove; // 释放TimerWheel中的weak_ptr
bool _cancel; // 为true表示要取消任务,为false表示正常执行任务
public:
TimerTask(uint64_t id, uint32_t timeout, const func_t& task)
: _id(id)
, _timeout(timeout)
, _task(task)
, _cancel(false)
{}
~TimerTask()
{
// 析构函数进行超时任务以及weak_ptr释放函数的执行(如果没有取消任务,才执行释放函数)
if(_cancel == false)
_task();
_remove();
}
uint64_t get_id() { return _id; }
uint32_t get_timeout() { return _timeout; }
void set_remove(const remove_t& remove) { _remove = remove; }
void set_cancel() { _cancel = true; }
};
// 时间轮与timerfd的整合类
class TimerWheel
{
using shared_t = std::shared_ptr<TimerTask>;
using weak_t = std::weak_ptr<TimerTask>;
private:
int _tick; // 当前的时间轮秒数,每一秒就往后走一步
int _capacity; // 时间轮数组大小,即时间轮的周期
std::vector<std::vector<shared_t>> _wheel; // 时间轮数组
std::unordered_map<uint64_t, weak_t> _table; // 保存所有定时任务对象的weak_ptr,这样才能在不影响shared_ptr计数器的同时,获取其shared_ptr
EventLoop* _loop; // 为了初始化_timer_channel和找到当前TimerWheel对应的EventLoop
int _timerfd; // 定时器描述符
std::unique_ptr<Channel> _timer_channel; // 对上面的定时器描述符进行事件管理
public:
TimerWheel(EventLoop* loop)
: _tick(0)
, _capacity(60)
, _wheel(_capacity)
, _loop(loop)
, _timerfd(create_timerfd())
, _timer_channel(new Channel(_timerfd, _loop))
{
// 进行定时器的可读事件设置,当触发可读事件的时候,即超时之后,则进行定时任务的删除
_timer_channel->set_read_callback(std::bind(&TimerWheel::timer_read, this));
_timer_channel->enable_read();
}
/* 定时器中有个_table成员,定时器信息的操作有可能在多线程中进行,因此需要考虑线程安全问题 */
/* 如果不想加锁,那就把对定期的所有操作,都放到一个线程中进行 */
// 将添加超时任务操作添加到对应EventLoop的任务队列中
void add_timertask_in_thread(uint64_t id, uint32_t timeout, const func_t& task);
// 将刷新超时任务操作添加到对应EventLoop的任务队列中
void refresh_timertask_in_thread(uint64_t id);
// 将删除超时任务操作添加到对应EventLoop的任务队列中
void cancel_timertask_in_thread(uint64_t id);
// 判断是否存在定时任务(存在线程安全问题,只能在一个线程中使用)
bool has_timertask(uint64_t id)
{
auto it = _table.find(id);
if(it == _table.end())
return false;
return true;
}
private:
static int create_timerfd()
{
// 1. 创建定时器描述符
int timerfd = timerfd_create(CLOCK_MONOTONIC, 0);
if(timerfd == -1)
{
ELOG("timerfd_create error");
abort();
}
// 设置定时器
struct itimerspec newtimer;
newtimer.it_value.tv_sec = 1; // 设置第一次超时的时间
newtimer.it_value.tv_nsec = 0;
newtimer.it_interval.tv_sec = 1; // 设置第一次超时后每次的超时间隔时间
newtimer.it_interval.tv_nsec = 0;
timerfd_settime(timerfd, 0, &newtimer, nullptr);
DLOG("create_timerfd success, the timerfd is %d", timerfd);
return timerfd;
}
// 定时器超时之后的处理函数
void timer_read()
{
// 1. 读取计数器内容,即清空计数器
// 有可能因为其他描述符的事件处理花费事件比较长,然后在处理定时器描述符事件的时候,有可能就已经超时了很多次
// read读取到的数据times就是从上一次read之后超时的次数
uint64_t times;
int ret = read(_timerfd, ×, 8);
if (ret < 0) {
ELOG("READ TIMEFD FAILED!");
abort();
}
// 2. 进行定时任务的删除
for(int i = 0; i < times; ++i)
run_timer();
}
// 时间运行函数
void run_timer()
{
// 一秒钟走一步,每次将到达的位置处的shared_ptr进行清空,如果是最后一次任务的话会自动调用其析构函数进行释放
_tick = (_tick + 1) % _capacity;
_wheel[_tick].clear();
}
// 在哈希表中去除并且释放weak_ptr
void remove_timer(uint64_t id)
{
// 先判断在不在哈希表中
if(!has_timertask(id))
return;
_table.erase(id);
}
// 添加定时任务
void add_timertask(uint64_t id, uint32_t timeout, const func_t& task)
{
// 1. 创建一个定时任务,由智能指针管理
shared_t newtask(new TimerTask(id, timeout, task));
if(newtask.get() == nullptr)
return;
// 2. 设置释放函数
newtask->set_remove(std::bind(&TimerWheel::remove_timer, this, id));
// 3. 向时间轮数组中添加定时任务
int pos = (_tick + timeout) % _capacity; // 注意需要取模,防止越界
_wheel[pos].push_back(newtask);
// 4. 将定时任务交给哈希表管理,记得要使用weak_ptr才不会导致计数增加
_table[id] = weak_t(newtask);
}
// 刷新定时任务
void refresh_timertask(uint64_t id)
{
// 1. 首先通过哈希表找到保存的超时任务的weak_ptr
auto it = _table.find(id);
if(it == _table.end())
return;
// 2. 通过weak_ptr构造一个shared_ptr出来
shared_t refresh_task(it->second.lock());
// 3. 将刷新任务添加到时间轮数组中
int pos = (_tick + refresh_task->get_timeout()) % _capacity; // 注意需要取模,防止越界
_wheel[pos].push_back(refresh_task);
}
// 取消定时任务
void cancel_timertask(uint64_t id)
{
// 先判断在不在哈希表中
auto it = _table.find(id);
if(it == _table.end())
return;
// 先拿到shared_ptr,再通过其取消任务
shared_t st(it->second.lock());
if(st.get() != nullptr)
st->set_cancel();
}
};Ⅳ. 代码实现
1、基本框架
根据上面给出的设计,我们可以先定义出 EventLoop 的基本框架:
using functor = std::function<void()>;
class EventLoop
{
private:
std::thread::id _tid; // EventLoop对应的线程ID
int _eventfd; // 用于唤醒线程等待事件就绪时候导致的阻塞
std::unique_ptr<Channel> _eventchannel; // 用Channel对象维护上面的_eventfd事件
Poller _poller; // 事件监控管理对象
TimerWheel _tw; // 定时任务管理对象
std::vector<functor> _tasks; // 任务队列(实际是一个数组,方便后面执行操作,减少队列的加锁消耗!)
std::mutex _mtx; // 互斥锁,保护任务队列操作
public:
EventLoop()
: _tid(std::this_thread::get_id())
, _eventfd(create_eventfd())
, _eventchannel(new Channel(_eventfd, this))
, _tw(this)
{
// 启动eventfd的可读事件监控(启动可写事件监控无意义)
_eventchannel->set_read_callback(std::bind(&EventLoop::event_read, this));
_eventchannel->enable_read();
}
// 三步走:启动监控、就绪处理、执行任务总函数
void start();
// 判断将要执行的任务是否处于当前线程中,如果是则直接执行,否则入队列
void run_in_thread(const functor& callback);
// 将任务入队列
void push(const functor& callback);
// 用于执行任务队列中的任务
void run_all_tasks();
// 添加/修改事件监控
void update_event(Channel* channel);
// 移除事件监控
void remove_event(Channel* channel);
// 添加定时任务
void add_timer(uint64_t id, uint32_t timeout, const func_t& task);
// 刷新定时任务
void refresh_timer(uint64_t id);
// 删除定时任务
void cancel_timer(uint64_t id);
// 判断是否存在定时任务
bool has_timer(uint64_t id);
private:
// 判断当前线程是否是EventLoop对应的线程
bool is_in_thread();
// 创建eventfd的函数
static int create_eventfd();
// eventfd的可读事件就绪处理函数
void event_read();
// 唤醒eventfd
void wakeup_eventfd();
}; 可以看到上面的一些函数其实不难实现,就是调用成员变量中已经实现好的接口罢了!
2、EventLoop 完整代码实现
在外部使用 EventLoop 的话,只需要调用 start() 函数就能启动其监听事件、就绪处理、执行任务队列中的任务的操作!并且使用者可以添加、修改、删除监听事件或者定时任务!
using functor = std::function<void()>;
class EventLoop
{
private:
std::thread::id _tid; // EventLoop对应的线程ID
int _eventfd; // 用于唤醒线程等待事件就绪时候导致的阻塞
std::unique_ptr<Channel> _eventchannel; // 用Channel对象维护上面的_eventfd事件
Poller _poller; // 事件监控管理对象
TimerWheel _tw; // 定时任务管理对象
std::vector<functor> _tasks; // 任务队列(实际是一个数组,方便后面执行操作,减少队列的加锁消耗!)
std::mutex _mtx; // 互斥锁,保护任务队列操作
public:
EventLoop()
: _tid(std::this_thread::get_id())
, _eventfd(create_eventfd())
, _eventchannel(new Channel(_eventfd, this))
, _tw(this)
{
// 启动eventfd的可读事件监控(启动可写事件监控无意义)
_eventchannel->set_read_callback(std::bind(&EventLoop::event_read, this));
_eventchannel->enable_read();
}
// 启动监控、就绪处理、执行任务总函数
void start()
{
while(true)
{
// 1. 启动监控
std::vector<Channel*> actives;
_poller.start_event(&actives);
// 2. 处理就绪事件
for(int i = 0; i < actives.size(); ++i)
actives[i]->handler();
// 3. 执行任务队列中的任务
run_all_tasks();
}
}
// 判断将要执行的任务是否处于当前线程中,如果是则直接执行,否则入队列
void run_in_thread(const functor& callback)
{
if(is_in_thread())
callback();
else
push(callback);
}
// 将任务入队列
void push(const functor& callback)
{
{
// 入队列要进行加锁
std::unique_lock<std::mutex> lock(_mtx);
_tasks.push_back(callback);
}
// 唤醒有可能因为没有事件就绪,而导致的epoll阻塞(其实很简单,就是给eventfd写入一条数据就能唤醒,因为触发了可读事件!)
wakeup_eventfd();
}
// 添加/修改事件监控
void update_event(Channel* channel) { return _poller.update_event(channel); }
// 移除事件监控
void remove_event(Channel* channel) { return _poller.remove_event(channel); }
// 添加定时任务
void add_timer(uint64_t id, uint32_t timeout, const func_t& task) { _tw.add_timertask_in_thread(id, timeout, task); }
// 刷新定时任务
void refresh_timer(uint64_t id) { _tw.refresh_timertask_in_thread(id); }
// 删除定时任务
void cancel_timer(uint64_t id) { _tw.cancel_timertask_in_thread(id); }
// 判断是否存在定时任务
bool has_timer(uint64_t id) { return _tw.has_timertask(id); }
public:
// 用于执行任务队列中的任务,该函数不给外界使用
void run_all_tasks()
{
// 开辟一个空的临时数组,将其与任务队列中的数据进行交换,任务队列就变空了
std::vector<functor> tmp;
{
// 交换过程要进行加锁
std::unique_lock<std::mutex> lock(_mtx);
_tasks.swap(tmp);
}
// 剩下的执行就交给临时数组即可,不需要考虑加锁问题
for(int i = 0; i < tmp.size(); ++i)
tmp[i]();
}
// 判断当前线程是否是EventLoop对应的线程
bool is_in_thread() { return _tid == std::this_thread::get_id(); }
private:
static int create_eventfd()
{
int efd = eventfd(0, EFD_CLOEXEC | EFD_NONBLOCK);
if(efd < 0)
{
ELOG("create_eventfd error, 原因:%s", strerror(errno));
abort();
}
DLOG("create_eventfd success, the eventfd is %d", efd);
return efd;
}
// eventfd的可读事件就绪处理函数
void event_read()
{
// 只需要做简单的读取,将内核计数器置零即可!
uint64_t val = 0;
int ret = read(_eventfd, &val, 8);
if(ret <= 0)
{
if(errno == EINTR || errno == EAGAIN) // 如果被打断或者缓冲区为空的话,不算是错误
return;
ELOG("read eventfd fail!");
abort();
}
}
// 唤醒eventfd
void wakeup_eventfd()
{
// 其实很简单,就是给eventfd写入一条数据就能唤醒,因为触发了可读事件!
uint64_t val = 1;
int ret = write(_eventfd, &val, sizeof(val));
if(ret <= 0)
{
if(errno == EINTR) // 如果被打断的话,不算是错误
return;
ELOG("read eventfd fail!");
abort();
}
}
};
void Channel::update() { _eventpoller->update_event(this); }
void Channel::remove() { _eventpoller->remove_event(this); }
void TimerWheel::add_timertask_in_thread(uint64_t id, uint32_t timeout, const func_t& task) { _loop->run_in_thread(std::bind(&TimerWheel::add_timertask, this, id, timeout, task)); }
void TimerWheel::refresh_timertask_in_thread(uint64_t id) { _loop->run_in_thread(std::bind(&TimerWheel::refresh_timertask, this, id)); }
void TimerWheel::cancel_timertask_in_thread(uint64_t id) { _loop->run_in_thread(std::bind(&TimerWheel::cancel_timertask, this, id)); }3、Channel 完整代码实现
只需要将之前测试用的 Poller 类改为 EventLoop 即可!
class EventLoop;
using eventcallback_t = std::function<void()>; // 事件触发的函数类型
class Channel
{
private:
int _fd; // 文件描述符
uint32_t _events; // 当前需要监控的事件
uint32_t _revents; // 当前触发或者就绪的事件(由外部设置)
eventcallback_t _read_callback; // 可读事件被触发的回调函数
eventcallback_t _write_callback; // 可写事件被触发的回调函数
eventcallback_t _error_callback; // 错误事件被触发的回调函数
eventcallback_t _close_callback; // 关闭事件被触发的回调函数
eventcallback_t _arbitrary_callback; // 任意事件被触发的回调函数
EventLoop* _eventpoller;
public:
Channel(int fd, EventLoop* eventpoller)
: _fd(fd), _events(0), _revents(0), _eventpoller(eventpoller)
{}
~Channel()
{
close(_fd); // 记得要释放文件描述符
}
int get_fd() { return _fd; } // 获取文件描述符
uint32_t get_events() { return _events; } // 获取当前监控的事件
void set_revents(uint32_t revents) { _revents = revents; } // 设置实际就绪的事件
// 设置对应触发事件的回调函数
void set_read_callback(const eventcallback_t& cb) { _read_callback = cb; }
void set_write_callback(const eventcallback_t& cb) { _write_callback = cb; }
void set_error_callback(const eventcallback_t& cb) { _error_callback = cb; }
void set_close_callback(const eventcallback_t& cb) { _close_callback = cb; }
void set_arbitrary_callback(const eventcallback_t& cb) { _arbitrary_callback = cb; }
bool is_read_able() { return (_events & EPOLLIN); } // 当前是否监控了可读
bool is_write_able() { return (_events & EPOLLOUT); } // 当前是否监控了可写
// 启动读事件监控
void enable_read() { _events |= EPOLLIN; update(); }
// 启动写事件监控
void enable_write() { _events |= EPOLLOUT; update(); }
// 关闭读事件监控
void disable_read() { _events &= (~EPOLLIN); update(); }
// 关闭写事件监控
void disable_write() { _events &= (~EPOLLOUT); update(); }
// 关闭所有事件监控
void disable_all() { _events = 0; update(); }
// 清除所有的回调函数
void clear_callback()
{
_read_callback = _write_callback = _error_callback = _close_callback = _arbitrary_callback = nullptr;
}
// 事件总处理函数。一旦触发了事件,就调用这个函数,而触发了什么事件如何处理由连接管理者决定
void handler()
{
// 下面因为错误和关闭事件触发的时候会释放连接,此时就不能再调用_arbitrary_callback了,所以需要提前先调用
if(_arbitrary_callback)
_arbitrary_callback(); // 不管任何事件,都调用的回调函数
if((_revents & EPOLLIN) || (_revents & EPOLLRDHUP) ||(_revents & EPOLLPRI))
{
// 如果是有数据可读、对端关闭写入、有带外数据的事件触发的话,则都属于是可读事件处理
//DLOG("read_callback");
if(_read_callback)
_read_callback();
}
// 下面的三个事件有可能会释放连接,所以只能处理一个,要用else if连接
if(_revents & EPOLLOUT)
{
//DLOG("write_callback");
if(_write_callback)
_write_callback(); // 可读事件触发的处理
}
else if(_revents & EPOLLERR)
{
//DLOG("error_callback");
if(_error_callback)
_error_callback(); // 错误事件触发的处理
}
else if(_revents & EPOLLHUP)
{
//DLOG("close_callback");
if(_close_callback)
_close_callback(); // 关闭事件触发的处理
}
}
// 添加或者修改事件监控
void update();
// 移除事件监控
void remove();
};4、测试代码
还是一样,我们用一个服务端和客户端来进行测试,只需要在之前的服务端和客户端代码上稍加修改,加入定时任务的添加,即添加新连接的定时任务为删除函数,这样子就能保证过段时间自动处理完非活跃连接,以及每当触发一次事件之后,都要刷新一下连接的过期时间!
然后模拟客户端发送五次消息后,过十秒之后自动被删除的现象,下面给出客户端代码:
#include "../source/server.hpp"
int main()
{
// 创建客户端套接字
Socket client_sock;
client_sock.create_client(8080, "127.0.0.1");
// 做五次简单的发送和回响,所以会刷新五次连接
for(int i = 0; i < 5; ++i)
{
std::string str = "lirendada";
client_sock.Send(str.c_str(), str.size());
char buf[1024] = { 0 };
client_sock.Recv(buf, sizeof(buf) - 1);
DLOG("%s", buf);
sleep(1);
}
// 进入死循环
while(1) sleep(1);
client_sock.Close();
return 0;
} 然后就是服务端的测试代码:
#include "../source/server.hpp"
void CloseEvent(Channel* channel)
{
DLOG("close:%d", channel->get_fd());
if(channel == nullptr || channel->get_fd() < 0)
return;
channel->clear_callback();
channel->remove(); // 移除监控
delete channel;
}
void ReadEvent(Channel* channel)
{
// 这里读事件处理,我们就做简单的打印、启动可写事件监控即可
int fd = channel->get_fd();
char buffer[1024] = { 0 };
int n = recv(fd, buffer, sizeof(buffer) - 1, 0);
if(n > 0)
{
buffer[n] = 0;
DLOG("接收到:%s", buffer);
// 接收到数据之后,启动可写事件监控
channel->enable_write();
}
else
CloseEvent(channel); // 其实不应该释放,但是因为当前只是测试,所以需要关闭
}
void WriteEvent(Channel* channel)
{
// 这里做个简单的发送即可
int fd = channel->get_fd();
const char* data = "lirendada 你好呀!";
int n = send(fd, data, strlen(data), 0);
if(n < 0)
{
return CloseEvent(channel); // 错误的话释放该对象
}
channel->disable_write(); // 然后关闭可写事件监控
}
void ErrorEvent(Channel* channel)
{
CloseEvent(channel); // 错误的话释放该对象
}
void ArbitraryEvent(Channel* channel, EventLoop* loop, uint64_t timerid)
{
// 刷新非活跃连接
loop->refresh_timer(timerid);
}
void Acceptor(Channel* listen_channel, EventLoop* loop)
{
// 获取新链接
int newfd = accept(listen_channel->get_fd(), nullptr, nullptr);
if(newfd < 0)
{
ELOG("accept error, the error is: ", strerror(errno));
return;
}
// 设置新链接的回调函数
uint64_t id = rand() % 10000;
Channel* channel = new Channel(newfd, loop);
channel->set_read_callback(std::bind(ReadEvent, channel));
channel->set_write_callback(std::bind(WriteEvent, channel));
channel->set_close_callback(std::bind(CloseEvent, channel));
channel->set_error_callback(std::bind(ErrorEvent, channel));
channel->set_arbitrary_callback(std::bind(ArbitraryEvent, channel, loop, id));
// 添加定时任务,即对新连接进行过期删除操作
loop->add_timer(id, 10, std::bind(CloseEvent, channel));
// 启动新链接的可读事件监控
channel->enable_read();
}
int main()
{
srand(time(nullptr));
// 创建服务器套接字
Socket server;
server.create_server(8080);
// 创建一个EventLoop对象
EventLoop loop;
// 创建一个用于监听套接字的Channel对象,然后利用bind函数设置可读回调函数,并且启动可读监控
Channel listen_channel(server.get_fd(), &loop);
listen_channel.set_read_callback(std::bind(Acceptor, &listen_channel, &loop));
listen_channel.enable_read();
while(true)
{
loop.start();
}
server.Close();
return 0;
}