hex-workqueue: fix race condition in seqn getting out of sync with idx_read/write (#29785)

This commit is contained in:
Max Krasnyansky
2026-10-01 08:35:52 +03:00
committed by GitHub
parent b0aca3c653
commit 3aa0ce9bca
+16 -17
View File
@@ -31,7 +31,7 @@ struct work_queue_task_s {
struct work_queue_s {
atomic_uint seqn; // seqno used to detect new jobs
atomic_uint idx_read; // Updated by producer (pop/reclaim)
unsigned int idx_write; // Updated by producer (push)
atomic_uint idx_write; // Updated by producer (push)
uint32_t idx_mask;
uint32_t capacity;
@@ -41,7 +41,7 @@ struct work_queue_s {
unsigned int n_threads; // total threads (workers + main)
unsigned int n_workers; // number of active threads (just workers)
atomic_bool active; // workers are polling/active
atomic_uint active; // workers are polling/active
atomic_bool killed; // threads need to exit
bool external_mem; // memory owned externally
@@ -62,7 +62,7 @@ static void work_queue_thread(void * context) {
if (atomic_load_explicit(&q->active, memory_order_relaxed)) {
hex_pause();
} else {
qurt_futex_wait(&q->seqn, prev_seqn);
qurt_futex_wait(&q->active, 0);
}
continue;
}
@@ -71,7 +71,7 @@ static void work_queue_thread(void * context) {
// Process all active tasks in the queue
unsigned int ir = atomic_load_explicit(&q->idx_read, memory_order_relaxed);
unsigned int iw = q->idx_write;
unsigned int iw = atomic_load_explicit(&q->idx_write, memory_order_relaxed);
while (ir != iw) {
struct work_queue_task_s * task = &q->queue[ir];
@@ -101,8 +101,8 @@ bool work_queue_run_async(work_queue_t q, work_queue_func_t func, void * data, u
return false;
}
unsigned int ir = atomic_load_explicit(&q->idx_read, memory_order_relaxed);
unsigned int iw = q->idx_write;
unsigned int ir = atomic_load_explicit(&q->idx_read, memory_order_relaxed);
unsigned int iw = atomic_load_explicit(&q->idx_write, memory_order_relaxed);
if (((iw + 1) & q->idx_mask) == ir) {
FARF(ERROR, "work-queue-push: queue is full\n");
@@ -115,7 +115,7 @@ bool work_queue_run_async(work_queue_t q, work_queue_func_t func, void * data, u
task->n_threads = n;
atomic_store_explicit(&task->barrier, n, memory_order_relaxed);
q->idx_write = (iw + 1) & q->idx_mask;
atomic_store_explicit(&q->idx_write, (iw + 1) & q->idx_mask, memory_order_relaxed);
// publish job to workers (already awake and polling)
atomic_fetch_add_explicit(&q->seqn, 1, memory_order_release);
@@ -168,10 +168,10 @@ work_queue_t work_queue_init(void * ptr, uint32_t n_threads, uint32_t capacity,
q->context[i].queue = q;
}
atomic_init(&q->idx_read, 0);
atomic_init(&q->seqn, 0);
atomic_init(&q->active, false);
q->idx_write = 0;
atomic_init(&q->idx_read, 0);
atomic_init(&q->idx_write, 0);
atomic_init(&q->seqn, 0);
atomic_init(&q->active, 0);
q->idx_mask = capacity - 1;
q->killed = 0;
for (int i = 0; i < (int) capacity; i++) {
@@ -219,8 +219,9 @@ void work_queue_free(work_queue_t q) {
if (!q) { return; }
atomic_store_explicit(&q->killed, 1, memory_order_relaxed);
atomic_store_explicit(&q->active, 1, memory_order_release);
atomic_fetch_add_explicit(&q->seqn, 1, memory_order_release);
qurt_futex_wake(&q->seqn, q->n_workers);
qurt_futex_wake(&q->active, q->n_workers);
for (unsigned int i = 0; i < q->n_workers; i++) {
if (q->thread[i]) {
@@ -232,13 +233,11 @@ void work_queue_free(work_queue_t q) {
void work_queue_wakeup(work_queue_t q) {
if (!atomic_load_explicit(&q->active, memory_order_relaxed)) {
atomic_store_explicit(&q->active, true, memory_order_release);
// Increment seqn and wake workers to transition them out of sleep
atomic_fetch_add_explicit(&q->seqn, 1, memory_order_release);
qurt_futex_wake(&q->seqn, q->n_workers);
atomic_store_explicit(&q->active, 1, memory_order_release);
qurt_futex_wake(&q->active, q->n_workers);
}
}
void work_queue_suspend(work_queue_t q) {
atomic_store_explicit(&q->active, false, memory_order_release);
atomic_store_explicit(&q->active, 0, memory_order_release);
}