while (1) {
struct util_queue_job job;
- pipe_semaphore_wait(&queue->queued);
- if (queue->kill_threads)
+ pipe_mutex_lock(queue->lock);
+ assert(queue->num_queued >= 0 && queue->num_queued <= queue->max_jobs);
+
+ /* wait if the queue is empty */
+ while (!queue->kill_threads && queue->num_queued == 0)
+ pipe_condvar_wait(queue->has_queued_cond, queue->lock);
+
+ if (queue->kill_threads) {
+ pipe_mutex_unlock(queue->lock);
break;
+ }
- pipe_mutex_lock(queue->lock);
job = queue->jobs[queue->read_idx];
queue->jobs[queue->read_idx].job = NULL;
queue->read_idx = (queue->read_idx + 1) % queue->max_jobs;
- pipe_mutex_unlock(queue->lock);
- pipe_semaphore_signal(&queue->has_space);
+ queue->num_queued--;
+ pipe_condvar_signal(queue->has_space_cond);
+ pipe_mutex_unlock(queue->lock);
if (job.job) {
queue->execute_job(job.job, thread_index);
queue->execute_job = execute_job;
pipe_mutex_init(queue->lock);
- pipe_semaphore_init(&queue->has_space, max_jobs);
- pipe_semaphore_init(&queue->queued, 0);
+
+ queue->num_queued = 0;
+ pipe_condvar_init(queue->has_queued_cond);
+ pipe_condvar_init(queue->has_space_cond);
queue->threads = (pipe_thread*)CALLOC(num_threads, sizeof(pipe_thread));
if (!queue->threads)
FREE(queue->threads);
if (queue->jobs) {
- pipe_semaphore_destroy(&queue->has_space);
- pipe_semaphore_destroy(&queue->queued);
+ pipe_condvar_destroy(queue->has_space_cond);
+ pipe_condvar_destroy(queue->has_queued_cond);
pipe_mutex_destroy(queue->lock);
FREE(queue->jobs);
}
unsigned i;
/* Signal all threads to terminate. */
- pipe_mutex_lock(queue->queued.mutex);
+ pipe_mutex_lock(queue->lock);
queue->kill_threads = 1;
- queue->queued.counter = queue->num_threads;
- pipe_condvar_broadcast(queue->queued.cond);
- pipe_mutex_unlock(queue->queued.mutex);
+ pipe_condvar_broadcast(queue->has_queued_cond);
+ pipe_mutex_unlock(queue->lock);
for (i = 0; i < queue->num_threads; i++)
pipe_thread_wait(queue->threads[i]);
- pipe_semaphore_destroy(&queue->has_space);
- pipe_semaphore_destroy(&queue->queued);
+ pipe_condvar_destroy(queue->has_space_cond);
+ pipe_condvar_destroy(queue->has_queued_cond);
pipe_mutex_destroy(queue->lock);
FREE(queue->jobs);
FREE(queue->threads);
assert(fence->signalled);
fence->signalled = false;
+ pipe_mutex_lock(queue->lock);
+ assert(queue->num_queued >= 0 && queue->num_queued <= queue->max_jobs);
+
/* if the queue is full, wait until there is space */
- pipe_semaphore_wait(&queue->has_space);
+ while (queue->num_queued == queue->max_jobs)
+ pipe_condvar_wait(queue->has_space_cond, queue->lock);
- pipe_mutex_lock(queue->lock);
ptr = &queue->jobs[queue->write_idx];
assert(ptr->job == NULL);
ptr->job = job;
ptr->fence = fence;
queue->write_idx = (queue->write_idx + 1) % queue->max_jobs;
+
+ queue->num_queued++;
+ pipe_condvar_signal(queue->has_queued_cond);
pipe_mutex_unlock(queue->lock);
- pipe_semaphore_signal(&queue->queued);
}