otsdaq-mu2e-stm  5.02.01
queue_old.cc
1 #include <atomic>
2 #include <cinttypes>
3 #include <iostream>
4 #include <thread>
5 
6 // Queue buffer header
7 #include "STMDAQ-TestBeam/utils/queue.hh"
8 
9 // Constructor
10 queue_buffer::queue_buffer()
11 {
12  // For each channel, set read and write pointers to zero
13  for(int chan = 0; chan < CHNUM; chan++)
14  {
15  write[chan].store(0);
16  read[chan].store(0);
17  occupancy[chan].store(0);
18  buffer[chan] = new int16_t[RING_BUFFER_LEN];
19  // pull_max = UDPsocket::RECVMMSG_NUM*UDPsocket::MAX_UDP_LEN;
20  }
21 }
22 
23 // Try to push data to queue
24 uint64_t queue_buffer::try_push(int chan, int16_t* data, uint64_t n, uint64_t index)
25 {
26  // Initialise the return value
27  uint64_t retval = 0;
28 
29  // Load the current tail
30  current_tail[chan] = write[chan].load();
31 
32  // Load the amount in buffer
33  write_num[chan] = occupancy[chan].load();
34 
35  // Get the available space in the buffer
36  uint64_t space = RING_BUFFER_LEN - write_num[chan];
37 
38  // If space is zero
39  if(space == 0)
40  {
41  // Return zero
42  return retval;
43  }
44  // Else if n < space
45  else if(n < space)
46  {
47  // Push all n
48  retval = n;
49  }
50  // Else if n > space
51  else if(n > space)
52  {
53  // Push only avaliable space
54  retval = space;
55  }
56 
57  // If the increase wraps the buffer around
58  if(increment(current_tail[chan], retval) < current_tail[chan])
59  {
60  // Get the distance to the end of the buffer
61  uint64_t end = RING_BUFFER_LEN - current_tail[chan];
62 
63  // Memcpy to the end of buffer
64  memcpy(&buffer[chan][current_tail[chan]], &data[index], end * sizeof(int16_t));
65 
66  // Memcpy the rest to the start of the buffer
67  memcpy(&buffer[chan][0], &data[index + end], (retval - end) * sizeof(int16_t));
68  }
69  // Else if if it doesn't wrap around
70  else
71  {
72  // Memcpy retval datagrams
73  memcpy(&buffer[chan][current_tail[chan]], &data[index], retval * sizeof(int16_t));
74  }
75 
76  // Store the updated number in buffer
77  occupancy[chan].fetch_add(retval);
78 
79  // Update the write pointer
80  write[chan].store(increment(current_tail[chan], retval));
81 
82  // Return retval
83  return retval;
84 }
85 
86 // Push data to queue
87 void queue_buffer::push(int chan, int16_t* data, uint64_t n)
88 {
89  // We want to push n datagrams
90  uint64_t retval = n;
91 
92  // While datagrams to push > 0
93  while(retval > 0)
94  {
95  // Find the index in the array to push from
96  uint64_t index = n - retval;
97  // Push, return the amount pushed and reclculate
98  // how many left to push
99  retval -= try_push(chan, data, retval, index);
100  std::cout << "Pushed " << n - retval << "/" << n << std::endl;
101  };
102 }
103 
104 // Try to pull data from queue
105 uint64_t queue_buffer::try_pull(int chan, int16_t*& data)
106 {
107  // Initialise retval
108  uint64_t retval = 0;
109 
110  // Get current read pointer
111  current_head[chan] = read[chan].load();
112 
113  // Get number to read in buffer
114  read_num[chan] = occupancy[chan].load();
115 
116  // If nothing to read in the buffer
117  if((retval = read_num[chan]) == 0)
118  {
119  // Return unsuccesful pull
120  return retval;
121  }
122 
123  // Make sure we don't try and pull more than the push_max
124  if(retval > pull_max)
125  retval = pull_max;
126 
127  // If the increase wraps the buffer around
128  if(increment(current_head[chan], retval) < current_head[chan])
129  {
130  // Get the distance to the end of the buffer
131  uint64_t end = RING_BUFFER_LEN - current_head[chan];
132 
133  // Memcpy from end of buffer
134  memcpy(data, &buffer[chan][current_head[chan]], end * sizeof(int16_t));
135  // Memcpy from the start of the buffer
136  memcpy(&data[end], &buffer[chan][0], (retval - end) * sizeof(int16_t));
137  }
138  // Else if if it doesn't wrap around
139  else
140  {
141  // Memcpy retval datagrams
142  memcpy(data, &buffer[chan][current_head[chan]], retval * sizeof(int16_t));
143  }
144 
145  // Store the updated number in buffer
146  occupancy[chan].fetch_sub(retval);
147 
148  // Store the updated read pointer
149  read[chan].store(increment(current_head[chan], retval));
150 
151  // Return successful pull
152  return retval;
153 }
154 
155 // Pull data from queue
156 uint64_t queue_buffer::pull(bool* timeout, int chan, int16_t*& data)
157 {
158  uint64_t retval = 0;
159 
160  // Wait until...
161  while((retval = try_pull(chan, data)) == 0 && !*timeout) {};
162 
163  if(*timeout and (write[chan].load() - current_head[chan]) != 0)
164  {
165  cout << "ERROR: Still data in queue!!!!" << endl;
166  exit(0);
167  }
168 
169  // Return pulled data
170  return retval;
171 }
Definition: data.hh:4