From 249ac7a8cf8928a25b061dc15fc191d5fa54f98b Mon Sep 17 00:00:00 2001 From: Max Wash Date: Sat, 6 Jun 2026 17:50:23 +0100 Subject: [PATCH] ds: add a ringbuffer data structure --- ds/ringbuf.c | 323 +++++++++++++++++++++++++++++++++++++++ include/kernel/ringbuf.h | 45 ++++++ 2 files changed, 368 insertions(+) create mode 100644 ds/ringbuf.c create mode 100644 include/kernel/ringbuf.h diff --git a/ds/ringbuf.c b/ds/ringbuf.c new file mode 100644 index 0000000..44dd5a5 --- /dev/null +++ b/ds/ringbuf.c @@ -0,0 +1,323 @@ +#include +#include +#include +#include +#include + +#define BUF_LOCKED(buf) \ + (((buf)->buf_flags & (RINGBUF_READ_LOCKED | RINGBUF_WRITE_LOCKED)) != 0) + +enum ringbuf_flags { + RINGBUF_READ_LOCKED = 0x01u, + RINGBUF_WRITE_LOCKED = 0x02u, +}; + +kern_status_t ringbuf_lock(struct ringbuf *buf, unsigned long *flags) +{ + spin_lock_irqsave(&buf->buf_lock, flags); + return KERN_OK; +} + +kern_status_t ringbuf_unlock(struct ringbuf *buf, unsigned long flags) +{ + spin_unlock_irqrestore(&buf->buf_lock, flags); + return KERN_OK; +} + +static kern_status_t ringbuf_clear(struct ringbuf *buf) +{ + buf->buf_read = buf->buf_write = 0; + return KERN_OK; +} + +static size_t ringbuf_write_capacity_remaining(const struct ringbuf *buf) +{ + if (buf->buf_read > buf->buf_write) { + return buf->buf_read - buf->buf_write - 1; + } else { + return buf->buf_capacity - buf->buf_write + buf->buf_read - 1; + } +} + +static size_t ringbuf_available_data_remaining(const struct ringbuf *buf) +{ + if (buf->buf_read < buf->buf_write) { + return buf->buf_write - buf->buf_read; + } else if (buf->buf_read > buf->buf_write) { + return buf->buf_capacity - buf->buf_read + buf->buf_write; + } else { + return 0; + } +} + +static kern_status_t ringbuf_open_read_buffer( + struct ringbuf *buf, + const void **ptr, + size_t *length) +{ + if (BUF_LOCKED(buf)) { + return KERN_BUSY; + } + + size_t contiguous_capacity = 0; + if (buf->buf_read > buf->buf_write) { + contiguous_capacity = buf->buf_capacity - buf->buf_read; + } else { + contiguous_capacity = buf->buf_write - buf->buf_read; + } + + if (contiguous_capacity == 0) { + return KERN_NO_DATA; + } + + buf->buf_opened_ptr = (unsigned char *)buf->buf_ptr + buf->buf_read; + buf->buf_opened_capacity = contiguous_capacity; + + buf->buf_flags |= RINGBUF_READ_LOCKED; + *ptr = buf->buf_opened_ptr; + *length = contiguous_capacity; + + return KERN_OK; +} + +static kern_status_t ringbuf_close_read_buffer( + struct ringbuf *buf, + const void **ptr, + size_t nbuf_read) +{ + if (!(buf->buf_flags & RINGBUF_READ_LOCKED)) { + return KERN_BAD_STATE; + } + + if (*ptr != buf->buf_opened_ptr) { + return KERN_INVALID_ARGUMENT; + } + + if (nbuf_read > buf->buf_opened_capacity) { + return KERN_INVALID_ARGUMENT; + } + + buf->buf_read += nbuf_read; + if (buf->buf_read >= buf->buf_capacity) { + buf->buf_read = 0; + } + + if (buf->buf_read == buf->buf_write) { + /* the ringbuf is now empty. set both pointers to 0. + * this ensures that the whole buffer will be available + * contiguously to the next call to open_write_buffer */ + buf->buf_read = 0; + buf->buf_write = 0; + } + + buf->buf_opened_ptr = NULL; + buf->buf_opened_capacity = 0; + buf->buf_flags &= ~RINGBUF_READ_LOCKED; + + return KERN_OK; +} + +static kern_status_t ringbuf_open_write_buffer( + struct ringbuf *buf, + void **ptr, + size_t *capacity) +{ + if (BUF_LOCKED(buf)) { + return KERN_BUSY; + } + + size_t contiguous_capacity = 0; + if (buf->buf_write >= buf->buf_read) { + contiguous_capacity = buf->buf_capacity - buf->buf_write - 1; + + if (buf->buf_read > 0) { + contiguous_capacity++; + } + } else { + contiguous_capacity = buf->buf_read - buf->buf_write - 1; + } + + if (contiguous_capacity == 0) { + return KERN_NO_SPACE; + } + + buf->buf_opened_ptr = (unsigned char *)buf->buf_ptr + buf->buf_write; + buf->buf_opened_capacity = contiguous_capacity; + + buf->buf_flags |= RINGBUF_WRITE_LOCKED; + *ptr = buf->buf_opened_ptr; + *capacity = contiguous_capacity; + + return KERN_OK; +} + +static kern_status_t ringbuf_close_write_buffer( + struct ringbuf *buf, + void **ptr, + size_t nbuf_written) +{ + if (!(buf->buf_flags & RINGBUF_WRITE_LOCKED)) { + return KERN_BAD_STATE; + } + + if (*ptr != buf->buf_opened_ptr) { + return KERN_INVALID_ARGUMENT; + } + + if (nbuf_written > buf->buf_opened_capacity) { + return KERN_INVALID_ARGUMENT; + } + + buf->buf_write += nbuf_written; + if (buf->buf_write >= buf->buf_capacity) { + buf->buf_write = 0; + } + + buf->buf_opened_ptr = NULL; + buf->buf_opened_capacity = 0; + buf->buf_flags &= ~RINGBUF_WRITE_LOCKED; + + return KERN_OK; +} + +static void wait_for_data(struct ringbuf *buf, unsigned long *flags) +{ + struct thread *self = get_current_thread(); + struct wait_item waiter; + wait_item_init(&waiter, self); + for (;;) { + thread_wait_begin(&waiter, &buf->buf_read_queue); + if (ringbuf_available_data_remaining(buf) > 0) { + break; + } + ringbuf_unlock(buf, *flags); + schedule(SCHED_NORMAL); + ringbuf_lock(buf, flags); + } + thread_wait_end(&waiter, &buf->buf_read_queue); + put_current_thread(self); +} + +static void wait_for_capacity(struct ringbuf *buf, unsigned long *flags) +{ + struct thread *self = get_current_thread(); + struct wait_item waiter; + wait_item_init(&waiter, self); + for (;;) { + thread_wait_begin(&waiter, &buf->buf_write_queue); + if (ringbuf_write_capacity_remaining(buf) > 0) { + break; + } + ringbuf_unlock(buf, *flags); + schedule(SCHED_NORMAL); + ringbuf_lock(buf, flags); + } + thread_wait_end(&waiter, &buf->buf_write_queue); + put_current_thread(self); +} + +kern_status_t ringbuf_read( + struct ringbuf *buf, + void *p, + size_t count, + size_t *nbuf_read, + unsigned long *irq_flags) +{ + if (BUF_LOCKED(buf)) { + return KERN_BUSY; + } + + size_t r = 0; + unsigned char *dest = p; + size_t remaining = count; + kern_status_t status = KERN_OK; + + while (remaining > 0) { + const void *src; + size_t available; + wait_for_data(buf, irq_flags); + + status = ringbuf_open_read_buffer(buf, &src, &available); + + if (status != KERN_OK) { + break; + } + + size_t to_copy = remaining; + if (to_copy > available) { + to_copy = available; + } + + memcpy(dest, src, to_copy); + + remaining -= to_copy; + dest += to_copy; + r += to_copy; + + ringbuf_close_read_buffer(buf, &src, to_copy); + } + + wakeup_queue(&buf->buf_write_queue); + if (nbuf_read) { + *nbuf_read = r; + } + + if (status == KERN_NO_DATA && r > 0) { + status = KERN_OK; + } + + return KERN_OK; +} + +kern_status_t ringbuf_write( + struct ringbuf *buf, + const void *p, + size_t count, + size_t *nbuf_written, + unsigned long *irq_flags) +{ + if (BUF_LOCKED(buf)) { + return KERN_BUSY; + } + + size_t w = 0; + const unsigned char *src = p; + size_t remaining = count; + kern_status_t status = KERN_OK; + + while (remaining > 0) { + void *dest; + size_t available; + wait_for_capacity(buf, irq_flags); + + status = ringbuf_open_write_buffer(buf, &dest, &available); + + if (status == KERN_NO_SPACE) { + break; + } + + size_t to_copy = remaining; + if (to_copy > available) { + to_copy = available; + } + + memcpy(dest, src, to_copy); + + remaining -= to_copy; + src += to_copy; + w += to_copy; + + ringbuf_close_write_buffer(buf, &dest, to_copy); + } + + wakeup_queue(&buf->buf_read_queue); + if (nbuf_written) { + *nbuf_written = w; + } + + if (status == KERN_NO_SPACE && w > 0) { + status = KERN_OK; + } + + return status; +} diff --git a/include/kernel/ringbuf.h b/include/kernel/ringbuf.h new file mode 100644 index 0000000..944cdf9 --- /dev/null +++ b/include/kernel/ringbuf.h @@ -0,0 +1,45 @@ +#ifndef KERNEL_RINGBUF_H_ +#define KERNEL_RINGBUF_H_ + +#include +#include + +struct ringbuf { + unsigned int buf_read; + unsigned int buf_write; + unsigned int buf_capacity; + + unsigned int buf_flags; + unsigned char *buf_ptr; + struct waitqueue buf_read_queue; + struct waitqueue buf_write_queue; + spin_lock_t buf_lock; + + unsigned char *buf_opened_ptr; + unsigned int buf_opened_capacity; +}; + +#define RINGBUF_DECLARE(name, size) \ + static unsigned char __buf_##name[size] = {0}; \ + static struct ringbuf name = { \ + .buf_ptr = __buf_##name, \ + .buf_capacity = size, \ + } + +extern kern_status_t ringbuf_lock(struct ringbuf *buf, unsigned long *flags); +extern kern_status_t ringbuf_unlock(struct ringbuf *buf, unsigned long flags); + +extern kern_status_t ringbuf_read( + struct ringbuf *buf, + void *out, + size_t count, + size_t *nr_read, + unsigned long *irq_flags); +extern kern_status_t ringbuf_write( + struct ringbuf *buf, + const void *out, + size_t count, + size_t *nr_written, + unsigned long *irq_flags); + +#endif