#include "../include/al/ring_buffer.h" #include "../include/al/atomic.h" #include "../include/al/macros.h" // https://github.com/MusicPlayerDaemon/MPD/blob/master/src/util/RingBuffer.hxx // // This is no longer an implementation of a "contiguous" ring buffer but I used this as reference. // Wrapping is handled in read() and write(). // https://andrea.lattuada.me/blog/2019/the-design-and-implementation-of-a-lock-free-ring-buffer-with-contiguous-reservations.html void al_ring_buffer_init(struct al_ring_buffer *buf, u8 *data, ptrdiff_t length) { buf->start = data; buf->end = data + length; al_ring_buffer_reset(buf); } static inline u8 *previous(struct al_ring_buffer *buf, u8 *ptr) { if (ptr == buf->start) { ptr = buf->end; } return --ptr; } static inline void add(struct al_ring_buffer *buf, u8 **v, u8 *ptr, ptrdiff_t n) { ptr += n; al_assert(ptr <= buf->end); if (ptr == buf->end) ptr = buf->start; al_atomic_store(void)(v, ptr, AL_ATOMIC_RELEASE); } ptrdiff_t al_ring_buffer_space(struct al_ring_buffer *buf) { u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_RELAXED); u8 *rp = previous(buf, (u8 *)al_atomic_load(void)(&buf->read, AL_ATOMIC_RELAXED)); return (wp <= rp) ? rp - wp : (buf->end - wp) + (rp - buf->start); } u8 *al_ring_buffer_write_chunk(struct al_ring_buffer *buf, ptrdiff_t *size) { u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_ACQUIRE); u8 *rp = previous(buf, al_atomic_load(void)(&buf->read, AL_ATOMIC_RELAXED)); *size = (wp <= rp ? rp : buf->end) - wp; return wp; } void al_ring_buffer_append(struct al_ring_buffer *buf, u8 *ptr, ptrdiff_t n) { add(buf, &buf->write, ptr, n); } ptrdiff_t al_ring_buffer_occupied(struct al_ring_buffer *buf) { u8 *rp = al_atomic_load(void)(&buf->read, AL_ATOMIC_RELAXED); u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_RELAXED); return (rp <= wp) ? wp - rp : (buf->end - rp) + (wp - buf->start); } u8 *al_ring_buffer_read_chunk(struct al_ring_buffer *buf, ptrdiff_t *size) { u8 *rp = al_atomic_load(void)(&buf->read, AL_ATOMIC_ACQUIRE); u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_RELAXED); *size = (rp <= wp ? wp : buf->end) - rp; return rp; } void al_ring_buffer_consume(struct al_ring_buffer *buf, u8 *ptr, ptrdiff_t n) { add(buf, &buf->read, ptr, n); } ptrdiff_t al_ring_buffer_write(struct al_ring_buffer *buf, u8 *data, ptrdiff_t n) { u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_ACQUIRE); u8 *rp = previous(buf, al_atomic_load(void)(&buf->read, AL_ATOMIC_RELAXED)); ptrdiff_t size, wrap; size = MIN((wp <= rp ? rp : buf->end) - wp, n); al_memcpy(wp, data, (size_t)size); wp += size; if (wp >= buf->end) { wrap = MIN(rp - buf->start, n - size); al_memcpy(buf->start, data + size, (size_t)wrap); wp = buf->start + wrap; size += wrap; } al_atomic_store(void)(&buf->write, wp, AL_ATOMIC_RELEASE); return size; } ptrdiff_t al_ring_buffer_read(struct al_ring_buffer *buf, u8 *ptr, ptrdiff_t n) { u8 *rp = al_atomic_load(void)(&buf->read, AL_ATOMIC_ACQUIRE); u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_RELAXED); ptrdiff_t size, wrap; size = MIN((rp <= wp ? wp : buf->end) - rp, n); al_memcpy(ptr, rp, (size_t)size); rp += size; if (rp >= buf->end) { wrap = MIN(wp - buf->start, n - size); al_memcpy(ptr + size, buf->start, (size_t)wrap); rp = buf->start + wrap; size += wrap; } al_atomic_store(void)(&buf->read, rp, AL_ATOMIC_RELEASE); return size; } ptrdiff_t al_ring_buffer_peek(struct al_ring_buffer *buf, u8 *ptr, ptrdiff_t offset, ptrdiff_t n) { u8 *rp = al_atomic_load(void)(&buf->read, AL_ATOMIC_RELAXED); u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_RELAXED); ptrdiff_t size, wrap; if (rp <= wp) { rp += offset; size = MIN(wp - rp, n); if (size < 0) { return 0; } al_memcpy(ptr, rp, (size_t)size); } else { rp += offset; size = MIN(buf->end - rp, n); if (size > 0) { al_memcpy(ptr, rp, (size_t)size); n -= size; rp = buf->start; } else { rp = buf->start - size; size = 0; } wrap = MIN(wp - rp, n); if (wrap > 0) { al_memcpy(ptr + size, rp, (size_t)wrap); size += wrap; } } return size; } ptrdiff_t al_ring_buffer_discard(struct al_ring_buffer *buf, ptrdiff_t n) { u8 *rp = al_atomic_load(void)(&buf->read, AL_ATOMIC_ACQUIRE); u8 *wp = al_atomic_load(void)(&buf->write, AL_ATOMIC_RELAXED); ptrdiff_t discard, wrap; discard = MIN((rp <= wp ? wp : buf->end) - rp, n); rp += discard; if (rp >= buf->end) { wrap = MIN(wp - buf->start, n - discard); rp = buf->start + wrap; discard += wrap; } al_atomic_store(void)(&buf->read, rp, AL_ATOMIC_RELEASE); return discard; } // Not thread-safe. void al_ring_buffer_reset(struct al_ring_buffer *buf) { al_atomic_store(void)(&buf->read, buf->start, AL_ATOMIC_RELAXED); al_atomic_store(void)(&buf->write, buf->start, AL_ATOMIC_RELAXED); }