Directory: | cvmfs/ |
---|---|
File: | cvmfs/util/tube.h |
Date: | 2025-10-12 02:35:38 |
Exec | Total | Coverage | |
---|---|---|---|
Lines: | 130 | 130 | 100.0% |
Branches: | 50 | 74 | 67.6% |
Line | Branch | Exec | Source |
---|---|---|---|
1 | /** | ||
2 | * This file is part of the CernVM File System. | ||
3 | */ | ||
4 | |||
5 | #ifndef CVMFS_UTIL_TUBE_H_ | ||
6 | #define CVMFS_UTIL_TUBE_H_ | ||
7 | |||
8 | #include <pthread.h> | ||
9 | #include <stdint.h> | ||
10 | |||
11 | #include <cassert> | ||
12 | #include <cstddef> | ||
13 | #include <vector> | ||
14 | |||
15 | #include "util/atomic.h" | ||
16 | #include "util/mutex.h" | ||
17 | #include "util/single_copy.h" | ||
18 | |||
19 | /** | ||
20 | * A thread-safe, doubly linked list of links containing pointers to ItemT. The | ||
21 | * ItemT elements are not owned by the Tube. FIFO or LIFO semantics. Using | ||
22 | * Slice(), items at arbitrary locations in the tube can be removed, too. | ||
23 | * | ||
24 | * | ||
25 | * The layout of the linked list is as follows: | ||
26 | * | ||
27 | * -------------------------------------------------------------- | ||
28 | * | | | ||
29 | * --> I$n$ (back) <--> I2 <--> ... <--> I1 (front) <--> HEAD <-- | ||
30 | * | ||
31 | * The tube links the steps in the file processing pipeline. It connects | ||
32 | * multiple producers to multiple consumers and can throttle the producers if a | ||
33 | * limit for the tube size is set. | ||
34 | * | ||
35 | * Internally, uses conditional variables to block when threads try to pop from | ||
36 | * the empty tube or insert into the full tube. | ||
37 | */ | ||
38 | template<class ItemT> | ||
39 | class Tube : SingleCopy { | ||
40 | public: | ||
41 | class Link : SingleCopy { | ||
42 | friend class Tube<ItemT>; | ||
43 | |||
44 | public: | ||
45 | 158781305 | explicit Link(ItemT *item) : item_(item), next_(NULL), prev_(NULL) { } | |
46 | 18 | ItemT *item() { return item_; } | |
47 | |||
48 | private: | ||
49 | ItemT *item_; | ||
50 | Link *next_; | ||
51 | Link *prev_; | ||
52 | }; | ||
53 | |||
54 | 202905 | Tube() : limit_(uint64_t(-1)), size_(0) { Init(); } | |
55 | 1634 | explicit Tube(uint64_t limit) : limit_(limit), size_(0) { Init(); } | |
56 | 208800 | ~Tube() { | |
57 | 208876 | Link *cursor = head_; | |
58 | do { | ||
59 | 208885 | Link *prev = cursor->prev_; | |
60 |
1/2✓ Branch 0 taken 105956 times.
✗ Branch 1 not taken.
|
208885 | delete cursor; |
61 | 208885 | cursor = prev; | |
62 |
2/2✓ Branch 0 taken 9 times.
✓ Branch 1 taken 105947 times.
|
208885 | } while (cursor != head_); |
63 | 208876 | pthread_cond_destroy(&cond_populated_); | |
64 | 208876 | pthread_cond_destroy(&cond_capacious_); | |
65 | 208876 | pthread_cond_destroy(&cond_empty_); | |
66 | 208876 | pthread_mutex_destroy(&lock_); | |
67 | 208876 | } | |
68 | |||
69 | /** | ||
70 | * Push an item to the back of the queue. Block if queue is currently full. | ||
71 | */ | ||
72 | 155926150 | Link *EnqueueBack(ItemT *item) { | |
73 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 84792611 times.
|
155926150 | assert(item != NULL); |
74 | 155926150 | MutexLockGuard lock_guard(&lock_); | |
75 |
2/2✓ Branch 0 taken 3772923 times.
✓ Branch 1 taken 84316067 times.
|
162511747 | while (size_ == limit_) |
76 |
1/2✓ Branch 1 taken 3755661 times.
✗ Branch 2 not taken.
|
7545846 | pthread_cond_wait(&cond_capacious_, &lock_); |
77 | |||
78 |
1/2✓ Branch 1 taken 84552871 times.
✗ Branch 2 not taken.
|
154965901 | Link *link = new Link(item); |
79 | 154098901 | link->next_ = head_->next_; | |
80 | 154098901 | link->prev_ = head_; | |
81 | 154098901 | head_->next_->prev_ = link; | |
82 | 154098901 | head_->next_ = link; | |
83 | 154098901 | size_++; | |
84 | 154098901 | int retval = pthread_cond_signal(&cond_populated_); | |
85 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 85141304 times.
|
156616375 | assert(retval == 0); |
86 | 156023987 | return link; | |
87 | 156616375 | } | |
88 | |||
89 | /** | ||
90 | * Push an item to the front of the queue. Block if queue currently full. | ||
91 | */ | ||
92 | 3638533 | Link *EnqueueFront(ItemT *item) { | |
93 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 3637899 times.
|
3638533 | assert(item != NULL); |
94 | 3638533 | MutexLockGuard lock_guard(&lock_); | |
95 |
2/2✓ Branch 0 taken 18 times.
✓ Branch 1 taken 3637881 times.
|
3638551 | while (size_ == limit_) |
96 |
1/2✓ Branch 1 taken 18 times.
✗ Branch 2 not taken.
|
36 | pthread_cond_wait(&cond_capacious_, &lock_); |
97 | |||
98 |
1/2✓ Branch 1 taken 3637872 times.
✗ Branch 2 not taken.
|
3638515 | Link *link = new Link(item); |
99 | 3638497 | link->next_ = head_; | |
100 | 3638497 | link->prev_ = head_->prev_; | |
101 | 3638497 | head_->prev_->next_ = link; | |
102 | 3638497 | head_->prev_ = link; | |
103 | 3638497 | size_++; | |
104 | 3638497 | int retval = pthread_cond_signal(&cond_populated_); | |
105 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 3637854 times.
|
3638488 | assert(retval == 0); |
106 | 3638515 | return link; | |
107 | 3638488 | } | |
108 | |||
109 | /** | ||
110 | * Remove any link from the queue and return its item, including first/last | ||
111 | * element. | ||
112 | */ | ||
113 | 29 | ItemT *Slice(Link *link) { | |
114 | 29 | MutexLockGuard lock_guard(&lock_); | |
115 | 29 | return SliceUnlocked(link); | |
116 | 29 | } | |
117 | |||
118 | /** | ||
119 | * Remove and return the first element from the queue. Block if tube is | ||
120 | * empty. | ||
121 | */ | ||
122 | 158905082 | ItemT *PopFront() { | |
123 | 158905082 | MutexLockGuard lock_guard(&lock_); | |
124 |
2/2✓ Branch 0 taken 31370402 times.
✓ Branch 1 taken 87916631 times.
|
220500145 | while (size_ == 0) |
125 |
1/2✓ Branch 1 taken 31436038 times.
✗ Branch 2 not taken.
|
61663264 | pthread_cond_wait(&cond_populated_, &lock_); |
126 | 316990900 | return SliceUnlocked(head_->prev_); | |
127 | 157884265 | } | |
128 | |||
129 | /** | ||
130 | * Remove and return the first element from the queue if there is any. | ||
131 | * Equivalent to an antomic | ||
132 | * ItemT item = NULL; | ||
133 | * if (!IsEmpty()) | ||
134 | * item = PopFront(); | ||
135 | */ | ||
136 | 3635887 | ItemT *TryPopFront() { | |
137 | 3635887 | MutexLockGuard lock_guard(&lock_); | |
138 | // Note that we don't need to wait for a signal to arrive | ||
139 |
2/2✓ Branch 0 taken 3330202 times.
✓ Branch 1 taken 306936 times.
|
3637138 | if (size_ == 0) |
140 | 3330202 | return NULL; | |
141 | 306936 | return SliceUnlocked(head_->prev_); | |
142 | 3637138 | } | |
143 | |||
144 | /** | ||
145 | * Remove and return the last element from the queue. Block if tube is | ||
146 | * empty. | ||
147 | */ | ||
148 | 1411 | ItemT *PopBack() { | |
149 | 1411 | MutexLockGuard lock_guard(&lock_); | |
150 |
2/2✓ Branch 0 taken 153 times.
✓ Branch 1 taken 777 times.
|
1717 | while (size_ == 0) |
151 |
1/2✓ Branch 1 taken 153 times.
✗ Branch 2 not taken.
|
306 | pthread_cond_wait(&cond_populated_, &lock_); |
152 | 2822 | return SliceUnlocked(head_->next_); | |
153 | 1411 | } | |
154 | |||
155 | /** | ||
156 | * Blocks until the tube is empty | ||
157 | */ | ||
158 | 2043 | void Wait() { | |
159 | 2043 | MutexLockGuard lock_guard(&lock_); | |
160 |
2/2✓ Branch 0 taken 431 times.
✓ Branch 1 taken 2043 times.
|
2474 | while (size_ > 0) |
161 |
1/2✓ Branch 1 taken 431 times.
✗ Branch 2 not taken.
|
431 | pthread_cond_wait(&cond_empty_, &lock_); |
162 | 2043 | } | |
163 | |||
164 | 11843 | bool IsEmpty() { | |
165 | 11843 | MutexLockGuard lock_guard(&lock_); | |
166 | 23686 | return size_ == 0; | |
167 | 11843 | } | |
168 | |||
169 | 247 | uint64_t size() { | |
170 | 247 | MutexLockGuard lock_guard(&lock_); | |
171 | 494 | return size_; | |
172 | 247 | } | |
173 | |||
174 | private: | ||
175 | 205869 | void Init() { | |
176 | 205869 | Link *sentinel = new Link(NULL); | |
177 | 205869 | head_ = sentinel; | |
178 | 205869 | head_->next_ = head_->prev_ = sentinel; | |
179 | |||
180 | 205869 | int retval = pthread_mutex_init(&lock_, NULL); | |
181 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 106023 times.
|
205869 | assert(retval == 0); |
182 | 205869 | retval = pthread_cond_init(&cond_populated_, NULL); | |
183 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 106023 times.
|
205869 | assert(retval == 0); |
184 | 205869 | retval = pthread_cond_init(&cond_capacious_, NULL); | |
185 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 106023 times.
|
205869 | assert(retval == 0); |
186 | 205869 | retval = pthread_cond_init(&cond_empty_, NULL); | |
187 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 106023 times.
|
205869 | assert(retval == 0); |
188 | 205869 | } | |
189 | |||
190 | 158660729 | ItemT *SliceUnlocked(Link *link) { | |
191 | // Cannot delete the sentinel link | ||
192 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 87981792 times.
|
158660729 | assert(link != head_); |
193 | 158660729 | link->prev_->next_ = link->next_; | |
194 | 158660729 | link->next_->prev_ = link->prev_; | |
195 | 158660729 | ItemT *item = link->item_; | |
196 |
2/2✓ Branch 0 taken 87682987 times.
✓ Branch 1 taken 298805 times.
|
158660729 | delete link; |
197 | 160123324 | size_--; | |
198 | 160123324 | int retval = pthread_cond_signal(&cond_capacious_); | |
199 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 88486254 times.
|
159669682 | assert(retval == 0); |
200 |
2/2✓ Branch 0 taken 35444830 times.
✓ Branch 1 taken 53041424 times.
|
159669682 | if (size_ == 0) { |
201 | 69068929 | retval = pthread_cond_broadcast(&cond_empty_); | |
202 |
2/2✓ Branch 0 taken 349188 times.
✓ Branch 1 taken 35107815 times.
|
69093275 | assert(retval == 0); |
203 | } | ||
204 | 158995652 | return item; | |
205 | } | ||
206 | |||
207 | |||
208 | /** | ||
209 | * Adding new item blocks as long as limit_ == size_ | ||
210 | */ | ||
211 | uint64_t limit_; | ||
212 | /** | ||
213 | * The current number of links in the list | ||
214 | */ | ||
215 | uint64_t size_; | ||
216 | /** | ||
217 | * Sentinel element in front of the first (front) element | ||
218 | */ | ||
219 | Link *head_; | ||
220 | /** | ||
221 | * Protects all internal state | ||
222 | */ | ||
223 | pthread_mutex_t lock_; | ||
224 | /** | ||
225 | * Signals if there are items enqueued | ||
226 | */ | ||
227 | pthread_cond_t cond_populated_; | ||
228 | /** | ||
229 | * Signals if there is space to enqueue more items | ||
230 | */ | ||
231 | pthread_cond_t cond_capacious_; | ||
232 | /** | ||
233 | * Signals if the queue runs empty | ||
234 | */ | ||
235 | pthread_cond_t cond_empty_; | ||
236 | }; | ||
237 | |||
238 | |||
239 | /** | ||
240 | * A tube group manages a fixed set of Tubes and dispatches items among them in | ||
241 | * such a way that items with the same tag (a positive integer) are all sent | ||
242 | * to the same tube. | ||
243 | */ | ||
244 | template<class ItemT> | ||
245 | class TubeGroup : SingleCopy { | ||
246 | public: | ||
247 | 16303 | TubeGroup() : is_active_(false) { atomic_init32(&round_robin_); } | |
248 | |||
249 | 16292 | ~TubeGroup() { | |
250 |
2/2✓ Branch 1 taken 98736 times.
✓ Branch 2 taken 9583 times.
|
210476 | for (unsigned i = 0; i < tubes_.size(); ++i) |
251 |
1/2✓ Branch 1 taken 98736 times.
✗ Branch 2 not taken.
|
194184 | delete tubes_[i]; |
252 | 16292 | } | |
253 | |||
254 | 194329 | void TakeTube(Tube<ItemT> *t) { | |
255 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 98809 times.
|
194329 | assert(!is_active_); |
256 | 194329 | tubes_.push_back(t); | |
257 | 194329 | } | |
258 | |||
259 | 16303 | void Activate() { | |
260 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 9589 times.
|
16303 | assert(!is_active_); |
261 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 9589 times.
|
16303 | assert(!tubes_.empty()); |
262 | 16303 | is_active_ = true; | |
263 | 16303 | } | |
264 | |||
265 | /** | ||
266 | * Like Tube::EnqueueBack(), but pick a tube according to ItemT::tag() | ||
267 | */ | ||
268 | 63849799 | typename Tube<ItemT>::Link *Dispatch(ItemT *item) { | |
269 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 63849799 times.
|
63849799 | assert(is_active_); |
270 |
2/2✓ Branch 1 taken 49935720 times.
✓ Branch 2 taken 13785958 times.
|
63849799 | unsigned tube_idx = (tubes_.size() == 1) ? 0 |
271 | 49935720 | : (item->tag() % tubes_.size()); | |
272 | 63222067 | return tubes_[tube_idx]->EnqueueBack(item); | |
273 | } | ||
274 | |||
275 | /** | ||
276 | * Like Tube::EnqueueBack(), use tubes one after another | ||
277 | */ | ||
278 | 5252645 | typename Tube<ItemT>::Link *DispatchAny(ItemT *item) { | |
279 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 5252645 times.
|
5252645 | assert(is_active_); |
280 |
2/2✓ Branch 1 taken 5252629 times.
✓ Branch 2 taken 16 times.
|
5252645 | unsigned tube_idx = (tubes_.size() == 1) |
281 | ? 0 | ||
282 | 5252629 | : (atomic_xadd32(&round_robin_, 1) % tubes_.size()); | |
283 | 5252645 | return tubes_[tube_idx]->EnqueueBack(item); | |
284 | } | ||
285 | |||
286 | private: | ||
287 | bool is_active_; | ||
288 | std::vector<Tube<ItemT> *> tubes_; | ||
289 | atomic_int32 round_robin_; | ||
290 | }; | ||
291 | |||
292 | #endif // CVMFS_UTIL_TUBE_H_ | ||
293 |