Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
timeoutQueue.C
Go to the documentation of this file.
1/*
2 Aleph_w
3
4 Data structures & Algorithms
5 version 2.0.0b
6 https://github.com/lrleon/Aleph-w
7
8 This file is part of Aleph-w library
9
10 Copyright (c) 2002-2026 Leandro Rabindranath Leon
11
12 Permission is hereby granted, free of charge, to any person obtaining a copy
13 of this software and associated documentation files (the "Software"), to deal
14 in the Software without restriction, including without limitation the rights
15 to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
16 copies of the Software, and to permit persons to whom the Software is
17 furnished to do so, subject to the following conditions:
18
19 The above copyright notice and this permission notice shall be included in all
20 copies or substantial portions of the Software.
21
22 THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
23 IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
24 FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
25 AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
26 LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
27 OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
28 SOFTWARE.
29*/
30
31# include <cstdio>
32# include <typeinfo>
33# include <timeoutQueue.H>
34# include <ah-errors.H>
35
36using namespace std::chrono;
37
38// Initialize static event ID counter
39std::atomic<TimeoutQueue::Event::EventId> TimeoutQueue::Event::nextId{0};
40
41// Convert POSIX timespec to std::chrono::time_point (system_clock)
42static auto timespec_to_timepoint(const Time & t)
43{
44 return system_clock::time_point(duration_cast<system_clock::duration>(
45 seconds(t.tv_sec) + nanoseconds(t.tv_nsec)));
46}
47
49{
50 workerThread = std::thread(&TimeoutQueue::triggerEvent, this);
51}
52
54{ {
55 std::lock_guard<std::mutex> lock(mtx);
56 if (not isShutdown)
57 {
58#ifndef NDEBUG
59 ah_warning(std::cerr)
60 << "TimeoutQueue destructor called without prior shutdown(). "
61 << "Invoking shutdown() automatically." << std::endl;
62#endif
64 }
65 }
66
67 if (workerThread.joinable())
68 workerThread.join();
69}
70
73{
74 ah_invalid_argument_if(event == nullptr)
75 << "nullptr event";
76
77 event->set_trigger_time(trigger_time);
78 schedule_event(event);
79}
80
82{
83 ah_invalid_argument_if(event == nullptr)
84 << "nullptr event";
85 ah_domain_error_if(event->time_key().tv_nsec < 0 or event->time_key().tv_nsec >= NSEC)
86 << "event nsec out of range: " << event->time_key().tv_nsec;
87
88 std::lock_guard<std::mutex> lock(mtx);
89
90 event_registry.insert(event);
91
93 << "Event has already been inserted in timeoutQueue";
94
95 if (isShutdown)
96 return;
97
98 event->set_execution_status(Event::In_Queue);
99
100 prio_queue.insert(event);
101 event_map[event->get_id()] = event;
102
103 cond.notify_one();
104}
105
107{
108 ah_invalid_argument_if(event == nullptr)
109 << "nullptr event";
110
112 bool became_empty = false; {
113 std::lock_guard<std::mutex> lock(mtx);
114
115 // If the event is no longer known to the queue, treat it as already
116 // completed/canceled instead of throwing. This can happen if cancellation
117 // races with execution completion.
118 if (not event_registry.contains(event))
119 return false;
120
122 return false;
123
124 callback = event->on_completed;
125
126 prio_queue.remove(event);
127 if (event_map.contains(event->get_id()))
128 event_map.remove(event->get_id());
129 event_registry.remove(event);
130
131 event->set_execution_status(Event::Canceled);
133
135 }
136
137 if (callback)
138 callback(event, Event::Canceled);
139
140 if (became_empty)
141 emptyCondition.notify_all();
142
143 cond.notify_one();
144
145 return true;
146}
147
149{
150 if (event == nullptr)
151 return;
152
153 Event *local = event;
156 bool became_empty = false; {
157 bool was_in_queue = false;
158 std::lock_guard<std::mutex> lock(mtx);
159
161 << "Event " << local << " not found in timeoutQueue";
162
164 {
165 prio_queue.remove(local);
166 if (event_map.contains(local->get_id()))
167 event_map.remove(local->get_id());
168 event_registry.remove(local);
170 was_in_queue = true;
171 }
172
174 {
175 // Worker thread will invoke callback and delete after EventFct() returns
177 event_registry.remove(local);
178
179 // Also remove from event_map to prevent find_by_id() from returning a pointer
180 // that is about to be deleted by the worker thread.
181 if (event_map.contains(local->get_id()))
182 event_map.remove(local->get_id());
183
184 event = nullptr;
185 cond.notify_one();
186 return;
187 }
188
189 callback = local->on_completed;
190
191 if (was_in_queue)
193
196 }
197
198 delete local;
199 event = nullptr;
200
201 if (callback)
202 callback(nullptr, final_status);
203
204 if (became_empty)
205 emptyCondition.notify_all();
206
207 cond.notify_one();
208}
209
211 TimeoutQueue::Event *event)
212{
213 ah_invalid_argument_if(event == nullptr)
214 << "nullptr event";
215
216 std::lock_guard<std::mutex> lock(mtx);
217
219 << "Event " << event << " not found in timeoutQueue";
220
221 if (isShutdown)
222 return;
223
225 prio_queue.remove(event); // eventMap entry stays, ID doesn't change
226
227 event->set_trigger_time(trigger_time);
228
229 event->set_execution_status(Event::In_Queue);
230
231 prio_queue.insert(event);
232 event_map[event->get_id()] = event; // Re-add in case it wasn't there
233
234 cond.notify_one();
235}
236
238{
239 std::unique_lock<std::mutex> lock(mtx);
240
241 while (true)
242 {
243 // Sleep if there are no events or if paused
244 while ((prio_queue.size() == 0 or isPaused) and not isShutdown)
245 cond.wait(lock);
246
247 if (isShutdown)
248 break;
249
250 // Read the soonest event
251 auto *event_to_schedule = static_cast<Event *>(prio_queue.top());
252
253 // Compute time when the event must be triggered (wall clock)
256
257 // Anchor both clocks under the same lock to avoid skew
258 const auto sys_now = system_clock::now();
259 const auto steady_now = steady_clock::now();
260
261 // Compute delta in system_clock domain, clamp negative to zero
262 auto delta = trigger_sys - sys_now;
263 if (delta < system_clock::duration::zero())
264 delta = system_clock::duration::zero();
265
266 // Convert delta to steady_clock and build a steady deadline
267 const auto deadline_steady =
269
270 // Wait until deadline or notification (immune to wall-clock jumps)
271 const auto wait_result = cond.wait_until(lock, deadline_steady);
272
273 if (isShutdown)
274 break;
275
276 // If paused, go back to waiting
277 if (isPaused)
278 continue;
279
280 if (wait_result == std::cv_status::timeout)
281 {
282 if (prio_queue.size() == 0)
283 continue;
284
285 // Peek at the soonest event without extracting it
286 // If the top changed (original was canceled/rescheduled) and
287 // the new top is in the future, go back to wait for it
288 if (const auto *next = static_cast<Event *>(prio_queue.top()); next->time_key() > original_trigger_time)
289 continue;
290
291 // Now extract the event we are going to execute
292 auto *event_to_execute = static_cast<Event *>(prio_queue.getMin());
293
295
296 lock.unlock();
297
298 try { event_to_execute->EventFct(); }
299 catch (...)
300 {
301 ah_warning(std::cerr) << "Uncaught exception in TimeoutQueue event execution (ID "
302 << event_to_execute->get_id() << ", name: '"
303 << event_to_execute->get_name() << "')" << std::endl;
304 }
305
306 lock.lock();
307
309
310 const auto current_status = event_to_execute->get_execution_status();
311 const Event::CompletionCallback callback = event_to_execute->on_completed;
313
315 {
317 event_to_execute->set_execution_status(Event::Deleted);
318 }
320 {
321 // Event was rescheduled during EventFct() - still in queue, don't touch
323 }
324 else
325 {
326 event_map.remove(event_to_execute->get_id());
328 event_to_execute->set_execution_status(Event::Executed);
329 }
330
331 const bool became_empty = prio_queue.size() == 0;
332
333 lock.unlock();
334
336 {
337 delete event_to_execute;
338 if (callback)
339 callback(nullptr, final_status);
340 }
341 else if (callback)
343
344 if (became_empty)
345 emptyCondition.notify_all();
346
347 lock.lock();
348 }
349 }
350
351 // Shutdown requested - cancel all pending events
352 while (prio_queue.size() > 0)
353 {
354 auto *event = static_cast<Event *>(prio_queue.getMin());
355 const Event::CompletionCallback callback = event->on_completed;
356 event_map.remove(event->get_id());
357 event_registry.remove(event);
358 // Set final status BEFORE callback to avoid use-after-free:
359 // User code may delete the event in the callback
360 event->set_execution_status(Event::Canceled);
362 lock.unlock();
363 if (callback)
364 callback(event, Event::Canceled);
365 lock.lock();
366 }
367
368 emptyCondition.notify_all();
369}
370
372{
373 std::lock_guard<std::mutex> lock(mtx);
375}
376
377size_t TimeoutQueue::size() const
378{
379 std::lock_guard<std::mutex> lock(mtx);
380 return prio_queue.size();
381}
382
384{
385 std::lock_guard<std::mutex> lock(mtx);
386 return prio_queue.size() == 0;
387}
388
390{
391 std::lock_guard<std::mutex> lock(mtx);
392 return not isShutdown;
393}
394
400
402{
403 std::lock_guard<std::mutex> lock(mtx);
404 if (prio_queue.size() == 0)
405 return {0, 0};
406 return static_cast<Event *>(prio_queue.top())->time_key();
407}
408
410{
411 std::unique_lock<std::mutex> lock(mtx);
412
413 size_t count = 0;
414 while (prio_queue.size() > 0)
415 {
416 auto *event = static_cast<Event *>(prio_queue.getMin());
417 const Event::CompletionCallback callback = event->on_completed;
418 event_map.remove(event->get_id());
419 event_registry.remove(event);
420 event->set_execution_status(Event::Canceled);
421 ++count;
423 lock.unlock();
424 if (callback)
425 callback(event, Event::Canceled);
426 lock.lock();
427 }
428
429 cond.notify_one();
430 emptyCondition.notify_all();
431
432 return count;
433}
434
436{
437 std::lock_guard<std::mutex> lock(mtx);
438 return executedCount;
439}
440
442{
443 std::lock_guard<std::mutex> lock(mtx);
444 return canceledCount;
445}
446
448{
449 std::lock_guard<std::mutex> lock(mtx);
450 executedCount = 0;
451 canceledCount = 0;
452}
453
455{
456 std::lock_guard<std::mutex> lock(mtx);
457 isPaused = true;
458}
459
461{
462 std::lock_guard<std::mutex> lock(mtx);
463 isPaused = false;
464 cond.notify_one();
465}
466
468{
469 std::lock_guard<std::mutex> lock(mtx);
470 return isPaused;
471}
472
474{
475 std::unique_lock<std::mutex> lock(mtx);
476
477 if (prio_queue.size() == 0)
478 return true;
479
480 if (timeout_ms <= 0)
481 {
482 emptyCondition.wait(lock, [this]() { return prio_queue.size() == 0 or isShutdown; });
483 return prio_queue.size() == 0;
484 }
485
486 return emptyCondition.wait_for(lock, std::chrono::milliseconds(timeout_ms),
487 [this]() { return prio_queue.size() == 0 or isShutdown; });
488}
489
491{
492 if (id == Event::InvalidId)
493 return nullptr;
494
495 std::lock_guard<std::mutex> lock(mtx);
496
497 if (const auto ptr_pair = event_map.search(id);
498 ptr_pair != nullptr and ptr_pair->second->get_execution_status() == Event::In_Queue)
499 return ptr_pair->second;
500
501 return nullptr;
502}
503
505{
506 if (id == Event::InvalidId)
507 return false;
508
510 Event *event = nullptr;
511 bool became_empty = false; {
512 std::lock_guard<std::mutex> lock(mtx);
513
514 const auto event_pair = event_map.search(id);
515 if (event_pair == nullptr)
516 return false;
517
518 event = event_pair->second;
519 if (event->get_execution_status() != Event::In_Queue)
520 return false;
521
522 callback = event->on_completed;
523 prio_queue.remove(event);
524 event_map.remove(id);
525 event_registry.remove(event);
526 event->set_execution_status(Event::Canceled);
528
530 }
531
532 if (callback)
533 callback(event, Event::Canceled);
534 cond.notify_one();
535
536 if (became_empty)
537 emptyCondition.notify_all();
538
539 return true;
540}
Exception handling system with formatted messages for Aleph-w.
#define ah_warning(out)
Emits an unconditional warning to a stream.
Definition ah-errors.H:177
#define ah_invalid_argument_unless(C)
Throws std::invalid_argument if condition does NOT hold.
Definition ah-errors.H:660
#define ah_domain_error_if(C)
Throws std::domain_error if condition holds.
Definition ah-errors.H:527
#define ah_invalid_argument_if(C)
Throws std::invalid_argument if condition holds.
Definition ah-errors.H:644
Time time_plus_msec(const Time &current_time, const int &msec)
Definition ah-time.H:112
Time read_current_time()
Definition ah-time.H:103
#define NSEC
Definition ah-time.H:55
struct timespec Time
Definition ah-time.H:50
Node * getMin()
Removes the node with the lowest priority from the heap.
Node * remove(Node *node)
Removes node from the heap.
Node * insert(Node *p) noexcept
Inserts a node into a heap.
Node * top()
Returns the node with the lowest priority according to the comparison criterion specified in the decl...
const size_t & size() const noexcept
Base class for scheduled events.
CompletionCallback on_completed
void set_execution_status(Execution_Status status)
Execution_Status
Possible states of an event in its lifecycle.
@ Executing
Currently executing EventFct()
@ Deleted
Memory freed.
@ Canceled
Removed from queue before execution.
@ Executed
Completed execution.
@ In_Queue
Scheduled and waiting for trigger time.
@ To_Delete
Marked for cleanup.
Execution_Status get_execution_status() const
static std::atomic< EventId > nextId
std::function< void(Event *, Execution_Status)> CompletionCallback
Optional callback invoked after an event completes, is canceled, or is deleted.
uint64_t EventId
Type for unique event identifiers.
static constexpr EventId InvalidId
Invalid/null event ID.
EventId get_id() const
Get the unique event ID (auto-generated on construction)
std::condition_variable cond
size_t executed_count() const
Get count of events that have been executed.
void reschedule_event(const Time &trigger_time, Event *event)
Reschedule an event to a new time.
void shutdown_locked()
void reset_stats()
Reset statistics counters to zero.
DynMapTree< Event::EventId, Event * > event_map
size_t size() const
Get the number of pending events in the queue.
void schedule_event(const Time &trigger_time, Event *)
Schedule an event at a specific time.
~TimeoutQueue()
Destructor - shuts down the queue and joins the thread.
bool is_running() const
Check if the queue is running (not shut down)
std::mutex mtx
Time next_event_time() const
Get the trigger time of the next (soonest) event.
bool cancel_by_id(Event::EventId id)
Cancel a scheduled event by its ID.
void resume()
Resume event execution after pause.
size_t canceled_count() const
Get count of events that have been canceled.
void pause()
Pause event execution (events remain scheduled but won't execute)
std::thread workerThread
size_t executedCount
void cancel_delete_event(Event *&event)
Cancel and delete an event.
void shutdown()
Shut down the queue and stop the background thread.
bool wait_until_empty(int timeout_ms=0)
Block until all pending events have been executed or canceled.
void schedule_after_ms(int ms_from_now, Event *event)
Schedule an event relative to the current time.
void triggerEvent()
bool is_paused() const
Check if the queue is paused.
Event * find_by_id(const Event::EventId id) const
Find a scheduled event by its ID.
BinHeapVtl< Time > prio_queue
bool cancel_event(Event *event)
Cancel a scheduled event.
size_t clear_all()
Cancel all pending events in the queue.
TimeoutQueue()
Default constructor - starts the background thread.
bool is_empty() const
Check if the queue has no pending events.
std::condition_variable emptyCondition
DynSetTree< Event * > event_registry
size_t canceledCount
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Definition Blossom.H:466
and
Check uniqueness with explicit hash + equality functors.
void next()
Advance all underlying iterators (bounds-checked).
Definition ah-zip.H:171
Itor::difference_type count(const Itor &beg, const Itor &end, const T &value)
Count elements equal to a value.
Definition ahAlgo.H:127
static auto timespec_to_timepoint(const Time &t)
Priority queue for scheduling timed events.