Aleph-w 3.0
A C++ Library for Data Structures and Algorithms
Loading...
Searching...
No Matches
timeoutQueue_test.cc
Go to the documentation of this file.
1
11# include <gtest/gtest.h>
12# include <stdexcept>
13# include <atomic>
14# include <chrono>
15# include <thread>
16# include <vector>
17# include <mutex>
18# include <condition_variable>
19# include <timeoutQueue.H>
20
21using namespace std;
22
23// Global queue (singleton)
24static TimeoutQueue* g_queue = nullptr;
25
26// Helper to get the current time plus milliseconds
28{
30}
31
39static bool wait_for_status(
40 const TimeoutQueue::Event* event,
42 const chrono::milliseconds timeout)
43{
44 const auto deadline = chrono::steady_clock::now() + timeout;
45 while (chrono::steady_clock::now() < deadline)
46 {
47 if (event->get_execution_status() == expected)
48 return true;
49 this_thread::sleep_for(chrono::milliseconds(10));
50 }
51 return event->get_execution_status() == expected;
52}
53
68template <typename Pred>
69static bool wait_until(Pred pred, const chrono::milliseconds timeout)
70{
71 const auto deadline = chrono::steady_clock::now() + timeout;
72 while (chrono::steady_clock::now() < deadline)
73 {
74 if (pred())
75 return true;
76 this_thread::sleep_for(chrono::milliseconds(2));
77 }
78 return pred();
79}
80
81// Test event that tracks execution
83{
84public:
85 atomic<bool> executed{false};
86 atomic<int> execution_count{0};
87 function<void()> callback;
88
89 TestEvent(const Time& t) : Event(t) {}
90 TestEvent(const Time& t, function<void()> cb) : Event(t), callback(move(cb)) {}
91
92 void EventFct() override
93 {
94 executed = true;
96 if (callback) callback();
97 }
98};
99
100// Test event that signals a condition variable
102{
103public:
104 mutex& mtx;
106 bool& flag;
107
108 SignalingEvent(const Time& t, mutex& m, condition_variable& c, bool& f)
109 : Event(t), mtx(m), cv(c), flag(f) {}
110
111 void EventFct() override
112 {
114 flag = true;
115 cv.notify_all();
116 }
117};
118
119// Test event that records execution time
121{
122public:
123 chrono::steady_clock::time_point scheduled_at;
124 chrono::steady_clock::time_point executed_at;
125 atomic<bool> executed{false};
126
127 TimingEvent(const Time& t) : Event(t), scheduled_at(chrono::steady_clock::now()) {}
128
129 void EventFct() override
130 {
131 executed_at = chrono::steady_clock::now();
132 executed = true;
133 }
134
135 [[nodiscard]] int elapsed_ms() const
136 {
137 return chrono::duration_cast<chrono::milliseconds>(
138 executed_at - scheduled_at).count();
139 }
140};
141
142// Test event that can reschedule itself
167
168// =============================================================================
169// Test Environment - manages the singleton TimeoutQueue
170// =============================================================================
171
172class TimeoutQueueEnvironment : public ::testing::Environment
173{
174public:
175 void SetUp() override
176 {
177 g_queue = new TimeoutQueue();
178 }
179
180 void TearDown() override
181 {
182 if (g_queue)
183 {
184 g_queue->shutdown();
185 this_thread::sleep_for(chrono::milliseconds(100));
186 delete g_queue;
187 g_queue = nullptr;
188 }
189 }
190};
191
192// =============================================================================
193// Basic Functionality Tests
194// =============================================================================
195
197{
198 mutex mtx;
200 bool executed = false;
201
202 bool completed = false;
203 bool callback_done = false;
204
205 auto* event = new SignalingEvent(time_from_now_ms(100), mtx, cv, executed);
206 event->set_completion_callback([&](TimeoutQueue::Event*, TimeoutQueue::Event::Execution_Status status) {
207 lock_guard<mutex> lock(mtx);
209 callback_done = true;
210 cv.notify_all();
211 });
212 g_queue->schedule_event(event);
213
214 unique_lock<mutex> lock(mtx);
215 ASSERT_TRUE(cv.wait_for(lock, chrono::milliseconds(500),
216 [&]{ return callback_done; }));
217 EXPECT_TRUE(executed);
219
220 delete event;
221}
222
224{
225 const auto completion_timeout = chrono::seconds(30);
226
227 auto* event = new TestEvent(time_from_now_ms(50));
228
229 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::Out_Queue);
230
231 g_queue->schedule_event(event);
232 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::In_Queue);
233
234 // Wait on the lifecycle state instead of a fixed sleep. A fixed delay is
235 // unreliable on heavily loaded CI runners where the worker thread may not run
236 // within the window, leaving the event unexecuted (and turning the following
237 // delete into a use-after-free). wait_for_status returns as soon as the event
238 // completes and only blocks longer when scheduling is genuinely delayed. Once
239 // the status is Executed the worker no longer references the event (no
240 // completion callback is set), so deleting it afterwards is safe.
241 const bool executed =
243 if (not executed)
244 {
245 // The worker never ran the event: it is still In_Queue (or Executing).
246 // Deleting it directly would abort, so hand it back to the queue for a
247 // safe teardown before failing.
248 const auto status = event->get_execution_status();
249 TimeoutQueue::Event* pending = event;
250 try
251 {
253 }
254 catch (const invalid_argument&)
255 {
256 // The worker executed the event between the timeout check and this
257 // call, so it is no longer in the registry; delete it ourselves.
258 delete event;
259 }
260 ADD_FAILURE() << "event did not reach Executed before timeout; "
261 << "last status was " << status;
262 return;
263 }
264
265 EXPECT_TRUE(event->executed);
266
267 delete event;
268}
269
271{
272 const auto completion_timeout = chrono::seconds(30);
273
274 auto* event = new TestEvent(time_from_now_ms(1000)); // Will be overridden
276
278
279 // Wait on the lifecycle state rather than a fixed sleep: the worker thread may
280 // not run within a fixed window on a loaded CI runner, which would both fail
281 // the assertion and turn the delete below into a use-after-free.
282 const bool executed =
284 if (not executed)
285 {
286 const auto status = event->get_execution_status();
287 TimeoutQueue::Event* pending = event;
288 try
289 {
291 }
292 catch (const invalid_argument&)
293 {
294 delete event;
295 }
296 ADD_FAILURE() << "event did not reach Executed before timeout; "
297 << "last status was " << status;
298 return;
299 }
300
301 EXPECT_TRUE(event->executed);
302
303 delete event;
304}
305
307{
308 Time t = time_from_now_ms(100);
309 auto* event = new TestEvent(t);
310
311 Time event_time = event->getAbsoluteTime();
312 EXPECT_EQ(event_time.tv_sec, t.tv_sec);
313 EXPECT_EQ(event_time.tv_nsec, t.tv_nsec);
314
315 delete event;
316}
317
318// =============================================================================
319// Cancellation Tests
320// =============================================================================
321
323{
324 auto* event = new TestEvent(time_from_now_ms(500));
325
326 g_queue->schedule_event(event);
327 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::In_Queue);
328
329 bool canceled = g_queue->cancel_event(event);
331 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::Canceled);
332
333 this_thread::sleep_for(chrono::milliseconds(100));
334 EXPECT_FALSE(event->executed);
335
336 delete event;
337}
338
340{
341 auto* event = new TestEvent(time_from_now_ms(100));
342
343 // Now returns false if event is not in registry
345
346 delete event;
347}
348
350{
352
353 g_queue->schedule_event(event);
355
356 EXPECT_EQ(event, nullptr);
357}
358
360{
361 TimeoutQueue::Event* event = nullptr;
363 EXPECT_EQ(event, nullptr);
364}
365
366// =============================================================================
367// Rescheduling Tests
368// =============================================================================
369
371{
372 const int original_delay_ms = 60000;
373 const int rescheduled_delay_ms = 200;
374 const auto completion_timeout = chrono::seconds(30);
375
377 g_queue->schedule_event(event);
378 this_thread::sleep_for(chrono::milliseconds(100));
379
381
382 const bool completed = wait_for_status(
384 if (not completed)
385 {
386 const auto status = event->get_execution_status();
388 try
389 {
391 }
392 catch (const invalid_argument&)
393 {
394 if (event->get_execution_status() != TimeoutQueue::Event::Executed)
395 throw;
396 delete event;
397 }
398 ADD_FAILURE() << "rescheduled event did not complete before timeout; "
399 << "last status was " << status;
400 return;
401 }
402
403 EXPECT_TRUE(event->executed);
404 EXPECT_LT(event->elapsed_ms(), original_delay_ms);
405
406 delete event;
407}
408
410{
411 auto* event = new TestEvent(time_from_now_ms(100));
412
413 // Now throws exception if event is not in registry
414 EXPECT_THROW(g_queue->reschedule_event(time_from_now_ms(50), event), std::invalid_argument);
415
416 delete event;
417}
418
420{
421 // Pasos de 150 ms y espera total de 1000 ms para absorber jitter en
422 // runners de CI virtualizados (macOS arm64 Debug en particular).
423 // El test sigue verificando la lógica de auto-reprogramación: el evento
424 // debe ejecutarse exactamente max_reschedules+1 = 3 veces.
425 const int step_ms = 150;
426 const int max_reschedules = 2;
428 max_reschedules, step_ms);
429
430 g_queue->schedule_event(event);
431
432 this_thread::sleep_for(chrono::milliseconds(1000));
433 EXPECT_EQ(event->execution_count, max_reschedules + 1);
434
435 delete event;
436}
437
438// =============================================================================
439// Multiple Events Tests
440// =============================================================================
441
443{
444 vector<int> execution_order;
445 mutex order_mutex;
446
447 auto make_event = [&](int id, int delay_ms) {
448 return new TestEvent(time_from_now_ms(delay_ms), [&, id]() {
450 execution_order.push_back(id);
451 });
452 };
453
454 auto* e1 = make_event(1, 150);
455 auto* e2 = make_event(2, 50);
456 auto* e3 = make_event(3, 100);
457
461
462 this_thread::sleep_for(chrono::milliseconds(400));
463
464 {
466 ASSERT_EQ(execution_order.size(), 3u);
470 }
471
472 delete e1;
473 delete e2;
474 delete e3;
475}
476
478{
479 const int num_events = 30;
480 vector<TestEvent*> events;
481 atomic<int> total_executed{0};
482
483 for (int i = 0; i < num_events; ++i)
484 {
485 auto* e = new TestEvent(time_from_now_ms(50 + i * 10), [&]() {
487 });
488 events.push_back(e);
490 }
491
492 this_thread::sleep_for(chrono::milliseconds(600));
493
495
496 for (auto* e : events)
497 delete e;
498}
499
500// =============================================================================
501// Edge Cases and Error Handling
502// =============================================================================
503
505{
507 auto* event = new TestEvent(now);
508
509 g_queue->schedule_event(event);
510
511 this_thread::sleep_for(chrono::milliseconds(100));
512 EXPECT_TRUE(event->executed);
513
514 delete event;
515}
516
518{
519 auto* throwing_event = new TestEvent(time_from_now_ms(50), []() {
520 throw runtime_error("Test exception");
521 });
522
523 auto* normal_event = new TestEvent(time_from_now_ms(100));
524
527
528 this_thread::sleep_for(chrono::milliseconds(300));
529
530 EXPECT_TRUE(throwing_event->executed);
531 EXPECT_TRUE(normal_event->executed);
532
533 delete throwing_event;
534 delete normal_event;
535}
536
538{
539 auto* event = new TestEvent(time_from_now_ms(100));
540
541 event->set_for_deletion();
542 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::To_Delete);
543
544 delete event;
545}
546
547// =============================================================================
548// Timing Accuracy Tests
549// =============================================================================
550
552{
553 // Tolerancias asimétricas: la queue no debe disparar mucho antes de tiempo
554 // (límite inferior estricto), pero en VMs de CI el jitter del scheduler
555 // puede añadir varios cientos de ms al disparo (límite superior generoso).
556 const int expected_delay = 200;
557 const int lower_tolerance = 30;
558 const int upper_tolerance = 250;
559
560 auto* event = new TimingEvent(time_from_now_ms(expected_delay));
561
562 g_queue->schedule_event(event);
563
564 this_thread::sleep_for(
565 chrono::milliseconds(expected_delay + upper_tolerance + 100));
566 ASSERT_TRUE(event->executed);
567
568 int actual_delay = event->elapsed_ms();
571
572 delete event;
573}
574
575// =============================================================================
576// Thread Safety Tests
577// =============================================================================
578
580{
581 atomic<int> executed_count{0};
582 const int num_threads = 4;
583 const int events_per_thread = 5;
584 vector<thread> threads;
586 mutex events_mutex;
587
588 for (int t = 0; t < num_threads; ++t)
589 {
590 threads.emplace_back([&, t]() {
591 for (int i = 0; i < events_per_thread; ++i)
592 {
593 auto* e = new TestEvent(time_from_now_ms(50 + i * 20), [&]() {
594 ++executed_count;
595 });
596 {
598 all_events.push_back(e);
599 }
601 }
602 });
603 }
604
605 for (auto& t : threads)
606 t.join();
607
608 this_thread::sleep_for(chrono::milliseconds(400));
609
610 EXPECT_EQ(executed_count, num_threads * events_per_thread);
611
612 for (auto* e : all_events)
613 delete e;
614}
615
617{
618 const int num_events = 10;
619 vector<TestEvent*> events;
620 atomic<int> canceled_count{0};
621
622 for (int i = 0; i < num_events; ++i)
623 {
624 auto* e = new TestEvent(time_from_now_ms(500));
625 events.push_back(e);
627 }
628
629 vector<thread> threads;
630 for (int t = 0; t < 2; ++t)
631 {
632 threads.emplace_back([&, t]() {
633 for (size_t i = t; i < events.size(); i += 2)
634 {
635 try {
636 if (g_queue->cancel_event(events[i]))
637 ++canceled_count;
638 } catch (const std::invalid_argument&) {
639 // Ignore if already canceled/deleted
640 }
641 }
642 });
643 }
644
645 for (auto& t : threads)
646 t.join();
647
648 EXPECT_EQ(canceled_count, num_events);
649
650 for (auto* e : events)
651 delete e;
652}
653
654// =============================================================================
655// Utility Methods Tests (new features)
656// =============================================================================
657
659{
661 EXPECT_EQ(g_queue->size(), 0u);
662
663 auto* e1 = new TestEvent(time_from_now_ms(500));
664 auto* e2 = new TestEvent(time_from_now_ms(600));
665
668 EXPECT_EQ(g_queue->size(), 1u);
669
671 EXPECT_EQ(g_queue->size(), 2u);
672
674 EXPECT_EQ(g_queue->size(), 1u);
675
678
679 delete e1;
680 delete e2;
681}
682
687
689{
690 // Tolerancias asimétricas para acomodar el jitter del scheduler en
691 // runners de CI virtualizados (en particular macOS Debug). El test
692 // sigue verificando que el evento dispara cerca del delay pedido, no
693 // sólo que "eventualmente se ejecuta".
694 const int expected_delay = 100;
695 const int lower_tolerance = 20;
696 const int upper_tolerance = 250;
697
698 auto* event = new TimingEvent(time_from_now_ms(1000)); // Will be overridden
699
701
702 this_thread::sleep_for(
703 chrono::milliseconds(expected_delay + upper_tolerance + 100));
704 ASSERT_TRUE(event->executed);
705
706 const int actual_delay = event->elapsed_ms();
709
710 delete event;
711}
712
714{
716 EXPECT_EQ(empty_time.tv_sec, 0);
717 EXPECT_EQ(empty_time.tv_nsec, 0);
718
719 auto* e1 = new TestEvent(time_from_now_ms(500));
720 auto* e2 = new TestEvent(time_from_now_ms(200));
721
724 EXPECT_EQ(t1.tv_sec, e1->getAbsoluteTime().tv_sec);
725
728 EXPECT_EQ(t2.tv_sec, e2->getAbsoluteTime().tv_sec); // e2 is sooner
729
732
733 delete e1;
734 delete e2;
735}
736
738{
739 auto* event = new TestEvent(time_from_now_ms(50));
740
741 g_queue->schedule_event(event);
742 this_thread::sleep_for(chrono::milliseconds(200));
743
744 EXPECT_TRUE(event->executed);
745 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::Executed);
746
747 delete event;
748}
749
751{
752 vector<TestEvent*> events;
753 for (int i = 0; i < 5; ++i)
754 {
755 auto* e = new TestEvent(time_from_now_ms(500 + i * 100));
756 events.push_back(e);
758 }
759
760 EXPECT_EQ(g_queue->size(), 5u);
761
762 size_t cleared = g_queue->clear_all();
763 EXPECT_EQ(cleared, 5u);
765
766 for (auto* e : events)
767 {
768 EXPECT_EQ(e->get_execution_status(), TimeoutQueue::Event::Canceled);
769 delete e;
770 }
771}
772
774{
776
781
782 // Execute some events
783 mutex mtx;
785 atomic<int> executed_callbacks{0};
786 auto* e1 = new TestEvent(time_from_now_ms(50));
787 auto* e2 = new TestEvent(time_from_now_ms(100));
788 const auto on_executed =
790 {
791 if (status == TimeoutQueue::Event::Executed)
793 cv.notify_all();
794 };
795 e1->set_completion_callback(on_executed);
796 e2->set_completion_callback(on_executed);
797
800
801 unique_lock<mutex> lock(mtx);
802 ASSERT_TRUE(cv.wait_for(lock, chrono::seconds(30),
803 [&] { return executed_callbacks.load() == 2; }));
804 lock.unlock();
805
807
808 // Cancel some events
809 auto* e3 = new TestEvent(time_from_now_ms(500));
810 auto* e4 = new TestEvent(time_from_now_ms(600));
813
816
818
819 delete e1;
820 delete e2;
821 delete e3;
822 delete e4;
823}
824
826{
827 // Ensure some stats exist
828 mutex mtx;
830 atomic<bool> completed{false};
831
832 auto* e = new TestEvent(time_from_now_ms(50));
833 e->set_completion_callback(
835 {
837 cv.notify_all();
838 });
839
841
842 unique_lock<mutex> lock(mtx);
843 ASSERT_TRUE(cv.wait_for(lock, chrono::seconds(30),
844 [&] { return completed.load(); }));
845 lock.unlock();
846
847 // Reset and verify
851
852 delete e;
853}
854
855// =============================================================================
856// New Features Tests
857// =============================================================================
858
860{
863
864 auto* e1 = new TestEvent(time_from_now_ms(100));
866
867 // Pause before event triggers
868 g_queue->pause();
870
871 // Wait past trigger time - event should NOT execute
872 this_thread::sleep_for(chrono::milliseconds(200));
873 EXPECT_FALSE(e1->executed);
874
875 // Resume - event should execute now
876 g_queue->resume();
878
879 this_thread::sleep_for(chrono::milliseconds(150));
880 EXPECT_TRUE(e1->executed);
881
882 delete e1;
883}
884
886{
887 auto* e1 = new TestEvent(time_from_now_ms(100));
888 auto* e2 = new TestEvent(time_from_now_ms(150));
889
892
894
895 // Wait for all events to complete
899 EXPECT_TRUE(e1->executed);
900 EXPECT_TRUE(e2->executed);
901
902 delete e1;
903 delete e2;
904}
905
907{
908 auto* e = new TestEvent(time_from_now_ms(500));
910
911 // Wait with short timeout - should timeout
915
916 // Cancel and cleanup
918 delete e;
919}
920
933
945
947{
948 // Create event with name
949 class NamedEvent : public TimeoutQueue::Event
950 {
951 public:
952 NamedEvent(const Time& t, const string& name) : Event(t, name) {}
953 void EventFct() override {}
954 };
955
956 auto* e = new NamedEvent(time_from_now_ms(100), "TestEventName");
957 EXPECT_EQ(e->get_name(), "TestEventName");
958
959 e->set_name("NewName");
960 EXPECT_EQ(e->get_name(), "NewName");
961
962 delete e;
963}
964
965TEST(TimeoutQueueTest, CompletionCallback)
966{
967 atomic<bool> callback_called{false};
968 atomic<int> final_status{-1};
969
970 auto* e = new TestEvent(time_from_now_ms(50));
971 e->set_completion_callback([&](TimeoutQueue::Event* ev, TimeoutQueue::Event::Execution_Status status) {
972 (void) ev;
973 callback_called = true;
974 final_status = static_cast<int>(status);
975 });
976
978 this_thread::sleep_for(chrono::milliseconds(200));
979
982
983 delete e;
984}
985
987{
988 atomic<bool> callback_called{false};
989 atomic<int> final_status{-1};
990
991 auto* e = new TestEvent(time_from_now_ms(500));
992 e->set_completion_callback([&](TimeoutQueue::Event*, TimeoutQueue::Event::Execution_Status status) {
993 callback_called = true;
994 final_status = static_cast<int>(status);
995 });
996
999
1000 this_thread::sleep_for(chrono::milliseconds(50));
1001
1004
1005 delete e;
1006}
1007
1009{
1010 auto* e1 = new TestEvent(time_from_now_ms(500));
1011 auto* e2 = new TestEvent(time_from_now_ms(600));
1012
1013 // Each event should have a unique ID
1016 EXPECT_NE(e1->get_id(), e2->get_id());
1017
1018 delete e1;
1019 delete e2;
1020}
1021
1023{
1024 auto* e1 = new TestEvent(time_from_now_ms(500));
1025 auto* e2 = new TestEvent(time_from_now_ms(600));
1026
1027 auto id1 = e1->get_id();
1028 auto id2 = e2->get_id();
1029
1032
1033 // Should find scheduled events
1036
1037 // Invalid ID should return nullptr
1039 EXPECT_EQ(g_queue->find_by_id(999999), nullptr);
1040
1043
1044 // After cancel, should not find
1045 EXPECT_EQ(g_queue->find_by_id(id1), nullptr);
1046
1047 delete e1;
1048 delete e2;
1049}
1050
1052{
1053 auto* e1 = new TestEvent(time_from_now_ms(500));
1054 auto* e2 = new TestEvent(time_from_now_ms(600));
1055
1056 auto id1 = e1->get_id();
1057 auto id2 = e2->get_id();
1058
1061
1062 EXPECT_EQ(g_queue->size(), 2u);
1063
1064 // Cancel by ID
1066 EXPECT_EQ(g_queue->size(), 1u);
1067 EXPECT_EQ(e1->get_execution_status(), TimeoutQueue::Event::Canceled);
1068
1069 // Cancel same ID again should fail
1071
1072 // Cancel invalid ID should fail
1075
1076 // Cancel second event
1079
1080 delete e1;
1081 delete e2;
1082}
1083
1085{
1086 atomic<bool> callback_called{false};
1087 auto* e = new TestEvent(time_from_now_ms(500));
1088 auto id = e->get_id();
1089
1090 e->set_completion_callback([&](TimeoutQueue::Event*, TimeoutQueue::Event::Execution_Status status) {
1091 callback_called = true;
1093 });
1094
1097
1099
1100 delete e;
1101}
1102
1103// =============================================================================
1104// Regression Tests for Bug Fixes
1105// =============================================================================
1106
1108{
1109 // Test that destructor auto-shutdowns if shutdown() wasn't called
1110 // In Debug builds, this should print a warning to stderr
1111 // In Release builds, it should silently auto-shutdown
1112 testing::internal::CaptureStderr();
1113
1114 TimeoutQueue* queue = new TimeoutQueue();
1115 auto* event = new TestEvent(time_from_now_ms(1000));
1116 queue->schedule_event(event);
1117
1118 // Delete without calling shutdown() - should auto-shutdown
1119 delete queue;
1120
1121 string output = testing::internal::GetCapturedStderr();
1122# ifndef NDEBUG
1123 // In debug builds, expect warning message
1124 EXPECT_NE(output.find("WARNING"), string::npos);
1125 EXPECT_NE(output.find("shutdown"), string::npos);
1126# else
1127 // In release builds, no warning should be printed
1128 EXPECT_EQ(output.find("WARNING"), string::npos);
1129# endif
1130
1131 delete event;
1132}
1133
1135{
1136 // Test that canceling an event before deletion is safe (no error)
1137 testing::internal::CaptureStderr();
1138
1139 auto* event = new TestEvent(time_from_now_ms(1000));
1140 g_queue->schedule_event(event);
1141
1142 EXPECT_EQ(event->get_execution_status(), TimeoutQueue::Event::In_Queue);
1143
1144 // Cancel first to remove from the queue, then delete
1145 g_queue->cancel_event(event);
1146
1147 // Now it's safe to delete
1148 delete event;
1149
1150 string output = testing::internal::GetCapturedStderr();
1151 // Should not have warning/error since we canceled first
1152 EXPECT_EQ(output, "");
1153}
1154
1155// Death test: Verify that deleting an In_Queue event throws fatal error
1156// Disabled under ThreadSanitizer: TSAN does not support fork() after threads
1157// have been created, and ASSERT_DEATH uses fork().
1158#if defined(__SANITIZE_THREAD__)
1159 // GCC defines __SANITIZE_THREAD__ when compiled with -fsanitize=thread
1160# define TSAN_ENABLED 1
1161#endif
1162
1163#ifdef __clang__
1164 // Clang uses __has_feature(thread_sanitizer)
1165# if __has_feature(thread_sanitizer)
1166# define TSAN_ENABLED 1
1167# endif
1168#endif
1169
1170#ifdef TSAN_ENABLED
1172#else
1174#endif
1175{
1176 // This test verifies fail-fast behavior to prevent use-after-free.
1177 // When Event::~Event() aborts from destructor, std::terminate() is called.
1178 ASSERT_DEATH({
1179 TimeoutQueue queue;
1180 auto* event = new TestEvent(time_from_now_ms(1000));
1181 queue.schedule_event(event);
1182 // Deleting without cancel should abort, causing process termination
1183 delete event;
1184 queue.shutdown();
1185 }, "In_Queue.*use-after-free");
1186}
1187
1189{
1190 // Regression test for getMin() bug: an event canceled while the worker
1191 // is inside wait_until must not cause the next event to be lost.
1192 //
1193 // Timing mirrors RescheduleDuringTimeout. GitHub-hosted macOS runners
1194 // (both macos-15 arm64 and Rosetta-emulated macos-15-intel) routinely
1195 // see sleep_for jitter exceeding 200 ms when ctest runs --parallel.
1196 // Two margins matter:
1197 //
1198 // 1. e1's trigger (500 ms) must NOT fire before we issue the cancel
1199 // (after a 200 ms sleep) — 300 ms of slack absorbs severe drift.
1200 // If e1 fires first, the worker moves it to Executing and the
1201 // subsequent cancel returns false, breaking every assertion below.
1202 // 2. e2's trigger (1500 ms) stays well past the cancel point so the
1203 // cancel demonstrably happens during the worker's wait_until on
1204 // e1's deadline, and the final 2000 ms wait clears e2 with margin.
1206
1207 auto* e1 = new TestEvent(time_from_now_ms(500));
1208 auto* e2 = new TestEvent(time_from_now_ms(1500));
1209
1212
1213 // Cancel e1 while the worker is inside wait_until on its deadline.
1214 this_thread::sleep_for(chrono::milliseconds(200));
1217
1218 // e2 (≈1500 ms total, ≈1300 ms from here) must still fire; the extra
1219 // ~700 ms of slack absorbs scheduler drift on slow runners.
1220 this_thread::sleep_for(chrono::milliseconds(2000));
1221 EXPECT_FALSE(e1->executed);
1222 EXPECT_TRUE(e2->executed);
1223
1226
1227 delete e1;
1228 delete e2;
1229}
1230
1232{
1233 // Reschedule the top event while the worker is waiting for it.
1234 //
1235 // Timing is calibrated for slow virtualized CI runners (macos-15-intel
1236 // runs x86_64 emulated under Rosetta, where sleep_for jitter routinely
1237 // exceeds 200 ms). Two margins are critical:
1238 //
1239 // 1. e1's original trigger (500 ms) must NOT fire before we issue the
1240 // reschedule (after a 200 ms sleep) — 300 ms of slack absorbs even
1241 // severe scheduler drift.
1242 // 2. e1's rescheduled trigger (200 + 2000 = 2200 ms) must stay well
1243 // ahead of the mid-test assertion (~1500 ms) and well behind the
1244 // final assertion (~3000 ms), so the order of events is observable
1245 // regardless of jitter.
1246 auto* e1 = new TestEvent(time_from_now_ms(500));
1247 auto* e2 = new TestEvent(time_from_now_ms(1200));
1248
1251
1252 // Reschedule e1 to fire much later (absolute trigger ≈ t=2200 ms)
1253 this_thread::sleep_for(chrono::milliseconds(200));
1255
1256 // e2 should have executed at its original time (≈1200 ms),
1257 // e1 should still be waiting on its new trigger (≈2200 ms).
1258 this_thread::sleep_for(chrono::milliseconds(1300));
1259 EXPECT_TRUE(e2->executed);
1260 EXPECT_FALSE(e1->executed);
1261
1262 // e1 should now have fired
1263 this_thread::sleep_for(chrono::milliseconds(1500));
1264 EXPECT_TRUE(e1->executed);
1265
1266 delete e1;
1267 delete e2;
1268}
1269
1271{
1272 // Test that passing nullptr throws exception (not assertion failure in Release)
1273 Time t = time_from_now_ms(100);
1274
1275 EXPECT_THROW(g_queue->schedule_event(nullptr), std::invalid_argument);
1276 EXPECT_THROW(g_queue->schedule_event(t, nullptr), std::invalid_argument);
1277 EXPECT_THROW(g_queue->reschedule_event(t, nullptr), std::invalid_argument);
1278}
1279
1281{
1282 // The implementation enforces tv_nsec bounds with assertions inside Event::set_trigger_time()
1283 // and with ah_domain_error_if inside schedule_event(Event*).
1284 // In Debug builds, the assert triggers first (process death).
1285 // In Release builds, assertions are disabled but ah_domain_error_if throws std::domain_error.
1286# ifndef NDEBUG
1288 {
1289 TimeoutQueue q;
1290 auto* e = new TestEvent(time_from_now_ms(1000));
1292 bad.tv_nsec = 2000000000; // Invalid: >= 1e9
1293 q.schedule_event(bad, e);
1294 },
1295 "");
1296# else
1297 TimeoutQueue q;
1298 auto* e = new TestEvent(time_from_now_ms(1000));
1300 bad.tv_nsec = 2000000000; // Invalid: >= 1e9
1301 EXPECT_THROW(q.schedule_event(bad, e), std::domain_error);
1302 delete e; // Clean up since schedule_event threw before taking ownership
1303 q.shutdown();
1304# endif
1305}
1306
1308{
1309 auto* e = new TestEvent(time_from_now_ms(500));
1311 EXPECT_THROW(g_queue->schedule_event(e), std::invalid_argument);
1313 delete e;
1314}
1315
1317{
1318 TimeoutQueue q;
1319
1320 std::mutex m;
1321 std::condition_variable cv;
1322 int callbacks = 0;
1323
1324 auto* e = new TestEvent(time_from_now_ms(1000));
1325 e->set_completion_callback([&](TimeoutQueue::Event* ev, TimeoutQueue::Event::Execution_Status st) {
1326 std::lock_guard<std::mutex> lk(m);
1327 EXPECT_EQ(ev->get_execution_status(), st);
1329 ++callbacks;
1330 cv.notify_all();
1331 });
1332
1333 q.schedule_event(e);
1334 q.shutdown();
1335
1336 std::unique_lock<std::mutex> lk(m);
1337 ASSERT_TRUE(cv.wait_for(lk, std::chrono::milliseconds(500), [&]{ return callbacks == 1; }));
1338
1340
1341 delete e;
1342}
1343
1345{
1346 // Test that cancel_delete_event invokes completion callback
1347 atomic<bool> callback_called{false};
1348 atomic<int> callback_status{-1};
1349
1351
1352 event->set_completion_callback([&](TimeoutQueue::Event* ev,
1354 EXPECT_EQ(ev, nullptr); // Event already destroyed when status is Deleted
1355 callback_called = true;
1356 callback_status = static_cast<int>(status);
1357 });
1358
1359 g_queue->schedule_event(event);
1361
1364 EXPECT_EQ(event, nullptr);
1365}
1366
1368{
1369 // Test cancel_delete_event on an event that's currently Executing
1370 mutex mtx;
1372 atomic<bool> event_started{false};
1373 atomic<bool> can_finish{false};
1374
1375 class BlockingEvent : public TimeoutQueue::Event
1376 {
1377 public:
1378 atomic<bool>& started;
1379 atomic<bool>& finish_flag;
1380
1381 BlockingEvent(const Time& t, atomic<bool>& s, atomic<bool>& f)
1382 : Event(t), started(s), finish_flag(f) {}
1383
1384 void EventFct() override
1385 {
1386 started = true;
1387 while (!finish_flag)
1388 this_thread::sleep_for(chrono::milliseconds(10));
1389 }
1390 };
1391
1394
1395 atomic<bool> callback_called{false};
1396 atomic<TimeoutQueue::Event*> callback_ev{reinterpret_cast<TimeoutQueue::Event*>(0x1)}; // sentinel
1397 event->set_completion_callback([&](auto* ev, auto status) {
1398 callback_ev = ev;
1399 callback_called = true;
1400 EXPECT_EQ(ev, nullptr); // Event already destroyed when status is Deleted
1402 });
1403
1404 g_queue->schedule_event(event);
1405
1406 // Wait for event to start executing
1407 while (!event_started)
1408 this_thread::sleep_for(chrono::milliseconds(10));
1409
1410 // Try to cancel_delete while it's executing
1411 // Should mark as To_Delete and worker will delete it
1413 EXPECT_EQ(event, nullptr);
1414
1415 // Let event finish
1416 can_finish = true;
1417
1418 // Wait deterministically for the worker to leave EventFct(), delete the event,
1419 // and invoke the callback. A fixed sleep here is unreliable under CPU
1420 // starvation and can leak unfinished worker activity into the next test.
1421 EXPECT_TRUE(wait_until([&] { return callback_called.load(); },
1422 chrono::seconds(5)))
1423 << "Worker did not invoke the deletion callback in time";
1424
1425 // Callback should have been called by worker thread
1427 EXPECT_EQ(callback_ev.load(), nullptr);
1428}
1429
1431{
1432 // Verify that Deleted status yields nullptr, while Executed/Canceled yield
1433 // a valid pointer.
1437 atomic<bool> exec_done{false};
1438
1439 // 1. Executed path: callback must receive non-null pointer.
1440 // Synchronize on the callback itself instead of a fixed sleep: under CPU
1441 // starvation the worker thread may take far longer than any fixed delay to
1442 // run the event, which previously caused this test to fail intermittently
1443 // (and to delete an event the worker still owned).
1445 e1->set_completion_callback([&](TimeoutQueue::Event* ev, auto status) {
1447 EXPECT_NE(ev, nullptr);
1448 exec_ptr = ev;
1449 exec_done = true; // set last: signals the worker is done touching the event
1450 });
1452
1453 const bool executed =
1454 wait_until([&] { return exec_done.load(); }, chrono::seconds(5));
1455 EXPECT_TRUE(executed) << "Executed completion callback was not invoked in time";
1456 EXPECT_NE(exec_ptr.load(), nullptr);
1457 if (executed)
1458 delete e1; // safe: status is Executed and the worker no longer references it
1459 else
1460 g_queue->cancel_delete_event(e1); // best-effort safe removal before return
1461
1462 // 2. Canceled path (cancel_event): the callback runs synchronously in this
1463 // thread, so no waiting is required.
1464 auto* e2 = new TestEvent(time_from_now_ms(500));
1465 e2->set_completion_callback([&](TimeoutQueue::Event* ev, auto status) {
1467 EXPECT_NE(ev, nullptr);
1468 cancel_ptr = ev;
1469 });
1472 EXPECT_NE(cancel_ptr.load(), nullptr);
1473 delete e2;
1474
1475 // 3. Deleted path (cancel_delete_event): callback must receive nullptr
1477 e3->set_completion_callback([&](TimeoutQueue::Event* ev, auto status) {
1479 EXPECT_EQ(ev, nullptr);
1480 delete_ptr = ev;
1481 });
1484 EXPECT_EQ(delete_ptr.load(), nullptr);
1485 EXPECT_EQ(e3, nullptr);
1486}
1487
1489{
1490 // Test that completion callback is called AFTER status is set
1491 // (regression test for callback/status ordering inconsistency)
1492 atomic<bool> callback_called{false};
1495
1496 auto* event = new TestEvent(time_from_now_ms(50));
1497
1498 event->set_completion_callback([&](TimeoutQueue::Event* ev,
1500 // Status should already be set when callback is invoked
1501 status_in_callback = ev->get_execution_status();
1502 EXPECT_EQ(status, ev->get_execution_status());
1503 callback_called = true; // Must be last: main thread deletes after seeing this
1504 });
1505
1506 g_queue->schedule_event(event);
1507
1508 // Spin until callback completes to avoid racing with delete
1509 while (!callback_called)
1510 this_thread::sleep_for(chrono::milliseconds(10));
1511
1513
1514 delete event;
1515}
1516
1518{
1519 atomic<bool> callback_called{false};
1520
1521 auto* e1 = new TestEvent(time_from_now_ms(50), [&]() { g_queue->clear_all(); });
1522 auto* e2 = new TestEvent(time_from_now_ms(200));
1523
1524 e1->set_completion_callback([&](auto*, auto) { callback_called = true; });
1525
1528
1529 this_thread::sleep_for(chrono::milliseconds(400));
1530
1532 EXPECT_TRUE(e1->executed);
1533 EXPECT_FALSE(e2->executed); // should have been canceled by clear_all()
1535
1536 delete e1;
1537 delete e2;
1538}
1539
1541{
1542 // Test that multiple events scheduled for the same time all execute
1543 const int num_events = 5;
1544 vector<TestEvent*> events;
1545 atomic<int> executed_count{0};
1546
1548
1549 for (int i = 0; i < num_events; ++i)
1550 {
1551 auto* e = new TestEvent(same_time, [&]() { ++executed_count; });
1552 events.push_back(e);
1554 }
1555
1556 this_thread::sleep_for(chrono::milliseconds(300));
1557
1558 EXPECT_EQ(executed_count, num_events);
1559
1560 for (auto* e : events)
1561 {
1562 EXPECT_TRUE(e->executed);
1563 delete e;
1564 }
1565}
1566
1567// =============================================================================
1568// Main
1569// =============================================================================
1570
1571int main(int argc, char **argv)
1572{
1573 ::testing::InitGoogleTest(&argc, argv);
1574 ::testing::AddGlobalTestEnvironment(new TimeoutQueueEnvironment());
1575 return RUN_ALL_TESTS();
1576}
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
struct timespec Time
Definition ah-time.H:50
int main()
Minimal std::expected-style result type for C++20.
ReschedulingEvent(const Time &t, TimeoutQueue *q, int max_resc, int step_ms=50)
atomic< int > execution_count
void EventFct() override
Event handler function to be overridden.
condition_variable & cv
void EventFct() override
Event handler function to be overridden.
SignalingEvent(const Time &t, mutex &m, condition_variable &c, bool &f)
function< void()> callback
TestEvent(const Time &t)
void EventFct() override
Event handler function to be overridden.
TestEvent(const Time &t, function< void()> cb)
atomic< bool > executed
atomic< int > execution_count
Base class for scheduled events.
Event(const Time &t, std::string name="")
Execution_Status
Possible states of an event in its lifecycle.
@ Deleted
Memory freed.
@ Out_Queue
Not currently in any queue.
@ 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
void set_completion_callback(CompletionCallback cb)
Set the completion callback (called after EventFct completes or the event is canceled)
virtual void EventFct()=0
Event handler function to be overridden.
static constexpr EventId InvalidId
Invalid/null event ID.
Thread-safe priority queue for scheduling timed events.
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 reset_stats()
Reset statistics counters to zero.
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.
bool is_running() const
Check if the queue is running (not shut down)
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)
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.
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.
bool cancel_event(Event *event)
Cancel a scheduled event.
size_t clear_all()
Cancel all pending events in the queue.
bool is_empty() const
Check if the queue has no pending events.
atomic< bool > executed
chrono::steady_clock::time_point scheduled_at
chrono::steady_clock::time_point executed_at
TimingEvent(const Time &t)
void EventFct() override
Event handler function to be overridden.
int elapsed_ms() const
#define TEST(name)
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
Freq_Node * pred
Predecessor node in level-order traversal.
bool completed() const noexcept
Return true if all underlying iterators are finished.
Definition ah-zip.H:136
STL namespace.
FooMap m(5, fst_unit_pair_hash, snd_unit_pair_hash)
Priority queue for scheduling timed events.
static bool wait_for_status(const TimeoutQueue::Event *event, const TimeoutQueue::Event::Execution_Status expected, const chrono::milliseconds timeout)
Wait until an event reaches an expected lifecycle state.
static TimeoutQueue * g_queue
static bool wait_until(Pred pred, const chrono::milliseconds timeout)
Spin-wait until a predicate becomes true or a timeout elapses.
static Time time_from_now_ms(int ms)
ofstream output
Definition writeHeap.C:215