Analyzing Linearizability in Relativistic Distributed Systems Kahbod Aeini
Wojciech Golab
[email protected] University of Waterloo Electrical and Computer Engineering Waterloo, Canada
[email protected] University of Waterloo Electrical and Computer Engineering Waterloo, Canada
arXiv:2606.30419v1 [cs.DC] 29 Jun 2026
ABSTRACT Einstein’s theory of relativity correctly predicted that time is relative, and subject to both kinematic and gravitational dilation. Therefore, executions of distributed systems cannot always be modeled as sequences of events totally ordered according to wall clock time. To address this fundamental problem, Gilbert and Golab formulated a generalization of Herlihy and Wing’s linearizability property for shared objects, which they called relativistic linearizability, and introduced a collection of theoretical tools to facilitate rigorous analysis. While they conjectured that several widely-studied classically linearizable algorithms are also relativistically linearizable, their work stopped short of presenting formal proofs of correctness, as pointed out recently by Jayanti. In this paper, we explain how Gilbert and Golab’s techniques can be used to establish relativistic linearizability for a replicated state machine, as well as variations of the widely studied read/write register construction of Attiya, Bar-Noy and Dolev (ABD). Our results establish a stronger form of relativistic linearizability than Jayanti’s central theorem for these asynchronous algorithms.
CCS CONCEPTS • Theory of computation → Distributed computing models; • Software and its engineering → Distributed systems organizing principles; • Applied computing → Physics.
KEYWORDS linearizability; consistency; inter-planetary computing; relativity; simultaneity; time dilation
1
INTRODUCTION
As humanity continues to expand scientific activity beyond Earth, distributed computing will play an increasingly important role in enabling reliable communication, coordination, and autonomy among spacecraft, satellites, robotic explorers, and, ultimately, human outposts. The extreme physical conditions of space fundamentally challenge conventional assumptions underlying terrestrial distributed systems. In particular, the strange predictions of Einstein’s theory of relativity become readily observable due to the vast distances involved, the high relative velocities of communicating entities, and heterogeneous gravitational fields. These factors give rise to relativistic effects such as time dilation and length contraction, which directly impact how time is perceived and measured from the reference frame of different system components.
This work is licensed under a Creative Commons Attribution-NonCommercialNoDerivatives 4.0 International License.
Because relativity predicts that time and simultaneity are relative, reasoning about collections of interdependent events becomes substantially more complex. Consequently, the classical Newtonian model of distributed system executions—where events are totally ordered in real time—no longer yields a canonical representation that can be agreed upon by all possible observers. This is not a farfetched, science-fiction concern. Hafele and Keating’s prediction [6] and famous experiment [5] vividly illustrated this phenomenon as early as the 1970s, without ever leaving the planet: two atomic clocks were flown around the world in opposite directions and compared to a reference clock on the ground, yielding measurable differences in elapsed time due to relativistic effects. If commercial aircraft suffice to produce such discrepancies, it is hardly a stretch to imagine the components of a distributed system finding themselves in comparable or more extreme situations, whether on satellites in distinct orbits, on different planets, or aboard spacecraft moving at high relative velocities. In such settings, the clocks of these components experience time dilation and accumulate discrepancies in local time, challenging the assumption of a shared notion of time that underlies the design and analysis of distributed algorithms. This, in turn, necessitates careful rethinking of long-cherished concepts such as Herlihy and Wing’s linearizability property [7], which states that operations on a shared object appear to take effect in some linear order consistent with their precedence in real time. To address this fundamental problem, Gilbert and Golab [4] formulated relativistic linearizability—a generalization of classic (i.e., Herlihy and Wing’s) linearizability to executions where events are only partially ordered in a physical sense. Furthermore, they introduced a collection of theoretical tools to facilitate rigorous analysis of shared object implementations. While they conjectured that several widely-studied classically linearizable algorithms are also relativistically linearizable, and moreover that this point could be proved using their specific technique, their foundational work stopped short of presenting formal proofs of correctness. In response, Jayanti [8] recently proved that all classically linearizable asynchronous algorithms are in fact relativistically linearizable. Intuitively, asynchronous algorithms are immune to the effects of time dilation as they rely only on physical causality (by way of program order and message exchange) for coordination, which is invariant across different frames of reference. While it may appear on first impression that Jayanti’s theorem confirms Gilbert and Golab’s conjecture, the two works in fact refer to distinct notions of relativistic linearizability. Gilbert and Golab formulate three such notions: R1-linearizability as the property that an execution appears linearizable in some frame of reference, R2linearizability as a special case of R1 where the execution appears linearizable in every frame of reference, and R3-linearizability as a special case of R2 where a single linearization applies in all frames
Aeini and Golab
of reference. Their conjecture is stated specifically with respect to R3 (which implies R2 and R1), whereas Jayanti’s definition of relativistic linearizability is equivalent under general relativity to the (strictly) weaker R2 property. As a result, Jayanti’s theorem neither contradicts nor proves the Gilbert and Golab conjecture. In this paper, we advance the understanding of relativistic distributed systems along several directions. (1) We partly confirm Gilbert and Golab’s conjecture. First, we present a detailed analysis of the Raft replicated state machine [13], which can be used to implement a shared object of arbitrary sequential specification, in Section 4. We then analyze the Attiya, Bar-Noy and Dolev (ABD) shared register construction [1] in Section 5. An informal discussion of quorum-replicated key-value stores follows in Section 6. We establish R3-linearizability in general for Raft, and in special cases for the other algorithms. (2) Our analysis yields the first proofs of R3-linearizability for concrete distributed algorithms, and also the first examples of applying Gilbert and Golab’s proof technique [4]. (3) Whereas Jayanti’s theorem [8] allows the linearization order for an algorithm to be relative (i.e., inherently dependent on the observer’s frame of reference), our results for common asynchronous algorithms demonstrate that all observers can in fact agree on a single global linearization order—a result of both practical and theoretical importance.
2
MODEL
Our model is based very closely on [4]. A set of 𝑁 processors, denoted P = {𝑝 1, 𝑝 2, ..., 𝑝 𝑁 }, is subject to relativistic effects due to motion and gravitational fields [3]. Since it is impossible for all observers (i.e., processors in their own frames of references) to agree on a total order of events, an execution is modeled as a history 𝐻 = (𝐸, <𝐸 ) where 𝐸 is a finite set of events and <𝐸 is an irreflexive partial order over 𝐸. Each event has a position in Einstein’s fourdimensional spacetime, and represents a primitive action (sending/receiving a message, or local computation). The partial order <𝐸 captures physical causality, which is invariant across frames of reference. Physical causality is a refinement of Lamport’s “happens before” relation [10, 11], which captures only program order and message passing—the sole means of coordination in purely asynchronous algorithms. A classic history 𝐻 = (𝐸, <𝐸 ) is one where <𝐸 is a total order, meaning that events are observed from the frame of reference of a specific processor. Different processors may perceive the same execution through distinct classic histories with conflicting total orders over the events, but all such total orders are refinements of the underlying partial order of physical causality. The processors execute algorithms that implement shared objects. Following Herlihy and Wing [7], processors interact with the shared object implementations by invoking operations and receiving their responses. An invocation event of operation 𝑜𝑝 by processor 𝑝 on object 𝑋 is denoted (In, 𝑝, 𝑋, 𝑜𝑝). A response event with return value 𝑟𝑒𝑡 is denoted (Re, 𝑝, 𝑋, 𝑟𝑒𝑡). A response is matching with respect to an invocation if both refer to the same processor and object. An operation execution (or op-ex) comprises an invocation and its matching response, if it exists. For any history 𝐻 = (𝐸, <𝐸 ), an op-ex 𝑜𝑥 is pending if it lacks a matching response,
and is complete otherwise. We denote by 𝑖𝑛𝑣 (𝑜𝑥) and 𝑟𝑒𝑠 (𝑜𝑥) the invocation and response events of 𝑜𝑥, respectively. If 𝑜𝑥 is complete, then 𝑖𝑛𝑣 (𝑜𝑥) <𝐸 𝑟𝑒𝑠 (𝑜𝑥). We denote by compl(𝐻 ) the subsequence of 𝐻 comprising the events of complete op-ex’s in 𝐻 . Given a history 𝐻 = (𝐸, <𝐸 ) and processor 𝑝, we denote by 𝐻 |𝑝 the subhistory comprising events applied at processor 𝑝 and the corresponding subset of <𝐸 . Two histories 𝐻 and 𝐻 ′ are equivalent if for every processor 𝑝, 𝐻 |𝑝 = 𝐻 ′ |𝑝. A history 𝐻 is well-formed if for every processor 𝑝, if 𝐻 |𝑝 is non-empty, then the events in 𝐻 |𝑝 are totally ordered and form an alternating sequence of invocation and response events, starting with an invocation. Lamport [11] defines two temporal relations over pairs of operation executions, 𝑜𝑥 1 and 𝑜𝑥 2 in a history 𝐻 = (𝐸, <𝐸 ): 𝑜𝑥 1 →𝐻 𝑜𝑥 2 denotes that 𝑜𝑥 1 is complete and its response causally precedes the invocation of 𝑜𝑥 2 ; and 𝑜𝑥 1 d𝐻 𝑜𝑥 2 denotes that the invocation of 𝑜𝑥 1 causally precedes some event of 𝑜𝑥 2 . We say that 𝐻 is sequential if every pair of distinct op-ex’s is related by →𝐻 . The correct behavior of an object under sequential access is defined by its type 𝜏 = (S, 𝑠𝑖𝑛𝑖𝑡 , O, R, 𝛿) where S is the set of states, 𝑠𝑖𝑛𝑖𝑡 ∈ S is the initial state, O is a set of operation types, R is the set of responses, and 𝛿 : S × O → S × R is a (one-to-many) state transition mapping. Specifically, if a processor applies an operation of type 𝑜𝑡 to an object of type 𝜏 that is in state 𝑠, then the object may return a response 𝑟 and change its state to 𝑠 ′ if and only if (𝑠 ′, 𝑟 ) ∈ 𝛿 (𝑠, 𝑜𝑡). For a sequential history 𝐻 , we say that an object 𝑋 conforms to its type 𝜏 = (S, 𝑠𝑖𝑛𝑖𝑡 , O, R, 𝛿) in 𝐻 if 𝐻 |𝑋 is consistent with some sequence of transitions of 𝛿 starting from state 𝑠𝑖𝑛𝑖𝑡 , in the following sense: letting 𝑜𝑡𝑖 and 𝑟𝑒𝑡𝑖 denote the operation type and response of the 𝑖’th op-ex on 𝑋 in 𝐻 |𝑋 out of 𝑘, there exists a sequence ⟨𝑠 0, 𝑠 1, 𝑠 2, ..., 𝑠𝑘 ⟩ of states in S such that 𝑠 0 = 𝑠𝑖𝑛𝑖𝑡 , and (𝑠𝑖 , 𝑟𝑒𝑡𝑖 ) ∈ 𝛿 (𝑠𝑖 −1, 𝑜𝑡𝑖 ) holds for all 𝑖 ≤ 𝑘. We say that 𝐻 is legal if for every object 𝑋 accessed in 𝐻 , 𝑋 conforms to its type in 𝐻 .
3
BACKGROUND
This section reviews the formal definitions of classic and relativistic linearizability. We start with classic linearizability as conceived by Herlihy and Wing [7] and explained in [4]: Definition 3.1. A classic history 𝐻 = (𝐸, <𝐸 ) is linearizable if there exists a classic history 𝐻 ′ = (𝐸 ′, <𝐸 ′ ) (called the completion of 𝐻 ) such that: L1 𝐻 ⊆ 𝐻 ′ (i.e., 𝐸 ⊆ 𝐸 ′ and <𝐸 ⊆<𝐸 ′ ); L2 𝐸 ′ is obtained from 𝐸 by adding a set 𝑀 of matching responses for a subset of pending operation executions in 𝐻 , and each event of 𝐸 precedes every event of 𝑀 in <𝐸 ′ ; L3 compl(𝐻 ′ ) is equivalent to some legal sequential history 𝑆; and L4 →𝑆 extends →𝐻 , meaning that for any pair of operation executions 𝑜𝑥, 𝑜𝑥 ′ in 𝑆, if 𝑜𝑥 →𝑆 𝑜𝑥 ′ then 𝑜𝑥 ′′ →𝐻 𝑜𝑥 is false where 𝑜𝑥 ′′ denotes the (possibly pending) counterpart of 𝑜𝑥 ′ in 𝐻 . The sequential history 𝑆 referred to by property L3 is called a linearization (in this paper a classic linearization) of 𝐻 . Next, we restate Gilbert and Golab’s [4] three variations of relativistic linearizability. Given an execution history 𝐻 = (𝐸, <𝐸 ) with partially ordered events, the definitions realize observations of 𝐻
Analyzing Linearizability in Relativistic Distributed Systems
from a specific frame of reference using a total order <𝑇 over the events in 𝐸 that refines <𝐸 (i.e., <𝐸 ⊆<𝑇 ). Definition 3.2 (R1-linearizability). For any execution history 𝐻 = (𝐸, <𝐸 ), we say that 𝐻 is R1-linearizable if and only if there exists a total order <𝑇 over 𝐸 that refines <𝐸 , such that the classic history 𝐻𝑇 = (𝐸, <𝑇 ) is linearizable. Any linearization of 𝐻𝑇 is called an R1-linearization of 𝐻 . Definition 3.3 (R2-linearizability). For any execution history 𝐻 = (𝐸, <𝐸 ), we say that 𝐻 is R2-linearizable if and only if for every total order <𝑇 over 𝐸 that refines <𝐸 , the classic history 𝐻𝑇 = (𝐸, <𝑇 ) is linearizable. Any linearization 𝑆 of any such 𝐻𝑇 is called an R2-linearization of 𝐻 . Definition 3.4 (R3-linearizability). For any execution history 𝐻 = (𝐸, <𝐸 ), we say that 𝐻 is R3-linearizable if and only if there exists an R1-linearization 𝑆 of 𝐻 such that for every total order <𝑇 over 𝐸 that refines <𝐸 , 𝑆 is a linearization of the classic history 𝐻𝑇 = (𝐸, <𝑇 ). Any such history 𝑆 is called an R3-linearization of 𝐻 . Intuitively, R1 states that a history 𝐻 appears linearizable in some frame of reference, R2 states that 𝐻 appears linearizable in every frame of reference, and R3 states that 𝐻 is not only linearizable in every frame of reference but all observers furthermore agree on a common linearization. R3 implies R2, which implies R1. R2 is the equivalent of Herlihy and Wing’s classic linearizability property in the relativistic model, and naturally inherits the locality property: a history 𝐻 involving multiple implemented objects is linearizable if and only if the projection of 𝐻 onto each implemented object 𝑂 (denoted 𝐻 |𝑂) is individually linearizable. As shown in [4], R3 is not local. The proof of this result exhibits a multi-object execution history that is R2-linearizable but not R3-linearizable, meaning that every observer perceives classical linearizability and yet there exist observers who fundamentally disagree on the linearization order. Another example of this phenomenon is presented in [8]. Ideally, a history 𝐻 would satisfy R2 linearizability and each projection 𝐻 |𝑂 onto an object 𝑂 accessed in 𝐻 would satisfy the stronger R3 property. Gilbert and Golab’s proof techniques for R2 and R3 linearizability are based on the insight that coordination mechanisms used by shared object implementations induce certain structural properties on execution histories. These properties are formalized in terms of connectedness, which characterizes causal relationships: Definition 3.5. For any history 𝐻 = (𝐸, <𝐸 ), and any distinct operation executions 𝑜𝑥, 𝑜𝑥 ′ in 𝐻 , we say that 𝑜𝑥 and 𝑜𝑥 ′ are: • strongly connected in 𝐻 if 𝑜𝑥 d𝐻 𝑜𝑥 ′ and 𝑜𝑥 ′ d𝐻 𝑜𝑥 • weakly connected in 𝐻 if either 𝑜𝑥 d𝐻 𝑜𝑥 ′ or 𝑜𝑥 ′ d𝐻 𝑜𝑥 (but not both) • connected in 𝐻 if they are strongly or weakly connected in 𝐻 • disconnected in 𝐻 if they are not connected in 𝐻 The history 𝐻 is called connected if every pair of operation executions in 𝐻 is connected. Using this notion, Gilbert and Golab show that some R1-linearizable implementations automatically satisfy R2:
Definition 3.6. For any history 𝐻 = (𝐸, <𝐸 ), an R1-linearization 𝑆 of 𝐻 is called R2-conducive if for every pair of operation executions 𝑜𝑥, 𝑜𝑥 ′ in 𝐻 that have counterparts in 𝑆, the following hold: (1) if 𝑜𝑥 and 𝑜𝑥 ′ are weakly connected and 𝑜𝑥 d𝐻 𝑜𝑥 ′ , then 𝑜𝑥 →𝑆 𝑜𝑥 ′ ; and (2) if 𝑜𝑥 and 𝑜𝑥 ′ are disconnected then both are executions of read-only operation types (i.e., ones that always cause trivial state transitions), and the same holds for all operation executions that appear between 𝑜𝑥 and 𝑜𝑥 ′ in 𝑆. Theorem 3.7. Let 𝐻 = (𝐸, <𝐸 ) be any history and suppose that 𝑆 is an R2-conducive R1-linearization of 𝐻 . Then 𝐻 is R2-linearizable. Furthermore, if 𝐻 is connected then 𝑆 itself is an R2-linearization of 𝐻 (in any frame of reference), hence 𝐻 is R3-linearizable.1 Gilbert and Golab conclude their paper by conjecturing that many classic distributed shared object implementations, particularly those based on majority quorums (e.g., ABD [1], replicated state machines and key-value stores), generate connected execution histories with R2-conducive R1-linearizations, and hence satisfy both R2 and R3 linearizability via Theorem 3.7.
4
RAFT REPLICATED STATE MACHINE
Raft is a consensus algorithm designed to manage a replicated log across a cluster of servers (i.e., processors in our model), ensuring that each server agrees on the same sequence of state machine commands even in the presence of failures [13]. Unlike Paxos, which is notoriously difficult to reason about, Raft was explicitly designed for understandability by decomposing the consensus problem into three relatively independent subproblems: leader election, log replication, and safety. At any given time, each server in a Raft cluster is in one of three states—leader, follower, or candidate—and time is divided into terms of arbitrary length, each identified by a monotonically increasing integer. A term begins with an election, and if a leader is successfully elected, it serves for the remainder of that term, coordinating all client interactions and log replication. The algorithm proceeds roughly as follows. In normal operation, the leader receives client requests, appends them as entries to its local log, and replicates those entries to the other servers in parallel via remote procedure calls (RPCs). Once a majority of servers have acknowledged a given entry, the leader considers it committed and applies it to its state machine. If the leader fails or becomes unreachable, followers that stop receiving heartbeats will time out, transition to the candidate state, and initiate a new election by incrementing the term and requesting votes from their peers. A candidate must receive votes from a majority of the cluster to become the new leader, and Raft’s voting rules guarantee that any elected leader already contains all previously committed entries. This majority-based mechanism—used in both elections and commit decisions—is what ensures that committed entries are never lost, even when servers crash and recover. Figure 1 depicts the structure of Raft’s log. Each cell represents a log entry (i.e., a command) labeled with the term in which it was created. Row 𝑖 represents the local log of server 𝑆𝑖 . The leader’s local log (𝑆 1 in Figure 1) has the largest number of entries in its leadership term, as it is the first node (i.e., server) that receives the 1 This result succinctly combines Theorems 5 and 6 in [4].
Aeini and Golab
time
log index crash 1
2
3
4
5
6
7
t1
t1
t1
t2
t3
t3
t3
×
𝑆1
×
(leader)
𝑆1 (leader)
heartbeat
election timeout
𝑆 2 elected leader
𝑆2
t1
t1
t1
t2
t3
t3
𝑆2 term++ ReqestVote
𝑆3
t1
t1
t1
t2
majority (𝑆 2 + 𝑆 3 ) VoteGranted
t3 𝑆3
𝑆4
t1
t1
t1
t2 RPC request
𝑆5
t1
× lost/no reply
t1 committed (index 5) term 1
term 2
term 3
Figure 1: Replicated log structure across a five-server Raft cluster during term 3. Each cell represents a log entry labeled with the term in which it was created. The dashed line marks the commit index: entries through index 5 have been replicated to a majority and are considered committed.
clients’ commands. Whenever an entry is logged into the leader’s local log, the leader replicates it to other servers concurrently, and once the majority of the servers (in this example, at least 2 follower servers) acknowledge the entry, it will be considered as committed. In Figure 1, 5 commands have been committed, which means in the case of failure, any server that would become the next leader will have the first five entries in the log. The Raft algorithm carries an inherent causal relationship among its operations. This can be observed from the algorithm’s behavior, where the state of the cluster is determined by the operations of nodes. E.g., a leader becomes unavailable, which causes some other node to identify itself as a candidate, which results in a new term, and possibly a new leader. Moreover, a majority of servers must acknowledge an entry before it is committed and applied to the state machine. This kind of behavior of the algorithm leads us to believe that a unique linearization order exists over the operations, as determined by the unique positions of the corresponding commands in the replicated log. Therefore, we expect Raft to be R3-linearizable. However, this observation does not guarantee that R3-linearizability of Raft can be proved by the technique introduced by Gilbert and Golab, and Jayanti’s theorem proves Raft R2-linearizability but does not settle the question for R3-linearizability. Figure 2 represents an execution in which the leader (server 𝑆 1 ) crashes and fails to send heartbeat RPCs. When the election timeout is reached, 𝑆 2 increments the term and requests votes from other servers, in order to become the new leader. The vote is granted to 𝑆 2 by 𝑆 3 and 𝑆 2 becomes the new leader.2 Raft’s voting scheme ensures that the elected server has a log that is at least as up-to-date as that of every server in some majority quorum. This is guaranteed because when a server (in this example 𝑆 2 ) requests votes from 2𝑆
RPC response
2 votes for itself, hence only one additional vote is required to form a majority in the three-server case.
Figure 2: Leader failure and election in Raft. After 𝑆 1 crashes, 𝑆 2 ’s election timer expires and it transitions to the candidate state, incrementing its term. 𝑆 2 sends ReqestVote RPCs to the other servers; 𝑆 1 does not respond (dashed arrow), but 𝑆 3 grants its vote. With a majority of votes (𝑆 2 + 𝑆 3 ), 𝑆 2 becomes the new leader for the current term.
other servers by sending RPCs to them, it also attaches the index and term of its last log entry along with the new term number to the RPC, and a voter endorses the candidate only if the candidate’s log is at least as up-to-date as its own. Once 𝑆 2 ’s election is confirmed, the clients send their commands to 𝑆 2 , and it coordinates their replication. One can easily see the many causal relationships among the events in Figure 2. For example, 𝑆 2 would not have increased the term and started a leader election if 𝑆 1 had not stopped sending heartbeats, at least assuming a timeout value appropriately tuned to the network and processing delays. Note that the heartbeat timeout is measured locally by each server, and therefore is immune to the time dilation and length contraction effects. The leader transition exemplifies even more causal relationships in Raft, independently of the chosen heartbeat timeout. In the scenario of Figure 2, for 𝑆 2 to become the leader after 𝑆 1 , the voting scheme requires 𝑆 2 ’s log to already contain any command corresponding to 𝑆 1 ’s term and earlier terms if that command was already committed or will be committed in the future. If not, then the majority of the servers that voted for 𝑆 2 lack the command and will no longer accept it after casting their votes for 𝑆 2 because the Raft protocol rejects RPCs from obsolete terms. Raft’s log matching property then ensures that any such command is replicated to a majority of servers before 𝑆 2 commits any new command in its own term. This enforces causality between the operations corresponding to commands committed in earlier terms and commands committed in 𝑆 2 ’s term. Thus, the causal relationships yield a unique linearization order, recorded in some prefix of the leader’s local log, and independent of the frame of reference. Moreover, this committed prefix (equivalently, the state machine state it produces) serves as a canonical record of the execution that every observer agrees upon, regardless of the frame of reference from which it is observed. Hence, we expect the Raft algorithm to be R3-linearizable. In order to prove R3-linearizability formally using the GilbertGolab technique, we first show that all pairs of operations are
Analyzing Linearizability in Relativistic Distributed Systems
connected, then exhibit an R2-conducive R1-linearization, and finally apply Theorem 3.7. Consider any history 𝐻 = (𝐸, <𝐸 ) of Raft. The completion 𝐻 ′ = (𝐸 ′, <𝐸 ′ ) of 𝐻 is constructed by completing the pending op-exes with committed commands, and removing all other pending ones from 𝐻 . For completed op-exes, the matching response event carries a return value determined by simulating the execution of the corresponding command based on the sequence of commands in the replicated log. A unique coordinate in spacetime is assigned to each matching response event 𝑒 such that every event of 𝐻 causally precedes 𝑒. Next, we construct a canonical R1linearization of 𝐻 from 𝐻 ′ . Fix a total order <𝑇 that refines <𝐸 ′ and corresponds to the reference frame of the final Raft leader (i.e., the leader for the maximum term that occurs in 𝐻 ).3 Let 𝐻𝑇 = (𝐸 ′, <𝑇 ), and let 𝑆 denote a sequential ordering of operations in 𝐻 ′ based on the ultimate order of the corresponding commands in the final leader’s local log. In other words, if the op-exes 𝑜𝑥 and 𝑜𝑥 ′ both appear in 𝐻 ′ and the command for 𝑜𝑥 precedes the command for 𝑜𝑥 ′ in the leader’s log at the end of 𝐻 , then 𝑜𝑥 →𝑆 𝑜𝑥 ′ . Since commands for op-exes in 𝐻 ′ are all committed at a majority of servers in Raft in 𝐻 , it follows that →𝑆 is consistent with →𝐻𝑇 . Furthermore, the history 𝑆 is legal since it reflects the final leader’s execution of commands in the correct order. Thus, 𝑆 is an R1-linearization of 𝐻 . Lemma 4.1. Let 𝐻 = (𝐸, <𝐸 ) be any history of Raft and let 𝐻 ′ be its completion. Then, any two operation executions 𝑜𝑥 and 𝑜𝑥 ′ in 𝐻 that have counterparts in 𝐻 ′ are connected. Proof. Let 𝑜𝑥 and 𝑜𝑥 ′ be operation executions in 𝐻 representing state machine commands 𝑐𝑚𝑑 and 𝑐𝑚𝑑 ′ , respectively. Since we assume both op-exes have counterparts in 𝐻 ′ , it follows that both have corresponding commands committed to a majority of servers. The two commands are committed either in the same leadership term or during distinct ones. Suppose that the two commands are committed in the same term, and without the loss of generality, 𝑐𝑚𝑑 precedes 𝑐𝑚𝑑 ′ in the replicated log. Therefore, the event in which the leader appends 𝑐𝑚𝑑 causally precedes the one appending 𝑐𝑚𝑑 ′ due to the leader’s program order. This implies that 𝑜𝑥 d𝐻 ′ 𝑜𝑥 ′ , as required. The other case to consider is when the two commands are executed during distinct terms. Without the loss of generality assume 𝑐𝑚𝑑 and 𝑐𝑚𝑑 ′ are in terms 𝑡 and 𝑡 ′ , respectively, and 𝑡 < 𝑡 ′ . The leader election at the beginning of term 𝑡 ′ uses majority voting and establishes a causal relationship between operation executions in different terms, as explained earlier in this section. Specifically, the local log of the term 𝑡 ′ leader records all committed commands for term 𝑡 before any committed command for term 𝑡 ′ . Furthermore, Raft’s log matching property ensures that all the term 𝑡 commands are committed before any term 𝑡 ′ commands. This implies that 𝑜𝑥 d𝐻 ′ 𝑜𝑥 ′ , as required. □ Lemma 4.2. Let 𝐻 = (𝐸, <𝐸 ) be any history of Raft, let 𝐻 ′ be its completion, and let 𝑆 be the corresponding canonical R1-linearization of 𝐻 . Then, 𝑆 is R2-conducive. Proof. Consider two op-exes 𝑜𝑥 and 𝑜𝑥 ′ in 𝐻 that have counterparts in 𝐻 ′ . If the command 𝑐𝑚𝑑 of an op-ex 𝑜𝑥 appears before the command 𝑐𝑚𝑑 ′ of an op-ex 𝑜𝑥 ′ in the replicated log, then some
event of 𝑜𝑥 precedes some event of 𝑜𝑥 ′ and 𝑜𝑥 d𝐻 𝑜𝑥 ′ holds, as explained in the proof of Lemma 4.1. Thus, the two op-exes are either weakly or strongly connected. In both cases, 𝑆 linearizes the op-exes based on the order of appearance of the corresponding commands in the replicated log. If 𝑜𝑥 and 𝑜𝑥 ′ are only weakly connected then 𝑜𝑥 d𝐻 𝑜𝑥 ′ implies that 𝑐𝑚𝑑 precedes 𝑐𝑚𝑑 ′ in the log, hence 𝑜𝑥 →𝑆 𝑜𝑥 ′ . Therefore, the first clause of Theorem 3.6 holds for 𝑆. From Lemma 4.1 we know that any two op-exes in 𝐻 ′ are connected, and therefore, the second clause of Theorem 3.6 holds trivially. Thus, the lemma is proved. □ Theorem 4.1. Let 𝐻 = (𝐸, <𝐸 ) be any history of Raft, let 𝐻 ′ be its completion, and let 𝑆 be the corresponding canonical R1-linearization of 𝐻 . Then, 𝐻 is R3-linearizable, and 𝑆 is its R3-linearization. Proof. It follows from Lemma 4.2 that 𝑆 is R2-conducive. Hence, by Theorem 3.7, 𝐻 is R2-linearizable. Moreover, from Lemma 4.1 we know that 𝐻 ′ is connected. Hence, by Theorem 3.7, 𝑆 is an R2-linearization of 𝐻 in any frame of reference. Thus, 𝐻 is R3linearizable and 𝑆 is its R3-linearization. □
5
ABD SIMULATION
The ABD simulation (Figure 3) implements a classically linearizable single-writer multi-reader shared register object. The state of the register is replicated across the 𝑁 processors using majority quorums. The values written to the register are tagged with monotonically increasing version numbers, generated by the writer. A write operation sends the new value-version pair to all processors and awaits a reply from a majority. In case messages are reordered by the network, each processor compares the incoming value-version pair against its local state, and overwrites its previous value only if the incoming one has a larger version. A read operation sends a request to retrieve value-version pairs from all processors and awaits replies from a majority. If all the reply messages have equal versions, then their values agree and the common value is returned. If not, then the reader executes a write-back phase whereby the received value with the highest version is sent to all processors (as in a write). The read returns this same value after a majority of processors acknowledges receipt of the write-back message. The construction can tolerate the permanent crash failure of any minority of processors. Consider any history 𝐻 = (𝐸, <𝐸 ) of ABD. The completion 𝐻 ′ = (𝐸 ′, <𝐸 ′ ) of 𝐻 is constructed as follows. First, discard any pending reads, as well as any pending write whose value is not returned by any read operation. If a pending write in 𝐻 has its value4 returned by a read, then “append” a matching response. A unique coordinate in spacetime is assigned to each matching response event 𝑒 such that every event of 𝐻 causally precedes 𝑒. Next, fix a total order <𝑇 that refines <𝐸 ′ . Let 𝐻𝑇 = (𝐸 ′, <𝑇 ), and let 𝑆 denote a sequential ordering of operations in 𝐻𝑇 according to the following rules. First, arrange writes in increasing order of the corresponding versions (i.e., in program order). Next, let 𝑅𝑘 denote the set of read operations that return the value 𝑣 obtained from a value-version pair ⟨𝑣, 𝑘⟩. Insert each read in 𝑅𝑘 after the write operation that creates version 𝑘 and before the next write. Within the set 𝑅𝑘 , the 4 For analysis, we assume that each read is mapped to a unique write via the version
3 This server is guaranteed to have the most complete copy of the replicated log.
number even if two writes may assign the same value 𝑣 .
Aeini and Golab
Variables: • 𝑙𝑎𝑠𝑡_𝑣𝑒𝑟𝑠𝑖𝑜𝑛: integer, initially zero (writer only) • 𝑣𝑎𝑙𝑖 , 𝑣𝑒𝑟𝑠𝑖𝑜𝑛𝑖 : processor 𝑖’s value (initially equal to the initial value of the simulated read/write register) and its corresponding version (initially zero) Procedure Write(v) for writer processor 𝑙𝑎𝑠𝑡_𝑣𝑒𝑟𝑠𝑖𝑜𝑛 := 𝑙𝑎𝑠𝑡_𝑣𝑒𝑟𝑠𝑖𝑜𝑛 + 1 3 send ⟨𝑣, 𝑙𝑎𝑠𝑡_𝑣𝑒𝑟𝑠𝑖𝑜𝑛⟩ to all processors 4 wait for ACK from majority
1
2
Procedure Read() for processor 𝑖 request 𝑣𝑎𝑙 𝑗 , 𝑣𝑒𝑟𝑠𝑖𝑜𝑛 𝑗 from each processor 𝑗 7 𝑅 := responses from majority of processors ⟨𝑣, 𝑣𝑒𝑟𝑠𝑖𝑜𝑛⟩ := element of 𝑅 with largest version 8 9 if |𝑅| > 1 then 10 send ⟨𝑣, 𝑣𝑒𝑟𝑠𝑖𝑜𝑛⟩ to all processors 11 wait for ACK from majority 5
Case B: 𝑜𝑥 1 and 𝑜𝑥 2 are both writes. Since there is only a single writer, both operations are executed by the same processor. Once again, this establishes causality. Case C: 𝑜𝑥 1 is a read and 𝑜𝑥 2 is a write. If the write 𝑜𝑥 2 is complete in 𝐻 then both operations are complete in 𝐻 . In that case both operations access a majority of processors, which establishes causality as in Case A. On the other hand, if 𝑜𝑥 2 is pending, then by construction of 𝐻 ′ it follows that all events of 𝑜𝑥 1 causally precede the response of 𝑜𝑥 2 . Thus, 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . □ To complete the proof of R3-linearizability, it suffices to show that the R1-linearization 𝑆 is R2-conducive and then apply Theorem 3.7. We demonstrate this first for the single-reader case in Section 5.1, and then for the multi-reader case in Section 5.2.
6
12
return 𝑣
Procedure Receive ⟨𝑣, 𝑣𝑒𝑟𝑠𝑖𝑜𝑛⟩, processor 𝑖 if 𝑣𝑒𝑟𝑠𝑖𝑜𝑛 > 𝑣𝑒𝑟𝑠𝑖𝑜𝑛𝑖 then 15 𝑣𝑎𝑙𝑖 := 𝑣 16 𝑣𝑒𝑟𝑠𝑖𝑜𝑛𝑖 := 𝑣𝑒𝑟𝑠𝑖𝑜𝑛
5.1
Single-Reader Case
For simplicity, we first complete the analysis in the special case where only a single processor can issue read operations.5 Intuitively, we expect the ABD algorithm instantiated in this setting to satisfy R3-linearizability since there can only be one valid linearization for a history 𝐻 in any frame of reference (once the completion 𝐻 ′ is fixed) as the ambiguity regarding the relative order of reads in our construction of 𝑆 is removed.
13
14
17
send back ACK
Procedure Receive request for ⟨𝑣𝑎𝑙𝑖 , 𝑣𝑒𝑟𝑠𝑖𝑜𝑛𝑖 ⟩, processor 𝑖 19 send back ⟨𝑣𝑎𝑙𝑖 , 𝑣𝑒𝑟𝑠𝑖𝑜𝑛𝑖 ⟩
18
Figure 3: Simplified pseudocode for the ABD simulation [1].
read operations can be ordered in an arbitrary manner consistent with →𝐻𝑇 . Classical analysis of the ABD algorithm shows that the constructed linearization 𝑆 is a legal sequential history and a linearization of 𝐻𝑇 , which implies R1-linearizability of 𝐻 . To deduce R2 and R3-linearizability from R1, we establish additional structural properties of the history 𝐻 . Intuitively, the use of majority quorums guarantees that the completion 𝐻 ′ is connected for any history 𝐻 of ABD: Lemma 5.1. Let 𝐻 = (𝐸, <𝐸 ) be a history of ABD, and let 𝐻 ′ be its completion. Then any two operation executions 𝑜𝑥 1 and 𝑜𝑥 2 in 𝐻 that have counterparts in 𝐻 ′ are connected. Proof. We proceed by an exhaustive case analysis on the possible pairs of operations in 𝐻 ′ . Case A: 𝑜𝑥 1 and 𝑜𝑥 2 are both reads. Since only complete reads are included in 𝐻 ′ , both operations access a majority quorum of processors. Since the majority quorums intersect, both operations access some common processor 𝑝, which establishes causality between some event of 𝑜𝑥 1 and some event of 𝑜𝑥 2 .
Lemma 5.2. Let 𝐻 = (𝐸, <𝐸 ) be a history of single-reader ABD, and let 𝐻 ′ be its completion. Fix a total order <𝑇 that refines <𝐸 , and let 𝑆 be the corresponding R1-linearization. Then for every pair of operation executions 𝑜𝑥 1 and 𝑜𝑥 2 in 𝐻 ′ , if 𝑜𝑥 1 and 𝑜𝑥 2 are weakly connected with 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 , then 𝑜𝑥 1 →𝑆 𝑜𝑥 2 . Proof. Let 𝑜𝑥 1 and 𝑜𝑥 2 be two op-exes in 𝐻 ′ that are weakly connected and satisfy 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . To show that 𝑜𝑥 1 →𝑆 𝑜𝑥 2 , we proceed by an exhaustive case analysis on the possible pairs of operations in 𝐻 ′ . Case A: 𝑜𝑥 1 and 𝑜𝑥 2 are both writes. Since there is only a single writer, both operations are executed by the same processor. Then 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 implies that 𝑟𝑒𝑠 (𝑜𝑥 1 ) <𝐸 𝑖𝑛𝑣 (𝑜𝑥 2 ), hence 𝑜𝑥 1 →𝑆 𝑜𝑥 2 by construction of 𝑆 irrespective of <𝑇 . Case B: 𝑜𝑥 1 and 𝑜𝑥 2 are both reads. The analysis is analogous to Case A since there is only one reader. Case C: 𝑜𝑥 1 is a read and 𝑜𝑥 2 is a write. We proceed by subcases on the relative position of 𝑜𝑥 1 and 𝑜𝑥 2 in the linearization order. Subcase (i): 𝑜𝑥 1 returns the value assigned by a write that precedes 𝑜𝑥 2 in 𝑆. Then 𝑜𝑥 1 →𝑆 𝑜𝑥 2 by construction of 𝑆. Subcase (ii): 𝑜𝑥 1 returns the exact value 𝑣 written by 𝑜𝑥 2 . Then 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 1 holds since 𝑜𝑥 1 must retrieve 𝑣 from some processor 𝑝 after 𝑝 receives it in a message generated by 𝑜𝑥 2 . This is a contradiction since the lemma assumes that 𝑜𝑥 1 and 𝑜𝑥 2 are weakly connected with 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . Subcase (iii): 𝑜𝑥 1 returns the value 𝑣 assigned by a write that follows 𝑜𝑥 2 in 𝑆. Let 𝑜𝑥 3 be that write. Then 𝑜𝑥 3 d𝐻 ′ 𝑜𝑥 1 holds since 𝑜𝑥 1 must retrieve 𝑣 from some processor 𝑝 after 𝑝 receives it in a message generated by 𝑜𝑥 3 . Furthermore, 𝑜𝑥 2 is complete and 𝑜𝑥 2 →𝐻 𝑜𝑥 3 . This implies that 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 1 . We obtain a contradiction as in Subcase (ii). ≥ 3 must hold for fault tolerance, the single reader assumption implies that at least one processor participates in the replication protocol but does not invoke operations on the implemented read/write register. 5 Since 𝑁
Analyzing Linearizability in Relativistic Distributed Systems
Case D: 𝑜𝑥 1 is a write and 𝑜𝑥 2 is a read. We proceed by subcases on the relative position of 𝑜𝑥 1 and 𝑜𝑥 2 in the linearization order. Subcase (i): 𝑜𝑥 2 returns the value 𝑣 assigned by a write that precedes 𝑜𝑥 1 in 𝑆. If 𝑜𝑥 1 is pending in 𝐻 then 𝑖𝑛𝑣 (𝑜𝑥 2 ) <𝐸 ′ 𝑟𝑒𝑠 (𝑜𝑥 1 ) by construction of 𝐻 ′ . Thus, 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 1 holds, which is a contradiction since the lemma assumes that 𝑜𝑥 1 and 𝑜𝑥 2 are weakly connected with 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . On the other hand, if 𝑜𝑥 1 is complete in 𝐻 , then there is a processor 𝑝 at the intersection of the quorums used by 𝑜𝑥 1 and 𝑜𝑥 2 such that 𝑜𝑥 2 reads 𝑝’s value before 𝑜𝑥 1 overwrites it. Then 𝑜𝑥 2 d𝐻 𝑜𝑥 1 holds, hence 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 1 , which once again leads to a contradiction. Subcase (ii): 𝑜𝑥 2 returns the value written by 𝑜𝑥 1 . Then 𝑜𝑥 1 →𝑆 𝑜𝑥 2 by construction of 𝑆. Subcase (iii): 𝑜𝑥 2 returns the value assigned by a write that follows 𝑜𝑥 1 in 𝑆. Let 𝑜𝑥 3 be that write. Then 𝑜𝑥 1 →𝑆 𝑜𝑥 3 and 𝑜𝑥 3 →𝑆 𝑜𝑥 2 by construction. This implies 𝑜𝑥 1 →𝑆 𝑜𝑥 2 . □ Theorem 5.1. Let 𝐻 = (𝐸, <𝐸 ) be a history of single-reader ABD, and let 𝐻 ′ be its completion. Fix any total order <𝑇 that refines <𝐸 , and let 𝑆 be the corresponding R1-linearization. Then 𝐻 is R3-linearizable, and 𝑆 its R3-linearization. Proof. The completion 𝐻 ′ is connected by Lemma 5.1, and the R1-linearization 𝑆 is R2-conducive by Lemma 5.2. R3-linearizability with respect to 𝑆 follows by Theorem 3.7. □
5.2
Multi-Reader Case, 𝑁 = 3
In the presence of multiple readers, the construction of the linearization 𝑆 becomes more complex as the set 𝑅𝑘 of read operations that return the value 𝑣 obtained from a value-version pair ⟨𝑣, 𝑘⟩ is no longer totally ordered by physical causality via program order. Since distinct observers may experience Herlihy and Wing’s “happens before” relation differently due to relativity of simultaneity, one observer may perceive two operations as overlapping in time while another observer perceives one operation producing a response before the other is invoked. For R3-linearizability, the operations in the set 𝑅𝑘 must therefore be ordered in a manner that respects Herlihy and Wing’s “happens before” relation in every conceivable frame of reference. We address this problem by fixing the order of reads in 𝑅𝑘 so that if 𝑜𝑥 1, 𝑜𝑥 2 ∈ 𝑅𝑘 are weakly connected and 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 then 𝑜𝑥 1 →𝑆 𝑜𝑥 2 . To see why this is sufficient, consider the different ways in which 𝑜𝑥 1 and 𝑜𝑥 2 can be related by causality. If 𝑜𝑥 1 and 𝑜𝑥 2 are strongly connected, then they appear concurrent in all frames of reference, and can be ordered arbitrarily in 𝑆. On the other hand, if 𝑜𝑥 1 and 𝑜𝑥 2 are only weakly connected, then each observer either perceives the two operations as concurrent or perceives 𝑜𝑥 1 happening before 𝑜𝑥 2 since 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . In both cases, the ordering 𝑜𝑥 1 →𝑆 𝑜𝑥 2 is consistent with the observer’s subjective perception of time. It remains to show that the linearization 𝑆 exists, which is tantamount to the statement that the d𝐻 ′ relation cannot create a cycle when applied to weakly connected op-ex pairs in a given set 𝑅𝑘 . We present the proof below for 𝑁 = 3. The result is nontrivial since d𝐻 ′ is not a transitive relation in the general context, in contrast to Herlihy and Wing’s “happens before” relation.
Lemma 5.3. Let 𝐻 = (𝐸, <𝐸 ) be a history of ABD, and let 𝐻 ′ be its completion. For any 𝑘, let 𝐺𝑘 be the directed graph whose vertices are the read operation executions in 𝑅𝑘 and where an edge exists from 𝑜𝑥 1 to 𝑜𝑥 2 if and only if they are weakly connected with 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . Then 𝐺𝑘 is acyclic for all 𝑘.
Proof. Suppose for contradiction that for some 𝑘 the graph 𝐺𝑘 has a cycle 𝐶. The cycle must have at least three edges since an operation execution cannot be weakly connected to itself (in a one-cycle or loop), and similarly two op-exes cannot be weakly connected to each other. Let 𝑜𝑥 1 , 𝑜𝑥 2 , and 𝑜𝑥 3 be three consecutive read operation executions along the cycle such that 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 and 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 3 , and each consecutive pair is weakly connected. Without loss of generality, assume that 𝐶 is a cycle of minimal length in 𝐺𝑘 . Next, note that the operations in 𝐶 must be executed by distinct processors, as otherwise 𝐶 can be shortened. That is, if 𝑜𝑥𝑖 and 𝑜𝑥 𝑗 are executed by the same processor and 𝑜𝑥𝑖 precedes 𝑜𝑥 𝑗 in program order, then letting 𝑜𝑥ℎ denote the predecessor of 𝑜𝑥𝑖 along 𝐶, 𝑜𝑥ℎ and 𝑜𝑥 𝑗 are weakly connected with 𝑜𝑥ℎ d𝐻 ′ 𝑜𝑥 𝑗 , which makes it possible to remove 𝑜𝑥𝑖 (and any op-exes strictly between 𝑜𝑥𝑖 and 𝑜𝑥 𝑗 ) from the cycle. To see this, first note that because 𝑜𝑥ℎ d𝐻 ′ 𝑜𝑥𝑖 and 𝑜𝑥𝑖 →𝐻 ′ 𝑜𝑥 𝑗 , 𝑜𝑥ℎ d𝐻 ′ 𝑜𝑥 𝑗 holds. Furthermore, 𝑜𝑥 𝑗 d𝐻 ′ 𝑜𝑥ℎ is false, as otherwise 𝑜𝑥𝑖 →𝐻 ′ 𝑜𝑥 𝑗 implies 𝑜𝑥𝑖 d𝐻 ′ 𝑜𝑥ℎ , which contradicts 𝑜𝑥ℎ and 𝑜𝑥𝑖 being weakly connected with 𝑜𝑥ℎ d𝐻 ′ 𝑜𝑥𝑖 . Without loss of generality, suppose that processes are numbered so that 𝑝𝑖 executes 𝑜𝑥𝑖 for 𝑖 ∈ {1, 2, 3}. Since we have shown that 𝐶 has length at least three, and that all operations in the cycle are executed by distinct processors, it follows from the assumption 𝑁 = 3 that 𝐶 has length exactly three. Consequently, 𝑜𝑥 3 and 𝑜𝑥 1 are weakly connected with 𝑜𝑥 3 d𝐻 ′ 𝑜𝑥 1 . Now consider the causal relationships among 𝑜𝑥 1 , 𝑜𝑥 2 , and 𝑜𝑥 3 . Due to the use of majority quorums in ABD, 𝑜𝑥 1 and 𝑜𝑥 2 must both 𝑝 fetch value-version pairs from a common processor 𝑝. Let 𝑒 1 be the event where 𝑝 receives the value request sent by 𝑜𝑥 1 , and let 𝑝 𝑒 2 be the event where 𝑝 receives the value request sent by 𝑜𝑥 2 . 𝑝 𝑝 Since 𝑝 is sequential, 𝑒 1 and 𝑒 2 are related by <𝐸 via program order. 𝑝 𝑝 Moreover, it follows that 𝑒 1 <𝐸 𝑒 2 as otherwise there is a chain 𝑝 𝑝 of causality 𝑖𝑛𝑣 (𝑜𝑥 2 ) <𝐸 𝑒 2 <𝐸 𝑒 1 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 1 ) that contradicts 𝑜𝑥 1 and 𝑜𝑥 2 being weakly connected with 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 . Thus, 𝑝 𝑝 𝑖𝑛𝑣 (𝑜𝑥 1 ) <𝐸 𝑒 1 <𝐸 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 2 ) holds. Next, consider the action of 𝑜𝑥 3 . Due to the use of majority quorums in ABD, 𝑜𝑥 3 must fetch a value-version pair from at least one other processor 𝑞. Since 𝑁 = 3, this processor 𝑞 must be 𝑝 1 or 𝑞 𝑝 2 . Let 𝑒 1 be the event where 𝑞 receives the value request sent by 𝑞 𝑜𝑥 3 , and let 𝑒 2 be the event where 𝑞 sends the response received by 𝑞 𝑞 𝑜𝑥 3 . Observe that 𝑖𝑛𝑣 (𝑜𝑥 3 ) <𝐸 𝑒 1 <𝐸 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 3 ) holds. 𝑞 𝑞 Case A: 𝑞 = 𝑝 1 , hence 𝑒 1 and 𝑒 2 are both related by program 𝑞 order to 𝑜𝑥 1 . It follows that 𝑒 2 <𝐸 𝑖𝑛𝑣 (𝑜𝑥 1 ) holds, as otherwise 𝑞 𝑖𝑛𝑣 (𝑜𝑥 1 ) <𝐸 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 3 ), which contradicts 𝑜𝑥 1 and 𝑜𝑥 3 being 𝑞 weakly connected with 𝑜𝑥 3 d𝐻 ′ 𝑜𝑥 1 . Thus, 𝑖𝑛𝑣 (𝑜𝑥 3 ) <𝐸 𝑒 1 <𝐸 𝑞 𝑒 2 <𝐸 𝑖𝑛𝑣 (𝑜𝑥 1 ) holds. Since we showed earlier that 𝑖𝑛𝑣 (𝑜𝑥 1 ) <𝐸 𝑝 𝑝 𝑒 1 <𝐸 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 2 ) holds, it follows from the transitivity of <𝐸 that 𝑖𝑛𝑣 (𝑜𝑥 3 ) <𝐸 𝑟𝑒𝑠 (𝑜𝑥 2 ). This contradicts 𝑜𝑥 2 and 𝑜𝑥 3 being weakly connected with 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 3 .
Aeini and Golab
𝑞
𝑞
Case B: 𝑞 = 𝑝 2 , hence 𝑒 1 and 𝑒 2 are both related by program 𝑞 order to 𝑜𝑥 2 . It follows that 𝑟𝑒𝑠 (𝑜𝑥 2 ) <𝐸 𝑒 2 holds, as otherwise 𝑞 𝑞 𝑖𝑛𝑣 (𝑜𝑥 3 ) <𝐸 𝑒 1 <𝐸 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 2 ), which contradicts 𝑜𝑥 2 and 𝑜𝑥 3 being weakly connected with 𝑜𝑥 2 d𝐻 ′ 𝑜𝑥 3 . Thus, 𝑟𝑒𝑠 (𝑜𝑥 2 ) <𝐸 𝑞 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 3 ) holds. Since we showed earlier that 𝑖𝑛𝑣 (𝑜𝑥 1 ) <𝐸 𝑝 𝑝 𝑒 1 <𝐸 𝑒 2 <𝐸 𝑟𝑒𝑠 (𝑜𝑥 2 ) holds, it follows from the transitivity of <𝐸 that 𝑖𝑛𝑣 (𝑜𝑥 1 ) <𝐸 𝑟𝑒𝑠 (𝑜𝑥 3 ). This contradicts 𝑜𝑥 1 and 𝑜𝑥 3 being weakly connected with 𝑜𝑥 3 d𝐻 ′ 𝑜𝑥 1 . □ We now complete the analysis of R3-linearizability, noting that Lemma 5.3 implies the existence of the R1-linearization 𝑆 for the completion 𝐻 ′ of 𝐻 . Lemma 5.4. Let 𝐻 = (𝐸, <𝐸 ) be a history of ABD for 𝑁 = 3 processors, and let 𝐻 ′ be its completion. Fix a total order <𝑇 that refines <𝐸 , and let 𝑆 be the corresponding R1-linearization. Then for every pair of operation executions 𝑜𝑥 1 and 𝑜𝑥 2 in 𝐻 ′ , if 𝑜𝑥 1 and 𝑜𝑥 2 are weakly connected with 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 , then 𝑜𝑥 1 →𝑆 𝑜𝑥 2 . Proof. We proceed by an exhaustive case analysis on the possible pairs of operations in 𝐻 ′ . The case analysis is identical to the proof of Lemma 5.2, except for the proof in Case B, where 𝑜𝑥 1 and 𝑜𝑥 2 are both reads. Here we appeal to the special way in which 𝑆 is constructed for the multi-reader ABD algorithm based on the d𝐻 ′ relation, whereby 𝑜𝑥 1 d𝐻 ′ 𝑜𝑥 2 implies 𝑜𝑥 1 →𝑆 𝑜𝑥 2 for weakly connected operation execution pairs. □ Theorem 5.2. Let 𝐻 = (𝐸, <𝐸 ) be a history of multi-reader ABD for 𝑁 = 3 processors, and let 𝐻 ′ be its completion. Fix any total order <𝑇 that refines <𝐸 , and let 𝑆 be the corresponding R1-linearization. Then 𝐻 is R3-linearizable, and 𝑆 its R3-linearization. Proof. The completion 𝐻 ′ is connected by Lemma 5.1, and the R1-linearization 𝑆 is R2-conducive by Lemma 5.4. R3-linearizability with respect to 𝑆 follows by Theorem 3.7. □
6
QUORUM-REPLICATED KEY-VALUE STORES
Gilbert and Golab’s conjecture covers key-value storage systems that use majority quorums for replication, and they cite as an example Amazon’s Dynamo key-value store [2], which has a number of widely-used derivatives including Apache Cassandra [9]. Such systems store collections of key-value pairs, and provide get and put operations to read and write the value of a given key. Each keyvalue pair can be modelled as a shared read/write register object, with the size of a read or write quorum tunable individually for each get or put operation. To a first approximation, if the key-value protocol is configured with majority quorums then the protocol behaves similarly to the ABD simulation described in Section 5, though it uses only a single round of communication for both gets and puts.6 The shared object corresponding to each key-value pair therefore satisfies the semantics of Lamport’s regular register [11]. A closer look at Dynamo reveals that the question of linearizability is more nuanced. Get operations sometimes return a collection of conflicting data versions, and therefore cannot be mapped directly to read operations on a register object. Some conflicts are 6 The majority quorums must be strict, as opposed to “sloppy” quorums with hinted
hand-off used to maintain full write availability under network partitions.
resolved automatically using Dynamo’s vector clock (syntactic reconciliation), but the version tree occasionally forks into parallel branches that must be merged using application-specific business logic (semantic reconciliation), such as a “last write wins” rule over uniquely timestamped versions. Applications with stateful sessions can further enforce monotonic reads, avoiding old-new inversions within a session. Thus, the exact implementation of a read/write register from Dynamo get/put operations depends on application-side session management and reconciliation policies, leaving Gilbert and Golab’s conjecture somewhat open to interpretation. We believe that a Dynamo-style key-value store configured with strict majority quorums and the “last write wins” reconciliation policy can yield R3-linearizable histories in special cases. Assume a single reader maintaining a continuous session, which ensures monotonicity across all reads, and a single writer who assigns data timestamps using a local clock, which avoids clock synchronization issues. As in our analysis of ABD in Section 5, majority quorums guarantee connected histories, and under the singlewriter single-reader assumption the histories have R2-conducive R1-linearizations, allowing the application of Theorem 3.7. Modern implementations of the Dynamo design also provide mechanisms for strong consistency beyond session guarantees. Notably, recent versions of Apache Cassandra [9] implement lightweight transactions using Paxos consensus [12], yielding classically linearizable behaviour for multiple writers and readers; by arguments similar to our Raft analysis in Section 4, we believe such systems are R3-linearizable.
7
CONCLUSION
Our analysis of relativistic linearizability in this paper validates Gilbert and Golab’s proof technique [4], and partly confirms their conjecture with respect to state machine replication, the ABD shared register construction, and quorum-replicated key-value stores. We establish R3-linearizability for all three algorithms, albeit only in special cases for ABD and key-value stores. To our knowledge, these are the first rigorous proofs of R3-linearizability for any distributed algorithm, and the first applications of Gilbert and Golab’s proof technique. Jayanti’s recent result [8] that all classically linearizable asynchronous algorithms are R2-linearizable complements ours but does not subsume our R3-linearizability proofs, a strictly stronger property. We leave open a broader analysis of the Gilbert-Golab conjecture across more interpretations and special cases (e.g., multi-reader ABD for 𝑁 > 3), as well as a systematic study of relativistic linearizability for algorithms that break the asynchronous assumption by using physical clocks for coordination.
ACKNOWLEDGMENTS We acknowledge the support of the Natural Sciences and Engineering Research Council of Canada (NSERC).
Analyzing Linearizability in Relativistic Distributed Systems
REFERENCES [1] Hagit Attiya, Amotz Bar-Noy, and Danny Dolev. 1995. Sharing Memory Robustly in Message-passing Systems. J. ACM 42, 1 (1995), 124–142. [2] Giuseppe DeCandia, Deniz Hastorun, Madan Jampani, Gunavardhan Kakulapati, Avinash Lakshman, Alex Pilchin, Swami Sivasubramanian, Peter Vosshall, and Werner Vogels. 2007. Dynamo: Amazon’s Highly Available Key-Value Store. In Proc. of the 21st ACM Symposium on Operating Systems Principles. [3] A. Einstein. 2003. The Meaning of Relativity (6th ed.). Taylor & Francis. [4] Seth Gilbert and Wojciech Golab. 2014. Making Sense of Relativistic Distributed Systems. In Proc. of the 28th International Symposium on Distributed Computing (DISC). 361–375. [5] J. C. Hafele and Richard E. Keating. 1972. Around-the-World Atomic Clocks: Observed Relativistic Time Gains. Science 177, 4044 (1972), 168– 170. arXiv:https://www.science.org/doi/pdf/10.1126/science.177.4044.168 https: //www.science.org/doi/abs/10.1126/science.177.4044.168 [6] J. C. Hafele and Richard E. Keating. 1972. Around-the-World Atomic Clocks: Predicted Relativistic Time Gains. Science 177, 4044 (1972), 166– 168. arXiv:https://www.science.org/doi/pdf/10.1126/science.177.4044.166 https:
//www.science.org/doi/abs/10.1126/science.177.4044.166 [7] M. Herlihy and J. M. Wing. 1990. Linearizability: A correctness condition for concurrent objects. ACM TOPLAS 12, 3 (1990), 463–492. [8] Siddhartha Jayanti. 2025. On Interplanetary and Relativistic Distributed Computing. In Proc. of the 44th ACM Symposium on Principles of Distributed Computing (PODC). 192–202. [9] Avinash Lakshman and Prashant Malik. 2010. Cassandra: a decentralized structured storage system. ACM SIGOPS Operating Systems Review 44, 2 (2010), 35–40. [10] Leslie Lamport. 1978. Time, Clocks, and the Ordering of Events in a Distributed System. Commun. ACM 21, 7 (1978), 558–565. [11] L. Lamport. 1986. On Interprocess Communication, Part I: Basic Formalism. Distributed Computing 1, 2 (1986), 77–85. [12] Leslie Lamport. 1998. The Part-time Parliament. ACM Trans. Comput. Syst. 16, 2 (1998), 133–169. [13] Diego Ongaro and John Ousterhout. 2014. In Search of an Understandable Consensus Algorithm. In Proc. of the USENIX Annual Technical Conference (ATC). 305–319.