otsdaq-mu2e-stm  5.02.01
queue.cc
1 #include <atomic>
2 #include <cinttypes>
3 #include <iostream>
4 #include <thread>
5 
6 // Queue buffer header
7 #include "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_SIZE];
19  }
20  // // Get the maximum number of datagrams to pull based on the data type
21  // if (is_same<T,UDPsocket::packet>::value){
22  // pull_max = UDPsocket::RECVMMSG_NUM;
23  // }
24  // else if (is_same<T,dataVars::off_event>::value){
25  // pull_max = UDPsocket::SENDMMSG_NUM;
26  // }
27 }
28 
29 // Try to push data to queue
30 int queue_buffer::try_push(int chan, int16_t* data, int n)
31 {
32  // Initialise the return value
33  uint64_t retval = 0;
34 
35  // Load the current tail
36  current_tail[chan] = write[chan].load();
37 
38  // Load the amount in buffer
39  write_num[chan] = occupancy[chan].load();
40 
41  // Get the available space in the buffer
42  uint64_t space = size - write_num[chan];
43 
44  // If no space for data to push
45  if(space <= n)
46  {
47  // Return zero
48  return retval;
49  }
50  // Else if space
51  else if(space > n)
52  {
53  // Push all n
54  retval = n;
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_SIZE - current_tail[chan];
62 
63  // Memcpy to the end of buffer
64  memcpy(&buffer[chan][current_tail[chan]], data, end * sizeof(int16_t));
65 
66  // Memcpy the rest to the start of the buffer
67  memcpy(&buffer[chan][0], &data[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, 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, int n)
88 {
89  // retval
90  uint64_t retval = 0;
91 
92  // Wait until pushed..
93  while(retval == 0)
94  {
95  // Try push
96  retval = try_push(chan, data, n);
97  };
98 }
99 
100 // Try to pull data from queue
101 int queue_buffer::try_pull(int chan, int16_t*& data)
102 {
103  // Initialise retval
104  uint64_t retval = 0;
105 
106  // Get current read pointer
107  current_head[chan] = read[chan].load();
108 
109  // Get number to read in buffer
110  read_num[chan] = occupancy[chan].load();
111 
112  // If nothing to read in the buffer
113  if((retval = read_num[chan]) == 0)
114  {
115  // Return unsuccesful pull
116  return retval;
117  }
118 
119  // If the increase wraps the buffer around
120  if(increment(current_head[chan], retval) < current_head[chan])
121  {
122  // Get the distance to the end of the buffer
123  uint64_t end = RING_BUFFER_SIZE - current_head[chan];
124 
125  // Memcpy from end of buffer
126  memcpy(data, &buffer[chan][current_head[chan]], end * sizeof(int16_t));
127  // Memcpy from the start of the buffer
128  memcpy(&data[end], &buffer[chan][0], (retval - end) * sizeof(int16_t));
129  }
130  // Else if if it doesn't wrap around
131  else
132  {
133  // Memcpy retval datagrams
134  memcpy(data, &buffer[chan][current_head[chan]], retval * sizeof(int16_t));
135  }
136 
137  // Store the updated number in buffer
138  occupancy[chan].fetch_sub(retval);
139 
140  // Store the updated read pointer
141  read[chan].store(increment(current_head[chan], retval));
142 
143  // Return successful pull
144  return retval;
145 }
146 
147 // Pull data from queue
148 int queue_buffer::pull(bool* timeout, int chan, int16_t*& data)
149 {
150  int retval = 0;
151 
152  // Wait until...
153  while((retval = try_pull(chan, data)) == 0 && !*timeout) {};
154 
155  if(*timeout and (write[chan].load() - current_head[chan]) != 0)
156  {
157  std::cout << "ERROR: Still data in queue!!!!" << std::endl;
158  exit(0);
159  }
160 
161  // Return pulled data
162  return retval;
163 }
Definition: data.hh:4