| git.druid.rocks | index | militantprogrammer | astral_canvas | src/ | lib/ | queue.c |
src/lib/queue.c
#include "lib/queue.h"
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <sys/eventfd.h>
Job *jobCreate(JobFn fn, void *data, int own)
{
Job *job = calloc(1, sizeof(*job));
if (!job) {
if (own && data)
free(data);
return NULL;
}
job->fn = fn;
job->data = data;
job->own = own;
return job;
}
void jobFree(Job *job)
{
if (!job)
return;
if (job->own && job->data)
free(job->data);
free(job);
}
int queueInit(JobQueue *q)
{
if (!q)
return -1;
memset(q, 0, sizeof(*q));
q->wake_fd = -1;
if (pthread_mutex_init(&q->mu, NULL) != 0)
return -1;
q->wake_fd = eventfd(0, EFD_CLOEXEC | EFD_NONBLOCK);
if (q->wake_fd < 0) {
pthread_mutex_destroy(&q->mu);
return -1;
}
q->ready = 1;
return 0;
}
void queueDestroy(JobQueue *q)
{
Job *job;
if (!q || !q->ready)
return;
pthread_mutex_lock(&q->mu);
while (q->head) {
job = q->head;
q->head = job->next;
jobFree(job);
}
q->tail = NULL;
pthread_mutex_unlock(&q->mu);
pthread_mutex_destroy(&q->mu);
if (q->wake_fd >= 0)
close(q->wake_fd);
q->wake_fd = -1;
q->ready = 0;
}
int queuePush(JobQueue *q, Job *job)
{
uint64_t one = 1;
if (!q || !job || !q->ready)
return -1;
pthread_mutex_lock(&q->mu);
job->next = NULL;
if (q->tail)
q->tail->next = job;
else
q->head = job;
q->tail = job;
pthread_mutex_unlock(&q->mu);
if (write(q->wake_fd, &one, sizeof one) < 0) {
/* counter is saturated or fd closed; pop still sees the node */
}
return 0;
}
Job *queuePop(JobQueue *q)
{
Job *job;
if (!q || !q->ready)
return NULL;
pthread_mutex_lock(&q->mu);
job = q->head;
if (job) {
q->head = job->next;
if (!q->head)
q->tail = NULL;
job->next = NULL;
}
pthread_mutex_unlock(&q->mu);
return job;
}
void queueDrainWake(JobQueue *q)
{
uint64_t value;
if (!q || q->wake_fd < 0)
return;
while (read(q->wake_fd, &value, sizeof value) > 0) {
}
}
int queueWakeFd(const JobQueue *q)
{
if (!q)
return -1;
return q->wake_fd;
}