73 #include "spmc_queue.hpp" 74 #include "notifier.hpp" 75 #include "observer.hpp" 76 #include "taskflow.hpp" 93 std::optional<Node*> cache;
177 template<
typename P,
typename C>
205 template<
typename Observer,
typename... Args>
219 unsigned _num_topologies {0};
236 unsigned _find_victim(
unsigned);
238 PerThread& _per_thread()
const;
240 bool _wait_for_task(
unsigned, std::optional<Node*>&);
242 void _spawn(
unsigned);
243 void _exploit_task(
unsigned, std::optional<Node*>&);
244 void _explore_task(
unsigned, std::optional<Node*>&);
245 void _schedule(Node*,
bool);
246 void _schedule(PassiveVector<Node*>&);
247 void _invoke(
unsigned, Node*);
248 void _invoke_static_work(
unsigned, Node*);
249 void _invoke_dynamic_work(
unsigned, Node*,
Subflow&);
250 void _init_module_node(Node*);
251 void _tear_down_topology(Topology*);
252 void _increment_topology();
253 void _decrement_topology();
254 void _decrement_topology_and_notify();
261 _notifier {_waiters} {
273 _notifier.notify(
true);
275 for(
auto& t : _threads){
282 return _workers.size();
286 inline Executor::PerThread& Executor::_per_thread()
const {
287 thread_local PerThread pt;
292 inline void Executor::_spawn(
unsigned N) {
295 for(
unsigned i=0; i<N; ++i) {
296 _threads.emplace_back([
this, i] () ->
void {
298 PerThread& pt = _per_thread();
302 std::optional<Node*> t;
311 if(_wait_for_task(i, t) ==
false) {
321 inline unsigned Executor::_find_victim(
unsigned thief) {
343 for(
unsigned vtm=0; vtm<_workers.size(); ++vtm){
344 if((thief == vtm && !_queue.
empty()) ||
345 (thief != vtm && !_workers[vtm].queue.empty())) {
350 return _workers.size();
354 inline void Executor::_explore_task(
unsigned thief, std::optional<Node*>& t) {
359 const unsigned l = 0;
360 const unsigned r = _workers.size() - 1;
362 const size_t F = (_workers.size() + 1) << 1;
363 const size_t Y = 100;
372 _workers[thief].rdgen
375 t = (vtm == thief) ? _queue.
steal() : _workers[vtm].queue.steal();
406 inline void Executor::_exploit_task(
unsigned i, std::optional<Node*>& t) {
408 assert(!_workers[i].cache);
411 auto& worker = _workers[i];
412 if(_num_actives.fetch_add(1) == 0 && _num_thieves == 0) {
413 _notifier.notify(
false);
420 worker.cache = std::nullopt;
423 t = worker.queue.pop();
433 inline bool Executor::_wait_for_task(
unsigned me, std::optional<Node*>& t) {
443 if(_explore_task(me, t); t) {
444 if(
auto N = _num_thieves.fetch_sub(1); N == 1) {
445 _notifier.notify(
false);
450 _notifier.prepare_wait(&_waiters[me]);
453 if(!_queue.
empty()) {
455 _notifier.cancel_wait(&_waiters[me]);
458 if(t = _queue.
steal(); t) {
459 if(
auto N = _num_thieves.fetch_sub(1); N == 1) {
460 _notifier.notify(
false);
470 _notifier.cancel_wait(&_waiters[me]);
471 _notifier.notify(
true);
476 if(_num_thieves.fetch_sub(1) == 1 && _num_actives) {
477 _notifier.cancel_wait(&_waiters[me]);
482 _notifier.commit_wait(&_waiters[me]);
488 template<
typename Observer,
typename... Args>
491 auto tmp = std::make_unique<Observer>(std::forward<Args>(args)...);
492 tmp->set_up(_workers.size());
493 _observer = std::move(tmp);
494 return static_cast<Observer*
>(_observer.get());
505 inline void Executor::_schedule(Node* node,
bool bypass) {
507 assert(_workers.size() != 0);
510 if(node->_module !=
nullptr && !node->_module->empty() && !node->is_spawned()) {
511 _init_module_node(node);
515 if(
auto& pt = _per_thread(); pt.pool ==
this) {
517 _workers[pt.worker_id].queue.push(node);
520 assert(!_workers[pt.worker_id].cache);
521 _workers[pt.worker_id].cache = node;
528 std::scoped_lock lock(_queue_mutex);
532 _notifier.notify(
false);
538 inline void Executor::_schedule(PassiveVector<Node*>& nodes) {
540 assert(_workers.size() != 0);
544 const auto num_nodes = nodes.size();
550 for(
auto node : nodes) {
551 if(node->_module !=
nullptr && !node->_module->empty() && !node->is_spawned()) {
552 _init_module_node(node);
557 if(
auto& pt = _per_thread(); pt.pool ==
this) {
558 for(
size_t i=0; i<num_nodes; ++i) {
559 _workers[pt.worker_id].queue.push(nodes[i]);
566 std::scoped_lock lock(_queue_mutex);
567 for(
size_t k=0; k<num_nodes; ++k) {
568 _queue.
push(nodes[k]);
572 if(num_nodes >= _workers.size()) {
573 _notifier.notify(
true);
576 for(
size_t k=0; k<num_nodes; ++k) {
577 _notifier.notify(
false);
583 inline void Executor::_init_module_node(Node* node) {
585 node->_work = [node=node,
this, tgt{PassiveVector<Node*>()}] ()
mutable {
588 if(node->is_spawned()) {
589 node->_dependents.resize(node->_dependents.size()-tgt.size());
591 t->_successors.clear();
599 PassiveVector<Node*> src;
601 for(
auto& n: node->_module->_graph.nodes()) {
602 n->_topology = node->_topology;
603 if(n->num_dependents() == 0) {
604 src.push_back(n.get());
606 if(n->num_successors() == 0) {
608 tgt.push_back(n.get());
617 inline void Executor::_invoke(
unsigned me, Node* node) {
619 assert(_workers.size() != 0);
623 const auto num_successors = node->num_successors();
627 if(
auto index=node->_work.index(); index == 1) {
628 if(node->_module !=
nullptr) {
629 bool first_time = !node->is_spawned();
630 _invoke_static_work(me, node);
636 _invoke_static_work(me, node);
640 else if (index == 2){
643 if(!node->is_spawned()) {
644 if(node->_subgraph) {
645 node->_subgraph->clear();
648 node->_subgraph.emplace();
652 Subflow fb(*(node->_subgraph));
654 _invoke_dynamic_work(me, node, fb);
657 if(!node->is_spawned()) {
659 if(!node->_subgraph->empty()) {
661 PassiveVector<Node*> src;
662 for(
auto& n: node->_subgraph->nodes()) {
663 n->_topology = node->_topology;
665 if(n->num_successors() == 0) {
667 node->_topology->_num_sinks++;
673 if(n->num_dependents() == 0) {
674 src.push_back(n.get());
691 if(!node->is_subtask()) {
694 if(node->_work.index() == 2 && !node->_subgraph->empty()) {
695 while(!node->_dependents.empty() && node->_dependents.back()->is_subtask()) {
696 node->_dependents.pop_back();
699 node->_num_dependents =
static_cast<int>(node->_dependents.size());
700 node->unset_spawned();
704 Node* cache {
nullptr};
706 for(
size_t i=0; i<num_successors; ++i) {
707 if(--(node->_successors[i]->_num_dependents) == 0) {
709 _schedule(cache,
false);
711 cache = node->_successors[i];
716 _schedule(cache,
true);
720 if(num_successors == 0) {
721 if(--(node->_topology->_num_sinks) == 0) {
722 _tear_down_topology(node->_topology);
728 inline void Executor::_invoke_static_work(
unsigned me, Node* node) {
730 _observer->on_entry(me,
TaskView(node));
731 std::invoke(std::get<Node::StaticWork>(node->_work));
732 _observer->on_exit(me,
TaskView(node));
735 std::invoke(std::get<Node::StaticWork>(node->_work));
740 inline void Executor::_invoke_dynamic_work(
unsigned me, Node* node,
Subflow& sf) {
742 _observer->on_entry(me,
TaskView(node));
743 std::invoke(std::get<Node::DynamicWork>(node->_work), sf);
744 _observer->on_exit(me,
TaskView(node));
747 std::invoke(std::get<Node::DynamicWork>(node->_work), sf);
753 return run_n(f, 1, [](){});
757 template <
typename C>
759 static_assert(std::is_invocable<C>::value);
760 return run_n(f, 1, std::forward<C>(c));
765 return run_n(f, repeat, [](){});
769 template <
typename C>
771 return run_until(f, [repeat]()
mutable {
return repeat-- == 0; }, std::forward<C>(c));
777 return run_until(f, std::forward<P>(pred), [](){});
781 inline void Executor::_tear_down_topology(Topology* tpg) {
783 auto &f = tpg->_taskflow;
788 if(!std::invoke(tpg->_pred)) {
789 tpg->_recover_num_sinks();
790 _schedule(tpg->_sources);
795 if(tpg->_call !=
nullptr) {
796 std::invoke(tpg->_call);
802 if(f._topologies.size() > 1) {
805 tpg->_promise.set_value();
806 f._topologies.pop_front();
810 _decrement_topology();
812 f._topologies.front()._bind(f._graph);
813 _schedule(f._topologies.front()._sources);
816 assert(f._topologies.size() == 1);
820 auto p {std::move(tpg->_promise)};
822 f._topologies.pop_front();
829 _decrement_topology_and_notify();
835 template <
typename P,
typename C>
839 static_assert(std::is_invocable_v<C> && std::is_invocable_v<P>);
841 _increment_topology();
844 if(f.
empty() || std::invoke(pred)) {
847 _decrement_topology_and_notify();
848 return promise.get_future();
851 if(_workers.size() == 0) {
852 TF_THROW(Error::EXECUTOR,
"no workers to execute the graph");
889 bool run_now {
false};
894 std::scoped_lock lock(f._mtx);
897 tpg = &(f._topologies.emplace_back(f, std::forward<P>(pred), std::forward<C>(c)));
898 future = tpg->_promise.get_future();
900 if(f._topologies.size() == 1) {
910 tpg->_bind(f._graph);
911 _schedule(tpg->_sources);
918 inline void Executor::_increment_topology() {
919 std::scoped_lock lock(_topology_mutex);
924 inline void Executor::_decrement_topology_and_notify() {
925 std::scoped_lock lock(_topology_mutex);
926 if(--_num_topologies == 0) {
927 _topology_cv.notify_all();
932 inline void Executor::_decrement_topology() {
933 std::scoped_lock lock(_topology_mutex);
940 _topology_cv.wait(lock, [&](){
return _num_topologies == 0; });
std::future< void > run(Taskflow &taskflow)
runs the taskflow once
Definition: executor.hpp:752
void remove_observer()
removes the associated observer
Definition: executor.hpp:498
bool empty() const noexcept
queries if the queue is empty at the time of this call
Definition: spmc_queue.hpp:172
std::future< void > run_until(Taskflow &taskflow, P &&pred)
runs the taskflow multiple times until the predicate becomes true and then invokes a callback ...
Definition: executor.hpp:776
~Executor()
destructs the executor
Definition: executor.hpp:266
void push(O &&item)
inserts an item to the queue
Definition: spmc_queue.hpp:189
Definition: taskflow.hpp:5
T hardware_concurrency(T... args)
bool detached() const
queries if the subflow will be detached from its parent task
Definition: flow_builder.hpp:869
Observer * make_observer(Args &&... args)
constructs an observer to inspect the activities of worker threads
Definition: executor.hpp:489
the class to create a task dependency graph
Definition: core/taskflow.hpp:15
A constant wrapper class to a task node, mainly used in the tf::ExecutorObserver interface.
Definition: task.hpp:370
bool empty() const
queries the emptiness of the taskflow
Definition: core/taskflow.hpp:122
bool joined() const
queries if the subflow will join its parent task
Definition: flow_builder.hpp:874
Lock-free unbounded single-producer multiple-consumer queue.
Definition: spmc_queue.hpp:29
The executor class to run a taskflow graph.
Definition: executor.hpp:88
size_t num_workers() const
queries the number of worker threads (can be zero)
Definition: executor.hpp:281
std::optional< T > steal()
steals an item from the queue
Definition: spmc_queue.hpp:239
Executor(unsigned n=std::thread::hardware_concurrency())
constructs the executor with N worker threads
Definition: executor.hpp:258
The building blocks of dynamic tasking.
Definition: flow_builder.hpp:817
std::future< void > run_n(Taskflow &taskflow, size_t N)
runs the taskflow for N times
Definition: executor.hpp:764
void wait_for_all()
wait for all pending graphs to complete
Definition: executor.hpp:938