Learn Labs
9. The Trouble with Distributed Systems

9.5 Knowledge, Truth, and Lies

A node in the network CANNOT KNOW ANYTHING FOR SURE about other nodes — it can only make guesses based on the messages it receives (or doesn't receive). A node can find out another node's state only by exchanging messages with it. If a remote node doesn't respond, THERE IS NO WAY OF KNOWING ITS STATE, because PROBLEMS IN THE NETWORK CANNOT RELIABLY BE DISTINGUISHED FROM PROBLEMS AT A NODE.

"Discussions of these systems border on the philosophical: What do we know to be true or false in our system? How sure can we be of that knowledge, if the mechanisms for perception and measurement are unreliable?"

Fortunately, we don't need to figure out the meaning of life. We can STATE THE ASSUMPTIONS we are making (the SYSTEM MODEL) and design the system to meet those assumptions. Algorithms can be PROVED to function correctly within a certain system model — meaning RELIABLE BEHAVIOR IS ACHIEVABLE EVEN IF THE UNDERLYING MODEL PROVIDES VERY FEW GUARANTEES.

5.1 The majority rules

Three parables:

  1. The asymmetric fault. A node Receives all messages but its Outgoing messages are dropped. It is working perfectly and receiving requests, but nobody hears its responses. After a timeout it is declared dead.

    “The semi-disconnected node is dragged to the graveyard, kicking and screaming ‘I’m not dead!’ — but since nobody can hear its screaming, The funeral procession continues with stoic determination.”

  2. The slightly less nightmarish version. The node Notices its messages aren’t being acknowledged and realizes there’s a network fault. “Nevertheless, the node is Wrongly declared dead, and it’s unable to do anything about it.”
  3. The pause. A node pauses for one minute; others declare it dead. The pause finishes and its threads continue As if nothing had happened. “The supposedly dead node suddenly Raises its head out of the coffin, in full health, and starts cheerfully chatting with bystanders. At first, the paused node doesn’t even realize that an entire minute has passed and that it was declared dead.”

The moral: A NODE CANNOT NECESSARILY TRUST ITS OWN JUDGMENT OF A SITUATION.

A distributed system cannot exclusively rely on a single node, because a node may fail at any time, potentially leaving the system stuck. Instead, many algorithms rely on a QUORUM: decisions require a minimum number of votes from several nodes.

That includes decisions about declaring nodes dead. IF A QUORUM OF NODES DECLARES ANOTHER NODE DEAD, THEN IT MUST BE CONSIDERED DEAD, EVEN IF THAT NODE STILL VERY MUCH FEELS ALIVE. The individual node MUST ABIDE BY THE QUORUM DECISION AND STEP DOWN.

Why a majority: allows the system to continue if a MINORITY are faulty (3 nodes → tolerate 1; 5 nodes → tolerate 2), and it is SAFE because THERE CAN BE ONLY ONE MAJORITY — there cannot be two majorities with conflicting decisions at the same time.

5.2 Distributed locks and leases

Locks and leases in distributed applications are PRONE TO MISUSE and are A COMMON SOURCE OF BUGS.

Where you need "only one of some thing":

UseConsequence of two holders
Only one node is leader for a shardSplit brain — lost or corrupted data. SERIOUS
Only one client updates a resourceCorruption by concurrent writes. SERIOUS
Only one node processes an input fileOnly some wasted computational resources. NOT A BIG DEAL

Two ways it goes wrong — and note the second has NO pause at all:

The pause — HBase actually had this bug
client 1lock serviceclient 2fileacquire leasegrantedthen a long pauselease expired → grantedstarts writingwakes up, still believes the lease is validwrites clash · file corrupted
The delayed request — no pause, just a slow network
client 1lock serviceclient 2storagewrite request — delayed a minute+client 1 crasheslease times out → grantedissues its own writethe ancient write finally arrivessame corruption
Figure 9.5.1Two ways it goes wrong — and note the second has NO pause at all
Fencing off zombies

A ZOMBIE is a former leaseholder that has not yet found out that it lost the lease and is still acting as if it were the current leaseholder. SINCE WE CANNOT RULE OUT ZOMBIES ENTIRELY, we have to instead ensure THEY CAN'T DO ANY DAMAGE. This is called FENCING OFF the zombie.

The bad approach — STONITH ("shoot the other node in the head"): disconnect it from the network, shut down the VM, or physically power down the machine.

It is NOT PARTICULARLY EFFECTIVE: it does not protect against the LARGE NETWORK DELAYS of case ②; ALL THE NODES COULD SHUT ONE ANOTHER DOWN; and by the time a zombie has been detected and shut down, IT MAY BE TOO LATE AND DATA MAY ALREADY BE CORRUPTED.

The good approach — FENCING TOKENS:

Every time the lock service grants a lease it also returns a fencing token — a number that increases every time a lock is granted. Every write request to the storage service must include the client's current token.

client 1lock serviceclient 2storageacquire leasetoken = 33long pause; lease expirestoken = 34write(token = 34)accepted; remembers 34wakes, write(token = 33)rejected

Storage refuses the late write because it has “already processed a write with a higher token number”. Which is why “a client that has just acquired the lease must immediately make a write to the storage service, and once that write has completed, any zombies are fenced off.”

Figure 9.5.2The good approach — FENCING TOKENS

Fencing tokens under other names:

SystemName
Chubby (Google's lock service)sequencers
Kafkaepoch numbers
Paxosballot number
Raftterm number
ZooKeeperthe transaction ID zxid or node version cversion
etcdthe revision number along with the lease ID
HazelcastFencedLock API generates one explicitly

Fencing is similar to OPTIMISTIC CONCURRENCY CONTROL (Ch 8) — except that FENCING IS PERMANENT, while concurrency control failures can be retried.

What the storage service must support: either a check for an outdated token, or simply a write that succeeds only if the object has not been written by another client since the current client last read it — an atomic CAS. Object stores support this: S3 calls it CONDITIONAL WRITES, Azure Blob Storage CONDITIONAL HEADERS, Google Cloud Storage REQUEST PRECONDITIONS.

The subtle point about needing a lock service at all:

If your clients write to only ONE storage service that supports conditional writes, THE LOCK SERVICE IS SOMEWHAT REDUNDANT — the lease assignment could have been implemented directly on that storage service. However, ONCE YOU HAVE A FENCING TOKEN, YOU CAN USE IT WITH MULTIPLE SERVICES OR REPLICAS and ensure the old leaseholder is fenced off ON ALL OF THEM.

Fencing a leaderless replicated store:

Put the fencing token in the most significant digits of the last-write-wins timestamp:

client 1 (token 33) writes timestamps  33xxxxxxxx
client 2 (token 34) writes timestamps  34xxxxxxxx   ← always greater

“Any timestamp generated by the new leaseholder will be greater than any timestamp from the old leaseholder, even if the old leaseholder's writes happened later.”

So suppose client 2 writes to a quorum but cannot reach replica 3, and zombie client 1's write does succeed at replica 3. That is fine: a subsequent quorum read prefers client 2's greater timestamp, and read repair or anti-entropy will eventually overwrite client 1's value.

5.3 Byzantine faults

Fencing tokens can detect and block a node acting in error INADVERTENTLY. However, if the node DELIBERATELY wanted to subvert the system's guarantees, it could easily do so by SENDING MESSAGES WITH A FAKE FENCING TOKEN.

In this book we assume that nodes are UNRELIABLE BUT HONEST. They may be slow or never respond, and their state may be outdated — but we assume that IF A NODE DOES RESPOND, IT IS TELLING THE "TRUTH."

A BYZANTINE FAULT is a node "lying" — sending arbitrary faulty or corrupted responses, e.g. casting MULTIPLE CONTRADICTORY VOTES IN THE SAME ELECTION.

(Etymology: it generalizes the two generals problem. In the Byzantine version, n generals need to agree, hampered by TRAITORS in their midst; it is not known in advance who the traitors are. The name comes from "byzantine" in the sense of excessively complicated, bureaucratic, devious — used in politics long before computers. Lamport wanted a nationality that would not offend readers, and was advised that calling it The Albanian Generals Problem was not such a good idea.)

Where Byzantine fault tolerance IS relevant:

  • Aerospace: data in memory or a CPU register could be corrupted by RADIATION, leading a node to respond in arbitrarily unpredictable ways. Since failure is very expensive (an aircraft crashing, a rocket colliding with the ISS), flight control systems must tolerate Byzantine faults
  • Multiple mutually untrusting parties: cryptocurrencies and blockchains are a way of getting mutually untrusting parties to agree on whether a transaction happened, WITHOUT RELYING ON A CENTRAL AUTHORITY

Where it is NOT, and why:

In a datacenter, all nodes are controlled by your organization (so they can hopefully be trusted), and RADIATION LEVELS ARE LOW ENOUGH that memory corruption is not a major problem. (Although datacenters in orbit are being considered.) Multitenant systems have mutually untrusting tenants, but they are ISOLATED VIA FIREWALLS, VIRTUALIZATION, AND ACCESS CONTROL POLICIES, not Byzantine fault tolerance. BFT protocols are QUITE EXPENSIVE — in most server-side data systems, THE COST MAKES THEM IMPRACTICABLE.

Two things BFT explicitly CANNOT save you from:

  1. A software bug. "A bug could be regarded as a Byzantine fault, but IF YOU DEPLOY THE SAME SOFTWARE TO ALL NODES, THEN A BYZANTINE FAULT-TOLERANT ALGORITHM CANNOT SAVE YOU." Most BFT algorithms need a supermajority of more than two-thirds functioning correctly — to use this against bugs you would need FOUR INDEPENDENT IMPLEMENTATIONS of the same software and hope a given bug appears in only one.
  2. Security compromise. "In most systems, IF AN ATTACKER CAN COMPROMISE ONE NODE, THEY CAN PROBABLY COMPROMISE ALL OF THEM, because the nodes are probably running the same software. Thus, TRADITIONAL MECHANISMS — authentication, access control, encryption, firewalls — CONTINUE TO BE THE MAIN PROTECTION."

(Web applications do need to expect arbitrary and malicious behavior from clients under end-user control — which is why input validation, sanitization, and output escaping matter. But we don't use BFT protocols; we simply make the server THE AUTHORITY on what client behavior is allowed. In peer-to-peer networks with no central authority, BFT is more relevant.)

Weak forms of lying — cheap, pragmatic, worth doing:

GuardAgainst
Application-level checksums (or TLS)Corrupted packets that EVADE TCP/UDP checksums — this does happen
Input sanitization: escaping to prevent SQL injection, range checks, limiting string size to prevent denial of service through large memory allocationsMalicious or malformed input. An internal service behind a firewall may get away with less, but basic checks in protocol parsers are still a good idea
Multiple NTP servers, estimating errors and checking that a majority agree on a time rangeA misconfigured NTP server reporting incorrect time is detected as an OUTLIER and excluded

On this page