Erasure coding

Upcoming. This page describes the master branch. Erasure coding belongs to network storage, which is not part of release 0.0.8, and details may change.

Erasure coding is how Hedgehog keeps a stored file alive on machines run by strangers. A file is encrypted, cut into pieces and spread over the network together with extra “parity” pieces, so the original can be rebuilt even when a good share of the machines holding it have gone offline or lost their data. It lets the network store files durably without keeping full copies, and without any gridnode seeing the content.

This page covers the coding itself. Where the pieces are placed, fetched and repaired is described in Network storage, and the surrounding system in Architecture overview.

Why not just copy the file

Gridnodes restart, lose disks and leave. The obvious answer is replication: keep three copies and survive the loss of any two, at three times the storage. An erasure code does better. Cut a block into k equal pieces and compute m further pieces from them, so that any k of the k + m pieces are enough to recompute the block. Losing a piece is then only a known gap that the reader solves for, and the code tolerates the loss of any m pieces, whichever they are.

Hedgehog uses Reed-Solomon coding, implemented in the project itself with no external library, and applies it in two layers. Together they store 2.25 bytes per file byte at the default settings, and survive more loss than three replicas at 3.00.

How a file is stored

A 100 MiB file at default settings: 101 data chunks in four stripes, outer parity per stripe, and one chunk
coded into 32 fragment slots

  1. Encrypt. The file is split into chunks of about 1 MiB, and each is encrypted with AES-256-GCM. Every upload draws a fresh random secret, the fingerprint, and all keys derive from it. Without it, nothing the network holds can be read or tied to a file.
  2. Outer code, across chunks. Consecutive chunks form a stripe of up to 32 data chunks, and parity chunks (50 percent, so 16 for a full stripe) are computed over the encrypted data. Only the owner knows which chunks belong together, so only the owner can apply this layer. It covers the loss of whole chunks.
  3. Inner code, within a chunk. Each chunk, parity chunks included, is cut into 16 slices of 64 KiB and extended to up to 32 fragments. This layer is public: any gridnode holding 16 verified fragments of a chunk can rebuild the others, which is how the network repairs itself while the owner is offline.
  4. Make it verifiable. Every chunk gets its own Ed25519 key. Its fragments are committed to in a Merkle tree whose root is signed in a small descriptor that each fragment carries, so a gridnode can check any fragment on its own. A repaired fragment is byte for byte the one the owner sealed.
  5. Keep a manifest. The file size and the layout used are stored encrypted in three small extra groups, so a reader can rebuild the same shape later even if the network settings have changed since.

Reading reverses this. Normally a reader fetches 16 verified fragments per chunk and decrypts, and the outer code is not touched. If a chunk cannot be opened, the reader fetches parity chunks and rebuilds it in memory.

What survives

Level What can be lost Stored per byte
One chunk, inner code Any 8 of its 24 guaranteed fragments 1.50
One full stripe, outer code Any 16 of its 48 chunks 1.50
Both layers together Both of the above 2.25
Three replicas, for comparison Any 2 of 3 copies 3.00
  • Never wrong bytes. Corrupt or forged fragments are detected and dropped. A file that has lost too much is reported as lost, and is never returned with wrong content.
  • The price is repair traffic. Rebuilding one lost piece means reading 16 others, so over time the network traffic of repair, rather than disk space, is the main cost of redundancy.

Key facts

Fact Default
Chunk size 1 MiB (1,048,576 bytes), 16 of them the authentication tag
Fragment size 64 KiB of data, 65,838 bytes encoded with descriptor and proof
Inner code 16 data and 8 guaranteed parity fragments, plus up to 8 extra: 32 slots
Outer code Up to 32 data chunks per stripe, with 50 percent parity
Manifest copies 3
Stored for a 100 MiB file 155 chunk groups, about 245 MB, or 2.34 bytes per file byte
Stored for a large file About 2.26 bytes per file byte, proofs included
Small files A fixed overhead of three manifests and a one-chunk stripe: a 1 KiB file stores about 7.9 MB
Arithmetic Reed-Solomon over GF(2^8), at most 255 pieces per code

These values are fields of the storage spork, which the network can change; see Grid sporks. A file records the layout it was written with, so files written under earlier settings stay readable.

In the source

  • application/src/main/java/org/unigrid/hedgehog/model/storage/erasure/ReedSolomon.java is the codec both layers use.
  • application/src/main/java/org/unigrid/hedgehog/model/storage/erasure/GaloisField.java is the byte arithmetic beneath it.
  • application/src/main/java/org/unigrid/hedgehog/model/storage/ChunkGroups.java seals a chunk, opens it and rebuilds fragments.
  • application/src/main/java/org/unigrid/hedgehog/model/storage/LayoutParameters.java derives every count and size.
  • application/src/main/java/org/unigrid/hedgehog/service/storage/StorageUpload.java applies the outer code on upload.
  • application/src/main/java/org/unigrid/hedgehog/service/storage/Retrieval.java reads a file back, including degraded reads.
  • documentation/storage-white-paper/sharded-and-redundant-storage.tex holds the design and its analysis.