Skip to content

Architecture: Shared-Nothing, Thread-per-Core

Pion is a shared-nothing, thread-per-core database engine written in Mojo. Each worker thread owns all of its resources — no locks, no shared state on the hot path.


Thread Model

main()
  ├── create_listen_socket(port)        # One shared listen fd (no SO_REUSEPORT)
  └── pion_spawn_workers(N) → pthreads  # N workers (16 MB stacks), each pinned to a CPU core (capped at 4 on macOS Apple Silicon)
        └── worker_task(i)
              ├── set_thread_affinity(i)
              ├── SlabAllocator[ListNode](10M)
              ├── Pion.__init__()
              │     ├── Step 1: NetworkEngine.__init__() + server.listen(shared_fd)
              │     │            ← Socket is LIVE before heavy init
              │     └── Step 2: SlabHashMap(10M) + HNSWGraph + WAL + RaftNode
              └── Pion.run_server()
                    └── NetworkEngine.run_server_kqueue()   # macOS: kqueue event loop
                        NetworkEngine.run_server_epoll()    # Linux --epoll: best for P=1 w=1
                        NetworkEngine.run_server_uring()    # Linux default: io_uring event loop
                        NetworkEngine.run_server_xdp()      # Linux --xdp: AF_XDP kernel bypass

Each worker owns privately: - SlabHashMap (10M slots) — the keyspace - HNSWGraph — the vector index (or borrows from SharedHNSWView after FT.OPTIMIZE) - SlabAllocator[ListNode] — list node slab - ObjectPool[SlabHashMap/SlabSkipList/SlabList] — recycled data structure instances - WAL, RaftNode — persistence and replication state - KVCacheStore, LayerStore, AttentionIndex — externalized attention state (when --kvcache) - SemanticRouter — semantic router state (when --kvcache) - SpeculativeRAG — per-session trajectory tracking + speculative cache (when --kvcache) - Network socket + kqueue fd + recv buffer + response buffer

Shared across all workers: - One listen socket fd (created before worker spawn; workers call accept() independently — one wins per connection) - SharedHNSWView — read-only pointers to the index built by the first worker to run FT.OPTIMIZE; other workers lazy-borrow these pointers on the first FT.SEARCH


Event Loop (macOS: kqueue, Linux: epoll / io_uring / XDP)

while True:
    nevents = kevent_batch(kq, pending_changes, pending_count, events, 1024)
    for event in events[0..nevents]:
        if event.filter == EVFILT_WRITE:        # Back-pressure drain
            writer.flush_response(fd, server, kq)
        elif fd == server.fd:                   # New connection
            new_fd = server.accept()
            server.set_nonblocking(new_fd)
            server.set_tcp_nodelay(new_fd)
            server.kevent_add_read(kq, new_fd)  # Immediate (not batched)
        else:                                   # Client data
            while True:
                n = server.recv(fd, buffer + stored, avail)
                consumed = fast_path.process_data_plane(...)
                if consumed == 0:
                    consumed = slow_path.process_slow_path(...)
                # memmove leftover; break when recv returns < buffer_size

EVFILT_READ is registered immediately on accept (not batched) to avoid missing data that arrives before the next kevent_batch() call.

Linux backends

Four backends are available on Linux, selected by CLI flag:

Backend Flag Syscalls per batch (K fds) Best for
epoll --epoll 2K+1 P=1 w=1 (91-96K RPS, Redis parity)
io_uring --iouring (default) 2 Multi-connection production
io_uring SQPOLL --sqpoll 0-1 Experimental (hangs under load)
XDP/AF_XDP --xdp 0 (kernel bypass) +13% vs io_uring at P=1 (154K over real NIC); multi-worker + P>1 ready

One event loop per worker

Each worker runs NetworkEngine.run_server_*() directly on its own pthread; there is no inner task scheduler. Mojo green threads do not yield at OS syscalls, so a second task per worker would block the event loop.

Throughput anchors

  • 14.0M ops/sec peak (io_uring w=32 P=50 d=256) -- 10.2x Redis
  • 6.23M ops/sec (epoll w=16 P=50) -- 4.7x Redis
  • P=1 parity with Redis via --epoll: 91-96K vs 95-96K, wins LRANGE (+1-4%)

Per-Request Data Flow

TCP recv
    │
    ├──► FastPathHandler.process_data_plane()     # zero-alloc, 28 hot commands
    │         returns consumed > 0 on success, 0 to fall through
    │
    └──► SlowPathHandler.process_slow_path()      # RESP3 fallback, FT.* commands
              │
              ▼
    ResponseWriter.flush_response(fd, server, kq)
              │
              ├── send() succeeds → done
              └── EAGAIN → copy to pending_buffers[fd], register EVFILT_WRITE

Shared-Nothing Data Model

Each worker's keyspace is completely independent. SO_REUSEPORT is not used — instead a single shared listen fd distributes connections via kernel accept() scheduling. A client connected to worker 0's fd stays on worker 0 for its entire lifetime.

This is observable as a correctness property, so it is fenced. -w N is N independent keyspaces, not one keyspace served by N threads: a SET acknowledged +OK on a connection that landed on worker 1 is invisible to a GET on a connection that landed on worker 3. Every default connection pool (redis-py, Jedis, go-redis, ioredis) opens more than one connection, so the violation is silent — nils, no error. Measured at -w 4 with 16 concurrently-opened connections: 41 of 90 GETs of a just-acked key returned nil.

Two consequences for anyone measuring this:

  • The default is -w 1, and -w N > 1 refuses to start without --independent-workers, which prints the semantics at startup. The flag is an acknowledgement, not a feature switch — it changes no behaviour beyond letting the server boot.
  • Connections opened SERIALLY all land on one worker and hide the split completely (0/12 nils in the same session that measured 41/90 concurrently). Any probe of cross-worker behaviour must open its connections concurrently, or it proves nothing.

The same split runs through the rest of the surface: cross-worker PUBLISH delivers to zero subscribers, FLUSHALL clears one worker's slice, and replication covers worker 0 only.

Vector data is an exception: after FT.OPTIMIZE, the index-building worker publishes read-only pointers via SharedHNSWView. Other workers borrow these pointers lazily on the first FT.SEARCH. All mutable per-search state (visited_map, query_int8, cur_num) remains per-worker.


Key Design Decisions

Decision Rationale
Single shared listen fd Avoids SO_REUSEPORT kernel hash imbalance; all workers compete fairly for connections
EVFILT_READ registered immediately on accept Prevents data loss if client sends before next kevent_batch()
4 MB shared response buffer per worker Pre-allocated; no alloc on hot path
Ziplist (≤1024 entries) + Quicklist (>1024) for lists Contiguous memory for small lists; two-sided 256-element segmented arrays for large lists
Generation-counter visited-set in HNSW O(1) mark/check without clearing 50K-entry array between queries
INT8-INT8 batch-8 search kernel Graph build and search use the same metric; no dequantization per neighbor
Binary protocol on port+1 0xCA5E framing avoids RESP parsing overhead for bulk tensor ops (3.67M tok/s ATTEND.STORE)

Externalized Attention Modules

When --kvcache is enabled, the following modules are activated:

Module File Role
KVCacheStore src/network/kv_cache_store.mojo Per-session, per-layer KV cache tensor storage
LayerStore src/network/layer_store.mojo Binary blob storage for per-layer tensors (LAYER.STORE/FETCH)
AttentionIndex src/network/attention_index.mojo Per-layer HNSW index for O(N log k) attention approximation
BinaryProtocol src/network/binary_protocol.mojo 0xCA5E framed wire protocol on port 1975 (main_port + 1)
kv_cache commands src/commands/kv_cache.mojo KV.STORE, KV.FETCH, KV.INFO handlers
attend commands src/commands/attend.mojo ATTEND.CREATE, ATTEND.STORE, ATTEND.FINALIZE, ATTEND.QUERY, ATTEND.INFO handlers

Data flow (externalized attention):

Inference engine (vLLM/MLX) → PionAttentionClient (vllm-pion/)
    │
    ├── ATTEND.CREATE session_id key_dim value_dim
    ├── ATTEND.STORE session_id layer_id num_tokens keys_fp32 values_fp32  (3.67M tok/s binary)
    ├── ATTEND.FINALIZE session_id layer_id  (builds HNSW, 1.0s for 128K tokens)
    └── ATTEND.QUERY session_id layer_id k query_fp32  (86us per layer at 128K)
         └── Returns top-k values (cosine 1.0 vs full attention at k=64)

Performance (binary protocol, 128K tokens, 128d keys): - Store: 3,667,377 tok/s - Query: 86us per layer at 128K tokens - 40-layer total: 3.4ms - Build: 1.0s for 128K tokens (finalize) - Memory: ~2.5 GB for 40 layers at 128K (staging freed after finalize)


Semantic Router Modules

When --kvcache is enabled, the semantic routing modules are activated:

Module File Role
SemanticRouter src/network/semantic_router.mojo HNSW/FP32 routing table indexed by node centroid embeddings
route commands src/commands/route.mojo AI.ROUTE.REGISTER, AI.ROUTE.UPDATE, AI.ROUTE, AI.ROUTE.REMOVE, AI.ROUTE.INFO handlers

Routing strategy: FP32 brute-force cosine similarity for <=16 nodes (zero recall loss); HNSW O(log N) for >16 nodes. Per-node capacity limits, exclude filters. 88% routing accuracy, 0.14ms/route, 7K QPS.


Speculative RAG Modules

When --kvcache is enabled, the speculative RAG modules are activated:

Module File Role
SpeculativeRAG src/network/speculative_rag.mojo Per-session trajectory ring buffer (last 10 embeddings), momentum prediction, speculative HNSW cache
speculative commands src/commands/speculative.mojo RAG.SPECULATE.ENABLE, RAG.QUERY, RAG.SPECULATE.INFO handlers

Data flow (speculative RAG):

RAG.SPECULATE.ENABLE session_id
    └── Creates per-session trajectory ring buffer

RAG.QUERY session_id query_embedding
    ├── Append embedding to trajectory ring buffer
    ├── Check speculative cache (cosine > 0.9 match against predictions)
    │     └── HIT: return pre-computed HNSW results (0.2ms)
    ├── MISS: execute live HNSW search
    ├── Generate 3 predictions: current + alpha * (current - previous), alpha=[0.5, 1.0, 1.5]
    └── Pre-execute HNSW search for each prediction → store in speculative cache

Performance: 75% hit rate on linear query trajectories (6/8 queries), 0.2ms prediction+lookup latency.