Naut/tests/unit/test_mpmc.c
ookami125 2178d6a70c Initial commit: Naut-Torrent — from-scratch 10 GbE BitTorrent client
A maintainable, extensible BitTorrent client (C11, Linux/io_uring) targeting
10 GbE saturation. All torrent functionality is built from scratch; liburing
is the only linked third-party dependency on the data path.

Implements Phases 1-7 of the roadmap:
- core: page-aligned buffer pool, MPMC/Treiber queues, bitfields, worker pool
- crypto: SHA-1/256 (SHA-NI + scalar), Merkle (BEP-52), RC4 (MSE)
- bencode/metainfo: zero-copy parser, v1/v2/hybrid .torrent + magnet
- peer: sans-IO wire codec, MSE/PE handshake state machine, BEP-10, ut_metadata, PEX
- piece/storage: block-level multi-peer engine, rarest-first + endgame,
  per-file completion events + single-file relocate (move-as-you-finish)
- tracker/dht: HTTP + UDP (BEP-15) trackers, BEP-5 KRPC iterative lookup
- platform: io_uring reactor (SQPOLL, registered buffers, SEND_ZC)
- surface: versioned RPC, native plugin ABI, sandboxed Lua scripting, nautd/nautctl

Verified against libtorrent (single/multi/hybrid, MSE, magnet-via-DHT, swarm);
unit + interop tests green; ASan/UBSan/TSan clean. Scripting reference in
docs/scripting.md.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-15 12:12:00 -04:00

63 lines
1.9 KiB
C

#include "naut/mpmc.h"
#include "test.h"
#include <pthread.h>
/* 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();
}