-
-
Notifications
You must be signed in to change notification settings - Fork 0
/
task_node.hpp
127 lines (91 loc) · 2.45 KB
/
task_node.hpp
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
// === (C) 2020-2024 === parallel_f (tasks, queues, lists in parallel threads)
// Written by Denis Oliver Kropp <[email protected]>
#pragma once
#include <any>
#include <map>
#include <memory>
#include <mutex>
#include <queue>
#include <tuple>
#include "log.hpp"
#include "stats.hpp"
#include "system.hpp"
#include "vthread.hpp"
#include "Event.hxx"
#include "task.hpp"
#include "joinable.hpp"
namespace parallel_f {
class task_node : public lli::EventListener
{
private:
std::shared_ptr<task_base> task;
unsigned int wait;
bool managed;
std::shared_ptr<vthread> thread;
std::mutex lock;
std::condition_variable cond;
bool finished;
public:
task_node(std::string name, std::shared_ptr<task_base> task, unsigned int wait, bool managed = true)
:
task(task),
wait(wait),
managed(managed),
finished(false)
{
LOG_DEBUG("task_node::task_node(%p, '%s', %u)\n", this, name.c_str(), wait);
thread = std::make_shared<vthread>(name);
task->finished.Attach(this, [this](int) {
std::unique_lock<std::mutex> l(lock);
finished = true;
cond.notify_all();
});
}
~task_node()
{
LOG_DEBUG("task_node::~task_node(%p '%s')\n", this, get_name().c_str());
}
void add_to_notify(std::shared_ptr<task_node> node)
{
LOG_DEBUG("task_node::add_to_notify(%p '%s', %p)\n", this, get_name().c_str(), node.get());
if (finished)
throw std::runtime_error("finished");
task->finished.Attach(node.get(), [node](int) {
node->notify();
});
}
void notify()
{
LOG_DEBUG("task_node::notify(%p '%s')...\n", this, get_name().c_str());
std::unique_lock<std::mutex> l(lock);
LOG_DEBUG("task_node::notify(%p '%s') wait count %u -> %u\n", this, get_name().c_str(), wait, wait-1);
if (!wait)
throw std::runtime_error("zero wait count");
if (!--wait) {
std::shared_ptr<task_base> t = task;
thread->start([t]() {
bool finished = t->finish();
}, managed);
}
LOG_DEBUG("task_node::notify(%p '%s') done.\n", this, get_name().c_str());
}
void join()
{
LOG_DEBUG("task_node::join(%p '%s')...\n", this, get_name().c_str());
std::unique_lock<std::mutex> l(lock);
while (!finished) {
LOG_DEBUG("task_node::join(%p '%s') waiting...\n", this, get_name().c_str());
vthread::wait(cond, l);
}
LOG_DEBUG("task_node::join(%p '%s') done.\n", this, get_name().c_str());
}
std::thread::id get_thread_id()
{
return thread->get_id();
}
std::string get_name()
{
return thread->get_name();
}
};
}