scheduler.cpp raw
1 // Copyright (c) 2015-present The Bitcoin Core developers
2 // Distributed under the MIT software license, see the accompanying
3 // file COPYING or http://www.opensource.org/licenses/mit-license.php.
4
5 #include <scheduler.h>
6
7 #include <sync.h>
8 #include <util/time.h>
9
10 #include <cassert>
11 #include <functional>
12 #include <utility>
13
14 CScheduler::CScheduler() = default;
15
16 CScheduler::~CScheduler()
17 {
18 assert(nThreadsServicingQueue == 0);
19 if (stopWhenEmpty) assert(taskQueue.empty());
20 }
21
22
23 void CScheduler::serviceQueue()
24 {
25 WAIT_LOCK(newTaskMutex, lock);
26 ++nThreadsServicingQueue;
27
28 // newTaskMutex is locked throughout this loop EXCEPT
29 // when the thread is waiting or when the user's function
30 // is called.
31 while (!shouldStop()) {
32 try {
33 while (!shouldStop() && taskQueue.empty()) {
34 // Wait until there is something to do.
35 newTaskScheduled.wait(lock);
36 }
37
38 // Wait until either there is a new task, or until
39 // the time of the first item on the queue:
40
41 while (!shouldStop() && !taskQueue.empty()) {
42 std::chrono::steady_clock::time_point timeToWaitFor = taskQueue.begin()->first;
43 if (newTaskScheduled.wait_until(lock, timeToWaitFor) == std::cv_status::timeout) {
44 break; // Exit loop after timeout, it means we reached the time of the event
45 }
46 }
47
48 // If there are multiple threads, the queue can empty while we're waiting (another
49 // thread may service the task we were waiting on).
50 if (shouldStop() || taskQueue.empty())
51 continue;
52
53 Function f = taskQueue.begin()->second;
54 taskQueue.erase(taskQueue.begin());
55
56 {
57 // Unlock before calling f, so it can reschedule itself or another task
58 // without deadlocking:
59 REVERSE_LOCK(lock, newTaskMutex);
60 f();
61 }
62 } catch (...) {
63 --nThreadsServicingQueue;
64 throw;
65 }
66 }
67 --nThreadsServicingQueue;
68 newTaskScheduled.notify_one();
69 }
70
71 void CScheduler::schedule(CScheduler::Function f, std::chrono::steady_clock::time_point t)
72 {
73 {
74 LOCK(newTaskMutex);
75 taskQueue.insert(std::make_pair(t, f));
76 }
77 newTaskScheduled.notify_one();
78 }
79
80 void CScheduler::MockForward(std::chrono::seconds delta_seconds)
81 {
82 assert(delta_seconds > 0s && delta_seconds <= 1h);
83
84 {
85 LOCK(newTaskMutex);
86
87 // use temp_queue to maintain updated schedule
88 std::multimap<std::chrono::steady_clock::time_point, Function> temp_queue;
89
90 for (const auto& element : taskQueue) {
91 temp_queue.emplace_hint(temp_queue.cend(), element.first - delta_seconds, element.second);
92 }
93
94 // point taskQueue to temp_queue
95 taskQueue = std::move(temp_queue);
96 }
97
98 // notify that the taskQueue needs to be processed
99 newTaskScheduled.notify_one();
100 }
101
102 static void Repeat(CScheduler& s, CScheduler::Function f, std::chrono::milliseconds delta)
103 {
104 f();
105 s.scheduleFromNow([=, &s] { Repeat(s, f, delta); }, delta);
106 }
107
108 void CScheduler::scheduleEvery(CScheduler::Function f, std::chrono::milliseconds delta)
109 {
110 scheduleFromNow([this, f, delta] { Repeat(*this, f, delta); }, delta);
111 }
112
113 size_t CScheduler::getQueueInfo(std::chrono::steady_clock::time_point& first,
114 std::chrono::steady_clock::time_point& last) const
115 {
116 LOCK(newTaskMutex);
117 size_t result = taskQueue.size();
118 if (!taskQueue.empty()) {
119 first = taskQueue.begin()->first;
120 last = taskQueue.rbegin()->first;
121 }
122 return result;
123 }
124
125 bool CScheduler::AreThreadsServicingQueue() const
126 {
127 LOCK(newTaskMutex);
128 return nThreadsServicingQueue;
129 }
130
131
132 void SerialTaskRunner::MaybeScheduleProcessQueue()
133 {
134 {
135 LOCK(m_callbacks_mutex);
136 // Try to avoid scheduling too many copies here, but if we
137 // accidentally have two ProcessQueue's scheduled at once its
138 // not a big deal.
139 if (m_are_callbacks_running) return;
140 if (m_callbacks_pending.empty()) return;
141 }
142 m_scheduler.schedule([this] { this->ProcessQueue(); }, std::chrono::steady_clock::now());
143 }
144
145 void SerialTaskRunner::ProcessQueue()
146 {
147 std::function<void()> callback;
148 {
149 LOCK(m_callbacks_mutex);
150 if (m_are_callbacks_running) return;
151 if (m_callbacks_pending.empty()) return;
152 m_are_callbacks_running = true;
153
154 callback = std::move(m_callbacks_pending.front());
155 m_callbacks_pending.pop_front();
156 }
157
158 // RAII the setting of fCallbacksRunning and calling MaybeScheduleProcessQueue
159 // to ensure both happen safely even if callback() throws.
160 struct RAIICallbacksRunning {
161 SerialTaskRunner* instance;
162 explicit RAIICallbacksRunning(SerialTaskRunner* _instance) : instance(_instance) {}
163 ~RAIICallbacksRunning()
164 {
165 {
166 LOCK(instance->m_callbacks_mutex);
167 instance->m_are_callbacks_running = false;
168 }
169 instance->MaybeScheduleProcessQueue();
170 }
171 } raiicallbacksrunning(this);
172
173 callback();
174 }
175
176 void SerialTaskRunner::insert(std::function<void()> func)
177 {
178 {
179 LOCK(m_callbacks_mutex);
180 m_callbacks_pending.emplace_back(std::move(func));
181 }
182 MaybeScheduleProcessQueue();
183 }
184
185 void SerialTaskRunner::flush()
186 {
187 assert(!m_scheduler.AreThreadsServicingQueue());
188 bool should_continue = true;
189 while (should_continue) {
190 ProcessQueue();
191 LOCK(m_callbacks_mutex);
192 should_continue = !m_callbacks_pending.empty();
193 }
194 }
195
196 size_t SerialTaskRunner::size()
197 {
198 LOCK(m_callbacks_mutex);
199 return m_callbacks_pending.size();
200 }
201