17#ifndef RESONANCE_AUDIO_UTILS_THREADSAFE_FIFO_H_
18#define RESONANCE_AUDIO_UTILS_THREADSAFE_FIFO_H_
22#include <condition_variable>
28#include "base/logging.h"
60 T* AcquireInputObject();
64 void ReleaseInputObject(
const T*
object);
70 T* AcquireOutputObject();
73 void ReleaseOutputObject(
const T*
object);
79 bool SleepUntilInputObjectIsAvailable()
const;
85 bool SleepUntilOutputObjectIsAvailable()
const;
89 void EnableBlockingSleepUntilMethods(
bool enable);
106 mutable std::mutex fifo_empty_mutex_;
107 mutable std::condition_variable fifo_empty_conditional_;
109 mutable std::mutex fifo_full_mutex_;
110 mutable std::condition_variable fifo_full_conditional_;
113 std::vector<T> fifo_;
118 std::atomic<size_t> fifo_size_;
120 std::atomic<bool> enable_sleeping_;
125 : fifo_(max_objects),
129 enable_sleeping_(true) {
130 CHECK_GT(max_objects, 0) <<
"FIFO size must be greater than zero";
134ThreadsafeFifo<T>::ThreadsafeFifo(
size_t max_objects,
const T& init)
135 : ThreadsafeFifo(max_objects) {
136 for (
auto&
object : fifo_) {
142T* ThreadsafeFifo<T>::AcquireInputObject() {
146 CHECK_LT(fifo_size_, fifo_.size());
149 return &fifo_[write_pos_];
153void ThreadsafeFifo<T>::ReleaseInputObject(
const T*
object) {
154 DCHECK_EQ(
object, &fifo_[write_pos_]);
157 write_pos_ = write_pos_ % fifo_.size();
158 if (fifo_size_.fetch_add(1) == 0) {
163 std::lock_guard<std::mutex> lock(fifo_empty_mutex_);
166 fifo_empty_conditional_.notify_one();
171T* ThreadsafeFifo<T>::AcquireOutputObject() {
175 CHECK_GT(fifo_size_, 0);
176 return &fifo_[read_pos_];
180void ThreadsafeFifo<T>::ReleaseOutputObject(
const T*
object) {
181 DCHECK_EQ(
object, &fifo_[read_pos_]);
184 read_pos_ = read_pos_ % fifo_.size();
186 if (fifo_size_.fetch_sub(1) == fifo_.size()) {
191 std::lock_guard<std::mutex> lock(fifo_full_mutex_);
194 fifo_full_conditional_.notify_one();
199bool ThreadsafeFifo<T>::SleepUntilInputObjectIsAvailable()
const {
202 std::unique_lock<std::mutex> lock(fifo_full_mutex_);
203 fifo_full_conditional_.wait(lock, [
this]() {
204 return fifo_size_.load() < fifo_.size() || !enable_sleeping_.load();
206 return fifo_size_.load() < fifo_.size();
210bool ThreadsafeFifo<T>::SleepUntilOutputObjectIsAvailable()
const {
212 std::unique_lock<std::mutex> lock(fifo_empty_mutex_);
213 fifo_empty_conditional_.wait(lock, [
this]() {
214 return fifo_size_.load() > 0 || !enable_sleeping_.load();
216 return fifo_size_.load() > 0;
220void ThreadsafeFifo<T>::EnableBlockingSleepUntilMethods(
bool enable) {
221 enable_sleeping_ = enable;
225 { std::lock_guard<std::mutex> lock(fifo_empty_mutex_); }
226 { std::lock_guard<std::mutex> lock(fifo_full_mutex_); }
227 fifo_empty_conditional_.notify_one();
228 fifo_full_conditional_.notify_one();
232size_t ThreadsafeFifo<T>::Size()
const {
233 return fifo_size_.load();
237bool ThreadsafeFifo<T>::Empty()
const {
238 return fifo_size_.load() == 0;
242bool ThreadsafeFifo<T>::Full()
const {
243 return fifo_size_.load() == fifo_.size();
247void ThreadsafeFifo<T>::Clear() {
249 T* output = AcquireOutputObject();
250 if (output !=
nullptr) {
251 ReleaseOutputObject(output);
Definition threadsafe_fifo.h:37