NexusForce 1.0.0
A rigorously engineered full-stack C++ backend library.
载入中...
搜索中...
未找到
timer.hpp
浏览该文件的文档.
1#ifndef NEFORCE_CORE_ASYNC_TIMER_HPP__
2#define NEFORCE_CORE_ASYNC_TIMER_HPP__
3
11
20NEFORCE_BEGIN_NAMESPACE__
21
27
33
42template <typename Clock>
44public:
45 using clock_type = Clock;
46 using time_point = typename clock_type::time_point;
47 using duration = typename clock_type::duration;
48 using token = size_t;
49 using handler_type = function<void()>;
50
51private:
58 struct node {
59 time_point expire;
60 token id;
61 handler_type handler;
62
63 node(time_point exp, const token tid, handler_type&& h) :
64 expire(exp),
65 id(tid),
66 handler(move(h)) {}
67
68 node(const node&) = default;
69 node& operator=(const node&) = default;
70 node(node&&) = default;
71 node& operator=(node&&) = default;
72
73 ~node() = default;
74
78 bool operator<(const node& other) const {
79 if (expire < other.expire) {
80 return true;
81 }
82 if (expire > other.expire) {
83 return false;
84 }
85 return id < other.id;
86 }
87 };
88
89 set<node> nodes_;
93
94 thread thread_;
95 mutable mutex mutex_;
97 token next_id_{1};
98 atomic<bool> stopped_{false};
99
100 friend class thread_pool;
101
102private:
111 void run() {
112 while (!stopped_.load()) {
113 unique_lock<mutex> lock(mutex_);
114
115 if (nodes_.empty()) {
116 cv_.wait_for(lock, 500_ms, [this] { return stopped_.load() || !nodes_.empty(); });
117 if (stopped_.load()) {
118 break;
119 }
120 }
121
122 time_point now = clock_type::now();
123 while (!nodes_.empty() && nodes_.begin()->expire <= now) {
124 auto it = nodes_.begin();
125 node current_node = *it;
126 nodes_.erase(it);
127 node_map_.erase(current_node.id);
128
129 auto flag_it = cancel_flags_.find(current_node.id);
130 const bool cancelled = (flag_it != cancel_flags_.end() && flag_it->second->load(memory_order_acquire));
131
132 lock.unlock_quiet();
133 if (!stopped_.load() && !cancelled) {
134 current_node.handler();
135 }
136 lock.lock_quiet();
137
138 cancel_flags_.erase(current_node.id);
139 promises_.erase(current_node.id);
140 now = clock_type::now();
141 }
142
143 if (!nodes_.empty()) {
144 auto expire_time = nodes_.begin()->expire;
145 cv_.wait_until(lock, expire_time, [this] {
146 return stopped_.load(memory_order_acquire) || nodes_.empty() ||
147 nodes_.begin()->expire <= clock_type::now();
148 });
149 }
150 }
151 }
152
153public:
157 timer_scheduler() { thread_ = thread(&timer_scheduler::run, this); }
158
163
167 void stop() {
168 stopped_.store(true);
169 cv_.notify_one();
170 if (thread_.joinable()) {
171 thread_.join();
172 }
173 }
174
175 timer_scheduler(const timer_scheduler&) = delete;
176 timer_scheduler& operator=(const timer_scheduler&) = delete;
178 timer_scheduler& operator=(timer_scheduler&&) = delete;
179
189 unique_lock<mutex> lock(mutex_);
190 token id = next_id_++;
191
192 auto flag = make_shared<atomic<bool>>(false);
194 handler_type wrapped = [flag, h = move(handler), promise]() {
195 if (!flag->load(memory_order_acquire)) {
196 h();
197 }
199 };
200
201 const bool is_earliest = nodes_.empty() || expire < nodes_.begin()->expire;
202
203 node new_node(expire, id, _NEFORCE move(wrapped));
204 auto result = nodes_.insert(new_node);
205 node_map_[id] = result.first;
206 cancel_flags_[id] = move(flag);
207 promises_[id] = promise;
208
209 if (is_earliest) {
210 cv_.notify_one();
211 }
212 lock.unlock_quiet();
213 return id;
214 }
215
223 bool cancel(token id) {
224 unique_lock<mutex> lock(mutex_);
225
226 const auto flag_it = cancel_flags_.find(id);
227 if (flag_it == cancel_flags_.end()) {
228 return false;
229 }
230 flag_it->second->store(true, memory_order_release);
231
232 auto node_it = node_map_.find(id);
233 if (node_it != node_map_.end()) {
234 const bool is_earliest = (node_it->second == nodes_.begin());
235 nodes_.erase(node_it->second);
236 node_map_.erase(node_it);
237
238 if (is_earliest) {
239 cv_.notify_one();
240 }
241
242 cancel_flags_.erase(flag_it);
243 promises_.erase(id);
244 } else {
245 const auto prom_it = promises_.find(id);
247 if (prom_it != promises_.end()) {
248 prom = prom_it->second;
249 promises_.erase(prom_it);
250 }
251 lock.unlock_quiet();
252
253 if (prom) {
254 prom->get_future().wait();
255 }
256 }
257
258 return true;
259 }
260
264 void cancel_all() {
265 unique_lock<mutex> lock(mutex_);
266 for (const auto& cancel_flag: cancel_flags_) {
267 cancel_flag.second->store(true, memory_order_release);
268 }
269 nodes_.clear();
270 node_map_.clear();
271 cancel_flags_.clear();
272 lock.unlock_quiet();
273 cv_.notify_one();
274 }
275
280 NEFORCE_NODISCARD size_t size() const {
281 lock<mutex> lock(mutex_);
282 return nodes_.size();
283 }
284
290 NEFORCE_NODISCARD bool is_pending(token id) const {
291 lock<mutex> lock(mutex_);
292 return node_map_.find(id) != node_map_.end();
293 }
294};
295
304template <typename Clock>
305class basic_timer {
306public:
307 using clock_type = Clock;
308 using time_point = typename clock_type::time_point;
309 using duration = typename clock_type::duration;
312
313private:
315 token task_id_{0};
316 time_point expire_{clock_type::now()};
317
318public:
319 basic_timer() :
320 scheduler_(make_shared<timer_scheduler<Clock>>()) {}
321
326 if (scheduler_) {
327 cancel();
328 }
329 }
330
331 basic_timer(const basic_timer&) = delete;
332 basic_timer& operator=(const basic_timer&) = delete;
333
337 basic_timer(basic_timer&& other) noexcept :
338 scheduler_(_NEFORCE move(other.scheduler_)),
339 task_id_(other.task_id_),
340 expire_(other.expire_) {
341 other.task_id_ = 0;
342 }
343
347 basic_timer& operator=(basic_timer&& other) noexcept {
348 if (_NEFORCE addressof(other) == this) {
349 return *this;
350 }
351
352 cancel();
353 scheduler_ = _NEFORCE move(other.scheduler_);
354 task_id_ = other.task_id_;
355 expire_ = other.expire_;
356 other.task_id_ = 0;
357
358 return *this;
359 }
360
367 void expires_at(const time_point& expiry_time) {
368 cancel();
369 expire_ = expiry_time;
370 }
371
378 void expires_after(const duration& expiry_duration) {
379 cancel();
380 expire_ = clock_type::now() + expiry_duration;
381 }
382
388
393 NEFORCE_NODISCARD time_point expiry() const { return expire_; }
394
399 NEFORCE_NODISCARD bool is_active() const { return task_id_ != 0 && scheduler_ && scheduler_->is_pending(task_id_); }
400
409 template <typename WaitHandler>
410 void async_wait(WaitHandler&& handler) {
411 cancel();
412 task_id_ = scheduler_->add_task(expire_, handler_type(_NEFORCE forward<WaitHandler>(handler)));
413 }
414
420 void cancel() {
421 if (scheduler_ && task_id_ != 0) {
422 scheduler_->cancel(task_id_);
423 task_id_ = 0;
424 }
425 }
426};
427
432
437 // AsyncTimer
439 // AsyncComponents
441
442NEFORCE_END_NAMESPACE__
443#endif // NEFORCE_CORE_ASYNC_TIMER_HPP__
原子类型完整实现
基本定时器
typename clock_type::time_point time_point
时间点类型
time_point expiry() const
获取到期时间点
Clock clock_type
时钟类型
basic_timer & operator=(basic_timer &&other) noexcept
移动赋值运算符
basic_timer(basic_timer &&other) noexcept
移动构造函数
typename clock_type::duration duration
时长类型
~basic_timer()
析构函数,自动取消未完成的任务
void expires_from_now(const int64_t ms)
设置从当前时间开始的毫秒数
void async_wait(WaitHandler &&handler)
异步等待定时器到期
typename timer_scheduler< Clock >::handler_type handler_type
回调函数类型
void expires_at(const time_point &expiry_time)
设置绝对到期时间
bool is_active() const
检查定时器是否活跃(有待执行的任务)
typename timer_scheduler< Clock >::token token
任务标识符类型
void expires_after(const duration &expiry_duration)
设置相对到期时间
cv_status wait_until(unique_lock< mutex > &lock, const time_point< steady_clock, Dur > &util)
等待直到稳定时钟时间点
cv_status wait_for(unique_lock< mutex > &lock, const duration< Rep, Period > &rest)
等待指定的持续时间
函数包装器主模板声明
锁管理器模板
映射容器
定义 map.hpp:37
iterator end() noexcept
获取结束迭代器
iterator find(const key_type &key)
查找具有指定键的元素
void erase(iterator position) noexcept(noexcept(tree_.erase(position)))
删除指定位置的元素
非递归互斥锁
Promise类模板
void set_value(Res value)
设置结果值
iterator begin() noexcept
获取起始迭代器
void erase(iterator position) noexcept(noexcept(tree_.erase(position)))
删除指定位置的元素
bool empty() const noexcept
检查是否为空
共享智能指针类模板
定时任务调度器
void stop()
停止调度线程并等待其结束
Clock clock_type
时钟类型
void cancel_all()
取消所有定时任务
token add_task(time_point expire, handler_type &&handler)
添加定时任务
typename clock_type::duration duration
时长类型
size_t size() const
获取当前待处理的任务数量
~timer_scheduler()
析构函数,停止调度线程并等待其结束
size_t token
任务标识符类型
timer_scheduler()
构造函数,启动调度线程
bool is_pending(token id) const
检查任务是否仍在等待或执行中
function< void()> handler_type
回调函数类型
bool cancel(token id)
取消定时任务
typename clock_type::time_point time_point
时间点类型
独占锁管理器模板
条件变量行为
集合容器
通用函数包装器
constexpr T * addressof(T &x) noexcept
获取对象的地址
constexpr T && forward(remove_reference_t< T > &x) noexcept
完美转发左值
basic_timer< system_clock > system_timer
基于系统时钟的定时器
basic_timer< steady_clock > steady_timer
基于稳定时钟的定时器
long int64_t
64位有符号整数类型
duration< int64_t, milli > milliseconds
毫秒持续时间
constexpr auto memory_order_release
释放内存顺序常量
constexpr auto memory_order_acquire
获取内存顺序常量
uint64_t size_t
无符号大小类型
enable_if_t<!is_unbounded_array_v< T > &&is_constructible_v< T, Args... >, shared_ptr< T > > make_shared(Args &&... args)
融合分配创建共享指针
constexpr Iterator2 move(Iterator1 first, Iterator1 last, Iterator2 result) noexcept(noexcept(inner::__move_aux(first, last, result)))
移动范围元素
映射容器
NeForce 异步任务包装器
共享智能指针实现
static time_point now() noexcept
获取当前时间点
时间点类模板
线程管理类