Skip to content

Repository files navigation

async-worker-pool (AWP)

Sharded low-latency dispatch worker pool in C — preallocated pthread workers, bounded per-worker queues, stable hash sharding for per-(feed, symbol) FIFO, blocking backpressure (zero drops), fault-isolated process callbacks, supervisor heartbeats, and bounded shutdown.

Designed as the C equivalent of a permanent-worker market-data dispatch stage. Local microbench target: p99 submit→process-return ≤ 5 ms (closed-loop burst with light simulated work; not open-loop publisher-accept SLA).

Features

  • N fixed workers created once (never per message)
  • Producer-side shard: FNV-1a(feed || 0x1F || symbol) → worker index or fast 64-bit keyed submission (awp_submit_keyed)
  • Bounded atomic ringsSPSC / MPSC / SPMC / MPMC (ring_mode), sequence protocol, spin/yield backpressure, never drop
  • Zero-Copy Claim & Commit API (awp_claim_frame / awp_commit_frame) for in-place writing directly into ring slabs with zero memcpy
  • Page-Aligned Ring Slabs (4KB) — zero heap allocation / zero CAS lock contention on the hot path
  • Dedicated broadcast workers for mark-price / funding-style feeds
  • Soft fault isolation: process() errors recycle the frame and continue
  • Supervisor: restarts dead/stalled workers; per-worker metrics
  • Bounded shutdown wait then quarantine stuck callbacks (no cancel/detach)
  • Rust FFI Bindings: Safe RAII wrapper and Zero-Copy crate in bindings/rust/
  • Zig 0.16 Parallel Project: Native SIMD @Vector implementation with ArenaAllocator in async-worker-pool_zig

Table of Contents


Cross-Language Benchmark Comparison (1,000,000 Messages)

Workload: 1,000,000 messages processed asynchronously (Zig: 4 Pinned Workers on Apple Silicon P-Cores; C11: 32 Workers).

Implementation Mode / API Throughput Median (p50) p99 Latency Mean Latency
Zig 0.16 (async-worker-pool_zig) Multi-Threaded Async (4 Pinned Workers) 5.38 M msg/s 🚀 < 100 ns 1.00 µs (1,000 ns) 547.0 ns (0.55 µs)
Zig 0.16 (Pure SPSC Ring) Concurrent SPSC (0 CAS) 171.76 M ops/s 🚀 < 6 ns < 8 ns 5.82 ns
C11 (async-worker-pool) Zero-Copy Claim/Commit (32 Workers) 0.52 M msg/s 3.46 µs (3,458 ns) 1.11 ms (1,110,000 ns) 2.11 µs (2,109 ns)
C11 (Raw SPSC Ring) Lock-Free Push/Pop 62.50 M ops/s N/A (Bulk) N/A (Bulk) 16.00 ns (avg)
Rust (awp-rs) Safe FFI Zero-Copy (v0.3.0) 0.53 M msg/s 3.35 µs (3,350 ns) 1.15 ms (1,150,000 ns) 1.87 µs (1,870 ns)

Detailed Tail Latencies Breakdown (1,000,000 Messages)

Percentile Zig 0.16 Engine (Phase 1 Final) C11 Engine (async-worker-pool) Delta / Notes
Min (Observed Floor) 15 ns (0.015 µs) 83 ns (0.083 µs) Observed Single-Hop Floor
p50 (Median) < 100 ns 3.46 µs (3,458 ns) Zig is > 34x lower latency 🚀
p90 1.00 µs (1,000 ns) 11.17 µs (11,167 ns) Zig is 11.2x lower latency 🚀
p99 (Tail) 1.00 µs (1,000 ns) 1.11 ms (1,110,000 ns) Zig is 1,110x lower tail jitter 🚀
p99.9 96.0 µs (96,000 ns) 1.27 ms (1,270,000 ns) Zig is 13.2x lower tail jitter 🚀
Max 128.0 µs (128,000 ns) 1.67 ms (1,670,000 ns) Zig is 13.0x lower peak jitter 🚀
Pure SPSC Throughput 171.76 Million ops/sec 62.50 Million ops/sec Zig is 2.75x faster (5.82 ns/op) 🚀

Throughput Comparison SPSC Comparison

Tail Latencies Distribution

Full benchmark reports and architecture evolution reasoning:

Lifetime contract (read this before production use)

Rule Meaning
Destroy exactly once One owner thread; after every other handle user has finished
External quiescence Join producers and concurrent shutdown waiters before destroy
Shutdown deadline Absolute wait budget only — does not kill process()
Quarantine Sticky; pool storage may leak; treat as process recycle
cfg.user lifetime Must remain valid until process exit if any callback may still run

Canonical owner sequence:

  1. Stop publishing the handle to new work; set the producer-stop condition.
  2. Call awp_pool_shutdown() from the designated owner; record its return value and metrics.
  3. Join every producer, metrics reader, and concurrent shutdown caller.
  4. Call awp_pool_destroy() exactly once.
  5. If shutdown returned > 0, preserve callback-owned state and normally terminate/recycle the process. Do not create replacement pools indefinitely in the same process.

Quick start

make          # libawp.a + tests + bench + examples
make check    # functional tests only (no latency gates)
make check-all # functional + benches + examples
#include "awp/awp.h"

static int on_frame(const awp_frame_t *f, void *user) {
    /* e.g. local publisher enqueue — do not retain f after return */
    (void)user; (void)f;
    return 0; /* non-zero = soft error; worker continues */
}

int main(void) {
    awp_config_t cfg;
    awp_pool_t *pool = NULL;
    int shut;

    awp_config_init(&cfg);
    cfg.n_workers = 32;           /* skew headroom, not core count */
    cfg.queue_capacity = 256;
    cfg.frame_pool_size = 4096;
    cfg.process = on_frame;

    awp_pool_create(&cfg, &pool);
    awp_submit(pool, "trades", "BTCUSDT", payload, len, 0);

    /* Real services: stop producers, shutdown (wakes blocked submits), then join. */
    shut = awp_pool_shutdown(pool);
    /* join producers / metrics threads here if any */
    awp_pool_destroy(pool);
    if (shut > 0)
        return 2; /* quarantined: recycle process in production */
    return 0;
}

See examples/simple_publish.c for multi-reader usage.


Rust FFI Bindings (awp-rs)

A memory-safe, idiomatic Rust crate with Zero-Copy Claim & Commit API and RAII lifecycle is provided in bindings/rust/:

cd bindings/rust
cargo test                                       # Run Rust tests
cargo run --release --example bench_throughput   # Benchmark 1,000,000 messages

Rust Example (Zero-Copy In-Place Dispatch)

use awp_rs::{AsyncWorkerPool, AwpRingMode};

fn main() -> Result<(), i32> {
    // Initialize pool with 16 workers and a thread-safe callback
    let pool = AsyncWorkerPool::new(16, 2048, AwpRingMode::Mpsc, |frame| {
        let data = frame.payload();
        // Process message in-place without heap allocations
        0
    })?;

    // Zero-Copy Claim: allocate slot directly in ring slab
    let mut guard = loop {
        match pool.claim(0) {
            Ok(g) => break g,
            Err(_) => std::thread::yield_now(),
        }
    };

    // In-place serialization with 0 memcpy overhead
    let buf = guard.payload_mut();
    buf[..32].fill(0x7F);
    guard.set_payload_len(32);
    guard.commit()?;

    Ok(())
}

Layout

include/awp/awp.h     Public API
src/                  ring, frame pool, shard, worker, supervisor, pool
tests/                unit + supervisor + e2e + lifecycle + contract drills
bench/                closed-loop microbench + open-loop mock harness
examples/             mock publish demo + per-mode demos
docs/DESIGN.md        Architecture, lifecycle contract, test matrix
docs/DIAGRAMS.md      Architecture / lifecycle / ring / supervisor diagrams
docs/BENCHMARKS.md    Local latency & throughput results
docs/diagrams/        Rendered PNG diagrams

Design notes (short)

Topic Choice
Queue Atomic sequence ring — SPSC/MPSC/SPMC/MPMC via ring_mode
N workers Config knob for hash skew headroom (e.g. 32), not #cores
Ordering Stable hash ⇒ one worker per key ⇒ FIFO by construction
Backpressure Block producer when full; drops must stay 0
Shutdown Quiesce → close rings/pool → join under wait budget → quarantine stuck callbacks

Documentation

Document Contents
docs/DESIGN.md Architecture, lifecycle contract, mitigations
docs/DIAGRAMS.md Architecture / submit / lifecycle diagrams
docs/BENCHMARKS.md Local latency & throughput captures
docs/PERFORMANCE_COMPARISON.md AWP vs market queues / pools / HFT stacks
docs/KNOWN_ISSUES.md Residual S3 nits · GitHub issues

Build, test, install

make lib
make check                 # functional correctness
make check-sanitize        # ASan+UBSan (Clang/GCC)
make check-bench           # optional microbench (not CI gate)
make install PREFIX=/usr/local
pkg-config --cflags --libs awp
Artifact Covers
test_unit / test_unit_modes FIFO, backpressure, faults × ring modes
test_ring_modes Raw ring stress with exact ID accounting
test_e2e* / test_supervisor Multi-reader, restart, sticky quarantine
test_e2e_lifecycle Drain + concurrent shutdown
test_teardown_contract Clean vs quarantined teardown drills
test_restart_create_fail Deterministic restart pthread_create failure
bench_dispatch / bench_all_modes / bench_ring Closed-loop microbench
bench_openloop Open-loop schedule + mock accept (not a real-publisher SLA)

cfg.ring_mode = AWP_RING_SPSC | MPSC | SPMC | MPMC — match actual producer/consumer counts.

License

MIT — see LICENSE.

About

High-performance lock-free async worker pool in C11 & Rust FFI — zero-copy claim/commit, cacheline-isolated rings, sub-microsecond latency

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages