64 : m_TerminationDeferral(), m_Buffer(buffer.beginOperation() ? &buffer :
nullptr) {}
68 m_Buffer->endOperation();
72 explicit operator bool()
const {
73 return m_Buffer !=
nullptr;
103 explicit Mailbox(
size_t ringSize) : m_Ring(ringSize) {}
119 m_WriteClosed =
true;
121 m_ReadCondition.broadcast();
123 if (m_WriteWaiters) {
124 m_WriteCondition.broadcast();
126 m_Notifications.closed();
129 while (m_ActiveOperations) {
130 m_DrainCondition.waitForCompletion(m_Lock);
145 if (!m_Closing || m_ActiveOperations || !m_Ring.count()) {
154 Error
write(
const T& obj, Time::Timestamp& timeout) {
160 const Error error = waitForLocked(MailboxWait::Writing, timeout);
161 if (error == NoError) {
167 Error write(
const T& obj) {
168 Time::Timestamp timeout = Time::Infinity;
169 return write(obj, timeout);
186 m_WriteClosed =
true;
187 m_FinalPending =
true;
188 if (m_WriteWaiters) {
189 m_WriteCondition.broadcast();
191 while (!m_Closing && m_Ring.count() >= m_Ring.capacity()) {
193 m_WriteCondition.waitForCompletion(m_Lock);
202 m_FinalPending =
false;
204 m_ReadCondition.broadcast();
222 if (!m_Lock.tryAcquire()) {
226 if (m_Closing || m_WriteClosed) {
231 if (m_Ring.count() >= m_Ring.capacity()) {
242 size_t write(
const T* obj,
size_t n, Time::Timestamp& timeout) {
248 if (n > m_Ring.capacity()) {
249 n = m_Ring.capacity();
253 while (written < n) {
254 if (waitForLocked(MailboxWait::Writing, timeout) != NoError) {
257 written += pushLocked(obj + written, n - written);
262 size_t write(
const T* obj,
size_t n) {
263 Time::Timestamp timeout = Time::Infinity;
264 return write(obj, n, timeout);
285 if (!timeout && !m_Ring.count() && !readsClosed()) {
289 error = waitForLocked(MailboxWait::Reading, timeout);
290 if (error != NoError) {
298 Time::Timestamp timeout = 0;
299 return read(out, timeout, error);
303 size_t read(T* out,
size_t n, Time::Timestamp& timeout) {
309 if (n > m_Ring.capacity()) {
310 n = m_Ring.capacity();
314 while (read < n && timeout > 0) {
315 if (waitForLocked(MailboxWait::Reading, timeout) != NoError) {
319 read += popLocked(out + read, n - read);
324 size_t read(T* out,
size_t n) {
325 Time::Timestamp timeout = Time::Infinity;
326 return read(out, n, timeout);
336 return m_Ring.count() > 0;
346 return !m_Closing && !m_WriteClosed && m_Ring.count() < m_Ring.capacity();
358 error = waitForLocked(wait, timeout);
359 return error == NoError;
362 bool waitFor(MailboxWait::WaitType wait, Time::Timestamp& timeout) {
363 Error error = NoError;
364 return waitFor(wait, timeout, error);
367 bool waitFor(MailboxWait::WaitType wait) {
368 Time::Timestamp timeout = Time::Infinity;
369 return waitFor(wait, timeout);
373 Error waitForLocked(MailboxWait::WaitType direction, Time::Timestamp& timeout) {
374 const bool writing = direction == MailboxWait::Writing;
376 if (writing ? (m_Closing || m_WriteClosed) : readsClosed()) {
379 if (writing ? m_Ring.count() < m_Ring.capacity() : m_Ring.count() != 0) {
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);
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);
395 m_Notifications.changed(wasEmpty,
false);
398 m_ReadCondition.signal();
400 m_ReadCondition.broadcast();
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);
411 m_Notifications.changed(
false, wasFull);
412 if (m_WriteWaiters) {
414 m_WriteCondition.signal();
416 m_WriteCondition.broadcast();
423 bool waitForChange(
ConditionVariable& condition,
size_t& waiters, Time::Timestamp& timeout,
424 ConditionVariable::Error& error) {
426 const bool result = condition.
wait(m_Lock, timeout, error);
431 bool beginOperation() {
438 ++m_ActiveOperations;
442 void endOperation() {
443 assert(m_ActiveOperations);
444 --m_ActiveOperations;
445 if (m_Closing && !m_ActiveOperations) {
446 m_DrainCondition.broadcast();
451 Error errorFromConditionVariable(ConditionVariable::Error err) {
453 case ConditionVariable::TimedOut:
455 case ConditionVariable::Interrupted:
462 Thread::InterruptedBySignal);
465 case ConditionVariable::TerminationDeferred:
466 return ThreadTerminating;
468 FATAL(
"invalid ConditionVariable::Error enum value for Mailbox");
475 bool readsClosed()
const {
476 return m_Closing || (m_WriteClosed && !m_FinalPending && !m_Ring.count());
480 [[no_unique_address]] Notifications m_Notifications;
481 bool m_Closing =
false;
482 bool m_WriteClosed =
false;
484 bool m_FinalPending =
false;
491 size_t m_ActiveOperations = 0;
493 size_t m_ReadWaiters = 0;
494 size_t m_WriteWaiters = 0;