Learn Labs
9. The Trouble with Distributed Systems

9.6 System Model and Reality

A SYSTEM MODEL is an abstraction describing an algorithm's assumptions.

6.1 Timing models

ModelAssumptionRealism
SynchronousBounded network delay, bounded process pauses, bounded clock error. (Not exactly synchronized clocks or zero delay — just that they NEVER EXCEED A FIXED UPPER BOUND)NOT a realistic model of most practical systems, because unbounded delays and pauses DO occur
Partially synchronous ✔Behaves like a synchronous system MOST of the time, but SOMETIMES exceeds the bounds. When it happens, delay, pauses, and clock error may become arbitrarily largeA REALISTIC MODEL OF MANY SYSTEMS. "Most of the time, networks and processes are quite well behaved — otherwise we would never be able to get anything done — but we have to reckon with the fact that ANY TIMING ASSUMPTIONS MAY BE SHATTERED OCCASIONALLY"
AsynchronousNo timing assumptions at all — it does not even have a clock, so it cannot use timeoutsSome algorithms can be designed for it, but it is VERY RESTRICTIVE

6.2 Node failure models

ModelAssumption
Crash-stop (fail-stop)A node can fail in ONLY ONE WAY: by crashing. It suddenly stops responding and THEREAFTER IS GONE FOREVER — IT NEVER COMES BACK
Crash-recovery ✔Nodes may crash at any moment and perhaps start responding again after an unknown time. Nodes have STABLE STORAGE preserved across crashes, while in-memory state is lost
Degraded performance / partial functionalityNodes SLOW DOWN. They may still respond to health checks while being TOO SLOW TO GET ANY REAL WORK DONE. Called a LIMPING NODE, GRAY FAILURE, or FAIL-SLOW — and it can be EVEN MORE DIFFICULT TO DEAL WITH THAN A CLEANLY FAILED NODE
Byzantine (arbitrary)Nodes may do absolutely anything, including trying to trick and deceive other nodes

Real causes of fail-slow, worth memorizing: "A Gigabit network interface could suddenly drop to 1 Kb/s throughput because of a driver bug; a process under memory pressure may spend most of its time performing garbage collection; worn-out SSDs can have erratic performance; and hardware can be affected by high temperature, loose connectors, mechanical vibration, power supply problems, firmware bugs." Plus: a process stops doing SOME of the things it is supposed to do while other aspects continue working — because a background thread has crashed or deadlocked.

For modeling real systems, the PARTIALLY SYNCHRONOUS model with CRASH-RECOVERY faults is generally the most useful.

6.3 Safety vs liveness

Example properties for a fencing-token generator:

PropertyDefinitionKind
UniquenessNo two requests for a fencing token return the same valueSAFETY
Monotonic sequenceIf request x returned tₓ and y returned t_y, and x completed before y began, then tₓ < t_ySAFETY
AvailabilityA node that requests a token and does not crash EVENTUALLY receives a responseLIVENESS

A giveaway: LIVENESS PROPERTIES OFTEN INCLUDE THE WORD "EVENTUALLY." (And yes — eventual consistency is a liveness property.)

The precise definitions:

Safety — “nothing bad happens”
If violated, We can point to the particular point in time it was broken. After a safety property has been violated, the violation cannot be undone — the damage is already done.
Liveness — “something good eventually happens”
It May not hold at a certain point in time (a node sent a request but hasn’t received a response), But there is always hope it may be satisfied in the future.

Why the distinction matters: for distributed algorithms it is common to require SAFETY PROPERTIES ALWAYS HOLD, IN ALL POSSIBLE SITUATIONS OF A SYSTEM MODEL. Even if all nodes crash, or the entire network fails, THE ALGORITHM MUST NEVER RETURN A WRONG RESULT.

With LIVENESS properties we ARE allowed to make caveats — e.g. a request needs to receive a response only if a MAJORITY OF NODES HAVE NOT CRASHED, and only if THE NETWORK EVENTUALLY RECOVERS. (The partially synchronous model requires exactly this: any period of interruption lasts only a finite duration and is then repaired.)

6.4 Where the model meets reality

Algorithms in the crash-recovery model assume DATA IN STABLE STORAGE SURVIVES CRASHES. But what if the data on disk is CORRUPTED OR WIPED OUT by hardware error or misconfiguration? What if a server has a FIRMWARE BUG and FAILS TO RECOGNIZE ITS HARD DRIVES ON REBOOT even though they're correctly attached?

Quorum algorithms rely on A NODE REMEMBERING THE DATA IT CLAIMS TO HAVE STORED. IF A NODE MAY SUFFER FROM AMNESIA and forget previously stored data, THAT BREAKS THE QUORUM CONDITION AND THUS BREAKS THE CORRECTNESS OF THE ALGORITHM. Perhaps a new model is needed in which stable storage mostly survives — but that model then becomes harder to reason about.

A real implementation may still have to include code to handle the case of something happening that was ASSUMED TO BE IMPOSSIBLE — even if that handling boils down to printf("Sucks to be you"); exit(666); — that is, LETTING A HUMAN OPERATOR CLEAN UP THE MESS. (THIS IS ONE DIFFERENCE BETWEEN COMPUTER SCIENCE AND SOFTWARE ENGINEERING.)

That is not to say abstract system models are worthless — QUITE THE OPPOSITE. They are incredibly helpful for DISTILLING DOWN the complexity of real systems to a manageable set of faults we can reason about.


On this page