Design
Durability and recovery
CLIO Core's default durability contract follows ext4. An acknowledged write is fast and survives a crash of the application. It is guaranteed to survive power loss only after fsync (or the periodic flush), and it is unavailable while its node is down unless replication is configured. This page explains how each layer meets that contract, and where it deliberately does not go further.
Four levels of "safe"
| Level | What it takes | Default behavior |
|---|---|---|
| Acknowledged | The owner has applied the write. For deferred puts, the write has been accepted into staging. | Survives an application crash once it reaches the runtime. Deferred-put errors surface at the next flush or close. |
| Process-crash safe | Data is on a persistent device (or in the OS page cache for it), and metadata log records have been written to the OS. | Holds for data on persistent tiers. Every metadata log record is flushed to the OS as it is written. |
| Power-loss safe | Device data and logs have been fsynced. | Requires fsync (a tag sync), or happens through the periodic flushes: data every 10 s, the metadata log every 5 s. |
| Node-loss safe | A copy exists on another node. | Off by default (replication factor 1). Configure remote copies, and failover to serve them. |
Data on a volatile tier (RAM, HBM, pinned memory) reaches only the first level until a flush moves it to a persistent tier.
What is persisted, and where
Persistent state lives under a single storage root
(CLIO_STORAGE_ROOT or clio_run start --disk,
default ~/.clio):
| Artifact | Owner | Contents | Sync policy |
|---|---|---|---|
| Block-device data file | bdev (file) | Blob bytes on the disk tier | fdatasync on sync. Otherwise kernel writeback. |
| Allocation log | bdev (file) | Allocation and free records, 40 bytes each | Flushed per call. Fsynced every 50 ms and on sync. Compacted atomically. |
| Metadata snapshot | CTE core | Every tag and blob: layouts, scores, flags, aliases | Written to a temporary file, fsynced, then renamed into place. |
| Metadata write-ahead log (one shard per worker) | CTE core | Create, extend, clear and delete blob; tag create, delete and identity; transform and droppable marks | Flushed per record. Fsynced every 5 s and on sync. |
| Pool log | Runtime | Which durable pools exist, with their parameters | Fsync of the file and its directory on every append. |
| Stream size log | Stream ChiMod | Logical file sizes, append plans and completions | Fsynced at least every 5 s and on file fsync. |
| Search index snapshot and log | Indexer | Term statistics for semantic search | Derived state, which can be rebuilt by rescanning. |
| Hand-back log | Replication | Changes made while standing in for a dead owner | Persisted only when a path is configured. |
| Filesystem namespace | clio-fs | Directory blocks, inode records, extended attributes | Stored as ordinary CTE blobs, so it inherits the CTE's durability. |
Without a metadata_log_path, the core persists nothing. The
bytes on disk would survive, but nothing would record which blob they
belong to. The shipped configuration sets the path.
Ordering invariants
Most of the reliability design comes down to doing things in the right order. Each rule below exists because the opposite order lost data at scale or under a file-system stress test.
| Rule | Prevents |
|---|---|
| Allocation is logged before blocks are returned; frees are logged before space is reused. | Handing out live bytes twice after a crash. The worst a crash can do is leak space. |
| Data is written before its layout is logged. | A recovered blob pointing at a recycled extent that holds someone else's old bytes. |
| A new copy is placed before the old one is freed (moves, flushes, reorganization). | A failed move or flush destroying the only copy. |
| The zero-IPC mirror entry is withdrawn before blocks are freed. | A reader in another process copying from an extent already reused by another blob. |
| Device data is synced before the allocation log. | A logged block whose bytes were lost. |
| Snapshots are written beside the live file, then renamed; the log is trimmed only up to the snapshot's sequence number. | Losing records written while the snapshot was being taken. |
| Filesystem inode records are stored before the name is published. | A directory entry pointing at a missing inode. |
| Link counts are raised before a new name is inserted. | A dangling name after a crash. The worst outcome is an over-count, a leak. |
The metadata log
The metadata log is sharded by worker, so concurrent workers do not contend on one file. Every record carries a sequence number from a single per-container counter. Replay merges all shards and sorts by that sequence before applying anything. Replaying each shard in file order once let an older full-layout record overwrite a newer one and shorten a blob after a crash. Records that carry a blob's layout are full-state rather than deltas, so replaying one is idempotent.
Snapshots are expensive (about a second per 18k blobs before the fix that made them occasional). A full snapshot is therefore rebuilt only when one of these holds:
- the log has grown past its capacity (32 MB by default);
- unlogged changes such as scores are pending and the last snapshot is over a minute old;
- a snapshot is explicitly requested.
Snapshotting yields regularly so it does not stall a worker.
Reliability tradeoff. Records are length-prefixed but carry no checksum. A record whose length runs past the end of the file ends the replay cleanly, which covers torn tails. A corrupted payload of valid length, however, would be applied. If restore hits a corrupt log, startup logs the problem and continues with the state it could recover. That favors availability over strictness, and the operator is told to move the log aside to start clean.
What a restart does
clio_run start always recovers existing state.
--fresh is the only way to discard it.
- Pools come back. The runtime replays its per-node pool log, in creation order, and recreates every durable pool through its module's restart path. Pools from the server's compose configuration are treated as durable.
- Block devices rebuild their allocators from the allocation log. Gaps below the highest live byte become free space, and the heap resumes after the last live byte.
- The core restores metadata from the snapshot, then replays the merged log.
- Volatile data is accounted as lost, not as zeros. Blocks that lived on RAM, HBM or pinned memory are dropped. A blob keeps the prefix before its first lost block, and the rest is marked lost. Reads into the lost range fail with an I/O error. The replication layer can serve those bytes from another copy if one exists.
- Sizes and caches are made consistent. Tag sizes are recomputed in one linear pass (an earlier quadratic pass kept a 6-node cluster from ever recovering). Node-local cache copies are discarded, because invalidations may have been missed while the node was down.
- Restored space is reserved on persistent tiers that do not log their own allocations, so new writes cannot land on restored data. The cost is that holes below the reserved high-water mark stay unusable.
- Higher layers reconcile.
- The stream module re-queues appends staged before the crash and finishes interrupted merges. It holds size operations until the filesystem confirms each file's size, or two minutes pass.
- The filesystem resumes inode numbering past its durable reservation and destroys orphaned inodes.
- The replication layer pulls back any changes another node made while standing in for this one.
What survives what
| State | Application crash | Runtime crash (kill -9) | Power loss | Node loss |
|---|---|---|---|---|
| Data on a persistent tier | Yes | Yes | After sync or periodic flush | Only with remote copies |
| Data on a volatile tier | Yes | No; reported as lost | No | Only with remote copies |
| Blob and tag metadata | Yes | Yes (each log record reaches the OS) | Up to ~5 s of log may be lost; less after sync | Shadow copies under failover |
| Blob scores | Yes | Up to ~1 minute of changes lost | Same | n/a |
| Allocator state | Yes | Yes | ~50 ms window; immediate after sync | n/a |
| Pool registry | Yes | Yes | Yes | n/a |
| File sizes (clio-fs) | Yes | Yes, with a stream log configured | Up to ~5 s; immediate after fsync | Via failover |
| Access counts, cache registrations | Yes | No (rebuilt) | No | No |
What fsync means end to end
Through the FUSE filesystem, fsync on a file does the
following, in order:
- Drains the client's write-behind pages.
- Publishes the file's size.
- Waits for any pending asynchronous replication of the file.
- Syncs the file's tag: data is moved to a persistent tier and its devices are fsynced, then the metadata log is fsynced.
- Fsyncs the stream size log.
- Syncs the parent directory's tag, so the name is durable too.
If a node that may have held this file's unsynced writes died since the
last sync, fsync returns an I/O error rather than claiming
durability it cannot vouch for. Measured cost: syncing to the device adds
roughly 23–35% to fsync-heavy workloads such as sqlite, tar and appends.
Front ends differ. The LD_PRELOAD
interceptors (POSIX, STDIO) implement fsync as "drain the
deferred writes and report any latched error". That makes the write
visible and surfaces errors, but it does not sync the tag to a device.
MPI-IO's sync is a no-op. Applications that need power-loss durability
through these paths depend on the periodic flush. See
I/O adapters.
Node failure
- Default: a single copy. With the default replication factor of 1, a node's blobs are unavailable while it is down. Requests for them fail, after a network timeout or immediately in fail-fast mode.
- Remote copies. With factor N, the owner writes N−1 copies to the next nodes in ring order before acknowledging. A dead successor is skipped: the write succeeds with one copy fewer, and nothing re-replicates it later.
- Failover. With
failover_to_successor, requests for a dead owner's blobs go to its first live successor, which serves its shadow copy and records every change. When the owner restarts, it pulls those changes back before serving. The successor also pushes them periodically. - Peer restarts are detected by incarnation. Every message carries the sender's incarnation number. When a peer comes back as a new process, requests that were in flight to the old one fail promptly and are retried, rather than waiting forever for a reply that will never come.
- Membership is conservative. Liveness comes from probing peers that have outstanding requests or have gone quiet. The gossip-based failure detector is off by default, because at 256 nodes wide collective operations starved its probes and healthy nodes were declared dead. As a consequence, automatic container recovery seldom triggers. Durability across node loss is the replication layer's job, not the runtime's.
Tradeoffs at a glance
| Decision | Buys | Costs |
|---|---|---|
| ext4-style "durable on fsync" | RAM-speed writes. Applications choose when to pay for durability. | Seconds of exposure to power loss for un-synced data. |
| Worker-sharded log, globally sequenced | No single-file contention on the write path. | Replay must merge and sort all shards. |
| Snapshot only when needed | Low background cost at large blob counts. | Score changes can lag disk by a minute. |
| Continue on a corrupt log | The node comes back with whatever it could recover. | Partial state is served. There are no per-record checksums. |
| Lost volatile data is an error | Applications never silently read zeros. | Readers must handle I/O errors after a restart unless replicas exist. |
| Single copy by default | No write amplification or network cost. | A down node's data is unavailable until it returns. |
| Reserve restored space by high-water mark | Restored data can never be overwritten. | Holes below the mark are lost capacity on tiers without an allocation log. |