otsdaq-mu2e-stm  5.02.01
thread_buffer_management.cpp
1 #include <atomic>
2 #include <condition_variable>
3 #include <cstring>
4 #include <iostream>
5 #include <memory>
6 #include <mutex>
7 #include <queue>
8 #include <thread>
9 #include <vector>
10 
11 // Data structure representing a buffer with metadata
12 struct DataStruct
13 {
14  int16_t data[60 * 1024 * 1024 / sizeof(int16_t)]; // 60 MB worth of int16_t data
15  size_t size; // Actual size of the data used
16  std::string metadata; // Additional metadata
17 };
18 
19 // Thread-safe queue for buffer management
20 template<typename T>
22 {
23  private:
24  std::queue<std::shared_ptr<T>> queue;
25  std::mutex mutex;
26  std::condition_variable cv;
27  bool stop = false;
28  size_t max_size; // Maximum queue size
29 
30  public:
31  explicit ThreadSafeQueue(size_t max_size) : max_size(max_size) {}
32 
33  void push(const std::shared_ptr<T>& item)
34  {
35  {
36  std::unique_lock<std::mutex> lock(mutex);
37  cv.wait(lock, [this]() { return stop || queue.size() < max_size; });
38  if(stop)
39  return; // Stop pushing if the queue is shutting down
40  queue.push(item);
41  }
42  cv.notify_all();
43  }
44 
45  std::shared_ptr<T> pop()
46  {
47  std::unique_lock<std::mutex> lock(mutex);
48  cv.wait(lock, [this]() { return stop || !queue.empty(); });
49  if(queue.empty())
50  {
51  return nullptr;
52  }
53  auto item = queue.front();
54  queue.pop();
55  cv.notify_all();
56  return item;
57  }
58 
59  void shutdown()
60  {
61  {
62  std::lock_guard<std::mutex> lock(mutex);
63  stop = true;
64  }
65  cv.notify_all();
66  }
67 };
68 
69 // Function for Kernel Buffer and UDP Server
70 auto kernelBufferThread = [](ThreadSafeQueue<DataStruct>& bufferQueue) {
71  while(true)
72  {
73  auto buffer = std::make_shared<DataStruct>();
74  buffer->size = (rand() % (60 * 1024 * 1024 / sizeof(int16_t))) +
75  1000; // Simulate variable size between 1,000 and max capacity
76  std::fill(buffer->data,
77  buffer->data + buffer->size,
78  rand() % 100); // Fill with random data
79  buffer->metadata = "Generated by kernelBufferThread"; // Example metadata
80  bufferQueue.push(buffer);
81  std::this_thread::sleep_for(std::chrono::milliseconds(10));
82  }
83 };
84 
85 // Function for Data Checking
86 auto checkDataThread = [](ThreadSafeQueue<DataStruct>& inputQueue,
87  ThreadSafeQueue<DataStruct>& outputQueue) {
88  while(true)
89  {
90  auto buffer = inputQueue.pop();
91  if(!buffer)
92  break; // Exit if queue is shutting down
93  // Simulate data checking (e.g., verifying integrity)
94  buffer->metadata += " | Checked";
95  outputQueue.push(buffer);
96  }
97 };
98 
99 // Function for Forming Events
100 auto formEventsThread = [](ThreadSafeQueue<DataStruct>& inputQueue,
101  ThreadSafeQueue<DataStruct>& outputQueue) {
102  while(true)
103  {
104  auto buffer = inputQueue.pop();
105  if(!buffer)
106  break; // Exit if queue is shutting down
107  // Simulate event formation (e.g., extracting relevant data)
108  buffer->metadata += " | Events Formed";
109  outputQueue.push(buffer);
110  }
111 };
112 
113 // Function for Processing Algorithms
114 auto processAlgorithmThread = [](ThreadSafeQueue<DataStruct>& inputQueue,
115  ThreadSafeQueue<DataStruct>& outputQueue) {
116  while(true)
117  {
118  auto buffer = inputQueue.pop();
119  if(!buffer)
120  break; // Exit if queue is shutting down
121  // Simulate algorithm processing (prescale, ZS, MWD, etc.)
122  buffer->metadata += " | Algorithm Processed";
123  outputQueue.push(buffer);
124  }
125 };
126 
127 int main()
128 {
129  const size_t max_queue_size = 10; // Example maximum size for each queue
130 
131  // Define thread-safe queues for each step
132  ThreadSafeQueue<DataStruct> kernelBufferQueue(max_queue_size);
133  ThreadSafeQueue<DataStruct> checkDataQueue(max_queue_size);
134  ThreadSafeQueue<DataStruct> formEventsQueue(max_queue_size);
135  ThreadSafeQueue<DataStruct> rawBufferQueue(max_queue_size);
136  ThreadSafeQueue<DataStruct> zsBufferQueue(max_queue_size);
137  ThreadSafeQueue<DataStruct> mwdBufferQueue(max_queue_size);
138 
139  // Launch threads for data processing
140  std::thread kernelThread(kernelBufferThread, std::ref(kernelBufferQueue));
141  std::thread checkData(
142  checkDataThread, std::ref(kernelBufferQueue), std::ref(checkDataQueue));
143  std::thread formEvents(
144  formEventsThread, std::ref(checkDataQueue), std::ref(formEventsQueue));
145  std::thread rawProcess(
146  processAlgorithmThread, std::ref(formEventsQueue), std::ref(rawBufferQueue));
147  std::thread zsProcess(
148  processAlgorithmThread, std::ref(rawBufferQueue), std::ref(zsBufferQueue));
149  std::thread mwdProcess(
150  processAlgorithmThread, std::ref(zsBufferQueue), std::ref(mwdBufferQueue));
151 
152  // Simulate running for a limited time
153  std::this_thread::sleep_for(std::chrono::seconds(10));
154 
155  // Shutdown all queues
156  kernelBufferQueue.shutdown();
157  checkDataQueue.shutdown();
158  formEventsQueue.shutdown();
159  rawBufferQueue.shutdown();
160  zsBufferQueue.shutdown();
161  mwdBufferQueue.shutdown();
162 
163  // Join threads
164  kernelThread.join();
165  checkData.join();
166  formEvents.join();
167  rawProcess.join();
168  zsProcess.join();
169  mwdProcess.join();
170 
171  return 0;
172 }
Definition: data.hh:4
Definition: queue.hh:42