#include "naut/mpmc.h" #include "test.h" #include /* Many producers and many consumers pass tagged integers through the queue; * verify nothing is lost or duplicated. */ #define CAP 1024 #define NPROD 4 #define NCONS 4 #define PER_PROD 100000 static naut_mpmc q; static _Atomic long produced_sum, consumed_sum; static _Atomic int consumed_cnt; static _Atomic int prod_done; static void *producer(void *arg) { long base = (long)(intptr_t)arg * PER_PROD + 1; for (long i = 0; i < PER_PROD; i++) { long v = base + i; while (!naut_mpmc_push(&q, (void *)(intptr_t)v)) sched_yield(); atomic_fetch_add(&produced_sum, v); } atomic_fetch_add(&prod_done, 1); return NULL; } static void *consumer(void *arg) { (void)arg; void *p; for (;;) { if (naut_mpmc_pop(&q, &p)) { atomic_fetch_add(&consumed_sum, (long)(intptr_t)p); atomic_fetch_add(&consumed_cnt, 1); } else if (atomic_load(&prod_done) == NPROD) { if (!naut_mpmc_pop(&q, &p)) break; /* drained */ atomic_fetch_add(&consumed_sum, (long)(intptr_t)p); atomic_fetch_add(&consumed_cnt, 1); } else { sched_yield(); } } return NULL; } int main(void) { CHECK(naut_mpmc_init(&q, CAP) == NAUT_OK); CHECK(naut_mpmc_init(&q, 1000) == NAUT_ERR_INVAL); /* not pow2 */ pthread_t pr[NPROD], co[NCONS]; for (int i = 0; i < NCONS; i++) pthread_create(&co[i], NULL, consumer, NULL); for (int i = 0; i < NPROD; i++) pthread_create(&pr[i], NULL, producer, (void *)(intptr_t)i); for (int i = 0; i < NPROD; i++) pthread_join(pr[i], NULL); for (int i = 0; i < NCONS; i++) pthread_join(co[i], NULL); CHECK_EQ(atomic_load(&consumed_cnt), NPROD * PER_PROD); CHECK_EQ(atomic_load(&consumed_sum), atomic_load(&produced_sum)); naut_mpmc_destroy(&q); TEST_MAIN_END(); }