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_SIZE];
30 int queue_buffer::try_push(
int chan, int16_t*
data,
int n)
36 current_tail[chan] = write[chan].load();
39 write_num[chan] = occupancy[chan].load();
42 uint64_t space = size - write_num[chan];
58 if(increment(current_tail[chan], retval) < current_tail[chan])
61 uint64_t end = RING_BUFFER_SIZE - current_tail[chan];
64 memcpy(&buffer[chan][current_tail[chan]],
data, end *
sizeof(int16_t));
67 memcpy(&buffer[chan][0], &
data[end], (retval - end) *
sizeof(int16_t));
73 memcpy(&buffer[chan][current_tail[chan]],
data, 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,
int n)
96 retval = try_push(chan,
data, n);
101 int queue_buffer::try_pull(
int chan, int16_t*&
data)
107 current_head[chan] = read[chan].load();
110 read_num[chan] = occupancy[chan].load();
113 if((retval = read_num[chan]) == 0)
120 if(increment(current_head[chan], retval) < current_head[chan])
123 uint64_t end = RING_BUFFER_SIZE - current_head[chan];
126 memcpy(
data, &buffer[chan][current_head[chan]], end *
sizeof(int16_t));
128 memcpy(&
data[end], &buffer[chan][0], (retval - end) *
sizeof(int16_t));
134 memcpy(
data, &buffer[chan][current_head[chan]], retval *
sizeof(int16_t));
138 occupancy[chan].fetch_sub(retval);
141 read[chan].store(increment(current_head[chan], retval));
148 int queue_buffer::pull(
bool* timeout,
int chan, int16_t*&
data)
153 while((retval = try_pull(chan,
data)) == 0 && !*timeout) {};
155 if(*timeout and (write[chan].load() - current_head[chan]) != 0)
157 std::cout <<
"ERROR: Still data in queue!!!!" << std::endl;