From 996546c0cdf4811191c800063c541e55ab2e77f6 Mon Sep 17 00:00:00 2001 From: Justin Mclean Date: Wed, 30 Sep 2026 16:31:35 +1000 Subject: [PATCH] fix(docs): correct the architecture page against server 0.9.0 TLS and WSS connections are handed off before the handshake, not terminated on shard 0. Also corrects the storage tree, checkpoint, superblock and segment size descriptions. --- content/docs/introduction/architecture.mdx | 21 ++++++++++----------- src/components/architecture-diagrams.tsx | 6 +++--- 2 files changed, 13 insertions(+), 14 deletions(-) diff --git a/content/docs/introduction/architecture.mdx b/content/docs/introduction/architecture.mdx index 5968a61898..1c444d3aec 100644 --- a/content/docs/introduction/architecture.mdx +++ b/content/docs/introduction/architecture.mdx @@ -13,7 +13,7 @@ Before diving into the architecture details, here's an overview of a message's j ## Thread per core (shared nothing) + io_uring -Iggy uses a **thread-per-core shared nothing architecture** combined with `io_uring` for maximum performance. This design has been proven by systems like ScyllaDB and Redpanda, and is inspired by the Seastar framework. +Iggy uses a **thread-per-core shared nothing architecture** combined with `io_uring` for maximum performance. This design has been proven by systems like ScyllaDB and Redpanda, both built on the Seastar framework. @@ -25,8 +25,8 @@ Each configured **shard** (an instance of `IggyShard`) has its own single-thread Shard 0 has a special role: it binds **every listener** - the replica plane and all client transports (TCP, QUIC, WebSocket, HTTP). Connections are then spread across shards at accept time: -- **Plaintext TCP and WebSocket** connections are handed off round-robin to peer shards. Shard 0's coordinator duplicates the socket's file descriptor, ships a connection-setup frame to the target shard, and drops its own handle, so the owning shard serves the connection from then on - before a single byte is read. -- **QUIC, TLS-wrapped TCP and secure WebSocket (WSS)** connections terminate on shard 0, because their per-connection state cannot be moved between shards. HTTP is also served on shard 0. +- **TCP and WebSocket** connections, plaintext or TLS, are handed off round-robin across the shards. Shard 0's coordinator duplicates the socket's file descriptor, ships a connection-setup frame to the target shard, and drops its own handle, so the owning shard serves the connection from then on - before a single byte is read. For TLS-wrapped TCP and secure WebSocket (WSS), the target shard runs the TLS handshake. Shard 0 is part of the rotation by default; `skip_shard_zero_for_clients = true` under `[cluster.coordinator]` takes it out. +- **QUIC** connections stay on shard 0, installed through its UDP endpoint. HTTP is also served on shard 0. All shards, **including shard 0**, own partitions and serve partition requests. @@ -35,7 +35,7 @@ All shards, **including shard 0**, own partitions and serve partition requests. Requests are routed between shards using **message passing** (via `crossfire` bounded mpsc channels), so partition state remains on its owning shard. The routing logic splits operations into two planes: - **Metadata operations** (create/delete stream/topic/user etc.) always execute on **shard 0** - it is the only shard that commits metadata -- **Partition operations** (send_messages, poll_messages, store_consumer_offset) are routed to the shard owning that partition via a lock-free concurrent map lookup. A request that lands on a non-owning shard rides the inter-shard message bus to the owner +- **Partition operations** (send_messages, poll_messages, store_consumer_offset) are routed to the shard owning that partition via a lock-free concurrent map lookup. The request is sent into the owning shard's channel, whichever shard received it Partition-to-shard assignment is **deterministic**: the packed `IggyNamespace` is hashed with `Murmur3`, and the upper 16 bits of the hash are taken modulo the shard count (the upper bits are used because Murmur3 has weak lower bits for small integer inputs). @@ -53,18 +53,18 @@ The sharding system supports multiple allocation modes via the `cpu_allocation` - A numeric value (e.g. `4`) - exactly N shards, pinned to the first N CPUs in the process's allowed CPU set when pinning is enabled - A range (e.g. `"5..8"`) - shards on CPUs 5, 6 and 7 when pinning is enabled; those CPUs must be allowed for the process - `"numa:auto"` - automatically detect NUMA topology and select physical cores, avoiding sibling hyperthreads -- `"numa:nodes=0,1;cores=4;no_ht=true"` - fine-grained NUMA control per node with hyperthread avoidance +- `"numa:nodes=0,1;cores=4;no_ht=true"` - explicit NUMA nodes, with the same number of cores taken from each listed node and hyperthread avoidance ### io_uring and compio -Traditional async runtimes like tokio use `epoll` which is **readiness-based** - you ask the kernel "is this file descriptor ready?" and then perform the I/O yourself. [Regular files cannot be registered with epoll](https://man7.org/linux/man-pages/man2/epoll_ctl.2.html). Tokio runs file I/O on a blocking thread pool (512 threads by default, configurable). This does not scale well. +Traditional async runtimes like tokio use `epoll` which is **readiness-based** - you ask the kernel "is this file descriptor ready?" and then perform the I/O yourself. [Regular files cannot be registered with epoll](https://man7.org/linux/man-pages/man2/epoll_ctl.2.html). Tokio runs file I/O on a blocking thread pool that grows on demand (up to 512 threads by default, configurable). This does not scale well. `io_uring` is **completion-based** - you submit I/O requests to a submission queue (SQ), and the kernel completes them and places results in a completion queue (CQ). Both queues are shared ring buffers between user space and kernel. Submissions and completions can be batched to reduce syscalls. This is fundamentally better for disk I/O. Iggy uses **compio** as its async runtime, which provides a driver-disaggregated architecture on top of io_uring (Linux) and IOCP (Windows). Each shard gets its own compio executor configured with: -- Capacity: 4096 concurrent I/O operations (by default) +- Capacity: an io_uring submission queue of 4096 entries (by default) - Event interval: poll the I/O driver after 128 scheduler ticks (roughly task polls) by default - Cooperative task running enabled @@ -99,7 +99,7 @@ local_data/ ├── logs/ │ └── iggy-server.log ├── state/ -│ └── log/ +│ └── log └── streams/ └── 0/ └── topics/ @@ -109,13 +109,12 @@ local_data/ ├── 00000000000000000000.index ├── 00000000000000000000.log ├── superblock.a - ├── superblock.b └── offsets/ ├── consumers/ └── groups/ ``` -This example shows a fresh partition with the default paths. The stream, topic and partition directories are named after their numeric IDs, **assigned from 0**. The `metadata/journal.wal` file is the VSR write-ahead log that persists all metadata operations (stream/topic/user creation, etc.). The `runtime/current_config.toml` file captures the configuration the server actually booted with. Segment files are named by their 20-digit start offset. Recovery can create an empty active segment at a reserved offset beyond the last stored message; partition superblocks record the recovery frontiers. The `.index` file is created automatically and speeds up searches by keeping track of the offsets and timestamps of the records. The `offsets/` directory holds the server-side consumer and consumer group offsets. +This example shows a fresh partition with the default paths. The stream, topic and partition directories are named after their numeric IDs, **assigned from 0**. The `metadata/journal.wal` file is the VSR write-ahead log for metadata operations (stream/topic/user creation, etc.). Its index has 1024 slots by default (`journal_slots` under `[metadata]`). Before they run out, the server checkpoints the metadata into `metadata/snapshot.bin`, records the checkpoint in a superblock under `metadata/`, and drops the checkpointed operations from the log. On a single node, neither file exists until the first checkpoint. The `runtime/current_config.toml` file is rewritten at every boot with the configuration the server booted with and the addresses it actually bound. Secrets (the JWT secrets, the encryption key and the cluster shared secret) are left out. The `state/log` file is empty; the server creates it at boot and does not otherwise use it. Segment files are named by their 20-digit start offset. Recovery can create an empty active segment at a reserved offset beyond the last stored message; partition superblocks record the recovery frontiers. A superblock is written to two files in turn, `superblock.a` and `superblock.b`. A fresh partition has only `superblock.a`; `superblock.b` appears at the next superblock write, for example after a restart or a topic purge. The `.index` file is created automatically and speeds up searches by keeping track of the offsets and timestamps of the records. The `offsets/` directory holds the server-side consumer and consumer group offsets. ## Memory pool @@ -129,6 +128,6 @@ Messages flow through a multi-stage write pipeline: 2. Partition VSR replicates prepares. With `durability=persisted`, a multi-replica group requires recoverable prepare-WAL copies at the replication quorum before commit. 3. Awaited writes apply committed operations before success is returned. A singleton with `durability=persisted` synchronizes local segment state before replying. 4. Ordinary segment writes use per-topic count and byte thresholds (defaults: 1024 messages, 1 MiB); required persistence, capacity pressure, and lifecycle work can flush earlier. The `MessagesWriter` uses **vectored I/O** in chunks of up to 1024 buffers. Partial writes can require more than one I/O submission. -5. When a segment reaches the topic's segment size (default 1 GiB), it is **sealed** and a new segment is created. +5. When a segment reaches or passes the topic's segment size (default 1 GiB), it is **sealed** and a new segment is created. The batch that crosses the size is written whole, so a sealed segment can exceed it by up to one batch. Message `durability` and `consumer_offset_durability` default independently to `replicated`. Both policies write data to disk; `persisted` adds a stable-storage requirement at completion. See [Durability](/docs/server/durability). diff --git a/src/components/architecture-diagrams.tsx b/src/components/architecture-diagrams.tsx index 0e0bb7bac1..f465a94102 100644 --- a/src/components/architecture-diagrams.tsx +++ b/src/components/architecture-diagrams.tsx @@ -126,21 +126,21 @@ export function ShardDiagram() { id: 0, label: "Shard 0 (Coordinator)", color: "var(--color-fd-primary)", - features: ["Binds all listeners: TCP, QUIC, HTTP, WS", "Replica plane listener", "Metadata plane (left-right write handle)", "QUIC + TCP-TLS + WSS + HTTP terminate here", "Hands plaintext TCP/WS to peers (fd transfer)"], + features: ["Binds all listeners: TCP, QUIC, HTTP, WS", "Replica plane listener", "Metadata plane (left-right write handle)", "QUIC + HTTP terminate here", "Hands TCP/WS, plaintext or TLS, to shards (fd transfer)"], partitions: ["P0", "P3", "P6"], }, { id: 1, label: "Shard 1", color: "#3b82f6", - features: ["Metadata read handle (left-right)", "Plaintext TCP/WS via fd transfer", "Owns partitions"], + features: ["Metadata read handle (left-right)", "TCP/WS, plaintext or TLS, via fd transfer", "Owns partitions"], partitions: ["P1", "P4", "P7"], }, { id: 2, label: "Shard 2", color: "#8b5cf6", - features: ["Metadata read handle (left-right)", "Plaintext TCP/WS via fd transfer", "Owns partitions"], + features: ["Metadata read handle (left-right)", "TCP/WS, plaintext or TLS, via fd transfer", "Owns partitions"], partitions: ["P2", "P5", "P8"], }, ];