Michel Raynal

dblp:r/MichelRaynal · DBLP profile ↗
← Back
392ranked-venue papers
55as first author
42since 2021 · last 2026
0000-0002-3355-8719ORCID · verified

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

Systems, architecture and hardware · 156 · 15 first-author · 8 since 2021Theory of computation · 89 · 15 first-author · 17 since 2021Security and privacy · 56 · 7 first-author · 7 since 2021Databases, data management, data science and information retrieval · 21 · 6 first-author · 1 since 2021Software engineering, systems software and programming languages · 18 · 4 first-authorApplied, interdisciplinary, general and emerging computing · 9Computer networks · 6 · 1 first-author · 1 since 2021
YearPublicationVenuePosition
2026 Byzantine-tolerant distributed grow-only sets: specification and applications
abstract
In order to formalize Distributed Ledger Technologies and their interconnections, recent research has introduced the concept of a Distributed Ledger Object (denoted $$\mathcal {O}^L$$ ), a concurrent abstraction that maintains a totally ordered sequence of records, capturing the essence of blockchains and distributed ledgers. In this work, we introduce the Distributed Grow-only Set object (denoted $$\mathcal {O}^{GS}$$ ), a novel abstraction that, unlike the $$\mathcal {O}^L$$ , maintains an immutable set of records by supporting only Add and Get operations. This object is inspired by the Grow-only Set (G-Set) a well-known Conflict-free Replicated Data Type (CRDT). We formally define the $$\mathcal {O}^{GS}$$ and present a Byzantine-tolerant, consensus-free implementation (denoted as $$\mathcal {O}^{GS}_B$$ ) that ensures eventual consistency. Building on this implementation, we propose consensus-free algorithmic solutions to two fundamental problems: the Atomic Appends problem, which concerns atomically appending multiple records to distinct ledgers, and the Atomic Adds problem, its counterpart in the context of G-Sets. Additionally, we show how the $$\mathcal {O}^{GS}_B$$ can be leveraged to construct a consensus-free, Single-Writer Byzantine-tolerant $$\mathcal {O}^L$$ . We argue that the applicability of the $$\mathcal {O}^{GS}_B$$ extends well beyond these specific use cases, offering a lightweight and efficient foundation for a variety of distributed applications.
Vicent Cholvi, Antonio Fernández 0001, Chryssis Georgiou, Nicolas C. Nicolaou, Michel Raynal, Antonio Russo 0004
Distributed Comput.5
2025 Contention-Aware Cooperation
abstract
As shown by Reliable Broadcast and Consensus, cooperation among a set of independent computing entities (sequential processes) is a central issue in distributed computing. Considering $n$-process asynchronous message-passing systems where some processes can be Byzantine, this paper introduces a new cooperation abstraction denoted Context-Adaptive Cooperation (CAC). While Reliable Broadcast is a one-to-$n$ cooperation abstraction and Consensus is an $n$-to-$n$ cooperation abstraction, CAC is a $d$-to-$n$ cooperation abstraction where the parameter $d$ ($1\leq d\leq n$) depends on the run and remains unknown to the processes. Moreover, the correct processes accept the same set of $\ell$ pairs $\langle v,i\rangle$ ($v$ is the value proposed by $p_i$) from the $d$ proposer processes, where $1 \leq \ell \leq d$ and, as $d$, $\ell$ remains unknown to the processes (except in specific cases). Those $\ell$ values are accepted one at a time in different orders at each process. Furthermore, CAC provides the processes with an imperfect oracle that gives information about the values that they may accept in the future. In a very interesting way, the CAC abstraction is particularly efficient in favorable circumstances. To illustrate its practical use, the paper describes in detail two applications that benefit from the abstraction: a fast consensus implementation under low contention (named Cascading Consensus), and a novel naming problem.
Timothé Albouy, Davide Frey, Mathieu Gestin, Michel Raynal, François Taïani
OPODIS4
2025 Solving Tasks with Fewer Registers Than Processes
abstract
This paper studies distributed-computing tasks through the lens of space complexity in the read/write wait-free model, defined as the number of multi-reader-multi-writer atomic read/write registers needed to solve a task using a wait-free algorithm. Surprisingly, even though the read/write wait-free model is at the foundation of distributed computing, previous work on space complexity has focused on synchronization primitives stronger than read/write registers or on weaker progress conditions. The paper reveals that the read/write wait-free model offers a rich space-complexity landscape: (1) assuming non-anonymous processes, it shows that there is an infinite hierarchy of tasks of increasing space complexity; (2) it shows that space complexity separates anonymous from non-anonymous memory; (3) regardless of process or register anonymity, it exhibits a task of space complexity two, which is the minimal non-trivial space complexity; (4) finally, it shows that subcases of the adopt-commit task have different space complexity in non-anonymous memory under bounded wait-freedom.
Eli Gafni, Giuliano Losa, Michel Raynal, Gadi Taubenfeld
OPODIS3
2025 Brief Announcement: Stranger-Free Tasks
abstract
Delporte-Gallet et al. show that, in a system of n processes, it is both necessary and sufficient to use n multi-writer multi-reader (MWMR) registers, that are not pre-allocated, to emulate with non-blocking progress n single-writer multi-reader (SWMR) registers that are uniquely pre-allocated. They conclude with the significant result that n MWMR registers are sufficient to solve any task solvable read-write wait-free. However, they mistakenly claim—likely inadvertently—that n MWMR registers are also necessary to solve any task solvable read-write wait-free (a counterexample is the splitter task, which is solvable for any number of processes with just 2 MWMR registers).
Eli Gafni, Giuliano Losa, Michel Raynal, Gadi Taubenfeld
PODC3
2025 Self-stabilizing multivalued consensus in the presence of Byzantine faults and asynchrony
abstract
Consensus, abstracting myriad problems in which processes must agree on a single value, is one of the most celebrated problems of fault-tolerant distributed computing. Consensus applications include fundamental services for the Cloud and Blockchain environments, and in such challenging environments, malicious behaviors are often modeled as adversarial Byzantine faults. At OPODIS 2010, Mostéfaoui and Raynal (in short, MR) presented a Byzantine-tolerant solution to consensus in which the decided value cannot be proposed only by Byzantine processes. MR has optimal resilience coping with up to t < n / 3 Byzantine nodes over n processes. MR provides this multivalued consensus object (which accepts proposals taken from a finite set of values), assuming the availability of a single binary consensus object (which accepts proposals taken from the set { 0 , 1 } ). This work, which focuses on multivalued consensus, aims to design an even more robust solution than MR. Our proposal expands MR's fault-model with self-stabilization, a vigorous notion of fault-tolerance. In addition to tolerating Byzantine, self-stabilizing systems can automatically recover after arbitrary transient-faults occur. These faults represent any violation of the assumptions according to which the system was designed to operate (provided that the algorithm code remains intact). To the best of our knowledge, we propose the first self-stabilizing solution for multivalued consensus for asynchronous message-passing systems prone to Byzantine failures. Our solution has an O ( t ) stabilization time from arbitrary transient faults.
Romaric Duvignau, Michel Raynal, Elad Michael Schiller
Theor. Comput. Sci.2
2024 Near-Optimal Communication Byzantine Reliable Broadcast Under a Message Adversary
abstract
We address the problem of Reliable Broadcast in asynchronous message-passing systems with n nodes, of which up to t are malicious (faulty), in addition to a message adversary that can drop some of the messages sent by correct (non-faulty) nodes. We present a Message-Adversary-Tolerant Byzantine Reliable Broadcast (MBRB) algorithm that communicates O(|m|+nκ) bits per node, where |m| represents the length of the application message and κ = Ω(log n) is a security parameter. This communication complexity is optimal up to the parameter κ. This significantly improves upon the state-of-the-art MBRB solution (Albouy, Frey, Raynal, and Taïani, TCS 2023), which incurs communication of O(n|m|+n²κ) bits per node. Our solution sends at most 4n² messages overall, which is asymptotically optimal. Reduced communication is achieved by employing coding techniques that replace the need for all nodes to (re-)broadcast the entire application message m. Instead, nodes forward authenticated fragments of the encoding of m using an erasure-correcting code. Under the cryptographic assumptions of threshold signatures and vector commitments, and assuming n > 3t+2d, where the adversary drops at most d messages per broadcast, our algorithm allows at least 𝓁 = n - t - (1 + ε)d (for any arbitrarily low ε > 0) correct nodes to reconstruct m, despite missing fragments caused by the malicious nodes and the message adversary.
Timothé Albouy, Davide Frey, Ran Gelles, Carmit Hazay, Michel Raynal, Elad Michael Schiller, François Taïani, Vassilis Zikas
OPODIS5
2024 Better Sooner Rather Than Later
Anaïs Durand, Michel Raynal, Gadi Taubenfeld
SIROCCO2
2024 On Distributed Computing: A View, Physical Versus Logical Objects, and a Look at Fully Anonymous Systems
Michel Raynal
SSS1
2024 Brief Announcement: Towards Optimal Communication Byzantine Reliable Broadcast Under a Message Adversary
abstract
International audience
Timothé Albouy, Davide Frey, Ran Gelles, Carmit Hazay, Michel Raynal, Elad Michael Schiller, François Taïani, Vassilis Zikas
DISC5
2024 Good-case early-stopping latency of synchronous byzantine reliable broadcast: the deterministic case
Timothé Albouy, Davide Frey, Michel Raynal, François Taïani
Distributed Comput.3
2024 Process-commutative distributed objects: From cryptocurrencies to Byzantine-Fault-Tolerant CRDTs
Davide Frey, Lucie Guillou, Michel Raynal, François Taïani
Theor. Comput. Sci.3
2024 Self-stabilizing indulgent zero-degrading binary consensus
abstract
Guerraoui proposed an indulgent solution for the binary consensus problem. Namely, he showed that an arbitrary behavior of the failure detector never violates safety requirements even if it compromises liveness. Consensus implementations are often used in a repeated manner. Dutta and Guerraoui proposed a zero-degrading solution, i.e., during system runs in which the failure detector behaves perfectly, a node failure during one consensus instance has no impact on the performance of future instances. Our study, which focuses on indulgent zero-degrading binary consensus, aims at the design of an even more robust communication abstraction. We do so through the lenses of self-stabilization—a very strong notion of fault-tolerance. In addition to node and communication failures, self-stabilizing algorithms can recover after the occurrence of arbitrary transient faults; these faults represent any violation of the assumptions according to which the system was designed to operate (as long as the algorithm code stays intact). This work proposes the first, to the best of our knowledge, self-stabilizing algorithm for indulgent zero-degrading binary consensus for time-free message-passing systems prone to detectable process failures. The proposed algorithm recovers within a finite time after the occurrence of the last arbitrary transient fault. Since the proposed solution uses an Ω failure detector, we also present the first, to the best of our knowledge, self-stabilizing asynchronous Ω failure detector, which is a variation on the one by Mostéfaoui, Mourgaya, and Raynal.
Oskar Lundström, Michel Raynal, Elad Michael Schiller
Theor. Comput. Sci.2
2024 Self-stabilizing multivalued consensus in asynchronous crash-prone systems
abstract
The multivalued consensus problem is a fundamental issue in fault-tolerant distributed computing. It encompasses a wide range of agreement problems where processes must unanimously decide on a specific value v ∈ V , with | V | ≥ 2 . Existing solutions that handle process crash failures simplify the multivalued consensus problem by reducing it to the binary consensus problem. Examples of such solutions include Mostéfaoui-Raynal-Tronel [IPL 2000] and Zhang-Chen [IPL 2009]. In this work, we aim to design an even more reliable solution by leveraging the concept of self-stabilization , which provides a strong form of fault tolerance. Self-stabilizing algorithms can recover from transient faults, which represent any deviation from the system's intended behavior (as long as the algorithm code remains intact) in addition to process and communication failures. To the best of our knowledge, this work presents the first self-stabilizing algorithm for multivalued consensus in asynchronous message-passing systems susceptible to process failures and transient faults. Our solution uses, at most, n concurrent invocations of binary consensus. This is another way we advance state-of-the-art solutions compared to previous non-self-stabilizing ones. For example, Mostéfaoui-Raynal-Tronel's solution requires an unbounded number of sequential invocations of binary consensus.
Oskar Lundström, Michel Raynal, Elad Michael Schiller
Theor. Comput. Sci.2
2023 Brief Announcement: The MBroadcast Abstraction
abstract
This short article presents a new communication abstraction denoted Mutual Broadcast (in short MBroadcast). It provides each pair of processes with the following property (called mutual ordering): for any pair of processes p and p′, if p broadcasts a message m and p′ broadcasts a message m′, it is not possible for p to deliver first (its message) m and then m′ while p′ delivers first (its message) m′ and then m. The computability power of this broadcast abstraction is the same as the one of an atomic read/write register. Interestingly, it constitutes the first characterization of RW registers in terms of (binary) message patterns.
Mathilde Déprés, Achour Mostéfaoui, Matthieu Perrin, Michel Raynal
PODC4
2023 About Informatics, Distributed Computing, and Our Job: A Personal View
Michel Raynal
SIROCCO1
2023 Self-stabilizing Byzantine-Tolerant Recycling
Chryssis Georgiou, Michel Raynal, Elad Michael Schiller
SSS2
2023 Send/Receive Patterns Versus Read/Write Patterns in Crash-Prone Asynchronous Distributed Systems
Mathilde Déprés, Achour Mostéfaoui, Matthieu Perrin, Michel Raynal
DISC4
2023 The Synchronization Power (Consensus Number) of Access-Control Objects: the Case of AllowList and DenyList
abstract
This article studies the synchronization power of AllowList and DenyList objects under the lens provided by Herlihy’s consensus hierarchy. It specifies AllowList and DenyList as distributed objects and shows that, while they can both be seen as specializations of a more general object type, they inherently have different synchronization power. While the AllowList object does not require synchronization between participating processes, a DenyList object requires processes to reach consensus on a specific set of processes. These results are then applied to a more global analysis of anonymity-preserving systems that use AllowList and DenyList objects. First, a blind-signature-based e-voting is presented. Second, DenyList and AllowList objects are used to determine the consensus number of a specific decentralized key management system. Third, an anonymous money transfer algorithm using the association of AllowList and DenyList objects is presented. Finally, this analysis is used to study the properties of these application, and to highlight efficiency gains that they can achieve in message passing environment.
Davide Frey, Mathieu Gestin, Michel Raynal
DISC3
2023 Set-Linearizable Implementations from Read/Write Operations: Sets, Fetch &Increment, Stacks and Queues with Multiplicity
Armando Castañeda, Sergio Rajsbaum, Michel Raynal
Distributed Comput.3
2023 Corrigendum to "Mutual exclusion in fully anonymous shared memory systems" [Inf. Process. Lett. 158 (2020) 105938]
Michel Raynal, Gadi Taubenfeld
Inf. Process. Lett.1
2023 Atomic Appends in Asynchronous Byzantine Distributed Ledgers
Vicent Cholvi, Antonio Fernández 0001, Chryssis Georgiou, Nicolas C. Nicolaou, Michel Raynal, Antonio Russo 0004
J. Parallel Distributed Comput.5
2023 Asynchronous Byzantine reliable broadcast with a message adversary
Timothé Albouy, Davide Frey, Michel Raynal, François Taïani
Theor. Comput. Sci.3
2023 Optimal algorithms for synchronous Byzantine k-set agreement
Carole Delporte-Gallet, Hugues Fauconnier, Michel Raynal, Mouna Safir
Theor. Comput. Sci.3
2023 Reaching agreement in the presence of contention-related crash failures
Anaïs Durand, Michel Raynal, Gadi Taubenfeld
Theor. Comput. Sci.2
2023 Self-stabilizing Byzantine fault-tolerant repeated reliable broadcast
abstract
We study a well-known communication abstraction called Byzantine Reliable Broadcast (BRB). This abstraction is central in the design and implementation of fault-tolerant distributed systems, as many fault-tolerant distributed applications require communication with provable guarantees on message deliveries. Our study focuses on fault-tolerant implementations for message-passing systems that are prone to process-failures, such as crashes and malicious behavior. At PODC 1983, Bracha and Toueg, in short, BT, solved the BRB problem. BT has optimal resilience since it can deal with t
Romaric Duvignau, Michel Raynal, Elad Michael Schiller
Theor. Comput. Sci.2
2023 DMCSC: a fully distributed multi-coloring approach for scalable communication in synchronous broadcast networks
Youcef Imine, Hicham Lakhlef, Michel Raynal, François Taïani
J. Supercomput.3
2022 Message from the General Chairs: MSN 2022
abstract
Welcome to the 18th International Conference on Mobility, Sensing and Networking (MSN 2022), held in Guangzhou, China, on December 14-16, 2022. MSN provides an open forum for academic researchers and industry practitioners to present their research progresses, exchange new ideas, and identify future directions in the fields of mobility, sensing and networking.
Michel Raynal, Jie Wu 0001
MSN1
2022 A Modular Approach to Construct Signature-Free BRB Algorithms Under a Message Adversary
Timothé Albouy, Davide Frey, Michel Raynal, François Taïani
OPODIS3
2022 Election in Fully Anonymous Shared Memory Systems: Tight Space Bounds and Algorithms
Damien Imbs, Michel Raynal, Gadi Taubenfeld
SIROCCO2
2022 Optimal Algorithms for Synchronous Byzantine k-Set Agreement
Carole Delporte-Gallet, Hugues Fauconnier, Michel Raynal, Mouna Safir
SSS3
2022 Reaching Consensus in the Presence of Contention-Related Crash Failures
Anaïs Durand, Michel Raynal, Gadi Taubenfeld
SSS2
2022 Self-stabilizing Byzantine Fault-Tolerant Repeated Reliable Broadcast
Romaric Duvignau, Michel Raynal, Elad Michael Schiller
SSS2
2022 Brief Announcement: Self-stabilizing Total-Order Broadcast
Oskar Lundström, Michel Raynal, Elad Michael Schiller
SSS2
2022 Good-Case Early-Stopping Latency of Synchronous Byzantine Reliable Broadcast: The Deterministic Case
abstract
This paper considers the good-case latency of Byzantine Reliable Broadcast (BRB), i.e., the time taken by correct processes to deliver a message when the initial sender is correct, and an essential property for practical distributed systems. Although significant strides have been made in recent years on this question, progress has mainly focused on either asynchronous or randomized algorithms. By contrast, the good-case latency of deterministic synchronous BRB under a majority of Byzantine faults has been little studied. In particular, it was not known whether a good-case latency below the worst-case bound of t+1 rounds could be obtained under a Byzantine majority. In this work, we answer this open question positively and propose a deterministic synchronous Byzantine reliable broadcast that achieves a good-case latency of max(2,t+3-c) rounds, where t is the upper bound on the number of Byzantine processes, and c the number of effectively correct processes.
Timothé Albouy, Davide Frey, Michel Raynal, François Taïani
DISC3
2022 Distributed computability: Relating k-immediate snapshot and x-set agreement
Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
Inf. Comput.4
2022 Contention-related crash failures: Definitions, agreement algorithms, and impossibility results
Anaïs Durand, Michel Raynal, Gadi Taubenfeld
Theor. Comput. Sci.2
2022 A visit to mutual exclusion in seven dates
Michel Raynal, Gadi Taubenfeld
Theor. Comput. Sci.1
2021 From Incomplete to Complete Networks in Asynchronous Byzantine Systems
Michel Raynal, Jiannong Cao 0001
AINA (1)1
2021 Byzantine-Tolerant Reliable Broadcast in the Presence of Silent Churn
Timothé Albouy, Davide Frey, Michel Raynal, François Taïani
SSS3
2021 On the weakest information on failures to solve mutual exclusion and consensus in asynchronous crash-prone read/write systems
Carole Delporte-Gallet, Hugues Fauconnier, Michel Raynal
J. Parallel Distributed Comput.3
2021 Byzantine-tolerant causal broadcast
Alex Auvolat, Davide Frey, Michel Raynal, François Taïani
Theor. Comput. Sci.3
2021 Set-constrained delivery broadcast: A communication abstraction for read/write implementable distributed objects
Damien Imbs, Achour Mostéfaoui, Matthieu Perrin, Michel Raynal
Theor. Comput. Sci.4
2020 Self-Stabilizing Set-Constrained Delivery Broadcast (extended abstract)
abstract
Fault-tolerant distributed applications require communication abstractions with provable guarantees on message deliveries. For example, Set-Constrained Delivery Broadcast (SCD-broadcast) is a communication abstraction for broadcasting messages in a manner that, if a process delivers a set of messages that includes m and later delivers a set of messages that includes m , no process delivers first a set of messages that includes m′ and later a set of messages that includes m.Imbs et al. proposed this communication abstraction and its first implementation. They have demonstrated that SCD-broadcast has the computational power of read/write registers and allows for an easy building of distributed objects such as snapshot objects and consistent counters. Imbs et al. focused on fault-tolerant implementations for asynchronous message-passing systems that are prone to process crashes. This paper aims to design an even more robust SCD-broadcast communication abstraction, namely a self-stabilizing SCD-broadcast. In addition to process and communication failures, self-stabilizing algorithms can recover after the occurrence of arbitrary transient faults; these faults represent any violation of the assumptions according to which the system was designed to operate (as long as the algorithm code stays intact).This work proposes the first self-stabilizing SCD-broadcast algorithm for asynchronous message-passing systems that are prone to process crash failures. The proposed self-stabilizing SCD-broadcast algorithm has an $\mathcal{O}(1)$ stabilization time (in terms of asynchronous cycles). The communication costs of our algorithm are similar to the ones of the non-self-stabilizing state-of-the-art. The main differences are that our proposal considers repeated gossiping of $\mathcal{O}(1)$ bits messages and deals with bounded space (which is a prerequisite for self-stabilization). We advance the state-of-the-art also by two new self-stabilizing applications: an atomic construction of snapshot objects and sequentially consistent counters.
Oskar Lundström, Michel Raynal, Elad Michael Schiller
ICDCS2
2020 Relaxed Queues and Stacks from Read/Write Operations
abstract
Considering asynchronous shared memory systems in which any number of processes may crash, this work identifies and formally defines relaxations of queues and stacks that can be non-blocking or wait-free while being implemented using only read/write operations. Set-linearizability and Interval-linearizability are used to specify the relaxations formally, and precisely identify the subset of executions which preserve the original sequential behavior. The relaxations allow for an item to be returned more than once by different operations, but only in case of concurrency; we call such a property multiplicity. The stack implementation is wait-free, while the queue implementation is non-blocking. Interval-linearizability is used to describe a queue with multiplicity, with the additional relaxation that a dequeue operation can return weak-empty, which means that the queue might be empty. We present a read/write wait-free interval-linearizable algorithm of a concurrent queue. As far as we know, this work is the first that provides formalizations of the notions of multiplicity and weak-emptiness, which can be implemented on top of read/write registers only.
Armando Castañeda, Sergio Rajsbaum, Michel Raynal
OPODIS3
2020 From Bezout's Identity to Space-Optimal Election in Anonymous Memory Systems
abstract
An anonymous shared memory REG can be seen as an array of atomic registers such that there is no a priori agreement among the processes on the names of the registers. As an example a very same physical register can be known as REG[x] by a process p and as REG[y] (where y ≠ x) by another process q. Moreover, the register known as REG[a] by a process p and the register known as REG[b] by a process q can be the same physical register. It is assumed that each process has a unique identifier that can only be compared for equality. This article is on solving the d-election problem, in which it is required to elect at least one and at most d leaders, in such an anonymous shared memory system. We notice that the 1-election problem is the familiar leader election problem. Let n be the number of processes and m the size of the anonymous memory (number of atomic registers). The article shows that the condition gcd(m, n) ≤ d is necessary and sufficient for solving the d-election problem, where communication is through read/write or read+modify+write registers. The algorithm used to prove the sufficient condition relies on Bezout's Identity - a Diophantine equation relating numbers according to their Greatest Common Divisor. Furthermore, in the process of proving the sufficient condition, it is shown that 1-leader election can be solved using only a single read/write register (which refutes a 1989 conjecture stating that three non-anonymous registers are necessary), and that the exact d-election problem, where exactly d leaders must be elected, can be solved if and only if gcd(m, n) divides d.
Emmanuel Godard, Damien Imbs, Michel Raynal, Gadi Taubenfeld
PODC3
2020 k-Immediate Snapshot and x-Set Agreement: How Are They Related?
Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
SSS4
2020 Brief Announcement: Leader Election in the ADD Communication Model
Sergio Rajsbaum, Michel Raynal, Karla Vargas
SSS2
2020 Mutual exclusion in fully anonymous shared memory systems
Michel Raynal, Gadi Taubenfeld
Inf. Process. Lett.1
2020 Leader-based de-anonymization of an anonymous read/write memory
Emmanuel Godard, Damien Imbs, Michel Raynal, Gadi Taubenfeld
Theor. Comput. Sci.3
2020 Collisions Are Preferred: RFID-Based Stocktaking with a High Missing Rate
abstract
RFID-based stocktaking uses RFID technology to verify the presence of objects in a region e.g., a warehouse or a library, compared with an inventory list. The existing approaches for this purpose assume that the number of missing tags is small. This is not true in some cases. For example, for a handheld RFID reader, only the objects in a larger region (e.g., the warehouse) rather than in its interrogation region can be known as the inventory list, and hence many tags in the list are regarded as missing. The missing objects significantly increase the time required for stocktaking. In this paper, we propose an algorithm called CLS (Coarse-grained inventory list based stocktaking) to solve this problem. CLS enables multiple missing objects to hash to a single time slot and thus verifies them together. CLS also improves the existing approaches by utilizing more kinds of RFID collisions and reducing approximately one-fourth of the amount of data sent by the reader. Moreover, we observe that the missing rate constantly changes during the identification because some of tags are verified present or absent, which affects time efficiency; accordingly, we propose a hybrid stocktaking algorithm called DLS (Dynamic inventory list based stocktaking) to adapt to such changes for the first time. According to the results of extensive simulations, when the inventory list is 20 times that of actually present tags, the execution time of our approach is 36.3 percent that of the best existing algorithm.
Weiping Zhu 0004, Xing Meng, Xiaolei Peng, Jiannong Cao 0001, Michel Raynal
IEEE Trans. Mob. Comput.5
2019 On the Weakest Failure Detector for Read/Write-Based Mutual Exclusion
Carole Delporte-Gallet, Hugues Fauconnier, Michel Raynal
AINA3
2019 One for All and All for One: Scalable Consensus in a Hybrid Communication Model
abstract
This paper addresses consensus in an asynchronous model where the processes are partitioned into clusters. Inside each cluster, processes can communicate through a shared memory, which favors efficiency. Moreover, any pair of processes can also communicate through a message-passing communication system, which favors scalability. In such a “hybrid communication” context, the paper presents two simple binary consensus algorithms (one based on local coins, the other one based on a common coin). These algorithms are straightforward extensions of existing message-passing randomized round-based consensus algorithms. At each round, the processes of each cluster first agree on the same value (using an underlying shared memory consensus algorithm), and then use a message-passing algorithm to converge on the same decided value. The algorithms are such that, if all except one processes of a cluster crash, the surviving process acts as if all the processes of its cluster were alive (hence the motto “one for all and all for one”). As a consequence, the hybrid communication model allows us to obtain simple, efficient, and scalable fault-tolerant consensus algorithms. As an important side effect, according to the size of each cluster, consensus can be obtained even if a majority of processes crash.
Michel Raynal, Jiannong Cao 0001
ICDCS1
2019 Byzantine-Tolerant Set-Constrained Delivery Broadcast
Alex Auvolat, Michel Raynal, François Taïani
OPODIS2
2019 Optimal Memory-Anonymous Symmetric Deadlock-Free Mutual Exclusion
abstract
The notion of an anonymous shared memory, introduced by Taubenfeld in PODC 2017, considers that processes use different names for the same memory location. As an example, a location name A used by a process p and a location name B ≠ A used by another process q can correspond to the very same memory location X, and similarly for the names B used by p and A used by q which may (or may not) correspond to the same memory location Y ≠ X. In this context, the PODC paper presented a 2-process symmetric deadlock-free mutual exclusion (mutex) algorithm and a necessary condition on the size m of the anonymous memory for the existence of such an n-process algorithm. This condition states that m must be belongs to M(n) {1} where M(n)= {m: ∀ ℓ: (1) < ℓ ≤ n: gcd(ℓ,m)=1). Symmetric means here that,process identities define a specific data type which allows a process to check only if two identities are equal or not.
Zahra Aghazadeh, Damien Imbs, Michel Raynal, Gadi Taubenfeld, Philipp Woelfel
PODC3
2019 A New Insight into Local Coin-Based Randomized Consensus
abstract
This paper presents a binary randomized consensus algorithm for n-process asynchronous message-passing systems in which (1) up to tO(sqrt(n)) it is no longer possible to implement a randomized consensus algorithm ensuring a constant number of communication steps despite unfair channels.
Achour Mostéfaoui, Matthieu Perrin, Michel Raynal
PRDC3
2019 Anonymous Read/Write Memory: Leader Election and De-anonymization
Emmanuel Godard, Damien Imbs, Michel Raynal, Gadi Taubenfeld
SIROCCO3
2019 The Notion of Universality in Crash-Prone Asynchronous Message-Passing Systems: A Tutorial
abstract
The notion of a universal construction is central in computing science and technology: general solutions make life easier and the wheel has not to be reinvented each time a new problem appears. In the context of message-passing asynchronous distributed systems made up of n processes, where some of them may commit crash failures, a universal construction is an algorithm that is able to build any object defined by a sequential specification despite the adversary effect resulting from the combination of asynchrony and process crashes. The aim of this tutorial is to introduce the reader to the notion of a distributed universal construction (and universal objects these constructions rely on), and more precisely, explain what can be done, what cannot be done, and which assumptions on the environment are necessary in order objects with provably reliability properties can be built. Its aim is be a guided tour providing the reader with the basic knowledge needed to understand and master asynchronous message-passing fault-tolerant computing. Its spirit is not to be a catalog of constructions proposed so far, but an as simple as possible presentation of concepts and mechanisms that constitute the basis these universal constructions rely on.
Michel Raynal
SRDS1
2019 Brief Announcement: Fully Anonymous Shared Memory Algorithms
Michel Raynal, Gadi Taubenfeld
SSS1
2019 A speculation-friendly binary search tree
abstract
Summary We introduce the first concurrent data structure algorithm designed for speculative executions. Prior to this work, concurrent structures were mainly designed for their pessimistic (non‐speculative) accesses to have a predictable asymptotic complexity. Researchers tried to evaluate transactional memory using such structures whose prominent example is the red‐black tree library developed by Oracle Labs that is part of multiple benchmark distributions. Although well‐engineered, such structures remain badly suited for speculative accesses, whose step complexity might raise dramatically with contention. We propose a binary search tree data structure whose key novelty stems from the decoupling of update operations, ie, instead of performing an update operation in a single large transaction, it is split into one transaction that modifies the abstraction state and several other transactions that restructure the tree implementation in the background. This results in a speculation‐friendly tree (s‐tree) that outperforms previous HTM‐based and STM‐based trees by being transiently unbalanced during contention peaks and by rebalancing in quadratic time when contention disappears. In particular, the s‐tree is shown correct, reusable, and speeds up a transaction‐based travel reservation application by up to 3.5×.
Tyler Crain, Vincent Gramoli, Michel Raynal
Concurr. Comput. Pract. Exp.3
2019 Crash-tolerant causal broadcast in O(n) messages
Achour Mostéfaoui, Matthieu Perrin, Michel Raynal, Jiannong Cao 0001
Inf. Process. Lett.3
2019 Making Local Algorithms Wait-Free: the Case of Ring Coloring
Armando Castañeda, Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
Theory Comput. Syst.5
2019 Vertex Coloring with Communication Constraints in Synchronous Broadcast Networks
abstract
This paper considers distributed vertex-coloring in broadcast/receive networks suffering from conflicts and collisions. (A collision occurs when, during the same round, messages are sent to the same process by too many neighbors; a conflict occurs when a process and one of its neighbors broadcast during the same round.) More specifically, the paper focuses on multi-channel networks, in which a process may either broadcast a message to its neighbors or receive a message from at most γ of them. The paper first provides a new upper bound on the corresponding graph coloring problem (known as frugal coloring) in general graphs, proposes an exact bound for the problem in trees, and presents a deterministic, parallel, color-optimal, collision- and conflict-free distributed coloring algorithm for trees, and proves its correctness.
Hicham Lakhlef, Michel Raynal, François Taïani
IEEE Trans. Parallel Distributed Syst.2
2018 Time-Efficient RFID-Based Stocktaking with a Coarse-Grained Inventory List
abstract
RFID-based stocktaking uses RFID technology to verify the presence of objects in a region e.g., a warehouse or a library. The existing approaches for this purpose assume that an inventory list of objects in the interrogation region of an RFID reader is known. This is not true in some cases. For example, for a handheld RFID reader, only the objects in a larger region (e.g., the warehouse) rather than in its interrogation region can be known. The additional objects significantly increase the time required for stocktaking. In this paper, we propose a time-efficient stocktaking algorithm called CLS (Coarse-grained inventory list based stocktaking) to solve this problem. We transform the problem to a missing tag identification problem with a large missing rate. CLS enables multiple missing objects to hash to a single time slot and thus verifies them together. CLS also improves the existing approaches by utilizing more kinds of RFID collisions and reducing approximately one-fourth of the amount of data sent by the reader. Extensive simulations are performed and the results show CLS outperforms the best existing algorithm.
Weiping Zhu 0004, Xing Meng, Xiaolei Peng, Jiannong Cao 0001, Michel Raynal
IWQoS5
2018 A Pleasant Stroll Through the Land of Distributed Machines, Computation, and Universality
Michel Raynal, Jiannong Cao 0001
MCU1
2018 DBFT: Efficient Leaderless Byzantine Consensus and its Application to Blockchains
abstract
This paper introduces a new leaderless Byzantine consensus called the Democratic Byzantine Fault Tolerance (DBFT) for blockchains. While most blockchain consensus protocols rely on a correct leader or coordinator to terminate, our algorithm can terminate even when its coordinator is faulty. The key idea is to allow processes to complete asynchronous rounds as soon as they receive a threshold of messages, instead of having to wait for a message from a coordinator that may be slow. The resulting decentralization is particularly appealing for blockchains for two reasons: (i) each node plays a similar role in the execution of the consensus, hence making the decision inherently “democratic” (ii) decentralization avoids bottlenecks by balancing the load, making the solution scalable. DBFT is deterministic, assumes partial synchrony, is resilience optimal, time optimal and does not need signatures. We first present a simple safe binary Byzantine consensus algorithm, modify it to ensure termination, and finally present an optimized reduction from multivalue consensus to binary consensus whose fast path terminates in 4 message delays.
Tyler Crain, Vincent Gramoli, Mikel Larrea, Michel Raynal
NCA4
2018 Set Agreement and Renaming in the Presence of Contention-Related Crash Failures
Anaïs Durand, Michel Raynal, Gadi Taubenfeld
SSS2
2018 Bee's Strategy Against Byzantines Replacing Byzantine Participants - (Extended Abstract)
Amitay Shaer, Shlomi Dolev, Silvia Bonomi, Michel Raynal, Roberto Baldoni
SSS4
2018 Agent-based broadcast protocols for wireless heterogeneous node networks
Hicham Lakhlef, Abdelmadjid Bouabdallah, Michel Raynal, Julien Bourgeois
Comput. Commun.3
2018 Anonymous obstruction-free (n, k)-set agreement with n-k+1 atomic read/write registers
Zohir Bouzid, Michel Raynal, Pierre Sutra
Distributed Comput.2
2018 Unifying Concurrent Objects and Distributed Tasks: Interval-Linearizability
abstract
Tasks and objects are two predominant ways of specifying distributed problems where processes should compute outputs based on their inputs. Roughly speaking, a task specifies, for each set of processes and each possible assignment of input values, their valid outputs. In contrast, an object is defined by a sequential specification. Also, an object can be invoked multiple times by each process, while a task is a one-shot problem. Each one requires its own implementation notion, stating when an execution satisfies the specification. For objects, linearizability is commonly used, while tasks implementation notions are less explored. The article introduces the notion of interval-sequential object, and the corresponding implementation notion of interval-linearizability , to encompass many problems that have no sequential specification as objects. It is shown that interval-sequential specifications are local , namely, one can consider interval-linearizable object implementations in isolation and compose them for free, without sacrificing interval-linearizability of the whole system. The article also introduces the notion of refined tasks and its corresponding satisfiability notion. In contrast to a task, a refined task can be invoked multiple times by each process. Also, objects that cannot be defined using tasks can be defined using refined tasks. In fact, a main result of the article is that interval-sequential objects and refined tasks have the same expressive power and both are complete in the sense that they are able to specify any prefix-closed set of well-formed executions. Interval-linearizability and refined tasks go beyond unifying objects and tasks; they shed new light on both of them. On the one hand, interval-linearizability brings to task the following benefits: an explicit operational semantics, a more precise implementation notion, a notion of state, and a locality property. On the other hand, refined tasks open new possibilities of applying topological techniques to objects.
Armando Castañeda, Sergio Rajsbaum, Michel Raynal
J. ACM3
2018 Randomized k-set agreement in crash-prone and Byzantine asynchronous systems
Achour Mostéfaoui, Moumen Hamouma, Michel Raynal
Theor. Comput. Sci.3
2018 Implementing Snapshot Objects on Top of Crash-Prone Asynchronous Message-Passing Systems
Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
IEEE Trans. Parallel Distributed Syst.4
2017 Providing Collision-Free and Conflict-Free Communication in General Synchronous Broadcast/Receive Networks
abstract
This work considers the problem of communication in dense and large scale wireless networks composed of resourcelimited nodes. In this kind of networks, a massive amount of data is becoming increasingly available, and consequently implementing protocols achieving error-free communication channels constitutes an important challenge. Indeed, in this kind of networks, the prevention of message conflicts and message collisions is a crucial issue. In terms of graph theory, solving this issue amounts to solve the distance-2 coloring problem in an arbitrary graph. The paper presents a distributed algorithm providing the processes with such a coloring. This algorithm is itself collision-free and conflict-free. It is particularly suited to wireless networks composed of nodes with communication or local memory constraints.
Abdelmadjid Bouabdallah, Hicham Lakhlef, Michel Raynal, François Taïani
AINA3
2017 Which Broadcast Abstraction Captures k-Set Agreement?
abstract
It is well-known that consensus (one-set agreement) and total order broadcast are equivalent in asynchronous systems prone to process crash failures. Considering wait-free systems, this article addresses and answers the following question: which is the communication abstraction that "captures" k-set agreement? To this end, it introduces a new broadcast communication abstraction, called k-BO-Broadcast, which restricts the disagreement on the local deliveries of the messages that have been broadcast (1-BO-Broadcast boils down to total order broadcast). Hence, in this context, k=1 is not a special number, but only the first integer in an increasing integer sequence. This establishes a new "correspondence" between distributed agreement problems and communication abstractions, which enriches our understanding of the relations linking fundamental issues of fault-tolerant distributed computing.
Damien Imbs, Achour Mostéfaoui, Matthieu Perrin, Michel Raynal
DISC4
2017 Signature-free asynchronous Byzantine systems: from multivalued to binary consensus with t< n/3, O(n2) messages, and constant time
Achour Mostéfaoui, Michel Raynal
Acta Informatica2
2017 A distributed leader election algorithm in crash-recovery and omissive systems
Christian Fernández-Campusano, Mikel Larrea, Roberto Cortiñas, Michel Raynal
Inf. Process. Lett.4
2017 Atomic Read/Write Memory in Signature-Free Byzantine Asynchronous Message-Passing Systems
Achour Mostéfaoui, Matoula Petrolia, Michel Raynal, Claude Jard
Theory Comput. Syst.3
2017 From wait-free to arbitrary concurrent solo executions in colorless distributed computing
Maurice Herlihy, Sergio Rajsbaum, Michel Raynal, Julien Stainer
Theor. Comput. Sci.3
2016 Efficient Broadcast Protocol for the Internet of Things
abstract
Internet of Things (IoT) is a network composed of a variety of heterogeneous things and objects such as Connected Wearable Devices (sensors, MEMS, microrobots, PDA, ...), Connected Cars, Connected Homes, Connected Cities, and the Industrial Internet. These things use generally wireless communication to interact and cooperate with each other to reach common services and goals. IoT(T, n) is a wireless network of things composed of T things with n items (information) distributed on it. The aim of the permutation routing is to route to each thing, its items, so it can accomplish its task. In this paper, we present an agent-based broadcast protocol for mobile Internet of Things that uses few communication channels. The main idea is to partition things into groups according to the number of channels. In each group, an agent manages a set of things. This new protocol performs efficiently with respect to the number of broadcast rounds and runs without conflict and collision on the communication channels. We give an estimation of the upper and the lower bounds of the number of broadcast rounds in the worst case. This paper is the first to present efficient broadcast protocol for the internet of things.
Hicham Lakhlef, Michel Raynal, Julien Bourgeois
AINA2
2016 Vertex Coloring with Communication and Local Memory Constraints in Synchronous Broadcast Networks
Hicham Lakhlef, Michel Raynal, François Taïani
ALGOSENSORS2
2016 Implementing Snapshot Objects on Top of Crash-Prone Asynchronous Message-Passing Systems
abstract
Distributed snapshots, as introduced by Chandy and Lamport in the context of asynchronous failure-free message-passing distributed systems, are consistent global states in which the observed distributed application might have passed through. It appears that two such distributed snapshots cannot necessarily be compared (in the sense of determining which one of them is the “first”). Differently, snapshots introduced in asynchronous crash-prone read/write distributed systems are totally ordered, which greatly simplify their use by upper layer applications. In order to benefit from shared memory snapshot objects, it is possible to simulate a read/write shared memory on top of an asynchronous crash-prone message-passing system, and build then snapshot objects on top of it. This algorithm stacking is costly in both time and messages. To circumvent this drawback, this paper presents algorithms building snapshot objects directly on top of asynchronous crash-prone message-passing system. “Directly” means here “without building an intermediate layer such as a read/write shared memory”. To the authors knowledge, the proposed algorithms are the first providing such constructions. Interestingly enough, these algorithms are efficient and relatively simple.
Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
ICA3PP4
2016 A Look at Basics of Distributed Computing
abstract
This tutorial presents concepts and basics of distributed computing which are important (at least from the author's point of view!), and should be known and mastered by Master students, researchers, and engineers. Those include: (a) a characterization of distributed computing (which is too much often confused with parallel computing); (b) the notion of a synchronous system and its associated notions of a local algorithm and message adversaries; (c) the notion of an asynchronous shared memory system and its associated notions of universality and progress conditions; and (d) the notion of an asynchronous messagepassing system with its associated broadcast and agreement abstractions, its impossibility results, and approaches to circumvent them. Hence, the tutorial can be seen as a guided tour to key elements that constitute basics of distributed computing.
Michel Raynal
ICDCS1
2016 Optimal Collision/Conflict-Free Distance-2 Coloring in Wireless Synchronous Broadcast/Receive Tree Networks
abstract
This article is on message-passing systems where communication is (a) synchronous and (b) based on the "broadcast/receive" pair of communication operations. "Synchronous" means that time is discrete and appears as a sequence of time slots (or rounds) such that each message is received in the very same round in which it is sent. "Broadcast/receive" means that during a round a process can either broadcast a message to its neighbors or receive a message from one of them. In such a communication model, no two neighbors of the same process, nor a process and any of its neighbors, must be allowed to broadcast during the same time slot (thereby preventing message collisions in the first case, and message conflicts in the second case). From a graph theory point of view, the allocation of slots to processes is known as the distance-2 coloring problem: a color must be associated with each process (defining the time slots in which it will be allowed to broadcast) in such a way that any two processes at distance at most 2 obtain different colors, while the total number of colors is "as small as possible". The paper presents a parallel message-passing distance-2 coloring algorithm suited to trees, whose roots are dynamically defined. This algorithm, which is itself collision-free and conflict-free, uses Δ+1 colors where Δ is the maximal degree of the graph (hence the algorithm is color-optimal). It does not require all processes to have different initial identities, and its time complexity is O(dΔ), where Δ is the depth of the tree. As far as we know, this is the first distributed distance-2 coloring algorithm designed for the broadcast/receive round-based communication model, which owns all the previous properties.
Davide Frey, Hicham Lakhlef, Michel Raynal
ICPP3
2016 Two-Bit Messages are Sufficient to Implement Atomic Read/Write Registers in Crash-prone Systems
abstract
Atomic registers are certainly the most basic objects of computing science. Their implementation on top of an n-process asynchronous message-passing system has received a lot of attention. It has been shown that t < n/2 (where t is the maximal number of processes that may crash) is a necessary and sufficient requirement to build an atomic register on top of a crash-prone asynchronous message-passing system. Considering such a context, this paper presents an algorithm which implements a single-writer multi-reader atomic register with four message types only, and where no message needs to carry control information in addition to its type. Hence, two bits are sufficient to capture all the control information carried by all the implementation messages. Moreover, the messages of two types need to carry a data value while the messages of the two other types carry no value at all. As far as we know, this algorithm is the first with such a sufficiency property on the size of control information carried by messages. It is also particularly efficient from a time complexity point of view.
Achour Mostéfaoui, Michel Raynal
PODC2
2016 t-Resilient Immediate Snapshot Is Impossible
Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
SIROCCO4
2016 Making Local Algorithms Wait-Free: The Case of Ring Coloring
Armando Castañeda, Carole Delporte-Gallet, Hugues Fauconnier, Sergio Rajsbaum, Michel Raynal
SSS5
2016 Are Byzantine Failures Really Different from Crash Failures?
Damien Imbs, Michel Raynal, Julien Stainer
DISC2
2016 Distributed Universality
Michel Raynal, Julien Stainer, Gadi Taubenfeld
Algorithmica1
2016 A necessary condition for Byzantine k-set agreement
Zohir Bouzid, Damien Imbs, Michel Raynal
Inf. Process. Lett.3
2016 Implementing set objects in dynamic distributed systems
Roberto Baldoni, Silvia Bonomi, Michel Raynal
J. Comput. Syst. Sci.3
2016 Read/write shared memory abstraction on top of asynchronous Byzantine message-passing systems
Damien Imbs, Sergio Rajsbaum, Michel Raynal, Julien Stainer
J. Parallel Distributed Comput.3
2016 Generalized Symmetry Breaking Tasks and Nondeterminism in Concurrent Objects
abstract
Processes in a concurrent system need to coordinate using an underlying shared memory or a message-passing system in order to solve agreement tasks such as, for example, consensus or set agreement. However, coordination is often needed to break the symmetry of processes that are initially in the same state---for example, to get exclusive access to a shared resource, to get distinct names, or to elect a leader. This paper introduces and studies the family of generalized symmetry breaking (GSB) tasks, which includes election, renaming, and many other symmetry breaking tasks, and studies how nondeterminism properties of objects solving tasks affects the computability power of GSB tasks. The aim is to develop the understanding of symmetry breaking tasks and their relation with agreement tasks and to study nondeterminism properties of objects solving tasks and how these properties affect the computability power of symmetry breaking tasks. Among various results characterizing the family of GSB tasks, it is shown that perfect renaming, i.e., $(n,n)$-renaming, is universal for all GSB tasks. The paper also shows that there is a large family of GSB tasks, which includes perfect renaming, that is strictly more powerful than $(n,n-1)$-set agreement. Some of these tasks are equivalent to perfect renaming, while others lie strictly between perfect renaming and $(n,n+1)$-renaming. Results comparing renaming and set agreement are proved, and the results in this paper complement known results. This paper sheds new light on the relations linking set agreement and symmetry breaking. The proofs are based on combinatorial topology techniques and new ideas about different notions of nondeterminism that can be associated with shared objects.
Armando Castañeda, Damien Imbs, Sergio Rajsbaum, Michel Raynal
SIAM J. Comput.4
2016 Predicate Detection in Asynchronous Distributed Systems: A Probabilistic Approach
abstract
In an asynchronous distributed system, a number of processes communicate with each other via message passing that has a finite but arbitrary long delay. There is no global clock in that system. Predicates, denoting the states of processes and their relations, are often used to specify the information of interest in such a system. Due to the lack of a global clock, the temporal relations between the states at different processes cannot be uniquely determined, but have multiple possible circumstances. Existing works of predicate detection are based on the definitely modality or the possibly modality, denoting that a predicate holds in all of the possible circumstances or in one of them, respectively. No information is provided about the probability that a predicate will hold, which hinders the taking of countermeasures for different situations. Moreover, the detection is based on single occurrence of a predicate, so the results are heavily affected by environmental noise and detection errors. In this paper, we propose a new approach to predicate detection to address these two issues. We generalize the definitely and possibly modalities to an occurrence probability to provide more detailed information, and further investigate how to detect multiple occurrences of a predicate. We propose a unified algorithm framework for detecting various types of predicates and demonstrate the use of it for three typical types of predicates, including simple predicates, simple sequences, and interval-constrained sequences. Theoretical proofs and simulation results show that our approach is effective and outperforms existing approaches.
Weiping Zhu 0004, Jiannong Cao 0001, Michel Raynal
IEEE Trans. Computers3
2016 Distributed Slicing in Dynamic Systems
abstract
Peer to peer (P2P) systems have moved from application specific architectures to a generic service oriented design philosophy. This raised interesting problems in connection with providing useful P2P middleware services capable of dealing with resource assignment and management in a large-scale, heterogeneous and unreliable environment. The slicing problem consists of partitioning a P2P network into$k$groups (slices) of a given portion of the network nodes that share similar resource values. As the network is large and dynamic this partitioning is continuously updated without any node knowing the network size. In this paper, we propose the first algorithm to solve the slicing problem. We introduce the metric of slice disorder and show that the existing ordering algorithm cannot nullify this disorder. We propose a new algorithm that speeds up the existing ordering algorithm but that suffers from the same inaccuracy. Then, we propose another algorithm based on ranking that is provably convergent under reasonable assumptions. In particular, we notice experimentally that ordering algorithms suffer from resource-correlated churn while the ranking algorithm can cope with it. These algorithms are proved viable theoretically and experimentally.
Antonio Fernández 0001, Vincent Gramoli, Ernesto Jiménez, Anne-Marie Kermarrec, Michel Raynal
IEEE Trans. Parallel Distributed Syst.5
2016 Intrusion-Tolerant Broadcast and Agreement Abstractions in the Presence of Byzantine Processes
abstract
A process commits a Byzantine failure when its behavior does not comply with the algorithm it is assumed to execute. Considering asynchronous message-passing systems, this paper presents distributed abstractions, and associated algorithms, that allow non-faulty processes to correctly cooperate, despite the uncertainty created by the net effect of asynchrony and Byzantine failures. These abstractions are broadcast abstractions (namely, no-duplicity broadcast, reliable broadcast, and validated broadcast), and agreement abstraction (namely, consensus). While no-duplicity broadcast and reliable broadcast are well-known one-to-all communication abstractions, validated broadcast is a new all-to-all communication abstraction designed to address agreement problems. After having introduced these abstractions, the paper presents an algorithm implementing validated broadcast on top of reliable broadcast. Then the paper presents two consensus algorithms, which are reductions of multivalued consensus to binary consensus. The first one is a generic algorithm, that can be instantiated with unreliable broadcast or no-duplicity broadcast, while the second is a consensus algorithm based on validated broadcast. Finally, a third algorithm is presented that solves the binary consensus problem. This algorithm is a randomized algorithm based on validated broadcast and a common coin. The presentation of all the abstractions and their algorithms is done incrementally.
Achour Mostéfaoui, Michel Raynal
IEEE Trans. Parallel Distributed Syst.2
2015 A Simple Predicate to Expedite the Termination of a Randomized Consensus Algorithm
abstract
Consensus is one of the most important problems encountered in fault-tolerant distributed computing. Basically, consensus allows processes to agree on a common value. Unfortunately, no deterministic algorithm can solve this problem in an asynchronous message-passing system prone to process crash failures. One way to circumvent this impossibility, consists in enriching the system with random numbers and design a randomized algorithm. This paper considers such a consensus algorithm and presents a simple predicate that allows to expedite its termination.
Michel Raynal
AINA1
2015 Concurrent Systems: Hybrid Object Implementations and Abortable Objects
Michel Raynal
Euro-Par1
2015 Anonymous Obstruction-Free (n, k)-Set Agreement with n-k+1 Atomic Read/Write Registers
abstract
The k-set agreement problem is a generalization of the consensus problem. Namely, assuming that each process proposes a value, every non-faulty process should decide one of the proposed values, and no more than k different values should be decided. This is a hard problem in the sense that we cannot solve it in an asynchronous system, as soon as k or more processes may crash. One way to sidestep this impossibility result consists in weakening the termination property, requiring that a process must decide a value only if it executes alone during a long enough period of time. This is the well-known obstruction-freedom progress condition. Consider a system of n anonymous asynchronous processes that communicate through atomic read/write registers, and such that any number of them may crash. In this paper, we address and solve the challenging open problem of designing an obstruction-free k-set agreement algorithm using only (n-k+1) atomic registers. From a shared memory cost point of view, our algorithm is the best algorithm known so far, thereby establishing a new upper bound on the number of registers needed to solve the problem, and in comparison to the previous upper bound, its gain is (n-k) registers. We then extend this algorithm into a space-optimal solution for the repeated version of k-set agreement, and an x-obstruction-free solution that employs 0(n-k+x) atomic registers (with 1 <= x <= k < n).
Zohir Bouzid, Michel Raynal, Pierre Sutra
OPODIS2
2015 Signature-Free Communication and Agreement in the Presence of Byzantine Processes (Tutorial)
abstract
Communication and agreement are fundamental abstractions in any distributed system. (If the computing entities do not need to communicate or agree in one way or another, the system is not a distributed system!) This tutorial was devoted to the design of such abstractions built on top of signature-free asynchronous distributed systems prone to Byzantine process failures. It is made up of three parts, each devoted to an abstraction and algorithms that implement it.
Michel Raynal
OPODIS1
2015 Stabilizing Server-Based Storage in Byzantine Asynchronous Message-Passing Systems: Extended abstract
abstract
A stabilizing Byzantine single-writer single-reader (SWSR) regular register, which stabilizes after the first invoked write operation, is first presented. Then, new/old ordering inversions are eliminated by the use of a (bounded) sequence number for writes, obtaining a practically stabilizing SWSR atomic register. A practically stabilizing Byzantine single-writer multi-reader (SWMR) atomic register is then obtained by using several copies of SWSR atomic registers. Finally, bounded time-stamps, with a time-stamp per writer, together with SWMR atomic registers, are used to construct a practically stabilizing Byzantine multi-writer multi-reader (MWMR) atomic register. In a system of n servers implementing an atomic register, and in addition to transient failures, the constructions tolerate t<n/8 Byzantine servers if communication is asynchronous, and t
Silvia Bonomi, Shlomi Dolev, Maria Potop-Butucaru, Michel Raynal
PODC4
2015 Minimal Synchrony for Byzantine Consensus
abstract
Solving the consensus problem requires in one way or another that the underlying system satisfies some synchrony assumption. Considering an asynchronous message-passing system of n processes where (a) up to t< n/3 may commit Byzantine failures, and (b) each pair of processes is connected by two uni-directional channels (with possibly different timing properties), this paper investigates the synchrony assumption required to solve consensus, and presents a signature-free consensus algorithm whose synchrony requirement is the existence of a process that is an eventual {t+1}bisource. Such a process p is a correct process that eventually has (a) timely input channels from t correct processes and (b) timely output channels to t correct processes (these input and output channels can connect p to different subsets of processes). As this synchrony condition was shown to be necessary and sufficient in the stronger asynchronous system model (a) enriched with message authentication, and (b) where the channels are bidirectional and have the same timing properties in both directions, it follows that it is also necessary and sufficient in the weaker system model considered in the paper. In addition to the fact that it closes a long-lasting problem related to Byzantine agreement, a noteworthy feature of the proposed algorithm lies in its design simplicity, which is a first-class property.
Zohir Bouzid, Achour Mostéfaoui, Michel Raynal
PODC3
2015 Eventual Leader Election Despite Crash-Recovery and Omission Failures
abstract
In this work we consider the problem of leader election, abstraction used by many distributed services to select a unique process for coordinating actions. We propose an eventual leader election algorithm for partially synchronous systems prone to concurrent crash-recovery and omission failures where any process may suffer failures forever as long as a majority of processes meet some weak connectivity and reliability conditions.
Christian Fernández-Campusano, Mikel Larrea, Roberto Cortiñas, Michel Raynal
PRDC4
2015 Signature-Free Asynchronous Byzantine Systems: From Multivalued to Binary Consensus with t < n/3, O(n2) Messages, and Constant Time
Achour Mostéfaoui, Michel Raynal
SIROCCO2
2015 Communication Patterns and Input Patterns in Distributed Computing - (Invited Talk)
Michel Raynal
SIROCCO1
2015 Specifying Concurrent Problems: Beyond Linearizability and up to Tasks - (Extended Abstract)
Armando Castañeda, Sergio Rajsbaum, Michel Raynal
DISC3
2015 Signature-Free Asynchronous Binary Byzantine Consensus with t < n/3, O(n2) Messages, and O(1) Expected Time
abstract
This article is on broadcast and agreement in asynchronous message-passing systems made up of n processes, and where up to t processes may have a Byzantine Behavior. Its first contribution is a powerful, yet simple, all-to-all broadcast communication abstraction suited to binary values. This abstraction, which copes with up to t < n /3 Byzantine processes, allows each process to broadcast a binary value, and obtain a set of values such that (1) no value broadcast only by Byzantine processes can belong to the set of a correct process, and (2) if the set obtained by a correct process contains a single value v , then the set obtained by any correct process contains v . The second contribution of this article is a new round-based asynchronous consensus algorithm that copes with up to t < n /3 Byzantine processes. This algorithm is based on the previous binary broadcast abstraction and a weak common coin. In addition to being signature-free and optimal with respect to the value of t , this consensus algorithm has several noteworthy properties: the expected number of rounds to decide is constant; each round is composed of a constant number of communication steps and involves O ( n 2) messages; each message is composed of a round number plus a constant number of bits. Moreover, the algorithm tolerates message reordering by the adversary (i.e., the Byzantine processes).
Achour Mostéfaoui, Moumen Hamouma, Michel Raynal
J. ACM3
2015 Failure detectors in homonymous distributed systems (with an application to consensus)
Sergio Arévalo, Antonio Fernández 0001, Damien Imbs, Ernesto Jiménez, Michel Raynal
J. Parallel Distributed Comput.5
2015 Special Issue on Distributed Computing and Networking
Michel Raynal, Franck Petit
Theor. Comput. Sci.1
2014 A Simple Broadcast Algorithm for Recurrent Dynamic Systems
abstract
This paper presents a simple broadcast algorithm suited to dynamic systems where links can repeatedly appear and disappear. The algorithm is proved correct and a simple improvement is introduced, that reduces the number and the size of control messages. As it extends in a simple way a classical network traversal algorithm to the dynamic context, the proposed algorithm has also pedagogical flavor.
Michel Raynal, Julien Stainer, Jiannong Cao 0001, Weigang Wu
AINA1
2014 Simple Deadlock Detection for the And-Communication Model
abstract
The advent of multicore architectures is a good incentive to revisit base synchronization mechanisms. Among them, the AND communication model is particularly attractive. This communication model provides the processes with a receive operation denoted receive (DS) where DS is a dynamically defined set of processes (DS stands for "dependency set"). The receive operation blocks the invoking process until it has received a message from each process appearing in the dynamically defined set DS. When this occurs, the invoking process consumes these messages and continues its execution. While it simplifies concurrent programming, this high-level communication operation is deadlock-prone. This paper presents a very simple algorithm which allows to detect on the fly communication deadlock in the AND communication model.
Michel Raynal
CISIS1
2014 Computing in the Presence of Concurrent Solo Executions
Maurice Herlihy, Sergio Rajsbaum, Michel Raynal, Julien Stainer
LATIN3
2014 Distributed Universality
Michel Raynal, Julien Stainer, Gadi Taubenfeld
OPODIS1
2014 Signature-free asynchronous byzantine consensus with t 2<n/3 and o(n2) messages
abstract
This paper presents a new round-based asynchronous consensus algorithm that copes with up to t
Achour Mostéfaoui, Moumen Hamouma, Michel Raynal
PODC3
2014 Brief announcement: distributed universality: contention-awareness; wait-freedom; object progress, and other properties
abstract
A notion of a universal construction suited to distributed computing has been introduced by M. Herlihy in his celebrated paper "Wait-free synchronization" (ACM TOPLAS, 1991). A universal construction is an algorithm that can be used to wait-free implement any object defined by a sequential specification. Herlihy's paper shows that the basic system model, which supports only atomic read/write registers, has to be enriched with consensus objects to allow the design of universal constructions.
Michel Raynal, Julien Stainer, Gadi Taubenfeld
PODC1
2014 Reliable Shared Memory Abstraction on Top of Asynchronous Byzantine Message-Passing Systems
Damien Imbs, Sergio Rajsbaum, Michel Raynal, Julien Stainer
SIROCCO3
2014 Fair Synchronization in the Presence of Process Crashes and its Weakest Failure Detector
abstract
A non-blocking implementation of a concurrent object is an implementation that does not prevent concurrent accesses to the internal representation of the object, while guaranteeing the deadlock-freedom progress condition without using locks. Considering a failure free context, G. Taubenfeld has introduced (DISC 2013) a simple modular approach, captured under a new problem called the it fair synchronization problem, to transform a non-blocking implementation into a starvation-free implementation satisfying a strong fairness requirement. This paper extends this approach in several directions. It first generalizes the fair synchronization problem to read/write asynchronous systems where any number of processes may crash. Then, it introduces a new failure detector and uses it to solve the fair synchronization problem when processes may crash. This failure detector, denoted QP (Quasi Perfect), is very close to, but strictly weaker than, the perfect failure detector. Last but not least, the paper shows that the proposed failure detector QP is optimal in the sense that the information on failures it provides to the processes can be extracted from any algorithm solving the fair synchronization problem in the presence of any number of process crash failures.
Carole Delporte-Gallet, Hugues Fauconnier, Michel Raynal
SRDS3
2013 Coordination and Computation in Distributed Intelligent MEMS
abstract
Over the last decades, research on microelectromechanical systems (MEMS) has focused on the engineering process which has led to major advances. Future challenges will consist in adding embedded intelligence to MEMS systems to obtain distributed intelligent MEMS. One intrinsic characteristic of MEMS is their ability to be mass-produced. This, however, poses scalability problems because a significant number of MEMS can be placed in a small volume. Managing this scalability requires paradigm-shifts both in hardware and software parts. Furthermore, the need for actuated synchronization, programming, communication and mobility management raises new challenges in both control and programming. Finally, MEMS are prone to faulty behaviors as they are mechanical systems and they are issued from a batch fabrication process. A new programming paradigm which can meet these challenges is therefore needed. In this article, we present CO2Dim, which stands for Coordination and Computation in Distributed Intelligent MEMS. CO2DIM is a common project between France and Hong-Kong building a new programming environment which includes a language, based on a joint development of programming and control capabilities, a simulator and real hardware.
Julien Bourgeois, Jiannong Cao 0001, Michel Raynal, Dominique Dhoutaut, Benoît Piranda, Eugen Dedu, Ahmed Mostefaoui, Hakim Mabed
AINA3
2013 A Short Introduction to Synchronous Communication
abstract
The advent of multicore architectures is a good incentive to better understand base synchronization mechanisms. This paper, which can be considered as a simple introduction to the topic, presents (with a pedagogical flavor) the concept of rendezvous (also called interaction, synchronous communication, or logically instantaneous communication) and several implementations of it. This abstraction adds synchronization to communication, namely, it requires that, for a message to be sent by a process, the receiver has to be ready to receive it. From an external observer point view, the message transmission looks like instantaneous: the sending and the reception of a message appear as being a single event (and the sense of the communication could have been in the other direction). From an operational point of view, we have the following: for each pair of processes, the first process that wants to communicate - be it the sender or the receiver - has to. wait until the other process is ready to communicate.
Michel Raynal
AINA1
2013 A Contention-Friendly Binary Search Tree
Tyler Crain, Vincent Gramoli, Michel Raynal
Euro-Par3
2013 No Hot Spot Non-blocking Skip List
abstract
This paper presents a new non-blocking skip list algorithm. The algorithm alleviates contention by localizing synchronization at the least contended part of the structure without altering consistency of the implemented abstraction. The key idea lies in decoupling a modification to the structure into two stages: an eager abstract modification that returns quickly and whose update affects only the bottom of the structure, and a lazy selective adaptation updating potentially the entire structure but executed continuously in the background. On SPECjbb as well as on micro-benchmarks, we compared the performance of our new non-blocking skip list against the performance of the JDK non-blocking skip list. The results indicate that our implementation can me more than twice as fast as the JDK skip list.
Tyler Crain, Vincent Gramoli, Michel Raynal
ICDCS3
2013 A Generalized Mutual Exclusion Problem and Its Algorithm
abstract
Mutual exclusion (ME) is a fundamental problem for resource allocation in distributed systems, It is concerned with how the various processes access shared resources in a mutually exclusive way. Besides the classic ME problem, several variant problems have been proposed and studied. In this paper, drawing inspiration from the scenario of controlling autonomous vehicles at intersections, we have defined a new ME problem, called Local Group Mutual Exclusion (LGME), w here mutual exclusion is necessary only among the processes requesting overlap but not the same set of resources. In comparison with other variant problems of ME, LGME is more general but also more challenging. To solve the LGME problem, we propose a novel notion called strong coterie, which can handle the complex process relationship in LGME. Based on strong coterie, we have designed an ME algorithm, which can handle concurrent CS execution and message asynchrony. The correctness of our algorithm is rigorously proved.
Aoxue Luo, Weigang Wu, Jiannong Cao 0001, Michel Raynal
ICPP4
2013 Agreement via Symmetry Breaking: On the Structure of Weak Subconsensus Tasks
abstract
This paper is on the relative power and the relations linking two important synchronization problems in n-process wait-free shared memory models, namely, set agreement and renaming, which are two of the most studied subconsensus tasks. Since the 2006 seminal paper of Gafni, Rajsbaum and Herlihy, it is known that some renaming instances are strictly weaker than set agreement. Indeed, it was later on shown that not even (n + 1)-renaming (the strongest task in the renaming family, after perfect n-renaming) can implement (n - 1)-set agreement (the weakest non-trivial task in the set agreement family). These and other results seem to imply that renaming and, more generally, the tasks called generalized symmetry breaking tasks (GSB) are weaker than agreement tasks. This paper shows that this is not the case, namely, it shows that there is a large family of GSB tasks that are more powerful than (n - 1)-set agreement. Some of these tasks are equivalent to n-renaming, while others lie strictly between n-renaming and (n+1)-renaming. Moreover, none of these GSB tasks can solve (n - 2)-set agreement. Hence, these subconsensus tasks have a rich structure and are interesting in their own. The proofs of these results are based on algebraic topology techniques and new ideas about different notions of nondeterminism that can be associated with shared objects. Interestingly, this paper sheds a new light on the relations linking set agreement and renaming.
Armando Castañeda, Sergio Rajsbaum, Michel Raynal
IPDPS3
2013 Synchrony weakened by message adversaries vs asynchrony restricted by failure detectors
abstract
A message adversary is a daemon that suppresses messages in round-based message-passing synchronous systems in which no process crashes. A property imposed on a message adversary defines a subset of messages that cannot be eliminated by the adversary. It has recently been shown that when a message adversary is constrained by a property denoted TOUR (for tournament), the corresponding synchronous system and the asynchronous crash-prone read/write system have the same computability power for task solvability.
Michel Raynal, Julien Stainer
PODC1
2013 Fault-Tolerant Leader Election in Mobile Dynamic Distributed Systems
abstract
This paper addresses the leader election problem in dynamic distributed systems with mobile processes. To do so, it is assumed that the system alternates periods of good and bad behavior, in the line of the timed asynchronous model of Cristian and Fetzer. We extend the eventual leadership properties recently proposed by Larrea et al. for non-mobile dynamic systems, defining two new properties that take into account graph joins/fragmentations due to process mobility. We also propose a new leader election algorithm in a weak mobile dynamic distributed system model. Using a categorization framework, we compare our system model with a number of models proposed in the literature, showing that our leader election algorithm works in a model which is weaker than the rest.
Carlos Gómez-Calzado, Alberto Lafuente, Mikel Larrea, Michel Raynal
PRDC4
2013 Simultaneous Consensus vs Set Agreement: A Message-Passing-Sensitive Hierarchy of Agreement Problems
Michel Raynal, Julien Stainer
SIROCCO1
2013 Anonymous asynchronous systems: the case of failure detectors
François Bonnet 0001, Michel Raynal
Distributed Comput.2
2013 Towards a universal construction for transaction-based multiprocess programs
Tyler Crain, Damien Imbs, Michel Raynal
Theor. Comput. Sci.3
2013 Trust-aware peer sampling: Performance and privacy tradeoffs
abstract
The ability to identify people that share one’s own interests is one of the most interesting promises of the Web 2.0 driving user-centric applications such as recommendation systems or collaborative marketplaces. To be truly useful, however, information about other users also needs to be associated with some notion of trust. Consider a user wishing to sell a concert ticket. Not only must she find someone who is interested in the concert, but she must also make sure she can trust this person to pay for it. This paper addresses the need for trust in user-centric applications by proposing two novel distributed protocols that combine interest-based connections between users with explicit links obtained from social networks à-la Facebook. Both protocols build trusted multi-hop paths between users in an explicit social network supporting the creation of semantic overlays backed up by social trust. The first protocol, TAPS2, extends our previous work on TAPS (Trust-Aware Peer Sampling), by improving the ability to locate trusted nodes. Yet, it remains vulnerable to attackers wishing to learn about trust values between arbitrary pairs of users. The second protocol, PTAPS (Private TAPS), improves TAPS2 with provable privacy guarantees by preventing users from revealing their friendship links to users that are more than two hops away in the social network. In addition to proving this privacy property, we evaluate the performance of our protocols through event-based simulations, showing significant improvements over the state of the art.
Davide Frey, Arnaud Jégou, Anne-Marie Kermarrec, Michel Raynal, Julien Stainer
Theor. Comput. Sci.4
2013 Power and limits of distributed computing shared memory models
Maurice Herlihy, Sergio Rajsbaum, Michel Raynal
Theor. Comput. Sci.3
2013 The weakest failure detector to implement a register in asynchronous systems with hybrid communication
Damien Imbs, Michel Raynal
Theor. Comput. Sci.2
2012 Trying to Unify the LL/SC Synchronization Primitive and the Notion of a Timed Register
abstract
The aim of this short paper is to show that both the LL/SC (Linked Load/Store Conditional) synchronization primitive and the notion of a Timed register can be seen as two instances of a general abstract synchronization object type that we called a predicate-based read/write synchronization object. More precisely, LL/SC corresponds to its time-free instance while a timed register corresponds to its timed instance. It follows that the notion of a predicate-based read/write synchronization object constitutes a unifying notion that allows for a deeper insight into synchronization objects proposed for multicore architectures.
Damien Imbs, Michel Raynal
AINA2
2012 Leader Election: From Higham-Przytycka's Algorithm to a Gracefully Degrading Algorithm
abstract
The leader election problem consists in selecting a process (called leader) in a group of processes. Several leader election algorithms have been proposed in the past for ring networks, tree networks, fully connected networks or regular networks (such as tori and hypercubes). As far as ring networks are concerned, it has been shown that the number of messages that processes have to exchange to elect a leader is Ω(n log n). The algorithm proposed by Higham and Przytycka is the best leader algorithm known so far for ring networks in terms of message complexity, which is 1.271 n log n + O(n). This algorithm uses round numbers and assumes that all processes start with the same round number. More precisely, when round numbers are not initially equal, the algorithm has runs that do not terminate. This paper presents an algorithm, based on Higham-Przytycka's technique, which allows processes to start with different round numbers. This extension is motivated by fault-tolerance with respect to initial values. While the algorithm always terminates, its message complexity is optimal, i.e., O(n log n), when the processes start with the same round number and increases up to O(n2) when all processes start with different round number values. We call graceful degradation this additional property that combines fault-tolerance (with respect to initial values) and efficiency.
Itziar Arrieta-Salinas, Federico Fariña, José Ramón González de Mendívil, Michel Raynal
CISIS4
2012 A Simple Asynchronous Shared Memory Consensus Algorithm Based on Omega and Closing Sets
abstract
This paper is on the design of a consensus object in the context of asynchronous shared memory systems where any number of process can suffer a crash failure. These systems are becoming more and more important with the advent of multicore architectures. To circumvent the impossibility of implementing a consensus object in such a context, the paper considers that the base read/write system model is enriched with an eventual leader failure detector (traditionally denoted Ω). This failure detector can easily be used to ensure that all the invocations of the consensus object issued by processes that do not crash eventually terminate(wait-freedom termination property). Hence, when one has to implement a consensus object in such an enriched system model, the main issue consists in designing an object (from base atomic read/write registers) on which the implementation can rely to ensure that no two different values can be decided from the consensus object. This paper presents such an object, called closing set. The main feature of this object is that it takes advantage of the system asynchrony by reducing the number of values that can be deposited: only concurrent deposits of values in an empty set are successful. The paper presents then a simple consensus algorithm based on closing sets. This algorithm is round-based and uses a closing set per round.
Michel Raynal, Julien Stainer
CISIS1
2012 From a Store-Collect Object and Ω to Efficient Asynchronous Consensus
Michel Raynal, Julien Stainer
Euro-Par1
2012 STM Systems: Enforcing Strong Isolation between Transactions and Non-transactional Code
Tyler Crain, Eleni Kanellou, Michel Raynal
ICA3PP (1)3
2012 Failure Detectors in Homonymous Distributed Systems (with an Application to Consensus)
abstract
This paper is on homonymous distributed systems where processes are prone to crash failures and have no initial knowledge of the system membership (``homonymous'' means that several processes may have the same identifier). New classes of failure detectors suited to these systems are first defined. Among them, the classes $\HO$ and $\HS$ are introduced that are the homonymous counterparts of the classes $\Omega$ and $\Sigma$, respectively. (Recall that the pair $\langle \Omega, \Sigma\rangle$ defines the weakest failure detector to solve consensus.) Then, the paper shows how $\HO$ and $\HS$ can be implemented in homonymous systems without membership knowledge (under different synchrony requirements). Finally, two algorithms are presented that use these failure detectors to solve consensus in homonymous asynchronous systems where there is no initial knowledge of the membership. One algorithm solves consensus with $\langle \HO, \HS\rangle$, while the other uses only $\HO$, but needs a majority of correct processes. Observe that the systems with unique identifiers and anonymous systems are extreme cases of homonymous systems from which follows that all these results also apply to these systems. Interestingly, the new failure detector class $\HO$ can be implemented with partial synchrony, while the analogous class $\AO$ defined for anonymous systems can not be implemented (even in synchronous systems). Hence, the paper provides us with the first proof showing that consensus can be solved in anonymous systems with only partial synchrony (and a majority of correct processes).
Sergio Arévalo, Antonio Fernández 0001, Damien Imbs, Ernesto Jiménez, Michel Raynal
ICDCS5
2012 Renaming Is Weaker Than Set Agreement But for Perfect Renaming: A Map of Sub-consensus Tasks
Armando Castañeda, Damien Imbs, Sergio Rajsbaum, Michel Raynal
LATIN4
2012 Chasing the Weakest Failure Detector for k-Set Agreement in Message-Passing Systems
abstract
This paper continues our quest for the weakest failure detector which allows the k-set agreement problem to be solved in asynchronous message-passing systems prone to process failures. It has two main contributions which will be instrumental to complete this quest. The first contribution is a new failure detector (denoted PiSigma(x, y)) that has several noteworthy properties. (a) It is stronger than Sigma(k) which has been shown to be necessary. (b) It is equivalent to the pair (Sigma, Omega) when x=y=1 (optimal to solve consensus). (c) It is equivalent to the pair (Sigma(n-1), Omega(n-1)) when x=y=n-1 (optimal for (n-1)-set agreement). (d) It is strictly weaker than the pair (Sigma(x), anti-Omega(y)) (which has been investigated in previous works). (e) It is operational: the paper presents a PiSigma(x, y)-based algorithm that solves k-set agreement for k greater or equal to xy (intuitively, x refers to the maximum number of isolated groups of processes and y to the number of leaders in each of these groups). The second contribution of the paper is a proof that, for k strictly between 1 and n-1, the eventual leaders failure detector Omega(k) (which eventually provides each process with the same set of k process identities, this set including at least one correct process) is not necessary to solve k-set agreement problem.
Achour Mostéfaoui, Michel Raynal, Julien Stainer
NCA2
2012 When and How Process Groups Can Be Used to Reduce the Renaming Space
Armando Castañeda, Michel Raynal, Julien Stainer
OPODIS2
2012 Brief announcement: there are plenty of tasks weaker than perfect renaming and stronger than set agreement
abstract
In the asynchronous wait-free shared memory model, two families of tasks play a central role because of their implications in theory and in practice: k-set agreement and M-renaming. Let n denote the number of processes in the system. Previous research shows that (n-1)-set agreement can solve (2n-2)-renaming, for any value of n, while (2n-2)-renaming cannot solve (n-1)-set agreement, when n is odd. It is also known that, for every n ≥ 3, n-renaming, also called perfect renaming, is strictly stronger than (n-1)-set agreement. This paper shows that when n ≥ 4, there is a family of tasks that are strictly stronger than (n-1)-set agreement and strictly weaker than perfect renaming. This enlarges our view of both the nature and the structure of what are distributed computing tasks.
Armando Castañeda, Sergio Rajsbaum, Michel Raynal
PODC3
2012 Brief announcement: increasing the power of the iterated immediate snapshot model with failure detectors
abstract
This short paper shows how to capture failure detectors so that the base asynchronous read/wite model and the distributed iterated model have the same computational power when both are enriched with the same failure detector. To that end it introduces the notion of a "strongly correct" process and presents simulations that prove the computational equivalence when both models are enriched with the same failure detector. Interestingly, these simulations, which work for a large family of failure detector classes, can be easily extended to the case where the wait-freedom requirement is replaced by the notion of t-resilience. A noteworthy and first class feature of the proposed approach lies in its simplicity.
Michel Raynal, Julien Stainer
PODC1
2012 A speculation-friendly binary search tree
abstract
We introduce the first binary search tree algorithm designed for speculative executions. Prior to this work, tree structures were mainly designed for their pessimistic (non-speculative) accesses to have a bounded complexity. Researchers tried to evaluate transactional memory using such tree structures whose prominent example is the red-black tree library developed by Oracle Labs that is part of multiple benchmark distributions. Although well-engineered, such structures remain badly suited for speculative accesses, whose step complexity might raise dramatically with contention.
Tyler Crain, Vincent Gramoli, Michel Raynal
PPoPP3
2012 Increasing the Power of the Iterated Immediate Snapshot Model with Failure Detectors
Michel Raynal, Julien Stainer
SIROCCO1
2012 Brief Announcement: A Contention-Friendly, Non-blocking Skip List
Tyler Crain, Vincent Gramoli, Michel Raynal
DISC3
2012 No double discount: Condition-based simultaneity yields limited gain
Yoram Moses, Michel Raynal
Inf. Comput.2
2012 From the Happened-Before Relation to the Causal Ordered Set Abstraction
Saúl E. Pomares Hernández, Jose Roberto Perez Cruz, Michel Raynal
J. Parallel Distributed Comput.3
2012 Help when needed, but no more: Efficient read/write partial snapshot
Damien Imbs, Michel Raynal
J. Parallel Distributed Comput.2
2012 Virtual world consistency: A condition for STM systems (with a versatile protocol with invisible read operations)
Damien Imbs, Michel Raynal
Theor. Comput. Sci.2
2012 Implementing a Regular Register in an Eventually Synchronous Distributed System Prone to Continuous Churn
abstract
Due to their capability to hide the complexity generated by the messages exchanged between processes, shared objects are one of the main abstractions provided to developers of distributed applications. Implementations of such objects, in modern distributed systems, have to take into account the fact that almost all services, implemented on top of distributed infrastructures, are no longer fully managed due to either their size or their maintenance cost. Therefore, these infrastructures exhibit several autonomic behaviors in order to, for example, tolerate failures and continuous arrival and departure of nodes (churn phenomenon). Among all the shared objects, the register object is a fundamental one. Several protocols have been proposed to build fault resilient registers on top of message-passing system, but, unfortunately, failures are not the only challenge in modern distributed systems and new issues arise in the presence of churn. This paper addresses the construction of a multiwriter/multireader regular register in an eventually synchronous distributed system affected by the continuous arrival/departure of participants. In particular, a general protocol implementing a regular register is proposed and feasibility conditions associated with the arrival and departure of the processes are given. The protocol is proved correct under the assumption that a constraint on the churn is satisfied.
Roberto Baldoni, Silvia Bonomi, Michel Raynal
IEEE Trans. Parallel Distributed Syst.3
2011 A Theory-Oriented Introduction to Wait-Free Synchronization Based on the Adaptive Renaming Problem
abstract
The recent deployment of multiprocessors (such as multicores) as the mainstream computing platform has given rise to a new concurrent programming impetus. In such a context it becomes extremely important to be able to design shared objects that can cope with the net effect of asynchrony and process crashes. This paper is a theory-oriented introduction to wait-free synchronization for such systems. It uses the adaptive renaming problem as a paradigm to explain the difficulties and subtleties of synchronization in presence of process crashes. Renaming is one of the most famous coordination problems studied in distributed computability. It consists in assigning new names to processes in such a way that no two processes obtain the same new name and the new name space be as small as possible. The paper visits the problem by presenting three solutions. This paper, that has a strong survey/short tutorial flavor, can consequently be considered as an introduction to both there naming problem and progress conditions for synchronization in presence of process crashes in the context of multiprocessor systems.
Sergio Rajsbaum, Michel Raynal
AINA2
2011 Read Invisibility, Virtual World Consistency and Probabilistic Permissiveness are Compatible
Tyler Crain, Damien Imbs, Michel Raynal
ICA3PP (1)3
2011 The universe of symmetry breaking tasks
abstract
This brief announcement introduces the family of generalized symmetry breaking (GSB) tasks, that includes election, renaming and many other symmetry breaking tasks. Differently from agreement tasks, a GSB task is "inputless", in the sense that processes do not propose values; the task specifies only the symmetry breaking requirement, independently of the system's initial state (where processes differ only on their identifiers). Among various results characterizing the family of GSB tasks, it is shown that (non adaptive) perfect renaming is universal for all GSB tasks.
Damien Imbs, Sergio Rajsbaum, Michel Raynal
PODC3
2011 The Universe of Symmetry Breaking Tasks
Damien Imbs, Sergio Rajsbaum, Michel Raynal
SIROCCO3
2011 A Survey on Some Recent Advances in Shared Memory Models
Sergio Rajsbaum, Michel Raynal
SIROCCO2
2011 Brief announcement: read invisibility, virtual world consistency and permissiveness are compatible
abstract
This brief announcement studies the relation between two STM properties (read invisibility and permissiveness) and two consistency conditions for STM systems, namely, opacity and virtual world consistency. A read operation issued by a transaction is invisible if it does not entail shared memory modifications. An STM system is permissive with respect to a consistency condition if it accepts every history that satisfies the condition. The brief announcement first shows that read invisibility, permissiveness and opacity are incompatible. It then shows that invisibility, permissiveness and virtual world consistency are compatible.
Tyler Crain, Damien Imbs, Michel Raynal
SPAA3
2011 The Weakest Failure Detector to Implement a Register in Asynchronous Systems with Hybrid Communication
Damien Imbs, Michel Raynal
SSS2
2011 Relations Linking Failure Detectors Associated with k-Set Agreement in Message-Passing Systems
Achour Mostéfaoui, Michel Raynal, Julien Stainer
SSS2
2011 Brief Announcement: ΔΩ: Specifying an Eventual Leader Service for Dynamic Systems
Mikel Larrea, Michel Raynal
DISC2
2011 A liveness condition for concurrent objects: x-wait-freedom
abstract
SUMMARY The liveness of concurrent objects despite asynchrony and failures is a fundamental problem. To that end several progress conditions have been proposed. Wait‐freedom is the strongest of these conditions: it states that any object operation must terminate if the invoking process does not crash. Obstruction‐freedom is a weaker progress condition as it requires progress only when a process executes in isolation for a long enough period. This paper explores progress conditions in n‐process asynchronous read/write systems enriched with base objects with consensus number x, 1
Damien Imbs, Michel Raynal
Concurr. Comput. Pract. Exp.2
2011 The Price of Anonymity: Optimal Consensus Despite Asynchrony, Crash, and Anonymity
abstract
This article addresses the consensus problem in asynchronous systems prone to process crashes, where additionally the processes are anonymous (they cannot be distinguished one from the other: they have no name and execute the same code). To circumvent the three computational adversaries (asynchrony, failures, and anonymity) each process is provided with a failure detector of a class denoted ψ , that gives it an upper bound on the number of processes that are currently alive (in a nonanonymous system, the classes ψ and P ---the class of perfect failure detectors---are equivalent). The article first presents a simple ψ -based consensus algorithm where the processes decide in 2 t + 1 asynchronous rounds (where t is an upper bound on the number of faulty processes). It then shows one of its main results, namely 2 t + 1 is a lower bound for consensus in the anonymous systems equipped with ψ . The second contribution addresses early-decision. The article presents and proves correct an early-deciding algorithm where the processes decide in min(2 f + 2, 2 t + 1) asynchronous rounds (where f is the actual number of process failures). This leads us to think that anonymity doubles the cost (with respect to synchronous systems) and it is conjectured that min(2 f + 2, 2 t + 1) is the corresponding lower bound. The article finally considers the k -set agreement problem in anonymous systems. It first shows that the previous ψ -based consensus algorithm solves the k -set agreement problem in Rt = 2⌊t k⌋ + 1 asynchronous rounds. Then, considering a family of failure detector classes { ψℓ }0 ≤ ℓ < k that generalizes the class ψ (= ψ 0 ), the article presents an algorithm that solves the k -set agreement in Rt,ℓ = 2 ⌊ t k − ℓ ⌋ + 1 asynchronous rounds. This last formula relates the cost ( Rt,ℓ ) the coordination degree of the problem ( k ), the maximum number of failures ( t ), and the the strength ( ℓ ) of the underlying failure detector. Finally the article concludes by presenting problems that remain open.
François Bonnet 0001, Michel Raynal
ACM Trans. Auton. Adapt. Syst.2
2011 On the road to the weakest failure detector for k-set agreement in message-passing systems
François Bonnet 0001, Michel Raynal
Theor. Comput. Sci.2
2011 Software transactional memories: an approach for multicore programming
Damien Imbs, Michel Raynal
J. Supercomput.2
2010 Consensus in Anonymous Distributed Systems: Is There a Weakest Failure Detector?
abstract
This paper is on failure detectors to solve the consensus problem in asynchronous systems made up of anonymous processes prone to crash and connected by asynchronous reliable channels. Anonymity means that any two processes cannot be distinguished one from the other: they have no name and execute the same code. The paper has several contributions. It first introduces two new classes of failures detectors, denoted AP and AOmega, and presents an AP-based algorithm and an AOmega-based algorithm that solve the consensus problem despite the three computational adversaries that are asynchrony, failures and anonymity. Then, the paper shows that, in crash-prone non-anonymous systems, (a) AP and the class of perfect failure detectors denoted P) are equivalent, and (b) AOmega and the class of eventual leader failure detectors (denoted Omega) are also equivalent. Finally, the paper addresses the question of the weakest failure detector to solve consensus in an asynchronous crash-prone anonymous system. In non-anonymous systems, the class P of perfect failure detectors is strictly stronger than the class Omega of eventual leader failure detectors that has been shown to be the weakest failure detector class for consensus in asynchronous crash-prone system. Quite surprisingly, the paper shows that their anonymous counterparts cannot be compared.
François Bonnet 0001, Michel Raynal
AINA2
2010 Value-Based Sequential Consistency for Set Objects in Dynamic Distributed Systems
Roberto Baldoni, Silvia Bonomi, Michel Raynal
Euro-Par (1)3
2010 The x-Wait-Freedom Progress Condition
Damien Imbs, Michel Raynal
Euro-Par (1)2
2010 Signature-Free Broadcast-Based Intrusion Tolerance: Never Decide a Byzantine Value
Achour Mostéfaoui, Michel Raynal
OPODIS2
2010 The multiplicative power of consensus numbers
abstract
The Borowsky-Gafni (BG) simulation algorithm is a powerful reduction algorithm that shows that t-resilience of decision tasks can be fully characterized in terms of wait-freedom. Said in another way, the BG simulation shows that the crucial parameter is not the number n of processes but the upper bound t on the number of processes that are allowed to crash. The BG algorithm considers colorless decision tasks in the base read/write shared memory model. (Colorless means that if, process decides a value, any other process is allowed to decide the very same value.)
Damien Imbs, Michel Raynal
PODC2
2010 On asymmetric progress conditions
abstract
Wait-freedom and obstruction-freedom have received a lot of attention in the literature. These are symmetric progress conditions in the sense that they consider all processes as being "equal". Wait-freedom has allowed to rank the synchronization power of objects in presence of process failures, while (the weaker) obstruction-freedom allows for simpler and more efficient object implementations.
Damien Imbs, Michel Raynal, Gadi Taubenfeld
PODC2
2010 On Adaptive Renaming under Eventually Limited Contention
Damien Imbs, Michel Raynal
SSS2
2010 The 2010 Edsger W. Dijkstra Prize in Distributed Computing
Marcos K. Aguilera, Michel Raynal
DISC2
2010 Anonymous Asynchronous Systems: The Case of Failure Detectors
François Bonnet 0001, Michel Raynal
DISC2
2010 A Timing Assumption and Two t-Resilient Protocols for Implementing an Eventual Leader Service in Asynchronous Shared Memory Systems
Antonio Fernández 0001, Ernesto Jiménez, Michel Raynal, Gilles Trédan
Algorithmica3
2010 A Methodological Construction of an Efficient Sequentially Consistent Distributed Shared Memory
abstract
The paper proposes a simple protocol that ensures sequential consistency. The protocol assumes that the shared memory abstraction is supported by the local memories of nodes that can communicate only by exchanging messages through reliable channels. Unlike other sequential consistency protocols, the one proposed here does not rely on a strong synchronization mechanism, such as an atomic broadcast primitive or a central node managing a copy of every shared object. From a methodological point of view, the protocol is built incrementally starting from the very definition of sequential consistency. It has the noteworthy property that a process that issues a write operation never has to wait for other processes. Depending on the current local state, most read operations issued also have the same property.
Vicent Cholvi, Antonio Fernández 0001, Ernesto Jiménez, Pilar Manzano-Hernandez, Michel Raynal
Comput. J.5
2010 The k-simultaneous consensus problem
Yehuda Afek, Eli Gafni, Sergio Rajsbaum, Michel Raynal, Corentin Travers
Distributed Comput.4
2010 A simple proof of the necessity of the failure detector Sigma to implement an atomic register in asynchronous message-passing systems
François Bonnet 0001, Michel Raynal
Inf. Process. Lett.2
2010 Eventual Leader Election with Weak Assumptions on Initial Knowledge, Communication Reliability, and Synchrony
Antonio Fernández 0001, Ernesto Jiménez, Michel Raynal
J. Comput. Sci. Technol.3
2010 Strongly Terminating Early-Stopping k-Set Agreement in Synchronous Systems with General Omission Failures
Philippe Raipin Parvédy, Michel Raynal, Corentin Travers
Theory Comput. Syst.2
2010 Narrowing power vs efficiency in synchronous set agreement: Relationship, algorithms and lower bound
Achour Mostéfaoui, Michel Raynal, Corentin Travers
Theor. Comput. Sci.2
2010 From an Asynchronous Intermittent Rotating Star to an Eventual Leader
abstract
Considering an asynchronous system made up of n processes and where up to t of them can crash, finding the weakest assumption that such a system has to satisfy for a common leader to be eventually elected is one of the holy grail quests of fault-tolerant asynchronous computing. This paper is a step in that direction. It has two contributions. Considering a simple and general asynchronous system model where processes generate asynchronous pulses during which they send and receive messages, it first introduces an additional assumption that allows to elect an eventual leader in all the runs that satisfy that assumption. That assumption is captured by the notion of asynchronous intermittent rotating t-star. An x-star is made up of one process p (the center of the star) plus a sequence of sets of x processes (the successive points of the star), which satisfies some properties. Intuitively, the intermittent rotating t-star assumption means that there are a process p, a subset of pulse numbers pn, and associated sets of processes Q(pn) such that each process of Q(pn) receives from p a message sent in pulse pn in a timely manner or among the first (n-t) messages tagged pn it ever receives. The t-star is called rotating because the set Q(pn) is allowed to change with pn; it is intermittent because it can disappear during finite periods; it is asynchronous because the points of a star are not required to be simultaneously at the same pulse. (This assumption combines and generalizes several synchrony and time-free assumptions that have been previously proposed to elect an eventual leader, e.g., eventual t-source, eventual t-moving source, and message pattern assumption.) The second contribution of the paper is an algorithm that eventually elects a common leader in the systems that satisfy the asynchronous intermittent rotating t-star assumption. This algorithm enjoys, among others, two noteworthy properties. First, from a design point of view, it is simple. Second, from a cost point of view, only the pulse numbers increase without bound. This means that, even in infinite executions, be links timely or not (or have the corresponding sender crashed or not), all the other local variables (including the timers) and message fields have a finite domain.
Antonio Fernández 0001, Michel Raynal
IEEE Trans. Parallel Distributed Syst.2
2009 Shared Memory Synchronization in Presence of Failures: An Exercise-Based
abstract
In the recent past, lots of papers have addressed synchronization in asynchronous shared memory systems prone to process crashes. Unfortunately, to date, nearly all these results have appeared only in theory-oriented journals and conferences, very few being presented and studied in textbooks. This aim of this paper is to give a flavor of a few of these fundamental results. To that end, it considers three problems and presents solutions proposed to solve them, emphasizing the basic concepts and techniques these solutions rely on. These problems have been selected because they address distinct facets of synchronization in presence of failures. So, the spirit of this introductory paper is mainly pedagogical (with an algorithmic taste).
Michel Raynal
CISIS1
2009 Implementing a Register in a Dynamic Distributed System
abstract
Providing distributed processes with concurrent objects is a fundamental service that has to be offered by any distributed system. The classical shared read/write register is one of the most basic ones. Several protocols have been proposed that build an atomic register on top of an asynchronous message-passing system prone to process crashes. In the same spirit, this paper addresses the implementation of a regular register (a weakened form of an atomic register) in an asynchronous dynamic message-passing system. The aim is here to cope with the net effect of the adversaries that are asynchrony and dynamicity (the fact that processes can enter and leave the system). The paper focuses on the class of dynamic systems the churn rate c of which is constant. It presents two protocols, one applicable to synchronous dynamic message passing systems, the other one to eventually synchronous dynamic systems. Both protocols rely on an appropriate broadcast communication service (similar to a reliable broadcast). Each requires a specific constraint on the churn rate c. Both protocols are first presented in an as intuitive as possible way, and are then proved correct.
Roberto Baldoni, Silvia Bonomi, Anne-Marie Kermarrec, Michel Raynal
ICDCS4
2009 Brief announcement: the price of anonymity: optimal consensus despite asynchrony, crash and anonymity
abstract
This paper proposes the first study of the consensus problem in the anonymous crash-prone message-passing systems.
François Bonnet 0001, Michel Raynal
PODC2
2009 Brief announcement: virtual world consistency: a new condition for STM systems
abstract
This BA presents a general consistency condition for software transactionnal memories.
Damien Imbs, José Ramón González de Mendívil, Michel Raynal
PODC3
2009 Regular Register: An Implementation in a Churn Prone Environment
Roberto Baldoni, Silvia Bonomi, Michel Raynal
SIROCCO3
2009 A Versatile STM Protocol with Invisible Read Operations That Satisfies the Virtual World Consistency Condition
Damien Imbs, Michel Raynal
SIROCCO2
2009 Looking for the Weakest Failure Detector for k-Set Agreement in Message-Passing Systems: Is ${\it \Pi}_k${\it \Pi}_k the End of the Road?
François Bonnet 0001, Michel Raynal
SSS2
2009 Visiting Gafni's Reduction Land: From the BG Simulation to the Extended BG Simulation
Damien Imbs, Michel Raynal
SSS2
2009 The Price of Anonymity: Optimal Consensus Despite Asynchrony, Crash and Anonymity
François Bonnet 0001, Michel Raynal
DISC2
2009 Help When Needed, But No More: Efficient Read/Write Partial Snapshot
Damien Imbs, Michel Raynal
DISC2
2009 A note on atomicity: Boosting Test&Set to solve consensus
Damien Imbs, Michel Raynal
Inf. Process. Lett.2
2009 Conditions for Set Agreement with an Application to Synchronous Systems
François Bonnet 0001, Michel Raynal
J. Comput. Sci. Technol.2
2009 Revisiting simultaneous consensus with crash failures
Yoram Moses, Michel Raynal
J. Parallel Distributed Comput.2
2009 From adaptive renaming to set agreement
Eli Gafni, Achour Mostéfaoui, Michel Raynal, Corentin Travers
Theor. Comput. Sci.3
2009 On the Fly Estimation of the Processes that Are Alive in an Asynchronous Message-Passing System
abstract
It is well known that in an asynchronous system where processes are prone to crash, it is impossible to design a protocol that provides each process with the set of processes that are currently alive. Basically, this comes from the fact that it is impossible to distinguish a crashed process from a process that is very slow or with which communications are very slow. Nevertheless, designing protocols that provide the processes with good approximations of the set of processes that are currently alive remains a real challenge in fault-tolerant-distributed computing. This paper proposes such a protocol, plus a second protocol that allows to cope with heterogeneous communication networks. These protocols consider a realistic computation model where the processes are provided with nonsynchronized local clocks and a function \alpha () that takes a local duration \Delta as a parameter, and returns an integer that is an estimate of the number of processes that could have crashed during that duration \Delta. A simulation-based experimental evaluation of the proposed protocols is also presented. These experiments show that the protocols are practically relevant.
Achour Mostéfaoui, Michel Raynal, Gilles Trédan
IEEE Trans. Parallel Distributed Syst.2
2009 Eventual Clusterer: A Modular Approach to Designing Hierarchical Consensus Protocols in MANETs
abstract
This paper proposes a modular approach to the design of hierarchical consensus protocols for the mobile ad hoc network with a static and known set of hosts. A two-layer hierarchy is imposed on the network by grouping mobile hosts into clusters, each with a clusterhead. The messages from and to the hosts in the same cluster are merged/unmerged by the clusterhead so as to reduce the message cost and improve the scalability. The proposed modular approach separates the concerns of clustering hosts from achieving consensus. A clustering function, called eventual clusterer (denoted as diamC), is designed for constructing and maintaining the two-layer hierarchy. Similar to unreliable failure detectors, diamC greatly facilitates the design of hierarchical protocols by providing the fault-tolerant clustering function transparently. We propose an implementation of diamC based on the failure detector diamS. Using diamC, we design a new hierarchical consensus protocol. As shown by the performance evaluation results, the proposed consensus protocol can save both message cost and time cost. Our proposed modular design is therefore effective and can lead to efficient solutions to achieving consensus in mobile ad hoc networks.
Weigang Wu, Jiannong Cao 0001, Michel Raynal
IEEE Trans. Parallel Distributed Syst.3
2008 Synchronization is Coming Back, But is it the Same?
abstract
This invited talk surveys notions related to synchronization in presence of asynchrony and failures. To the author knowledge, there is currently no textbook in which these notions are pieced together, unified, and presented in a homogeneous way. This talk (that pretends to be neither exhaustive, nor objective) is only a first endeavor in that direction. The notions that are presented and discussed are listed in the keyword list.
Michel Raynal
AINA1
2008 The Iterated Restricted Immediate Snapshot Model
Sergio Rajsbaum, Michel Raynal, Corentin Travers
COCOON2
2008 Conditions for Set Agreement with an Application to Synchronous Systems
abstract
The k-set agreement problem is a generalization of the consensus problem: considering a system made up of n processes where each process proposes a value, each non-faulty process has to decide a value such that a decided value is a proposed value, and no more than k different values are decided. While this problem cannot be solved in an asynchronous system prone to t process crashes when t \geq k, it can always be solved in a synchronous system; \lfloor \frac{t}{k} \rfloor +1 is then a lower bound on the number of rounds (consecutive communication steps) for the non-faulty processes to decide. The {\it condition-based} approach has been introduced in the consensus context. Its aim was to both circumvent the consensus impossibility in asynchronous systems, and allow for more efficient consensus algorithms in synchronous systems. This paper addresses the condition-based approach in the context of the k-set agreement problem. It has two main contributions. The first is the definition of a framework that allows defining conditions suited to the \ell$-set agreement problem. More precisely, a condition is defined as a set of input vectors such that each of its input vectors can be seen as "encoding" \ell values, namely, the values that can be decided from that vector. A condition is characterized by the parameters t, \ell, and a parameter denoted d such that the greater d+\ell, the least constraining the condition (i.e., it includes more and more input vectors when d+\ell increases, and there is a condition that includes all the input vectors when d+\ell≫t$). The conditions characterized by the triple of parameters t, d and \ell define the class of conditions denoted ${\cal S}_t^{d,\ell}$, $0\leq d\leq t$, $1\leq \ell \leq n-1 $. The properties of the sets ${\cal S}_t^{d,\ell}$ are investigated, and it is shown that they have a lattice structure. The second contribution is a generic synchronous k-set agreement algorithm based on a condition $C\in {\cal S}_t^{d,\ell}$, i.e., a condition suited to the $\ell$-set agreement problem, for $\ell \leq k$. This algorithm requires at most $\left\lfloor \frac{d-1+\ell}{k} \right\rfloor +1$ rounds when the input vector belongs to $C$, and $\left\lfloor \frac{t}{k} \right\rfloor +1 rounds otherwise. (Interestingly, this algorithm includes as particular cases the classical synchronous k-set agreement algorithm that requires $\left\lfloor \frac{t}{k} \right\rfloor+1 rounds (case $d=t$ and $\ell=1$), and the synchronous consensus condition-based algorithm that terminates in d+1 rounds when the input vector belongs to the condition, and in t+1 rounds otherwise (case $k=\ell=1$).)
François Bonnet 0001, Michel Raynal
ICDCS2
2008 On Modeling Fault Tolerance of Gossip-Based Reliable Multicast Protocols
abstract
Gossiping has been widely used for disseminating data in large scale networks. Existing works have mainly focused on the design of gossip-based protocols but few have been reported on developing models for analyzing the fault tolerance property of these protocols. In this paper, we propose a general gossiping algorithm and develop a mathematical model based on generalized random graphs for evaluating the reliability of gossiping, i.e., to what extent gossip-based protocols can tolerate node failures, yet guarantee the specified message delivery. We analytically derive the maximum ratio of failed nodes that can be tolerated without reducing the required degree of reliability. We also investigate the impact of the parameters, namely the fanout distribution and the nonfailed member ratio, on the protocol reliability. Simulations have been carried out to validate the effectiveness of our analytic model in terms of the reliability of gossiping and the success of gossiping. The results obtained can be used to guide the design of fault tolerant gossip-based protocols.
Xiaopeng Fan 0002, Jiannong Cao 0001, Weigang Wu, Michel Raynal
ICPP4
2008 On the Solvability of Anonymous Partial Grids Exploration by Mobile Robots
Roberto Baldoni, François Bonnet 0001, Alessia Milani, Michel Raynal
OPODIS4
2008 A Lock-Based STM Protocol That Satisfies Opacity and Progressiveness
Damien Imbs, Michel Raynal
OPODIS2
2008 Looking for the optimal conditions for solving set agreement
abstract
This BA extends the condition-based approach to the ll-set agreement problem.
François Bonnet 0001, Michel Raynal
PODC2
2008 From anarchy to geometric structuring: the power of virtual coordinates
abstract
This note define self-structuring in a large-scale networked system as the ability of the participating entities to collaboratively impose a geometric structure to the network. This refers to assigning virtual coordinates to participating entities and to dividing the entities in several partitions, in such a way that each entity knows to which partition it belongs.
Anne-Marie Kermarrec, Achour Mostéfaoui, Michel Raynal, Gilles Trédan, Aline Carneiro Viana
PODC3
2008 Brief Announcement: On the Solvability of Anonymous Partial Grids Exploration by Mobile Robots
Roberto Baldoni, François Bonnet 0001, Alessia Milani, Michel Raynal
DISC4
2008 No Double Discount: Condition-Based Simultaneity Yields Limited Gain
Yoram Moses, Michel Raynal
DISC2
2008 On the computability power and the robustness of set agreement-oriented failure detector classes
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Corentin Travers
Distributed Comput.3
2008 Anonymous graph exploration without collision by mobile robots
Roberto Baldoni, François Bonnet 0001, Alessia Milani, Michel Raynal
Inf. Process. Lett.4
2008 An impossibility about failure detectors in the iterated immediate snapshot model
Sergio Rajsbaum, Michel Raynal, Corentin Travers
Inf. Process. Lett.2
2008 Using asynchrony and zero degradation to speed up indulgent consensus protocols
Weigang Wu, Jiannong Cao 0001, Jin Yang 0005, Michel Raynal
J. Parallel Distributed Comput.4
2008 The Combined Power of Conditions and Information on Failures to Solve Asynchronous Set Agreement
abstract
To cope with the impossibility of solving agreement problems in asynchronous systems made up of n processes and prone to t process crashes, system designers tailor their algorithms to run fast in “normal” circumstances. Two orthogonal notions of “normality” have been studied in the past through failure detectors that give processes information about process crashes, and through conditions that restrict the inputs to an agreement problem. This paper investigates how the two approaches can benefit from each other to solve the k-set agreement problem, where processes must agree on at most k of their input values (when $k=1$ we have the famous consensus problem). It proposes novel failure detectors for solving k-set agreement and a protocol that combines them with conditions, establishing a new bridge among asynchronous, synchronous, and partially synchronous systems with respect to agreement problems. The paper also proves a lower bound when solving the k-set agreement problem with a condition.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Corentin Travers
SIAM J. Comput.3
2007 A Universal Construction for Concurrent Objects
abstract
A concurrent object is an object that can be concurrently accessed by several processes. A wait-free implementation of an object is such that any operation issued by a non-faulty process terminates in a finite number of its own steps, whatever the behavior of the other processes (that can be very slow or even have crashed). An object type is universal if objects of that type, together with atomic registers, allows implementing any concurrent object defined by a sequential specification. A universal construction is a wait-free algorithm, based only on atomic registers and universal objects, that, given any sequential object type T, provides the processes with a wait-free concurrent object of the type T. In a famous paper (titled "Wait-free synchronization") Herlihy has shown that consensus objects are universal, and has presented a consensus-based universal construction. We present here a new universal construction. That construction, that is built incrementally, is particularly simple. While, in addition to consensus objects, Herlihy's universal construction uses low-level objects such as pointers, the design of the construction presented here is based on the simple and well-known state machine replication paradigm. Its proof is also simple and consequently allows to better understand not only the power of consensus objects but also the subtleties of wait-free computations and the way the consensus objects allow coping with both process failures and non-determinism. In that sense, this paper has a pedagogical flavor.
Rachid Guerraoui, Michel Raynal
ARES2
2007 Workshop on Dependable Application Support for Self-Organizing Networks (DASSON 2007)
abstract
This report gives an overview of the workshop on "Dependable Application Support for Self-Organising Networks" held in conjunction with DSN 2007. The principal objective of the workshop is to facilitate a forum for researchers to explore, examine and address dependability related challenges in hosting distributed applications in self-organised networks such as MANETs, sensor and P2P networks.
Paul D. Ezhilchelvan, Michel Raynal, Ajoy K. Datta
DSN2
2007 Electing an Eventual Leader in an Asynchronous Shared Memory System
abstract
This paper considers the problem of electing an eventual leader in an asynchronous shared memory system. While this problem has received a lot of attention in message- passing systems, very few solutions have been proposed for shared memory systems. As an eventual leader cannot be elected in a pure asynchronous system prone to process crashes, the paper first proposes to enrich the asynchronous system model with an additional assumption. That assumption, denoted AWB, requires that after some time (1) there is a process whose write accesses to some shared variables are timely, and (2) the timers of the other processes are asymptotically well-behaved. The asymptotically well-behaved timer notion is a new notion that generalizes and weakens the traditional notion of timers whose durations are required to monotonically increase when the values they are set to increase. Then, the paper presents two A WB-based algorithms that elect an eventual leader. Both algorithms are independent of the value of t (the maximal number of processes that may crash). The first algorithm enjoys the following noteworthy properties: after some time only the elected leader has to write the shared memory, and all but one shared variables have a bounded domain, be the execution finite or infinite. This algorithm is consequently optimal with respect to the number of processes that have to write the shared memory. The second algorithm enjoys the following property: all the shared variables have a bounded domain. This is obtained at the following additional price: all the processes are required to forever write the shared memory. A theorem is proved which states that this price has to be paid by any algorithm that elects an eventual leader in a bounded shared memory model. This second algorithm is consequently optimal with respect to the number of processes that have to write in such a constrained memory model. In a very interesting way, these algorithms show an inherent tradeoff relating the number of processes that have to write the shared memory and the bounded/unbounded attribute of that memory.
Antonio Fernández 0001, Ernesto Jiménez, Michel Raynal
DSN3
2007 Distributed Slicing in Dynamic Systems
abstract
Peer to peer (P2P) systems are moving from application specific architectures to a generic service oriented design philosophy. This raises interesting problems in connection with providing useful P2P middleware services capable of dealing with resource assignment and management in a large-scale, heterogeneous and unreliable environment. The slicing service, has been proposed to allow for an automatic partitioning of P2P networks into groups (slices) that represent a controllable amount of some resource and that are also relatively homogeneous with respect to that resource. In this paper we propose two gossip-based algorithms to solve the distributed slicing problem. The first algorithm speeds up an existing algorithm sorting a set of uniform random numbers. The second algorithm statistically approximates the rank of nodes in the ordering. The scalability, efficiency and resilience to dynamics of both algorithms rely on their gossip-based models. These algorithms are proved viable theoretically and experimentally.
Antonio Fernández 0001, Vincent Gramoli, Ernesto Jiménez, Anne-Marie Kermarrec, Michel Raynal
ICDCS5
2007 A Timing Assumption and a t-Resilient Protocol for Implementing an Eventual Leader Service in Asynchronous Shared Memory Systems
abstract
While electing an eventual common leader, despite process crashes, in a shared memory system where the processes communicate only by reading and writing shared registers is possible when the processes progress synchronously, this problem becomes impossible to solve as soon as the processes can progress in a fully asynchronous way. So, an important problem consists in finding additional behavioral assumptions that are, at the same time, "as weak as possible" (in order they are practically always satisfied), and "strong enough" in order to allow implementing an eventual leader service despite the net effect of asynchrony and failures. This paper focuses on this dilemma. More explicitly, it investigates a timing assumption that allows implementing an eventual leader in presence of partial asynchrony and process crashes. The proposed timing assumptions are particularly weak. They are the following: after some time (i) there is a process that behaves synchronously, and (ii) (t - f) other processes have timers that work correctly (t is the maximal number of processes that may crash, and f the actual number of process crashes; a timer works incorrectly when it expires too early with respect to the value it has been set). Then, the paper proposes a t-resilient protocol that elects an eventual common leader in any shared memory system that satisfies the previous assumption. Interestingly, this protocol is based on simple design principles.
Antonio Fernández 0001, Ernesto Jiménez, Michel Raynal, Gilles Trédan
ISORC3
2007 A Dual-Token-Based Fault Tolerant Mutual Exclusion Algorithm for MANETs
Weigang Wu, Jiannong Cao 0001, Michel Raynal
MSN3
2007 Eventual Leader Service in Unreliable Asynchronous Systems: Why? How?
abstract
Providing processes with an eventual leader service is an important issue when one has to design and implement a middleware layer on top of a failure-prone asynchronous distributed system. This invited lecture investigates this problem. It first shows that such a service cannot be built if the underlying system is fully asynchronous. Then, the paper visits several additional behavioral assumptions that have been proposed in the literature to cope with this impossibility and presents corresponding eventual leader election protocols. This lecture can be seen as a guided tour of the eventual leader service problem, whose aim is to benefit researchers and system engineers working in distributed middleware built on top of asynchronous networks.
Michel Raynal
NCA1
2007 From an Intermittent Rotating Star to a Leader
Antonio Fernández 0001, Michel Raynal
OPODIS2
2007 Small-World Networks: From Theoretical Bounds to Practical Systems
François Bonnet 0001, Anne-Marie Kermarrec, Michel Raynal
OPODIS3
2007 Timed Quorum Systems for Large-Scale and Dynamic Environments
Vincent Gramoli, Michel Raynal
OPODIS2
2007 Failure detectors are schedulers
abstract
No abstract available.
Alejandro Cornejo, Sergio Rajsbaum, Michel Raynal, Corentin Travers
PODC3
2007 From an intermittent rotating star to a leader
abstract
No abstract available.
Antonio Fernández 0001, Michel Raynal
PODC2
2007 The Eventual Leadership in Dynamic Mobile Networking Environments
abstract
Eventual leadership has been identified as a basic building block to solve synchronization or coordination problems in distributed computing systems. However, it is a challenging task to implement the eventual leadership facility, especially in dynamic distributed systems, where the global system structure is unknown to the processes and can vary over time. This paper studies the implementation of a leadership facility in infrastructured mobile networks, where an unbounded set of mobile hosts arbitrarily move in the area covered by fixed mobile support stations. Mobile hosts can crash and suffer from disconnections. We develop an eventual leadership protocol based on a time-free approach. The mobile support stations exchange queries and responses on behalf of mobile hosts. With assumptions on the message exchange flow, a correct mobile host is eventually elected as the unique leader. Since no time property is assumed on the communication channels, the proposed protocol is especially effective and efficient in mobile environments, where time-based properties are difficult to satisfy due to the dynamics of the network.
Jiannong Cao 0001, Michel Raynal, Corentin Travers, Weigang Wu
PRDC2
2007 From Renaming to Set Agreement
Achour Mostéfaoui, Michel Raynal, Corentin Travers
SIROCCO2
2007 The notion of a timed register and its application to indulgent synchronization
abstract
A new type of shared object, called timed register, is proposed and used to design indulgent timing-based algorithms. A timed register generalizes the notion of an atomic register as follows: if a process invokes two consecutive operations on the same timed register which are a read followed by a write, then the write operation is executed only if it is invoked at most d time units after the read operation, where d is defined as part of the read operation. In this context, a timing-based algorithm is an algorithm whose correctness relies on the existence of a bound Δ such that any pair of consecutive constrained read and write operations issued by the same process on the same timed register are separated by at most Δ time units. An indulgent algorithm is an algorithm that always guarantees the safety properties, and ensures the liveness property as soon as the timing assumptions are satisfied. The usefulness of this new type of shared object is demonstrated by presenting simple and elegant indulgent timing-based algorithms that solve the mutual exclusion, l-exclusion, adaptive renaming,test&set, and consensus problems. Interestingly, timed registers are universal objects in systems with process crashes and transient timing failures (i.e., they allow building any concurrent object with a sequential specification). The paper also suggests connections with schedulers and contention managers.
Michel Raynal, Gadi Taubenfeld
SPAA1
2007 Test & Set, Adaptive Renaming and Set Agreement: a Guided Visit to Asynchronous Computability
abstract
An important issue in fault-tolerant asynchronous computing is the respective power of an object type with respect to another object type. This question has received a lot of attention, mainly in the context of the consensus problem where a major advance has been the introduction of the consensus number notion that allows ranking the synchronization power of base object types (atomic registers, queues, test&set objects, compare&swap objects, etc.) with respect to the consensus problem. This has given rise to the well-known Herlihy's hierarchy. Due to its very definition, the consensus number notion is irrelevant for studying the respective power of object types that are too weak to solve consensus for an arbitrary number of processes (these objects are usually called subconsensus objects). Considering an asynchonous system made up of n processes prone to crash, this paper addresses the power of such object types, namely, the k-test&set object type, the k-set agreement object type, and the adaptive M-renaming object type for M = 2p - [P/N] and M = min(2p - 1,p + k - 1), where p < n is the number of processes that want to acquire a new name. It investigates their respective power stating the necessary and sufficient conditions to build objects of any of these types from objects of any of the other types. More precisely, the paper shows that (1) these object types define a strict hierarchy when k ne1,n - 1, (2) they all are equivalent when k = n - 1, and (3) they all are equivalent except k-set agreement that is stronger when k = 1 ne n - 1 (a side effect of these results is that that the consensus number of the renaming problem is 2.)
Eli Gafni, Michel Raynal, Corentin Travers
SRDS2
2007 The Eventual Clusterer Oracle and Its Application to Consensus in MANETs
abstract
This paper studies the design of hierarchical consensus protocols for mobile ad hoc networks. A two-layer hierarchy is imposed on the mobile hosts by grouping them into clusters, each with a clusterhead. The messages from and to the hosts in the same cluster are merged/unmerged by the clusterhead so as to reduce the message cost and improve the scalability. We adopt a modular method in the design, separating clustering from achieving consensus using the clusters. The clustering function, named eventual clusterer (denoted as diamC), is designed to construct a cluster-based hierarchy over the mobile hosts in the network. Since diamC provides the fault tolerant clustering function transparently, it can be used as a new oracle (i.e. an abstract tool to provide some kind of information about the state of the system) for the design of hierarchical consensus protocols. Based on diamC, we design a new consensus protocol, which can significantly reduce the message cost of achieving consensus. We also propose an implementation of the diamC oracle based on the failure detector diamS.
Weigang Wu, Jiannong Cao 0001, Michel Raynal
SRDS3
2007 A Subjective Visit to Selected Topics in Distributed Computing
Michel Raynal
DISC1
2007 DISC at Its 20th Anniversary (Stockholm, 2006)
Michel Raynal, Sam Toueg, Shmuel Zaks
DISC1
2007 The Alpha of Indulgent Consensus
abstract
This paper presents a simple framework unifying a family of consensus algorithms that can tolerate process crash failures and asynchronous periods of the network, also called indulgent consensus algorithms. Key to the framework is a new abstraction we introduce here, called Alpha, and which precisely captures consensus safety. Implementations of Alpha in shared memory, storage area network, message passing and active disk systems are presented, leading to directly derived consensus algorithms suited to these communication media. The paper also considers the case where the number of processes is unknown and can be arbitrarily large.
Rachid Guerraoui, Michel Raynal
Comput. J.2
2007 From omega to Omega: A simple bounded quiescent reliable broadcast-based transformation
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Corentin Travers
J. Parallel Distributed Comput.3
2007 Asynchronous Agreement and Its Relation with Error-Correcting Codes
abstract
The condition-based approach identifies sets of input vectors, called conditions, for which it is possible to design an asynchronous protocol solving a distributed problem despite process crashes. This paper establishes a direct correlation between distributed agreement problems and error-correcting codes. In particular, crash failures in distributed agreement problems correspond to erasure failures in error-correcting codes and Byzantine and value domain faults correspond to corruption errors. This correlation is exemplified by concentrating on two well-known agreement problems, namely, consensus and interactive consistency, in the context of the condition-based approach. Specifically, the paper presents the following results: first, it shows that the conditions that allow interactive consistency to be solved despite fccrashes and fcvalue domain faults correspond exactly to the set of error-correcting codes capable of recovering from fcerasures and fccorruptions. Second, the paper proves that consensus can be solved despite fccrash failures if the condition corresponds to a code whose Hamming distance is fc+ 1 and Byzantine consensus can be solved despite fbByzantine faults if the Hamming distance of the code is 2 fb+ 1. Finally, the paper uses the above relations to establish several results in distributed agreement that are derived from known results in error-correcting codes and vice versa.
Roy Friedman 0001, Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
IEEE Trans. Computers4
2007 Design and Performance Evaluation of Efficient Consensus Protocols for Mobile Ad Hoc Networks
abstract
Designing protocols for solving the consensus problem faces new challenges in mobile computing environments. Among others, how we can achieve message efficiency for saving resource consumption has been the focus of research. In this paper, we present the HC protocol, a message efficient consensus protocol for MANETs. We consider the widely used system model where the hosts fail by crashes and the system is equipped with Chandra-Toueg's unreliable failure detectors. Unlike existing consensus protocols, the HC protocol uses a two-layer hierarchy based on clusters to achieve message efficiency. The messages from and to the hosts in the same cluster are merged so as to reduce the message cost. However, adding such a hierarchy is not trivial. Due to host movements and failures, the hierarchy changes from time to time and this may cause message loss. In designing HC, we also propose methods to handle such message losses. Extensive simulations have been carried out to evaluate and compare the performance of the HC protocol and similar protocols in a MANET environment. Simulation results show that, in most cases, our protocol can significantly reduce both the message cost and time cost. With increases in the system scale or the percentage of faulty hosts, the advantage of our protocol becomes more obvious.
Weigang Wu, Jiannong Cao 0001, Jin Yang 0005, Michel Raynal
IEEE Trans. Computers4
2007 Preface
Andrzej Pelc, David Peleg, Michel Raynal
Theor. Comput. Sci.3
2007 An Adaptive Programming Model for Fault-Tolerant Distributed Computing
abstract
The capability of dynamically adapting to distinct runtime conditions is an important issue when designing distributed systems where negotiated quality of service (QoS) cannot always be delivered between processes. Providing fault tolerance for such dynamic environments is a challenging task. Considering such a context, this paper proposes an adaptive programming model for fault-tolerant distributed computing, which provides upper-layer applications with process state information according to the current system synchrony (or QoS). The underlying system model is hybrid, composed by a synchronous part (where there are time bounds on processing speed and message delay) and an asynchronous part (where there is no time bound). However, such a composition can vary over time, and, in particular, the system may become totally asynchronous (e.g., when the underlying system QoS degrade) or totally synchronous. Moreover, processes are not required to share the same view of the system synchrony at a given time. To illustrate what can be done in this programming model and how to use it, the consensus problem is taken as a benchmark problem. This paper also presents an implementation of the model that relies on a negotiated quality of service (QoS) for communication channels
Sérgio Gorender, Raimundo José de Araújo Macêdo, Michel Raynal
IEEE Trans. Dependable Secur. Comput.3
2007 On the Respective Power of *P and *S to Solve One-Shot Agreement Problems
abstract
Unreliable failure detectors are abstract devices that, when added to asynchronous distributed systems, enable solving distributed computing problems (e.g., consensus) that otherwise would be impossible to solve in these systems. This paper focuses on two classes of failure detectors defined by Chandra and Toueg, namely, the classes denoted diamP (eventually perfect) and diamS (eventually strong). Both classes include failure detectors that eventually detect permanently all process crashes, but while the failure detectors of diamP eventually make no erroneous suspicions, the failure detectors of diamS are only required to eventually not suspect a single correct process. Informally, in a one-shot agreement problem, a new problem instance is created each time the processes propose new values to be decided on (e.g., consensus is one-shot). In such a context, this paper addresses the following question related to the comparative power of these classes, namely: "Are there one-shot agreement problems that can be solved in asynchronous distributed systems with reliable links but prone to process crash failures augmented with op, but cannot be solved when those systems are augmented with diamS?" Surprisingly, the paper shows that the answer to this question is "no." An important consequence of this result is that diamP cannot be the weakest class of failure detectors that enables solving one-shot agreement problems in unreliable asynchronous distributed systems
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Parallel Distributed Syst.3
2006 From Failure Detectors with Limited Scope Accuracy to System-wide Leadership
abstract
A failure detector is a device that provides the processes with information on failures. The accuracy property of a failure detector defines the type of mistakes it is not allowed to make. The limited scope of the accuracy property restricts it to only a part of the system. /spl diams/S/sub k/ is a class of unreliable failure detectors with a limited scope accuracy. Eventually each process that crashes is suspected by every correct process, and there is a time after which some correct process is never suspected by only k processes. An eventual leader facility (usually denoted /spl Omega/)is a device that eventually provides all the processes with the identity of one of them that is correct. Such a facility is used as a basic service in a lot of fault-tolerant distributed protocols (e.g., asynchronous consensus protocols). This paper proposes a protocol that builds an eventual leader service from any unreliable failure detector of the class /spl diams/S/sub t+1/ where t is the maximum number of processes that can crash during a run. The fact that /spl diams/S/sub t+1/ is easier to build than /spl diams/S or /spl Omega/ and the design simplicity of the proposed protocol makes it attractive.
Achour Mostéfaoui, Michel Raynal, Corentin Travers, Sergio Rajsbaum
AINA (1)2
2006 Eventual Leader Election with Weak Assumptions on Initial Knowledge, Communication Reliability, and Synchrony
abstract
This paper considers the eventual leader election problem in asynchronous message-passing systems where an arbitrary number t of processes can crash (t
Antonio Fernández 0001, Ernesto Jiménez, Michel Raynal
DSN3
2006 The Power and Limit of Adding Synchronization Messages for Synchronous Agreement
abstract
This paper investigates the use of additional synchronization messages in round-based message-passing synchronous systems. It first presents a synchronous computation model allowing a process to send such messages. The difference with respect to the traditional round-based synchronous model lies in the sending phase, where a process can first send a data message to each other process, and then, without a break, a synchronization message (their sendings can be pipelined). This model is suited to the class of local area networks where communication channels are reliable. (It is not for networks where unreliable communication requires message retransmission.) To illustrate the model, the paper presents a uniform consensus algorithm suited to this model. This algorithm, based on the rotating coordinator paradigm, allows the processes to decide in at most f + 1 rounds where f is the actual number of processes that crash in the corresponding run. (This improves the f + 2 lower bound of the traditional synchronous model.) In addition to its efficiency, the algorithm enjoys another first class property, namely, design simplicity. The paper focuses also on lower bound results, and shows that any uniform consensus algorithm designed for the proposed model, requires at least f + 1 rounds in the worst case. The proposed algorithm is consequently optimal. In that sense the paper has to be seen as an investigation of both the power and the limit of adding synchronization messages to synchronous systems built on top of local networks with reliable communication
Jiannong Cao 0001, Michel Raynal, Xianbing Wang, Weigang Wu
ICPP2
2006 The Committee Decision Problem
Eli Gafni, Sergio Rajsbaum, Michel Raynal, Corentin Travers
LATIN3
2006 In Search of the Holy Grail: Looking for the Weakest Failure Detector for Wait-Free Set Agreement
Michel Raynal, Corentin Travers
OPODIS1
2006 A Hierarchical Consensus Protocol for Mobile Ad Hoc Networks
abstract
Mobile ad hoc networks (MANETs) raise new challenges in designing protocols for solving the consensus problem. Among the others, how to design message efficient protocols so as to save resource consumption, has been the focus of research. In this paper, we present the design of such an efficient consensus protocol. We consider the system model for MANETs with host crashes, but equipped with Chandra-Toueg's unreliable failure detectors of class /spl diams/P. At most f hosts can crash where f < n/2 (n is the total number of the hosts). The protocol adopts the coordinator rotation paradigm to achieve consensus. Unlike existing consensus protocols, the proposed protocol is based on a two-layer hierarchy with hosts associated with proxies. At least f + 1 hosts act as proxies and each host is associated with one proxy host. The messages from and/or to the local hosts of the same proxy are merged so as to reduce the message cost. Moreover, the hierarchical approach can improve the scalability of the consensus protocol. Performance analysis shows that the proposed protocol can significantly save cost compared existing protocols.
Weigang Wu, Jiannong Cao 0001, Jin Yang 0005, Michel Raynal
PDP4
2006 Irreducibility and additivity of set agreement-oriented failure detector classes
abstract
Solving agreement problems (such as consensus and k-set agreement) in asynchronous distributed systems prone to process failures has been shown to be impossible. To circumvent this impossibility, distributed oracles (also called unreliable failure detectors) have been introduced. A failure detector provides information on failures, and a failure detector class is defined by a set of abstract properties that encapsulate (and hide) synchrony assumptions. Some failure detector classes have been shown to be the weakest to solve some agreement problems (e.g., Ω is the weakest class of failure detectors that allow solving the consensus problem in asynchronous systems where a majority of processes do not crash).This paper considers several failure detector classes and focuses on their additivity or their irreducibility. It mainly investigates two families of failure detector classes (denoted ◊ Sx and ◊ φy, 0≤ x, y ≤ n), shows that they can be "added" to provide a failure detector of the class Ωz (a generalization of Ω). It also characterizes the power of such an "addition", namely, ◊ Sx + ◊ φy ➝ Ωz ⇔ x+y+z>t+1, where t is the maximum number of processes that can crash in a run. As an example, the paper shows that, while ◊ St allows solving 2-set agreement (and not consensus) and ◊ φ1 allows solving t-set agreement (but not (t-1)-set agreement), their "addition" allows solving consensus. More generally, the paper studies the failure detector classes ◊ Sx, ◊ φy and Ωz, and shows which reductions among these classes are possible and which are not. The paper presents also an Ωk-based k-set agreement protocol. In that sense, it can be seen as a step toward the characterization of the weakest failure detector that allows solving the k-set agreement problem.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Corentin Travers
PODC3
2006 On the fly estimation of the processes that are alive/crashed in an asynchronous message-passing system
abstract
It is well-known that, in an asynchronous system where processes are prone to crash, it is impossible to design a protocol that provides each process with the set of processes that are currently alive. Basically, this comes from the fact that it is impossible to distinguish a crashed process from a process that is very slow or with which communications are very slow. Nevertheless, designing protocols that provide the processes with good approximations of the set of processes that are currently alive remains a real challenge in fault-tolerant distributed computing. This paper proposes such a protocol. To that end, it considers a realistic computation model where the processes are provided with non-synchronized local clocks and a function alpha(). That function takes a local duration as a parameter, and returns an integer that is an estimate of the number of processes that can crash during that duration. A simulation-based experimental evaluation of the protocol is also presented. The experiments show that the protocol is practically relevant
Achour Mostéfaoui, Michel Raynal, Gilles Trédan
PRDC2
2006 Synchronous Set Agreement: a Concise Guided Tour (including a new algorithm and a list of open problems)
abstract
The k-set agreement problem is a paradigm of coordination problems encountered in distributed computing. The parameter k defines the coordination degree we are interested in. (The case k=1 corresponds to the well-known uniform consensus problem.) More precisely, the k-set agreement problem considers a system made up of n processes where each process proposes a value. It requires that each non-faulty process decides a value such that a decided value is a proposed value, and no more than k different values are decided. This paper visits the k-set agreement problem in synchronous systems where up to t processes can experience failures. Three failure models are explored: the crash failure model, the send omission failure model, and the general omission failure model. Lower bounds and protocols are presented for each model. Open problems for the general omission failure model are stated. This paper can be seen as a short tutorial whose aim is to make the reader familiar with the k-set agreement problem in synchrony models with increasing fault severity. An important concern of the paper is simplicity. In addition to its survey flavor, several results and protocols that are presented are new
Michel Raynal, Corentin Travers
PRDC1
2006 Strongly Terminating Early-Stopping k-Set Agreement in Synchronous Systems with General Omission Failures
Philippe Raipin Parvédy, Michel Raynal, Corentin Travers
SIROCCO2
2006 Exploring Gafni's Reduction Land: From Omegak to Wait-Free Adaptive (2p-[p/k])-Renaming Via k-Set Agreement
Achour Mostéfaoui, Michel Raynal, Corentin Travers
DISC2
2006 Synchronous condition-based consensus
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
Distributed Comput.3
2006 Time-Free and Timer-Based Assumptions Can Be Combined to Obtain Eventual Leadership
abstract
Leader-based protocols rest on a primitive able to provide the processes with the same unique leader. Such protocols are very common in distributed computing to solve synchronization or coordination problems. Unfortunately, providing such a primitive is far from being trivial in asynchronous distributed systems prone to process crashes. (It is even impossible in fault-prone purely asynchronous systems.) To circumvent this difficulty, several protocols have been proposed that build a leader facility on top of an asynchronous distributed system enriched with additional assumptions. The protocols proposed so far consider either additional assumptions based on synchrony or additional assumptions on the pattern of the messages that are exchanged. Considering systems with n processes and up to f process crashes, 1lesf
Achour Mostéfaoui, Michel Raynal, Corentin Travers
IEEE Trans. Parallel Distributed Syst.2
2005 A Simple Protocol Offering Both Atomic Consistent Read Operations and Sequentially Consistent Read Operations
abstract
A concurrent object is an object that can be concurrently accessed by several processes. Two well-known consistency criteria for such objects are atomic consistency (also called linearizability) and sequential consistency. Both criteria require that all the operations on the concurrent objects can be totally ordered in such a way that each read operation obtains the last value written into the corresponding object. They differ in the meaning of the word "last" that refers to physical time for atomic consistency, and to logical time for sequential consistency. This paper investigates the merging of these consistency criteria in a multiprocess program. The proposed combination offers two read operations to the processes, namely, an atomic read operation and a sequentially consistent read operation. While the first provides a process with the last "physical" value of an object, the second provides it with a value that is approximate with respect to real-time but whose semantics is perfectly well defined. A protocol that implements the combination on top of an asynchronous distributed system is described. The protocol provides a better understanding of the similarities and differences between these consistency criteria. Moreover, the protocol is generic in the sense that it can be tailored to provide only one of these consistency criteria.
Michel Raynal, Matthieu Roy, Ciprian Tutu
AINA1
2005 A Hybrid and Adaptive Model for Fault-Tolerant Distributed Computing
abstract
The capability of dynamically adapting to distinct runtime conditions is an important issue when designing distributed systems where negotiated quality of service (QoS) cannot always be delivered between processes. Providing fault-tolerance for such dynamic environments is a challenging task. Considering such a context, this paper proposes an adaptive model for fault-tolerant distributed computing. This model encompasses both the synchronous model (where there are time bounds on processing speed and message delay) and the asynchronous model (where there is no time bound). To illustrate what can be done in this model and how to use it, the consensus problem is taken as a benchmark problem. An implementation of the model is also described. This implementation relies on a negotiated quality of service (QoS) for channels, that can be timely or untimely. Moreover, the QoS of a channel can be lost during the execution (i.e., dynamically modified from timely to untimely), thereby adding uncertainty into the system.
Sérgio Gorender, Raimundo José de Araújo Macêdo, Michel Raynal
DSN3
2005 Mixed Consistency Model: Meeting Data Sharing Needs of Heterogeneous Users
abstract
Heterogeneous users usually have different requirements as far as consistency of shared data is concerned. This paper proposes and investigates a mixed consistency model to meet this heterogeneity challenge in large scale distributed systems that support shared objects. This model allows combining strong (Sequential) consistency and weak (Causal) consistency. The paper defines the model, motivates it and proposes a protocol implementing it.
Zhiyuan Zhan, Mustaque Ahamad, Michel Raynal
ICDCS3
2005 Building Responsive TMR-Based Servers in Presence of Timing Constraints
abstract
This paper is on the construction of a fault-tolerant and responsive server subsystem in an application context where the subsystem is accessed through an asynchronous network by a large number of clients. The server is made fault-tolerant by the triple modular redundancy (TMR) technique: at least two server processes behave correctly, while the third one can behave arbitrarily. An essential requirement for process replication is that the client inputs be delivered to server replicas for processing in an identical order. Moreover, in order to cope with process' memory requirement, a time bound constraint is imposed: no client input can stay in the local memory of server process more than /spl Sigma/ units of time. Based on known technologies, two assumptions are made: (1) the network delivers a given client input to any two server processes within a known bounded time (D), and (2), there is an Ordered Timed Atomic Broadcast protocol built on top of the TMR system with timeliness A. The paper presents two results. The first is a protocol that delivers an ordered stream of client inputs, such that every client input is delivered exactly once to each correct server, thus eliminating redundant verification. It works under the assumption /spl Sigma/>D+/spl Delta/. The second is an impossibility result, namely there can be no ordering protocol when /spl Sigma/<D+/spl Delta/_, where /spl Delta/_ is the minimum timeliness of any reliable broadcast protocol that can be implemented on top of the TMR server (/spl Delta/_ < /spl Delta/).
Paul D. Ezhilchelvan, Jean-Michel Hélary, Michel Raynal
ISORC3
2005 Two Abstractions for Implementing Atomic Objects in Dynamic Systems
Roy Friedman 0001, Michel Raynal, Corentin Travers
OPODIS2
2005 Brief announcement: abstractions for implementing atomic objects in dynamic systems
abstract
No abstract available.
Roy Friedman 0001, Michel Raynal, Corentin Travers
PODC2
2005 The combined power of conditions and failure detectors to solve asynchronous set agreement
abstract
An approach to cope with the impossibility of solving agreement problems in asynchronous systems made up of n processes and prone to t process crashes is to use failure detectors. An orthogonal approach that has been used is to consider conditions that restrict the possible inputs to such a problem. This paper considers a system with both failure detectors and conditions. The aim is to identify the failure detector class that abstracts away the synchrony needed to solve k-set agreement for a given condition.Three main contributions are presented. The first is a new class of failure detectors denoted Φty, 0≤ y≤ t. The processes can invoke a primitive queryy(S) with a set of process ids S. Roughly speaking, queryy(S) returns true only when all processes in S have crashed, provided t-y
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
PODC3
2005 Intersecting Sets: a Basic Abstraction for Asynchronous Agreement Problems
abstract
Defining good abstractions is a central issue when one wants to understand the deep structure and basic principles that underlie computing mechanisms. This paper introduces a basic and particularly simple distributed computing abstraction suited to asynchronous distributed agreement problems. This abstraction, called intersecting sets, requires each process to deposit a value and allows each non-faulty process to obtain a subset of these values such that any two such sets have a non-empty intersection. This simple abstraction captures an essential part of distributed agreement problems. After having introduced and motivated this abstraction, the paper investigates its properties, its power and its benefit when solving distributed agreement problems.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
PRDC3
2005 Decision Optimal Early-Stopping k-set Agreement in Synchronous Systems Prone to Send Omission Failures
abstract
The k-set agreement problem is a generalization of the consensus problem: each process proposes a value, and each non-faulty process has to decide a value such that a decided value is a proposed value, and no more than k different values are decided. This paper focuses on the k-set agreement problem in the context of synchronous systems where up to t < n processes can experience crash or send omission failures (n being the total number of processes). The paper presents a k-set agreement protocol for this failure model (the first to our knowledge) which has two main outstanding features. (1) It provides the following early deciding and stopping property: no process decides or halts after the round min(/spl lfloor/f/k/spl rfloor/ + 2, /spl lfloor/t/k/spl rfloor/ + 1) where f is the number of actual crashes (0 /spl les/ f /spl les/ t). (2) It is decision-optimal. This new optimality criterion, suited to the omission failure model, concerns the number of processes that decide, namely, the protocol forces all the processes that do not crash to decide (regardless of whether they commit omission faults or not). It is noteworthy that each of these properties (early deciding/stopping vs decision-optimality) is not obtained at the detriment of the other. Last but not least, the protocol enjoys another first-class property, namely, simplicity.
Philippe Raipin Parvédy, Michel Raynal, Corentin Travers
PRDC2
2005 A Note on a Simple Equivalence between Round-based Synchronous and Asynchronous Models
abstract
This short paper characterizes a round-based synchronous (timely) computing model that is equivalent to the popular crash prone round-based asynchronous (time-free) distributed computing model. Equivalence means here that any problem that can be solved by a protocol in one model can be solved by the same protocol in the other model. The style of this note is voluntarily informal. Its aim is mainly pedagogical. Its ambition is to help better understand relations linking synchronous and asynchronous distributed computing systems, and the nature of failures that make them difficult to master.
Michel Raynal, Matthieu Roy
PRDC1
2005 From Static Distributed Systems to Dynamic Systems
abstract
A noteworthy advance in distributed computing is due to the recent development of peer-to-peer systems. These systems are essentially dynamic in the sense that no process can get a global knowledge on the system structure. They mainly allow processes to look up for data that can be dynamically added/suppressed in a permanently evolving set of nodes. Although protocols have been developed for such dynamic systems, to our knowledge, up to date no computation model for dynamic systems has been proposed. Nevertheless, there is a strong demand for the definition of such models as soon as one wants to develop provably correct protocols suited to dynamic systems. This paper proposes a model for (a class of) dynamic systems. That dynamic model is defined by (1) a parameter (an integer denoted a) and (2) two basic communication abstractions (query-response and persistent reliable broadcast). The new parameter is a threshold value introduced to capture the liveness part of the system (it is the counterpart of the minimal number of processes that do not crash in a static system). To show the relevance of the model, the paper adapts an eventual leader protocol designed for the static model, and proves that the resulting protocol is correct within the proposed dynamic model. In that sense, the paper has also a methodological flavor, as it shows that simple modifications to existing protocols can allow them to work in dynamic systems.
Achour Mostéfaoui, Michel Raynal, Corentin Travers, Stacy Patterson, Divyakant Agrawal, Amr El Abbadi
SRDS2
2005 Wait-free computing: an introductory lecture
Michel Raynal
Future Gener. Comput. Syst.1
2005 Stabilizing mobile philosophers
Ajoy K. Datta, Maria Potop-Butucaru, Michel Raynal
Inf. Process. Lett.3
2005 Asynchronous bounded lifetime failure detectors
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.3
2005 Simple and Efficient Oracle-Based Consensus Protocols for Asynchronous Byzantine Systems
abstract
This paper is on the consensus problem in asynchronous distributed systems where (up to f) processes (among n) can exhibit a Byzantine behavior, i.e., can deviate arbitrarily from their specification. One way to solve the consensus problem in such a context consists of enriching the system with additional oracles that are powerful enough to cope with the uncertainty and unpredictability created by the combined effect of Byzantine behavior and asynchrony. This paper presents two kinds of Byzantine asynchronous consensus protocols using two types of oracles, namely, a common coin that provides processes with random values and a failure detector oracle. Both allow the processes to decide in one communication step in favorable circumstances. The first is a randomized protocol for an oblivious scheduler model that assumes n > 6f. The second one is a failure detector-based protocol that assumes n > tif. These protocols are designed to be particularly simple and efficient in terms of communication steps, the number of messages they generate in each step, and the size of messages. So, although they are not optimal in the number of Byzantine processes that can be tolerated, they are particularly efficient when we consider the number of communication steps they require to decide and the number and size of the messages they use. In that sense, they are practically appealing.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Dependable Secur. Comput.3
2004 A Distributed Implementation of Sequential Consistency with Multi-Object Operations
abstract
Sequential consistency is a consistency criterion for concurrent objects stating that the execution of a multiprocess program is correct if it could have been produced by executing the program on a mono-processor system, preserving the order of the operations of each individual process. Several protocols implementing sequential consistency on top of asynchronous distributed systems have been proposed. They assume that the processes access the shared objects through basic read and write operations. We consider the case where the processes can invoke multiobject operations which can read or write several objects in a single operation atomically. It proposes a particularly simple protocol that guarantees sequentially consistent executions in such a context. The previous sequential consistency protocols, in addition to considering only unary operations, assume either full replication or a central manager storing copies of all the objects. In contrast, the proposed protocol has the noteworthy feature that each object has a separate manager. Interestingly, this provides the protocol with a versatility dimension that allows deriving simple protocols providing sequential consistency or atomic consistency when each operation is on a single object.
Michel Raynal, Krishnamurthy Vidyasankar
ICDCS1
2004 A Methodological Construction of an Efficient Sequential Consistency Protocol
abstract
A concurrent object is an object that can be concurrently accessed by several processes. Sequential consistency is a consistency criterion for such objects. Informally, it states that a multiprocess program executes correctly if its results could have been produced by executing that program on a single processor system. (Sequential consistency is weaker than atomic consistency -the usual consistency criterion- as it does not refer to real-time.) The paper proposes a simple protocol that ensures sequential consistency when the shared memory abstraction is supported by the local memories of nodes that can communicate only by exchanging messages through reliable channels. Differently from other sequential consistency protocols, the proposed protocol does not rely on a strong synchronization mechanism such as an atomic broadcast primitive or a central node managing a copy of every shared object. From a methodological point of view, the protocol is built incrementally starting from the very definition of sequential consistency. It lies the noteworthy property of providing fast writes operations (i.e., a process has never to wait when it writes a new value in a shared object). According to the current local state, some read operations can also be fast. An experimental evaluation of the protocol is also presented. The proposed protocol could be used to manage Web page caching.
Vicent Cholvi, Antonio Fernández 0001, Ernesto Jiménez, Michel Raynal
NCA4
2004 Brief announcement: veto number and the respective power of eventual failure detectors
abstract
No abstract available.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
PODC3
2004 Brief announcement: the synchronous condition-based consensus hierarchy
abstract
No abstract available.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
PODC3
2004 A Hybrid Approach for Building Eventually Accurate Failure Detectors
abstract
Unreliable failure detectors introduced by Chandra and Toueg are abstract mechanisms that provide information about process crashes. On the one hand, failure detectors allow a statement of the minimal requirements on process failures that allow solutions to problems that cannot otherwise be solved in purely asynchronous systems. However, on the other hand, they cannot be implemented in such systems: their implementation requires that the underlying distributed system be enriched with additional assumptions. Classic failure detector implementations rely on additional synchrony assumptions such as partial synchrony. More recently, a new approach for implementing failure detectors has been proposed: it relies on behavioral properties on the flow of messages exchanged. This shows that these approaches are not antagonistic and can be advantageously combined. A hybrid protocol (the first to our knowledge) implementing failure detectors with eventual accuracy properties is presented. Interestingly, this protocol benefits from the best of both worlds in the sense that it converges (i.e., provides the required failure detector) as soon as either the system behaves synchronously or the required message exchange pattern is satisfied. This shows that, to expedite convergence, it can be interesting to consider that the underlying system can satisfy several alternative assumptions.
Achour Mostéfaoui, David Powell, Michel Raynal
PRDC3
2004 Optimal early stopping uniform consensus in synchronous systems with process omission failures
abstract
Consensus is a central problem of fault-tolerant distributed computing that, in the context of synchronous distributed systems, has received a lot of attention in the crash failure model and in the Byzantine failure model. This paper considers synchronous distributed systems made up of n processes, where up to t can commit failures by crashing or omitting to send or receive messages when they should ("process omission" failure model). It presents a protocol solving uniform consensus in such a context. This protocol has several noteworthy features. First, it is particularly simple. Then, it is optimal both in (1) the number of communication steps needed for processes to decide and stop, namely, min(f+2,t+1) where f is the actual number of faulty processes, and (2) the number of processes that can be faulty, namely t
Philippe Raipin Parvédy, Michel Raynal
SPAA2
2004 Simple and Efficient Oracle-Based Consensus Protocols for Asynchronous Byzantine Systems
abstract
This paper is on the consensus problem in asynchronous distributed systems where (up to f) processes (among n) can exhibit a Byzantine behavior, i.e., can deviate arbitrarily from their specification. A way to solve the consensus problem in such a context consists of enriching the system with additional oracles that are powerful enough to cope with the uncertainty and unpredictability created by the combined effect of Byzantine behavior and asynchrony. Considering two types of such oracles, namely, an oracle that provides processes with random values, and a failure detector oracle, the paper presents two families of Byzantine asynchronous consensus protocols. Two of these protocols are particularly noteworthy: they allow the processes to decide in one communication step in favorable circumstances. The first is a randomized protocol that assumes n > 5f. The second one is a failure detector-based protocol that assumes n > 6f. These protocols are designed to be particularly simple and efficient in terms of communication steps, the number of messages they generate in each step, and the size of messages. So, although they are not optimal in the number of Byzantine processes that can be tolerated, they are particularly efficient when we consider the number of communication steps they require to decide, and the number and size of the messages they use. In that sense, they are practically appealing.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
SRDS3
2004 Crash-Resilient Time-Free Eventual Leadership
abstract
Leader-based protocols rest on a primitive able to provide the processes with the same unique leader. Such protocols are very common in distributed computing to solve synchronization or coordination problems. Unfortunately, providing such a primitive is far from being trivial in asynchronous distributed systems prone to process crashes. (It is even impossible in fault-prone purely asynchronous systems.) To circumvent this difficulty, several protocols have been proposed that build a leader facility on top of an asynchronous distributed system enriched with synchrony assumptions. This paper consider another approach to build a leader facility, namely, it considers a behavioral property on the flow of messages that are exchanged. This property has the noteworthy feature not to involve timing assumptions. Two protocols based on this time-free property that implement a leader primitive are described. The first one uses potentially unbounded counters, while the second one (which is a little more involved) requires only finite memory. These protocols rely on simple design principles that make them attractive, easy to understand and provably correct.
Achour Mostéfaoui, Michel Raynal, Corentin Travers
SRDS2
2004 The Notion of Veto Number and the Respective Power of OP and OS to Solve One-Shot Agreement Problems
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
DISC3
2004 The Synchronous Condition-Based Consensus Hierarchy
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC3
2004 Condition-based consensus solvability: a hierarchy of conditions and efficient protocols
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
Distributed Comput.3
2004 A weakest failure detector-based asynchronous consensus protocol for f<n
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.3
2004 A necessary and sufficient condition for transforming limited accuracy failure detectors
Emmanuelle Anceaume, Antonio Fernández 0001, Achour Mostéfaoui, Gil Neiger, Michel Raynal
J. Comput. Syst. Sci.5
2004 The Information Structure of Indulgent Consensus
abstract
To solve consensus, distributed systems have to be equipped with oracles such as a failure detector, a leader capability, or a random number generator. For each oracle, various consensus algorithms have been devised. Some of these algorithms are indulgent toward their oracle in the sense that they never violate consensus safety, no matter how the underlying oracle behaves. We present a simple and generic indulgent consensus algorithm that can be instantiated with any specific oracle and be as efficient as any ad hoc consensus algorithm initially devised with that oracle in mind. The key to combining genericity and efficiency is to factor out the information structure of indulgent consensus executions within a new distributed abstraction, which we call "Lambda". Interestingly, identifying this information structure also promotes a fine-grained study of the inherent complexity of indulgent consensus. We show that instantiations of our generic algorithm with specific oracles, or combinations of them, match lower bounds on oracle-efficiency, zero-degradation, and one-step-decision. We show, however, that no leader or failure detector-based consensus algorithm can be, at the same time, zero-degrading and configuration-efficient. Moreover, we show that leader-based consensus algorithms that are oracle-efficient are inherently zero-degrading, but some failure detector-based consensus algorithms can be both oracle-efficient and configuration-efficient. These results highlight some of the fundamental trade offs underlying each oracle.
Rachid Guerraoui, Michel Raynal
IEEE Trans. Computers2
2003 Token-Based Sequential Consistency in Asynchronous Distributed Systems
abstract
A concurrent object is an object that can be concurrently accessed by several processes. Sequential consistency is a consistency criterion for such objects. It informally states that a multiprocess program executes correctly if its results could have been produced by executing that program on a single processor system. (Sequential consistency is weaker than atomic consistency-the usual consistency criterion-as it does not refer to real-time.) The paper proposes a new, surprisingly simple protocol that ensures sequential consistency when the shared memory abstraction is supported by the local memories of nodes that can communicate only by exchanging messages through reliable channels. The protocol nicely combines, in a simple way, the use a of token with cached values. It has the noteworthy property to never invalidate cached values, thereby providing fast read operations (i.e., a process has never to wait to get a correct value of a shared object). Additionally, The paper presents a simple token navigation protocol.
Michel Raynal
AINA1
2003 Evaluating the Condition-Based Approach to Solve Consensus
abstract
Several approaches have been proposed to circumvent the impossibility to solve consensus in asynchronous distributed systems prone to process crash failures. Among them, randomization, unreliable failure detectors, and leader oracles have been particularly investigated. Recently a new approach (called “condition-based”) has been proposed. Let an input vector be a vector whose i-th entry contains the value proposed by process pi. The conditionbased approach consists in stating conditions on input vectors that make consensus solvable despite up to f process crashes. Several conditions have been proposed. (As an example, one of them requires that the greatest value in an input vector appears more than f times.) This paper presents an evaluation of the condition-based approach to solve consensus. It shows that this approach is particularly attractive and very efficient when the probability of process crashes is low (a common fact in practice). In these cases, the probability for the condition-based protocol to terminate is practically equal to 1.
Achour Mostéfaoui, Eric Mourgaya, Philippe Raipin Parvédy, Michel Raynal
DSN4
2003 Asynchronous Implementation of Failure Detectors
abstract
Unreliable failure detectors introduced by Chandra and Toueg are abstract mechanisms that provide information on process failures. On the one hand, failure detectors allow to state the minimal requirements on process failures that allow to solve problems that cannot be solved in purely asynchronous systems. But, on the other hand, they cannot be implemented in such systems: their implementation requires that the underlying distributed system be enriched with additional assumptions. The usual failure detector implementations rely on additional synchrony assumptions (e.g., partial synchrony). This paper proposes a new look at the implementation of failure detectors and more specifically at Chandra-Toueg’s failure detectors. The proposed approach does not rely on synchrony assumptions (e.g., it allows the communication delays to always increase). It is based on a query-response mechanism and assumes that the query/response messages exchanged obey a pattern where the responses from some processes to a query arrive among the (n − f ) first ones (n being the total number of processes, f the maximum number of them that can crash, with 1 ≤ f< n). When we consider the particular case f =1 , and the implementation of a failure detector of the class denoted S (the weakest class that allows to solve the consensus problem), the additional assumption the underlying system has to satisfy boils down to a simple channel property, namely, there is eventually a pair of processes (pi ,p j) such that the channel connecting them is never the slowest among the channels connecting pi or pj to the other processes. A probabilistic analysis shows that this requirement is practically met in asynchronous distributed systems.
Achour Mostéfaoui, Eric Mourgaya, Michel Raynal
DSN3
2003 Elastic Vector Time
abstract
In recent years there has been an increasing demand to build "soft" real-time applications on top of asynchronous distributed systems. Designing and implementing such applications is a non-trivial task and application designers are often faced with the need to circumvent impossibility results. In this paper we discuss how to ensure that actions are executed in the correct order even in the face of failures. We propose a novel time base and a new synchronization mechanism for the design of distributed "soft" real-time applications. We demonstrate (1) how this time base can be used to enforce an externally consistent ordering, and (2) how it permits to circumvent impossibility results by sketching how to solve the leader election and perfect failure detection problem.
Christof Fetzer, Michel Raynal
ICDCS2
2003 A Generic Framework for Indulgent Consensus
abstract
Consensus is a fundamental distributed agreement problem that has to be solved when one has to design or implement reliable applications. As consensus cannot be solved in pure asynchronous distributed systems, those systems have to be equipped with appropriate oracles to circumvent the impossibility. Several oracles (unreliable failure detector leader capability, random number generator) have been proposed, and consensus protocols based on such ad hoc oracles have been designed This paper presents a generic consensus framework that can be instantiated with any oracle, or combination of oracles, that satisfies a set of properties. This generic framework provides indulgent consensus protocols that are particularly simple and efficient both in well-behaved runs (i.e., when there are no failures), and in stable runs (i.e., when there is no failure during the execution although some processes can be initially crashed). In those runs, the protocols terminate in two communication steps (which is optimal). Indulgence means that the resulting protocol never violates its safety property even when the underlying oracle behaves arbitrarily. Interestingly, the protocol can also allow processes to decide in one communication step in some specific configurations.
Rachid Guerraoui, Michel Raynal
ICDCS2
2003 Nested Invocation Protocol for Object-Based Systems
abstract
We discuss how to invoke a method on multiple object replicas in a quorum-based way. Suppose each instance of a method t on replicas of an object x invokes another method u on replicas in a quorum of an object y. Here, the method u is redundantly invoked multiple times on some replicas of the object y. If each instance of the method t issues a method u to its own quorum, more number of replicas are manipulated than the quorum number This is quorum expansion. We discuss a protocol to invoke methods on replicas in a nested manner without the redundant invocation and quorum expansion. We evaluate the protocol on how many replicas are manipulated and requests are issued.
Kenichi Hori, Tomoya Enokido, Makoto Takizawa 0001, Michel Raynal
ISORC4
2003 Brief announcement: early decision despite general process omission failures
abstract
No abstract available.
Fabrice Le Fessant, Philippe Raipin Parvédy, Michel Raynal
PODC3
2003 Using Conditions to Expedite Consensus in Synchronous Distributed Systems
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC3
2003 Conditions on input vectors for consensus solvability in asynchronous distributed systems
abstract
This article introduces and explores the condition-based approach to solve the consensus problem in asynchronous systems. The approach studies conditions that identify sets of input vectors for which it is possible to solve consensus despite the occurrence of up to f process crashes. The first main result defines acceptable conditions and shows that these are exactly the conditions for which a consensus protocol exists. Two examples of realistic acceptable conditions are presented, and proved to be maximal, in the sense that they cannot be extended and remain acceptable. The second main result is a generic consensus shared-memory protocol for any acceptable condition. The protocol always guarantees agreement and validity, and terminates (at least) when the inputs satisfy the condition with which the protocol has been instantiated, or when there are no crashes. An efficient version of the protocol is then designed for the message passing model that works when f < n /2, and it is shown that no such protocol exists when f ≥ n /2. It is also shown how the protocol's safety can be traded for its liveness.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
J. ACM3
2003 Efficient Causality-Tracking Timestamping
abstract
Vector clocks are the appropriate mechanism used to track causality among the events produced by a distributed computation. Traditional implementations of vector clocks require application messages to piggyback a vector of n integers (where n is the number of processes). This paper investigates the tracking of the causality relation on a subset of events (namely, the events that are defined as "relevant" from the application point of view) in a context where communication channels are not required to be FIFO, and where there is no a priori information on the connectivity of the communication graph or the communication pattern. More specifically, the paper proposes a suite of simple and efficient implementations of vector clocks that address the reduction of the size of message timestamps, i.e., they do their best to have message timestamps whose size is less than n. The relevance of such a suite of protocols is twofold. From a practical side, it constitutes the core of an adaptive timestamping software layer that can used by underlying applications. From a theoretical side, it provides a comprehensive view that helps better understand distributed causality-tracking mechanisms.
Jean-Michel Hélary, Michel Raynal, Giovanna Melideo, Roberto Baldoni
IEEE Trans. Knowl. Data Eng.2
2003 Atomic Broadcast in Asynchronous Crash-Recovery Distributed Systems and Its Use in Quorum-Based Replication
abstract
Atomic broadcast is a fundamental problem of distributed systems: It states that messages must be delivered in the same order to their destination processes. This paper describes a solution to this problem in asynchronous distributed systems in which processes can crash and recover. A consensus-based solution to atomic broadcast problem has been designed by Chandra and Toueg for asynchronous distributed systems where crashed processes do not recover. We extend this approach: it transforms any consensus protocol suited to the crash-recovery model into an atomic broadcast protocol suited to the same model. We show that atomic broadcast can be implemented requiring few additional log operations in excess of those required by the consensus. The paper also discusses how additional log operations can improve the protocol in terms of faster recovery and better throughput. To illustrate the use of the protocol, the paper also describes a solution to the replica management problem in asynchronous distributed systems in which processes can crash and recover. The proposed technique makes a bridge between established results on weighted voting and recent results on the consensus problem.
Luís E. T. Rodrigues, Michel Raynal
IEEE Trans. Knowl. Data Eng.2
2003 Early Stopping in Global Data Computation
abstract
No abstract available.
Carole Delporte-Gallet, Hugues Fauconnier, Jean-Michel Hélary, Michel Raynal
IEEE Trans. Parallel Distributed Syst.4
2002 A Versatile and Modular Consensus Protoco
abstract
Investigates a modular and versatile approach to solve the consensus problem in asynchronous distributed systems in which up to f processes may crash (f<n/2), but equipped with appropriate oracles. It presents a generic protocol that proceeds by consecutive asynchronous rounds. Each round follows a "two-phase" pattern. The modularity and the versatility of the protocol appear at each phase of a round. The first phase is a selection phase that allows to use any combination merging random oracle, leader oracle and condition. Its aim is to ensure termination by allowing the processes to start the second phase with the same value. The aim of the second phase is to ensure that the agreement property cannot be violated. Its cost depends on the value of f: two communication steps when f<n/2, that reduce to a single communication step when f
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DSN3
2002 Early stopping in aglobal data computation
Carole Delporte-Gallet, Hugues Fauconnier, Jean-Michel Hélary, Michel Raynal
PODC4
2002 Building responseive TMR-based servers in presence of timing constraints
abstract
No abstract available.
Paul D. Ezhilchelvan, Jean-Michel Hélary, Michel Raynal
PODC3
2002 Asynchronous interactive consistency and its relation with error-correcting codes
abstract
No abstract available.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
PODC3
2002 Towards a formal model for view maintenance in data warehouses
abstract
No abstract available.
Achour Mostéfaoui, Michel Raynal, Matthieu Roy, Divyakant Agrawal, Amr El Abbadi
PODC2
2002 An Introduction to the Renaming Problem
abstract
The aim of this paper is to provide a brief introduction to the renaming problem for unfamiliar readers. In the renaming problem the processes have to acquire new names from a small bounded space despite possible process crashes and asynchrony. The problem is first introduced. Then two solutions are presented. One considers the shared memory model, while the second considers the message-passing model.
Michel Raynal
PRDC1
2002 Consensus in Synchronous Systems: A Concise Guided Tour
abstract
This paper considers consensus protocols for synchronous systems where processes can commit crash failures, omission failures or Byzantine failures. It presents and revisits consensus protocols coping with such failures in an increasing order of difficulty. The paper can be seen as a short tutorial whose aim is to make the reader familiar with synchrony assumptions, different definitions of the consensus problem, and a hierarchy of process failure models. An important concern of the paper lies in simplicity. In addition to the survey flavor of the paper, several results that are presented are new, among which the ones concerning the omission failure model.
Michel Raynal
PRDC1
2002 Tracking immediate predecessors in distributed computations
abstract
A distributed computation is usually modeled as a partially ordered set of relevant events (the relevant events are a subset of the primitive events produced by the computation). An important causality-related distributed computing problem, that we call the Immediate Predecessors Tracking (IPT) problem, consists in associating with each relevant event, on the fly and without using additional control messages, the set of relevant events that are its immediate predecessors in the partial order. So, IPT is the on-the-fly computation of the transitive reduction (i.e., Hasse diagram) of the causality relation defined by a distributed computation. This paper addresses the IPT problem: it presents a family of protocols that provides each relevant event with a timestamp that exactly identifies its immediate predecessors. The family is defined by a general condition that allows application messages to piggyback control information whose size can be smaller than $n$ (the number of processes). In that sense, this family defines message size-efficient IPT protocols. According to the way the general condition is implemented, different IPT protocols can be obtained. Two of them are exhibited.
Emmanuelle Anceaume, Jean-Michel Hélary, Michel Raynal
SPAA3
2002 Sequential consistency as lazy linearizability
abstract
This revue shows that, from an implementation point of view, sequential consistency can be considered as a form of lazy linearizability. This claim is supported by a versatile protocol that can be tailored to implement any of them.
Michel Raynal
SPAA1
2002 The Lord of the Rings: Efficient Maintenance of Views at Data Warehouses
Divyakant Agrawal, Amr El Abbadi, Achour Mostéfaoui, Michel Raynal, Matthieu Roy
DISC4
2002 Distributed Agreement and Its Relation with Error-Correcting Codes
Roy Friedman 0001, Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC4
2002 Condition-Based Protocols for Set Agreement Problems
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
DISC3
2002 An introduction to oracles for asynchronous distributed systems
Achour Mostéfaoui, Eric Mourgaya, Michel Raynal
Future Gener. Comput. Syst.3
2002 Interval Consistency of Asynchronous Distributed Computations
Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
J. Comput. Syst. Sci.3
2002 A Versatile Family of Consensus Protocols Based on Chandra-Toueg's Unreliable Failure Detectors
abstract
This paper is on consensus protocols for asynchronous distributed systems prone to process crashes, but equipped with Chandra-Toueg's (1996) unreliable failure detectors. It presents a unifying approach based on two orthogonal versatility dimensions. The first concerns the class of the underlying failure detector. An instantiation can consider any failure detector of the class S (provided that at least one process does not crash), or oS (provided that a majority of processes do not crash). The second versatility dimension concerns the message exchange pattern used during each round of the protocol. This pattern (and, consequently, the round message cost) can be defined for each round separately, varying from O(n) (centralized pattern) to O(n/sup 2/) (fully distributed pattern), n being the number of processes. The resulting versatile protocol has nice features and actually gives rise to a large and well-identified family of failure detector-based consensus protocols. Interestingly, this family includes at once new protocols and some well-known protocols (e.g., Chandra-Toueg's oS-based protocol). The approach is also interesting from a methodological point of view. It provides a precise characterization of the two sets of processes that, during a round, have to receive messages for a decision to be taken (liveness) and for a single value to be decided (safety), respectively. Interestingly, the versatility of the protocol is not restricted to failure detectors: a simple timer-based instance provides a consensus protocol suited to partially synchronous systems.
Michel Hurfin, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Computers3
2001 Building TMR-Based Reliable Servers Despite Bounded Input Lifetimes
Paul D. Ezhilchelvan, Jean-Michel Hélary, Michel Raynal
Euro-Par3
2001 Shared State Consistency for Time-Sensitive Distributed Applications
abstract
Distributed applications that share a dynamically changing state are increasingly being deployed in wide-area environments. Such applications must access the state in a consistent manner, but the consistency requirements vary significantly from other systems. For example, shared memory models, such as sequential consistency, focus on the ordering of operations, and the same level of consistency is provided to each process. In interactive distributed applications, the timeliness of updates becoming effective could be an extremely important consistency requirement, and it could be different across different users. We propose a system that provides both non-timed and time-sensitive read and write operations for dynamic shared state. For example, a timed read can be used by a process to read a recently written value, whereas a timed write can make a new value available to all readers within a certain amount of time. We develop a consistency model that precisely defines the semantics of timed and non-tinted read and write operations. A protocol that implements this model is also presented. We also describe an implementation and some performance measurements.
Vijaykumar Krishnaswamy, Mustaque Ahamad, Michel Raynal, David E. Bakken
ICDCS3
2001 A Condition for k-Set Agreement in Asynchronous Distributed Systems
abstract
The k-set agreement problem has no solution in fully asynchronous distributed systems made up of n processes (where at most f of them can crash), when f/spl ges/k. This paper presents a condition on the occurrence pattern of proposed values that allows to solve the problem whatever the value of f. More specifically, it is shown that if there is a set of more than (kn+f)/(k+1) processes that propose at most k different values, then the k-set agreement problem can be solved. When we consider the particular case of the consensus problem (k=1), this means that, whatever the value of f, this problem can be solved when more than (n+f)/2 processes propose the same value. As an example of the usefulness of this condition, it is used to improve Ben-Or's randomized consensus protocol.
Achour Mostéfaoui, Michel Raynal
IPDPS2
2001 Primary Component Asynchronous Group Membership as an Instance of a Generic Agreement Framework
abstract
Group-based computing is becoming more and more popular when one has to design middleware able to support reliable distributed applications. This paradigm is made of two basic services, namely, a group membership service and a group communication service. More generally, a group is a set of processes cooperating to carry out a common task (e.g., copies of a replicated server, participants in a transaction or users in a CSCW-based application). Due to the desire of new processes to join the group, to the desire of a group member to leave it, or to process crashes, the composition of a group can evolve dynamically. The set of processes that currently implements the group is called the current view of the group. This paper addresses the specification and the implementation of a primary component group membership service. Primary component means that the specification imposes to have a single view at any time. The paper first proposes a specification for the problem. Then it presents a protocol that implements that specification in asynchronous distributed systems equipped with failure detectors. This primary component group membership protocol is obtained as an appropriate instantiation of a general agreement framework.
Fabíola Greve, Michel Hurfin, Michel Raynal, Frédéric Tronel
ISADS3
2001 Randomized Multivalued Consensus
abstract
The consensus problem is a fundamental problem one has to solve to implement reliable services or applications on top of asynchronous distributed systems prone to failures. Unfortunately, this problem cannot be solved in those systems as soon as one process crashes (Fischer-Lynch-Paterson's impossibility result). Two approaches have been investigated to circumvent this impossibility result. Both consist in enriching the underlying system with appropriate "oracles". The unreliable failure detector concept proposed by Chandra and Toueg (1996) constitutes one family of such oracles. Since it has been proposed the failure detector based approach has given rise to several failure detector-based consensus protocols. The other family of oracles consists in allowing each process to use a random number generator. In that case, protocol termination is only probabilistic. A few randomized consensus protocols for message-passing asynchronous distributed systems have been proposed. Moreover, they consider that processes can only propose values from a binary set. This paper proposes a new randomized consensus protocol that allows processes to propose arbitrary values. Contrary to other randomized consensus protocols, the proposed protocol does not require the a priori knowledge of the set of values that can be proposed by processes. It relies on a relatively simple combination of randomization and reliable broadcast.
Paul D. Ezhilchelvan, Achour Mostéfaoui, Michel Raynal
ISORC3
2001 A hierarchy of conditions for consensus solvability
abstract
In a previous paper we introduced the condition-based approach, consisting of identifying sets of input vectors, called conditions, for which there exists an asynchronous protocol solving consensus despite the occurrence of up to f process crashes, and characterized this set of conditions, @@@@wkf. Here, we investigate @@@@wkf from the complexity perspective, and show that this class consists of a hierarchy of classes of conditions, @@@@[d]f, where d, 0 ⪇ d ⪇ f, is the degree of the condition, each one strictly contained in the previous one. The value f - d represents the “difficulty” of the class @@@@[d]f: we present a generic condition-based protocol that can be instantiated with any C ∈ @@@@[d]f, and solve consensus with (2n + 1) [log2([(f - d)/2] + 1)] shared memory read/write operations per process. For each d we present two natural conditions, C1[d]f and C2[d]f, that might be useful in practice, and we use them to show that the class containments stated above are strict. Various properties of the hierarchy are also derived. Mainly, it is shown that a class can be characterized in two equivalent but complementary ways: one is convenient for designing protocols while the other is for analyzing the class properties.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
PODC3
2001 An Adaptive Failure Detection Protocol
abstract
The detection of process failures is a crucial problem system designers have to cope with in order to build fault-tolerant distributed platforms. Unfortunately, it is impossible to distinguish with certainty a crashed process from a very slow process in a purely asynchronous distributed system. This prevents some problems from being solved in such systems. That is why failure detector oracles have been introduced to circumvent these impossibility results. The paper presents a relatively simple protocol that allows a process to "monitor" another process, and consequently to detect its crash. This protocol relies as much as possible on application messages to do this monitoring. Different from previous process crash detection protocols, it uses control messages only when no application message is sent by the monitoring process to the observed process. When the underlying system satisfies the partial synchrony assumption, it actually implements an eventually perfect failure detector (i.e., a failure detector of the class usually denoted OP). Moreover if the average observed transmission delay is finite and the upper layer application terminates within a bounded number of steps for any failure detector in OP after the failure detector becomes "perfect", then, when run with the proposed protocol, it also terminates correctly. These properties make the protocol inexpensive, implementable, and powerful. The paper also describes performance measurements of an implementation of the protocol.
Christof Fetzer, Michel Raynal, Frédéric Tronel
PRDC2
2001 Efficient Condition-Based Consensus
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
SIROCCO3
2001 Randomized k-set agreement
abstract
The k-Set Agreement problem generalizes the consensus problem (which corresponds to the case k = 1). The processes propose values and each correct process has to decide a value such that (1) a decided value is a proposed value, and (2) no more than k distinct values are decided. Let f be the maximum number of processes that can crash. It has first been shown that the consensus problem cannot be solved in asynchronous distributed systems when f > 0 (this is the well-known FLP's impossibility result). It has then been shown that this impossibility still holds for the k-set agreement problem when f ⪈ k.
Achour Mostéfaoui, Michel Raynal
SPAA2
2001 A Consensus Protocol Based on a Weak FailureDetector and a Sliding Round Window
abstract
The paper revisits the "sliding window" notion commonly encountered in communication protocols and applies it to the round numbers of round-based asynchronous protocols. This approach is novel. To illustrate its benefits, the paper presents an original weak failure detector-based consensus protocol that allows each process to be simultaneously involved in several rounds. The rounds in which a process is simultaneously involved defines "sliding round window". The proposed approach has several advantages. It fits better to the uncertainty created by the asynchrony and failures, and consequently permits one to design efficient round-based asynchronous protocols. Maybe more important, it also provides a better understanding of the global synchronization that manages the protocol progress from round to round. This appears clearly in the proposed failure detector-based consensus protocol, where the "sliding round window" allows one to dynamically define the message exchange pattern for each round separately.
Michel Hurfin, Achour Mostéfaoui, Michel Raynal, Raimundo José de Araújo Macêdo
SRDS3
2001 Conditions on input vectors for consensus solvability in asynchronous distributed systems
abstract
This paper introduces and explores a new condition based approach to solve the consensus problem in asynchronous systems. The approach consists of identifying sets of input vectors, called conditions, for which it is possible to design a protocol solving consensus despite the occurrence of up to f process crashes.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
STOC3
2001 Consistent Checkpointing for Transaction Systems
abstract
Whether it is for audit or for recovery purposes, data checkpointing is an important problem of transaction systems. Actually, transactions establish dependence relations on data checkpoints taken by data object managers. So, given an arbitrary set of data checkpoints (including at least a single data checkpoint from a data manager, and at most a data checkpoint from each data manager), an important question is the following one: ‘Can these data checkpoints be members of a same consistent global checkpoint?’ This paper answers this question by providing a necessary and sufficient condition suited to transaction systems. Moreover, to show its usefulness, two non-intrusive data checkpointing protocols are designed from this condition.
Roberto Baldoni, Francesco Quaglia, Michel Raynal
Comput. J.3
2001 The logically instantaneous communication mode: a communication abstraction
Achour Mostéfaoui, Michel Raynal, Paulo Veríssimo
Future Gener. Comput. Syst.2
2001 Rollback-Dependency Trackability: A Minimal Characterization and Its Protocol
Roberto Baldoni, Jean-Michel Hélary, Michel Raynal
Inf. Comput.3
2001 Impossibility of scalar clock-based communication-induced checkpointing protocols ensuring the RDT property
Roberto Baldoni, Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.4
2001 Consensus-Based Fault-Tolerant Total Order Multicast
abstract
While total order broadcast (or atomic broadcast) primitives have received a lot of attention, this paper concentrates on total order multicast to multiple groups in the context of asynchronous distributed systems in which processes may suffer crash failures. "Multicast to Multiple Groups" means that each message is sent to a subset of the process groups composing the system, distinct messages possibly having distinct destination groups. "Total Order" means that all message deliveries must be totally ordered. This paper investigates a consensus-based approach to solve this problem and proposes a corresponding protocol to implement this multicast primitive. This protocol is based on two underlying building blocks, namely, uniform reliable multicast and uniform consensus. Its design characteristics lie in the two following properties. The first one is a minimality property, more precisely, only the sender of a message and processes of its destination groups have to participate in the total order multicast of the message. The second property is a locality property: No execution of a consensus has to involve processes belonging to distinct groups (i.e., consensus is executed on a "per group" basis). This locality property is particularly useful when one is interested in using the total order multicast primitive in large-scale distributed systems. In addition to a correctness proof, an improvement that reduces the cost of the protocol is also suggested.
Udo Fritzke Jr., Philippe Ingels, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Parallel Distributed Syst.4
2000 From Crash Fault-Tolerance to Arbitrary-Fault Tolerance: Towards a Modular Approach
abstract
Presents a generic methodology to transform a protocol which is resilient to process crashes into one that is resilient to arbitrary failures in the case where processes run the same text and regularly exchange messages (i.e. the case of round-based protocols). The methodology follows a modular approach, encapsulating the detection of arbitrary failures in specific modules. This can be the starting point for designing tools that allow automatic transformation. We show an application of this methodology to the case of consensus.
Roberto Baldoni, Jean-Michel Hélary, Michel Raynal
DSN3
2000 The Best of Both Worlds: A Hybrid Approach to Solve Consensus
abstract
It is now well recognized that the consensus problem is a fundamental problem when one has to implement fault-tolerant distributed services in asynchronous distributed systems prone to process crash failures. This paper considers the binary consensus problem in such a system. Following an approach investigated by Aguilera and Toueg, it proposes a simple binary consensus protocol that combines failure detection and randomization. This protocol terminates deterministically when the failure detection mechanism works correctly; it terminates with probability 1, otherwise. A performance evaluation of the protocol is also provided. Last but not least, it is important to note that the proposed protocol is both efficient and simple. Additionally it can be simplified to give rise either to a deterministic failure detector-based consensus protocol or to a randomized consensus protocol.
Achour Mostéfaoui, Michel Raynal, Frédéric Tronel
DSN2
2000 Logical Instantaneity and Causal Order: Two "First Class" Communication Modes for Parallel Computing
Michel Raynal
Euro-Par1
2000 Quorum-Based Replication in Asynchronous Crash-Recovery Distributed Systems (Research Note)
Luís E. T. Rodrigues, Michel Raynal
Euro-Par2
2000 Computing Global Functions in Asynchronous Distributed Systems Prone to Process Crashes
abstract
Global data is a vector with one entry per process. Each entry must be filled with an appropriate value provided by the corresponding process. Several distributed computing problems amount to compute a function on global data. This paper proposes a protocol to solve such problems in the context of asynchronous distributed systems where processes may fail by crashing. The main problem that has to be solved lies in computing the global data and in providing each non-crashed process with a copy of it, despite the possible crash of some processes. To be consistent, the global data must contain (at least) all the values provided by the processes that do not crash. This defines the global data computation (GDC) problem. To solve this problem, processes execute a sequence of asynchronous rounds during which they construct (in a decentralized way) the value of the global data, and eventually each process gets a copy of it. To cope with process crashes, the protocol uses a perfect failure detector. The proposed protocol has been designed to be time-efficient. It allows early decisions. Let t be the maximum number of processes that may crash (t
Jean-Michel Hélary, Michel Hurfin, Achour Mostéfaoui, Michel Raynal, Frédéric Tronel
ICDCS4
2000 Atomic Broadcast in Asynchronous Crash-Recovery Distributed Systems
abstract
Atomic broadcast is a fundamental problem of distributed systems: it states that messages must be delivered in the same order to their destination processes. This paper describes a solution to this problem in asynchronous distributed systems in which processes can crash and recover. A consensus-based solution to atomic broadcast problem has been designed by Chandra and Toueg (1996) for asynchronous distributed systems where crashed processes do nor recover. Although our solution is based on different algorithmic principles, it follows the same approach: it transforms any consensus protocol suited to the crash-recovery model into an atomic broadcast protocol suited to the same model. We show that atomic broadcast can be implemented without requiring any additional log operations in excess of those required by the consensus. The paper also discusses how additional log operations can improve the protocol in terms of faster recovery and better throughput.
Luís E. T. Rodrigues, Michel Raynal
ICDCS2
2000 Consensus Based on Failure Detectors with a Perpetual Accuracy Property
abstract
This paper is on the Consensus problem, in the context of asynchronous distributed systems made of n processes, at most f of them may crash. A family of failure detector classes satisfying a Perpetual Accuracy property is first defined. This family includes the failure detector class S (the class of Strong failure detectors defined by Chandra and Toueg) central to the definition of a class (S/sub x/) where x is the minimum number (x/spl ges/1) of correct processes that can never be suspected to have crashed Then, a protocol that solves the Consensus problem is given. This protocol works with any failure detector class (S/sub x/) of this family. It is particularly simple and uses a Reliable Broadcast protocol as a skeleton. It requires n-x+1 communication steps, and its communication bit complexity is (n-x+1)(n-1)|/spl nu/| (where |/spl nu/| is the maximal size of an initial value a process can propose).
Achour Mostéfaoui, Michel Raynal
IPDPS2
2000 Deadline-Constrained Causal Order
abstract
A causal ordering protocol ensures that if two messages are causally related and have the same destination, they are delivered to the application in their sending order. Causal order strongly simplifies the development of distributed object oriented systems. To prevent causal order violation, either messages may be forced to wait for messages in their past, or late messages may have to be discarded. For a real time setting, the first approach is not suitable since when a message misses a deadline, all the messages that causally depend on it may also be forced to miss their deadlines. We propose a novel causal ordering abstraction that takes message deadlines into consideration. Two implementations are proposed in the context of multicast and broadcast communication that deliver as many messages as possible to the application. Examples of distributed soft real time applications that benefit from the use of a deadline-constrained causal ordering primitive are given.
Luís E. T. Rodrigues, Roberto Baldoni, Emmanuelle Anceaume, Michel Raynal
ISORC4
2000 Time and message-efficient S-based consensus (brief announcement)
abstract
The class of strong failure detectors (denoted S) includes all failure detectors that suspect all crashed processes and that do not suspect some (a priori unknown) process that never crashes. So, a failure detector that belongs to S is intrinsically unreliable as it can arbitrarily suspect correct processes.
Fabíola Greve, Michel Hurfin, Raimundo José de Araújo Macêdo, Michel Raynal
PODC4
2000 k-set agreement with limited accuracy failure detectors
abstract
Let the scope of the accuracy property of an unreliable failure detector be the number x of processes that may not suspect a correct process. The scope notion gives rise to new classes of failure detectors among which we consider Sx and ⋄Sx in this paper (Usual failure detectors consider an implicit scope equal to n, the total number of processes).
Achour Mostéfaoui, Michel Raynal
PODC2
2000 Low cost consensus-based Atomic Broadcast
abstract
Atomic Broadcast (all processes deliver the same set of messages in the same order) is a very powerful communication primitive when one is interested in building fault-tolerant distributed systems. Moreover, it has been shown that Atomic Broadcast and Consensus are equivalent problems in asynchronous distributed systems prone to process crash failures. Hence, several Consensus-based Atomic Broadcast protocols have been designed. This paper introduces a new and particularly efficient Consensus-based Atomic Broadcast protocol. The efficiency is obtained by limiting the use of the Consensus subroutine to the cases where asynchrony and crashes prevent processes from obtaining a simple agreement on the message delivery order. The protocol assumes n>2f (where n is the number of processes and f the maximum number of them that can crash). In the most favorable cases, it requires two communication steps for processes to determine a message batch. In the worst case it requires an additional Consensus execution. It is shown that, when n>3f, the protocol can be simplified. It then requires a single communication step in the most favorable cases. This exhibits an interesting tradeoff relating the cost of the protocol with the maximum number of process failures.
Achour Mostéfaoui, Michel Raynal
PRDC2
2000 Consensus in byzantine asynchronous systems
Roberto Baldoni, Jean-Michel Hélary, Michel Raynal, Lénaick Tanguy
SIROCCO3
2000 Tracking causality in distributed systems: a suite of efficient protocols
Jean-Michel Hélary, Giovanna Melideo, Michel Raynal
SIROCCO3
2000 Communication-Based Prevention of Useless Checkpoints in Fistributed Computations
Jean-Michel Hélary, Achour Mostéfaoui, Robert H. B. Netzer, Michel Raynal
Distributed Comput.4
2000 From Binary Consensus to Multivalued Consensus in asynchronous message-passing systems
Achour Mostéfaoui, Michel Raynal, Frédéric Tronel
Inf. Process. Lett.2
2000 Computing Global Functions in Asynchronous Distributed Systems with Perfect Failure Detectors
abstract
A Global Data is a vector with one entry per process. Each entry must be filled with an appropriate value provided by the corresponding process. Several distributed computing problems amount to compute a function on a global data. This paper proposes a protocol to solve such problems in the context of asynchronous distributed systems where processes may fail by crashing. The main problem that has to be solved lies in computing the global data and in providing each noncrashed process with a copy of it, despite the possible crash of some processes. To be consistent, the global data must contain, at least, all the values provided by the processes that do not crash. This defines the Global Data Computation (GDC) problem. To solve this problem, processes execute a sequence of asynchronous rounds during which they construct, in a decentralized way, the value of the global data and eventually each process gets a copy of it. To cope with process crashes, the protocol uses a perfect failure detector. The proposed protocol has been designed to be time efficient: it allows early decision. Let t be the maximum number of processes that may crash, t
Jean-Michel Hélary, Michel Hurfin, Achour Mostéfaoui, Michel Raynal, Frédéric Tronel
IEEE Trans. Parallel Distributed Syst.4
1999 Distributed Database Checkpointing
Roberto Baldoni, Francesco Quaglia, Michel Raynal
Euro-Par3
1999 Illustrating the Use of Vector Clocks in Property Detection: An Example and a Counter-Example
Michel Raynal
Euro-Par1
1999 Unreliable Failure Detectors with Limited Scope Accuracy and an Application to Consensus
Achour Mostéfaoui, Michel Raynal
FSTTCS2
1999 On Classes of Problems in Asynchronous Distributed Systems with Process Crashes
abstract
This paper is on classes of problems encountered in asynchronous distributed systems in which processes can crash but links are reliable. The hardness of a problem is defined with respect to the difficulty to solve it despite failures: a problem is easy if it can be solved in presence of failures, otherwise it is hard. Three classes of problems are defined: F, NF and NFC. F is the class of easy problems, namely, those that can be solved in presence of failures (e.g., reliable broadcast). The class NF includes harder problems, namely, the ones that can be solved in a non-faulty system (e.g., consensus). The class NFC (NF-complete) is a subset of NF that includes the problems that are the most difficult to solve in presence of failures. It is shown that the terminating reliable broadcast problem, the non-blocking atomic commitment problem and the construction of a perfect failure detector (problem P) are equivalent problems and belong to NFC. Moreover the consensus problem is not in NFC. The paper presents a general reduction protocol that reduces any problem of NF to P. This shows that P is a problem that lies at the core of distributed fault-tolerance.
Eddy Fromentin, Michel Raynal, Frédéric Tronel
ICDCS2
1999 Direct Dependency-Based Determination of Consistent GlobalCheckpoints
Roberto Baldoni, Michel Raynal, Giacomo Cioffi, Jean-Michel Hélary
OPODIS2
1999 Simple Vector Clocks are limited to Solve some Causallity Related Problems
Michel Raynal
OPODIS1
1999 Rollback-Dependency Trackability: Visible Characterizations
abstract
When we consider an asynchronous distributed computation on which local checkpoints have been defined (namely, a communication and checkpoint pattern -in brief, CCP), two types of dependencies between its local checkpoints can be observed.The first type is due to causal sequences of messages that establish on-line trackable dependencies.The second type is due to noncausal sequences of messages (called Z-paths) that establish "hidden" dependencies between local checkpoints (a dependency is "hidden" if it can not be tracked online).The Rollback Dependency Trackability (RDT) property, defined by Y.-M.Wang, has been introduced to study CCPs.A CCP satisfies the RDT property if every pair of local checkpoints that are connected by a "hidden" dependency are also connected by a causal sequence of messages.The RDT property has a great interest: CCPs that satisfy this property allow relatively simple solutions to a lot of practical problems.This paper first introduces the notion of RDT-compliant property.In a given CCP, an X-path is a Z-path that satisfies a property X.The property X is RDT-compliant if every CCP without X-paths satisfies the RDT property.Then, the paper presents a particular RDT-compliant property.This property enjoys several very interesting features.(1) It is "visible" (i.e., it can be tested on-line).(2) It is stronger than previously known RDT-compliant properties.Consequently, this property provides a characterization of RDT better than the previous ones.The question of the minimal characterization of the RDT property is finally investigated.1 Introduction Long running scientific applications and service providing facilities use rollback-recovery techniques to increase their fault tolerance and their availability.This is done by saving onto stable storage the state of processes (i.e., Permission to make digital or hard copies of all or part ofthis work for personal or c~assroon~ use is granted without fee provided that copies arc not made or distrihutcd for protit or commercial advantage and that copies bear this notice and the full citation 011 the first page.'I'0 copy other&c.to republish, to post on servers or to redistribute to lists.
Roberto Baldoni, Jean-Michel Hélary, Michel Raynal
PODC3
1999 Timed Consistency for Shared Distributed Objects
abstract
Ordering and time are two different aspects of consistency of shared objects in a distributed system. One avoids conflicts between operations, the other addresses how quickly the effects of an operation are perceived by the rest of the system. Consistency models such as sequential consistency and causal consistency do not consider the particular time at which an operation is executed to establish a valid order among all the operations of a computation. Timed consistency models require that if a write operation is executed at time t, it must be visible to all nodes by time t+D. Timed consistency generalizes several existing consistency criteria and it is well suited for interactive and collaborative applications, where the action of one user must be seen by others in a timely fashion.
Francisco J. Torres-Rojas, Mustaque Ahamad, Michel Raynal
PODC3
1999 A General Framework to Solve Agreement Problems
abstract
Agreement problems are among the most important problems designers of distributed systems have to cope with. A way to solve them is to first provide a solution to the Consensus problem and then to reduce each agreement problem to Consensus. This "run-time customizing" approach is particularly relevant when upper layer applications have to solve several distinct agreement problems. We investigate a "compile-time customizing" approach to automatically generate ad hoc agreement protocols. A general agreement framework, characterized by six "versatility" parameters, is defined. Appropriate instantiations of these parameters provide particular agreement protocols. This approach is particularly suited to generate efficient agreement protocols.
Michel Hurfin, Raimundo José de Araújo Macêdo, Michel Raynal, Frédéric Tronel
SRDS3
1999 Solving Consensus Using Chandra-Toueg's Unreliable Failure Detectors: A General Quorum-Based Approach
Achour Mostéfaoui, Michel Raynal
DISC2
1999 A Simple and Fast Asynchronous Consensus Protocol Based on a Weak Failure Detector
Michel Hurfin, Michel Raynal
Distributed Comput.2
1999 Restricted failure detectors: Definition and reduction protocols
Michel Raynal, Frédéric Tronel
Inf. Process. Lett.1
1999 Communication-Induced Determination of Consistent Snapshots
abstract
A classical way to determine consistent snapshots consists in using Chandy-Lamport's algorithm. This algorithm relies on specific control messages that allow processes to synchronize local checkpoint determination and message recording in order for the resulting snapshot to be consistent. This paper investigates a communication-induced approach to determine consistent snapshots. In such an approach, control information is carried out by application messages. Two abstract necessary and sufficient conditions are stated: one associated with global checkpoint consistency, the other associated with message recording. A general protocol is derived from these abstract conditions. Actually, this general protocol can be instantiated in distinct ways, giving rise to a family of communication-induced snapshot protocols. This general protocol shows there is an intrinsic trade-off between the number of forced checkpoints and the number of recorded messages. Finally, a particular instantiation of the general protocol is provided.
Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Parallel Distributed Syst.3
1999 Consistency Issues in Distributed Checkpoints
abstract
A global checkpoint is a set of local checkpoints, one per process. The traditional consistency criterion for global checkpoints states that a global checkpoint is consistent if it does not include messages received and not sent. The paper investigates other consistency criteria, transitlessness, and strong consistency. A global checkpoint is transitless if it does not exhibit messages sent and not received. Transitlessness can be seen as a dual of traditional consistency. Strong consistency is the addition of transitlessness to traditional consistency. The main result of the paper is a statement of the necessary and sufficient condition answering the following question: "given an arbitrary set of local checkpoints, can this set be extended to a global checkpoint that satisfies P" (where P is traditional consistency, transitlessness, or strong consistency). From a practical point of view, this condition, when applied to transitlessness, is particularly interesting as it helps characterize which messages do not need to be recorded by checkpointing protocols.
Jean-Michel Hélary, Robert H. B. Netzer, Michel Raynal
IEEE Trans. Software Eng.3
1998 An Adaptive Protocol for Implementing Causally Consistent Distributed Services
abstract
Distributed services that are accessed by widely distributed clients are becoming common place. Such services cannot be provided at the needed level of performance and availability without replicating the service at multiple nodes, and without allowing a relatively weak level of consistency among replicated copies of the state of a service. This paper explores causally consistent distributed services when multiple related services are replicated to meet performance and availability requirements. This consistency criterion is particularly well suited for some distributed services (e.g., cooperative document sharing), and it is attractive because of the efficient implementations allowed by it.
Mustaque Ahamad, Michel Raynal, Gérard Thia-Kime
ICDCS2
1998 Asynchronous Protocols to Meet Real-Time Constraints: Is It Really Sensible? How to Proceed?
abstract
This paper investigates the use of asynchronous protocols to design and build middleware whose aim is to provide run-time support for soft real-time applications. A simple and general framework is described. This framework allows to take into account timeliness constraints of upper layer applications, while using an asynchronous protocol at the underlying level. When a timeliness constraint is about to be violated, the application layer is informed and can take appropriate measures. The deadline period can then be extended if the constraint is soft enough; in the other case, a default value can be used as a result. This framework can be seen as a bridge from asynchronous systems to synchronous ones. The proposed approach is illustrated with the consensus problem. This approach is investigated in the ARGO system we are implementing. The target applications of the ARGO middleware are telecommunication applications.
Michel Hurfin, Michel Raynal
ISORC2
1998 Fault-Tolerant Total Order Multicast to Asynchronous Groups
abstract
While Total Order Broadcast (or Atomic Broadcast) primitives have received a lot of attention, the paper concentrates on Total Order Multicast to Multiple Groups in the context of asynchronous distributed systems in which processes may suffer crash failures. "Multicast to Multiple Groups" means that each message is sent to a subset of the process groups composing the system, distinct messages possibly having distinct destination groups. "Total Order" means that all message deliveries must be totally ordered. The paper proposes a protocol for such a multicast primitive. This protocol is based on two underlying building blocks, namely, Uniform Reliable Multicast and Uniform Consensus. Its design characteristics lie in the two following properties. The first one is a minimality property, more precisely, only the sender of a message and processes of its destination groups have to participate in the multicast of the message. The second property is a locality property: no execution of a consensus has to involve processes belonging to distinct groups (i.e., consensus are executed on a "per group" basis). This locality property is particularly useful when one is interested in using the Total Order Multicast primitive in large scale distributed systems. An improvement that reduces the cost of the protocol is also suggested.
Udo Fritzke Jr., Philippe Ingels, Achour Mostéfaoui, Michel Raynal
SRDS4
1998 Consensus in Asynchronous Systems Where Processes Can Crash and Recover
abstract
The consensus problem is now well identified as being one of the most important problems encountered in the design and the construction of fault-tolerant distributed systems. This problem is defined as follows: processes have to reach a common decision, which depends on their inputs, despite failures. We consider the consensus problem in asynchronous distributed systems augmented with unreliable failure detectors. Several protocols have been proposed for these systems, when process crashes are assumed to be definitive. This paper addresses the consensus problem in a more practical asynchronous system model, namely in a context where processes can crash and recover. As a process crash entails the loss of its volatile memory, each process is equipped with a stable storage. So, to be efficient a consensus protocol has to log as few critical data as possible. The proposed protocol uses a new class of failure detectors suited to the crash/recovery model. It is particularly efficient when, whether there are crashes or not, the underlying failure detector makes few mistakes. Additionally, the proposed protocol tolerates message duplication and copes with some message losses.
Michel Hurfin, Achour Mostéfaoui, Michel Raynal
SRDS3
1998 Lifetime Based Consistency Protocols for Distributed Objects
Francisco J. Torres-Rojas, Mustaque Ahamad, Michel Raynal
DISC3
1998 Consistent Records in Asynchronous Computations
Roberto Baldoni, Jean-Michel Hélary, Michel Raynal
Acta Informatica3
1998 k-Arbiter: A Safe and General Scheme for h-out of-k Mutual Exclusion
Yoshifumi Manabe, Roberto Baldoni, Michel Raynal, Shigemi Aoyagi
Theor. Comput. Sci.3
1998 Efficient Distributed Detection of Conjunctions of Local Predicates
abstract
Global predicate detection is a fundamental problem in distributed systems and finds applications in many domains such as testing and debugging distributed programs. This paper presents an efficient distributed algorithm to detect conjunctive-form global predicates in distributed systems. The algorithm detects the first consistent global state that satisfies the predicate even if the predicate is unstable. Unlike previously proposed run-time predicate detection algorithms, our algorithm does not require the exchange of control messages during the normal computation. All the necessary information to detect predicates is piggybacked on computation messages of application programs. The algorithm is distributed because the predicate detection efforts as well as the necessary information are equally distributed among the processes. We prove the correctness of the algorithm and compare its performance with respect to message, storage and computational complexities with that of the previously proposed run-time predicate detection algorithms.
Michel Hurfin, Masaaki Mizuno, Michel Raynal, Mukesh Singhal
IEEE Trans. Software Eng.3
1997 Cycle Prevention in Distributed Checkpointing
Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
OPODIS3
1997 Preventing Useless Checkpoints in Distributed Computations
abstract
A useless checkpoint is a local checkpoint that cannot be part of a consistent global checkpoint. The paper addresses the following important problem. Given a set of processes that take (basic) local checkpoints in an independent and unknown way, the problem is to design a communication induced checkpointing protocol that directs processes to take additional local (forced) checkpoints to ensure that no local checkpoint is useless. A general and efficient protocol answering this problem is proposed. It is shown that several existing protocols that solve the same problem are particular instances of it. The design of this general protocol is motivated by the use of communication induced checkpointing protocols in "consistent global checkpoint" based distributed applications. Detection of stable or unstable properties, rollback recovery and determination of distributed breakpoints are examples of such applications.
Jean-Michel Hélary, Achour Mostéfaoui, Robert H. B. Netzer, Michel Raynal
SRDS4
1997 Shared Global States in Distributed Computations
Eddy Fromentin, Michel Raynal
J. Comput. Syst. Sci.2
1997 An Adaptive Causal Ordering Algorithm Suited to Mobile Computing Environments
abstract
Causal message ordering is required for several distributed applications. In order to preserve causal ordering, only direct dependency information between messages, with respect to the destination process(es), need be sent with each message. By eliminating other kinds of control information from the messages, the communication overheads can be significantly reduced. In this paper we present an algorithm that uses this knowledge to efficiently enforce causal ordering of messages. The proposed algorithm does not require any prior knowledge of the network topology or communication pattern. As computation proceeds, it acquires knowledge of the communication pattern and is capable of handling dynamically changing multicast communication groups, and minimizing the communication overheads. With regard to communication overheads, the algorithm is optimal for the broadcast communication case. Extensive simulation experiments demonstrate that the algorithm imposes lower communication overheads than previous causal ordering algorithms. The algorithm can be employed in a variety of distributed computing environments. Its energy efficiency and low bandwidth requirement make it especially suitable for mobile computing systems. We show how to employ the algorithm for causally ordered multicasting of messages in mobile computing environments.
Ravi Prakash 0001, Michel Raynal, Mukesh Singhal
J. Parallel Distributed Comput.2
1996 An Efficient Causal Ordering Algorithm for Mobile Computing Environments
abstract
Causal message ordering is required for several distributed applications. In order to preserve causal ordering, only direct dependency information between messages with respect to the destination process(es) should be sent with each message. By eliminating other kinds of control information from the messages, the communication overheads can be significantly reduced. In this paper we present an algorithm that uses this knowledge to efficiently enforce causal ordering of messages. The proposed algorithm does not require any prior knowledge of the network or communication topology. As computation proceeds, it acquires knowledge of the logical communication topology and is capable of handling dynamically changing multicast communication groups. With regard to communication overheads, the algorithm is optimal for the broadcast communication case. Its energy efficiency and four bandwidth requirement make it suitable for mobile computing systems. We present a strategy that employs the algorithm for causally ordered multicasting of messages in mobile computing environments.
Ravi Prakash 0001, Michel Raynal, Mukesh Singhal
ICDCS2
1996 About State Recording in Asynchronous Computations (Abstract)
abstract
No abstract available.
Roberto Baldoni, Jean-Michel Hélary, Michel Raynal
PODC3
1996 Efficient Delta-Causal Broadcasting of Multimedia Applications (Abstract)
abstract
No abstract available.
Roberto Baldoni, Ravi Prakash 0001, Michel Raynal, Mukesh Singhal
PODC3
1996 From Serializable to Causal Transactions (Abstract)
abstract
No abstract available.
Michel Raynal, Gérard Thia-Kime, Mustaque Ahamad
PODC1
1996 Detecting Diamond Necklaces in Labeled Dags (A Problem from Distributed Debugging)
Michel Hurfin, Michel Raynal
WG2
1996 Erratum: Deadlock Models and a General Algorithm for Distributed Deadlock Detection
Jerzy Brzezinski, Jean-Michel Hélary, Michel Raynal, Mukesh Singhal
J. Parallel Distributed Comput.3
1996 A unified framework for the specification and run-time detection of dynamic properties in distributed computations
Özalp Babaoglu, Eddy Fromentin, Michel Raynal
J. Syst. Softw.3
1996 Causal Delivery of Messages with Real-Time Data in Unreliable Networks
Roberto Baldoni, Achour Mostéfaoui, Michel Raynal
Real Time Syst.3
1995 From Causal Consistency to Sequential Consistency in Shared Memory Systems
Michel Raynal, André Schiper
FSTTCS1
1995 Efficient Causally Ordered Communications for Multimedia Real-Time Applications
abstract
Multimedia real-time collaborative applications or groupware real-time applications require participants to exchange real-time audio and video information over a communication network. This flow of information must preserve the causal dependency even though part of the information can be lost or can be discarded if it violates the tinting constraints imposed by a real-time interaction. In this paper we propose a communication abstraction to cope with unreliable communication networks with real-time delivery constraints: messages have a lifetime, /spl Delta/, after which their contents can no longer be used, moreover some of them can be lost. This new abstraction, called /spl Delta/-causal order, requires to deliver as much messages as possible within their lifetime in such a way that these deliveries respect causal order. An efficient protocol is proposed in the case of one-to-one communications. A variation of this protocol weld suited to broadcast communications is also shown.
Roberto Baldoni, Achour Mostéfaoui, Michel Raynal
HPDC3
1995 Characterizing and Detecting The Set of Global States Seen by all Observers of a Distributed Computation
abstract
A consistent observation of a given distributed computation is a sequence of global states that could be produced by executing that computation on a monoprocessor system. Therefore a distributed execution generally accepts several consistent observations. This paper concentrates on what all these observations have in common. An abstraction called common global state is defined. A necessary and sufficient condition characterizing such states is given. A monitor-based algorithm that detects them is also presented and proved correct. Previous works on detection of unstable properties of distributed computations are revisited and explained with this abstraction. Moreover other uses of such particular states are sketched.
Eddy Fromentin, Michel Raynal
ICDCS2
1995 Debugging Distributed Executions by Using Language Recognition
Özalp Babaoglu, Eddy Fromentin, Michel Raynal
ICPP (2)3
1995 On-The-Fly Analysis of Distributed Computations
Eddy Fromentin, Claude Jard, Guy-Vincent Jourdan, Michel Raynal
Inf. Process. Lett.4
1995 Specification and Verification of Dynamic Properties in Distributed Computations
Özalp Babaoglu, Michel Raynal
J. Parallel Distributed Comput.2
1995 Deadlock Models and a General Algorithm for Distributed Deadlock Detection
Jerzy Brzezinski, Jean-Michel Hélary, Michel Raynal, Mukesh Singhal
J. Parallel Distributed Comput.3
1994 On the Fly Testing of Regular Patterns in Distributed Computations
abstract
A class of properties of distributed computations is described and an algorithm which detects them is presented. This class of properties called regular patterns allows the user to specify an expected (or unwanted) behavior of a computation as sequences of relevant events (or as sequences of local predicates that must be successively verified). The sequences are defined by a finite state automaton (hence the name regular patterns) A computation verifies the property if and only if one of its causal paths matches a sequence.
Eddy Fromentin, Michel Raynal, Vijay K. Garg, Alexander I. Tomlinson
ICPP (2)2
1994 Towards the Construction of Distributed Detection Programs, with an Application to Distributed Termination
Jean-Michel Hélary, Michel Raynal
Distributed Comput.2
1994 A General Scheme for Token- and Tree-Based Distributed Mutual Exclusion Algorithms
abstract
In a distributed context, mutual exclusion algorithms can be divided into two families according to their underlying algorithmic principles: those that are permission-based and those that are token-based. Within the latter family, a lot of algorithms use a rooted tree structure to move the requests and the unique token. This paper presents a very general information structure (and the associated generic algorithm) for token- and tree-based mutual exclusion algorithms. This general structure not only covers, as particular cases, several known algorithms, but also allows for the design of new ones that are well suited for various topology requirements.>
Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Parallel Distributed Syst.3
1993 Termination Detection in a Very General Distributed Computing Model
abstract
Termination detection constitutes one of the basic problems of distributed computing, and many distributed algorithms have been proposed to solve it, but all these algorithms consider a very simple model for the underlying application programs: for processes of such programs, nondeterministic constructs are allowed, but each 'receive' statement (request) concerns only one message at a time. A more realistic and very general model of distributed computing is first presented, allowing a request to be atomic on several messages and to obey AND/OR/AND-OR/k-out-of-n/etc. request types. Within this framework, two definitions of termination are proposed and discussed. Then, accordingly, two distributed algorithms for detecting these terminations are presented and evaluated; they differ in the information they use and in the time they need to claim termination.>
Jerzy Brzezinski, Jean-Michel Hélary, Michel Raynal
ICDCS3
1993 Debugging tool for distributed Estelle programs
Michel Hurfin, Noël Plouzeau, Michel Raynal
Comput. Commun.3
1992 A General Method to Define Quorums
abstract
Composition, a general method for constructing quorum sets, coteries, and bicoteries, is discussed. It is shown that composition provides a natural method for constructing quorum structures in an arbitrary network or even in a collection of interconnected networks, and that the resulting structures, called composite structures, can be efficiently evaluated. In particular, an efficient method for determining if a given set contains a quorum of a composite structure is presented. If this method is used, it is not necessary to compute and store all of the quorums of the composite structure in advance.>
Mitchell L. Neilsen, Masaaki Mizuno, Michel Raynal
ICDCS3
1992 Synchronization and Concurrency Measures for Distributed Computations
abstract
Several qualitative measures that quantify the degree of concurrency in a distributed computation are presented. The measures characterize the synchronization constraints inherent in a distributed computation and are independent of the underlying system running the computation. The measures are defined by two well defined abstractions, called cone and cylinder, to which simple measures can be associated: volume, weight, and height. Simple ways to compute the measures are proposed. The mechanism uses two types of vector clocks that trace the history of the computation. It is shown that the measures can be easily incorporated into any system to analyze distributed executions.>
Michel Raynal, Masaaki Mizuno, Mitchell L. Neilsen
ICDCS1
1991 The Causal Ordering Abstraction and a Simple Way to Implement it
Michel Raynal, André Schiper, Sam Toueg
Inf. Process. Lett.1
1989 Prime Numbers as a Tool to Design Distributed Algorithms
Michel Raynal
Inf. Process. Lett.1
1988 A Distributed Algorithm for Mutual Exclusion in an Arbitrary Network
abstract
A distributed algorithm for mutual exclusion is presented. No particular assumptions on the network topology are required, except connectivity; the communication graph may be arbitrary. The processes communicate by using messages only and there is no global controller. Furthermore, no process needs to know or learn the global network topology. In that sense, the algorithm is more general than the mutual exclusion algorithms which make use of an a priori knowledge of the network topology (for example either ring or complete network). A proof of the correctness of the algorithm is provided. The algorithm's complexity is examined by evaluating the number of messages required for the mutual exclusion protocol.
Jean-Michel Hélary, Noël Plouzeau, Michel Raynal
Comput. J.3
1987 Detection of Stable Properties in Distributed Applications
abstract
When evaluated to true, a stable property remains true forever.Such a stable property may characterize important states of a computation.This is the case of deadlocked or terminated computations.In this paper we expose a general algorithm for the distributed detection of stable properties in distributed applications or systems.This distributed algorithm deals with every stable property of a fairly general class : in this sense the algorithm is generic.This was achieved using a methodical approach, with a strong distinction between the computation and control activities in the problem.Moreover, the detection method used by the algorithm is based on an observational mechanism,
Jean-Michel Hélary, Claude Jard, Noël Plouzeau, Michel Raynal
PODC4
1987 A Distributed Algorithm to Prevent Mutual Drift Between n Logical Clocks
Michel Raynal
Inf. Process. Lett.1
1983 Structured Specification of Communicating Systems
abstract
Specification methods for distributed systems is the underlying theme of this paper. A model of communicating processes with rendezvous interactions is assumed as a basis for the discussion. The possible interactions by a process, and the interconnection between several subprocesses within a process are specified using the concept of ports, which are specified separately. Step-wise refinement of process specifications and associated verification rules are considered. The step-wise refinement of port specifications and associated interactions is considered as well. After the presentation of an introductory example, the paper discusses the basic concepts of the specification method. They are then applied to more complex examples. The step-wise wefinement of ports and interactions is demonstrated by a hardware interface for which an abstract specification and a more detailed implementation is given. Proof rules for verifying the consistency of detailed and more abstract specifications are discussed in some detail.
Gregor von Bochmann, Michel Raynal
IEEE Trans. Computers2
1981 An Experience in Implementing Abstract Data Types
abstract
Abstract The abstract data type concept appears to be a useful software structuring tool. A project, called ‘Système d'Objets Conservés’, which was developed at the University of Rennes, (France), gave some experience in implementing this concept. The possibility of including abstract data type into a pre‐existing compiler is demonstrated, and desirable properties of the host language are exhibited. Provision of external procedures and data makes some type checking extensions necessary: these features increase software reliability.
Michel Banâtre, André Couvert, D. Herman, Michel Raynal
Softw. Pract. Exp.4