7 #include "STMDAQ-TestBeam/utils/queue.hh"
10 queue_buffer::queue_buffer()
13 for(
int chan = 0; chan < CHNUM; chan++)
17 occupancy[chan].store(0);
18 buffer[chan] =
new int16_t[RING_BUFFER_LEN];
24 uint64_t queue_buffer::try_push(
int chan, int16_t*
data, uint64_t n, uint64_t index)
30 current_tail[chan] = write[chan].load();
33 write_num[chan] = occupancy[chan].load();
36 uint64_t space = RING_BUFFER_LEN - write_num[chan];
58 if(increment(current_tail[chan], retval) < current_tail[chan])
61 uint64_t end = RING_BUFFER_LEN - current_tail[chan];
64 memcpy(&buffer[chan][current_tail[chan]], &
data[index], end *
sizeof(int16_t));
67 memcpy(&buffer[chan][0], &
data[index + end], (retval - end) *
sizeof(int16_t));
73 memcpy(&buffer[chan][current_tail[chan]], &
data[index], retval *
sizeof(int16_t));
77 occupancy[chan].fetch_add(retval);
80 write[chan].store(increment(current_tail[chan], retval));
87 void queue_buffer::push(
int chan, int16_t*
data, uint64_t n)
96 uint64_t index = n - retval;
99 retval -= try_push(chan,
data, retval, index);
100 std::cout <<
"Pushed " << n - retval <<
"/" << n << std::endl;
105 uint64_t queue_buffer::try_pull(
int chan, int16_t*&
data)
111 current_head[chan] = read[chan].load();
114 read_num[chan] = occupancy[chan].load();
117 if((retval = read_num[chan]) == 0)
124 if(retval > pull_max)
128 if(increment(current_head[chan], retval) < current_head[chan])
131 uint64_t end = RING_BUFFER_LEN - current_head[chan];
134 memcpy(
data, &buffer[chan][current_head[chan]], end *
sizeof(int16_t));
136 memcpy(&
data[end], &buffer[chan][0], (retval - end) *
sizeof(int16_t));
142 memcpy(
data, &buffer[chan][current_head[chan]], retval *
sizeof(int16_t));
146 occupancy[chan].fetch_sub(retval);
149 read[chan].store(increment(current_head[chan], retval));
156 uint64_t queue_buffer::pull(
bool* timeout,
int chan, int16_t*&
data)
161 while((retval = try_pull(chan,
data)) == 0 && !*timeout) {};
163 if(*timeout and (write[chan].load() - current_head[chan]) != 0)
165 cout <<
"ERROR: Still data in queue!!!!" << endl;