Sam Toueg

dblp:t/SamToueg · DBLP profile ↗
← Back
105ranked-venue papers
11as first author
8since 2021 · last 2026
0009-0000-2014-1132ORCID · corroborated

Domains — the database's venue-derived domains; a paper can count in several

Systems, architecture and hardware · 51 · 4 first-author · 5 since 2021Theory of computation · 22 · 6 first-author · 1 since 2021Databases, data management, data science and information retrieval · 11 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 6Security and privacy · 4Software engineering, systems software and programming languages · 3Computer networks · 1 · 1 first-author
YearPublicationVenuePosition
2026 Generalized Compare-and-Swap and Space-Efficient Universal Constructions for the Infinite-Arrival Model
abstract
We introduce GCAS, a natural generalization of the well-known compare-and-swap (CAS) object. Intuitively, GCAS just replaces the fixed equality test of CAS with a parametrized comparator chosen from {<, =, >}. To showcase the utility of GCAS, we present two space-efficient wait-free universal constructions for systems where the number of participating processes is unknown and may be infinite (the infinite-arrival model). The first has space-complexity linear in the number of processes that have participated so far, while the second has space-complexity linear in the point contention but assumes bounded concurrency. To the best of our knowledge, these are the first wait-free universal constructions that achieve this space complexity in the infinite-arrival model. To achieve space complexity linear in the point contention, our second universal construction uses a novel memory recycling scheme that works in the infinite-arrival model with bounded concurrency. The ideas behind this recycling scheme could be of more general use.
Vassos Hadzilacos, Myles Thiessen, Sam Toueg
PODC3
2025 You can lie but not deny: SWMR registers with signature properties in systems with Byzantine processes
abstract
We define and show how to implement SWMR registers that provide properties of unforgeable digital signatures—without actually using such signatures—in systems with Byzantine processes. More precisely, we first define SWMR verifiable registers. Intuitively, processes can use these registers to write values as if they are "signed", such that these "signed values" can be "verified" by any process and "relayed" to any process. We give a signature-free implementation of such registers from plain SWMR registers in systems with n > 3f processes, f of which can be Byzantine. We also give a signature-free implementation of SWMR sticky registers from SWMR registers in systems with n > 3f processes. Once the writer p writes a value υ into a SWMR sticky register R, the register never changes its value. Note that the value υ can be considered "signed" by p: once p writes υ in R, p cannot change the value in R or deny that it wrote υ in R, and every reader can verify that p wrote υ just by reading R. This holds even if the writer p of R is Byzantine. We prove that our implementations are optimal in the number of Byzantine processes they can tolerate. Since SWMR registers can be implemented in message-passing systems with Byzantine processes and n > 3f [11], the results in this paper also show that one can implement verifiable registers and sticky registers in such systems.
Xing Hu 0009, Sam Toueg
PODC2
2025 Asymmetric Linearizable Local Reads
abstract
Many linearizable local read algorithms have been proposed to minimize the read latency of strongly consistent distributed databases deployed in geo-distributed networks. These algorithms do so by enabling reads to be performed immediately against any process' copy of the database in the best case. However, as our analysis shows, worst-case read latency at every process with all existing algorithms is at least the network's relative diameter in terms of the maximum message delay minus a known lower bound on message delay between any two processes. We then show that by leveraging the asymmetric message delays of geo-distributed networks, worst-case read latency can be below the network's relative diameter at processes close to the leader or the network's center by presenting two new linearizable local read algorithms. Our experimental evaluation shows that these new algorithms reduce worst-case read latency by up to 50x compared to existing ones.
Myles Thiessen, Guy Khazma, Sam Toueg, Eyal de Lara
Proc. VLDB Endow.3
2024 On implementing SWMR registers from SWSR registers in systems with Byzantine failures
Xing Hu 0009, Sam Toueg
Distributed Comput.2
2022 On Implementing SWMR Registers from SWSR Registers in Systems with Byzantine Failures
abstract
The implementation of registers from (potentially) weaker registers is a classical problem in the theory of distributed computing. Since Lamport’s pioneering work [Leslie Lamport, 1986], this problem has been extensively studied in the context of asynchronous processes with crash failures. In this paper, we investigate this problem in the context of Byzantine process failures, with and without process signatures. In particular, we first show a strong impossibility result, namely, that there is no wait-free linearizable implementation of a 1-writer n-reader register from atomic 1-writer (n-1)-reader registers. In fact, this impossibility result holds even if all the processes except the writer are given atomic 1-writer n-reader registers, and even if we assume that the writer can only crash and at most one reader is subject to Byzantine failures. In light of this impossibility result, we give two register implementations. The first one implements a 1-writer n-reader register from atomic 1-writer 1-reader registers. This implementation is linearizable (under any combination of Byzantine process failures), but it is wait-free only under the assumption that the writer is correct or no reader is Byzantine - thus matching the impossibility result. The second implementation assumes process signatures; it is wait-free and linearizable under any number and combination of Byzantine process failures.
Xing Hu 0009, Sam Toueg
DISC2
2022 On atomic registers and randomized consensus in M&M systems
abstract
Motivated by recent distributed systems technology, Aguilera et al. introduced a hybrid model of distributed computing, called the message-and-memory model or m&m model for short. In this model, processes can communicate by message passing and also by accessing some shared memory (e.g., through some RDMA connections). We first consider the basic problem of implementing an atomic single-writer multi-reader (SWMR) register shared by all the processes in m&m systems. Specifically, we give an algorithm that implements such a register in m&m systems and show that it is optimal in the number of process crashes that it tolerates. This generalizes the well-known ABD implementation of an atomic SWMR register in a pure message-passing system. We then combine our register implementation for m&m systems with a randomized consensus algorithm of Aspnes and Herlihy, and obtain a randomized consensus algorithm for m&m systems that is also optimal in the number of process crashes that it can tolerate. Finally, we determine the minimum number of RDMA connections that is sufficient to implement a SWMR register, or solve randomized consensus, in an m&m system with t process crashes, for any given t .
Vassos Hadzilacos, Xing Hu 0009, Sam Toueg
Distributed Comput.3
2022 Randomized consensus with regular registers
Vassos Hadzilacos, Xing Hu 0009, Sam Toueg
Inf. Process. Lett.3
2021 On Register Linearizability and Termination
abstract
It is well-known that, for deterministic algorithms, linearizable objects can be used as if they were atomic objects. As pointed out by Golab, Higham, and Woelfel, however, a randomized algorithm that works with atomic objects may lose some of its properties if we replace the atomic objects that it uses with objects that are only linearizable. It was not known whether the properties that can be lost include the all-important property of termination (with probability 1). In this paper, we first show that a randomized algorithm can indeed lose its termination property if we replace the atomic registers that it uses with linearizable ones.
Vassos Hadzilacos, Xing Hu 0009, Sam Toueg
PODC3
2020 Life beyond set agreement
David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
Distributed Comput.3
2020 Bounded disagreement
David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
Theor. Comput. Sci.3
2019 Optimal Register Construction in M&M Systems
abstract
Motivated by recent distributed systems technology, Aguilera et al. introduced a hybrid model of distributed computing, called message-and-memory model or m&m model for short [Marcos K. Aguilera et al., 2018]. In this model, processes can communicate by message passing and also by accessing some shared memory. We consider the basic problem of implementing an atomic single-writer multi-reader (SWMR) register shared by all the processes in m&m systems. Specifically, we give an algorithm that implements such a register in m&m systems and show that it is optimal in the number of process crashes that it can tolerate. This generalizes the well-known implementation of an atomic SWMR register in a pure message-passing system [Attiya et al., 1995].
Vassos Hadzilacos, Xing Hu 0009, Sam Toueg
OPODIS3
2019 On Deterministic Linearizable Set Agreement Objects
abstract
A recent work showed that, for all n and k, there is a linearizable (n,k)-set agreement object O_L that is equivalent to the (n,k)-set agreement task [David Yu Cheng Chan et al., 2017]: given O_L, it is possible to solve the (n,k)-set agreement task, and given any algorithm that solves the (n,k)-set agreement task (and registers), it is possible to implement O_L. This linearizable object O_L, however, is not deterministic. It turns out that there is also a deterministic (n,k)-set agreement object O_D that is equivalent to the (n,k)-set agreement task, but this deterministic object O_D is not linearizable. This raises the question whether there exists a deterministic and linearizable (n,k)-set agreement object that is equivalent to the (n,k)-set agreement task. Here we show that in general the answer is no: specifically, we prove that for all n ≥ 4, every deterministic linearizable (n,2)-set agreement object is strictly stronger than the (n,2)-set agreement task. We prove this by showing that, for all n ≥ 4, every deterministic and linearizable (n,2)-set agreement object (together with registers) can be used to solve 2-consensus, whereas it is known that the (n,2)-set agreement task cannot do so. For a natural subset of (n,2)-set agreement objects, we prove that this result holds even for n = 3.
Felipe de Azevedo Piovezan, Vassos Hadzilacos, Sam Toueg
OPODIS3
2018 Passing Messages while Sharing Memory
abstract
We introduce a new distributed computing model called m&m that allows processes to both pass messages and share memory. Motivated by recent hardware trends, we find that this model improves the power of the pure message-passing and shared-memory models. As we demonstrate by example with two fundamental problems---consensus and eventual leader election---the added power leads to new algorithms that are more robust against failures and asynchrony. Our consensus algorithm combines the superior scalability of message passing with the higher fault tolerance of shared memory, while our leader election algorithms reduce the system synchrony needed for correctness. These results point to a wide new space for future exploration of other problems, techniques, and benefits.
Marcos K. Aguilera, Naama Ben-David, Irina Calciu, Rachid Guerraoui, Erez Petrank, Sam Toueg
PODC6
2018 On the Classification of Deterministic Objects via Set Agreement Power
abstract
Since the early days of the shared memory model for distributed computing, researchers have sought a simple and precise characterization of an object's ability to implement other objects in a wait-free manner.
David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
PODC3
2017 Life Beyond Set Agreement
abstract
The set agreement power of a shared object O describes O's ability to solve set agreement problems: it is the sequence (n_1, n_2, ..., n_k, ...) such that, for every k >= 1, using O and registers one can solve the k-set agreement problem among at most n_k processes. It has been shown that the ability of an object O to implement other objects is not fully characterized by its consensus number the first component of its set agreement power) [1, 3, 14]. This raises the following natural question: is the ability of an object O to implement other objects fully characterized by its set agreement power? We prove that the answer is no: every level n >= 2 of Herlihy's consensus hierarchy has two objects that have the same set agreement power but are not equivalent, i.e., at least one cannot implement the other. We also show that every level n >= 2 of the consensus hierarchy contains a deterministic object O_n with some set agreement power (n_1, n_2, ..., n_k, ...) such that being able to solve the k-set agreement problems among n_k processes, for all k >= 1, is not enough to implement O_n.
David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
PODC3
2017 On the Number of Objects with Distinct Power and the Linearizability of Set Agreement Objects
abstract
We first prove that there are uncountably many objects with distinct computational powers. More precisely, we show that there is an uncountable set of objects such that for any two of them, at least one cannot be implemented from the other (and registers) in a wait-free manner. We then strengthen this result by showing that there are uncountably many linearizable objects with distinct computational powers. To do so, we prove that for all positive integers n and k, there is a linearizable object that is computationally equivalent to the k-set agreement task among n processes. To the best of our knowledge, these are the first linearizable objects proven to be computationally equivalent to set agreement tasks.
David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
DISC3
2016 Bounded Disagreement
David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
OPODIS3
2016 An Algorithm for Replicated Objects with Efficient Reads
abstract
The problem. We consider the problem of implementing a consistent replicated object in a partially synchronous message passing distributed system susceptible to process and communication failures. The object is a generic shared resource, such as a data structure, a file, or a lock. The processes implementing the replicated object access it by applying operations to it at unpredictable times and potentially concurrently.1 The object should be linearizable: it should behave as if each operation applied to it takes effect at a distinct instant in time during the interval between its invocation and its response.
Tushar Deepak Chandra, Vassos Hadzilacos, Sam Toueg
PODC3
2016 k-Abortable Objects: Progress Under High Contention
Naama Ben-David, David Yu Cheng Chan, Vassos Hadzilacos, Sam Toueg
DISC4
2015 A Separation of n-consensus and (n + 1)-consensus Based on Process Scheduling
Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
SIROCCO3
2013 On deterministic abortable objects
abstract
We define deterministic abortable (DA) objects, which guarantee that operations complete normally if executed solo, but may abort if executed concurrently with other operations. An operation that aborts has no effect on the object. This simple and attractive behavior is reminiscent of transactional memory, database transactions, and abortable mutual exclusion --- techniques in which a process can, under contention, ``bail out'' of the computation without leaving a trace.
Vassos Hadzilacos, Sam Toueg
PODC2
2012 Partial synchrony based on set timeliness
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
Distributed Comput.4
2012 The correctness proof of Ben-Or's randomized consensus algorithm
Marcos K. Aguilera, Sam Toueg
Distributed Comput.2
2012 The Weakest Failure Detectors to Solve Quittable Consensus and Nonblocking Atomic Commit
abstract
We define quittable consensus, a natural variation of the consensus problem, where processes have the option to agree on “quit” if failures occur, and we relate this problem to the well-known problem of nonblocking atomic commit. We then determine the weakest failure detectors for these two problems in all environments, regardless of the number of faulty processes.
Rachid Guerraoui, Vassos Hadzilacos, Petr Kuznetsov, Sam Toueg
SIAM J. Comput.4
2011 The minimum information about failures for solving non-local tasks in message-passing systems
Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
Distributed Comput.3
2010 Adaptive progress: a gracefully-degrading liveness property
Marcos K. Aguilera, Sam Toueg
Distributed Comput.2
2009 The Minimum Information about Failures for Solving Non-local Tasks in Message-Passing Systems
Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
OPODIS3
2009 Partial synchrony based on set timeliness
abstract
We introduce a new model of partial synchrony for read-write shared memory systems. This model is based on the notion of set timeliness--a natural and straightforward generalization of the seminal concept of timeliness in the partially synchrony model of Dwork, Lynch and Stockmeyer [8].
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
PODC4
2009 Brief Announcement: The Minimum Failure Detector for Non-Local Tasks in Message-Passing Systems
Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DISC3
2008 A robust and lightweight stable leader election service for dynamic systems
abstract
We describe the implementation and experimental evaluation of a fault-tolerant leader election service for dynamic systems. Intuitively, distributed applications can use this service to elect and maintain an operational leader for any group of processes which may dynamically change. If the leader of a group crashes, is temporarily disconnected, or voluntarily leaves the group, the service automatically re-elects a new group leader. The current version of the service implements two recent leader election algorithms, and users can select the one that fits their system better. Both algorithms ensure leader stability, a desirable feature that lacks in some other algorithms, but one is more robust in the face of extreme network disruptions, while the other is more scalable. The leader election service is flexible and easy to use. By using a stochastic failure detector and a link quality estimator, it provides some degree of QoS control and it adapts to changing network conditions. Our experimental evaluation indicates that it is also highly robust and inexpensive to run in practice.
Nicolas Schiper, Sam Toueg
DSN2
2008 With Finite Memory Consensus Is Easier Than Reliable Broadcast
Carole Delporte-Gallet, Stéphane Devismes, Hugues Fauconnier, Franck Petit, Sam Toueg
OPODIS5
2008 Timeliness-based wait-freedom: a gracefully degrading progress condition
abstract
We introduce a simple progress condition for shared object implementations that is gracefully degrading depending on the degree of synchrony in each run. This progress property, called timeliness-based wait-freedom, provides a gradual bridge between obstruction-freedom and wait-freedom in partially synchronous systems. We show that timeliness-based wait-freedom can be achieved with synchronization primitives that are very weak. More precisely, every object has a timeliness-based wait-free implementation that uses only abortable registers (which are weaker than safe registers). As part of this work, we present a new leader election primitive that processes can use to dynamically compete for leadership such that if there is at least one timely process among the current candidates for leadership, then a timely leader is eventually elected among the candidates. We also show that this primitive can be implemented using abortable registers.
Marcos K. Aguilera, Sam Toueg
PODC2
2008 Every problem has a weakest failure detector
abstract
Several basic problems that arise in fault-tolerant distributed computing were shown to have a weakest failure detector. We show here that every problem that is solvable with a failure detector has a weakest failure detector.
Prasad Jayanti, Sam Toueg
PODC2
2008 On implementing omega in systems with weak reliability and synchrony assumptions
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
Distributed Comput.4
2007 Abortable and query-abortable objects and their efficient implementation
abstract
We introduce abortable and query-abortable objects, intended for asynchronous shared-memory systems with low contention. These objects behave like ordinary objects when accessed sequentially, but may abort operations when accessed concurrently. An aborted operation may or may not take effect, i.e., cause a state transition, and it returns no indication of which possibility occurred. Since this uncertainty is problematic, a query-abortable object supports a QUERY operation that each process can use to determine its last non-QUERY operation on the object that caused a state transition, and the response associated with this state transition. Query-abortable objects can easily implement obstruction-free objects (introduced by Herlihy, Luchangco and Moir) and pausable objects (introduced by Attiya, Guerraoui and Kouznetsov).
Marcos K. Aguilera, Svend Frølund, Vassos Hadzilacos, Stephanie Lorraine Horn, Sam Toueg
PODC5
2007 DISC at Its 20th Anniversary (Stockholm, 2006)
Michel Raynal, Sam Toueg, Shmuel Zaks
DISC2
2007 The weakest failure detector to solve nonuniform consensus
Jonathan Eisler, Vassos Hadzilacos, Sam Toueg
Distributed Comput.3
2006 Consensus with Byzantine Failures and Little System Synchrony
abstract
We study consensus in a message-passing system where only some of the n2links exhibit some synchrony. This problem was previously studied for systems with process crashes; we now consider Byzantine failures. We show that consensus can be solved in a system where there is at least one non-faulty process whose links are eventually timely; all other links can be arbitrarily slow. We also show that, in terms of problem solvability, such a system is strictly weaker than one where all links are eventually timely
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DSN4
2006 Brief Announcement: Abortable and Query-Abortable Objects
Marcos K. Aguilera, Svend Frølund, Vassos Hadzilacos, Stephanie Lorraine Horn, Sam Toueg
DISC5
2006 From Set Membership to Group Membership: A Separation of Concerns
abstract
We revisit the well-known group membership problem and show how it can be considered a special case of a simple problem, the set membership problem. In the set membership problem, processes maintain a set whose elements are drawn from an arbitrary universe: They can request the addition or removal of elements to/from that set, and they agree on the current value of the set. Group membership corresponds to the special case where the elements of the set happen to be processes. We exploit this new way of looking at group membership to give a simple and succinct specification of this problem and to outline a simple implementation approach based on the state machine paradigm. This treatment of group membership separates several issues that are often mixed in existing specifications and/or implementations of group membership. We believe that this separation of concerns greatly simplifies the understanding of this problem.
André Schiper, Sam Toueg
IEEE Trans. Dependable Secur. Comput.2
2005 Fast fault-tolerant agreement algorithms
abstract
In the synchronous round-based model, a process crash is dirty if it occurs exactly while a process is sending messages in a round, and this causes the process to send to some, but not all, of the intended recipients for the given round. Dirty crashes are possible; however, they are unlikely to occur, since the time spent sending messages is usually very small compared to the maximum message delay (i.e., compared to the duration of a round). In this paper, we investigate how fast one can solve some agreement problems, namely consensus and terminating reliable broadcast (TRB), when the number of dirty crashes that occur is small. In particular, we describe some algorithms for the uniform and non-uniform versions of these problems, and provide some matching lower bounds. All our uniform algorithms are strictly better than conventional early-stopping algorithms, in the sense that they never take more rounds to decide or halt, and they take fewer rounds when the number of dirty crashes is small.
Carole Delporte-Gallet, Hugues Fauconnier, Stephanie Lorraine Horn, Sam Toueg
PODC4
2005 The weakest failure detector to solve nonuniform consensus
abstract
We determine the weakest failure detector to solve nonuniform consensus in any environment, i.e., regardless of the number of faulty processes. Together with previous results, this closes all aspects of the following question: What is the weakest failure detector to solve (uniform or nonuniform) consensus in any environment?
Jonathan Eisler, Vassos Hadzilacos, Sam Toueg
PODC3
2004 Communication-efficient leader election and consensus with limited link synchrony
abstract
We study the degree of synchrony required to implement the leader election failure detector Ω and to solve consensus in partially synchronous systems. We show that in a system with n processes and up to f process crashes, one can implement Ω and solve consensus provided there exists some (unknown) correct process with f outgoing links that are eventually timely. In the special case where f = 1 , an important case in practice, this implies that to implement Ω and solve consensus it is sufficient to have just one eventually timely link -- all the other links in the system, Θ(n2) of them, may be asynchronous. There is no need to know which link p → q is eventually timely, when it becomes timely, or what is its bound on message delay. Surprisingly, it is not even required that the source p or destination q of this link be correct: either p or q may actually crash, in which case the link p → q is eventually timely in a trivial way, and it is useless for sending messages. We show that these results are in a sense optimal: even if every process has f - 1 eventually timely links, neither Ω nor consensus can be solved. We also give an algorithm that implements Ω in systems where some correct process has f outgoing links that are eventually timely, such that eventually only f links carry messages, and we show that this is optimal. For f = 1 , this algorithm ensures that all the links, except for one, eventually become quiescent.
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
PODC4
2004 The weakest failure detectors to solve certain fundamental problems in distributed computing
abstract
We determine the weakest failure detectors to solve several fundamental problems in distributed message-passing systems, for all environments -- i.e., regardless of the number and timing of crashes. The problems that we consider are: implementing an atomic register, solving consensus, solving quittable consensus (a variant of consensus in which processes have the option to decide 'quit' if a failure occurs), and solving non-blocking atomic commit.
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Vassos Hadzilacos, Petr Kuznetsov, Sam Toueg
PODC6
2004 Generalized Irreducibility of Consensus and the Equivalence of t-Resilient and Wait-Free Implementations of Consensus
abstract
We study the consensus problem, which requires multiple processes with different input values to agree on one of these values, in the context of asynchronous shared memory systems. Prior research focussed either on t-resilient solutions of this problem (which must be correct even if up to t processes crash) or on wait-free solutions (which must be correct despite the crash of any number of processes). In this paper, we show that these two forms of solvability are closely related. Specifically, for all $n > t \ge 2$ and all sets ${\mathcal{S}}$ of shared object types (that include simple read/write registers), there is a t-resilient solution to n-process consensus using objects of types in ${\mathcal{S}}$ if and only if there is a wait-free solution to (t + 1)-process consensus using objects of types in ${\mathcal{S}}$. Our proof of this equivalence uses another result derived in this paper, which is of independent interest. Roughly speaking, this result states that a wait-free solution to (n - 1)-process consensus is never necessary in designing a wait-free solution to n-process consensus, regardless of the types of objects available. More precisely, for all $n \ge 2$ and all sets ${\mathcal{S}}$ of shared object types (that include simple read/write registers), if there is a wait-free solution to n-process consensus that uses a wait-free solution to (n - 1)-process consensus and objects of types in ${\mathcal{S}}$, then there is a wait-free solution to n-process consensus that uses only objects of types in ${\mathcal{S}}$.
Tushar Deepak Chandra, Vassos Hadzilacos, Prasad Jayanti, Sam Toueg
SIAM J. Comput.4
2003 On implementing omega with weak reliability and synchrony assumptions
abstract
We study the feasibility and cost of implementing Ω---a fundamental failure detector at the core of many algorithms---in systems with weak reliability and synchrony assumptions. Intuitively, Ω allows processes to eventually elect a common leader. We first give an algorithm that implements Ω in a weak system S where processes are synchronous, but: (a) any number of them may crash, and (b) only the output links of an unknown correct process are eventually timely (all other links can be asynchronous and/or lossy). This is in contrast to previous implementations of Ω which assume that a quadratic number of links are eventually timely, or systems that are strong enough to implement the eventually perfect failure detector P. We next show that implementing Ω in S is expensive: even if we want an implementation that tolerates just one process crash, all correct processes (except possibly one) must send messages forever; moreover, a quadratic number of links must carry messages forever. We then show that with a small additional assumption---the existence of some unknown correct process whose asynchronous links are lossy but fair---we can implement Ω efficiently: we give an algorithm for Ω such that eventually only one process (the elected leader) sends messages.
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
PODC4
2002 On the Impact of Fast Failure Detectors on Real-Time Fault-Tolerant Systems
Marcos K. Aguilera, Gérard Le Lann, Sam Toueg
DISC3
2002 On the Quality of Service of Failure Detectors
abstract
We study the quality of service (QoS) of failure detectors. By QoS, we mean a specification that quantifies: (1) how fast the failure detector detects actual failures and (2) how well it avoids false detections. We first propose a set of QoS metrics to specify failure detectors for systems with probabilistic behaviors, i.e., for systems where message delays and message losses follow some probability distributions. We then give a new failure detector algorithm and analyze its QoS in terms of the proposed metrics. We show that, among a large class of failure detectors, the new algorithm is optimal with respect to some of these QoS metrics. Given a set of failure detector QoS requirements, we show how to compute the parameters of our algorithm so that it satisfies these requirements and we show how this can be done even if the probabilistic behavior of the system is not known. We then present some simulation results that show that the new failure detector algorithm provides a better QoS than an algorithm that is commonly used in practice. Finally, we suggest some ways to make our failure detector adaptive to changes in the probabilistic behavior of the network.
Wei Chen 0013, Sam Toueg, Marcos K. Aguilera
IEEE Trans. Computers2
2002 On the Quality of Service of Failure Detectors
abstract
We study the quality of service (QoS) of failure detectors. By QoS, we mean a specification that quantifies 1) how fast the failure detector detects actual failures and 2) how well it avoids false detections. We first propose a set of QoS metrics to specify failure detectors for systems with probabilistic behaviors, i.e., for systems where message delays and message losses follow some probability distributions. We then give a new failure detector algorithm and analyze its QoS in terms of the proposed metrics. We show that, among a large class of failure detectors, the new algorithm is optimal with respect to some of these QoS metrics. Given a set of failure detector QoS requirements, we show how to compute the parameters of our algorithm so that it satisfies these requirements and we show how this can be done,even if the probabilistic behavior of the system is not known. We then present some simulation results that show that the new failure detector algorithm provides a better QoS than an algorithm that is commonly used in practice. Finally, we suggest some ways to make our failure detector adaptive to changes in the probabilistic behavior of the network.
Wei Chen 0013, Sam Toueg, Marcos K. Aguilera
IEEE Trans. Computers2
2001 Stable Leader Election
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DISC4
2000 Revisiting Safety and Liveness in the Context of Failures
Bernadette Charron-Bost, Sam Toueg, Anindya Basu
CONCUR2
2000 On the Quality of Service of Failure Detectors
abstract
Studies the quality of service (QoS) of failure detectors. By QoS, we mean a specification that quantifies (a) how fast the failure detector detects actual failures, and (b) how well it avoids false detections. We first propose a set of QoS metrics to specify failure detectors for systems with probabilistic behaviors, i.e. for systems where message delays and message losses follow some probability distributions. We then give a new failure detector algorithm and analyze its QoS in terms of the proposed metrics. We show that, among a large class of failure detectors, the new algorithm is optimal with respect to some of these QoS metrics. Given a set of failure detector QoS requirements, we show how to compute the parameters of our algorithm so that it satisfies these requirements, and we show how this can be done even if the probabilistic behavior of the system is not known. Finally, we briefly explain how to make our failure detector adaptive, so that it automatically reconfigures itself when there is a change in the probabilistic behavior of the network.
Wei Chen 0013, Sam Toueg, Marcos K. Aguilera
DSN2
2000 Thrifty Generic Broadcast
Marcos K. Aguilera, Carole Delporte-Gallet, Hugues Fauconnier, Sam Toueg
DISC4
2000 Failure Detection and Consensus in the Crash-Recovery Model
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
Distributed Comput.3
2000 On Quiescent Reliable Communication
abstract
We study the problem of achieving reliable communication with quiescent algorithms (i.e., algorithms that eventually stop sending messages) in asynchronous systems with process crashes and lossy links. We first show that it is impossible to solve this problem in asynchronous systems (with no failure detectors). We then show that, among failure detectors that output lists of suspects, the weakest one that can be used to solve this problem is $\diamond \cal P,$ a failure detector that cannot be implemented. To overcome this difficulty, we introduce an implementable failure detector called Heartbeat and show that it can be used to achieve quiescent reliable communication. Heartbeat is novel: in contrast to typical failure detectors, it does not output lists of suspects and it is implementable without timeouts. With Heartbeat, many existing algorithms that tolerate only process crashes can be transformed into quiescent algorithms that tolerate both process crashes and message losses. This can be applied to consensus, atomic broadcast, k-set agreement, atomic commitment, etc.
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
SIAM J. Comput.3
2000 Time and Space Lower Bounds for Nonblocking Implementations
abstract
We show the following time and space complexity lower bounds. Let $\cal{I}$ be any randomized nonblocking n-process implementation of any object in set A from any combination of objects in set B, where A = {increment, fetch&add, modulo k counter (for any $k \ge 2n$), LL/SC bit, k-valued compare&swap (for any $k \ge n$), single-writer snapshot}, and B = {resettable consensus} $\cup$ {historyless objects such as registers and swap registers}. The space complexity of $\cal{I}$ is at least n-1. Moreover, if $\cal{I}$ is deterministic, both its time and space complexity are at least n-1. These lower bounds hold even if objects used in the implementation are of unbounded size. This improves on some of the $\Omega(\sqrt{n})$ space complexity lower bounds of Fich, Herlihy, and Shavit [i Proceedings of the 12th Annual ACM Symposium on Principles of Distributed Computing, Ithaca, NY, 1993, pp. 241--249; J. Assoc. Comput. Mach., 45 (1998), pp. 843--862]. It also shows the near optimality of some known wait-free implementations in terms of space complexity.
Prasad Jayanti, King Tan, Sam Toueg
SIAM J. Comput.3
1999 Revising the Weakest Failure Detector for Uniform Reliable Broadcast
Marcos K. Aguilera, Sam Toueg, Borislav Deianov
DISC2
1999 A Simple Bivalency Proof that t-Resilient Consensus Requires t + 1 Rounds
abstract
We use a straightforward bivalency argument borrowed from Fischer et al. (1985) to show that in a synchronous system with up to t crash failures solving consensus requires at least t+1 rounds. The proof is simpler and more intuitive than the traditional one: It uses an easy forward induction rather than a more complex backward induction which needs the induction hypothesis several times.
Marcos K. Aguilera, Sam Toueg
Inf. Process. Lett.2
1999 The Cost of Graceful Degradation for Omission Failures
abstract
An implementation of a shared object O is t-tolerant if the object remains correct and wait-free even when up to t base objects (objects used in the implementation of O) fail. The implementation is gracefully degrading if, no matter how many base objects fail, O does not fail more severely than its base objects. For the omission failure mode, we derive a lower bound on the space complexity of a gracefully degrading t-tolerant implementation. This result lets us conclude that, for omission failures, graceful degradation can be achieved only at the cost of increased space complexity.
Prasad Jayanti, Tushar Deepak Chandra, Sam Toueg
Inf. Process. Lett.3
1999 Using the Heartbeat Failure Detector for Quiescent Reliable Communication and Consensus in Partitionable Networks
abstract
We consider partitionable networks with process crashes and lossy links, and focus on the problems of reliable communication and consensus for such networks. For both problems we seek algorithms that are quiescent, i.e., algorithms that eventually stop sending messages. We first tackle the problem of reliable communication for partitionable networks by extending the results of Aguilera et al. (1997). In particular, we generalize the specification of the heartbeat failure detector HB, show how to implement it, and show how to use it to achieve quiescent reliable communication. We then turn our attention to the problem of consensus for partitionable networks. We first show that, even though this problem can be solved using a natural extension of failure detector ♢ L, such solutions are not quiescent — in other words, ♢ L alone is not sufficient to achieve quiescent consensus in partitionable networks. We then solve this problem using ♢ L and the quiescent reliable communication primitives that we developed in the first part of the paper.
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
Theor. Comput. Sci.3
1998 Failure Detection and Consensus in the Crash-Recovery Model
Marcos K. Aguilera, Wei Chen 0013, Sam Toueg
DISC3
1998 Fault-Tolerant Wait-Free Shared Objects
abstract
Wait-free implementations of shared objects tolerate the failure of processes, but not the failure of base objects from which they are implemented. We consider the problem of implementing shared objects that tolerate the failure of both processes and base objects. We identify two classes of object failures: responsive and nonresponsive . With responsive failures, a faulty object responds to every operation, but its responses may be incorrect. With nonresponsive failures, a faulty object may also “hang” without responding. In each class, we define crash, omission, and arbitrary modes of failure. We show that all responsive failure modes can be tolerated. More precisely, for all responsive failure modes ℱ, object types T , and t ≥ 0, we show how to implement a shared object of type T which is t -tolerant for ℱ. Such an object remains correct and wait-free even if up to t base objects fail according to ℱ. In contrast to responsive failures, we show that even the most benign non-responsive failure mode cannot be tolerated. We also show that randomization can be used to circumvent this impossibility result. Graceful degradation is a desirable property of fault-tolerant implementations: the implemented object never fails more severely than the base objects it is derived from, even if all the base objects fail. For several failure modes, we show wheter this property can be achieved, and, if so, how.
Prasad Jayanti, Tushar Deepak Chandra, Sam Toueg
J. ACM3
1998 Failure Detection and Randomization: A Hybrid Approach to Solve Consensus
abstract
We present a consensus algorithm that combines unreliable failure detection and randomization, two well-known techniques for solving consensus in asynchronous systems with crash failures. This hybrid algorithm combines advantages from both approaches: it guarantees deterministic termination if the failure detector is accurate, and probabilistic termination otherwise. In executions with no failures or failure detector mistakes, the most likely ones in practice, consensus is reached in only two asynchronous rounds.
Marcos K. Aguilera, Sam Toueg
SIAM J. Comput.2
1996 Crash Failures vs. Crash + Link Failures (Abstract)
Anindya Basu, Bernadette Charron-Bost, Sam Toueg
PODC3
1996 On the Impossibility of Group Membership
abstract
Projet REFLECS
Tushar Deepak Chandra, Vassos Hadzilacos, Sam Toueg, Bernadette Charron-Bost
PODC3
1996 Time and Space Lower Bounds for Non-Blocking Implementations (Preliminary Version)
abstract
We show the following time and space complexity lower bounds.Let Z be any randomized nonbk, ;ng n-process implementation of any object in 1 from any combination of objects in set 1?, where A = {increment, store-conditional bit, compare@swap, bounded-counter, single-writer atomic snapshot, fetch~add}, and B = { resettable consensus, register, swap register}.The space complexity of ~is at least n -1.Moreover, if ~is deterministic, both its time and space complexit y are at least n -1.These lower bounds hold even if objects used in the implementation are of unbounded size.This improves on some of the C?(@) space com-plexit~lower bounds of Fich, Herlihy & Shavit [FHS93].It also shows the near optimality of Dome known wait-free implementations in terms of space complexity.
Prasad Jayanti, King Tan, Sam Toueg
PODC3
1996 The Weakest Failure Detector for Solving Consensus
abstract
We determine what information about failures is necessary and sufficient to solve Consensus in asynchronous distributed systems subject to crash failures. In Chandra and Toueg [1996], it is shown thatW, a failure detector that provides surprisingly little information about which processes have crashed, is sufficient to solve Consensus in asynchronous systems with a majority of correct processes. In this paper, we prove that to solve Consensus, any failure detector has to provide at least as much information as W. Thus, W is indeed the weakest failure detector for solving Consensus in asynchronous systems with a majority of correct processes.
Tushar Deepak Chandra, Vassos Hadzilacos, Sam Toueg
J. ACM3
1996 Unreliable Failure Detectors for Reliable Distributed Systems
abstract
We introduce the concept of unreliable failure detectors and study how they can be used to solve Consensus in asynchronous systems with crash failures. We characterise unreliable failure detectors in terms of two properties—completeness and accuracy. We show that Consensus can be solved even with unreliable failure detectors that make an infinite number of mistakes, and determine which ones can be used to solve Consensus despite any number of crashes, and which ones require a majority of correct processes. We prove that Consensus and Atomic Broadcast are reducible to each other in asynchronous systems with crash failures; thus, the above results also apply to Atomic Broadcast. A companion paper shows that one of the failure detectors introduced here is the weakest failure detector for solving Consensus [Chandra et al. 1992].
Tushar Deepak Chandra, Sam Toueg
J. ACM2
1994 Wait-Freedom vs. t-Resiliency and the Robustness of Wait-Free Hierarchies
abstract
We seek two properties in such a hierarchy:(1) If a type T is at level N, then, for all types T',
Tushar Deepak Chandra, Vassos Hadzilacos, Prasad Jayanti, Sam Toueg
PODC4
1993 Simulating Synchronized Clocks and Common Knowledge in Distributed Systems
abstract
Time and knowledge are studied in synchronous and asynchronous distributed systems. A large class of problems that can be solved using logical clocks as if they were perfectly synchronized clocks is formally characterized. For the same class of problems, a broadcast primitive that can be used as if it achieves common knowledge is also proposed. Thus, logical clocks and the broadcast primitive simplify the task of designing and verifying distributed algorithms: The designer can assume that processors have access to perfectly synchronized clocks and the ability to achieve common knowledge.
Gil Neiger, Sam Toueg
J. ACM2
1992 Fault-tolerant Wait-free Shared Objects
abstract
The authors classify object failures into two broad categories: responsive and non-responsive. They require that wait-free objects subject to responsive failures continue to respond (in finite time) to operation invocations. The responses may be incorrect. In contrast, wait-free objects subject to non-responsive failures are exempt from responding to operation invocations. Such objects may 'hang' on the invoking process. They divide responsive failures into three models: R-crash,R-omission, and R-arbitrary. They divide non-responsive failures into crash, omission, and arbitrary. An object subject to crash failure behaves correctly until it fails, and once it fails, it never responds to operation invocations. An object subject to omission failures may fail to respond to the invocations of an arbitrary subset of processes, but continue to respond to the invocations of the remaining processes (forever).>
Prasad Jayanti, Tushar Deepak Chandra, Sam Toueg
FOCS3
1992 The Weakest Failure Detector for Solving Consensus
abstract
We determine what information about failures is necessary and sufficient to solve Consensus in asynchronous distributed systems subject to crash failures.In [CT91], we proved that OVV, a failure
Tushar Deepak Chandra, Vassos Hadzilacos, Sam Toueg
PODC3
1991 Unreliable Failure Detectors for Asynchronous Systems (Preliminary Version)
Tushar Deepak Chandra, Sam Toueg
PODC2
1991 Inconsistency and Contamination (Preliminary Version)
abstract
Article Free Access Share on Inconsistency and contamination (preliminary version) Authors: Ajei Gopal Department of Computer Science, Cornell University, Ithaca, New York Department of Computer Science, Cornell University, Ithaca, New YorkView Profile , Sam Toueg Department of Computer Science, Cornell University, Ithaca, New York Department of Computer Science, Cornell University, Ithaca, New YorkView Profile Authors Info & Claims PODC '91: Proceedings of the tenth annual ACM symposium on Principles of distributed computingJuly 1991 Pages 257–272https://doi.org/10.1145/112600.112622Published:01 July 1991Publication History 12citation189DownloadsMetricsTotal Citations12Total Downloads189Last 12 Months3Last 6 weeks1 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteeReaderPDF
Ajei S. Gopal, Sam Toueg
PODC2
1991 The Causal Ordering Abstraction and a Simple Way to Implement it
Michel Raynal, André Schiper, Sam Toueg
Inf. Process. Lett.3
1990 Early-Delivery Atomic Broadcast
abstract
Article Free Access Share on Early-delivery atomic broadcast Authors: Ajei Gopal Cornell University Cornell UniversityView Profile , Ray Strong IBM ARC IBM ARCView Profile , Sam Toueg Cornell University Cornell UniversityView Profile , Flaviu Cristian IBM ARC IBM ARCView Profile Authors Info & Claims PODC '90: Proceedings of the ninth annual ACM symposium on Principles of distributed computingAugust 1990 Pages 297–309https://doi.org/10.1145/93385.93430Published:01 August 1990Publication History 16citation324DownloadsMetricsTotal Citations16Total Downloads324Last 12 Months9Last 6 weeks2 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my Alerts New Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteeReaderPDF
Ajei S. Gopal, Ray Strong, Sam Toueg, Flaviu Cristian
PODC3
1989 The Group Paradigm for Concurrency Control Protocols
abstract
The authors propose a paradigm for developing, describing, and proving the correctness of concurrency control protocols for replicated databases in the presence of failures or communication restrictions. The approach used is to hierarchically divide the problem of achieving one-copy serializability by introducing the notion of a 'group' that is a higher level of abstraction than transactions. Instead of dealing with the overall problem, the paradigm breaks it into two simpler ones: (1) a local policy for each group that ensures a total order of all transactions in that group; and (2) a global policy that ensures a correct serialization of all groups. The paradigm is used to demonstrate the similarities between several concurrency control protocols by comparing the way they achieve correctness.>
Amr El Abbadi, Sam Toueg
IEEE Trans. Knowl. Data Eng.2
1989 Maintaining Availability in Partitioned Replicated Databases
abstract
In a replicated database, a data item may have copies residing on several sites. A replica control protocol is necessary to ensure that data items with several copies behave as if they consist of a single copy, as far as users can tell. We describe a new replica control protocol that allows the accessing of data in spite of site failures and network partitioning. This protocol provides the database designer with a large degree of flexibility in deciding the degree of data availability, as well as the cost of accessing data.
Amr El Abbadi, Sam Toueg
ACM Trans. Database Syst.2
1988 Automatically Increasing the Fault-Tolerance of Distributed Systems
abstract
The design of fault-tolerant distributed systems is a costly and diflicult task.Its cost and difficulty increase dramatically with the severity of failures that a system must tolerate.We seek to simplify this task by developing methods to automatically translate protocols tolerant of "benign" failures to ones tolerant of more "severe" failures.This paper describes two new translation mechanisms for qr~hronous systems; one translates programs tolerant of crash failures into programs tolerant of general omission failures, and the other translates from gene& omiesion failures to arbitrary failures.Together these can be used to translate any program tolerant of the most benign failures to a program tolerant of the most severe.
Gil Neiger, Sam Toueg
PODC2
1988 The Group Paradigm for Concurrency Control Protocols
abstract
We propose a paradigm for developing, describing and proving the correctness of concurrency control protocols for replicated databases in the presence of failures or communication restrictions. Our approach is to hierarchically divide the problem of achieving one-copy serializability by introducing the notion of a “group” that is a higher level of abstraction than transactions. Instead of dealing with the overall problem of serializing all transactions, our paradigm divides the problem into two simpler ones. (1) A local policy for each group that ensures a total order of all transactions in that group. (2) A global policy that ensures a correct serialization of all groups. We use the paradigm to demonstrate the similarities between several concurrency control protocols by comparing the way they achieve correctness.
Amr El Abbadi, Sam Toueg
SIGMOD Conference2
1988 Effects of Message Loss on the Termination of Distributed Protocols
Richard Koo, Sam Toueg
Inf. Process. Lett.2
1987 Substituting for Real Time and Common Knowledge in Asynchronous Distributed Systems
abstract
We study time and knowledge in reliable distributed systems with asynchronous communication.We first describe an extension of Lamport's logical clocks that can be used as if they were perfectly synchronized real-time clocks in the solution of a large class of problems that we formally characterize.For this same class of problems, we also propose a broadcast primitive that can be used as if it achieves common knowledge.Our logical clocks and broadcast primitive are tools that considerably simplify the design of distributed algorithms: one can now design and prove them correct with the assumption that processors have access to real-time clocks and the ability to achieve common knowledge.The latter can be used to implement the abstraction of shared memory.Extensions to more synchronous systems are considered.
Gil Neiger, Sam Toueg
PODC2
1987 Distributed Deadlock Detection
Gabriel Bracha, Sam Toueg
Distributed Comput.2
1987 Simulating Authenticated Broadcasts to Derive Simple Fault-Tolerant Algorithms
T. K. Srikanth, Sam Toueg
Distributed Comput.2
1987 Optimal clock synchronization
abstract
We present a simple, efficient, and unified solution to the problems of synchronizing, initializing, and integrating clocks for systems with different types of failures: crash, omission, and arbitrary failures with and without message authentication. This is the first known solution that achieves optimal accuracy—the accuracy of synchronized clocks (with respect to real time) is as good as that specified for the underlying hardware clocks. The solution is also optimal with respect to the number of faulty processes that can be tolerated to achieve this accuracy.
T. K. Srikanth, Sam Toueg
J. ACM2
1987 Fast Distributed Agreement
abstract
We describe a Byzantine Agreement algorithm, with early stopping, for systems with arbitrary process failures. The algorithm presented is simpler and more efficient than those previously known. It was derived using a broadcast primitive that provides properties of message authentication and thus restricts the disruptive behavior of faulty processes. This primitive is a general tool for deriving fault-tolerant algorithms in the presence of arbitrary failures.
Sam Toueg, Kenneth J. Perry, T. K. Srikanth
SIAM J. Comput.1
1987 Checkpointing and Rollback-Recovery for Distributed Systems
abstract
We consider the problem of bringing a distributed system to a consistent state after transient failures. We address the two components of this problem by describing a distributed algorithm to create consistent checkpoints, as well as a rollback-recovery algorithm to recover the system to a consistent state. In contrast to previous algorithms, they tolerate failures that occur during their executions. Furthermore, when a process takes a checkpoint, a minimal number of additional processes are forced to take checkpoints. Similarly, when a process rolls back and restarts after a failure, a minimal number of additional processes are forced to roll back with it. Our algorithms require each process to store at most two checkpoints in stable storage. This storage requirement is shown to be minimal under general assumptions.
Richard Koo, Sam Toueg
IEEE Trans. Software Eng.2
1986 Availability in Partitioned Replicated Databases
abstract
Article Free Access Share on Availability in partitioned replicated databases Authors: Amr El Abbadi View Profile , Sam Toueg View Profile Authors Info & Claims PODS '86: Proceedings of the fifth ACM SIGACT-SIGMOD symposium on Principles of database systemsJune 1985 Pages 240–251https://doi.org/10.1145/6012.15418Published:01 June 1985Publication History 51citation369DownloadsMetricsTotal Citations51Total Downloads369Last 12 Months14Last 6 weeks4 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteeReaderPDF
Amr El Abbadi, Sam Toueg
PODS2
1986 State Machines and Assertions: An Integrated Approach to Modeling and Verification of Distributed Systems
Thomas A. Joseph, Thomas Räuchle, Sam Toueg
Sci. Comput. Program.3
1986 Distributed Agreement in the Presence of Processor and Communication Faults
abstract
A model of distributed computation is proposed in which processes may fail by not sending or receiving the message specified by a protocol. The solution to the Byzantine generals problem for this model is presented. The algorithm exhibits early stopping under conditions of less than maximum failure and is as efficient as the algorithm developed for the more restrictive crash-fault model in terms of time, message, and bit complexity. The authors show extant models to underestimate resiliency when faults in the communication medium are considered; the model outlined here is more accurate in this regard.
Kenneth J. Perry, Sam Toueg
IEEE Trans. Software Eng.2
1985 Optimal Clock Synchronization
abstract
We present a simple, efficient and unified solution to the problems of synchronizing clocks, initializing these clocks, and integrating new clocks, for both authenticated and nonauthenticated systems with arbitrary failures.This is the first known solution that achieves optimal accuracy, i.e., the accuracy of synchronized clocks (with respect to real time) is as good as that specified for the underlying hardware clocks.The algorithms presented are also optimal with respect to the number of faulty processes that can be tolerated to achieve this accuracy.....
T. K. Srikanth, Sam Toueg
PODC2
1985 Fast Distributed Agreement (Preliminary Version)
abstract
We describe a non-authenticated Byzantine Generals algorithm, with early stopping, for systems with arbitrary process failures.The algorithm presented is simpler, terminates earlier and has a lower communication complexity than those previously known.Surprisingly, the earlystopping algorithm is as efficient as previously proposed algorithms that do not exhibit the early-stopping property.It was derived using a broadcast primitive that simulates authentication and thus restricts the visible failure behavior of faulty processes.This primitive is a general tool for deriving fault-tolerant algorithms in the presence of arbitrary failures.
Sam Toueg, Kenneth J. Perry, T. K. Srikanth
PODC1
1985 Exposure to Deadlock for Communicating Processes is Hard to Detect
Thomas Räuchle, Sam Toueg
Inf. Process. Lett.2
1985 Asynchronous Consensus and Broadcast Protocols
abstract
A consensus protocol enables a system of n asynchronous processes, some of which are faulty, to reach agreement. There are two kinds of faulty processes: fail-stop processes that can only die and malicious processes that can also send false messages. The class of asynchronous systems with fair schedulers is defined, and consensus protocols that terminate with probability 1 for these systems are investigated. With fail-stop processes, it is shown that ⌈( n + 1)/2⌉ correct processes are necessary and sufficient to reach agreement. In the malicious case, it is shown that ⌈(2 n + 1)/3⌉ correct processes are necessary and sufficient to reach agreement. This is contrasted with an earlier result, stating that there is no consensus protocol for the fail-stop case that always terminates within a bounded number of steps, even if only one process can fail. The possibility of reliable broadcast (Byzantine Agreement) in asynchronous systems is also investigated. Asynchronous Byzantine Agreement is defined, and it is shown that ⌈(2 n + 1)/3⌉ correct processes are necessary and sufficient to achieve it.
Gabriel Bracha, Sam Toueg
J. ACM2
1984 A Distributed Algorithm for Generalized Deadlock Detection
abstract
An efficient distributed algorithm to detect deadlocks in distributed and dynamically changing systems is presented. In our model, processes can request any $N$ available resources from a pool of size $M$. This is a generalization of the well-known AND-OR request model. The algorithm is incrementally derived and proven correct. Its communication, computational, and space complexity compares favorably to those of previously known distributed AND-OR deadlock detection algorithms.
Gabriel Bracha, Sam Toueg
PODC2
1984 Randomized Byzantine Agreements
abstract
Randomized algorithms for reaching Byzantine Agreement were recently proposed in [Rabi83]. With these algorithms, agreement is reached within an expected number of phases that is a small constant independent of the number of processes n and the number of faulty processes t. The algorithms in [Rabi83] tolerate up to [(n-1)/10] faulty processes in asynchronous systems, and up to [(n-1)/4] faulty processes in synchronous systems. In this paper, using the same computation model as in [Rabi83], we describe algorithms that overcome up to [(n-1)/3] faulty processes in asynchronous systems, and up to [(n-1)/2] faulty processes in synchronous systems. With both proposed algorithms, agreement is reached within an expected number of phases that is a small constant independent of n and t, but the communication complexity is higher than in [Rabi83]. It is also shown that no Byzantine Agreement algorithm can overcome more than [(n-1)/3] faulty processes in asynchronous authenticated systems, and hence the asynchronous algorithm proposed here is optimal in this respect.
Sam Toueg
PODC1
1984 On the Optimum Checkpoint Selection Problem
abstract
We consider a model of computation consisting of a sequence of n tasks. In the absence of failures, each task i has a known completion time $t_i $ . Checkpoints can be placed between any two consecutive tasks. At a checkpoint, the state of the computation is saved on a reliable storage medium. Establishing a checkpoint immediately before task i is known to cost $s_i $. This is the time spent in saving the state of the computation. When a failure is detected, the computation is restarted at the most recent checkpoint. Restarting the computation at checkpoint i requires restoring the state to the previously saved value. The time necessary for this action is given by $r_i $. We derive an $O(n^3 )$ algorithm to select out of the $n - 1$ potential checkpoint locations those that result in the smallest expected time to complete all the tasks. An $O(n^2 )$ algorithm is described for the reasonable case where $s_i > s_j $ implies $r_i \geqslant r_j $. These algorithms are applied to two models of failure. In the first one, each task i has a given probability $p_i $ of completing without a failure, i.e., in time $t_i $. Furthermore, failures occur independently and are detected at the end of the task during which they occur. The second model admits a continuous time failure mode where the failure intervals are independent and identically distributed random variables drawn from any given distribution. In this model, failures are detected immediately. In both models, the algorithm also gives the expected value of the overall completion time and we show how to derive all the other moments.
Sam Toueg, Özalp Babaoglu
SIAM J. Comput.1
1983 Resilient Consensus Protocols
abstract
A consensus protocol enables a system of n asynchronous processes, some of which are faulty, to reach agreement. There are two kinds of faulty processes: fail-stop processes can only die, malicious processes can also send false messages. We investigate consensus protocols that terminate within finite time with probability 1 under certain assumptions on the behavior of the system. With fail-stop processes, we show that [(n + 1)/2] correct processes are necessary and sufficient to reach agreement. In the malicious case, we show that [(2n + 1)/3] correct processes are necessary and sufficient to reach agreement. This is contrasted with a recent result that there is no consensus protocol for the fail-stop case that always terminates within a bounded number of steps, even if only one process can fail.
Gabriel Bracha, Sam Toueg
PODC2
1982 The Complexity of Optimal Addressing in Radio Networks
abstract
We consider the complexity of finding optimal fixed- or variable-length unambiguous address codes for the nodes of a packet radio network. For fixed-length codes this problem is proved to be NP-complete, and its complexity for variable-length codes is still unknown. Some suboptimal heuristic algorithms are proposed.
Sam Toueg, Kenneth Steiglitz
IEEE Trans. Commun.1
1981 An all-pairs shortest-path distributed algorithm
Sam Toueg
Perform. Evaluation1
1981 Some Complexity Results in the Design of Deadlock-Free Packet Switching Networks
abstract
Deadlocks are very serious system failures and have been observed in existing packet switching networks (PSN’s). Several problems related to the design of deadlock-free PSN’s are investigated here. Polynomial-time algorithms are given for some of these problems, but most of them are shown to be NP-complete or NP-hard, and therefore polynomial-time algorithms are not likely to be found.
Sam Toueg, Kenneth Steiglitz
SIAM J. Comput.1
1981 Deadlock-Free Packet Switching Networks
abstract
Deadlock states have been observed in existing computer networks, emphasizing the need for carefully designed flow control procedures (controllers) to avoid deadlocks. Such a deadlock-free controller is readily found if we allow it global information about the overall network state. Generally, this assumption is not realistic, and we must resort to deadlock-free local controllers using only packet and node information. We present here several types of such controllers, we study their relationship and give a proof of their optimality with respect to deadlock-free controllers using the same set of local parameters.
Sam Toueg, Jeffrey D. Ullman
SIAM J. Comput.1
1980 Deadlock- and Livelock-Free Packet Switching Networks
abstract
A controller for a packet switching network is an algorithm to control the flow of packets through the network. A local controller is a controller executed independently by each node in the network, using only local information available to these nodes. A controller is deadlock- and livelock-free if it guarantees that every packet in the network reaches its destination within a finite amount of time. We present a local controller which is proved to be deadlock- and livelock-free.
Sam Toueg
STOC1
1979 Deadlock-Free Packet Switching Networks
abstract
Deadlock is one of the most serious system failures that can occur in a computer system or a network. Deadlock states have been observed in existing computer networks emphasizing the need for carefully designed flow control procedures (controllers) to avoid deadlocks. Such a deadlock-free controller is readily found if we allow it global information about the overall network state. Generally, this assumption is not realistic, and we must resort to deadlock free local controllers using only packet and node information. We present here several types of such controllers, we study their relationship and give a proof of their optimality with respect to deadlock free controllers using the same set of local parameters.
Sam Toueg, Jeffrey D. Ullman
STOC1
1979 The Design of Small-Diameter Networks by Local Search
abstract
A local search algorithm for the design of small-diameter networks is presented for both directed and undirected regular graphs. In all cases the resulting graphs are at least as good as any previously known, in the sense that they have at least as small a diameter and average shortest distance for a given number of nodes and degrees.
Sam Toueg, Kenneth Steiglitz
IEEE Trans. Computers1