otsdaq-mu2e-stm  5.02.01
ring_buffer.hh
1 #ifndef RING_BUFFER_hh_
2 #define RING_BUFFER_hh_
3 
4 #include <iostream>
5 #include <thread>
6 #include <mutex>
7 #include <condition_variable>
8 #include <vector>
9 #include <queue>
10 #include <atomic>
11 
12 // Global total data tracker
13 static std::atomic<uint64_t> globalTotalBytes{0};
14 //static const uint64_t targetTotalBytes = 10L * 1024 * 1024 * 1024; // 10 GB
15 static const uint64_t targetTotalBytes = 28800L * 1024 * 1024 * 1024; // 10 GB
16 
17 // Data structure representing a buffer with metadata
18 struct DataStruct {
19  std::vector<int16_t> data; // Dynamically sized data buffer
20  std::vector<int16_t> data2; // Dynamically sized data buffer
21  std::vector<int16_t> data3; // Dynamically sized data buffer
22  size_t size; // Actual size of the data used
23  size_t size2; // Actual size of the data used
24  size_t size3; // Actual size of the data used
25  char metadata[128]; // Fixed-size metadata array
26 
27  DataStruct(size_t buffer_size)
28  // : data(buffer_size), size(0) {
29  : data(buffer_size), data2(), data3(), size(0), size2(0), size3(0) {
30  std::fill(std::begin(metadata), std::end(metadata), '\0');
31  }
32 };
33 
34 // Ring buffer for thread-safe buffer management
35 template <typename T> class RingBuffer {
36 
37 private:
38 
39  // Circular buffer storage
40  std::vector<std::shared_ptr<T>> buffer;
41  // Index of next insertion point
42  size_t head = 0;
43  // Index of next removal point
44  size_t tail = 0;
45  // Maxaimum capacity of buffer
46  size_t capacity;
47  // Current size of buffer
48  std::atomic<size_t> size = 0;
49  // Mutex for synchronisation
50  std::mutex mutex;
51  // Conditional variable for producers
52  std::condition_variable cv_push;
53  // Conditiona variable for consumers
54  std::condition_variable cv_pop;
55  // Flag to stop buffer operations
56  bool stop = false;
57 
58 public:
59 
60  // Constructor to initialise buffer
61  explicit RingBuffer(size_t max_size) : capacity(max_size), buffer(max_size) {}
62 
63  // Destructor for logging (not actually needed with shared_ptr)
64  ~RingBuffer() {
65  std::cout << "RingBuffer destructor called.\n";
66  }
67 
68  // Add an item to the buffer
69  void push(const std::shared_ptr<T>& item);
70 
71  // Remove and return an item from the buffer
72  std::shared_ptr<T> pop();
73 
74  // Stop buffer operations
75  void shutdown();
76 
77  // Check if buffer is empty
78  bool is_empty();
79 
80 };
81 
82 #endif
Definition: data.hh:4