ds: add a ringbuffer data structure
This commit is contained in:
+323
@@ -0,0 +1,323 @@
|
|||||||
|
#include <kernel/ringbuf.h>
|
||||||
|
#include <kernel/sched.h>
|
||||||
|
#include <kernel/thread.h>
|
||||||
|
#include <kernel/wait.h>
|
||||||
|
#include <magenta/status.h>
|
||||||
|
|
||||||
|
#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;
|
||||||
|
}
|
||||||
@@ -0,0 +1,45 @@
|
|||||||
|
#ifndef KERNEL_RINGBUF_H_
|
||||||
|
#define KERNEL_RINGBUF_H_
|
||||||
|
|
||||||
|
#include <kernel/wait.h>
|
||||||
|
#include <magenta/types.h>
|
||||||
|
|
||||||
|
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
|
||||||
Reference in New Issue
Block a user