casacore
Loading...
Searching...
No Matches
Lane.h
Go to the documentation of this file.
1#ifndef CASACORE_AOCOMMON_LANE_11_H_
2#define CASACORE_AOCOMMON_LANE_11_H_
3
4#include <condition_variable>
5#include <cstring>
6#include <deque>
7#include <mutex>
8
17
18// #define LANE_DEBUG_MODE
19
20#ifdef LANE_DEBUG_MODE
21#include <cmath>
22#include <iostream>
23#include <sstream>
24#include <string>
25#endif
26
28
29#ifdef LANE_DEBUG_MODE
30#define set_lane_debug_name(lane, str) (lane).setDebugName(str)
31#define LANE_REGISTER_DEBUG_INFO registerDebugInfo()
32#define LANE_REGISTER_DEBUG_WRITE_WAIT registerDebugWriteWait()
33#define LANE_REGISTER_DEBUG_READ_WAIT registerDebugReadWait()
34#define LANE_REPORT_DEBUG_INFO reportDebugInfo()
35
36#else
37
38#define set_lane_debug_name(lane, str)
39#define LANE_REGISTER_DEBUG_INFO
40#define LANE_REGISTER_DEBUG_WRITE_WAIT
41#define LANE_REGISTER_DEBUG_READ_WAIT
42#define LANE_REPORT_DEBUG_INFO
43
44#endif
45
99template <typename Tp>
100class Lane {
101 public:
103 typedef std::size_t size_type;
104
106 typedef Tp value_type;
107
116 Lane() noexcept
117 : _buffer(nullptr),
118 _capacity(0),
122
128 explicit Lane(size_t capacity)
129 : _buffer(new Tp[capacity]),
134
135 Lane(const Lane<Tp>& source) = delete;
136
142 Lane(Lane<Tp>&& source) noexcept
143 : _buffer(nullptr),
144 _capacity(0),
148 swap(source);
149 }
150
156 delete[] _buffer;
157 }
158
159 Lane<Tp>& operator=(const Lane<Tp>& source) = delete;
160
167 Lane<Tp>& operator=(Lane<Tp>&& source) noexcept {
168 swap(source);
169 return *this;
170 }
171
176 void swap(Lane<Tp>& other) noexcept {
177 std::swap(_buffer, other._buffer);
178 std::swap(_capacity, other._capacity);
179 std::swap(_write_position, other._write_position);
180 std::swap(_free_write_space, other._free_write_space);
181 std::swap(_status, other._status);
182 }
183
191 void clear() noexcept {
192 _write_position = 0;
195 }
196
206 void write(const value_type& element) {
207 std::unique_lock<std::mutex> lock(_mutex);
209
210 while (_free_write_space == 0 && _status == status_normal) {
213 }
214 if (_status == status_normal) {
215 _buffer[_write_position] = element;
218 // Now that there is less free write space, there is more free read
219 // space and thus readers may continue.
220 _reading_possible_condition.notify_all();
221 }
222 }
223
235 template <typename... Args>
236 void emplace(Args&&... args) {
237 std::unique_lock<std::mutex> lock(_mutex);
239
240 while (_free_write_space == 0 && _status == status_normal) {
243 }
244
245 if (_status == status_normal) {
246 _buffer[_write_position] = value_type(std::forward<Args>(args)...);
249 // Now that there is less free write space, there is more free read
250 // space and thus readers can possibly continue.
251 _reading_possible_condition.notify_all();
252 }
253 }
254
263 void write(value_type&& element) {
264 std::unique_lock<std::mutex> lock(_mutex);
266
267 while (_free_write_space == 0 && _status == status_normal) {
270 }
271 if (_status == status_normal) {
272 _buffer[_write_position] = std::move(element);
275 // Now that there is less free write space, there is more free read
276 // space and thus readers can possibly continue.
277 _reading_possible_condition.notify_all();
278 }
279 }
280
281 void write(const value_type* elements, size_t n) { write_generic(elements, n); }
282
283 void move_write(value_type* elements, size_t n) { write_generic(elements, n); }
284
285 bool read(value_type& destination) {
286 std::unique_lock<std::mutex> lock(_mutex);
288 while (free_read_space() == 0 && _status == status_normal) {
291 }
292 if (free_read_space() == 0)
293 return false;
294 else {
295 destination = std::move(_buffer[read_position()]);
297 // Now that there is more free write space, writers can possibly continue.
298 _writing_possible_condition.notify_all();
299 return true;
300 }
301 }
302
303 size_t read(value_type* destinations, size_t n) {
304 size_t n_left = n;
305
306 std::unique_lock<std::mutex> lock(_mutex);
308
309 size_t free_space = free_read_space();
310 size_t read_size = free_space > n ? n : free_space;
311 immediate_read(destinations, read_size);
312 n_left -= read_size;
313
314 while (n_left != 0 && _status == status_normal) {
315 destinations += read_size;
316
317 do {
320 } while (free_read_space() == 0 && _status == status_normal);
321
322 free_space = free_read_space();
323 read_size = free_space > n_left ? n_left : free_space;
324 immediate_read(destinations, read_size);
325 n_left -= read_size;
326 }
327 return n - n_left;
328 }
329
335 size_t discard(size_t n) {
336 size_t n_left = n;
337
338 std::unique_lock<std::mutex> lock(_mutex);
340
341 size_t free_space = free_read_space();
342 size_t read_size = free_space > n ? n : free_space;
343 immediate_discard(read_size);
344 n_left -= read_size;
345
346 while (n_left != 0 && _status == status_normal) {
347 do {
350 } while (free_read_space() == 0 && _status == status_normal);
351
352 free_space = free_read_space();
353 read_size = free_space > n_left ? n_left : free_space;
354 immediate_discard(read_size);
355 n_left -= read_size;
356 }
357 return n - n_left;
358 }
359
360 void write_end() {
361 std::lock_guard<std::mutex> lock(_mutex);
364 _writing_possible_condition.notify_all();
365 _reading_possible_condition.notify_all();
366 }
367
368 size_t capacity() const noexcept { return _capacity; }
369
370 size_t size() const {
371 std::lock_guard<std::mutex> lock(_mutex);
373 }
374
375 bool empty() const {
376 std::lock_guard<std::mutex> lock(_mutex);
378 }
379
384 bool is_end() const {
385 std::lock_guard<std::mutex> lock(_mutex);
386 return _status == status_end;
387 }
388
392 bool is_end_and_empty() const {
393 std::lock_guard<std::mutex> lock(_mutex);
395 }
396
401 void resize(size_t new_capacity) {
402 Tp* new_buffer = new Tp[new_capacity];
403 delete[] _buffer;
404 _buffer = new_buffer;
405 _capacity = new_capacity;
406 _write_position = 0;
407 _free_write_space = new_capacity;
409 }
410
415 std::unique_lock<std::mutex> lock(_mutex);
416 while (_capacity != _free_write_space) {
418 }
419 }
420
421#ifdef LANE_DEBUG_MODE
428 void setDebugName(const std::string& nameStr) { _debugName = nameStr; }
429#endif
430 private:
432
433 size_t _capacity;
434
436
438
440
441 mutable std::mutex _mutex;
442
444
445 size_t read_position() const noexcept {
447 }
448
449 size_t free_read_space() const noexcept { return _capacity - _free_write_space; }
450
451 // This is a template to allow const and non-const (to be able to move)
452 template <typename T>
453 void write_generic(T* elements, size_t n) {
454 std::unique_lock<std::mutex> lock(_mutex);
456
457 if (_status == status_normal) {
458 size_t write_size = _free_write_space > n ? n : _free_write_space;
459 immediate_write(elements, write_size);
460 n -= write_size;
461
462 while (n != 0 && _status == status_normal) {
463 elements += write_size;
464
465 do {
468 } while (_free_write_space == 0 && _status == status_normal);
469
470 write_size = _free_write_space > n ? n : _free_write_space;
471 immediate_write(elements, write_size);
472 n -= write_size;
473 }
474 }
475 }
476
477 // This is a template to allow const and non-const (to be able to move)
478 template <typename T>
479 void immediate_write(T* elements, size_t n) noexcept {
480 // Split the writing in two ranges if needed. The first range fits in
481 // [_write_position, _capacity), the second range in [0, end). By doing
482 // so, we only have to calculate the modulo in the write position once.
483 if (n > 0) {
484 size_t nPart;
485 if (_write_position + n > _capacity) {
486 nPart = _capacity - _write_position;
487 } else {
488 nPart = n;
489 }
490 for (size_t i = 0; i < nPart; ++i, ++_write_position) {
491 _buffer[_write_position] = std::move(elements[i]);
492 }
493
495
496 for (size_t i = nPart; i < n; ++i, ++_write_position) {
497 _buffer[_write_position] = std::move(elements[i]);
498 }
499
501
502 // Now that there is less free write space, there is more free read
503 // space and thus readers may continue.
504 _reading_possible_condition.notify_all();
505 }
506 }
507
508 void immediate_read(value_type* elements, size_t n) noexcept {
509 // As with write, split in two ranges if needed. The first range fits in
510 // [read_position(), _capacity), the second range in [0, end).
511 if (n > 0) {
512 size_t nPart;
513 size_t position = read_position();
514 if (position + n > _capacity) {
515 nPart = _capacity - position;
516 } else {
517 nPart = n;
518 }
519 for (size_t i = 0; i < nPart; ++i, ++position) {
520 elements[i] = std::move(_buffer[position]);
521 }
522
523 position = position % _capacity;
524
525 for (size_t i = nPart; i < n; ++i, ++position) {
526 elements[i] = std::move(_buffer[position]);
527 }
528
530
531 // Now that there is more free write space, writers can possibly continue.
532 _writing_possible_condition.notify_all();
533 }
534 }
535
536 void immediate_discard(size_t n) noexcept {
537 if (n > 0) {
539
540 // Now that there is more free write space, writers can possibly continue.
541 _writing_possible_condition.notify_all();
542 }
543 }
544
545#ifdef LANE_DEBUG_MODE
546 void registerDebugInfo() noexcept {
547 _debugSummedSize += _capacity - _free_write_space;
548 _debugMeasureCount++;
549 }
550 void registerDebugReadWait() noexcept { ++_debugReadWaitCount; }
551 void registerDebugWriteWait() noexcept { ++_debugWriteWaitCount; }
552 void reportDebugInfo() {
553 if (!_debugName.empty()) {
554 std::stringstream str;
555 str << "*** Debug report for the following Lane: ***\n"
556 << "\"" << _debugName << "\"\n"
557 << "Capacity: " << _capacity << '\n'
558 << "Total read/write ops: " << _debugMeasureCount << '\n'
559 << "Average size of buffer, measured per read/write op.: "
560 << round(double(_debugSummedSize) * 100.0 / _debugMeasureCount) / 100.0 << '\n'
561 << "Number of wait events during reading: " << _debugReadWaitCount << '\n'
562 << "Number of wait events during writing: " << _debugWriteWaitCount << '\n';
563 std::cout << str.str();
564 }
565 }
566 std::string _debugName;
567 size_t _debugSummedSize = 0, _debugMeasureCount = 0, _debugReadWaitCount = 0,
568 _debugWriteWaitCount = 0;
569#endif
570};
571
572template <typename Tp>
574 first.swap(second);
575}
576
577} // namespace casacore::aocommon
578
579#endif // AO_LANE11_H
#define LANE_REGISTER_DEBUG_WRITE_WAIT
Definition Lane.h:40
#define LANE_REGISTER_DEBUG_INFO
Definition Lane.h:39
#define LANE_REGISTER_DEBUG_READ_WAIT
Definition Lane.h:41
#define LANE_REPORT_DEBUG_INFO
Definition Lane.h:42
The Lane is an efficient cyclic buffer that is synchronized.
Definition Lane.h:100
std::mutex _mutex
Definition Lane.h:441
std::condition_variable _writing_possible_condition
Definition Lane.h:443
void immediate_discard(size_t n) noexcept
Definition Lane.h:536
void emplace(Args &&... args)
Write a single element by constructing it.
Definition Lane.h:236
Lane(const Lane< Tp > &source)=delete
bool is_end() const
True when write_end() was called.
Definition Lane.h:384
void move_write(value_type *elements, size_t n)
Definition Lane.h:283
~Lane()
Destructor.
Definition Lane.h:154
void clear() noexcept
Clear the contents and reset the state of the Lane.
Definition Lane.h:191
void swap(Lane< Tp > &other) noexcept
Swap the contents of this Lane with another.
Definition Lane.h:176
Lane< Tp > & operator=(Lane< Tp > &&source) noexcept
Move assignment.
Definition Lane.h:167
void wait_for_empty()
Wait until this Lane is empty.
Definition Lane.h:414
Lane< Tp > & operator=(const Lane< Tp > &source)=delete
void resize(size_t new_capacity)
Change the capacity of the Lane.
Definition Lane.h:401
void write(value_type &&element)
Write a single element by moving it in.
Definition Lane.h:263
void immediate_write(T *elements, size_t n) noexcept
This is a template to allow const and non-const (to be able to move).
Definition Lane.h:479
Lane(Lane< Tp > &&source) noexcept
Move construct a Lane.
Definition Lane.h:142
size_t capacity() const noexcept
Definition Lane.h:368
Tp value_type
Type of elements stored in the Lane.
Definition Lane.h:106
size_t size() const
Definition Lane.h:370
void write(const value_type *elements, size_t n)
Definition Lane.h:281
enum casacore::aocommon::Lane::@165250377254016164334045006074227223211343376364 _status
size_t discard(size_t n)
This method does the same thing as read(buffer, n) but discards the data.
Definition Lane.h:335
std::size_t size_type
Integer type used to store size types.
Definition Lane.h:103
size_t read_position() const noexcept
Definition Lane.h:445
void write_generic(T *elements, size_t n)
This is a template to allow const and non-const (to be able to move).
Definition Lane.h:453
bool empty() const
Definition Lane.h:375
void write(const value_type &element)
Write a single element.
Definition Lane.h:206
size_t read(value_type *destinations, size_t n)
Definition Lane.h:303
std::condition_variable _reading_possible_condition
Definition Lane.h:443
Lane(size_t capacity)
Construct a Lane with the given capacity.
Definition Lane.h:128
size_t free_read_space() const noexcept
Definition Lane.h:449
bool is_end_and_empty() const
True when write_end() and the lane does not contain items.
Definition Lane.h:392
void immediate_read(value_type *elements, size_t n) noexcept
Definition Lane.h:508
bool read(value_type &destination)
Definition Lane.h:285
Lane() noexcept
Construct a Lane with zero elements.
Definition Lane.h:116
struct Node * first
Definition malloc.h:325
void swap(aocommon::Lane< Tp > &first, aocommon::Lane< Tp > &second) noexcept
Definition Lane.h:573
LatticeExprNode round(const LatticeExprNode &expr)