Design
The interposition chain
Replication, caching, indexing and compression are not built into the core. Each is a separate ChiMod that speaks the core's own interface and forwards to the next pool down. A deployment composes the stack it needs, and a client pointed at the top of the stack uses the plain core API unchanged.
The interposer pattern
An interposer is a ChiMod that accepts the core's method numbers and
task formats. It overrides the handful of operations it cares about,
usually put, get, size and batched put, and forwards everything else
untouched to the pool named by its next_pool_id. Its own
extra methods are numbered from 100 upward, so they cannot collide with a
core method. A pool must be composed after the pool it forwards to.
The core supplies the mechanisms the layers need: replica slots, shadow copies, cache-copy registration and invalidation, transform flags and fault handlers. Each layer supplies the policy: what to copy, where, and when.
Why this order
- Cache on top. The cache keeps an untransformed copy on the reader's node. Blobs compressed further down can still be served by zero-copy shared-memory reads, which refuse transformed data.
- Indexer above the compressor. The indexer re-reads blobs to tokenize them, and needs the logical, uncompressed bytes.
- Replication directly above the core. Replicas are copies of exactly what is stored, compressed or not, so every copy has the same format.
Replication: the reliability layer
Replication manages two independent kinds of copy.
| Remote copies | Local persistent replicas | |
|---|---|---|
| Protects against | Node loss | Device loss, and data loss on a volatile primary |
| Where | The next k nodes in ring order after the blob's owner | Durable tiers on the owner's own node |
| When written | Synchronously, before the put is acknowledged, one copy after another | Asynchronously by default: a sweep every 50 ms copies the current primary, so repeated overwrites are coalesced. Synchronous if the period is set to 0. |
| On failure | The dead or failing copy is skipped with a warning, and the put still succeeds | The blob stays dirty and the next sweep retries it |
| Configured by | replication_factor (default 1, so no remote copies) or remote_copies | num_replicas (default 0) |
Write path
At the owner, a put does the following in order:
- Waits for any pending hand-back after a restart.
- Locks the blob, so there is one writer.
- Refills a primary that lost bytes in a restart, so a write past its end cannot leave a zero-filled gap.
- Writes the primary first, so the primary is never staler than a replica.
- Updates local replicas, or marks them dirty for the sweep.
- Mirrors to remote copies.
fsync waits for any running sweep, flushes the tag's dirty
replicas, and then forwards the sync to the core. After an fsync, the
replicas are current.
Read path
A get tries copies in this order:
- The primary, if it covers the requested range.
- A local persistent replica. The whole replica is then re-cached into the primary.
- A remote copy, for example when the primary's device is gone or its RAM contents were lost in a restart. The primary is healed from it.
- Otherwise the core's normal "absent" semantics.
After a restart, the primary is first reconciled with its remote copy. The reported size accounts for bytes lost in the restart.
Failover and hand-back
With failover enabled in the core, a dead owner's blobs are served by its first live successor from that successor's shadow copy. The successor logs every change it makes and invalidates cached copies cluster-wide. Changes return to the owner two ways:
- Pull: a restarted owner pulls changes from every live node before it serves requests.
- Push: a sweep every 500 ms pushes changes to owners that have come back.
A hand-back log can make this survive a restart of the stand-in too. Hand-back runs even at replication factor 1, because the stand-in may hold the only copy of writes made during the outage. Hand-back pushes are not mirrored again, so a hand-back can never overwrite a newer copy. Puts forwarded to an owner that dies mid-flight are retried.
Tradeoffs. Replication is owner-ordered and best-effort, not quorum-based.
- Remote copies add one synchronous network write each, sent one after another, so put latency grows with the factor.
- A copy that fails while a node is down leaves the write with fewer copies. Nothing re-replicates it later; the only repair mechanism is hand-back to a returning owner.
- Placement follows ring order, with no awareness of racks or other failure domains.
Space per synced byte is roughly 1 + replicas + remote copies, plus any cache copy.
Cache: the locality layer
The cache keeps a node-local, untransformed copy of each blob a node uses but does not own. Requests to the cache are routed to the submitter's node, so a process touching its own data stays local, and only the authoritative hop crosses the network.
Its central invariant is that a cache copy that exists is complete and current. Reads can therefore try the local copy without first probing its size.
Writes: write-through with owner-driven invalidation
- On the owner's own node, nothing is cached. The primary is already local, and a copy there would sit outside the invalidation protocol.
- Elsewhere, the local copy is updated first. If no copy exists and the put covers the blob from offset 0, one is created speculatively.
- The authoritative put then goes to the owner. The owner invalidates every other registered copy and registers this one, atomically under the blob's write token. Because the local write happens before registration, any later write by another node is guaranteed to invalidate it.
- If the authoritative put fails, or the blob turns out to extend beyond what this put wrote, the local copy is discarded. A cache problem never fails a put.
The cache was originally write-back, with a dirty set and a flush sweep. It was changed to write-through, so an acknowledged write is always on the authoritative path, and so 4 KiB writes no longer trigger re-pushes of the whole blob.
Reads: fetch, populate, register with a version check
On a miss, the cache reads from the owner straight into the caller's buffer, then populates the local copy. It registers the copy conditionally on the version it read. If the owner's content has moved on in the meantime, registration is refused and the copy deleted. A population that fails partway deletes the whole copy, because a half-written copy once served another file's old extents.
Retention
The cache has no LRU of its own. Copies are ordinary core replicas, with a score floor (default 0.5) that the organizer will not demote below. Only real capacity pressure in the core reclaims them, and they are the first thing reclaimed.
Tradeoffs. A remote-owned put costs one extra local write. A miss costs a remote read, a local write and a registration round trip. Partial writes to blobs with no existing copy are cached only when they are next read. On multi-node clusters, batched puts are not mirrored into the cache, because batches carry no origin registration.
Indexer: semantic search
The indexer owns the core's semantic-search operation. Before it existed, every search re-read and re-tokenized every candidate blob.
- A forward index. Each container keeps, for each blob it owns, a map of term frequencies and a document length. It is not an inverted index. Each query computes BM25 statistics over the documents that match its tag and blob name filters, so the matching documents must be visited anyway. What the index removes is reading and tokenizing them.
- Asynchronous indexing. A successful mutation enqueues its key in a coalescing pending set, so a hot blob overwritten many times is re-tokenized once. A sweep, every 100 ms by default, re-reads current bytes and updates the index. The first write to an unseen tag schedules a one-time backfill of its existing blobs.
- Read-your-writes search. A search drains the pending set first, filters by tag and blob regular expressions, and scores with BM25. The broadcast results are merged into a global top-k.
- Durable but derived. With an index log configured, the index is kept as a snapshot plus a log, and restored on restart without rescanning storage. An explicit reindex scan repairs anything missed in a crash window.
Measured. Indexing synchronously on the put path cut 1 MiB put throughput from 8.1 GB/s to 347 MB/s. Asynchronous indexing restored it to 3.64 GB/s, and 4 KiB puts reached 31.7k IOPS against 38.2k with no indexer. The remaining gap is the sweep cost under sustained overwrites.
Limits. Matching is lexical (BM25 over alphanumeric tokens), with no embeddings. Each shard computes term statistics over its own slice, so merged scores only approximate a global BM25. The index is held entirely in memory.
Compressor: the encoding layer
The compressor encodes data on the way down and decodes it on the way up. It is off by default, both in the build and in the shipped configuration.
- Which writes are compressed. Whole-blob writes with a codec requested. Partial writes, replica writes and puts with no codec pass through raw.
- Stored format. A small header (codec, preset, original and compressed sizes) followed by the codec output.
- Kept only if worthwhile. The compressed form must save at least an eighth, otherwise the data is stored raw. Weakly compressible pages, such as FP8 tensors at 93% of their raw size, had doubled runtime.
- Marked in the core. A persistent transform flag records that the blob is compressed, so every layer knows the stored form differs from the logical one.
- Codecs. Lossless CPU codecs (zstd, lz4, zlib, bzip2, lzma, snappy, blosc2), lossy scientific codecs (SZ3, ZFP, FPZIP) and GPU codecs (nvCOMP family, cuSZ, ndzip, a SYCL ZFP).
Reads. A read of a compressed blob fetches the whole stored blob, decompresses it once, and copies out the requested ranges. Small random reads of compressed data therefore pay for a full decode. The cache layer above offsets this by keeping a raw copy near the reader.
Dynamic codec selection. A feature-based predictor scores a few candidate codecs on samples of the data (entropy, deviation, smoothness) and picks either the best ratio or the fastest end-to-end time for the target tier. It is opt-in through a dedicated client and should be treated as experimental. Pinning a specific codec is the supported path.
Stream: sizes and ordered appends
Stream is not an interposer. It is a pool of its own that turns a tag into a byte stream, and it is the filesystem's source of truth for file sizes. It owns three things:
- The page layout: 1 MiB page blobs named by page index.
- The authoritative logical size of each stream. It changes by raise-only updates from writes, reservations from synchronous appends, explicit sets from truncates, and drops.
- A cross-node append pipeline, described below.
Each stream has a home container (for files, the inode's home).
Deferred appends
- Stage. The origin node stamps the append with a per-node monotonic clock and a counter, and writes the bytes as a self-describing staged blob through the replication layer. Once that put returns, a crash cannot lose the append, and the writer is acknowledged.
- Ship. Every millisecond, queued appends are grouped by stream and shipped in chunks of up to 16 MiB to the stream's live home. A failed shipment is retried ahead of newer appends, so retries never reorder.
- Merge. The home removes duplicates, sorts by (clock, origin, counter), reserves the tail, logs the plan, and copies the bytes into page blobs. The plan stays open until every copy succeeds.
The result is a deterministic total order: exact within one writer, timestamp order across writers. Truncates and drops wait for an in-progress merge, so they are ordered after it and never interleave with it.
Durability. Sizes and plans are kept in a record log, fsynced at least every 5 s and on file fsync. After a restart, staged appends are re-queued and interrupted plans finished. Restored sizes are held until the filesystem confirms them, so a truncate made elsewhere during failover is not undone by replaying old plans at stale offsets.
Tradeoffs. Every append is written twice (once
staged, once merged), in exchange for low-latency, crash-safe
acknowledgement on the origin. Each stream's home serializes its merges,
so a single very hot stream is limited by one node. Without a
log_path, sizes are not persisted across restarts.
Checkpoint: lazy copies through fault handlers
The core lets a tag register a fault handler: a pool that is asked to produce a blob when a read (or put) finds it missing. The checkpoint ChiMod is such a handler. A "copy" of a dataset is therefore a new tag whose handler points at the source. The first access to each blob copies the whole source blob into the new tag and serves the read from it.
- A lazy copy takes about 0.4 ms regardless of size. A forced eager copy takes about 9.5 ms for 64 MiB and about 148 ms for 1 GiB.
- It is copy-on-first-access from the live source, not a point-in-time snapshot. A blob first touched after the source changes gets the new bytes. Use the eager mode when a true snapshot matters.
- Handler registrations are not persisted, so a lazy copy should be materialized before relying on it across a restart. Concurrent first accesses to the same blob are not serialized.
Its main user is the GPU vector's copy operation. This is a young component.
Experimental components
| Component | Idea | Status |
|---|---|---|
| UVM (GPU virtual memory) | Reserve a huge GPU virtual range and back pages on demand. Pages are evicted to pinned host memory (or to the CTE) and restored when touched. | Prototype. The caller must touch pages explicitly, there is no eviction policy, and the CTE backing is not yet wired in. |
| LMCache backend | Store LLM KV-cache records as CTE blobs: hashed keys, a self-describing record, deferred vectored puts, and batches for small records. | The most developed of the LLM hooks. Each hit costs three gets. |
| llama.cpp KV and weight paging | Prefix-KV offload, and layer-by-layer weight paging over UVM. | Prototypes. They depend on an out-of-tree llama.cpp fork. |
| Python bindings | The core client API, with async futures and the GIL released around every call. | Maintained. Other ChiMods are reached by pointing the core client at their pool, for example the indexer for search. |
Tradeoffs at a glance
| Decision | Buys | Costs |
|---|---|---|
| Features as stackable interposers | Pay only for what you compose. Each layer can be tested in isolation, and clients need no changes. | Each layer adds a forwarding hop, and the stack's order is part of its correctness. |
| Write-through cache with owner invalidation | No dirty data outside the authoritative path. Readers never see stale copies. | Extra local writes, and coherence messages to the owner. |
| Synchronous remote copies that tolerate failures | Durable on ack in the common case, and available while degraded. | No quorum: degraded writes keep fewer copies and are not re-replicated. |
| Asynchronous indexing | Puts run near native speed. | A search pays the backlog drain, and the index is memory-resident. |
| Whole-blob compression | Simple and codec-agnostic, with a strong ratio. | Partial reads decode the whole blob, and partial writes are stored raw. |
| Stage-then-merge appends | Crash-safe appends from any node, in a deterministic order. | Bytes are written twice, and merging is serialized per stream. |