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"

LevelWhat it takesDefault behavior
AcknowledgedThe 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 safeData 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 safeDevice 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 safeA 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):

ArtifactOwnerContentsSync policy
Block-device data filebdev (file)Blob bytes on the disk tierfdatasync on sync. Otherwise kernel writeback.
Allocation logbdev (file)Allocation and free records, 40 bytes eachFlushed per call. Fsynced every 50 ms and on sync. Compacted atomically.
Metadata snapshotCTE coreEvery tag and blob: layouts, scores, flags, aliasesWritten to a temporary file, fsynced, then renamed into place.
Metadata write-ahead log (one shard per worker)CTE coreCreate, extend, clear and delete blob; tag create, delete and identity; transform and droppable marksFlushed per record. Fsynced every 5 s and on sync.
Pool logRuntimeWhich durable pools exist, with their parametersFsync of the file and its directory on every append.
Stream size logStream ChiModLogical file sizes, append plans and completionsFsynced at least every 5 s and on file fsync.
Search index snapshot and logIndexerTerm statistics for semantic searchDerived state, which can be rebuilt by rescanning.
Hand-back logReplicationChanges made while standing in for a dead ownerPersisted only when a path is configured.
Filesystem namespaceclio-fsDirectory blocks, inode records, extended attributesStored 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.

RulePrevents
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.

  1. 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.
  2. 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.
  3. The core restores metadata from the snapshot, then replays the merged log.
  4. 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.
  5. 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.
  6. 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.
  7. 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.
Pool logrecreate pools Alloc logrebuild heaps Snapshot + WALsequence-ordered Drop volatilemark lost bytes Reserverestored space Reconcilestream, fs, repl.
The order of a recovering start. Client requests that arrive early are held until the module that owns them has finished its restore.

What survives what

StateApplication crashRuntime crash (kill -9)Power lossNode loss
Data on a persistent tierYesYesAfter sync or periodic flushOnly with remote copies
Data on a volatile tierYesNo; reported as lostNoOnly with remote copies
Blob and tag metadataYesYes (each log record reaches the OS)Up to ~5 s of log may be lost; less after syncShadow copies under failover
Blob scoresYesUp to ~1 minute of changes lostSamen/a
Allocator stateYesYes~50 ms window; immediate after syncn/a
Pool registryYesYesYesn/a
File sizes (clio-fs)YesYes, with a stream log configuredUp to ~5 s; immediate after fsyncVia failover
Access counts, cache registrationsYesNo (rebuilt)NoNo

What fsync means end to end

Through the FUSE filesystem, fsync on a file does the following, in order:

  1. Drains the client's write-behind pages.
  2. Publishes the file's size.
  3. Waits for any pending asynchronous replication of the file.
  4. Syncs the file's tag: data is moved to a persistent tier and its devices are fsynced, then the metadata log is fsynced.
  5. Fsyncs the stream size log.
  6. 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

DecisionBuysCosts
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 sequencedNo single-file contention on the write path.Replay must merge and sort all shards.
Snapshot only when neededLow background cost at large blob counts.Score changes can lag disk by a minute.
Continue on a corrupt logThe node comes back with whatever it could recover.Partial state is served. There are no per-record checksums.
Lost volatile data is an errorApplications never silently read zeros.Readers must handle I/O errors after a restart unless replicas exist.
Single copy by defaultNo write amplification or network cost.A down node's data is unavailable until it returns.
Reserve restored space by high-water markRestored data can never be overwritten.Holes below the mark are lost capacity on tiers without an allocation log.