otsdaq-mu2e-stm  5.02.01
queue_safe.cc
1 #include <unistd.h>
2 #include <atomic>
3 #include <cinttypes>
4 #include <iostream>
5 #include <thread>
6 
7 // Queue buffer header
8 #include "STMDAQ-TestBeam/utils/queue.hh"
9 
10 // Constructor
11 queue_buffer::queue_buffer()
12 {
13  // For each channel, set read and write pointers to zero
14  for(int chan = 0; chan < CHNUM; chan++)
15  {
16  // Initialise tail pointer
17  // tail_[chan].store({0,0,0,0});
18  tail[chan].store(0);
19  // Initialise head pointers
20  // head_[chan].store({0,0,0,0});
21  head[chan].store(0);
22  // Initialise occupancy pointer
23  occupancy[chan].store(0);
24  // occupancy_size[chan].store(0);
25  // occupancy_data[chan].store(0);
26  // Initialise buffers
27  // size_buffer[chan] = new uint64_t [RING_BUFFER_NUM];
28  // header_buffer[chan] = new packet_header [RING_BUFFER_NUM];
29  data_buffer[chan] = new int16_t[RING_BUFFER_LEN];
30  }
31 }
32 
33 // Try to push data to queue
34 int queue_buffer::try_push(int chan, int16_t* data, int n)
35 {
36  // Initialise the push values
37  uint64_t push_size = 0;
38  uint64_t push_len = 0;
39 
40  // Load the current tail pointers
41  current_tail[chan] = tail[chan].load();
42  // // Current size buffer pointer
43  // current_tail_size[chan] = make_int32_t(current_tail[chan].size_0,
44  // current_tail[chan].size_1);
45  // // Current data buffer pointer
46  // current_tail_data[chan] = make_int32_t(current_tail[chan].data_0,
47  // current_tail[chan].data_1);
48 
49  // Load the amount in buffer
50  write_num[chan] = occupancy[chan].load();
51 
52  // Get the available space in the buffer
53  uint64_t space = RING_BUFFER_LEN - write_num[chan];
54 
55  // If space is zero
56  if(space < n)
57  {
58  std::cout << "Space full" << std::endl;
59  // Return zero
60  return push_len;
61  }
62  // Else if n > space
63  else
64  {
65  // Push only avaliable space
66  push_len = n;
67  }
68 
69  // Get data index, push size and push length
70  push_size = push_len * sizeof(int16_t);
71 
72  // // Copy to size buffer
73  // // If the increase wraps the size buffer around
74  // if(increment(current_tail_size[chan],push_num,RING_BUFFER_NUM) < current_tail_size[chan]){
75  // // Get the distance to the end of the buffer
76  // uint64_t end = RING_BUFFER_NUM - current_tail_size[chan];
77  // // Memcpy to the end of buffer
78  // memcpy(&size_buffer[chan][current_tail_size[chan]],
79  // &sizes[SIZE_INDEX][index],
80  // end*sizeof(uint64_t));
81  // // Memcpy the rest to the start of the buffer
82  // memcpy(&size_buffer[chan][0],
83  // &sizes[SIZE_INDEX][index+end],
84  // (push_num-end)*sizeof(uint64_t));
85  // }
86  // // Else if if it doesn't wrap around
87  // else{
88  // // Memcpy all datagrams
89  // memcpy(&size_buffer[chan][current_tail_size[chan]],
90  // &sizes[SIZE_INDEX][index],push_num*sizeof(uint64_t));
91  // }
92 
93  // Copy to data buffer
94  // If the increase wraps the data buffer around
95  if(increment(current_tail[chan], push_len) < current_tail[chan])
96  {
97  // Get the distance to the end of the buffer
98  uint64_t end = RING_BUFFER_LEN - current_tail[chan];
99  // Memcpy to the end of buffer
100  memcpy(&data_buffer[chan][current_tail[chan]], data, end * sizeof(int16_t));
101  // Memcpy the rest to the start of the buffer
102  memcpy(&data_buffer[chan][0], &data[end], (push_len - end) * sizeof(int16_t));
103  }
104  // Else if if it doesn't wrap around
105  else
106  {
107  // Memcpy all datagrams
108  memcpy(&data_buffer[chan][current_tail[chan]], data, push_size);
109  }
110 
111  // Store the updated occupancy numbers
112  occupancy[chan].fetch_add(push_len);
113  // occupancy_size[chan].fetch_add(push_num);
114 
115  // // Get new current size buffer pointer values
116  // current_tail_size[chan] = increment(current_tail_size[chan],push_num,RING_BUFFER_NUM);
117  // int16_t size0 = current_tail_size[chan] & 0xFFFF;
118  // int16_t size1 = current_tail_size[chan] >> 16;;
119 
120  // // Get new current data buffer pointer
121  // current_tail_data[chan] = increment(current_tail_data[chan],push_len,RING_BUFFER_LEN);
122  // int16_t data0 = current_tail_data[chan] & 0xFFFF;
123  // int16_t data1 = current_tail_data[chan] >> 16;;
124 
125  // // Update tail pointer
126  // tail_[chan].store({size0,size1,data0,data1});
127  tail[chan].store(increment(current_tail[chan], push_len));
128 
129  // Return push_len
130  return push_len;
131 }
132 
133 // Push data to queue
134 void queue_buffer::push(int chan, int16_t* data, int n)
135 {
136  // If no packets to push, exit
137  if(n == 0)
138  return;
139 
140  // retval
141  uint64_t retval = 0;
142 
143  // Wait until pushed..
144  while(retval == 0)
145  {
146  // Try push
147  retval = try_push(chan, data, n);
148  };
149 
150  // We want to push n datagrams
151  int push_num = n;
152 
153  // // While datagrams to push > 0
154  // while(push_num > 0) {
155  // // Find the index in the array to push from
156  // int index = n - push_num;
157  // // Push, return the amount pushed and reclculate
158  // // how many left to push
159  // push_num -= try_push(chan,sizes,data,push_num,index);
160  // };
161 }
162 
163 // Try to pull data from queue
164 uint64_t queue_buffer::try_pull(int chan, int16_t*& data)
165 {
166  // Initialise the pull values
167  uint64_t pull_size = 0;
168  uint64_t pull_len = 0;
169 
170  // Load the current head pointers
171  current_head[chan] = head[chan].load();
172  // // Current size buffer pointer
173  // current_head_size[chan] = make_int32_t(current_head[chan].size_0,
174  // current_head[chan].size_1);
175  // // Current data buffer pointer
176  // current_head_data[chan] = make_int32_t(current_head[chan].data_0,
177  // current_head[chan].data_1);
178 
179  // Get number to read in buffer
180  // read_num_size[chan] = occupancy_size[chan].load();
181  read_num[chan] = occupancy[chan].load();
182 
183  // // Get the available data in the buffer
184  // pull_num = read_num_size[chan];
185 
186  // If nothing to read in the buffer
187  if(read_num[chan] == 0)
188  {
189  // Return unsuccesful pull
190  return read_num[chan]; //{pull_num,pull_len};
191  }
192  // Get pull length and pull size
193  pull_len = read_num[chan];
194  pull_size = pull_len * sizeof(int16_t);
195 
196  // // Copy from size buffer
197  // // If the increase wraps the buffer around
198  // if(increment(current_head_size[chan],pull_num,RING_BUFFER_NUM) < current_head_size[chan]){
199 
200  // // Get the distance to the end of the buffer
201  // uint64_t end = RING_BUFFER_NUM - current_head_size[chan];
202 
203  // // Memcpy from end of buffer
204  // memcpy(sizes,&size_buffer[chan][current_head_size[chan]],end*sizeof(uint64_t));
205  // // Memcpy from the start of the buffer
206  // memcpy(&sizes[end],&size_buffer[chan][0],(pull_num-end)*sizeof(uint64_t));
207 
208  // }
209  // // Else if if it doesn't wrap around
210  // else{
211  // // Memcpy all datagrams
212  // memcpy(sizes,&size_buffer[chan][current_head_size[chan]],pull_num*sizeof(uint64_t));
213  // }
214 
215  // Copy from data buffer
216  // If the increase wraps the buffer around
217  if(increment(current_head[chan], pull_len) < current_head[chan])
218  {
219  // Get the distance to the end of the buffer
220  uint64_t end = RING_BUFFER_LEN - current_head[chan];
221 
222  // Memcpy from end of buffer
223  memcpy(data, &data_buffer[chan][current_head[chan]], end * sizeof(int16_t));
224  // Memcpy from the start of the buffer
225  memcpy(&data[end], &data_buffer[chan][0], (pull_len - end) * sizeof(int16_t));
226  }
227  // Else if if it doesn't wrap around
228  else
229  {
230  // Memcpy all datagrams
231  memcpy(data, &data_buffer[chan][current_head[chan]], pull_size);
232  }
233 
234  // Store the updated number in buffer
235  occupancy[chan].fetch_sub(pull_len);
236  // occupancy_size[chan].fetch_sub(pull_num);
237 
238  // // Get new current size buffer pointer values
239  // current_head_size[chan] = increment(current_head_size[chan],pull_num,RING_BUFFER_NUM);
240  // int16_t size0 = current_head_size[chan] & 0xFFFF;
241  // int16_t size1 = current_head_size[chan] >> 16;;
242 
243  // // Get new current data buffer pointer
244  // current_head_data[chan] = increment(current_head_data[chan],pull_len,RING_BUFFER_LEN);
245  // int16_t data0 = current_head_data[chan] & 0xFFFF;
246  // int16_t data1 = current_head_data[chan] >> 16;;
247 
248  // // Update head pointer
249  // head_[chan].store({size0,size1,data0,data1});
250 
251  // Store the updated read pointer
252  head[chan].store(increment(current_head[chan], pull_len));
253 
254  // Return successful pull
255  // return {pull_num,pull_len};
256  return pull_len;
257 }
258 
259 // Pull data from queue
260 uint64_t queue_buffer::pull(bool* timeout, int chan, int16_t*& data)
261 {
262  // Number of datagrams pulled
263  // std::pair<uint64_t,uint64_t> pull_num = {0,0};
264  uint64_t pull_len = 0;
265  // Wait until...
266  while((pull_len = try_pull(chan, data)) == 0 && !*timeout) {};
267 
268  // if (*timeout and (write[chan].load() - current_head[chan]) != 0){
269  // std::cout << "ERROR: Still data in queue!!!!" << std::endl;
270  // exit(0);
271  // }
272 
273  // Return pulled data
274  return pull_len;
275 }
Definition: data.hh:4