11# include <gtest/gtest.h>
18# include <condition_variable>
42 const chrono::milliseconds timeout)
44 const auto deadline = chrono::steady_clock::now() + timeout;
45 while (chrono::steady_clock::now() <
deadline)
49 this_thread::sleep_for(chrono::milliseconds(10));
51 return event->get_execution_status() ==
expected;
68template <
typename Pred>
71 const auto deadline = chrono::steady_clock::now() + timeout;
72 while (chrono::steady_clock::now() <
deadline)
76 this_thread::sleep_for(chrono::milliseconds(2));
137 return chrono::duration_cast<chrono::milliseconds>(
185 this_thread::sleep_for(chrono::milliseconds(100));
200 bool executed =
false;
215 ASSERT_TRUE(cv.wait_for(lock, chrono::milliseconds(500),
216 [&]{ return callback_done; }));
241 const bool executed =
248 const auto status =
event->get_execution_status();
254 catch (
const invalid_argument&)
260 ADD_FAILURE() <<
"event did not reach Executed before timeout; "
261 <<
"last status was " << status;
282 const bool executed =
286 const auto status =
event->get_execution_status();
292 catch (
const invalid_argument&)
296 ADD_FAILURE() <<
"event did not reach Executed before timeout; "
297 <<
"last status was " << status;
333 this_thread::sleep_for(chrono::milliseconds(100));
378 this_thread::sleep_for(chrono::milliseconds(100));
386 const auto status =
event->get_execution_status();
392 catch (
const invalid_argument&)
398 ADD_FAILURE() <<
"rescheduled event did not complete before timeout; "
399 <<
"last status was " << status;
426 const int max_reschedules = 2;
432 this_thread::sleep_for(chrono::milliseconds(1000));
433 EXPECT_EQ(event->execution_count, max_reschedules + 1);
462 this_thread::sleep_for(chrono::milliseconds(400));
492 this_thread::sleep_for(chrono::milliseconds(600));
496 for (
auto* e : events)
511 this_thread::sleep_for(chrono::milliseconds(100));
520 throw runtime_error(
"Test exception");
528 this_thread::sleep_for(chrono::milliseconds(300));
541 event->set_for_deletion();
564 this_thread::sleep_for(
581 atomic<int> executed_count{0};
582 const int num_threads = 4;
588 for (
int t = 0; t < num_threads; ++t)
590 threads.emplace_back([&, t]() {
605 for (
auto& t : threads)
608 this_thread::sleep_for(chrono::milliseconds(400));
620 atomic<int> canceled_count{0};
630 for (
int t = 0; t < 2; ++t)
632 threads.emplace_back([&, t]() {
633 for (
size_t i = t; i < events.size(); i += 2)
638 }
catch (
const std::invalid_argument&) {
645 for (
auto& t : threads)
650 for (
auto* e : events)
702 this_thread::sleep_for(
742 this_thread::sleep_for(chrono::milliseconds(200));
753 for (
int i = 0; i < 5; ++i)
766 for (
auto* e : events)
803 [&] { return executed_callbacks.load() == 2; }));
833 e->set_completion_callback(
844 [&] { return completed.load(); }));
872 this_thread::sleep_for(chrono::milliseconds(200));
879 this_thread::sleep_for(chrono::milliseconds(150));
957 EXPECT_EQ(e->get_name(),
"TestEventName");
959 e->set_name(
"NewName");
978 this_thread::sleep_for(chrono::milliseconds(200));
1000 this_thread::sleep_for(chrono::milliseconds(50));
1027 auto id1 =
e1->get_id();
1028 auto id2 =
e2->get_id();
1056 auto id1 =
e1->get_id();
1057 auto id2 =
e2->get_id();
1088 auto id = e->get_id();
1112 testing::internal::CaptureStderr();
1121 string output = testing::internal::GetCapturedStderr();
1137 testing::internal::CaptureStderr();
1150 string output = testing::internal::GetCapturedStderr();
1158#if defined(__SANITIZE_THREAD__)
1160# define TSAN_ENABLED 1
1165# if __has_feature(thread_sanitizer)
1166# define TSAN_ENABLED 1
1185 },
"In_Queue.*use-after-free");
1214 this_thread::sleep_for(chrono::milliseconds(200));
1220 this_thread::sleep_for(chrono::milliseconds(2000));
1253 this_thread::sleep_for(chrono::milliseconds(200));
1258 this_thread::sleep_for(chrono::milliseconds(1300));
1263 this_thread::sleep_for(chrono::milliseconds(1500));
1292 bad.tv_nsec = 2000000000;
1300 bad.tv_nsec = 2000000000;
1321 std::condition_variable cv;
1326 std::lock_guard<std::mutex>
lk(
m);
1336 std::unique_lock<std::mutex>
lk(
m);
1337 ASSERT_TRUE(cv.wait_for(
lk, std::chrono::milliseconds(500), [&]{ return callbacks == 1; }));
1388 this_thread::sleep_for(chrono::milliseconds(10));
1408 this_thread::sleep_for(chrono::milliseconds(10));
1422 chrono::seconds(5)))
1423 <<
"Worker did not invoke the deletion callback in time";
1453 const bool executed =
1455 EXPECT_TRUE(executed) <<
"Executed completion callback was not invoked in time";
1510 this_thread::sleep_for(chrono::milliseconds(10));
1529 this_thread::sleep_for(chrono::milliseconds(400));
1545 atomic<int> executed_count{0};
1552 events.push_back(e);
1556 this_thread::sleep_for(chrono::milliseconds(300));
1560 for (
auto* e : events)
1573 ::testing::InitGoogleTest(&
argc,
argv);
Time time_plus_msec(const Time ¤t_time, const int &msec)
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.
void EventFct() override
Event handler function to be overridden.
SignalingEvent(const Time &t, mutex &m, condition_variable &c, bool &f)
function< void()> callback
void EventFct() override
Event handler function to be overridden.
TestEvent(const Time &t, function< void()> cb)
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.
@ 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.
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.
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Freq_Node * pred
Predecessor node in level-order traversal.
bool completed() const noexcept
Return true if all underlying iterators are finished.
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)