Achour Mostéfaoui

dblp:m/AchourMostefaoui · DBLP profile ↗
← Back
123ranked-venue papers
59as first author
13since 2021 · last 2026
0000-0001-7208-4635ORCID · verified

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

Systems, architecture and hardware · 54 · 29 first-author · 6 since 2021Security and privacy · 21 · 11 first-author · 1 since 2021Theory of computation · 21 · 12 first-author · 2 since 2021Databases, data management, data science and information retrieval · 6 · 2 first-authorSoftware engineering, systems software and programming languages · 5 · 4 first-authorApplied, interdisciplinary, general and emerging computing · 3 · 2 first-author
YearPublicationVenuePosition
2026 Byzantine-tolerant privacy-preserving atomic register
abstract
This paper extends and improves upon our work [Kowalski et al., ICDCN 2025], which proposed the construction of a privacy-preserving single-writer multi-reader (SWMR) atomic register in a Byzantine-prone distributed model. Specifically, we consider a closed model in which one process can write values in the register and only a subset of the other processes are allowed to read them. The goal is to ensure that processes without the requisite read permission are unable to read the content of the register, even when they are Byzantine. This guarantees the privacy of the stored value. We achieve this privacy by encoding the value written by the writer using secret sharing, thereby splitting it into multiple shards that are disseminated among the participating reader processes. The technical challenge is then to coordinate the correct reader processes so as to achieve Byzantine linearizability without revealing the register’s contents. The main contribution of this work is an improved resilience bound of the linearizable read-write (R/W) privacy-preserving register construction from t < n 7 to t < n 5 , where t is the number of Byzantine processes and n denotes the total number of processes in the system. Despite being more resilient than the previous version, the new construction algorithm is significantly simpler and more appealing, and it comes with a clearer and more concise correctness proof.
Vincent Kowalski, Achour Mostéfaoui, Matthieu Perrin, Sinchan Sengupta
Theor. Comput. Sci.2
2025 Invited Paper: On the Equivalence of Snapshot/Append Objects and Broadcast Abstractions under Byzantine Failures
Vincent Kowalski, Achour Mostéfaoui, Matthieu Perrin, Jolan Riallo
SSS2
2024 No Symmetric Broadcast Abstraction Characterizes k-Set-Agreement in Message-Passing Systems
Sylvain Gay, Achour Mostéfaoui, Matthieu Perrin
OPODIS2
2024 Brief Announcement: No Broadcast Abstraction Characterizes k-Set-Agreement in Message-Passing Systems
abstract
This paper explores the relationship between broadcast abstractions and the k-set agreement (k-SA) problem in crash-prone asynchronous message-passing distributed systems. It specifically investigates whether any broadcast abstraction is computationally equivalent to k-SA in message-passing systems. A key contribution of the paper is the introduction of a clear definition of admissible broadcast abstractions, achieved by introducing two new symmetry properties: compositionality and content-neutrality. The paper's primary contribution is the demonstration that no broadcast abstraction, which is both content-neutral and compositional, is computationally equivalent to k-set agreement when 1 < k < n.
Sylvain Gay, Achour Mostéfaoui, Matthieu Perrin
PODC2
2024 Brief Announcement: Randomized Consensus: Common Coins Are not the Holy Grail!
abstract
This paper studies the round complexity of randomized binary consensus in crash-prone asynchronous distributed systems. While the Consensus problem cannot be solved deterministically, Ben-Or and Rabin showed that randomization allows solving the problem with probability 1. Moreover, while local coins may need an exponential number of rounds in n, a common coin that delivers the same random sequence to all processes allows termination within a constant mean number of rounds. This paper studies the round complexity and the optimality for different coins. Surprisingly, while the common coin is optimal when t > n/3, it is not when t ≤ n/3.
Achour Mostéfaoui, Matthieu Perrin, Julien Weibel
PODC1
2023 Atomic Register Abstractions for Byzantine-Prone Distributed Systems
Vincent Kowalski, Achour Mostéfaoui, Matthieu Perrin
OPODIS2
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
PODC2
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
DISC2
2023 Differentiated Consistency for Worldwide Gossips
abstract
Eventual consistency is a consistency model that favors liveness over safety. It is often used in large-scale distributed systems where models ensuring a stronger safety incur performance that are too low to be deemed practical. Eventual consistency tends to be uniformly applied within a system, but we argue a demand exists for differentiated eventual consistency, e.g. in blockchain systems. We propose update-query consistency with primaries and secondaries (UPS) to address this demand. UPS is a novel consistency mechanism that works in pair with our novel two-phase epidemic broadcast protocol gossip primary-secondary (GPS) to offer differentiated eventual consistency and delivery speed. We propose two complementary analyses of the broadcast protocol: a continuous analysis and a discrete analysis based on compartmental models used in epidemiology. Additionally, we propose the formal definition of a scalable consistency metric to measure the consistency trade-off at runtime. We evaluate UPS in two simulated worldwide settings: a one-million-node network and a network emulating that of the Ethereum blockchain. In both settings, UPS reduces inconsistencies experienced by a majority of the nodes and reduces the average message latency for the remaining nodes.
Davide Frey, Achour Mostéfaoui, Matthieu Perrin, Pierre-Louis Roman, François Taïani
IEEE Trans. Parallel Distributed Syst.2
2022 Extending the wait-free hierarchy to multi-threaded systems
Matthieu Perrin, Achour Mostéfaoui, Grégoire Bonin, Ludmila Courtillat-Piazza
Distributed Comput.2
2021 Wait-Free CAS-Based Algorithms: The Burden of the Past
abstract
Herlihy proved that CAS is universal in the classical computing system model composed of an a priori known number of processes. This means that CAS can implement, together with reads and writes, any object with a sequential specification. For this, he proposed the first universal construction capable of emulating any data structure. It has recently been proved that CAS is still universal in the infinite arrival computing model, a model where any number of processes can be created on the fly (e.g. multi-threaded systems). In this paper, we prove that CAS does not allow to implement wait-free and linearizable visible objects in the infinite model with a space complexity bounded by the number of active processes (i.e. ones that have operations in progress on this object). This paper also shows that this lower bound is tight, in the sense that this dependency can be made as low as desired (e.g. logarithmic) by proposing a wait-free and linearizable universal construction, using the compare-and-swap operation, whose space complexity in the number of ever issued operations is defined by a parameter that can be linked to any unbounded function.
Denis Bédin, François Lépine, Achour Mostéfaoui, Damien Perez, Matthieu Perrin
DISC3
2021 A scalable sequence encoding for collaborative editing
abstract
Summary Distributed real‐time editors made real‐time editing easy for millions of users. However, main stream editors rely on Cloud services to mediate sessions raising privacy and scalability issues. Decentralized editors tackle privacy issues, but scalability issues remain. We aim to build a decentralized editor that allows real‐time editing anytime, anywhere, whatever is the number of participants. In this study, we propose an approach based on a massively replicated sequence data structure that represents the shared document. We establish an original trade‐off on communication, time, and space complexity to maintain this sequence over a network of browsers. We prove a sublinear upper bound on communication complexity while preserving an affordable time and space complexity. To validate this trade‐off, we built a full working editor and measured its performance on large‐scale experiments involving up till 600 participants. As expected, the results show a traffic increasing as whereIis the number of insertions in the document, andRthe number of participants.
Brice Nédelec, Pascal Molli, Achour Mostéfaoui
Concurr. Comput. Pract. Exp.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.2
2020 Extending the Wait-free Hierarchy to Multi-Threaded Systems
abstract
In modern operating systems and programming languages adapted to multicore computer architectures, parallelism is abstracted by the notion of execution threads. Multi-threaded systems have two major specificities: 1) new threads can be created dynamically at runtime, so there is no bound on the number of threads participating in a long-running execution. 2) threads have access to a memory allocation mechanism that cannot allocate infinite arrays. This makes it challenging to adapt some algorithms to multi-threaded systems, especially those that assign one shared register per process.
Matthieu Perrin, Achour Mostéfaoui, Grégoire Bonin
PODC2
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
PRDC1
2019 Brief Announcement: Wait-Free Universality of Consensus in the Infinite Arrival Model
abstract
In classical asynchronous distributed systems composed of a fixed number n of processes where some proportion may fail by crashing, many objects do not have a wait-free linearizable implementation (e.g. stacks, queues, etc.). It has been proved that consensus is universal in such systems, which means that this system augmented with consensus objects allows to implement any object that has a sequential specification. In this paper, we consider a more general system model called infinite arrival model where infinitely many processes may arrive and leave or crash during a run. We prove that consensus is still universal in this more general model. For that, we propose a universal construction based on a weak log that can be implementated using consensus objects.
Grégoire Bonin, Achour Mostéfaoui, Matthieu Perrin
DISC2
2019 Crash-tolerant causal broadcast in O(n) messages
Achour Mostéfaoui, Matthieu Perrin, Michel Raynal, Jiannong Cao 0001
Inf. Process. Lett.1
2018 Causal Broadcast: How to Forget?
abstract
Causal broadcast constitutes a fundamental communication primitive of many distributed protocols and applications. However, state-of-the-art implementations fail to forget obsolete control information about already delivered messages. They do not scale in large and dynamic systems. In this paper, we propose a novel implementation of causal broadcast. We prove that all and only obsolete control information is safely removed, at cost of a few lightweight control messages. The local space complexity of this protocol does not monotonically increase and depends at each moment on the number of messages still in transit and the degree of the communication graph. Moreover, messages only carry a scalar clock. Our implementation constitutes a sustainable communication primitive for causal broadcast in large and dynamic systems.
Brice Nédelec, Pascal Molli, Achour Mostéfaoui
OPODIS3
2018 Breaking the Scalability Barrier of Causal Broadcast for Large and Dynamic Systems
abstract
Many distributed protocols and applications rely on causal broadcast to ensure consistency criteria. However, none of causality tracking state-of-the-art approaches scale in large and dynamic systems. This paper presents a new non-blocking causal broadcast protocol suited for such systems. The proposed protocol outperforms state-of-the-art in size of messages, execution time complexity, and local space complexity. Most importantly, messages piggyback control information the size of which is constant. We prove that for both static and dynamic systems. Consequently, large and dynamic systems can finally afford causal broadcast.
Brice Nédelec, Pascal Molli, Achour Mostéfaoui
SRDS3
2018 Randomized k-set agreement in crash-prone and Byzantine asynchronous systems
Achour Mostéfaoui, Moumen Hamouma, Michel Raynal
Theor. Comput. Sci.1
2018 An adaptive peer-sampling protocol for building networks of browsers
Brice Nédelec, Julian Tanke, Davide Frey, Pascal Molli, Achour Mostéfaoui
World Wide Web5
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
DISC2
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 Informatica1
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.1
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
PODC1
2016 Causal consistency: beyond memory
abstract
In distributed systems where strong consistency is costly when not impossible, causal consistency provides a valuable abstraction to represent program executions as partial orders. In addition to the sequential program order of each computing entity, causal order also contains the semantic links between the events that affect the shared objects -- messages emission and reception in a communication channel, reads and writes on a shared register. Usual approaches based on semantic links are very difficult to adapt to other data types such as queues or counters because they require a specific analysis of causal dependencies for each data type. This paper presents a new approach to define causal consistency for any abstract data type based on sequential specifications. It explores, formalizes and studies the differences between three variations of causal consistency and highlights them in the light of PRAM, eventual consistency and sequential consistency: weak causal consistency, that captures the notion of causality preservation when focusing on convergence; causal convergence that mixes weak causal consistency and convergence; and causal consistency, that coincides with causal memory when applied to shared memory.
Matthieu Perrin, Achour Mostéfaoui, Claude Jard
PPoPP2
2016 Speed for the Elite, Consistency for the Masses: Differentiating Eventual Consistency in Large-Scale Distributed Systems
abstract
Eventual consistency is a consistency model that emphasizes liveness over safety, it is often used for its ability to scale as distributed systems grow larger. Eventual consistency tends to be uniformly applied to an entire system, but we argue that there is a growing demand for differentiated eventual consistency requirements. We address this demand with UPS, a novel consistency mechanism that offers differentiated eventual consistency and delivery speed by working in pair with a two-phase epidemic broadcast protocol. We propose a closed-form analysis of our approach's delivery speed, and we evaluate our complete mechanism experimentally on a simulated network of one million nodes. To measure the consistency trade-off, we formally define a novel and scalable consistency metric that operates at runtime. In our simulations, UPS divides by more than 4 the inconsistencies experienced by a majority of the nodes, while reducing the average latency incurred by a small fraction of the nodes from 6 rounds down to 3 rounds.
Davide Frey, Achour Mostéfaoui, Matthieu Perrin, Pierre-Louis Roman, François Taïani
SRDS2
2016 On Composition and Implementation of Sequential Consistency
Matthieu Perrin, Matoula Petrolia, Achour Mostéfaoui, Claude Jard
DISC3
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.1
2015 Update Consistency for Wait-Free Concurrent Objects
abstract
In large scale systems such as the Internet, replicating data is an essential feature in order to provide availability and fault-tolerance. Attila and Welch proved that using strong consistency criteria such as atomicity is costly as each operation may need an execution time linear with the latency of the communication network. Weaker consistency criteria like causal consistency and PRAM consistency do not ensure convergence. The different replicas are not guaranteed to converge towards a unique state. Eventual consistency guarantees that all replicas eventually converge when the participants stop updating. However, it fails to fully specify the semantics of the operations on shared objects and requires additional non-intuitive and error-prone distributed specification techniques. This paper introduces and formalizes a new consistency criterion, called update consistency, that requires the state of a replicated object to be consistent with a linearization of all the updates. In other words, whereas atomicity imposes a linearization of all of the operations, this criterion imposes this only on updates. Consequently some read operations may return out-dated values. Update consistency is stronger than eventual consistency, so we can replace eventually consistent objects with update consistent ones in any program. Finally, we prove that update consistency is universal, in the sense that any object can be implemented under this criterion in a distributed system where any number of nodes may crash.
Matthieu Perrin, Achour Mostéfaoui, Claude Jard
IPDPS2
2015 A Message-Passing and Adaptive Implementation of the Randomized Test-and-Set Object
abstract
This paper presents a solution to the well-known Test-and-Set operation in asynchronous systems prone to process crashes. Test-and-Set is a synchronization operation that, when invoked by a set of processes, returns "yes" to a unique process and returns "no" to all the others. Recently many advances in implementing Test and Set objects have been achieved, however all of them uniquely target the shared memory model. In this paper we propose an implementation of a Test-and-Set object for message passing distributed systems. This implementation can be invoked by any number p of processes. It has an expected step complexity in O(p) and an expected message complexity in O(np), where n is the total number of processes in the system. The proposed Test and Set object is built atop a new basic building block that allows to select a winning group among two groups of processes.
Emmanuelle Anceaume, François Castella, Achour Mostéfaoui, Bruno Sericola
NCA3
2015 Efficiently Summarizing Data Streams over Sliding Windows
abstract
Estimating the frequency of any piece of information in large-scale distributed data streams became of utmost importance in the last decade (e.g., in the context of network monitoring, big data, etc.). If some elegant solutions have been proposed recently, their approximation is computed from the inception of the stream. In a runtime distributed context, one would prefer to gather information only about the recent past. This may be led by the need to save resources or by the fact that recent information is more relevant. In this paper, we consider the sliding window model and propose two different (on-line) algorithms that approximate the items frequency in the active window. More precisely, we determine a (ε, δ)-additive-approximation meaning that the error is greater than ε only with probability δ. These solutions use a very small amount of memory with respect to the size N of the window and the number n of distinct items of the stream, namely, O(1/ε log 1/δ (log N+log n)) and O(1/τε log 1/δ (log N+log n)) bits of space, where τ is a parameter limiting memory usage. We also provide their distributed variant, i.e., considering the sliding window functional monitoring model. We compared the proposed algorithms to each other and also to the state of the art through extensive experiments on synthetic traces and real data sets that validate the robustness and accuracy of our algorithms.
Nicolo Rivetti, Yann Busnel, Achour Mostéfaoui
NCA3
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
PODC2
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
SIROCCO1
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. ACM1
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
PODC1
2014 Update Consistency in Partitionable Systems
Matthieu Perrin, Achour Mostéfaoui, Claude Jard
DISC2
2013 LSEQ: an adaptive structure for sequences in distributed collaborative editing
abstract
Distributed collaborative editing systems allow users to work distributed in time, space and across organizations. Trending distributed collaborative editors such as Google Docs, Etherpad or Git have grown in popularity over the years. A new kind of distributed editors based on a family of distributed data structure replicated on several sites called Conflict-free Replicated Data Type (CRDT for short) appeared recently. This paper considers a CRDT that represents a distributed sequence of basic elements that can be lines, words or characters (sequence CRDT). The possible operations on this sequence are the insertion and the deletion of elements. Compared to the state of the art, this approach is more decentralized and better scales in terms of the number of participants. However, its space complexity is linear with respect to the total number of inserts and the insertion points in the document. This makes the overall performance of such editors dependent on the editing behaviour of users. This paper proposes and models LSEQ, an adaptive allocation strategy for a sequence CRDT. LSEQ achieves in the average a sub-linear spatial-complexity whatever is the editing behaviour. A series of experiments validates LSEQ showing that it outperforms existing approaches.
Brice Nédelec, Pascal Molli, Achour Mostéfaoui, Emmanuel Desmontils
ACM Symposium on Document Engineering3
2013 Topic 8: Distributed Systems and Algorithms - (Introduction)
Achour Mostéfaoui, Andreas Polze, Carlos Baquero, Paul D. Ezhilchelvan, Lars Lundberg
Euro-Par1
2013 Synchronous byzantine agreement with nearly a cubic number of communication bits: synchronous byzantine agreement with nearly a cubic number of communication bits
abstract
This paper studies the problem of Byzantine consensus in a synchronous message-passing system of n processes. The first deterministic algorithm, and also the simplest in its principles, was the Exponential Information Gathering protocol (EIG) proposed by Pease, Shostak and Lamport in [19]. The algorithm requires processes to send exponentially long messages. Many follow-up works reduced the cost of the algorithm. However, they had to either lower the maximum number of faulty processes t from the optimal range t < n/3 to some smaller range of t [4, 11, 18], or increase the maximum worst-case number of rounds needed for termination (the lower bound being t + 1) [3, 9, 20].
Dariusz R. Kowalski, Achour Mostéfaoui
PODC2
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
NCA1
2011 Relations Linking Failure Detectors Associated with k-Set Agreement in Message-Passing Systems
Achour Mostéfaoui, Michel Raynal, Julien Stainer
SSS1
2010 Time-Free Authenticated Byzantine Consensus
abstract
This paper presents a simple protocol that solves the authenticated Byzantine Consensus problem in asynchronous distributed systems. To circumvent the FLP impossibility result in a deterministic way, synchrony assumptions should be added. In the context of Byzantine failures for systems where at most t processes may exhibit a Byzantine behavior and where not all the system is assumed eventually synchronous, Moumen et al. provide the main result. They assume at least one correct process, called 2t-bisource, connected with 2t privileged neighbors with eventually timely outgoing and incoming links. The present paper shows that a deterministic solution for the authenticated byzantine consensus problem is possible if the system model satisfies an additional assumption that does not rely on physical time but on the pattern of messages that are exchanged. The basic message exchange between processes is the query-response mechanism. To solve the Consensus problem, we assume a correct process p, called 2t-winning process, and a set Q of 2t processes such that, eventually, for each query issued by p, any process q of Q receives a response from p among the (n - t) first responses to that query. The processes in the set Q can exhibit a Byzantine behavior and this set may change over time. Whereas many time-free solutions have been designed for the consensus problem in the crash model, this is, to our knowledge, the first time-free deterministic solution to the Byzantine consensus problem.
Moumen Hamouma, Achour Mostéfaoui
NCA2
2010 Signature-Free Broadcast-Based Intrusion Tolerance: Never Decide a Byzantine Value
Achour Mostéfaoui, Michel Raynal
OPODIS1
2010 Narrowing power vs efficiency in synchronous set agreement: Relationship, algorithms and lower bound
Achour Mostéfaoui, Michel Raynal, Corentin Travers
Theor. Comput. Sci.1
2009 What Agreement Problems Owe Michel
Achour Mostéfaoui
DISC1
2009 From adaptive renaming to set agreement
Eli Gafni, Achour Mostéfaoui, Michel Raynal, Corentin Travers
Theor. Comput. Sci.2
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.1
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
PODC2
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.1
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.1
2007 Topic 8 Distributed Systems and Algorithms
Luís E. T. Rodrigues, Achour Mostéfaoui, Christof Fetzer, Philippas Tsigas
Euro-Par2
2007 Byzantine Consensus with Few Synchronous Links
Moumen Hamouma, Achour Mostéfaoui, Gilles Trédan
OPODIS2
2007 Towards the minimal synchrony for byzantine consensus
abstract
No abstract available.
Achour Mostéfaoui, Gilles Trédan
PODC1
2007 From Renaming to Set Agreement
Achour Mostéfaoui, Michel Raynal, Corentin Travers
SIROCCO1
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.1
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. Computers2
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.2
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)1
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
PODC1
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
PRDC1
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
DISC1
2006 Synchronous condition-based consensus
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
Distributed Comput.1
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.1
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
PODC1
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
PRDC2
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
SRDS1
2005 Asynchronous bounded lifetime failure detectors
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.2
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.2
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
PODC2
2004 Brief announcement: the synchronous condition-based consensus hierarchy
abstract
No abstract available.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
PODC1
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
PRDC1
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
SRDS2
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
SRDS1
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
DISC2
2004 The Synchronous Condition-Based Consensus Hierarchy
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC1
2004 Condition-based consensus solvability: a hierarchy of conditions and efficient protocols
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
Distributed Comput.1
2004 A weakest failure detector-based asynchronous consensus protocol for f<n
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.2
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.3
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
DSN1
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
DSN1
2003 Single-Write Safe Consensus using Constrained Inputs
Matthieu Roy, Achour Mostéfaoui
SIROCCO2
2003 Using Conditions to Expedite Consensus in Synchronous Distributed Systems
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC1
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. ACM1
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
DSN1
2002 Asynchronous interactive consistency and its relation with error-correcting codes
abstract
No abstract available.
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
PODC1
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
PODC1
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
DISC3
2002 Distributed Agreement and Its Relation with Error-Correcting Codes
Roy Friedman 0001, Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC2
2002 Condition-Based Protocols for Set Agreement Problems
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
DISC1
2002 An introduction to oracles for asynchronous distributed systems
Achour Mostéfaoui, Eric Mourgaya, Michel Raynal
Future Gener. Comput. Syst.1
2002 Interval Consistency of Asynchronous Distributed Computations
Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
J. Comput. Syst. Sci.2
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. Computers2
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
IPDPS1
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
ISORC2
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
PODC1
2001 Efficient Condition-Based Consensus
Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal, Matthieu Roy
SIROCCO1
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
SPAA1
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
SRDS2
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
STOC1
2001 The logically instantaneous communication mode: a communication abstraction
Achour Mostéfaoui, Michel Raynal, Paulo Veríssimo
Future Gener. Comput. Syst.1
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.3
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.3
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
DSN1
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
ICDCS3
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
IPDPS1
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
PODC1
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
PRDC1
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.2
2000 From Binary Consensus to Multivalued Consensus in asynchronous message-passing systems
Achour Mostéfaoui, Michel Raynal, Frédéric Tronel
Inf. Process. Lett.1
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.3
1999 Unreliable Failure Detectors with Limited Scope Accuracy and an Application to Consensus
Achour Mostéfaoui, Michel Raynal
FSTTCS1
1999 Solving Consensus Using Chandra-Toueg's Unreliable Failure Detectors: A General Quorum-Based Approach
Achour Mostéfaoui, Michel Raynal
DISC1
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.2
1998 Shrinking Timestamp Sizes of Event Ordering Protocols
abstract
Almost all published work on causal ordering mechanisms assumes theoretically unbounded counters for timestamps, thus ignoring the real world problem that arises if one is actually interested in an operable implementation, since unbounded counters simply cannot be realized. An argument for its justification often encountered states, that the counter size can be chosen such that counters practically do not overflow or wrap around. For example, using matrix timestamps in a distributed computation involving not more than 50 processes and 32 bits per integer, results in a timestamp size of almost 10 K byte. We present a solution, called Factorized Timestamp Approach (FTA) that substantially reduces the amount of piggybacked control information. It is based on introducing the notion of phases in which much smaller timestamps are used. Simulation results given in the paper show the suitability of this approach.
Achour Mostéfaoui, Oliver E. Theel
ICPADS1
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
SRDS3
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
SRDS2
1997 Cycle Prevention in Distributed Checkpointing
Jean-Michel Hélary, Achour Mostéfaoui, Michel Raynal
OPODIS2
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
SRDS2
1996 Causal Delivery of Messages with Real-Time Data in Unreliable Networks
Roberto Baldoni, Achour Mostéfaoui, Michel Raynal
Real Time Syst.2
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
HPDC2
1994 A O(log2 n) Fault-Tolerant Distributed Mutual Exclusion Algorithm Based on Open-Cube Structure
abstract
A new distributed mutual exclusion algorithm, using a token and based upon an original rooted tree structure, is presented. The rooted tree introduced, named "open-cube", has noteworthy stability and locality properties, allowing the proposed algorithm to achieve good performances and high tolerance to node failures: the worst case message complexity per request is, in the absence of node failures, log/sub 2/n+1 where n is the number of nodes, whereas O(log/sub 2/n) extra messages in the average are necessary to tolerate each node failure. This algorithm is a particular instance of a general scheme for token and tree-based distributed mutual exclusion algorithms, previously presented in part by the authors; consequently, its safety and liveness properties are inherited from the general one.>
Jean-Michel Hélary, Achour Mostéfaoui
ICDCS2
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.2