Non-local redundancy: Erasure coding and dispersed replicas for robust retrieval in...
ETHBerlin·Mon, Jun 16, 2025, 04:45 PM · 24:26
.. the Swarm peer-to-peer network. In this presentation, we describe in detail how non-local redundancy is implemented in the Swarm peer-to-peer network. Among other things, we show how to apply classic Reed-Solomon erasure codes to files in Swarm, discuss security levels of data availability and derive their respective parameterisations, describe a construct that enables cross-neighbourhood redundancy for singleton chunks to complete erasure coding, and explore alternative retrieval strategies applicable to erasure-coded files and their impact on latency as well as price of retrieval.
Transcript
There's the idea of the World Computer from the beginning, so yeah, that's what we will talk a little bit about here, and now I have to turn around to read out this title, so the title will be Non-Local Redundancy, Erasure, Coding, and Dispersed Replicas for Robust Retrieval in the Swarm Peer-to-Peer Network. So yeah, please welcome Yuri to the stage. Thank you very much, everyone. I'm an uglier and slightly less knowledgeable version of Viktor, but I'll do my best. So this one is not working.
Nope. The lightsaber works. There we go, okay. So if you don't know Swarm, it is, as the previous talk mentioned, there's an attempt to make a world computer, Swarm is attempting to make a world hard drive, the idea being that anyone can chip in with their own storage device, and then it's a decentralized storage where all data gets chopped up into small chunks and gets scattered around the network, and then you can retrieve your data there, therefore secure. Anyone who hacks into a computer only sees little meaningless individual shards of data, so that is the idea.
And what I want to talk about today is redundancy, data redundancy in such a system. And in particular, we're going to be going over three topics, as you see, namely what erasure codes are, and once we understand that, how they are implemented in Swarm, and finally, what are some of their properties that are important for end-users in particular, such as how much extra storage overhead do you need to buy if you are trying to take advantage of this. All right, so what's wrong with good old data duplication? By that I mean, it's very easy to ensure data redundancy in a system. Let's say we have three little chunks of data.
These three little chunks represent one single file, for example, that's chopped up into three bits, scattered around the network, and then we could simply duplicate them and have it at that. What's wrong with this? Well, the problem is, first of all, this is very wasteful, because you have to duplicate all the data, and second of all, it's actually not very secure. The reason being is that, if you think about it, let's imagine that the computers storing these data are churned, or for some reason are not available, and by unlucky chance, the computers are down such that both of these chunks, three are missing, and you only have available chunks one and two. In that case, you still cannot recover your data, so not only is this system of simple-minded duplication very wasteful in terms of storage, it's also not as secure as it could be.
Enter the idea of erasure codes, which is actually an old idea, dating back to a famous paper by Reed and Solomon in the 1960s. The idea is the following. Let's have the exact same file that we had before, chopped up into three little chunks, and we will pad it out with a number of extra chunks. For illustration's sake, on this slide, two extra little chunks. You notice that I've denoted them with blue, and the idea is that these two blue chunks are so-called parity chunks, and now, the big idea is that any three out of these five are sufficient to recover the original data, so this erasure coding employs some form of encoding which allows you to reconstruct the data no matter which three out of these five you have.
So again, to repeat, erasure coding is about having a number of original data chunks, padding it out with another number of parity chunks such that any of the original number of chunks is enough to recover the data, so it doesn't matter whether I have access to two, four, and three, or one, four, and five, any three of them will be able to recover the data. I don't have time to go into the details of how this does its magic, but I'm happy to talk about that afterwards if people are interested. Okay, so this was what erasure codes are. The next question is how we can implement them in Swarm, and in order to understand that, let's first understand how the file system works in Swarm. So as I mentioned earlier, all files are chopped up into small four-kilobyte little chunks, and those are randomly scattered across the network, and the way the files are represented is through what is called the Swarm hash tree, and I'm going to illustrate that for you.
There we go. This is a Merkle tree, not a binary one, but a 128-ary one, so it's a tree with a balancing factor of 128. It all starts from—not that—it all starts—there we go—from a root hash, and then it points to various references, and the references point to other references and so on until we reach the actual data that's in the bottom row that we have there, the actual data chunks that we have. I always use the convention that these rounded corner rectangles are the actual chunks of data, so you see that we have slightly larger but still rounded rectangles here, those items there. Those are so-called packed address chunks, also known as intermediate chunks, so those are the chunks that are sort of in the network, and their role is to give the references to the levels—to the other levels, and then at the end we reach, of course, the actual data that we want to retrieve.
So this is the structure of files, so this is, again, a Merkle tree that we have. Each of these so-called packed address chunks that are pointed out by the arrows contain 128 references. Why 128 references? The symbol is very prosaic. Every reference takes up 32 bytes of memory, and then because every chunk has to be 4 kilobytes in size, as mentioned, that means that 128 references then fit into a single chunk to make up for the 4 kilobytes.
Parenthetical little comment. There's also an option to get encrypted information, where it's not just the address but also an encryption key that one needs to store in every reference. In that case, the size of every reference, instead of being 32 bytes, is increased to 64 bytes. Quick calculation will reveal that then, instead of 128 references, you can pack in 64 references for each of these intermediate packed address chunks. Parentheses closed.
I'm sometimes going to mention encrypted data, but mostly for the simplicity of the examples I will just stick to normal unencrypted information. All right, moving on. Here is the key idea for how we can implement erasure coding in a decentralized storage system such as Swarm, which works with these Merkle trees to store data. The idea is as follows. We take each of these intermediate packed address chunks that we've seen before, those ones.
They have 128 references, and the idea is very simple. Instead of using all the 128 of them to be references, we only use a couple of them, and the rest we reserve for parities such that then any combination of M references will be able to recover the original data. So just to give you a very simple example, let's assume that this is our packed address chunk here with all the 128 references, and let's assume that this number M is 100. That means 100 references will be, quote, real references, and the other 28 will be parities such that any 100 out of the 128 will be able to recover the references. There you go.
This is what it means, that you have 100 original and 28 parity chunks, or, well, references to chunks in this case. This is the key idea. This is how we implement erasure coding in Swarm, and then what that actually means is that the real structure in the presence of erasure coding of the Swarm hash tree is like this. We have the root hash, and then we still have these intermediate packed address chunks, but in each of them, some of the references are parity references. This also means that these lines that I used to connect various references with one another are slightly metaphorical because, of course, in the presence of erasure coding, it's no longer true that this reference strictly refers to that.
Instead, the references are sort of encoded in all of the data out there such that even if for some reason we're missing some of the information, we can always reconstruct the required information from just 100 of those references out of the 128. All right. There is one other one out that I will mention just very quickly, the root hash. The root is one single chunk, but it's not structured like the other ones. It doesn't have those references.
Therefore, it has to be treated in its own way. How can we do that? A very simple first idea might be that you have is that we can simply duplicate the root hash, have many copies of the same thing, except unfortunately in the context of the swarm system, that doesn't work. The reason is that chunks with the same hash are automatically deduplicated by the system. How can we solve this problem?
I'm not going to go into the gory details of how it is done, but to give you an idea of what's done, the sort of outline of the idea is quite simple, which is that we have to tag on a hash onto the actual payload, the actual information contained in those chunks, such that A, the owner is clearly identifiable, and B, it is such that these quasi-duplicate chunks are uniformly scattered around the network such that you have an equal probability for any node to receive them. We really don't want a situation where all of them go to one computer so that if one computer is down, then we lose the data. We want it to be nicely and uniformly distributed. This kind of solution is what we call dispersed replicas. These quasi-copies of the same chunk with the added on tag are what we call dispersed replicas in the system.
That is how we solve the problem of replicating the root hash. In the case of the root hash, really, these quasi-duplicates are good enough because it's a single chunk, so there's not much meaning to erasure coding beyond that because erasure coding is, of course, a situation that's useful when you have multiple chunks and you want to create redundancy. There, here, we just have one, the root hash, and we can create these quasi-duplicates of them. All right. How many parities should there be in the system to provide adequate security?
Well, in order to answer that question, we have to first agree on what we mean by adequate security. Somewhat arbitrarily, we agree that adequate security will be 99.9999, so six nines, you can keep that in mind, percent chance that we can correctly recover the data. In other words, probabilistically, I could also say that there's a 10 to the minus six probability of a fault going on somewhere. All right.
So that is our target. In order to keep this target, we have to make some assumptions about how likely it is that any single chunk cannot be correctly recovered. We assume that these probabilities are equal and independent from one another. That's, of course, a bit of an abstraction, but with those assumptions, we are able to work out how many parities out of 128 there should be in each content address chunk. Let me give you a little example to illustrate how this is supposed to work.
So example number one, we have a per-chunk probability of failure of one percent. Again, that's an assumption, or it might be just a measured piece of data on the network. It doesn't really matter. We assume that any one chunk, when you recover data, has a one percent chance of going awry. In that case, one can calculate, and you don't even need to work very hard for it, that you need nine parities in order to ensure 99.
9999 percent probability of correct data retrieval. On the other hand, if the per-chunk probability of failure is 50 percent, because say there's a giant schism that happens in the network, and half of the network is churned, and there's an apocalypse going on, then one can compute that you need actually 90 parities. So you'll have 38 original references padded by 90 parities, and that will actually protect you even against such a very harsh failure of the system. I know I have not gone into any details as to how I obtained these numbers, as if by these numbers nine and 90 and others. I am happy to talk about this ad nauseum.
Look for me after the talk, but I don't think it's particularly interesting for this context here. I'll just mention this for those of you who might be interested that really it's the binomial distribution that you need to look for here, because we assume that we have 128 independent references and each of them has the same probability of going wrong, and we're asking the question something bad happens, namely that's the number of parities that we're obtaining. Okay. Actually, let me go back one step. These are just two particular examples for a failure probability of 1% and 50% for any individual chunk.
But we can refine this picture with this graph. So in this graph, on the x-axis, you see the per-chunk error rate smoothly increasing from 0% to 75%. As we increase this, occasionally there's a jump in the number of required parities, which is an integer, and that's what we have over here, going from 0 to the maximum 128. So as we increase the per-chunk error rate, the required number of parities is also increasing. That is, the required number of parities to achieve a 99.
9999% probability of overall correct data retrieval. In SWARM, to make it easier for users, we're not using this whole graph. We're actually just picking up out a few of these slices of the graph, namely 0%, 1%, 5%, 10%, and 50% error rates. And this little table gives you an idea of how that works. So I've collected out those percentages here.
So, can I go back? There we go. Those values where these vertical lines cross are exactly the values that you see in the table down there, 0, 1, 5, 10, and 50%. Here you have how many parities there are for those values and how many non-parity chunks you have with that. The sum of those two should always add up to 128.
If they don't, I've made a mistake. So you can double-check. For the sake of SWARM users, we wanted to give these more memorable names than 1% chance of error or 5% chance of error. So these are the actual names that you will find in the system. No erasure coding.
So if you assume that there's no reason to think that chunks will go wrong, you can choose that, and then at your own risk, you can have no erasure coding, small parenthetical comment. Even then, you're somewhat protected because the way SWARM works is that nodes are organized into so-called neighborhoods, and in each neighborhood, the data should be just replicated across the nodes that belong to the same neighborhood. So you're still somewhat good, but just do this at your own risk, of course. And then we assume a 1% error rate, which leads to nine parities. That's what we call a medium encoding.
5% error rate, in which case you need 21 parities. That's strong. 10% with 31 parities. That's insane. And the 50% one, if you're worried about the network just going through a major schism and breaking in half, that's the so-called paranoid setting.
So you can get that, although you will need to buy quite a bit more data if you do that. Speaking of data overhead, let's actually see how much storage overhead there is because of using erasure coding. By overhead, what I mean is that if you didn't use erasure coding, you could use all of these references to represent actual data. If we do use erasure coding, then of course we can only use some of those references as The rest are parities. So we can compute how much overhead we have compared to not having used any of the erasure coding.
I've pre-computed those percentages there for you. So of course, there's none. If you don't have erasure coding, there's about 7.6%. If you have medial encoding, 20% for strong, 32% for insane, and 237% for paranoid.
That's quite a bit, but you're of course very safe in return. We can also do a bit more precise computation as for how much overhead there is because we know that this is the 128th airy tree that we need to look at. And therefore, at every level of the tree, you have 1 over 128 times the number of chunks that you had in the previous level, which simply means that because we have this tree with branching factor of 128, we simply, to compute the effective size of the file you need, you take the raw size. For example, you want one gigabyte, take the raw size, multiply that with the erasure overhead. That is 1 plus those percentages that were displayed over there.
And then we have to take care of the tree so that we have 1 over 128 plus 1 over 128 square, et cetera, until the depth of the network. We don't know how deep the network is. What we don't need to know because this series actually converges very, very fast. We can just take a conservative estimate and say it's infinitely deep. If we do that, actually the value of this happens to be exactly 128 over 127 or 1.
008. Small comment for encrypted, it's double that, so it would be 1.016 for encrypted content. For non-encrypted, it's 1.008.
That's our multiplying factor. In other words, this is the formula you use to actually compute your effective file size. With that, thank you very much for your attention. I would like to thank my co-authors. First and foremost, Victor Thrun, and then the other co-authors I've had.
I'd really like to thank the whole Swarm Research team and the Swarm B team who've made this possible and implemented these solutions. I'd like to thank the organizers and not least you all for listening, and I'm happy to take your questions. Thank you very much. All right. Thank you very much.
Yeah, maybe, Victor, you want to join on the stage? All right. We have three questions so far. Anyone can, of course, add their questions by scanning the QR code. I think first one, just a quick one, like by error rate, do you mean nodes going offline?
When you say error rate, do you mean nodes going offline? It could mean yes, for example, but there could be other reasons as well. I don't know. Some network error, but essentially nodes going offline. Yes, that would be the most typical kind of reason why you cannot recover a chunk.
I guess corruption of the storage of a given computer is also a possible reason. All right. Then we have, is the vision of Swarm still to be the storage layer of the world computer like it started out to be? I will have that too, Victor. Yes, if you want it to, everything.
It's ready and waiting for adoption. So adopt it and it will be then. That's one of the more difficult. It's not only the storage layer. In fact, it's, I don't know if you know the original Holy Trinity.
There was a bit of a talk about Whisper and the communication layer. So I don't know if those of you who were in Prague like two weeks ago, I gave a talk on the Web3 Privacy Now event about signal on Swarm, how to do a signal protocol on Swarm and how to even improve it. So basically, you can have a private, very strong privacy messaging protocol on Swarm. All right, thank you. I think this one is also just interesting in general.
I think a lot of projects have been around for a long time, have faced a lot of different challenges. So what has been the main challenges of building Swarm now for more than 10 years? Was it unexpected complexities or just new ideas all the time? I don't know if new ideas all the time is a problem, but yes, kind of. So indeed, unexpected complexity, that was first.
And I mean, indeed, there was a lot of things that had to be kind of invented in a way. And what we came up with in the end is a very, very robust system and in a way, very, very elegant in its structure. And indeed, it took a while to perfect it. So especially the incentive system that is now complete, the ratio coding and yeah, all the rest of the system is kind of, there's many, many moving parts. And that's why the clients are kind of more difficult to build.
Of course, there's only like the core layer, which is basically the disk protocol, which is the distributed immutable storage of chunks. And that layer is complemented by a second layer, which kind of talks about more meaningful units to users like files, and then directories, etc. And then all the messaging and all the useful endpoints are defined on the second level. And it has also a lot of innovative techniques that we used and came up with over time. Thanks.
Alright, thank you very much, both of you, for this and for the presentation. You're out of time, but yeah, thank you for the talk. Thank you. Thank you. And we'll do five minutes until the next talk.
So you can go grab a drink or whatever you need.
Automatic transcript — names and jargon may be misspelled.