The Pedigree Project 0.1
Mailbox.h
1/*
2 * Copyright (c) 2008-2014, Pedigree Developers
3 *
4 * Please see the CONTRIB file in the root of the source tree for a full
5 * list of contributors.
6 *
7 * Permission to use, copy, modify, and distribute this software for any
8 * purpose with or without fee is hereby granted, provided that the above
9 * copyright notice and this permission notice appear in all copies.
10 *
11 * THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
12 * WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
13 * MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
14 * ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
15 * WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
16 * ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
17 * OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
18 */
19
20#ifndef PEDIGREE_KERNEL_UTILITIES_MAILBOX_H
21#define PEDIGREE_KERNEL_UTILITIES_MAILBOX_H
22#include "pedigree/kernel/LockGuard.h"
23#include "pedigree/kernel/Log.h"
24#include "pedigree/kernel/compiler.h"
25#include "pedigree/kernel/process/ConditionVariable.h"
26#include "pedigree/kernel/process/Mutex.h"
27#include "pedigree/kernel/process/TerminationDeferral.h"
28#include "pedigree/kernel/processor/types.h"
29#include "pedigree/kernel/time/Time.h"
30#include "pedigree/kernel/utilities/RingQueue.h"
31#include "pedigree/kernel/utilities/assert.h"
32#include "pedigree/kernel/utilities/new"
33
34#include <config.h>
35
36#if THREADS
37#include "pedigree/kernel/process/Thread.h"
38#include "pedigree/kernel/processor/Processor.h"
39#include "pedigree/kernel/processor/ProcessorInformation.h"
40#endif
41
42namespace MailboxWait {
43enum WaitType { Reading, Writing };
44}
45
47 void changed(bool, bool) {}
48 void closed() {}
49};
50
56template <class T, size_t preallocatedSize = 0, class Notifications = MailboxNotifications>
57class EXPORTED_PUBLIC Mailbox {
58 protected:
59 // Owns admission and the mutex together. Condition waits temporarily release
60 // the mutex, but close() must still drain the admitted operation.
62 public:
63 explicit ActiveOperation(Mailbox& buffer)
64 : m_TerminationDeferral(), m_Buffer(buffer.beginOperation() ? &buffer : nullptr) {}
65
67 if (m_Buffer) {
68 m_Buffer->endOperation();
69 }
70 }
71
72 explicit operator bool() const {
73 return m_Buffer != nullptr;
74 }
75
76 private:
77 ActiveOperation(const ActiveOperation&) = delete;
78 ActiveOperation& operator=(const ActiveOperation&) = delete;
79
80 TerminationDeferral m_TerminationDeferral;
81 Mailbox* m_Buffer;
82 };
83
84 public:
85 enum Error {
86 NoError,
87
88 // Mailbox is empty and a zero timeout was specified
89 Empty,
90
91 // A nonblocking operation could not acquire the lock or found no space
92 WouldBlock,
93
94 // ConditionVariable failure modes
95 TimedOut,
96 Interrupted,
97 ThreadTerminating,
98
99 // Mailbox has closed and will not admit another operation.
100 Closed,
101 };
102
103 explicit Mailbox(size_t ringSize) : m_Ring(ringSize) {}
104
107 close();
108 }
109
114 void close() {
115 TerminationDeferral terminationDeferral;
116 m_Lock.acquire();
117 if (!m_Closing) {
118 m_Closing = true;
119 m_WriteClosed = true;
120 if (m_ReadWaiters) {
121 m_ReadCondition.broadcast();
122 }
123 if (m_WriteWaiters) {
124 m_WriteCondition.broadcast();
125 }
126 m_Notifications.closed();
127 }
128
129 while (m_ActiveOperations) {
130 m_DrainCondition.waitForCompletion(m_Lock);
131 }
132 m_Lock.release();
133 }
134
143 out = T();
144 LockGuard<Mutex> guard(m_Lock);
145 if (!m_Closing || m_ActiveOperations || !m_Ring.count()) {
146 return false;
147 }
148
149 popLocked(&out, 1);
150 return true;
151 }
152
154 Error write(const T& obj, Time::Timestamp& timeout) {
155 ActiveOperation operation(*this);
156 if (!operation) {
157 return Closed;
158 }
159
160 const Error error = waitForLocked(MailboxWait::Writing, timeout);
161 if (error == NoError) {
162 pushLocked(&obj, 1);
163 }
164 return error;
165 }
166
167 Error write(const T& obj) {
168 Time::Timestamp timeout = Time::Infinity;
169 return write(obj, timeout);
170 }
171
176 bool closeWritesWithFinal(const T& obj) {
177 ActiveOperation operation(*this);
178 if (!operation) {
179 return false;
180 }
181
182 if (m_WriteClosed) {
183 return false;
184 }
185
186 m_WriteClosed = true;
187 m_FinalPending = true;
188 if (m_WriteWaiters) {
189 m_WriteCondition.broadcast();
190 }
191 while (!m_Closing && m_Ring.count() >= m_Ring.capacity()) {
192 ++m_WriteWaiters;
193 m_WriteCondition.waitForCompletion(m_Lock);
194 --m_WriteWaiters;
195 }
196
197 if (m_Closing) {
198 return false;
199 }
200
201 pushLocked(&obj, 1);
202 m_FinalPending = false;
203 if (m_ReadWaiters) {
204 m_ReadCondition.broadcast();
205 }
206 return true;
207 }
208
213 Error tryWrite(const T& obj) {
214#if THREADS
215 // Waking a blocked reader enters scheduler locks and is not IRQ-safe.
217 return WouldBlock;
218 }
219#endif
220
221 TerminationDeferral terminationDeferral;
222 if (!m_Lock.tryAcquire()) {
223 return WouldBlock;
224 }
225
226 if (m_Closing || m_WriteClosed) {
227 m_Lock.release();
228 return Closed;
229 }
230
231 if (m_Ring.count() >= m_Ring.capacity()) {
232 m_Lock.release();
233 return WouldBlock;
234 }
235
236 pushLocked(&obj, 1);
237 m_Lock.release();
238 return NoError;
239 }
240
242 size_t write(const T* obj, size_t n, Time::Timestamp& timeout) {
243 ActiveOperation operation(*this);
244 if (!operation) {
245 return 0;
246 }
247
248 if (n > m_Ring.capacity()) {
249 n = m_Ring.capacity();
250 }
251
252 size_t written = 0;
253 while (written < n) {
254 if (waitForLocked(MailboxWait::Writing, timeout) != NoError) {
255 break;
256 }
257 written += pushLocked(obj + written, n - written);
258 }
259 return written;
260 }
261
262 size_t write(const T* obj, size_t n) {
263 Time::Timestamp timeout = Time::Infinity;
264 return write(obj, n, timeout);
265 }
266
276 MUST_USE_RESULT bool read(T& out, Time::Timestamp& timeout, Error& error) {
277 ActiveOperation operation(*this);
278 if (!operation) {
279 out = T();
280 error = Closed;
281 return false;
282 }
283
284 out = T();
285 if (!timeout && !m_Ring.count() && !readsClosed()) {
286 error = Empty;
287 return false;
288 }
289 error = waitForLocked(MailboxWait::Reading, timeout);
290 if (error != NoError) {
291 return false;
292 }
293 popLocked(&out, 1);
294 return true;
295 }
296
297 MUST_USE_RESULT bool read(T& out, Error& error) {
298 Time::Timestamp timeout = 0;
299 return read(out, timeout, error);
300 }
301
303 size_t read(T* out, size_t n, Time::Timestamp& timeout) {
304 ActiveOperation operation(*this);
305 if (!operation) {
306 return 0;
307 }
308
309 if (n > m_Ring.capacity()) {
310 n = m_Ring.capacity();
311 }
312
313 size_t read = 0;
314 while (read < n && timeout > 0) {
315 if (waitForLocked(MailboxWait::Reading, timeout) != NoError) {
316 out[read] = T();
317 break;
318 }
319 read += popLocked(out + read, n - read);
320 }
321 return read;
322 }
323
324 size_t read(T* out, size_t n) {
325 Time::Timestamp timeout = Time::Infinity;
326 return read(out, n, timeout);
327 }
328
330 bool dataReady() {
331 ActiveOperation operation(*this);
332 if (!operation) {
333 return false;
334 }
335
336 return m_Ring.count() > 0;
337 }
338
340 bool canWrite() {
341 ActiveOperation operation(*this);
342 if (!operation) {
343 return false;
344 }
345
346 return !m_Closing && !m_WriteClosed && m_Ring.count() < m_Ring.capacity();
347 }
348
350 MUST_USE_RESULT bool waitFor(MailboxWait::WaitType wait, Time::Timestamp& timeout, Error& error) {
351 error = NoError;
352 ActiveOperation operation(*this);
353 if (!operation) {
354 error = Closed;
355 return false;
356 }
357
358 error = waitForLocked(wait, timeout);
359 return error == NoError;
360 }
361
362 bool waitFor(MailboxWait::WaitType wait, Time::Timestamp& timeout) {
363 Error error = NoError;
364 return waitFor(wait, timeout, error);
365 }
366
367 bool waitFor(MailboxWait::WaitType wait) {
368 Time::Timestamp timeout = Time::Infinity;
369 return waitFor(wait, timeout);
370 }
371
372 private:
373 Error waitForLocked(MailboxWait::WaitType direction, Time::Timestamp& timeout) {
374 const bool writing = direction == MailboxWait::Writing;
375 while (true) {
376 if (writing ? (m_Closing || m_WriteClosed) : readsClosed()) {
377 return Closed;
378 }
379 if (writing ? m_Ring.count() < m_Ring.capacity() : m_Ring.count() != 0) {
380 return NoError;
381 }
382
383 ConditionVariable::Error error = ConditionVariable::NoError;
384 if (!waitForChange(writing ? m_WriteCondition : m_ReadCondition,
385 writing ? m_WriteWaiters : m_ReadWaiters, timeout, error)) {
386 return errorFromConditionVariable(error);
387 }
388 }
389 }
390
391 size_t pushLocked(const T* items, size_t count) {
392 const bool wasEmpty = !m_Ring.count();
393 const size_t written = m_Ring.push(items, count);
394 if (written) {
395 m_Notifications.changed(wasEmpty, false);
396 if (m_ReadWaiters) {
397 if (written == 1) {
398 m_ReadCondition.signal();
399 } else {
400 m_ReadCondition.broadcast();
401 }
402 }
403 }
404 return written;
405 }
406
407 size_t popLocked(T* items, size_t count) {
408 const bool wasFull = m_Ring.count() == m_Ring.capacity();
409 const size_t read = m_Ring.pop(items, count);
410 if (read) {
411 m_Notifications.changed(false, wasFull);
412 if (m_WriteWaiters) {
413 if (read == 1) {
414 m_WriteCondition.signal();
415 } else {
416 m_WriteCondition.broadcast();
417 }
418 }
419 }
420 return read;
421 }
422
423 bool waitForChange(ConditionVariable& condition, size_t& waiters, Time::Timestamp& timeout,
424 ConditionVariable::Error& error) {
425 ++waiters;
426 const bool result = condition.wait(m_Lock, timeout, error);
427 --waiters;
428 return result;
429 }
430
431 bool beginOperation() {
432 m_Lock.acquire();
433 if (m_Closing) {
434 m_Lock.release();
435 return false;
436 }
437
438 ++m_ActiveOperations;
439 return true;
440 }
441
442 void endOperation() {
443 assert(m_ActiveOperations);
444 --m_ActiveOperations;
445 if (m_Closing && !m_ActiveOperations) {
446 m_DrainCondition.broadcast();
447 }
448 m_Lock.release();
449 }
450
451 Error errorFromConditionVariable(ConditionVariable::Error err) {
452 switch (err) {
453 case ConditionVariable::TimedOut:
454 return TimedOut;
455 case ConditionVariable::Interrupted:
456#if THREADS
457 // ConditionVariable consumes the detailed reason as part of
458 // its result contract. Preserve it for callers, such as
459 // lwIP, whose compatibility API collapses all failures to a
460 // timeout sentinel.
461 Processor::information().getCurrentThread()->setInterruptionReason(
462 Thread::InterruptedBySignal);
463#endif
464 return Interrupted;
465 case ConditionVariable::TerminationDeferred:
466 return ThreadTerminating;
467 default:
468 FATAL("invalid ConditionVariable::Error enum value for Mailbox");
469 }
470
471 return NoError;
472 }
473
474 protected:
475 bool readsClosed() const {
476 return m_Closing || (m_WriteClosed && !m_FinalPending && !m_Ring.count());
477 }
478
480 [[no_unique_address]] Notifications m_Notifications;
481 bool m_Closing = false;
482 bool m_WriteClosed = false;
483 // Readers must not observe EOF while the final writer is waiting for space.
484 bool m_FinalPending = false;
485
486 private:
487 Mutex m_Lock;
488 ConditionVariable m_WriteCondition;
489 ConditionVariable m_ReadCondition;
490 ConditionVariable m_DrainCondition;
491 size_t m_ActiveOperations = 0;
492 // Protected by m_Lock, including enrolment and return from condition waits.
493 size_t m_ReadWaiters = 0;
494 size_t m_WriteWaiters = 0;
495};
496
497#endif
MUST_USE_RESULT bool wait(Mutex &mutex, Time::Timestamp &timeout, Error &error, WaitQueue::StackDiscardCleanup onStackDiscard=nullptr, void *stackDiscardContext=nullptr)
void close()
Definition Mailbox.h:114
MUST_USE_RESULT bool waitFor(MailboxWait::WaitType wait, Time::Timestamp &timeout, Error &error)
waitFor - block until the given condition is true (readable/writeable)
Definition Mailbox.h:350
bool canWrite()
canWrite - is it possible to write to the ring buffer without blocking?
Definition Mailbox.h:340
MUST_USE_RESULT bool takeAfterClose(T &out)
Definition Mailbox.h:142
size_t write(const T *obj, size_t n, Time::Timestamp &timeout)
Publishes each available chunk before waiting for more space.
Definition Mailbox.h:242
bool closeWritesWithFinal(const T &obj)
Definition Mailbox.h:176
MUST_USE_RESULT bool read(T &out, Time::Timestamp &timeout, Error &error)
Definition Mailbox.h:276
Error tryWrite(const T &obj)
Definition Mailbox.h:213
Error write(const T &obj, Time::Timestamp &timeout)
Writes one object, waiting for space up to the supplied timeout.
Definition Mailbox.h:154
size_t read(T *out, size_t n, Time::Timestamp &timeout)
Removes each available chunk before waiting for more data.
Definition Mailbox.h:303
bool dataReady()
dataReady - is data ready for reading from the ring buffer?
Definition Mailbox.h:330
~Mailbox()
Destructor - closes the ring and drains every admitted operation.
Definition Mailbox.h:106
Definition Mutex.h:56
static ProcessorInformation & information()
static bool inDeviceHardIrq()
Definition Processor.h:581
#define assert(x)
Definition assert.h:39