Efficient message ordering in dynamic networks
Idit Keidar, Danny Dolev · 1996
We present an algorithm for totally ordering messages in the face of network partitions and site failures.The algorithm aJways aJlows a majority of connected processors in the network to make progress (z.e. to order messages), if they remain connected for sufficiently long, regardless of past failures.Furthermore, our aJgorithm always allows processors to initiate messages, even when they are not members of a connected majority component in the network.Thus, messages can eventually become totally ordered even if their initiator is never a member of a majority component.The algorithm guarantees that when a majority is connected, each message is ordered within two communication rounds, if no failures occur during these rounds. 1 Introduction Consistent order is a powerful paradigm for the design of fault tolerant applications, e.g.consistent replication [Sch90, Kei94].We present an efficient algorithm for consistent message ordering in the face of network partitions and site failures, The network may partition into several components, and remerge.The algorithm is most adequate for dynamic networks where failures are transient.The algorithm uses an underlaying group communication service as a building block.Problem Definition Atomic broadcast deals with consistent message ordering.Informally, atomic broadcast requires that all the correct processors will deliver all the messages to the application in the same order and that they eventually deliver all messages sent by correct processors.In our