Design

Runtime and block devices

Everything in the storage stack runs as tasks inside the CLIO runtime: one process per node that hosts modules (ChiMods) and executes their operations as coroutines on a small pool of worker threads. This page covers the parts of the runtime that most affect the storage stack's performance and reliability, and the block-device module that every storage tier sits on.

Tasks and coroutines

Every operation (a put, a directory insert, a block write) is a task. A task names a pool, a method and a routing query, and carries input fields, output fields and optional bulk data buffers.

  • Tasks are serialized, not shared. A task lives in the submitting process's private memory. To cross into the runtime, its input fields are serialized into a buffer that is copied through a shared-memory ring or sent over a socket. The runtime rebuilds the task on its side, and serializes only the output fields on the way back. Bulk buffers that already live in shared memory travel as pointers, with no copy. Bulk data in private memory is copied once.
  • Handlers are coroutines. A handler can await other tasks (for example, the CTE awaiting block-device writes) without blocking its worker thread. This is what makes nested operations possible, such as a pool's creation awaiting the creation of the device pools beneath it. Two cooperative primitives exist: await a subtask, or yield for a while.
  • Lifetime is reference-counted. A future holds a reference to its task, so tasks are freed when the last holder lets go.

Getting a request into the runtime

A client process reaches the runtime over one of three transports, chosen automatically unless CLIO_IPC_MODE overrides the choice:

ModeHow it worksRelative cost
Shared memoryThe client attaches the runtime's segments, serializes each request into one of several inbound rings (sharded by thread, which keeps each thread's requests in order), and receives replies on its own response ring. Bulk data in the client's shared-memory segment is passed by pointer.Fastest. Chosen first when the runtime's segment exists.
Unix socketSerialized requests over a local socket; bulk bytes are sent inline.About 5× slower than shared memory. Used on the same machine when shared memory is unavailable.
TCPZeroMQ sockets, with the runtime dialing back to the client for replies; bulk bytes are sent inline.About 190× slower than shared memory. Works across machines and containers.

Workers drain their own inbound shard as part of their main loop, so there is no dedicated receive thread on the hot path. Code already running inside the runtime skips serialization entirely and hands the task over directly.

Dead clients and dead runtimes

  • The runtime survives dead clients. Shared-memory rings check whether the other side's process is still alive and skip transfers that stall. TCP replies to an unreachable client are retried briefly, then dropped. The number of reply connections the runtime keeps is bounded.
  • Clients survive a dead runtime. A heartbeat detects the loss, waits for the runtime to come back (or fails over to another node, if configured), reconnects and resends.
  • Dead clients' memory is reclaimed late. The shared-memory segments of clients that died are freed only when the runtime next starts.

Workers and scheduling

The runtime starts a fixed number of worker threads (four by default) and assigns them roles:

  • a scheduler and housekeeping worker;
  • a network-send worker and a network-receive worker, so outbound back-pressure cannot starve inbound replies;
  • a GPU-queue worker;
  • I/O workers in cost classes (quick, medium, heavy).

Tasks are mapped by predicted cost:

  • Large I/O is spread round-robin across I/O workers.
  • Cheap tasks run inline on the worker that spawned them, if it is not busy.
  • Everything else goes to the least-loaded worker in its cost class.

Predictions come from per-module models that learn from observed run times.

A coroutine that yields is parked in a blocked queue, bucketed by how long it wants to wait, and resumed on the same worker. Completions of subtasks wake their parent directly through an event queue. A monitor detects a worker stuck on one task for over a second and moves that worker's queue, together with its parked coroutines, to a freshly spawned worker. It also grows a cost class whose workers are all saturated, and retires idle extra workers.

Tradeoff: a small default worker count. With four threads, one worker handles all storage computation and I/O, and the GPU queue shares the scheduler worker. Elastic workers soften stalls, but throughput-oriented deployments should raise num_threads. The LMCache integration, for example, measured +35% going from four to six.

Cross-node traffic

Remote requests and replies are sent by ordinary periodic tasks of the admin module, running on the network-send worker. Each tick does the following:

  1. Drains a latency lane of small messages completely.
  2. Alternates requests and replies from a bulk lane, up to an 8 MiB budget.

Probes and acknowledgements therefore never queue behind megabyte puts. The price is that bulk bandwidth per node is bounded by the tick rate. Replies travel on per-request connections, so a small acknowledgement does not wait behind a large frame on a shared one.

Failures

  • Restarted peers. Every message carries the sender's incarnation number. When a peer restarts, requests in flight to its old incarnation fail promptly and are retried.
  • Silent peers. Peers with requests outstanding for too long, or that have gone quiet, are probed. Silence for 30 seconds marks them dead. A node whose own probe scan ran late re-arms instead of blaming others.
  • Gossip failure detection is off by default. A SWIM-style detector exists but is disabled: at 256 nodes, wide collective operations starved its probes and healthy nodes were declared dead. A dead peer therefore produces failed requests, after the timeout, or immediately in fail-fast mode. It does not produce automatic container migration. Data availability across node loss is handled by replication.

Pools, containers and routing

A pool is a named instance of a module spread across the cluster. A container is that pool's instance on one node.

  • Container ID equals node ID. The number of containers is fixed by the hostfile, not by how many nodes are currently alive. That keeps hashing stable even while nodes come and go.
  • Pool queries choose where a task runs:
    • Local: this node.
    • Direct: one container, by ID or by hash.
    • Range or Broadcast: many containers.
    • Dynamic: the module decides per task. The CTE, for example, routes each blob to its owner.
  • Interposition works at this level. Interposer modules override the routing decision too; the cache, for example, routes to the submitter's own node. That is how routing and layering compose.

Pools defined in the configuration's compose list are created at startup, in order. Each pool must come after the pool it forwards to.

The block-device module

Every storage tier is a block-device (bdev) pool. Its interface is deliberately minimal:

  • allocate blocks of a given total size;
  • write or read a list of blocks;
  • free blocks;
  • report statistics;
  • sync.

The bdev keeps no per-object metadata: recording which blocks belong to which blob is the CTE's job. For a local device, the CTE calls the allocator directly instead of sending a task, because the round trip was the dominant cost of small writes.

BackendStoragePersistence
ramOne sparse shared-memory mapping, readable directly by clients. This is what makes zero-IPC reads possible.Volatile
pinned, hbmPinned host memory, or GPU device memory (a failed allocation is fatal; it never silently falls back to host memory)Volatile
fileA backing file, using io_uring, libaio, POSIX AIO or IOCP depending on platformPersistent, with an allocation log
s3, gcsObject storagePersistent data; allocator state not persisted
noopDiscards dataFor benchmarking

The allocator

The allocator is a bump-pointer heap plus segregated free lists (512 B to 1 MB size classes), one set per worker. A worker can steal from other workers' lists. Once the heap is exhausted, a large request is assembled from several free blocks, so one allocation can return several extents.

The allocator never compacts or coalesces, and the heap pointer only moves forward. Long-running churn therefore fragments free space into blocks of 1 MB or less, and the CTE handles the resulting multi-block layouts. Memory tiers align to 512 bytes; 4 KiB alignment had inflated small blobs 3.6×. File tiers align to 4 KiB, so direct I/O can be used when buffers allow.

File tiers

  • Lazy growth. A new backing file starts at one growth unit (1 GiB by default) and is extended on demand as allocations reach its end.
  • Space reservation. On filesystems with native fallocate, each growth step reserves real blocks, so a full disk fails at allocation time rather than as an I/O error on a later write. On filesystems without it (some network, FUSE and virtualized mounts), the file is kept sparse. The C library's emulation of the reservation writes every byte, which once stalled pool creation for 155 seconds on a Docker Desktop bind mount, so it is not used. The tradeoff: on those filesystems, a full disk can surface as a write-time error.
  • Honest capacity. Reported free space is capped by what the host filesystem actually has free, so the CTE does not place data on a device configured larger than its disk. Several tiers sharing one disk each see the whole free space, so the figure is an upper bound.
  • I/O is per block. Within one task, blocks are written one at a time. Concurrency comes from many tasks running at once.

The allocation log

A file tier's allocator survives restarts through an append-only allocation log.

  • Allocations are logged before blocks are handed out.
  • Frees are logged before space is reused.
  • Flushing and sync. Records are flushed to the OS on every call and fsynced every 50 ms by a background thread.
  • Compaction happens when dead records outnumber live ones two to one, by writing a temporary file and renaming it into place.
  • On restart, live extents are replayed. Every gap below the highest live byte becomes free space, and the heap resumes after the last live byte.

A crash can leak space, but never hand out live bytes twice. Before this log was restored, a restarted heap began again at zero and overwrote live data.

Sync fdatasyncs the data file first, then the log, so a crash can leave synced bytes the log does not yet mention (the CTE's own log covers them), but never a logged block whose bytes were lost.

Durable pool registry

Each node keeps a small pool log of the durable pools it hosts and their parameters, fsynced on every change. clio_run start replays it in creation order, so a node comes back with the same storage tiers and module stack it had. --fresh discards it along with all other persisted state.

Tradeoffs at a glance

DecisionBuysCosts
Serialized tasks over shared-memory ringsProcess isolation with near-memcpy latency. Bulk data in shared memory moves by pointer.Fixed-size rings. A client with very many outstanding requests can wait on ring space.
Coroutines on few workersThousands of in-flight operations without thread overhead.A handler that blocks a thread stalls its worker until the monitor rescues it.
Geometry fixed by the hostfileStable hashing, no split-brain resharding.Adding nodes means a new hostfile and a restart, and dead nodes' slots remain.
Conservative failure detectionNo false "dead" verdicts under load.Slow detection (around 30 s) and no automatic container recovery.
Bump heap and free lists, no compactionConstant-time allocation with a simple, crash-safe log.Fragmentation under churn; blobs become multi-extent.
Sparse files where fallocate is missingNo multi-minute zero-fills on network or virtual filesystems.A full disk may surface as a write error instead of an allocation error.