#include "../include/al/lib.h" #include "../include/al/ring_buffer.h" // https://github.com/MusicPlayerDaemon/MPD/blob/master/src/util/RingBuffer.hxx // 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, size_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, volatile u8 **v, u8 *ptr, size_t n) { ptr += n; al_assert(ptr <= buf->end); if (ptr == buf->end) ptr = buf->start; al_atomic_ptr_store(v, ptr, AL_ATOMIC_RELEASE); } size_t al_ring_buffer_space(struct al_ring_buffer *buf) { u8 *wp = al_atomic_ptr_load(&buf->write, AL_ATOMIC_RELAXED); u8 *rp = previous(buf, (u8 *)al_atomic_ptr_load(&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, size_t *size) { u8 *wp = al_atomic_ptr_load(&buf->write, AL_ATOMIC_ACQUIRE); u8 *rp = previous(buf, al_atomic_ptr_load(&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, size_t n) { add(buf, &buf->write, ptr, n); } size_t al_ring_buffer_occupied(struct al_ring_buffer *buf) { u8 *rp = al_atomic_ptr_load(&buf->read, AL_ATOMIC_RELAXED); u8 *wp = al_atomic_ptr_load(&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, size_t *size) { u8 *rp = al_atomic_ptr_load(&buf->read, AL_ATOMIC_ACQUIRE); u8 *wp = al_atomic_ptr_load(&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, size_t n) { add(buf, &buf->read, ptr, n); } size_t al_ring_buffer_write(struct al_ring_buffer *buf, u8 *data, size_t n) { u8 *wp = al_atomic_ptr_load(&buf->write, AL_ATOMIC_ACQUIRE); u8 *rp = previous(buf, al_atomic_ptr_load(&buf->read, AL_ATOMIC_RELAXED)); size_t size = AL_MIN((wp <= rp ? rp : buf->end) - wp, (ptrdiff_t)n); al_memcpy(wp, data, size); wp += size; if (wp >= buf->end) { size_t wrap = AL_MIN(rp - buf->start, (ptrdiff_t)(n - size)); al_memcpy(buf->start, data + size, wrap); wp = buf->start + wrap; size += wrap; } al_atomic_ptr_store(&buf->write, wp, AL_ATOMIC_RELEASE); return size; } size_t al_ring_buffer_read(struct al_ring_buffer *buf, u8 *ptr, size_t n) { u8 *rp = al_atomic_ptr_load(&buf->read, AL_ATOMIC_ACQUIRE); u8 *wp = al_atomic_ptr_load(&buf->write, AL_ATOMIC_RELAXED); size_t size = AL_MIN((rp <= wp ? wp : buf->end) - rp, (ptrdiff_t)n); al_memcpy(ptr, rp, size); rp += size; if (rp >= buf->end) { size_t wrap = AL_MIN(wp - buf->start, (ptrdiff_t)(n - size)); al_memcpy(ptr + size, buf->start, wrap); rp = buf->start + wrap; size += wrap; } al_atomic_ptr_store(&buf->read, rp, AL_ATOMIC_RELEASE); return size; } size_t al_ring_buffer_peek(struct al_ring_buffer *buf, u8 *ptr, size_t offset, size_t n) { u8 *rp = al_atomic_ptr_load(&buf->read, AL_ATOMIC_RELAXED); u8 *wp = al_atomic_ptr_load(&buf->write, AL_ATOMIC_RELAXED); ptrdiff_t size; if (rp <= wp) { rp += offset; size = AL_MIN(wp - rp, (ptrdiff_t)n); if (size < 0) return 0; al_memcpy(ptr, rp, size); } else { rp += offset; size = AL_MIN(buf->end - rp, (ptrdiff_t)n); if (size > 0) { al_memcpy(ptr, rp, size); n -= size; rp = buf->start; } else { rp = buf->start - size; size = 0; } ptrdiff_t wrap = AL_MIN(wp - rp, (ptrdiff_t)n); if (wrap > 0) { al_memcpy(ptr + size, rp, wrap); size += wrap; } } return size; } size_t al_ring_buffer_discard(struct al_ring_buffer *buf, size_t n) { u8 *rp = al_atomic_ptr_load(&buf->read, AL_ATOMIC_ACQUIRE); u8 *wp = al_atomic_ptr_load(&buf->write, AL_ATOMIC_RELAXED); size_t discard = AL_MIN((rp <= wp ? wp : buf->end) - rp, (ptrdiff_t)n); rp += discard; if (rp >= buf->end) { size_t wrap = AL_MIN(wp - buf->start, (ptrdiff_t)(n - discard)); rp = buf->start + wrap; discard += wrap; } al_atomic_ptr_store(&buf->read, rp, AL_ATOMIC_RELEASE); return discard; } // Not thread-safe. void al_ring_buffer_reset(struct al_ring_buffer *buf) { al_atomic_ptr_store(&buf->read, buf->start, AL_ATOMIC_RELAXED); al_atomic_ptr_store(&buf->write, buf->start, AL_ATOMIC_RELAXED); }