-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathconcurrent_queue.h
More file actions
119 lines (100 loc) · 2.92 KB
/
Copy pathconcurrent_queue.h
File metadata and controls
119 lines (100 loc) · 2.92 KB
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
/*
* Copyright (C) 2018 IRT GmbH
*
* Author:
* Fabian Sattler
*
* This file is a part of IRT DAB library.
*
* This library is free software; you can redistribute it and/or
* modify it under the terms of the GNU Lesser General Public
* License as published by the Free Software Foundation; either
* version 2.1 of the License, or (at your option) any later version.
*
* This library is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
* Lesser General Public License for more details.
*
*/
#ifndef CONCURRENT_QUEUE
#define CONCURRENT_QUEUE
#include <queue>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <chrono>
template <typename T>
class ConcurrentQueue {
public:
T pop() {
std::unique_lock<std::mutex> mlock(m_mutex);
//effective against spurious wakes
while (m_queue.empty()) {
m_cond.wait(mlock);
}
auto item = m_queue.front();
m_queue.pop();
return item;
}
void pop(T& item) {
std::unique_lock<std::mutex> mlock(m_mutex);
while (m_queue.empty()) {
m_cond.wait(mlock);
}
item = m_queue.front();
m_queue.pop();
}
//tries to pop an element but unlocks after timeout
bool tryPop(T& item, std::chrono::milliseconds timeout) {
std::unique_lock<std::mutex> mlock(m_mutex);
//lambda inside for predicate
if(!m_cond.wait_for(mlock, timeout, [this] {return !m_queue.empty(); })) {
return false;
}
item = m_queue.front();
m_queue.pop();
return true;
}
void push(const T& item) {
std::unique_lock<std::mutex> mlock(m_mutex);
m_queue.push(item);
//unlock manually before notifying
mlock.unlock();
m_cond.notify_one();
}
void push(T&& item) {
//the item will be std::moved to the queue
std::unique_lock<std::mutex> mlock(m_mutex);
m_queue.push(std::move(item));
//unlock manually before notifying
mlock.unlock();
m_cond.notify_one();
}
int getSize() {
std::unique_lock<std::mutex> mlock(m_mutex);
return m_queue.size();
}
void getSize(int& size) {
std::unique_lock<std::mutex> mlock(m_mutex);
size = m_queue.size();
mlock.unlock();
}
bool isEmpty() {
std::unique_lock<std::mutex> mlock(m_mutex);
return m_queue.empty();
}
void clear() {
std::unique_lock<std::mutex> mlock(m_mutex);
std::queue<T> empty;
std::swap(m_queue, empty);
mlock.unlock();
//make sure every waiting thread gets its fair share
m_cond.notify_all();
}
public:
std::queue<T> m_queue;
std::mutex m_mutex;
std::condition_variable m_cond;
};
#endif // CONCURRENT_QUEUE