2016-10-07 16:56:34 -04:00
|
|
|
//
|
|
|
|
// AsyncTaskQueue.hpp
|
|
|
|
// Clock Signal
|
|
|
|
//
|
|
|
|
// Created by Thomas Harte on 07/10/2016.
|
2018-05-13 15:19:52 -04:00
|
|
|
// Copyright 2016 Thomas Harte. All rights reserved.
|
2016-10-07 16:56:34 -04:00
|
|
|
//
|
|
|
|
|
|
|
|
#ifndef AsyncTaskQueue_hpp
|
|
|
|
#define AsyncTaskQueue_hpp
|
|
|
|
|
2017-11-08 22:36:41 -05:00
|
|
|
#include <atomic>
|
|
|
|
#include <condition_variable>
|
|
|
|
#include <functional>
|
2016-10-07 16:56:34 -04:00
|
|
|
#include <thread>
|
2022-07-13 21:36:01 -04:00
|
|
|
#include <vector>
|
2016-10-07 16:56:34 -04:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
#include "../ClockReceiver/TimeTypes.hpp"
|
|
|
|
|
2016-10-07 16:56:34 -04:00
|
|
|
namespace Concurrency {
|
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
/// An implementation detail; provides the time-centric part of a TaskQueue with a real Performer.
|
|
|
|
template <typename Performer> struct TaskQueueStorage {
|
|
|
|
template <typename... Args> TaskQueueStorage(Args&&... args) :
|
|
|
|
performer(std::forward<Args>(args)...),
|
|
|
|
last_fired_(Time::nanos_now()) {}
|
2016-10-07 16:56:34 -04:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
Performer performer;
|
2016-10-07 17:18:46 -04:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
protected:
|
|
|
|
void update() {
|
|
|
|
auto time_now = Time::nanos_now();
|
|
|
|
performer.perform(time_now - last_fired_);
|
|
|
|
last_fired_ = time_now;
|
|
|
|
}
|
2016-10-07 16:56:34 -04:00
|
|
|
|
|
|
|
private:
|
2022-07-14 16:39:26 -04:00
|
|
|
Time::Nanos last_fired_;
|
2016-10-07 16:56:34 -04:00
|
|
|
};
|
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
/// An implementation detail; provides a no-op implementation of time advances for TaskQueues without a Performer.
|
2022-07-15 16:29:29 -04:00
|
|
|
template <> struct TaskQueueStorage<void> {
|
2022-07-14 16:39:26 -04:00
|
|
|
TaskQueueStorage() {}
|
2017-12-15 22:14:09 -05:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
protected:
|
|
|
|
void update() {}
|
|
|
|
};
|
2017-12-15 22:14:09 -05:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
/*!
|
2022-07-16 14:41:04 -04:00
|
|
|
A task queue allows a caller to enqueue @c void(void) functions. Those functions are guaranteed
|
2022-07-14 16:39:26 -04:00
|
|
|
to be performed serially and asynchronously from the caller.
|
2017-12-15 22:14:09 -05:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
If @c perform_automatically is true, functions will be performed as soon as is possible,
|
|
|
|
at the cost of thread synchronisation.
|
2017-12-15 22:14:09 -05:00
|
|
|
|
2022-07-16 14:41:04 -04:00
|
|
|
If @c perform_automatically is false, functions will be queued up but not dispatched
|
2022-07-14 16:39:26 -04:00
|
|
|
until a call to perform().
|
2017-12-15 22:14:09 -05:00
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
If a @c Performer type is supplied then a public member, @c performer will be constructed
|
2022-07-16 14:41:04 -04:00
|
|
|
with the arguments supplied to TaskQueue's constructor. That instance will receive calls of the
|
|
|
|
form @c .perform(nanos) before every batch of new actions, indicating how much time has
|
|
|
|
passed since the previous @c perform.
|
|
|
|
|
|
|
|
@note Even if @c perform_automatically is true, actions may be batched, when a long-running
|
|
|
|
action occupies the asynchronous thread for long enough. So it is not true that @c perform will be
|
|
|
|
called once per action.
|
2022-07-14 16:39:26 -04:00
|
|
|
*/
|
2022-07-16 14:41:04 -04:00
|
|
|
template <bool perform_automatically, typename Performer = void> class AsyncTaskQueue: public TaskQueueStorage<Performer> {
|
2022-07-14 16:39:26 -04:00
|
|
|
public:
|
2022-07-16 14:41:04 -04:00
|
|
|
template <typename... Args> AsyncTaskQueue(Args&&... args) :
|
2022-07-14 16:39:26 -04:00
|
|
|
TaskQueueStorage<Performer>(std::forward<Args>(args)...),
|
|
|
|
thread_{
|
|
|
|
[this] {
|
|
|
|
ActionVector actions;
|
|
|
|
|
2022-07-15 15:13:03 -04:00
|
|
|
while(!should_quit_) {
|
2022-07-14 16:39:26 -04:00
|
|
|
// Wait for new actions to be signalled, and grab them.
|
|
|
|
std::unique_lock lock(condition_mutex_);
|
|
|
|
while(actions_.empty()) {
|
|
|
|
condition_.wait(lock);
|
|
|
|
}
|
|
|
|
std::swap(actions, actions_);
|
|
|
|
lock.unlock();
|
|
|
|
|
|
|
|
// Update to now (which is possibly a no-op).
|
|
|
|
TaskQueueStorage<Performer>::update();
|
|
|
|
|
|
|
|
// Perform the actions and destroy them.
|
|
|
|
for(const auto &action: actions) {
|
|
|
|
action();
|
|
|
|
}
|
|
|
|
actions.clear();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
} {}
|
|
|
|
|
|
|
|
/// Enqueus @c post_action to be performed asynchronously at some point
|
|
|
|
/// in the future. If @c perform_automatically is @c true then the action
|
|
|
|
/// will be performed as soon as possible. Otherwise it will sit unsheculed until
|
|
|
|
/// a call to @c perform().
|
|
|
|
///
|
|
|
|
/// Actions may be elided.
|
|
|
|
///
|
|
|
|
/// If this TaskQueue has a @c Performer then the action will be performed
|
|
|
|
/// on the same thread as the performer, after the performer has been updated
|
|
|
|
/// to 'now'.
|
|
|
|
void enqueue(const std::function<void(void)> &post_action) {
|
|
|
|
std::lock_guard guard(condition_mutex_);
|
|
|
|
actions_.push_back(post_action);
|
|
|
|
|
|
|
|
if constexpr (perform_automatically) {
|
|
|
|
condition_.notify_all();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
/// Causes any enqueued actions that are not yet scheduled to be scheduled.
|
|
|
|
void perform() {
|
|
|
|
if(actions_.empty()) {
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
condition_.notify_all();
|
|
|
|
}
|
|
|
|
|
|
|
|
/// Permanently stops this task queue, blocking until that has happened.
|
|
|
|
/// All pending actions will be performed first.
|
|
|
|
///
|
|
|
|
/// The queue cannot be restarted; this is a destructive action.
|
|
|
|
void stop() {
|
|
|
|
if(thread_.joinable()) {
|
2022-07-15 15:13:03 -04:00
|
|
|
should_quit_ = true;
|
2022-07-14 16:39:26 -04:00
|
|
|
enqueue([] {});
|
|
|
|
if constexpr (!perform_automatically) {
|
|
|
|
perform();
|
|
|
|
}
|
|
|
|
thread_.join();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
/// Schedules any remaining unscheduled work, then blocks synchronously
|
|
|
|
/// until all scheduled work has been performed.
|
|
|
|
void flush() {
|
|
|
|
std::mutex flush_mutex;
|
|
|
|
std::condition_variable flush_condition;
|
|
|
|
bool has_run = false;
|
|
|
|
std::unique_lock lock(flush_mutex);
|
|
|
|
|
|
|
|
enqueue([&flush_mutex, &flush_condition, &has_run] () {
|
|
|
|
std::unique_lock inner_lock(flush_mutex);
|
|
|
|
has_run = true;
|
|
|
|
flush_condition.notify_all();
|
|
|
|
});
|
|
|
|
|
|
|
|
if constexpr (!perform_automatically) {
|
|
|
|
perform();
|
|
|
|
}
|
|
|
|
|
|
|
|
flush_condition.wait(lock, [&has_run] { return has_run; });
|
|
|
|
}
|
|
|
|
|
2022-07-16 14:41:04 -04:00
|
|
|
~AsyncTaskQueue() {
|
2022-07-14 16:39:26 -04:00
|
|
|
stop();
|
|
|
|
}
|
2022-06-02 12:50:45 -04:00
|
|
|
|
2017-12-15 22:14:09 -05:00
|
|
|
private:
|
2022-07-14 16:39:26 -04:00
|
|
|
// The list of actions waiting be performed. These will be elided,
|
|
|
|
// increasing their latency, if the emulation thread falls behind.
|
|
|
|
using ActionVector = std::vector<std::function<void(void)>>;
|
|
|
|
ActionVector actions_;
|
|
|
|
|
|
|
|
// Necessary synchronisation parts.
|
2022-07-15 15:13:03 -04:00
|
|
|
std::atomic<bool> should_quit_ = false;
|
2022-07-14 16:39:26 -04:00
|
|
|
std::mutex condition_mutex_;
|
|
|
|
std::condition_variable condition_;
|
|
|
|
|
|
|
|
// Ensure the thread isn't constructed until after the mutex
|
|
|
|
// and condition variable.
|
|
|
|
std::thread thread_;
|
2017-12-15 22:14:09 -05:00
|
|
|
};
|
|
|
|
|
2016-10-07 16:56:34 -04:00
|
|
|
}
|
|
|
|
|
2022-07-14 16:39:26 -04:00
|
|
|
#endif /* AsyncTaskQueue_hpp */
|