Design

CTE core: data model and data path

The core ChiMod is the storage engine at the bottom of the Context Transfer Engine. It names data, decides which device holds each byte, and keeps that mapping consistent under concurrency and failure. The layers above it (replication, caching, indexing, compression, streams, the filesystem) add policy on top. None of them bypass the core.

Data model

The core stores blobs inside tags, on storage targets.

Tags

A tag is a named container, roughly a bucket or a file. Its identifier combines the node that created it with a per-node counter, so tags can be created anywhere without a central allocator. Part of the counter space is reserved for client-minted identifiers, which lets batch creators (such as the filesystem's create path) propose IDs without a round trip. A tag records:

  • its canonical name and total size;
  • POSIX-style access, modify and change times;
  • a list of aliases (hard links);
  • an optional fault handler: another pool that is asked to produce a blob when a read misses. Lazy copies and checkpoints use this.

Hierarchical names are stored as chains. Each child tag is named relative to its parent's identifier, not by its full path, so renaming a parent rewrites one name rather than a whole subtree. A trigram index over tag names serves regular-expression tag queries without scanning every name.

Blobs

A blob is addressed by (tag, name). Its metadata is a list of blocks, each a (target, offset, logical size, physical capacity) extent on one storage device. A blob has no fixed maximum size and no fixed chunk size: it is as many blocks, on as many devices, as its size needs. Alongside the layout, a blob carries:

FieldPurpose
Score (0–1)How "hot" the blob is. It selects the tier during placement and drives reorganization.
Replica slotsIndependent block lists for extra copies: persistent replicas, node-local cache copies, and shadow copies kept for a failed owner. Each has its own score and flags.
Version (last-modified stamp)Used as the content version for conditional puts ("put only if still at version v") and for cache-coherence registration.
GenerationLets a reader wait for a writer to publish a specific generation (producer/consumer pipelines).
Transform flagsSticky marks such as "stored compressed", set by layers above and persisted in the log.
Droppable flagFixed when the blob is created. Marks the blob as a cache that may be evicted to make room. Authoritative data is never droppable.
Concurrency wordsA write token, a reader-pin counter, a placement generation and a content sequence number (see Concurrency).
Lost-bytes markerAfter a restart, the range that lived on volatile memory and is gone. Reads into it fail with an I/O error instead of returning zeros.

Each block records both its logical size and its physical capacity. The slack in the last allocation can then absorb a later append in place, without a new allocation, which makes small sequential appends cheap.

Storage targets

A target is a block-device pool (RAM, HBM, pinned host memory, file, S3, GCS, or a no-op device) registered with the core. Each target has a score that orders the tiers, a persistence level (volatile or durable), a capacity, measured performance figures, and a predicted device lifetime. The score can be set in the configuration. If it is not, the core derives it from measured bandwidth on a logarithmic scale.

Free space is tracked in two places. An authoritative counter is debited and credited lock-free on the data path. A mirror inside the vector that the placement engine reads is refreshed periodically. Placement always re-reads the authoritative counter first: a stale mirror once made a full tier look permanently full, even after its blobs were deleted.

Tag name, aliases, times size share (per container) optional fault handler Blob (tag, name) score, version, generation primary block list replica slots (persistent / cache / shadow) write token, reader pins transform + droppable flags many Block target, offset size / capacity Block target, offset size / capacity Target: RAM (score 1.0) volatile, node-local Target: NVMe (score 0.6) persistent Target: file / S3 (score 0.2) persistent, capacity tier
A blob's bytes can span several blocks on several tiers. The blob's score picks the preferred tier, and capacity decides whether it spills to others.

Distribution across a cluster

The core runs one container per node, and there is no metadata server.

  • Each blob has exactly one owner. The owner container is chosen by hashing (tag, blob name) over the number of containers. That container holds the blob's metadata and its primary layout, and every operation on the blob (put, get, delete, reorganize, truncate) is routed to it.
  • Tag names hash to an owner too. Get-or-create of a tag that is not already cached locally is routed to the container chosen by the name's hash. That keeps tag creation race-free without a central service.
  • Tag sizes are distributed. A container that holds some of a tag's blobs keeps a share of the tag's size. Asking for a tag's size is a broadcast that sums the shares. Dead nodes are skipped rather than waited on, so the answer can be low while a node is down.
  • Queries are broadcasts. Tag queries, blob queries, listing a tag's contents, semantic search and eviction fan out to every container and merge the replies.

Where the bytes go: the neighborhood

The owner decides placement, but it does not have to use only its own devices. Each container registers its own targets plus those of the next k−1 nodes in ring order (targets.neighborhood). Two rules constrain this:

  • Memory tiers are always node-local. RAM, HBM and pinned-memory tiers are never borrowed from a neighbor; only persistent tiers span the neighborhood.
  • Dead nodes' devices are filtered out before placement.

The shipped configuration uses a neighborhood of 1, so the owner writes only to its own node's devices and placement is entirely owner-local. A larger neighborhood pools persistent capacity across nodes, at the price of network hops on the data path.

Failover

When targets.failover_to_successor is on, requests for a dead owner's blobs are routed to the first live successor container. That successor holds a shadow copy written by the replication layer. Shadow copies are hidden from listings and size totals until their holder is actually standing in for the owner. The replication page describes how changes are handed back when the owner returns.

Tradeoff: single owner per blob. One container serializes every mutation of a blob, which makes coherence, versioning and crash ordering simple to reason about. The cost is that one hot blob cannot be written in parallel from many workers. Vectored puts, client-side write coalescing (the "sieve"), the deferred-put pipeline and optional worker-level batching all reduce how often that serialization is hit, rather than remove it.

The client side

Clients talk to the runtime through tasks (see Runtime). The core client adds four mechanisms that keep most small operations off the critical path:

MechanismWhat it doesWhy
Zero-IPC readsWhen a blob's metadata is mirrored in shared memory and its bytes sit in the local RAM tier, the client copies them out directly, without a task. It validates the placement generation and content sequence number before and after the copy, and any doubt falls back to a normal request.Hot reads cost a memory copy instead of a round trip (about 10 µs versus about 100 µs or more).
Deferred putsA put is copied into shared-memory staging and returns. A process-wide registry tracks in-flight puts, applies back-pressure based on staging capacity, serves reads of pending bytes from staging (read-your-writes), and latches errors for the next flush.Hides write latency the way a page cache does.
Write sieveCoalesces small writes into 64 KiB pages (up to 16 per blob and 1024 overall) and ships each page when it fills or on a short timer.Many 4 KiB writes become a few large puts.
Private-memory putsCallers can put from ordinary process memory. The runtime copies once into shared memory, or not at all when the caller is inside the runtime.Callers (including Python) never manage shared memory themselves.

The shared-memory metadata mirror is a derived, best-effort copy. Its tables are sized once at startup (64k tags and 256k blobs by default), do not grow, and hold at most 16 blocks per blob. Anything larger or overflowing is served by a normal request, never treated as absent. The mirror reserves on the order of 100 MB up front.

The put path

A put that reaches the owner proceeds as follows:

  1. Validate and resolve the score. An unspecified score keeps the existing blob's score. A brand-new blob defaults to 1.0, the hottest tier. Vectored puts (several disjoint segments) are normalized to their covering extent.
  2. Consult the fault handler if the blob is missing and the tag has one, so lazily copied data is materialized before it is overwritten.
  3. Create the blob if needed with insert-if-absent semantics, so concurrent creators converge on one blob instead of one silently replacing another.
  4. Acquire the blob's write token. This serializes all layout and size changes to the blob.
  5. Check conditions. Put-if-absent and put-if-version fail with distinct return codes when the condition does not hold.
  6. Drain readers. In-flight readers of the blob are allowed to finish, and the content sequence number is made odd for the duration of the write. Without this, a concurrent reader could see half old and half new bytes.
  7. Place. If the put grows the blob, the placement engine allocates new space (see Placement). A put that overwrites existing bytes skips placement entirely. If there is no room, the put evicts droppable data on the local node and retries once.
  8. Write. One asynchronous device write is issued per touched block, and all of them are awaited together. Writes past the current end of the blob zero-fill the gap, giving sparse-file semantics.
  9. Log after the data. The new layout is recorded in the write-ahead log after the bytes are on the device. Logging first could, after a crash, point the blob at a recycled extent that still holds someone else's old bytes.
  10. Roll back on failure. If anything fails after placement, the blob is shrunk back to its previous size and the rolled-back layout is logged. Only growth can be undone: a failed overwrite inside the old size leaves those bytes in whatever state the device reached.
  11. Publish and invalidate. The version is bumped, statistics are updated, the shared-memory mirror is refreshed, and every registered remote cache copy is invalidated before the acknowledgement.

All-or-nothing placement

When a blob grows, all the new blocks, the reuse of slack in its last block and its new size are published as one step, with no suspension point in between. If any part of the allocation fails, the partial allocation is freed and the blob is restored exactly. No reader ever sees a blob that is longer than its blocks.

Startup and slow operations

A put that arrives before the core has finished registering its targets waits (up to two minutes by default) instead of failing. This covers clients that start before a large cluster has finished starting up. Any put that takes longer than two seconds logs a phase breakdown (token wait, placement, write), so slow paths can be diagnosed in production.

The get path

  • Generational reads wait, with back-off, until the blob reaches the requested generation. The wait is bounded (two minutes by default).
  • Missing blobs go to the tag's fault handler if there is one. Otherwise, on request, an empty blob is created.
  • Reads never take the write token. They pin the blob, re-check that it still exists under the pin, snapshot the layout and read the blocks in parallel. Reads of the same blob proceed concurrently. They wait only for operations that free extents, and for in-place overwrites, both of which drain pins first.
  • Lost data is an error, not zeros. A read that falls into a range lost in a restart (for example, bytes that lived only in RAM) fails with an I/O error code.
  • The version is reported with the data. The version is read before the layout snapshot. A racing write can therefore only make a version-checked cache registration fail safely, never succeed with stale data.

Delete, truncate and links

Deleting a blob takes the write token and drains readers. It then does the following, in order:

  1. Invalidates remote cache copies.
  2. Folds replica blocks into the free list.
  3. Withdraws the shared-memory mirror entry, before any block is freed, so a zero-IPC reader can never copy from an extent that has already been reused.
  4. Frees the blocks and updates the tag's size share.

Deleting a tag distinguishes three cases:

  • Removing one alias, which leaves the tag in place.
  • Removing the canonical name while aliases remain. An alias is promoted to canonical, as in POSIX link semantics.
  • A full recursive delete. This uses a striped per-tag blob index so it does not scan every blob on the node, and it prunes ancestors left empty.

Truncation shares the resize path with replacing puts.

Concurrency

Handlers are coroutines. A blob's tasks all land on its owner container, but the runtime's elastic scheduler can run them on different worker threads at the same time. The core therefore never relies on "no suspension point, therefore atomic" alone. It uses four explicit mechanisms:

MechanismProtectsBehavior
Per-blob write tokenEvery change to one blob's layout or size: put, resize, delete, reorganize, flush to persistent tier, snapshot.Reentrant. Contenders poll with a cooperative yield (10 µs by default). That avoids lost wake-ups but gives no FIFO fairness.
Reader pins with a drain bitReaders against operations that free extents or overwrite in place.Readers never block writers. A writer that must exclude readers sets the drain bit, so new readers back off, and waits for the pin count to reach zero.
Placement generation and content sequence (seqlock)Zero-IPC readers in other processes.The placement generation is globally monotonic, so a delete followed by a re-create cannot be mistaken for the old blob. The sequence number is odd while a write is in progress. Readers compare both values before and after copying.
Self-locking maps holding shared pointersMetadata lookups.Each map operation locks internally. Erasing an entry never frees an object that another task is still using. Renames replace a tag's entry instead of mutating it in place.

Operations on different blobs, reads of the same blob, and the per-block device I/O within one request all run in parallel. Mutations of one blob are serialized.

Tradeoffs at a glance

DecisionBuysCosts
Hash-owned blobs, no metadata serverLinear metadata scaling, no single point of failure, no directory service hop.Tag size and listing are broadcasts. Scans such as eviction and organizer rounds are O(blobs) per node.
Blobs are variable-length extent listsNo fixed chunking, cheap appends into slack, a blob can straddle tiers.The shared-memory mirror caps fast-path blobs at 16 blocks. Fragmented blobs fall back to requests.
Log after data, place before freeNo crash can make a blob reference another blob's old bytes.Moves need transient double space, and log records trail the data slightly.
Readers never take the write tokenRead-heavy workloads scale with workers.In-place writes must drain readers, and token contenders spin-yield with no fairness.
Rollback covers growth onlyCheap, with no undo log of old bytes.A failed in-place overwrite is not atomic, so callers that need atomic replace must write elsewhere and switch.
New blobs default to the hottest tierFresh data is fast to read back.The top tier fills first and later writes spill down. The organizer or explicit scores must demote cold data.