VLDB 2026 Research / reviewers in the wild / expert
Rachid Guerraoui
dblp:g/RachidGuerraoui
· DBLP profile ↗
369ranked-venue papers
112as first author
63since 2021 · last 2026
0000-0002-4794-8902ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 156 · 46 first-author · 23 since 2021Software engineering, systems software and programming languages · 52 · 16 first-author · 6 since 2021Security and privacy · 41 · 11 first-author · 5 since 2021Theory of computation · 34 · 12 first-author · 2 since 2021Artificial intelligence and machine learning · 25 · 20 since 2021Databases, data management, data science and information retrieval · 17 · 9 first-author · 2 since 2021Computer networks · 11 · 5 first-authorApplied, interdisciplinary, general and emerging computing · 8 · 3 first-authorGraphics, computer vision, multimedia, augmented reality and games · 1 · 1 since 2021Human-computer interaction and ubiquitous computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Efficient Federated Search for Retrieval-Augmented Generation Using Lightweight Routing
Akash Balasaheb Dhasade, Rachid Guerraoui, Anne-Marie Kermarrec, Diana Petrescu, Rafael Pires 0001, Mathis Randl, Martijn de Vos |
DAIS | 2 |
| 2026 | Scalable Accountable Byzantine Agreement and Beyond
Pierre Civit, Daniel Collins 0001, Vincent Gramoli, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira, Pouriya Zarbafian |
SP | 4 |
| 2025 | Adaptive Gradient Clipping for Robust Federated LearningabstractRobust federated learning aims to maintain reliable performance despite the presence of adversarial or misbehaving workers. While state-of-the-art (SOTA) robust distributed gradient descent (Robust-DGD) methods were proven theoretically optimal, their empirical success has often relied on pre-aggregation gradient clipping.
However, existing static clipping strategies yield inconsistent results: enhancing robustness against some attacks while being ineffective or even detrimental against others.
To address this limitation, we propose a principled adaptive clipping strategy, Adaptive Robust Clipping (ARC), which dynamically adjusts clipping thresholds based on the input gradients. We prove that ARC not only preserves the theoretical robustness guarantees of SOTA Robust-DGD methods but also provably improves asymptotic convergence when the model is well-initialized. Extensive experiments on benchmark image classification tasks confirm these theoretical insights, demonstrating that ARC significantly enhances robustness, particularly in highly heterogeneous and adversarial settings. Youssef Allouah, Rachid Guerraoui, Nirupam Gupta, Ahmed Jellouli, Geovani Rizk, John Stephan |
ICLR | 2 |
| 2025 | The Utility and Complexity of In- and Out-of-Distribution Machine UnlearningabstractMachine unlearning, the process of selectively removing data from trained models, is increasingly crucial for addressing privacy concerns and knowledge gaps post-deployment. Despite this importance, existing approaches are often heuristic and lack formal guarantees. In this paper, we analyze the fundamental utility, time, and space complexity trade-offs of approximate unlearning, providing rigorous certification analogous to differential privacy. For in-distribution forget data—data similar to the retain set—we show that a surprisingly simple and general procedure, empirical risk minimization with output perturbation, achieves tight unlearning-utility-complexity trade-offs, addressing a previous theoretical gap on the separation from unlearning ``for free" via differential privacy, which inherently facilitates the removal of such data. However, such techniques fail with out-of-distribution forget data—data significantly different from the retain set—where unlearning time complexity can exceed that of retraining, even for a single sample. To address this, we propose a new robust and noisy gradient descent variant that provably amortizes unlearning time complexity without compromising utility. Youssef Allouah, Joshua Kazdan, Rachid Guerraoui, Oluwasanmi Koyejo |
ICLR | 3 |
| 2025 | Towards Trustworthy Federated Learning with Untrusted ParticipantsabstractResilience against malicious participants and data privacy are essential for trustworthy federated learning, yet achieving both with good utility typically requires the strong assumption of a trusted central server. This paper shows that a significantly weaker assumption suffices: each pair of participants shares a randomness seed unknown to others.
In a setting where malicious participants may collude with an untrusted server, we propose CafCor, an algorithm that integrates robust gradient aggregation with correlated noise injection, using shared randomness between participants.
We prove that CafCor achieves strong privacy-utility trade-offs, significantly outperforming local differential privacy (DP) methods, which do not make any trust assumption, while approaching central DP utility, where the server is fully trusted.
Empirical results on standard benchmarks validate CafCor's practicality, showing that privacy and robustness can coexist in distributed systems without sacrificing utility or trusting the server. Youssef Allouah, Rachid Guerraoui, John Stephan |
ICML | 2 |
| 2025 | Certified Unlearning for Neural NetworksabstractWe address the problem of machine unlearning, where the goal is to remove the influence of specific training data from a model upon request, motivated by privacy concerns and regulatory requirements such as the “right to be forgotten.” Unfortunately, existing methods rely on restrictive assumptions or lack formal guarantees.
To this end, we propose a novel method for certified machine unlearning, leveraging the connection between unlearning and privacy amplification by stochastic post-processing. Our method uses noisy fine-tuning on the retain data, i.e., data that does not need to be removed, to ensure provable unlearning guarantees. This approach requires no assumptions about the underlying loss function, making it broadly applicable across diverse settings. We analyze the theoretical trade-offs in efficiency and accuracy and demonstrate empirically that our method not only achieves formal unlearning guarantees but also performs effectively in practice, outperforming existing baselines. Anastasia Koloskova, Youssef Allouah, Animesh Jha, Rachid Guerraoui, Oluwasanmi Koyejo |
ICML | 4 |
| 2025 | Stabl: The Sensitivity of Blockchains to FailuresabstractBlockchains promise to make online services more fault tolerant because they are replicated on a distributed system of nodes. Their nodes typically run different implementations of the same protocol across different geo-distributed regions, making the protocol supposedly tolerant to various failures including isolated crashes, transient failures, network partitions or attacks. Unfortunately, their fault tolerance has never been compared. Vincent Gramoli, Rachid Guerraoui, Andrei Lebedev, Gauthier Voron |
Middleware | 2 |
| 2025 | Repeated Agreement is Cheap! On Weak Accountability and Multishot Byzantine AgreementabstractByzantine Agreement (BA) allows n processes to propose input values to reach consensus on a common, valid Lo-bit value, even in the presence of up to t < n faulty processes that can deviate arbitrarily from the protocol. Although strategies like randomization, adaptiveness, and batching have been extensively explored to mitigate the inherent limitations of one-shot agreement tasks, there has been limited progress on achieving good amortized performance for multi-shot agreement, despite its obvious relevance to long-lived functionalities such as state machine replication. Pierre Civit, Muhammad Ayaz Dzulfikar, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira |
PODC | 4 |
| 2025 | Partial Synchrony for Free: New Upper Bounds for Byzantine AgreementabstractByzantine agreement allows n processes to decide on a common value, in spite of arbitrary failures. The seminal Dolev-Reischuk bound states that any deterministic solution to Byzantine agreement exchanges Ω(n2) bits. In synchronous networks, with a known upper bound on message delays, solutions with optimal O (n2) bit complexity, optimal fault tolerance, and no cryptography have been established for over three decades. However, these solutions lack robustness under adverse network conditions. Therefore, research has increasingly focused on Byzantine agreement for partially synchronous networks, which behave synchronously only eventually and are thus more reflective of real-world conditions. Numerous solutions have been proposed for the partially synchronous setting. However, these solutions are notoriously hard to prove correct, and the most efficient cryptography-free algorithms still require O (n3) exchanged bits in the worst case. Even with cryptography, the state-of-the-art remains a κ-bit factor away from the Ω(n2) lower bound (where κ is the security parameter). This discrepancy between synchronous and partially synchronous solutions has remained unresolved for decades. Pierre Civit, Muhammad Ayaz Dzulfikar, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira, Igor Zablotchi |
SODA | 4 |
| 2025 | ARGO: Overcoming hardware dependence in distributed learning
Karim Boubouh, Amine Boussetta, Rachid Guerraoui, Alexandre Maurer |
Future Gener. Comput. Syst. | 3 |
| 2025 | Carbon: Scaling Trusted Payments With Untrusted machinesabstractThis paper introduces Carbon, a high-throughput system enabling asynchronous (safe) and consensus-free (efficient) payments and votes within a dynamic set of clients. Carbon is operated by a dynamic set of validators that may be reconfigured asynchronously, offering its clients eclipse resistance as well as lightweight bootstrap. Carbon offers clients the ability to select validators by voting them in and out of the system thanks to its novel asynchronous and stake-less voting mechanism. Carbon relies on an asynchronous and deterministic implementation of Byzantine reliable broadcast that uniquely leverages a permissionless set of untrusted servers, brokers, to slash the cost of client authentication inherent to Byzantine fault tolerant systems. Carbon is able to sustain a throughput of one million payments per second in a geo-distributed environment, outperforming the state of the art by three orders of magnitude with equivalent latencies. Martina Camaioni, Rachid Guerraoui, Jovan Komatovic, Matteo Monti, Pierre-Louis Roman, Manuel Vidigueira, Gauthier Voron |
IEEE Trans. Dependable Secur. Comput. | 2 |
| 2024 | Robust Sparse VotingabstractMany applications, such as content moderation and recommendation, require reviewing and scoring a large number of alternatives. Doing so robustly is however very challenging. Indeed, voters’ inputs are inevitably sparse: most alternatives are only scored by a small fraction of voters. This sparsity amplifies the effects of biased voters introducing unfairness, and of malicious voters seeking to hack the voting process by reporting dishonest scores. We give a precise definition of the problem of robust sparse voting, highlight its underlying technical challenges, and present a novel voting mechanism addressing the problem. We prove that, using this mechanism, no voter can have more than a small parameterizable effect on each alternative’s score; a property we call Lipschitz resilience. We also identify conditions of voters comparability under which any unanimous preferences can be recovered, even when each voter provides sparse scores, on a scale that is potentially very different from any other voter’s score scale. Proving these properties required us to introduce, analyze and carefully compose novel aggregation primitives which could be of independent interest. Youssef Allouah, Rachid Guerraoui, Lê-Nguyên Hoang, Oscar Villemaud |
AISTATS | 2 |
| 2024 | Accelerating Transfer Learning with Near-Data Computation on Cloud Object StoresabstractStorage disaggregation underlies today's cloud and is naturally complemented by pushing down some computation to storage, thus mitigating the potential network bottleneck between the storage and compute tiers. We show how ML training benefits from storage pushdowns by focusing on transfer learning (TL), the widespread technique that democratizes ML by reusing existing knowledge on related tasks. We propose HAPI, a new TL processing system centered around two complementary techniques that address challenges introduced by disaggregation. First, applications must carefully balance execution across tiers for performance. HAPI judiciously splits the TL computation during the feature extraction phase yielding pushdowns that not only improve network time but also improve total TL training time by overlapping the execution of consecutive training iterations across tiers. Second, operators want resource efficiency from the storage-side computational resources. HAPI employs storage-side batch size adaptation allowing increased storage-side pushdown concurrency without affecting training accuracy. HAPI yields up to 2.5× training speed-up while choosing in 86.8% of cases the best performing split point or one that is at most 5% off from the best. Diana Petrescu, Arsany Guirguis, Do Le Quoc, Javier Picorel, Rachid Guerraoui, Florin Dinu |
SoCC | 5 |
| 2024 | Byzantine-Robust Federated Learning: Impact of Client Subsampling and Local UpdatesabstractThe possibility of adversarial (a.k.a., Byzantine) clients makes federated learning (FL) prone to arbitrary manipulation. The natural approach to robustify FL against adversarial clients is to replace the simple averaging operation at the server in the standard $\mathsf{FedAvg}$ algorithm by a robust averaging rule. While a significant amount of work has been devoted to studying the convergence of federated robust averaging (which we denote by $\mathsf{FedRo}$), prior work has largely ignored the impact of client subsampling and local steps, two fundamental FL characteristics. While client subsampling increases the effective fraction of Byzantine clients, local steps increase the drift between the local updates computed by honest (i.e., non-Byzantine) clients. Consequently, a careless deployment of $\mathsf{FedRo}$ could yield poor performance. We validate this observation by presenting an in-depth analysis of $\mathsf{FedRo}$ tightly analyzing the impact of client subsampling and local steps. Specifically, we present a sufficient condition on client subsampling for nearly-optimal convergence of $\mathsf{FedRo}$ (for smooth non-convex loss). Also, we show that the rate of improvement in learning accuracy diminishes with respect to the number of clients subsampled, as soon as the sample size exceeds a threshold value. Interestingly, we also observe that under a careful choice of step-sizes, the learning error due to Byzantine clients decreases with the number of local steps. We validate our theory by experiments on the FEMNIST and CIFAR-$10$ image classification tasks. Youssef Allouah, Sadegh Farhadkhani, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot, Geovani Rizk, Sasha Voitovych |
ICML | 3 |
| 2024 | The Privacy Power of Correlated Noise in Decentralized LearningabstractDecentralized learning is appealing as it enables the scalable usage of large amounts of distributed data and resources without resorting to any central entity, while promoting privacy since every user minimizes the direct exposure of their data. Yet, without additional precautions, curious users can still leverage models obtained from their peers to violate privacy. In this paper, we propose Decor, a variant of decentralized SGD with differential privacy (DP) guarantees. Essentially, in Decor, users securely exchange randomness seeds in one communication round to generate pairwise-canceling correlated Gaussian noises, which are injected to protect local models at every communication round. We theoretically and empirically show that, for arbitrary connected graphs, Decor matches the central DP optimal privacy-utility trade-off. We do so under SecLDP, our new relaxation of local DP, which protects all user communications against an external eavesdropper and curious users, assuming that every pair of connected users shares a secret, i.e., an information hidden to all others. The main theoretical challenge is to control the accumulation of non-canceling correlated noise due to network sparsity. We also propose a companion SecLDP privacy accountant for public use. Youssef Allouah, Anastasia Koloskova, Aymane El Firdoussi, Martin Jaggi, Rachid Guerraoui |
ICML | 5 |
| 2024 | Towards Practical Homomorphic Aggregation in Byzantine-Resilient Distributed LearningabstractThe growing availability of distributed data has led to the increased use of machine learning (ML) algorithms in distributed topologies, where multiple nodes collaborate to train models under the coordination of a central server. However, distributed learning faces two significant challenges: the risk of Byzantine nodes corrupting the learning process by sending incorrect information, and the potential for a curious server to violate the privacy of individual nodes, even reconstructing their private data. While homomorphic encryption (HE) has been a promising solution for privacy preservation in distributed settings, its high computational cost, especially for high-dimensional ML models, has made it challenging to design robust (non-linear) Byzantine-resilient algorithms using HE. Antoine Choffrut, Rachid Guerraoui, Rafael Pinot, Renaud Sirdey, John Stephan, Martin Zuber |
Middleware | 2 |
| 2024 | Revisiting Ensembling in One-Shot Federated LearningabstractFederated Learning (FL) is an appealing approach to training machine learning models without sharing raw data. However, standard FL algorithms are iterative and thus induce a significant communication cost. One-Shot FL (OFL) trades the iterative exchange of models between clients and the server with a single round of communication, thereby saving substantially on communication costs. Not surprisingly, OFL exhibits a performance gap in terms of accuracy with respect to FL, especially under high data heterogeneity. We introduce Fens, a novel federated ensembling scheme that approaches the accuracy of FL with the communication efficiency of OFL. Learning in Fens proceeds in two phases: first, clients train models locally and send them to the server, similar to OFL; second, clients collaboratively train a lightweight prediction aggregator model using FL. We showcase the effectiveness of Fens through exhaustive experiments spanning several datasets and heterogeneity levels. In the particular case of heterogeneously distributed CIFAR-10 dataset, Fens achieves up to a $26.9$% higher accuracy over SOTA OFL, being only $3.1$% lower than FL. At the same time, Fens incurs at most $4.3\times$ more communication than OFL, whereas FL is at least $10.9\times$ more communication-intensive than Fens. Youssef Allouah, Akash Balasaheb Dhasade, Rachid Guerraoui, Nirupam Gupta, Anne-Marie Kermarrec, Rafael Pinot, Rafael Pires 0001, Rishi Sharma 0001 |
NeurIPS | 3 |
| 2024 | Fine-Tuning Personalization in Federated Learning to Mitigate Adversarial ClientsabstractFederated learning (FL) is an appealing paradigm that allows a group of machines
(a.k.a. clients) to learn collectively while keeping their data local. However, due
to the heterogeneity between the clients’ data distributions, the model obtained
through the use of FL algorithms may perform poorly on some client’s data.
Personalization addresses this issue by enabling each client to have a different
model tailored to their own data while simultaneously benefiting from the other
clients’ data. We consider an FL setting where some clients can be adversarial, and
we derive conditions under which full collaboration fails. Specifically, we analyze
the generalization performance of an interpolated personalized FL framework in the
presence of adversarial clients, and we precisely characterize situations when full
collaboration performs strictly worse than fine-tuned personalization. Our analysis
determines how much we should scale down the level of collaboration, according
to data heterogeneity and the tolerable fraction of adversarial clients. We support
our findings with empirical results on mean estimation and binary classification
problems, considering synthetic and benchmark image classification datasets Youssef Allouah, Abdellah El Mrini, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot |
NeurIPS | 3 |
| 2024 | DSig: Breaking the Barrier of Signatures in Data Centers
Marcos K. Aguilera, Clément Burgelin, Rachid Guerraoui, Antoine Murat, Athanasios Xygkis, Igor Zablotchi |
OSDI | 3 |
| 2024 | Chop Chop: Byzantine Atomic Broadcast to the Network Limit
Martina Camaioni, Rachid Guerraoui, Matteo Monti, Pierre-Louis Roman, Manuel Vidigueira, Gauthier Voron |
OSDI | 2 |
| 2024 | DARE to Agree: Byzantine Agreement With Optimal Resilience and Adaptive CommunicationabstractByzantine Agreement (BA) enables n processes to reach consensus on a common valid Lo-bit value, even in the presence of up to t < n faulty processes that can deviate arbitrarily from their prescribed protocol. Despite its significance, the optimal communication complexity for key variations of BA has not been determined within the honest majority regime (n = 2t +1), for both the worst-case scenario and the adaptive scenario, which accounts for the actual number f ≤ t of failures. We introduce ada-Dare (Adaptively Disperse, Agree, Retrieve), a novel universal approach to solve BA efficiently. Let κ represent the size of the cryptographic objects required to solve BA when t > n/3. Different instantiations of ada-Dare achieve near-optimal adaptive bit complexity of O(nLo + n(f + 1)κ) for both strong multi-valued validated BA (SMVBA) and interactive consistency (IC). By definition, for IC, Lo = nLin where Lin is the size of an input value. These results achieve optimal O(n(Lo + f)) word complexity and significantly improve the previous best results by up to a linear factor, depending on Lo and f. Pierre Civit, Muhammad Ayaz Dzulfikar, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira |
PODC | 4 |
| 2024 | All Byzantine Agreement Problems Are ExpensiveabstractByzantine agreement, arguably the most fundamental problem in distributed computing, operates among n processes, out of which t < n can exhibit arbitrary failures. The problem states that all correct (non-faulty) processes must eventually decide (termination) the same value (agreement) from a set of admissible values defined by the proposals of the processes (validity). Depending on the exact version of the validity property, Byzantine agreement comes in different forms, from Byzantine broadcast to strong and weak consensus, to modern variants of the problem introduced in today's blockchain systems. Regardless of the specific flavor of the agreement problem, its communication cost is a fundamental metric whose improvement has been the focus of decades of research. The Dolev-Reischuk bound, one of the most celebrated results in distributed computing, proved 40 years ago that, at least for Byzantine broadcast, no deterministic solution can do better than Ω(t2) exchanged messages in the worst case. Since then, it remained unknown whether the quadratic lower bound extends to all variants of Byzantine agreement. This paper answers the question in the affirmative, closing this long-standing open problem. Namely, we prove that any non-trivial agreement problem requires Ω(t2) messages to be exchanged in the worst case. To prove the general lower bound, we determine the weakest Byzantine agreement problem and show, via a novel indistinguishability argument, that it incurs Ω(t2) exchanged messages even in the general omission model. Pierre Civit, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Anton Paramonov, Manuel Vidigueira |
PODC | 3 |
| 2024 | Brief Announcement: A Case for Byzantine Machine LearningabstractThe success of machine learning (ML) has been intimately linked with the availability of large amounts of data, typically collected from heterogeneous sources and processed on vast networks of computing devices (also called workers). Beyond accuracy, the use of ML in critical domains such as healthcare and autonomous driving calls for robustness against data poisoning and faulty workers. The problem of Byzantine ML formalizes these robustness issues by considering a distributed ML environment in which workers (storing a portion of the global dataset) can deviate arbitrarily from the prescribed algorithm. Although the problem has attracted a lot of attention from a theoretical point of view, its practical importance for addressing realistic faults (where the behavior of any worker is locally constrained) remains unclear. It has been argued that the seemingly weaker threat model where only workers' local datasets get poisoned is more reasonable. We highlight here some important results on the efficacy of Byzantine robustness for tackling data poisoning. In particular, we discuss cases where, while tolerating a wider range of faulty behaviors, Byzantine ML yields solutions that are optimal even under the weaker threat model of data poisoning. Sadegh Farhadkhani, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot |
PODC | 2 |
| 2024 | SWARM: Replicating Shared Disaggregated-Memory Data in No TimeabstractMemory disaggregation is an emerging data center architecture that improves resource utilization and scalability. Replication is key to ensure the fault tolerance of applications, but replicating shared data in disaggregated memory is hard. We propose SWARM (Swift WAit-free Replication in disaggregated Memory), the first replication scheme for in-disaggregated-memory shared objects to provide (1) single-roundtrip reads and writes in the common case, (2) strong consistency (linearizability), and (3) strong liveness (wait-freedom). SWARM makes two independent contributions. The first is Safe-Guess, a novel wait-free replication protocol with single-roundtrip operations. The second is In-n-Out, a novel technique to provide conditional atomic update and atomic retrieval of large buffers in disaggregated memory in one roundtrip. Using SWARM, we build SWARM-KV, a low-latency, strongly consistent and highly available disaggregated key-value store. We evaluate SWARM-KV and find that it has marginal latency overhead compared to an unreplicated key-value store, and that it offers much lower latency and better availability than FUSEE, a state-of-the-art replicated disaggregated key-value store. Antoine Murat, Clément Burgelin, Athanasios Xygkis, Igor Zablotchi, Marcos K. Aguilera, Rachid Guerraoui |
SOSP | 6 |
| 2024 | Peerswap: A Peer-Sampler with Randomness GuaranteesabstractThe ability of a peer-to-peer (P2P) system to effectively host decentralized applications often relies on the availability of a peer-sampling service, which provides each participant with a random sample of other peers. Despite the practical effectiveness of existing peer samplers, their ability to produce random samples within a reasonable time frame remains poorly understood from a theoretical standpoint. This paper contributes to bridging this gap by introducing PeersWap,a peer-sampling protocol with provable randomness guarantees. We establish execution time bounds for PeerSwap, demonstrating its ability to scale effectively with the network size. We prove that PeerSwap maintains the fixed structure of the communication graph while allowing sequential peer position swaps within this graph. We do so by showing that PeerSwap is a specific instance of an interchange process, a renowned model for particle movement analysis. Leveraging this mapping, we derive execution time bounds, expressed as a function of the network size n. Depending on the network structure, this time can be as low as a polylogarithmic function of n, highlighting the efficiency of PeerSwap. We implement PeerSwap and conduct numerical evaluations using regular graphs with varying connectivity and containing up to 32768(215) peers. Our evaluation demonstrates that PeerSwap quickly provides peers with uniform random samples of other peers. Rachid Guerraoui, Anne-Marie Kermarrec, Anastasiia Kucherenko, Rafael Pinot, Martijn de Vos |
SRDS | 1 |
| 2024 | Efficient Signature-Free Validated Agreement
Pierre Civit, Muhammad Ayaz Dzulfikar, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira, Igor Zablotchi |
DISC | 4 |
| 2024 | Byzantine consensus is Θ (n2): the Dolev-Reischuk bound is tight even in partial synchrony!abstractAbstract The Dolev-Reischuk bound says that any deterministic Byzantine consensus protocol has (at least) quadratic (in the number of processes) communication complexity in the worst case: given a system with n processes and at most $$f < n / 3$$ f < n / 3 failures, any solution to Byzantine consensus exchanges $$\Omega \big (n^2\big )$$ Ω ( n 2 ) words, where a word contains a constant number of values and signatures. While it has been shown that the bound is tight in synchronous environments, it is still unknown whether a consensus protocol with quadratic communication complexity can be obtained in partial synchrony where the network alternates between (1) asynchronous periods, with unbounded message delays, and (2) synchronous periods, with $$\delta $$ δ -bounded message delays. Until now, the most efficient known solutions for Byzantine consensus in partially synchronous settings had cubic communication complexity (e.g., HotStuff, binary DBFT). This paper closes the existing gap by introducing SQuad , a partially synchronous Byzantine consensus protocol with $$O\big (n^2\big )$$ O ( n 2 ) worst-case communication complexity. In addition, SQuad is optimally-resilient (tolerating up to $$f < n / 3$$ f < n / 3 failures) and achieves $$O(f \cdot \delta )$$ O ( f · δ ) worst-case latency complexity. The key technical contribution underlying SQuad lies in the way we solve view synchronization , the problem of bringing all correct processes to the same view with a correct leader for sufficiently long. Concretely, we present RareSync , a view synchronization protocol with $$O\big (n^2\big )$$ O ( n 2 ) communication complexity and $$O(f \cdot \delta )$$ O ( f · δ ) latency complexity, which we utilize in order to obtain SQuad . Pierre Civit, Muhammad Ayaz Dzulfikar, Seth Gilbert, Vincent Gramoli, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira |
Distributed Comput. | 5 |
| 2023 | Fixing by Mixing: A Recipe for Optimal Byzantine ML under HeterogeneityabstractByzantine machine learning (ML) aims to ensure the resilience of distributed learning algorithms to misbehaving (or Byzantine) machines. Although this problem received significant attention, prior works often assume the data held by the machines to be homogeneous, which is seldom true in practical settings. Data heterogeneity makes Byzantine ML considerably more challenging, since a Byzantine machine can hardly be distinguished from a non-Byzantine outlier. A few solutions have been proposed to tackle this issue, but these provide suboptimal probabilistic guarantees and fare poorly in practice. This paper closes the theoretical gap, achieving optimality and inducing good empirical results. In fact, we show how to automatically adapt existing solutions for (homogeneous) Byzantine ML to the heterogeneous setting through a powerful mechanism, we call nearest neighbor mixing (NNM), which boosts any standard robust distributed gradient descent variant to yield optimal Byzantine resilience under heterogeneity. We obtain similar guarantees (in expectation) by plugging NNM in the distributed stochastic heavy ball method, a practical substitute to distributed gradient descent. We obtain empirical results that significantly outperform state-of-the-art Byzantine ML solutions. Youssef Allouah, Sadegh Farhadkhani, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot, John Stephan |
AISTATS | 3 |
| 2023 | On the Strategyproofness of the Geometric MedianabstractThe geometric median, an instrumental component of the secure machine learning toolbox, is known to be effective when robustly aggregating models (or gradients), gathered from potentially malicious (or strategic) users. What is less known is the extent to which the geometric median incentivizes dishonest behaviors. This paper addresses this fundamental question by quantifying its strategyproofness. While we observe that the geometric median is not even approximately strategyproof, we prove that it is asymptotically $\alpha$-strategyproof: when the number of users is large enough, a user that misbehaves can gain at most a multiplicative factor $\alpha$, which we compute as a function of the distribution followed by the users. We then generalize our results to the case where users actually care more about specific dimensions, determining how this impacts $\alpha$. We also show how the skewed geometric medians can be used to improve strategyproofness. El Mahdi El Mhamdi, Sadegh Farhadkhani, Rachid Guerraoui, Lê-Nguyên Hoang |
AISTATS | 3 |
| 2023 | uBFT: Microsecond-Scale BFT using Disaggregated MemoryabstractWe propose uBFT, the first State Machine Replication (SMR) system to achieve microsecond-scale latency in data centers, while using only 2f+1 replicas to tolerate f Byzantine failures. The Byzantine Fault Tolerance (BFT) provided by uBFT is essential as pure crashes appear to be a mere illusion with real-life systems reportedly failing in many unexpected ways. uBFT relies on a small non-tailored trusted computing base—disaggregated memory—and consumes a practically bounded amount of memory. uBFT is based on a novel abstraction called Consistent Tail Broadcast, which we use to prevent equivocation while bounding memory. We implement uBFT using RDMA-based disaggregated memory and obtain an end-to-end latency of as little as 10 us. This is at least 50× faster than MinBFT, a state-of-the-art 2f+1 BFT SMR based on Intel’s SGX. We use uBFT to replicate two KV-stores (Memcached and Redis), as well as a financial order matching engine (Liquibook). These applications have low latency (up to 20 us) and become Byzantine tolerant with as little as 10 us more. The price for uBFT is a small amount of reliable disaggregated memory (less than 1 MiB), which in our prototype consists of a small number of memory servers connected through RDMA and replicated for fault tolerance. Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Antoine Murat, Athanasios Xygkis, Igor Zablotchi |
ASPLOS (2) | 3 |
| 2023 | Diablo: A Benchmark Suite for BlockchainsabstractWith the recent advent of blockchains, we have witnessed a plethora of blockchain proposals. These proposals range from using work to using time, storage or stake in order to select blocks to be appended to the chain. As a drawback it makes it difficult for the application developer to choose the right blockchain to support their applications. In particular, the scalability and performance one can obtain from a specific blockchain is typically unknown. The claimed results are often obtained in isolation by the developers of the blockchain themselves. The experimental conditions corresponding to these results are generally missing and the lack of details make these results irreproducible. Vincent Gramoli, Rachid Guerraoui, Andrei Lebedev, Christopher Natoli, Gauthier Voron |
EuroSys | 2 |
| 2023 | On the Privacy-Robustness-Utility Trilemma in Distributed LearningabstractThe ubiquity of distributed machine learning (ML) in sensitive public domain applications calls for algorithms that protect data privacy, while being robust to faults and adversarial behaviors. Although privacy and robustness have been extensively studied independently in distributed ML, their synthesis remains poorly understood. We present the first tight analysis of the error incurred by any algorithm ensuring robustness against a fraction of adversarial machines, as well as differential privacy (DP) for honest machines' data against any other curious entity. Our analysis exhibits a fundamental trade-off between privacy, robustness, and utility. To prove our lower bound, we consider the case of mean estimation, subject to distributed DP and robustness constraints, and devise reductions to centralized estimation of one-way marginals. We prove our matching upper bound by presenting a new distributed ML algorithm using a high-dimensional robust aggregation rule. The latter amortizes the dependence on the dimension in the error (caused by adversarial workers and DP), while being agnostic to the statistical properties of the data. Youssef Allouah, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot, John Stephan |
ICML | 2 |
| 2023 | Robust Collaborative Learning with Linear Gradient OverheadabstractCollaborative learning algorithms, such as distributed SGD (or D-SGD), are prone to faulty machines that may deviate from their prescribed algorithm because of software or hardware bugs, poisoned data or malicious behaviors. While many solutions have been proposed to enhance the robustness of D-SGD to such machines, previous works either resort to strong assumptions (trusted server, homogeneous data, specific noise model) or impose a gradient computational cost that is several orders of magnitude higher than that of D-SGD. We present MoNNA, a new algorithm that (a) is provably robust under standard assumptions and (b) has a gradient computation overhead that is linear in the fraction of faulty machines, which is conjectured to be tight. Essentially, MoNNA uses Polyak’s momentum of local gradients for local updates and nearest-neighbor averaging (NNA) for global mixing, respectively. While MoNNA is rather simple to implement, its analysis has been more challenging and relies on two key elements that may be of independent interest. Specifically, we introduce the mixing criterion of $(\alpha, \lambda)$-reduction to analyze the non-linear mixing of non-faulty machines, and present a way to control the tension between the momentum and the model drifts. We validate our theory by experiments on image classification and make our code available at https://github.com/LPD-EPFL/robust-collaborative-learning. Sadegh Farhadkhani, Rachid Guerraoui, Nirupam Gupta, Lê-Nguyên Hoang, Rafael Pinot, John Stephan |
ICML | 2 |
| 2023 | Robust Distributed Learning: Tight Error Bounds and Breakdown Point under Data HeterogeneityabstractThe theory underlying robust distributed learning algorithms, designed to resist adversarial machines, matches empirical observations when data is homogeneous. Under data heterogeneity however, which is the norm in practical scenarios, established lower bounds on the learning error are essentially vacuous and greatly mismatch empirical observations. This is because the heterogeneity model considered is too restrictive and does not cover basic learning tasks such as least-squares regression. We consider in this paper a more realistic heterogeneity model, namely $(G,B)$-gradient dissimilarity, and show that it covers a larger class of learning problems than existing theory. Notably, we show that the breakdown point under heterogeneity is lower than the classical fraction $\frac{1}{2}$. We also prove a new lower bound on the learning error of any distributed learning algorithm. We derive a matching upper bound for a robust variant of distributed gradient descent, and empirically show that our analysis reduces the gap between theory and practice. Youssef Allouah, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot, Geovani Rizk |
NeurIPS | 2 |
| 2023 | Epidemic Learning: Boosting Decentralized Learning with Randomized CommunicationabstractWe present Epidemic Learning (EL), a simple yet powerful decentralized learning (DL) algorithm that leverages changing communication topologies to achieve faster model convergence compared to conventional DL approaches. At each round of EL, each node sends its model updates to a random sample of $s$ other nodes (in a system of $n$ nodes). We provide an extensive theoretical analysis of EL, demonstrating that its changing topology culminates in superior convergence properties compared to the state-of-the-art (static and dynamic) topologies. Considering smooth non-convex loss functions, the number of transient iterations for EL, i.e., the rounds required to achieve asymptotic linear speedup, is in $O(n^3/s^2)$ which outperforms the best-known bound $O(n^3)$ by a factor of $s^2$, indicating the benefit of randomized communication for DL. We empirically evaluate EL in a 96-node network and compare its performance with state-of-the-art DL approaches. Our results illustrate that EL converges up to $ 1.7\times$ quicker than baseline DL algorithms and attains $2.2 $\% higher accuracy for the same communication volume. Martijn de Vos, Sadegh Farhadkhani, Rachid Guerraoui, Anne-Marie Kermarrec, Rafael Pires 0001, Rishi Sharma 0001 |
NeurIPS | 3 |
| 2023 | On the Validity of ConsensusabstractThe Byzantine consensus problem involves n processes, out of which t < n could be faulty and behave arbitrarily. Three properties characterize consensus: (1) termination, requiring correct (non-faulty) processes to eventually reach a decision, (2) agreement, preventing them from deciding different values, and (3) validity, precluding "unreasonable" decisions. But, what is a reasonable decision? Strong validity, a classical property, stipulates that, if all correct processes propose the same value, only that value can be decided. Weak validity, another established property, stipulates that, if all processes are correct and they propose the same value, that value must be decided. The space of possible validity properties is vast. Yet, their impact on consensus algorithms remains unclear. Pierre Civit, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira |
PODC | 3 |
| 2023 | Every Bit Counts in ConsensusabstractConsensus enables n processes to agree on a common valid L-bit value, despite t < n/3 processes being faulty and acting arbitrarily. A long line of work has been dedicated to improving the worst-case communication complexity of consensus in partial synchrony. This has recently culminated in the worst-case word complexity of O(n^2). However, the worst-case bit complexity of the best solution is still O(n^2 L + n^2 kappa) (where kappa is the security parameter), far from the Ω(n L + n^2) lower bound. The gap is significant given the practical use of consensus primitives, where values typically consist of batches of large size (L > n). This paper shows how to narrow the aforementioned gap while achieving optimal linear latency. Namely, we present a new algorithm, DARE (Disperse, Agree, REtrieve), that improves upon the O(n^2 L) term via a novel dispersal primitive. DARE achieves O(n^{1.5} L + n^{2.5} kappa) bit complexity, an effective sqrt{n}-factor improvement over the state-of-the-art (when L > n kappa). Moreover, we show that employing heavier cryptographic primitives, namely STARK proofs, allows us to devise DARE-Stark, a version of DARE which achieves the near-optimal bit complexity of O(n L + n^2 poly(kappa)). Both DARE and DARE-Stark achieve optimal O(n) latency. Pierre Civit, Seth Gilbert, Rachid Guerraoui, Jovan Komatovic, Matteo Monti, Manuel Vidigueira |
DISC | 3 |
| 2023 | On the Inherent Anonymity of GossipingabstractDetecting the source of a gossip is a critical issue, related to identifying patient zero in an epidemic, or the origin of a rumor in a social network. Although it is widely acknowledged that random and local gossip communications make source identification difficult, there exists no general quantification of the level of anonymity provided to the source. This paper presents a principled method based on $\varepsilon$-differential privacy to analyze the inherent source anonymity of gossiping for a large class of graphs. First, we quantify the fundamental limit of source anonymity any gossip protocol can guarantee in an arbitrary communication graph. In particular, our result indicates that when the graph has poor connectivity, no gossip protocol can guarantee any meaningful level of differential privacy. This prompted us to further analyze graphs with controlled connectivity. We prove on these graphs that a large class of gossip protocols, namely cobra walks, offers tangible differential privacy guarantees to the source. In doing so, we introduce an original proof technique based on the reduction of a gossip protocol to what we call a random walk with probabilistic die out. This proof technique is of independent interest to the gossip community and readily extends to other protocols inherited from the security community, such as the Dandelion protocol. Interestingly, our tight analysis precisely captures the trade-off between dissemination time of a gossip protocol and its source anonymity. Rachid Guerraoui, Anne-Marie Kermarrec, Anastasiia Kucherenko, Rafael Pinot, Sasha Voitovych |
DISC | 1 |
| 2023 | Leaderless consensus
Karolos Antoniadis, Julien Benhaim, Antoine Desjardins, Elias Poroma, Vincent Gramoli, Rachid Guerraoui, Gauthier Voron, Igor Zablotchi |
J. Parallel Distributed Comput. | 6 |
| 2023 | As easy as ABC: Optimal (A)ccountable (B)yzantine (C)onsensus is easy!abstractIn a non-synchronous system with n processes, no t0-resilient (deterministic or probabilistic) Byzantine consensus protocol can prevent a disagreement among correct processes if the number of faulty processes is ≥n−2t0. Therefore, the community defined the accountable Byzantine consensus problem: the problem of (i) solving Byzantine consensus whenever possible (e.g., when the number of faulty processes does not exceed t0), and (ii) allowing correct processes to obtain proofs of culpability of n−2t0 faulty processes whenever a disagreement occurs. This paper presents ABC, a simple yet efficient transformation of any non-synchronous t0-resilient (deterministic or probabilistic) Byzantine consensus protocol into its accountable counterpart. In the common case (up to t0 faults), ABC introduces an additive overhead of two communication rounds and O(n2) exchanged bits. Whenever they disagree, correct processes detect culprits by exchanging O(n3) messages, which we prove optimal. Lastly, ABC is not limited to Byzantine consensus: ABC provides accountability for other essential distributed problems (e.g., reliable and consistent broadcast). Pierre Civit, Seth Gilbert, Vincent Gramoli, Rachid Guerraoui, Jovan Komatovic |
J. Parallel Distributed Comput. | 4 |
| 2023 | GoldFinger: Fast & Approximate Jaccard for Efficient KNN Graph ConstructionsabstractWe proposeGoldFinger, a newcompactandfast-to-computebinary representation of datasets to approximate Jaccard's index. We illustrate the effectiveness of GoldFinger on the emblematic big data problem of K-Nearest-Neighbor (KNN) graph construction and show that GoldFinger can drastically accelerate a large range of existing KNN algorithms with little to no overhead. As a side effect, we also show that the compact representation of the data protects users’ privacyfor freeby providingk-anonymity andl-diversity. Our extensive evaluation of the resulting approach on several realistic datasets shows that our approach reduces computation times by up to 78.9% compared to raw data while only incurring a negligible to moderate loss in terms of KNN quality. We also show that GoldFinger can be applied to KNN queries (a widely-used search technique) and delivers speedups of up to$\times 3.55$over one of the most efficient approaches to this problem. Rachid Guerraoui, Anne-Marie Kermarrec, Guilhem Niot, Olivier Ruas, François Taïani |
IEEE Trans. Knowl. Data Eng. | 1 |
| 2022 | Crime and Punishment in Distributed Byzantine Decision TasksabstractA decision task is a distributed input-output problem in which each process starts with its input value and eventually produces its output value. Examples of such decision tasks are broad and range from consensus to reliable broadcast to lattice agreement. A distributed protocol solves a decision task if it enables processes to produce admissible output values despite arbitrary (Byzantine) failures. Unfortunately, it has been known for decades that many decision tasks cannot be solved if the system is overly corrupted, i.e., safety of distributed protocols solving such tasks can be violated in unlucky scenarios.By contrast, only recently did the community discover that some of these distributed protocols can be made accountable by ensuring that correct processes irrevocably detect some faulty processes responsible for any safety violation. This realization is particularly surprising (and positive) given that accountability is a powerful tool to mitigate safety violations in distributed protocols. Indeed, exposing crimes and introducing punishments naturally incentivize exemplarity.In this paper, we propose a generic transformation, called τscr, of any non-synchronous distributed protocol solving a decision task into its accountable version. Our τscrtransformation is built upon the well-studied simulation of crash failures on top of Byzantine failures and increases the communication complexity by a quadratic multiplicative factor in the worst case. Pierre Civit, Seth Gilbert, Vincent Gramoli, Rachid Guerraoui, Jovan Komatovic, Zarko Milosevic 0001, Adi Seredinschi |
ICDCS | 4 |
| 2022 | Byzantine Machine Learning Made Easy By Resilient Averaging of MomentumsabstractByzantine resilience emerged as a prominent topic within the distributed machine learning community. Essentially, the goal is to enhance distributed optimization algorithms, such as distributed SGD, in a way that guarantees convergence despite the presence of some misbehaving (a.k.a., Byzantine) workers. Although a myriad of techniques addressing the problem have been proposed, the field arguably rests on fragile foundations. These techniques are hard to prove correct and rely on assumptions that are (a) quite unrealistic, i.e., often violated in practice, and (b) heterogeneous, i.e., making it difficult to compare approaches. We present RESAM (RESilient Averaging of Momentums), a unified framework that makes it simple to establish optimal Byzantine resilience, relying only on standard machine learning assumptions. Our framework is mainly composed of two operators: resilient averaging at the server and distributed momentum at the workers. We prove a general theorem stating the convergence of distributed SGD under RESAM. Interestingly, demonstrating and comparing the convergence of many existing techniques become direct corollaries of our theorem, without resorting to stringent assumptions. We also present an empirical evaluation of the practical relevance of RESAM. Sadegh Farhadkhani, Rachid Guerraoui, Nirupam Gupta, Rafael Pinot, John Stephan |
ICML | 2 |
| 2022 | An Equivalence Between Data Poisoning and Byzantine Gradient AttacksabstractTo study the resilience of distributed learning, the “Byzantine" literature considers a strong threat model where workers can report arbitrary gradients to the parameter server. Whereas this model helped obtain several fundamental results, it has sometimes been considered unrealistic, when the workers are mostly trustworthy machines. In this paper, we show a surprising equivalence between this model and data poisoning, a threat considered much more realistic. More specifically, we prove that every gradient attack can be reduced to data poisoning, in any personalized federated learning system with PAC guarantees (which we show are both desirable and realistic). This equivalence makes it possible to obtain new impossibility results on the resilience of any “robust” learning algorithm to data poisoning in highly heterogeneous applications, as corollaries of existing impossibility theorems on Byzantine machine learning. Moreover, using our equivalence, we derive a practical attack that we show (theoretically and empirically) can be very effective against classical personalized federated learning models. Sadegh Farhadkhani, Rachid Guerraoui, Lê-Nguyên Hoang, Oscar Villemaud |
ICML | 2 |
| 2022 | As easy as ABC: Optimal (A)ccountable (B)yzantine (C)onsensus is easy!abstractIt is known that the agreement property of the Byzantine consensus problem among$n$processes can be violated in a non-synchronous system if the number of faulty processes exceeds$t_{0}$= ┌$n$/3┐ − 1 [10], [19]. In this paper, we investigate the accountable Byzantine consensus problem in non-synchronous systems: the problem of solving Byzantine consensus whenever possible (e.g., when the number of faulty processes does not exceed$t_{0}$) and allowing correct processes to obtain proof of culpability of (at least)$t_{0}+ 1$faulty processes whenever correct processes disagree. We present four complementary contributions: 1) We introduce ABC: a simple yet efficient transformation of any Byzantine consensus protocol to an accountable one. ABC introduces an overhead of only two all-to-all communication rounds and$O(n^{2})$additional bits in executions with up to$t_{0}$faults (i.e., in the common case). 2) We define the accountability complexity, a complex-ity metric representing the number of accountability-specific messages that correct processes must send. Fur-thermore, we prove a tight lower bound. In particular, we show that any accountable Byzantine consensus protocol incurs cubic accountability complexity. Moreover, we illustrate that the bound is tight by applying the ABC transformation to any Byzantine consensus protocol. 3) We demonstrate that, when applied to an optimal Byzan-tine consensus protocol, ABC constructs an accountable Byzantine consensus protocol that is (1) optimal with respect to the communication complexity in solving consensus whenever consensus is solvable, and (2) op-timal with respect to the accountability complexity in obtaining accountability whenever disagreement occurs. 4) We generalize ABC to other distributed computing prob-lems besides the classic consensus problem. We charac-terize a class of agreement tasks, including reliable and consistent broadcast [5], that ABC renders accountable. Pierre Civit, Seth Gilbert, Vincent Gramoli, Rachid Guerraoui, Jovan Komatovic |
IPDPS | 4 |
| 2022 | The Universal Gossip FighterabstractThe notion of adversary is a staple of distributed computing. An adversary typically models “hostile” assumptions about the underlying distributed environment, e.g., a network that can drop messages, an operating system that can delay processes or an attacker that can hack machines. So far, the goal of distributed computing researchers has mainly been to develop a distributed algorithm that can face a given adversary, the abstraction characterizing worst-case scenarios. This paper initiates the study of the somehow opposite approach. Given a distributed algorithm, the adversary is the abstraction we seek to implement. More specifically, we consider the problem of controlling the spread of messages in a large-scale system, conveying the practical motivation of limiting the dissemination of fake news or viruses. Essentially, we assume a general class of gossip protocols, called all-to-all gossip protocols, and devise a practical method to hinder the dissemination. We present the Universal Gossip Fighter (UGF). Just like classical adversaries in distributed computing, UGF can observe the status of a dissemination and decide to stop some processes or delay some messages. The originality of UGF lies in the fact that it is universal, i.e., it applies to any all-to-all gossip protocol. We show that any gossip protocol attacked by UGF ends up exhibiting a quadratic message complexity (in the total number of processes) if it achieves sublinear time of dissemination. We also show that if a gossip protocol aims to achieve a message complexity$\alpha$times smaller than quadratic, then the time complexity rises exponentially in relation to$\alpha$. We convey the practical relevance of our theoretical findings by implementing UGF and conducting a set of empirical experiments that confirm some of our results. Anastasiia Gorbunova, Rachid Guerraoui, Anne-Marie Kermarrec, Anastasiia Kucherenko, Rafael Pinot |
IPDPS | 2 |
| 2022 | uKharon: A Membership Service for Microsecond Applications
Rachid Guerraoui, Antoine Murat, Javier Picorel, Athanasios Xygkis, Huabing Yan, Pengfei Zuo |
USENIX ATC | 1 |
| 2022 | Oracular Byzantine Reliable Broadcast
Martina Camaioni, Rachid Guerraoui, Matteo Monti, Manuel Vidigueira |
DISC | 2 |
| 2022 | Byzantine Consensus Is Θ(n²): The Dolev-Reischuk Bound Is Tight Even in Partial Synchrony!abstractThe Dolev-Reischuk bound says that any deterministic Byzantine consensus protocol has (at least) quadratic communication complexity in the worst case. While it has been shown that the bound is tight in synchronous environments, it is still unknown whether a consensus protocol with quadratic communication complexity can be obtained in partial synchrony. Until now, the most efficient known solutions for Byzantine consensus in partially synchronous settings had cubic communication complexity (e.g., HotStuff, binary DBFT). This paper closes the existing gap by introducing SQuad, a partially synchronous Byzantine consensus protocol with quadratic worst-case communication complexity. In addition, SQuad is optimally-resilient and achieves linear worst-case latency complexity. The key technical contribution underlying SQuad lies in the way we solve view synchronization, the problem of bringing all correct processes to the same view with a correct leader for sufficiently long. Concretely, we present RareSync, a view synchronization protocol with quadratic communication complexity and linear latency complexity, which we utilize in order to obtain SQuad. Pierre Civit, Muhammad Ayaz Dzulfikar, Seth Gilbert, Vincent Gramoli, Rachid Guerraoui, Jovan Komatovic, Manuel Vidigueira |
DISC | 5 |
| 2022 | Genuinely distributed Byzantine machine learningabstractAbstract Machine learning (ML) solutions are nowadays distributed, according to the so-called server/worker architecture. One server holds the model parameters while several workers train the model. Clearly, such architecture is prone to various types of component failures, which can be all encompassed within the spectrum of a Byzantine behavior. Several approaches have been proposed recently to tolerate Byzantine workers. Yet all require trusting a central parameter server. We initiate in this paper the study of the “general” Byzantine-resilient distributed machine learning problem where no individual component is trusted. In particular, we distribute the parameter server computation on several nodes. We show that this problem can be solved in an asynchronous system, despite the presence of $$\frac{1}{3}$$ 1 3 Byzantine parameter servers (i.e., $$n_{ps} > 3f_{ps}+1$$ n ps > 3 f ps + 1 ) and $$\frac{1}{3}$$ 1 3 Byzantine workers (i.e., $$n_w > 3f_w$$ n w > 3 f w ), which is asymptotically optimal. We present a new algorithm, ByzSGD, which solves the general Byzantine-resilient distributed machine learning problem by relying on three major schemes. The first, scatter/gather, is a communication scheme whose goal is to bound the maximum drift among models on correct servers. The second, distributed median contraction (DMC), leverages the geometric properties of the median in high dimensional spaces to bring parameters within the correct servers back close to each other, ensuring safe and lively learning. The third, Minimum-diameter averaging (MDA), is a statistically-robust gradient aggregation rule whose goal is to tolerate Byzantine workers. MDA requires a loose bound on the variance of non-Byzantine gradient estimates, compared to existing alternatives [e.g., Krum (Blanchard et al., in: Neural information processing systems, pp 118-128, 2017)]. Interestingly, ByzSGD ensures Byzantine resilience without adding communication rounds (on a normal path), compared to vanilla non-Byzantine alternatives. ByzSGD requires, however, a larger number of messages which, we show, can be reduced if we assume synchrony. We implemented ByzSGD on top of both TensorFlow and PyTorch, and we report on our evaluation results. In particular, we show that ByzSGD guarantees convergence with around 32% overhead compared to vanilla SGD. Furthermore, we show that ByzSGD’s throughput overhead is 24–176% in the synchronous case and 28–220% in the asynchronous case. El Mahdi El Mhamdi, Rachid Guerraoui, Arsany Guirguis, Lê-Nguyên Hoang, Sébastien Rouault |
Distributed Comput. | 2 |
| 2022 | The consensus number of a cryptocurrencyabstractAbstract Many blockchain-based algorithms, such as Bitcoin, implement a decentralized asset transfer system, often referred to as a cryptocurrency. As stated in the original paper by Nakamoto, at the heart of these systems lies the problem of preventing double-spending; this is usually solved by achieving consensus on the order of transfers among the participants. In this paper, we treat the asset transfer problem as a concurrent object and determine its consensus number, showing that consensus is, in fact, not necessary to prevent double-spending. We first consider the problem as defined by Nakamoto, where only a single process—the account owner—can withdraw from each account. Safety and liveness need to be ensured for correct account owners, whereas misbehaving account owners might be unable to perform transfers. We show that the consensus number of an asset transfer object is 1. We then consider a more general k-shared asset transfer object where up to k processes can atomically withdraw from the same account, and show that this object has consensus number k. We establish our results in the context of shared memory with benign faults, allowing us to properly understand the level of difficulty of the asset transfer problem. We also translate these results in the message passing setting with Byzantine players, a model that is more relevant in practice. In this model, we describe an asynchronous Byzantine fault-tolerant asset transfer implementation that is both simpler and more efficient than state-of-the-art consensus-based solutions. Our results are applicable to both the permissioned (private) and permissionless (public) setting, as normally their differentiation is hidden by the abstractions on top of which our algorithms are based. Rachid Guerraoui, Petr Kuznetsov, Matteo Monti, Matej Pavlovic, Dragos-Adrian Seredinschi |
Distributed Comput. | 1 |
| 2022 | Correction to: The consensus number of a cryptocurrency
Rachid Guerraoui, Petr Kuznetsov, Matteo Monti, Matej Pavlovic, Dragos-Adrian Seredinschi |
Distributed Comput. | 1 |
| 2022 | Removing algorithmic discrimination (with minimal individual error)
El Mahdi El Mhamdi, Rachid Guerraoui, Lê-Nguyên Hoang, Alexandre Maurer |
Theor. Comput. Sci. | 2 |
| 2022 | Byzantine-Resilient Multi-Agent SystemabstractWe consider the problem of making a multi-agent system (MAS) resilient to Byzantine failures through replication. We consider a very general model of MAS, where randomness can be involved in the behavior of each agent. We propose the first universal scheme to make such a MAS Byzantine-resilient. Rachid Guerraoui, Alexandre Maurer |
IEEE Trans. Dependable Secur. Comput. | 1 |
| 2022 | FLeet: Online Federated Learning via Staleness Awareness and Performance PredictionabstractFederated learning (FL) is very appealing for its privacy benefits: essentially, a global model is trained with updates computed on mobile devices while keeping the data of users local. Standard FL infrastructures are however designed to have no energy or performance impact on mobile devices, and are therefore not suitable for applications that require frequent ( online ) model updates, such as news recommenders. This article presents FLeet , the first Online FL system, acting as a middleware between the Android operating system and the machine learning application. FLeet combines the privacy of Standard FL with the precision of online learning thanks to two core components: (1) I-Prof , a new lightweight profiler that predicts and controls the impact of learning tasks on mobile devices, and (2) AdaSGD , a new adaptive learning algorithm that is resilient to delayed updates. Our extensive evaluation shows that Online FL, as implemented by FLeet , can deliver a 2.3× quality boost compared to Standard FL while only consuming 0.036% of the battery per day. I-Prof can accurately control the impact of learning tasks by improving the prediction accuracy by up to 3.6× in terms of computation time, and by up to 19× in terms of energy. AdaSGD outperforms alternative FL approaches by 18.4% in terms of convergence speed on heterogeneous data. Georgios Damaskinos, Rachid Guerraoui, Anne-Marie Kermarrec, Vlad Nitu, Rhicheek Patra, François Taïani |
ACM Trans. Intell. Syst. Technol. | 2 |
| 2021 | Differentially Private Stochastic Coordinate DescentabstractIn this paper we tackle the challenge of making the stochastic coordinate descent algorithm differentially private. Compared to the classical gradient descent algorithm where updates operate on a single model vector and controlled noise addition to this vector suffices to hide critical information about individuals, stochastic coordinate descent crucially relies on keeping auxiliary information in memory during training. This auxiliary information provides an additional privacy leak and poses the major challenge addressed in this work. Driven by the insight that under independent noise addition, the consistency of the auxiliary information holds in expectation, we present DP-SCD, the first differentially private stochastic coordinate descent algorithm. We analyze our new method theoretically and argue that decoupling and parallelizing coordinate updates is essential for its utility. On the empirical side we demonstrate competitive performance against the popular stochastic gradient descent alternative (DP-SGD) while requiring significantly less tuning. Georgios Damaskinos, Celestine Dünner, Rachid Guerraoui, Nikolaos Papandreou, Thomas P. Parnell |
AAAI | 3 |
| 2021 | GARFIELD: System Support for Byzantine Machine Learning (Regular Paper)abstractWe present GARFIELD, a library to transparently make machine learning (ML) applications, initially built with popular (but fragile) frameworks, e.g., TensorFlow and PyTorch, Byzantine-resilient. GARFIELD relies on a novel object-oriented design, reducing the coding effort, and addressing the vulnerability of the shared-graph architecture followed by classical ML frameworks. GARFIELD encompasses various communication patterns and supports computations on CPUs and GPUs, allowing addressing the general question of the practical cost of Byzantine resilience in ML applications. We report on the usage of GARFIELD on three main ML architectures: (a) a single server with multiple workers, (b) several servers and workers, and (c) peer-to-peer settings. Using GARFIELD, we highlight interesting facts about the cost of Byzantine resilience. In particular, (a) Byzantine resilience, unlike crash resilience, induces an accuracy loss, (b) the throughput overhead comes more from communication than from robust aggregation, and (c) tolerating Byzantine servers costs more than tolerating Byzantine workers. Rachid Guerraoui, Arsany Guirguis, Jérémy Plassmann, Anton Ragot, Sébastien Rouault |
DSN | 1 |
| 2021 | Leaderless ConsensusabstractClassical synchronous consensus algorithms are leaderless: processes exchange their proposals, retain the maximum value and decide when they see the same choice across a couple of rounds. Indulgent consensus algorithms are more robust in that they only require eventual synchrony, but are however typically leader-based. Intuitively, this is a weakness for a slow leader can delay any decision. This paper asks whether, under eventual synchrony, it is possible to deterministically solve consensus without a leader. The fact that the weakest failure detector to solve consensus is one that also eventually elects a leader seems to indicate that the answer to the question is negative. We prove in this paper that the answer is actually positive. We first give a precise definition of the very notion of a leaderless algorithm. Then we present three indulgent leaderless consensus algorithms, each we believe interesting in its own right: (i) for shared memory, (ii) for message passing with omission failures and (iii) for message passing with Byzantine failures (with and without authentication). Karolos Antoniadis, Antoine Desjardins, Vincent Gramoli, Rachid Guerraoui, Igor Zablotchi |
ICDCS | 4 |
| 2021 | Distributed Momentum for Byzantine-resilient Stochastic Gradient Descent
El Mahdi El Mhamdi, Rachid Guerraoui, Sébastien Rouault |
ICLR | 2 |
| 2021 | Collaborative Learning in the Jungle (Decentralized, Byzantine, Heterogeneous, Asynchronous and Nonconvex Learning)abstractWe study \emph{Byzantine collaborative learning}, where $n$ nodes seek to collectively learn from each others' local data. The data distribution may vary from one node to another. No node is trusted, and $f < n$ nodes can behave arbitrarily. We prove that collaborative learning is equivalent to a new form of agreement, which we call \emph{averaging agreement}. In this problem, nodes start each with an initial vector and seek to approximately agree on a common vector, which is close to the average of honest nodes' initial vectors. We present two asynchronous solutions to averaging agreement, each we prove optimal according to some dimension. The first, based on the minimum-diameter averaging, requires $n \geq 6f+1$, but achieves asymptotically the best-possible averaging constant up to a multiplicative constant. The second, based on reliable broadcast and coordinate-wise trimmed mean, achieves optimal Byzantine resilience, i.e., $n \geq 3f+1$. Each of these algorithms induces an optimal Byzantine collaborative learning protocol. In particular, our equivalence yields new impossibility theorems on what any collaborative learning algorithm can achieve in adversarial and heterogeneous environments. El Mahdi El Mhamdi, Sadegh Farhadkhani, Rachid Guerraoui, Arsany Guirguis, Lê-Nguyên Hoang, Sébastien Rouault |
NeurIPS | 3 |
| 2021 | Differential Privacy and Byzantine Resilience in SGD: Do They Add Up?abstractThis paper addresses the problem of combining Byzantine resilience with privacy in machine learning (ML). Specifically, we study if a distributed implementation of the renowned Stochastic Gradient Descent (SGD) learning algorithm is feasible withboth differential privacy (DP) and (α,f)-Byzantine resilience. To the best of our knowledge, this is the first work to tackle this problem from a theoretical point of view. A key finding of our analyses is that the classical approaches to these two (seemingly) orthogonal issues are incompatible. More precisely, we show that a direct composition of these techniques makes the guarantees of the resulting SGD algorithm depend unfavourably upon the number of parameters of the ML model, making the training of large models practically infeasible. We validate our theoretical results through numerical experiments on publicly-available datasets; showing that it is impractical to ensure DP and Byzantine resilience simultaneously. Rachid Guerraoui, Nirupam Gupta, Rafael Pinot, Sébastien Rouault, John Stephan |
PODC | 1 |
| 2021 | Frugal Byzantine ComputingabstractTraditional techniques for handling Byzantine failures are expensive: digital signatures are too costly, while using $3f{+}1$ replicas is uneconomical ($f$ denotes the maximum number of Byzantine processes). We seek algorithms that reduce the number of replicas to $2f{+}1$ and minimize the number of signatures. While the first goal can be achieved in the message-and-memory model, accomplishing the second goal simultaneously is challenging. We first address this challenge for the problem of broadcasting messages reliably. We consider two variants of this problem, Consistent Broadcast and Reliable Broadcast, typically considered very close. Perhaps surprisingly, we establish a separation between them in terms of signatures required. In particular, we show that Consistent Broadcast requires at least 1 signature in some execution, while Reliable Broadcast requires $O(n)$ signatures in some execution. We present matching upper bounds for both primitives within constant factors. We then turn to the problem of consensus and argue that this separation matters for solving consensus with Byzantine failures: we present a practical consensus algorithm that uses Consistent Broadcast as its main communication primitive. This algorithm works for $n=2f{+}1$ and avoids signatures in the common-case -- properties that have not been simultaneously achieved previously. Overall, our work approaches Byzantine computing in a frugal manner and motivates the use of Consistent Broadcast -- rather than Reliable Broadcast -- as a key primitive for reaching agreement. Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Dalia Papuc, Athanasios Xygkis, Igor Zablotchi |
DISC | 3 |
| 2021 | Probabilistic and temporal failure detectors for solving distributed problems
Rachid Guerraoui, David Kozhaya, Yvonne-Anne Pignolet |
J. Parallel Distributed Comput. | 1 |
| 2020 | Online Payments by Merely Broadcasting MessagesabstractWe address the problem of online payments, where users can transfer funds among themselves. We introduce Astro, a system solving this problem efficiently in a decentralized, deterministic, and completely asynchronous manner. Astro builds on the insight that consensus is unnecessary to prevent double-spending. Instead of consensus, Astro relies on a weaker primitive---Byzantine reliable broadcast---enabling a simpler and more efficient implementation than consensus-based payment systems. In terms of efficiency, Astro executes a payment by merely broadcasting a message. The distinguishing feature of Astro is that it can maintain performance robustly, i.e., remain unaffected by a fraction of replicas being compromised or slowed down by an adversary. Our experiments on a public cloud network show that Astro can achieve near-linear scalability in a sharded setup, going from 10K payments/sec (2 shards) to 20K payments/sec (4 shards). In a nutshell, Astro can match VISA-level average payment throughput, and achieves a 5× improvement over a state-of-the-art consensus-based solution, while exhibiting sub-second 95^th percentile latency. Daniel Collins 0001, Rachid Guerraoui, Jovan Komatovic, Petr Kuznetsov, Matteo Monti, Matej Pavlovic, Yvonne-Anne Pignolet, Dragos-Adrian Seredinschi, Andrei Tonkikh, Athanasios Xygkis |
DSN | 2 |
| 2020 | Thread-Placement LearningabstractIn a non-uniform memory access machine, the placement of software threads to hardware cores can have a significant effect on the performance of concurrent applications. Detecting the best possible placement for each application is a necessity for thread scheduling. Yet, due to the difficulty of this problem, operating-system schedulers do not really try to understand the needs of applications, but rather focus on (non-portable) scheduling heuristics. In this paper, we introduce thread-placement learning (TPLE), a technique for understanding the placement requirements of applications. TPLE utilizes machine learning and performance counters for choosing between different placement policies. To feed the machine learning model, TPLE requires a set of portable microbenchmarks that produce training data-i.e., performance counter measurements-for all the target placement policies. We use this data to train a classifier that is able to choose between these policies online in order to change the thread-placement of a running application. We demonstrate the practicality of TPLE by implementing a thread-placement algorithm, named Slate. Slate is able to automatically and online (i.e., in runtime) select between the two most commonly-used placement policies, namely locality and round-robin placement on the nodes of a multicore. To the best of our knowledge, Slate is the first online thread-placement algorithm that utilizes machine learning in combination with performance counters. We evaluate Slate and show that it achieves up to 93% accuracy in its decisions and outperforms the Linux scheduler by up to 16%. Karolos Antoniadis, Rachid Guerraoui, Vasileios Trigonakis |
ICDCS | 2 |
| 2020 | The Impossibility of Fast TransactionsabstractWe prove that transactions cannot be fast in an asynchronous fault-tolerant system. Our result holds in any system where we require transactions to ensure monotonic writes, or any stronger consistency model, such as, causal consistency. Thus, our result unveils an important, and so far unknown, limitation of fast transactions: they are impossible if we want to tolerate the failure of even one server. Karolos Antoniadis, Diego Didona, Rachid Guerraoui, Willy Zwaenepoel |
IPDPS | 3 |
| 2020 | FLeet: Online Federated Learning via Staleness Awareness and Performance PredictionabstractFederated Learning (FL) is very appealing for its privacy benefits: essentially, a global model is trained with updates computed on mobile devices while keeping the data of users local. Standard FL infrastructures are however designed to have no energy or performance impact on mobile devices, and are therefore not suitable for applications that require frequent (online) model updates, such as news recommenders. Georgios Damaskinos, Rachid Guerraoui, Anne-Marie Kermarrec, Vlad Nitu, Rhicheek Patra, François Taïani |
Middleware | 2 |
| 2020 | FeGAN: Scaling Distributed GANsabstractExisting approaches to distribute Generative Adversarial Networks (GANs) either (i) fail to scale for they typically put the two components of a GAN (the generator and the discriminator) on different machines, inducing significant communication overhead, or (ii) they face GAN training specific issues, exacerbated by distribution. Rachid Guerraoui, Arsany Guirguis, Anne-Marie Kermarrec, Erwan Le Merrer |
Middleware | 1 |
| 2020 | AKSEL: Fast Byzantine SGDabstractModern machine learning architectures distinguish servers and workers. Typically, a d-dimensional model is hosted by a server and trained by n workers, using a distributed stochastic gradient descent (SGD) optimization scheme. At each SGD step, the goal is to estimate the gradient of a cost function. The simplest way to do this is to average the gradients estimated by the workers. However, averaging is not resilient to even one single Byzantine failure of a worker. Many alternative gradient aggregation rules (GARs) have recently been proposed to tolerate a maximum number f of Byzantine workers. These GARs differ according to (1) the complexity of their computation time, (2) the maximal number of Byzantine workers despite which convergence can still be ensured (breakdown point), and (3) their accuracy, which can be captured by (3.1) their angular error, namely the angle with the true gradient, as well as (3.2) their ability to aggregate full gradients. In particular, many are not full gradients for they operate on each dimension separately, which results in a coordinate-wise blended gradient, leading to low accuracy in practical situations where the number (s) of workers that are actually Byzantine in an execution is small (s < < f). We propose Aksel, a new scalable median-based GAR with optimal time complexity (𝒪(nd)), optimal breakdown point (n > 2f) and the lowest upper bound on the expected angular error (𝒪(√d)) among full gradient approaches. We also study the actual angular error of Aksel when the gradient distribution is normal and show that it only grows in 𝒪(√dlog{n}), which is the first logarithmic upper bound ever proven on the number of workers n assuming an optimal breakdown point. We also report on an empirical evaluation of Aksel on various classification tasks, which we compare to alternative GARs against state-of-the-art attacks. Aksel is the only GAR reaching top accuracy when there is actually none or few Byzantine workers while maintaining a good defense even under the extreme case (s = f). For simplicity of presentation, we consider a scheme with a single server. However, as we explain in the paper, Aksel can also easily be adapted to multi-server architectures that tolerate the Byzantine behavior of a fraction of the servers. Amine Boussetta, El Mahdi El Mhamdi, Rachid Guerraoui, Alexandre Maurer, Sébastien Rouault |
OPODIS | 3 |
| 2020 | Dynamic Byzantine Reliable BroadcastabstractReliable broadcast is a communication primitive guaranteeing, intuitively, that all processes in a distributed system deliver the same set of messages. The reason why this primitive is appealing is twofold: (i) we can implement it deterministically in a completely asynchronous environment, unlike stronger primitives like consensus and total-order broadcast, and yet (ii) reliable broadcast is powerful enough to implement important applications like payment systems. The problem we tackle in this paper is that of dynamic reliable broadcast, i.e., enabling processes to join or leave the system. This property is desirable for long-lived applications (aiming to be highly available), yet has been precluded in previous asynchronous reliable broadcast protocols. We study this property in a general adversarial (i.e., Byzantine) environment. We introduce the first specification of a dynamic Byzantine reliable broadcast (DBRB) primitive that is amenable to an asynchronous implementation. We then present an algorithm implementing this specification in an asynchronous network. Our DBRB algorithm ensures that if any correct process in the system broadcasts a message, then every correct process delivers that message unless it leaves the system. Moreover, if a correct process delivers a message, then every correct process that has not expressed its will to leave the system delivers that message. We assume that more than $2/3$ of processes in the system are correct at all times, which is tight in our context. We also show that if only one process in the system can fail---and it can fail only by crashing---then it is impossible to implement a stronger primitive, ensuring that if any correct process in the system broadcasts or delivers a message, then every correct process in the system delivers that message---including those that leave. Rachid Guerraoui, Jovan Komatovic, Petr Kuznetsov, Yvonne-Anne Pignolet, Dragos-Adrian Seredinschi, Andrei Tonkikh |
OPODIS | 1 |
| 2020 | Microsecond Consensus for Microsecond Applications
Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Virendra J. Marathe, Athanasios Xygkis, Igor Zablotchi |
OSDI | 3 |
| 2020 | Genuinely Distributed Byzantine Machine LearningabstractMachine Learning (ML) solutions are nowadays distributed, according to the so-called server/worker architecture. One server holds the model parameters while several workers train the model. Clearly, such architecture is prone to various types of component failures, which can be all encompassed within the spectrum of a Byzantine behavior. Several approaches have been proposed recently to tolerate Byzantine workers. Yet all require trusting a central parameter server. We initiate in this paper the study of the "general" Byzantine-resilient distributed machine learning problem where no individual component is trusted. In particular, we distribute the parameter server computation on several nodes. El Mahdi El Mhamdi, Rachid Guerraoui, Arsany Guirguis, Lê-Nguyên Hoang, Sébastien Rouault |
PODC | 2 |
| 2020 | Robust P2P Personalized LearningabstractDecentralized machine learning over peer-to-peer networks is very appealing for it enables to learn personalized models without sharing users data, nor relying on any central server. Peers can improve upon their locally trained model across a network graph of other peers with similar objectives. Whilst they offer an inherently scalable scheme with a very simple cost-efficient learning model, peer-to-peer networks are also fragile. In particular, they can be very easily disrupted by unfairness, free-riding, and adversarial behaviors. In this paper, we present CDPL (Contribution Driven P2P Learning), a novel Byzantine-resilient distributed algorithm to train personalized models across similar peers. We convey theoretically and empirically the effectiveness of CDPL in terms of speed of convergence as well as robustness to Byzantine behavior. Karim Boubouh, Amine Boussetta, Yahya Benkaouz, Rachid Guerraoui |
SRDS | 4 |
| 2020 | Fast and Robust Distributed Learning in High DimensionabstractCould a gradient aggregation rule (GAR) for distributed machine learning be both robust and fast? This paper answers by the affirmative through Multi-Bulyan. Given n workers, f of which are arbitrary malicious (Byzantine) and m = n - f are not, we prove that Multi-Bulyan can ensure a strong form of Byzantine resilience, as well as an m/n slowdown, compared to averaging, the fastest (but non Byzantine resilient) rule for distributed machine learning. When m ≈ n (almost all workers are correct), Multi-Bulyan reaches the speed of averaging. We also prove that Multi-Bulyan's cost in local computation is O(d) (like averaging), an important feature for ML where d commonly reaches 109, while robust alternatives have at least quadratic cost in d. Our theoretical findings are complemented with an experimental evaluation which, in addition to supporting the linear O(d) complexity argument, conveys the fact that Multi-Bulyan's parallelisability further adds to its efficiency. El Mahdi El Mhamdi, Rachid Guerraoui, Sébastien Rouault |
SRDS | 2 |
| 2020 | Who Started This Rumor? Quantifying the Natural Differential Privacy of Gossip ProtocolsabstractGossip protocols (also called rumor spreading or epidemic protocols) are widely used to disseminate information in massive peer-to-peer networks. These protocols are often claimed to guarantee privacy because of the uncertainty they introduce on the node that started the dissemination. But is that claim really true? Can the source of a gossip safely hide in the crowd? This paper examines, for the first time, gossip protocols through a rigorous mathematical framework based on differential privacy to determine the extent to which the source of a gossip can be traceable. Considering the case of a complete graph in which a subset of the nodes are curious, we study a family of gossip protocols parameterized by a "muting" parameter s: nodes stop emitting after each communication with a fixed probability 1-s. We first prove that the standard push protocol, corresponding to the case s = 1, does not satisfy differential privacy for large graphs. In contrast, the protocol with s = 0 (nodes forward only once) achieves optimal privacy guarantees but at the cost of a drastic increase in the spreading time compared to standard push, revealing an interesting tension between privacy and spreading time. Yet, surprisingly, we show that some choices of the muting parameter s lead to protocols that achieve an optimal order of magnitude in both privacy and speed. Privacy guarantees are obtained by showing that only a small fraction of the possible observations by curious nodes have different probabilities when two different nodes start the gossip, since the source node rapidly stops emitting when s is small. The speed is established by analyzing the mean dynamics of the protocol, and leveraging concentration inequalities to bound the deviations from this mean behavior. We also confirm empirically that, with appropriate choices of s, we indeed obtain protocols that are very robust against concrete source location attacks (such as maximum a posteriori estimates) while spreading the information almost as fast as the standard (and non-private) push protocol. Aurélien Bellet, Rachid Guerraoui, Hadrien Hendrikx |
DISC | 2 |
| 2020 | Efficient Multi-Word Compare and SwapabstractAtomic lock-free multi-word compare-and-swap (MCAS) is a powerful tool for designing concurrent algorithms. Yet, its widespread usage has been limited because lock-free implementations of MCAS make heavy use of expensive compare-and-swap (CAS) instructions. Existing MCAS implementations indeed use at least 2k+1 CASes per k-CAS. This leads to the natural desire to minimize the number of CASes required to implement MCAS. We first prove in this paper that it is impossible to "pack" the information required to perform a k-word CAS (k-CAS) in less than k locations to be CASed. Then we present the first algorithm that requires k+1 CASes per call to k-CAS in the common uncontended case. We implement our algorithm and show that it outperforms a state-of-the-art baseline in a variety of benchmarks in most considered workloads. We also present a durably linearizable (persistent memory friendly) version of our MCAS algorithm using only 2 persistence fences per call, while still only requiring k+1 CASes per k-CAS. Rachid Guerraoui, Alex Kogan, Virendra J. Marathe, Igor Zablotchi |
DISC | 1 |
| 2020 | Smaller, Faster & Lighter KNN Graph ConstructionsabstractWe propose GoldFinger, a new compact and fast-to-compute binary representation of datasets to approximate Jaccard’s index. We illustrate the effectiveness of GoldFinger on the emblematic big data problem of K-Nearest-Neighbor (KNN) graph construction and show that GoldFinger can drastically accelerate a large range of existing KNN algorithms with little to no overhead. As a side effect, we also show that the compact representation of the data protects users’ privacy for free by providing k-anonymity and l-diversity. Our extensive evaluation of the resulting approach on several realistic datasets shows that our approach delivers speedups of up to 78.9% compared to the use of raw data while only incurring a negligible to moderate loss in terms of KNN quality. To convey the practical value of such a scheme, we apply it to item recommendation and show that the loss in recommendation quality is negligible. Rachid Guerraoui, Anne-Marie Kermarrec, Olivier Ruas, François Taïani |
WWW | 1 |
| 2020 | The Cost of Scaling a Reliable Interconnection TopologyabstractIn distributed computing, many papers try to evaluate the message complexity of a distributed system as a function of the number of nodes n. But what about the cost of building the distributed system itself? Assuming that we want to reliably connect n nodes, how does the total number of nodes of the network evolve with n? Addressing such a question lies at the heart of achieving scalability in cloud computing. In this paper, we give the explicit description of a distributed system of which any two of then nodes, for any n, remain connected (by a path of alive nodes and channels) with probability at least m, despite the very fact that (a) every other node or channel has an independent probability λ of failing, and (b) the number of channels connected to every node is physically bounded by a constant. We show however that if we also require any two of the n nodes to maintain a balanced message throughput with a constant probability, then O (nlog1+εn) additional intermediary nodes are sufficient, where ε > 0 is an arbitrarily small constant. Rachid Guerraoui, Alexandre Maurer |
IEEE Trans. Dependable Secur. Comput. | 1 |
| 2019 | Unified and Scalable Incremental Recommenders with Consumed Item Packs
Rachid Guerraoui, Erwan Le Merrer, Rhicheek Patra, Jean-Ronan Vigouroux |
Euro-Par | 1 |
| 2019 | Fingerprinting Big Data: The Case of KNN Graph ConstructionabstractWe propose fingerprinting, a new technique that consists in constructing compact, fast-to-compute and privacy-preserving binary representations of datasets. We illustrate the effectiveness of our approach on the emblematic big data problem of K-Nearest-Neighbor (KNN) graph construction and show that fingerprinting can drastically accelerate a large range of existing KNN algorithms, while efficiently obfuscating the original data, with little to no overhead. Our extensive evaluation of the resulting approach (dubbed GoldFinger) on several realistic datasets shows that our approach delivers speedups of up to 78.9% compared to the use of raw data while only incurring a negligible to moderate loss in terms of KNN quality. Rachid Guerraoui, Anne-Marie Kermarrec, Olivier Ruas, François Taïani |
ICDE | 1 |
| 2019 | Demystifying Bitcoin (Keynote Abstract)abstractThis talk will explain the bitcoin algorithm from the distributed computing perspective, precisely define the underlying double-payment problem, and present a much simpler alternative to solve the problem without relying on consensus and consuming so much energy. Rachid Guerraoui is professor in Computer Science at EPFL where he leads the Distributed Computing Laboratory. He worked in the past with École des Mines de Paris, CEA Saclay, HP Labs in Palo Alto and MIT. He has been elected ACM Fellow and Professor of the College de France. He was awarded a Senior ERC Grant and a Google Focused Award. Rachid Guerraoui |
OPODIS | 1 |
| 2019 | The Impact of RDMA on AgreementabstractRemote Direct Memory Access (RDMA) is becoming widely available in data centers. This technology allows a process to directly read and write the memory of a remote host, with a mechanism to control access permissions. In this paper, we study the fundamental power of these capabilities. We consider the well-known problem of achieving consensus despite failures, and find that RDMA can improve the inherent trade-off in distributed computing between failure resilience and performance. Specifically, we show that RDMA allows algorithms that simultaneously achieve high resilience and high performance, while traditional algorithms had to choose one or another. With Byzantine failures, we give an algorithm that only requires n \geq 2f_P + 1 processes (where f_P is the maximum number of faulty processes) and decides in two (network) delays in common executions. With crash failures, we give an algorithm that only requires n \geq f_P + 1 processes and also decides in two delays. Both algorithms tolerate a minority of memory failures inherent to RDMA, and they provide safety in asynchronous systems and liveness with standard additional assumptions. Marcos K. Aguilera, Naama Ben-David, Rachid Guerraoui, Virendra J. Marathe, Igor Zablotchi |
PODC | 3 |
| 2019 | The Consensus Number of a CryptocurrencyabstractMany blockchain-based algorithms, such as Bitcoin, implement a decentralized asset transfer system, often referred to as a cryptocurrency. As stated in the original paper by Nakamoto, at the heart of these systems lies the problem of preventing double-spending ; this is usually solved by achieving consensus on the order of transfers among the participants. By treating the asset transfer problem as a concurrent object and determining its consensus number, we show that consensus is not necessary to prevent double-spending. We first consider the problem as defined by Nakamoto, where only a single process---the account owner---can withdraw from each account. Safety and liveness need to be ensured for correct account owners, whereas misbehaving account owners might be unable to perform transfers. We show that the consensus number of an asset transfer object is 1. We then consider a more general k-shared asset transfer object where up to k processes can atomically withdraw from the same account, and show that this object has consensus number k. We first establish these these results in the context of shared memory with benign faults, in order to properly understand the level of difficulty of the asset transfer problem. Then, we translate our result in the more practically relevant message passing setting with Byzantine players. We describe an asynchronous Byzantine fault-tolerant asset transfer implementation that is both simpler and more efficient than state-of-the-art consensus-based solutions. Our results are applicable to both the permissioned (private) and permissionless (public) setting, as normally their differentiation is hidden by the abstractions on top of which our algorithms are based. Rachid Guerraoui, Petr Kuznetsov, Matteo Monti, Matej Pavlovic, Dragos-Adrian Seredinschi |
PODC | 1 |
| 2019 | Distributed Transactional Systems Cannot Be FastabstractWe prove that no fully transactional system can provide fast read transactions (including read-only ones that are considered the most frequent in practice). Specifically, to achieve fast read transactions, the system has to give up support of transactions that write more than one object. We prove this impossibility result for distributed storage systems that are causally consistent, i.e., they do not require to ensure any strong form of consistency. Therefore, our result holds also for any system that ensures a consistency level stronger than causal consistency, e.g., strict serializability. The impossibility result holds even for systems that store only two objects (and support at least two servers and at least four clients). It also holds for systems that are partially replicated. Our result justifies the design choices of state-of-the-art distributed transactional systems and insists that system designers should not put more effort to design fully-functional systems that support both fast read transactions and ensure causal or any stronger form of consistency. Diego Didona, Panagiota Fatourou, Rachid Guerraoui, Jingjing Wang 0007, Willy Zwaenepoel |
SPAA | 3 |
| 2019 | Scalable Byzantine Reliable BroadcastabstractByzantine reliable broadcast is a powerful primitive that allows a set of processes to agree on a message from a designated sender, even if some processes (including the sender) are Byzantine. Existing broadcast protocols for this setting scale poorly, as they typically build on quorum systems with strong intersection guarantees, which results in linear per-process communication and computation complexity. We generalize the Byzantine reliable broadcast abstraction to the probabilistic setting, allowing each of its properties to be violated with a fixed, arbitrarily small probability. We leverage these relaxed guarantees in a protocol where we replace quorums with stochastic samples. Compared to quorums, samples are significantly smaller in size, leading to a more scalable design. We obtain the first Byzantine reliable broadcast protocol with logarithmic per-process communication and computation complexity. We conduct a complete and thorough analysis of our protocol, deriving bounds on the probability of each of its properties being compromised. During our analysis, we introduce a novel general technique that we call adversary decorators. Adversary decorators allow us to make claims about the optimal strategy of the Byzantine adversary without imposing any additional assumptions. We also introduce Threshold Contagion, a model of message propagation through a system with Byzantine processes. To the best of our knowledge, this is the first formal analysis of a probabilistic broadcast protocol in the Byzantine fault model. We show numerically that practically negligible failure probabilities can be achieved with realistic security parameters. Rachid Guerraoui, Petr Kuznetsov, Matteo Monti, Matej Pavlovic, Dragos-Adrian Seredinschi |
DISC | 1 |
| 2019 | The weakest failure detector for eventual consistency
Swan Dubois, Rachid Guerraoui, Petr Kuznetsov, Franck Petit, Pierre Sens 0001 |
Distributed Comput. | 2 |
| 2019 | The PCL Theorem: Transactions cannot be Parallel, Consistent, and LiveabstractWe establish a theorem called the PCL theorem, which states that it is impossible to design a transactional memory algorithm that ensures (1) parallelism , i.e., transactions do not need to synchronize unless they access the same application objects, (2) very little consistency , i.e., a consistency condition, called weak adaptive consistency , introduced here and that is weaker than snapshot isolation, processor consistency, and any other consistency condition stronger than them (such as opacity, serializability, causal serializability, etc.), and (3) very little liveness , i.e., which transactions eventually commit if they run solo. Victor Bushkov, Dmytro Dziuma, Panagiota Fatourou, Rachid Guerraoui |
J. ACM | 4 |
| 2018 | Personalized and Private Peer-to-Peer Machine LearningabstractThe rise of connected personal devices together with privacy concerns call for machine learning algorithms capable of leveraging the data of a large number of agents to learn personalized models under strong privacy requirements. In this paper, we introduce an efficient algorithm to address the above problem in a fully decentralized (peer-to-peer) and asynchronous fashion, with provable convergence rate. We show how to make the algorithm differentially private to protect against the disclosure of information about the personal datasets, and formally analyze the trade-off between utility and privacy. Our experiments show that our approach dramatically outperforms previous work in the non-private case, and that under privacy constraints, we can significantly improve over models learned in isolation. Aurélien Bellet, Rachid Guerraoui, Mahsa Taziki, Marc Tommasi |
AISTATS | 2 |
| 2018 | Collaborative Filtering Under a Sybil Attack: Similarity Metrics do Matter!abstractRecommendation systems help users identify interesting content, but they also open new privacy threats. In this paper, we deeply analyze the effect of a Sybil attack that tries to infer information on users from a user-based collaborative-filtering recommendation systems. We discuss the impact of different similarity metrics used to identity users with similar tastes in the trade-off between recommendation quality and privacy. Finally, we propose and evaluate a novel similarity metric that combines the best of both worlds: a high recommendation quality with a low prediction accuracy for the attacker. Our results, on a state-of-the-art recommendation framework and on real datasets show that existing similarity metrics exhibit a wide range of behaviors in the presence of Sybil attacks, while our new similarity metric consistently achieves the best trade-off while outperforming state-of-the-art solutions. Antoine Boutet, Florestan De Moor, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Antoine Rault |
DSN | 4 |
| 2018 | Monotonic Prefix Consistency in Distributed Systems
Alain Girault, Gregor Gößler, Rachid Guerraoui, Jad Hamza, Dragos-Adrian Seredinschi |
FORTE | 3 |
| 2018 | Asynchronous Byzantine Machine Learning (the case of SGD)abstractAsynchronous distributed machine learning solutions have proven very effective so far, but always assuming perfectly functioning workers. In practice, some of the workers can however exhibit Byzantine behavior, caused by hardware failures, software bugs, corrupt data, or even malicious attacks. We introduce Kardam, the first distributed asynchronous stochastic gradient descent (SGD) algorithm that copes with Byzantine workers. Kardam consists of two complementary components: a filtering and a dampening component. The first is scalar-based and ensures resilience against 1/3 Byzantine workers. Essentially, this filter leverages the Lipschitzness of cost functions and acts as a self-stabilizer against Byzantine workers that would attempt to corrupt the progress of SGD. The dampening component bounds the convergence rate by adjusting to stale information through a generic gradient weighting scheme. We prove that Kardam guarantees almost sure convergence in the presence of asynchrony and Byzantine behavior, and we derive its convergence rate. We evaluate Kardam on the CIFAR100 and EMNIST datasets and measure its overhead with respect to non Byzantine-resilient solutions. We empirically show that Kardam does not introduce additional noise to the learning procedure but does induce a slowdown (the cost of Byzantine resilience) that we both theoretically and empirically show to be less than f/n, where f is the number of Byzantine failures tolerated and n the total number of workers. Interestingly, we also empirically observe that the dampening component is interesting in its own right for it enables to build an SGD algorithm that outperforms alternative staleness-aware asynchronous competitors in environments with honest workers. Georgios Damaskinos, El Mahdi El Mhamdi, Rachid Guerraoui, Rhicheek Patra, Mahsa Taziki |
ICML | 3 |
| 2018 | The Hidden Vulnerability of Distributed Learning in ByzantiumabstractWhile machine learning is going through an era of celebrated success, concerns have been raised about the vulnerability of its backbone: stochastic gradient descent (SGD). Recent approaches have been proposed to ensure the robustness of distributed SGD against adversarial (Byzantine) workers sending poisoned gradients during the training phase. Some of these approaches have been proven Byzantine–resilient: they ensure the convergence of SGD despite the presence of a minority of adversarial workers. We show in this paper that convergence is not enough. In high dimension $d \gg 1$, an adver\-sary can build on the loss function’s non–convexity to make SGD converge to ineffective models. More precisely, we bring to light that existing Byzantine–resilient schemes leave a margin of poisoning of $\bigOmega\left(f(d)\right)$, where $f(d)$ increases at least like $\sqrt[p]{d }$. Based on this leeway, we build a simple attack, and experimentally show its strong to utmost effectivity on CIFAR–10 and MNIST. We introduce Bulyan, and prove it significantly reduces the attackers leeway to a narrow $\bigO\,( \sfrac{1}{\sqrt{d }})$ bound. We empirically show that Bulyan does not suffer the fragility of existing aggregation rules and, at a reasonable cost in terms of required batch size, achieves convergence as if only non–Byzantine gradients had been used to update the model. El Mahdi El Mhamdi, Rachid Guerraoui, Sébastien Rouault |
ICML | 2 |
| 2018 | SPADE: Tuning scale-out OLTP on modern RDMA clusters
Georgios Chatzopoulos, Aleksandar Dragojevic, Rachid Guerraoui |
Middleware | 3 |
| 2018 | Passing Messages while Sharing MemoryabstractWe introduce a new distributed computing model called m&m that allows processes to both pass messages and share memory. Motivated by recent hardware trends, we find that this model improves the power of the pure message-passing and shared-memory models. As we demonstrate by example with two fundamental problems---consensus and eventual leader election---the added power leads to new algorithms that are more robust against failures and asynchrony. Our consensus algorithm combines the superior scalability of message passing with the higher fault tolerance of shared memory, while our leader election algorithms reduce the system synchrony needed for correctness. These results point to a wide new space for future exploration of other problems, techniques, and benefits. Marcos K. Aguilera, Naama Ben-David, Irina Calciu, Rachid Guerraoui, Erez Petrank, Sam Toueg |
PODC | 4 |
| 2018 | Locking Timestamps versus Locking Objects
Marcos K. Aguilera, Tudor David, Rachid Guerraoui, Junxiong Wang |
PODC | 3 |
| 2018 | The Inherent Cost of Remembering ConsistentlyabstractNon-volatile memory (NVM) promises fast, byte-addressable and durable storage, with raw access latencies in the same order of magnitude as DRAM. But in order to take advantage of the durability of NVM, programmers need to design \em persistent objects which maintain consistent state across system crashes and restarts. Concurrent implementations of persistent objects typically make heavy use of expensive persistent fence instructions to order NVM accesses, thus negating some of the performance benefits of NVM. This raises the question of the minimal number of persistent fence instructions required to implement a persistent object. We answer this question in the deterministic lock-free case by providing lower and upper bounds on the required number of fence instructions. We obtain our upper bound by presenting a new universal construction that implements durably any object using at most one persistent fence per update operation invoked. Our lower bound states that in the worst case, each process needs to issue at least one persistent fence per update operation invoked. Nachshon Cohen, Rachid Guerraoui, Igor Zablotchi |
SPAA | 2 |
| 2018 | Log-Free Concurrent Data Structures
Tudor David, Aleksandar Dragojevic, Rachid Guerraoui, Igor Zablotchi |
USENIX ATC | 3 |
| 2018 | State Machine Replication Is More Expensive Than Consensus
Karolos Antoniadis, Rachid Guerraoui, Dahlia Malkhi, Dragos-Adrian Seredinschi |
DISC | 2 |
| 2018 | The entropy of a distributed computation random number generation from memory interleaving
Karolos Antoniadis, Peva Blanchard, Rachid Guerraoui, Julien Stainer |
Distributed Comput. | 3 |
| 2018 | TM2C: a software transactional memory for many-cores
Vincent Gramoli, Rachid Guerraoui, Vasileios Trigonakis |
Distributed Comput. | 2 |
| 2018 | Causal Consistency and Latency Optimality: Friend or Foe?abstractCausal consistency is an attractive consistency model for geo-replicated data stores. It is provably the strongest model that tolerates network partitions. It avoids the long latencies associated with strong consistency, and, especially when using read-only transactions (ROTs), it prevents many of the anomalies of weaker consistency models. Recent work has shown that causal consistency allows "latency-optimal" ROTs, that are nonblocking, single-round and single-version in terms of communication. On the surface, this latency optimality is very appealing, as the vast majority of applications are assumed to have read-dominated workloads. In this paper, we show that such "latency-optimal" ROTs induce an extra overhead on writes that is so high that it actually jeopardizes performance even in read-dominated workloads. We show this result from a practical as well as from a theoretical angle. We present the Contrarian protocol that implements "almost latency-optimal" ROTs, but that does not impose on the writes any of the overheads incurred by latency-optimal protocols. In Contrarian, ROTs are nonblocking and single-version, but they require two rounds of client-server communication. We experimentally show that this protocol not only achieves higher throughput, but, surprisingly, also provides better latencies for all but the lowest loads and the most read-heavy workloads. We furthermore prove that the extra overhead imposed on writes by latency-optimal ROTs is inherent, i.e., it is not an artifact of the design we consider, and cannot be avoided by any implementation of latency-optimal ROTs. We show in particular that this overhead grows linearly with the number of clients. Diego Didona, Rachid Guerraoui, Jingjing Wang 0007, Willy Zwaenepoel |
Proc. VLDB Endow. | 2 |
| 2018 | Lock-Unlock: Is That All? A Pragmatic Analysis of Locking in Software SystemsabstractA plethora of optimized mutex lock algorithms have been designed over the past 25 years to mitigate performance bottlenecks related to critical sections and locks. Unfortunately, there is currently no broad study of the behavior of these optimized lock algorithms on realistic applications that consider different performance metrics, such as energy efficiency and tail latency. In this article, we perform a thorough and practical analysis of synchronization, with the goal of providing software developers with enough information to design fast, scalable, and energy-efficient synchronization in their systems. First, we perform a performance study of 28 state-of-the-art mutex lock algorithms, on 40 applications, on four different multicore machines. We consider not only throughput (traditionally the main performance metric) but also energy efficiency and tail latency, which are becoming increasingly important. Second, we present an in-depth analysis in which we summarize our findings for all the studied applications. In particular, we describe nine different lock-related performance bottlenecks, and we propose six guidelines helping software developers with their choice of a lock algorithm according to the different lock properties and the application characteristics. From our detailed analysis, we make several observations regarding locking algorithms and application behaviors, several of which have not been previously discovered: (i) applications stress not only the lock–unlock interface but also the full locking API (e.g., trylocks, condition variables); (ii) the memory footprint of a lock can directly affect the application performance; (iii) for many applications, the interaction between locks and scheduling is an important application performance factor; (vi) lock tail latencies may or may not affect application tail latency; (v) no single lock is systematically the best; (vi) choosing the best lock is difficult; and (vii) energy efficiency and throughput go hand in hand in the context of lock algorithms. These findings highlight that locking involves more considerations than the simple lock/unlock interface and call for further research on designing low-memory footprint adaptive locks that fully and efficiently support the full lock interface, and consider all performance metrics. Rachid Guerraoui, Hugo Guiroux, Renaud Lachaize, Vivien Quéma, Vasileios Trigonakis |
ACM Trans. Comput. Syst. | 1 |
| 2017 | I Know Nothing about You But Here is What You Might LikeabstractRecommenders widely use collaborative filtering schemes. These schemes, however, threaten privacy as user profiles are made available to the service provider hosting the recommender and can even be guessed by curious users who analyze the recommendations. Users can encrypt their profiles to hide them from the service provider and add noise to make them difficult to guess. These precautionary measures hamper latency and recommendation quality. In this paper, we present a novel recommender, X-REC, enabling an effective collaborative filtering scheme to ensure the privacy of users against the service provider (system-level privacy) or other users (user-level privacy). X-REC builds on two underlying services: X-HE, an encryption scheme designed for recommenders, and X-NN, a neighborhood selection protocol over encrypted profiles. We leverage uniform sampling to ensure differential privacy against curious users. Our extensive evaluation demonstrates that X-REC provides (1) recommendation quality similar to non-private recommenders, and (2) significant latency improvement over privacy-aware alternatives. Rachid Guerraoui, Anne-Marie Kermarrec, Rhicheek Patra, Mahammad Valiyev, Jingjing Wang 0007 |
DSN | 1 |
| 2017 | FloDB: Unlocking Memory in Persistent Key-Value StoresabstractLog-structured merge (LSM) data stores enable to store and process large volumes of data while maintaining good performance. They mitigate the I/O bottleneck by absorbing updates in a memory layer and transferring them to the disk layer in sequential batches. Yet, the LSM architecture fundamentally requires elements to be in sorted order. As the amount of data in memory grows, maintaining this sorted order becomes increasingly costly. Contrary to intuition, existing LSM systems could actually lose throughput with larger memory components. Oana Balmau, Rachid Guerraoui, Vasileios Trigonakis, Igor Zablotchi |
EuroSys | 2 |
| 2017 | Abstracting Multi-Core Topologies with MCTOPabstractPortability and efficiency are usually antagonists in multi-core computing. In order to develop efficient code, one needs to take into account the topology of the target multi-cores (e.g., for locality). This clearly hampers code portability. In this paper, we show that you can have the cake and eat it too. Georgios Chatzopoulos, Rachid Guerraoui, Tim Harris 0001, Vasileios Trigonakis |
EuroSys | 2 |
| 2017 | Capturing the Moment: Lightweight Similarity ComputationsabstractSimilarity computations are crucial in various web activities like advertisements, search or trust-distrust predictions. These similarities often vary with time as product perception and popularity constantly change with users' evolving inclination. The huge volume of user-generated data typically results in heavyweight computations for even a single similarity update. We present I-SIM, a novel similarity metric that enables lightweight similarity computations in an incremental and temporal manner. Incrementality enables updates with low latency whereas temporality captures users' evolving inclination. The main idea behind I-SIM is to disintegrate the similarity metric into mutually independent time-aware factors which can be updated incrementally. We illustrate the efficacy of I-SIM through a novel recommender (SWIFT) as well as through a trust-distrust predictor in Online Social Networks (I-TRUST). We experimentally show that I-SIM enables fast and accurate predictions in an energy-efficient manner. Georgios Damaskinos, Rachid Guerraoui, Rhicheek Patra |
ICDE | 2 |
| 2017 | When Neurons FailabstractNeural networks have been traditionally considered robust in the sense that their precision degrades gracefully with the failure of neurons and can be compensated by additional learning phases. Nevertheless, critical applications for which neural networks are now appealing solutions, cannot afford any additional learning at run-time. In this paper, we view a multilayer neural network as a distributed system of which neurons can fail independently, and we evaluate its robustness in the absence of any (recovery) learning phase. We give tight bounds on the number of neurons that can fail without harming the result of a computation. To determine our bounds, we leverage the fact that neural activation functions are Lipschitz-continuous. Our bound is given in the form of quantity, we call the Forward Error Propagation, computing this quantity only requires looking at the topology of the network, while experimentally assessing the robustness of a network requires the costly experiment of looking at all the possible inputs and testing all the possible configurations of the network corresponding to different failure situations, facing a discouraging combinatorial explosion. We distinguish the case of neurons that can fail and stop their activity (crashed neurons) from the case of neurons that can fail by transmitting arbitrary values (Byzantine neurons). In the crash case, our bound involves the number of neurons per layer, the Lipschitz constant of the neural activation function, the number of failing neurons, the synaptic weights and the depth of the layer where the failure occurred. In the case of Byzantine failures, our bound involves, in addition, the synaptic transmission capacity. Interestingly, as we show in the paper, our bound can easily be extended to the case where synapses can fail. We present three applications of our results. The first is a quantification of the effect of memory cost reduction on the accuracy of a neural network. The second is a quantification of the amount of information any neuron needs from its preceding layer, enabling thereby a boosting scheme that prevents neurons from waiting for unnecessary signals. Our third application is a quantification of the trade-off between neural networks robustness and learning cost. El Mahdi El Mhamdi, Rachid Guerraoui |
IPDPS | 2 |
| 2017 | Machine Learning with Adversaries: Byzantine Tolerant Gradient DescentabstractWe study the resilience to Byzantine failures of distributed implementations of Stochastic Gradient Descent (SGD). So far, distributed machine learning frameworks have largely ignored the possibility of failures, especially arbitrary (i.e., Byzantine) ones. Causes of failures include software bugs, network asynchrony, biases in local datasets, as well as attackers trying to compromise the entire system. Assuming a set of $n$ workers, up to $f$ being Byzantine, we ask how resilient can SGD be, without limiting the dimension, nor the size of the parameter space. We first show that no gradient aggregation rule based on a linear combination of the vectors proposed by the workers (i.e, current approaches) tolerates a single Byzantine failure. We then formulate a resilience property of the aggregation rule capturing the basic requirements to guarantee convergence despite $f$ Byzantine workers. We propose \emph{Krum}, an aggregation rule that satisfies our resilience property, which we argue is the first provably Byzantine-resilient algorithm for distributed SGD. We also report on experimental evaluations of Krum. Peva Blanchard, El Mahdi El Mhamdi, Rachid Guerraoui, Julien Stainer |
NIPS | 3 |
| 2017 | Dynamic Safe Interruptibility for Decentralized Multi-Agent Reinforcement LearningabstractIn reinforcement learning, agents learn by performing actions and observing their outcomes. Sometimes, it is desirable for a human operator to interrupt an agent in order to prevent dangerous situations from happening. Yet, as part of their learning process, agents may link these interruptions, that impact their reward, to specific states and deliberately avoid them. The situation is particularly challenging in a multi-agent context because agents might not only learn from their own past interruptions, but also from those of other agents. Orseau and Armstrong defined safe interruptibility for one learner, but their work does not naturally extend to multi-agent systems. This paper introduces dynamic safe interruptibility, an alternative definition more suited to decentralized learning problems, and studies this notion in two learning frameworks: joint action learners and independent learners. We give realistic sufficient conditions on the learning algorithm to enable dynamic safe interruptibility in the case of joint action learners, yet show that these conditions are not sufficient for independent learners. We show however that if agents can detect interruptions, it is possible to prune the observations to ensure dynamic safe interruptibility even for independent learners. El Mahdi El Mhamdi, Rachid Guerraoui, Hadrien Hendrikx, Alexandre Maurer |
NIPS | 2 |
| 2017 | Brief Announcement: Byzantine-Tolerant Machine LearningabstractWe report on Krum, the first provably Byzantine-tolerant aggregation rule for distributed Stochastic Gradient Descent (SGD). Krum guarantees the convergence of SGD even in a distributed setting where (asymptotically) up to half of the workers can be malicious adversaries trying to attack the learning system. Peva Blanchard, El Mahdi El Mhamdi, Rachid Guerraoui, Julien Stainer |
PODC | 3 |
| 2017 | How Fast can a Distributed Transaction Commit?abstractThe atomic commit problem lies at the heart of distributed database systems. The problem consists for a set of processes (database nodes) to agree on whether to commit or abort a transaction (agreement property). The commit decision can only be taken if all processes are initially willing to commit the transaction, and this decision must be taken if all processes are willing to commit and there is no failure (validity property). An atomic commit protocol is said to be non-blocking if every correct process (a database node that does not fail) eventually reaches a decision (commit or abort) even if there are failures elsewhere in the distributed database system (termination property). Rachid Guerraoui, Jingjing Wang 0007 |
PODS | 1 |
| 2017 | On verifying causal consistencyabstractCausal consistency is one of the most adopted consistency criteria for distributed implementations of data structures. It ensures that operations are executed at all sites according to their causal precedence. We address the issue of verifying automatically whether the executions of an implementation of a data structure are causally consistent. We consider two problems: (1) checking whether one single execution is causally consistent, which is relevant for developing testing and bug finding algorithms, and (2) verifying whether all the executions of an implementation are causally consistent. Ahmed Bouajjani, Constantin Enea, Rachid Guerraoui, Jad Hamza |
POPL | 3 |
| 2017 | The Utility and Privacy Effects of a ClickabstractRecommenders are becoming one of the main ways to navigate the Internet. They recommend appropriate items to users based on their clicks, i.e., likes, ratings, purchases, etc. These clicks are key to providing relevant recommendations and, in this sense, have a significant utility. Since clicks reflect the preferences of users, they also raise privacy concerns. At first glance, there seems to be an inherent trade-off between the utility and privacy effects of a click. Nevertheless, a closer look reveals that the situation is more subtle: some clicks do improve utility without compromising privacy, whereas others decrease utility while hampering privacy. Rachid Guerraoui, Anne-Marie Kermarrec, Mahsa Taziki |
SIGIR | 1 |
| 2017 | On the Smallest Grain of Salt to Get a Unique Identity
Peva Blanchard, Rachid Guerraoui |
SIROCCO | 2 |
| 2017 | On the Robustness of a Neural NetworkabstractWith the development of neural networks based machine learning and their usage in mission critical applications, voices are rising against the black box aspect of neural networks as it becomes crucial to understand their limits and capabilities. With the rise of neuromorphic hardware, it is even more critical to understand how a neural network, as a distributed system, tolerates the failures of its computing nodes, neurons, and its communication channels, synapses. Experimentally assessing the robustness of neural networks involves the quixotic venture of testing all the possible failures, on all the possible inputs, which ultimately hits a combinatorial explosion for the first, and the impossibility to gather all the possible inputs for the second.In this paper, we prove an upper bound on the expected error of the output when a subset of neurons crashes. This bound involves dependencies on the network parameters that can be seen as being too pessimistic in the average case. It involves a polynomial dependency on the Lipschitz coefficient of the neurons' activation function, and an exponential dependency on the depth of the layer where a failure occurs. We back up our theoretical results with experiments illustrating the extent to which our prediction matches the dependencies between the network parameters and robustness. Our results show that the robustness of neural networks to the average crash can be estimated without the need to neither test the network on all failure configurations, nor access the training set used to train the network, both of which are practically impossible requirements. El Mahdi El Mhamdi, Rachid Guerraoui, Sébastien Rouault |
SRDS | 2 |
| 2017 | TRIAD: Creating Synergies Between Memory, Disk and Log in Log Structured Key-Value Stores
Oana Balmau, Diego Didona, Rachid Guerraoui, Willy Zwaenepoel, Huapeng Yuan, Aashray Arora, Pavan Konka |
USENIX ATC | 3 |
| 2017 | Elastic transactions
Pascal Felber, Vincent Gramoli, Rachid Guerraoui |
J. Parallel Distributed Comput. | 3 |
| 2017 | Heterogeneous Recommendations: What You Might Like To Read After Watching InterstellarabstractRecommenders, as widely implemented nowadays by major e-commerce players like Netflix or Amazon, use collaborative filtering to suggest the most relevant items to their users. Clearly, the effectiveness of recommenders depends on the data they can exploit, i.e., the feedback of users conveying their preferences, typically based on their past ratings. As of today, most recommenders are homogeneous in the sense that they utilize one specific application at a time. In short, Alice will only get recommended a movie if she has been rating movies. But what if she has been only rating books and would like to get recommendations for a movie? Clearly, the multiplicity of web applications is calling for heterogeneous recommenders that could utilize ratings in one application to provide recommendations in another one. This paper presents X-M ap , a heterogeneous recommender. X-M ap leverages meta-paths between heterogeneous items over several application domains, based on users who rated across these domains. These meta-paths are then used in X-M ap to generate, for every user, a profile ( AlterEgo ) in a domain where the user might not have rated any item yet. Not surprisingly, leveraging meta-paths poses non-trivial issues of (a) meta-path-based inter-item similarity , in order to enable accurate predictions, (b) scalability , given the amount of computation required, as well as (c) privacy , given the need to aggregate information across multiple applications. We show in this paper how X-M ap addresses the above-mentioned issues to achieve accuracy, scalability and differential privacy. In short, X-M ap weights the meta-paths based on several factors to compute inter-item similarities, and ensures scalability through a layer-based pruning technique. X-M ap guarantees differential privacy using an exponential scheme that leverages the meta-path-based similarities while determining the probability of item selection to construct the AlterEgos. We present an exhaustive experimental evaluation of X-M ap using real traces from Amazon. We show that, in terms of accuracy, X-M ap outperforms alternative heterogeneous recommenders and, in terms of throughput, X-M ap achieves a linear speedup with an increasing number of machines. Rachid Guerraoui, Anne-Marie Kermarrec, Tao Lin 0004, Rhicheek Patra |
Proc. VLDB Endow. | 1 |
| 2017 | Operation-Level Wait-Free Transactional Memory with Support for Irrevocable OperationsabstractTransactional memory (TM) aims to be a general purpose concurrency mechanism. However, operations which cause side-effects cannot be easily managed by a TM system, in which transactions are executed optimistically. In particular, networking, I/O, and some system calls cannot be executed within a transaction that may abort and restart (e.g., due to conflicts). Thus, many TM systems let transactions become irrevocable, i.e., they are guaranteed to commit. Supporting this in TM is a challenge, but there exist fast and highly parallel TM systems that allow for irrevocable transactions. However, no such system so far provides guarantees that all transactional operations terminate in a finite time. In this paper, we show that support for irrevocable operations does not entail inherent waiting. We present a TM algorithm that guarantees wait-freedom for any transactional operation. The algorithm is based on the weakest synchronization primitive possible (test-and-set), and guarantees opacity and strong progressiveness. To experimentally evaluate the algorithm, we developed a proof-of-concept TM system and tested it using the STMBench7 benchmark. Jan Z. Konczak, Pawel T. Wojciechowski, Rachid Guerraoui |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2016 | ProteusTM: Abstraction Meets Performance in Transactional MemoryabstractThe Transactional Memory (TM) paradigm promises to greatly simplify the development of concurrent applications. This led, over the years, to the creation of a plethora of TM implementations delivering wide ranges of performance across workloads. Yet, no universal implementation fits each and every workload. In fact, the best TM in a given workload can reveal to be disastrous for another one. This forces developers to face the complex task of tuning TM implementations, which significantly hampers their wide adoption. In this paper, we address the challenge of automatically identifying the best TM implementation for a given workload. Our proposed system, ProteusTM, hides behind the TM interface a large library of implementations. Underneath, it leverages a novel multi-dimensional online optimization scheme, combining two popular learning techniques: Collaborative Filtering and Bayesian Optimization. Diego Didona, Nuno Diegues, Anne-Marie Kermarrec, Rachid Guerraoui, Ricardo Neves, Paolo Romano 0002 |
ASPLOS | 4 |
| 2016 | Frugal topology construction for stream aggregation in the cloudabstractAggregation of streamed data is key to the expansion of the Internet of Things. This paper addresses the problem of designing a topology for reliably aggregating data flows from many devices arriving at a datacenter. Reliability here means ensuring operation without data loss. We seek a frugal solution that prevents wasteful resource consumption (over-provisioning). This problem is salient when building an aggregation service out of components (here aggregation nodes) that exhibit hard constraints on the amount of information they can handle per unit of time. We first formalize the problem and provide an analysis of the relation between monitored devices (plus information they send), and the operations performed at aggregation nodes, in terms of data rates. Building on this rate analysis, we devise a novel algorithm, which we call CSA, that basically outputs an aggregation topology capable of handling those incoming data rates, preventing thereby empirical trial-and-error design. We analyze the algorithm, before validating it on the Amazon Kinesis platform, using a device dataset from a European telco operator. Rachid Guerraoui, Erwan Le Merrer, Rhicheek Patra, Bao Duy Tran |
INFOCOM | 1 |
| 2016 | Never Say Never - Probabilistic and Temporal Failure DetectorsabstractThe failure detector approach for solving distributed computing problems has been celebrated for its modularity. This approach allows the construction of algorithms using abstract failure detection mechanisms, defined by axiomatic properties, as building blocks. The minimal synchrony assumptions on communication, which enable to implement the failure detection mechanism, are studied separately. Such synchrony assumptions are typically expressed as eventual guarantees that need to hold, after some point in time, forever and deterministically. But in practice, they never do. Synchrony assumptions may hold only probabilistically and temporarily. In this paper, we study failure detectors in a realistic distributed system N, with asynchrony inflicted by probabilistic synchronous communication. We address the following paradox: an implementation of "consensus with probability 1" is possible in N without using randomness in the algorithm itself, while an implementation of "◇S with probability 1" is impossible to achieve in N (◇S being the weakest failure detector to solve the consensus problem and many equivalent problems). We circumvent this paradox by introducing a new failure detector ◇S*, a variant of ◇S with probabilistic and temporal accuracy. We prove that ◇S* is implementable in N and we provide an optimal ◇S* algorithm. Interestingly, we show that ◇S* can replace ◇S, in several existing deterministic consensus algorithms using ◇S, to yield an algorithm that solves "consensus with probability 1". In fact, we show that such result holds for all decisive problems (not only consensus) and also for failure detector ◇P (not only ◇S). The resulting algorithms combine the modularity of distributed computing practices with the practicality of networking ones. Dacfey Dzung, Rachid Guerraoui, David Kozhaya, Yvonne-Anne Pignolet |
IPDPS | 2 |
| 2016 | Locking Made Easy
Jelena Antic, Georgios Chatzopoulos, Rachid Guerraoui, Vasileios Trigonakis |
Middleware | 3 |
| 2016 | Atum: Scalable Group Communication Using Volatile Groups
Rachid Guerraoui, Anne-Marie Kermarrec, Matej Pavlovic, Dragos-Adrian Seredinschi |
Middleware | 1 |
| 2016 | Collision-Free Pattern Formation
Rachid Guerraoui, Alexandre Maurer |
OPODIS | 1 |
| 2016 | Incremental Consistency Guarantees for Replicated Objects
Rachid Guerraoui, Matej Pavlovic, Dragos-Adrian Seredinschi |
OSDI | 1 |
| 2016 | ESTIMA: extrapolating scalability of in-memory applicationsabstractThis paper presents ESTIMA, an easy-to-use tool for extrapolating the scalability of in-memory applications. ESTIMA is designed to perform a simple, yet important task: given the performance of an application on a small machine with a handful of cores, ESTIMA extrapolates its scalability to a larger machine with more cores, while requiring minimum input from the user. The key idea underlying ESTIMA is the use of stalled cycles (e.g. cycles that the processor spends waiting for various events, such as cache misses or waiting on a lock). ESTIMA measures stalled cycles on a few cores and extrapolates them to more cores, estimating the amount of waiting in the system. ESTIMA can be effectively used to predict the scalability of in-memory applications. For instance, using measurements of memcached and SQLite on a desktop machine, we obtain accurate predictions of their scalability on a server. Our extensive evaluation on a large number of in-memory benchmarks shows that ESTIMA has generally low prediction errors. Georgios Chatzopoulos, Aleksandar Dragojevic, Rachid Guerraoui |
PPoPP | 3 |
| 2016 | Optimistic concurrency with OPTIKabstractWe introduce OPTIK, a new practical design pattern for designing and implementing fast and scalable concurrent data structures. OPTIK relies on the commonly-used technique of version numbers for detecting conflicting concurrent operations. We show how to implement the OPTIK pattern using the novel concept of OPTIK locks. These locks enable the use of version numbers for implementing very efficient optimistic concurrent data structures. Existing state-of-the-art lock-based data structures acquire the lock and then check for conflicts. In contrast, with OPTIK locks, we merge the lock acquisition with the detection of conflicting concurrency in a single atomic step, similarly to lock-free algorithms. We illustrate the power of our OPTIK pattern and its implementation by introducing four new algorithms and by optimizing four state-of-the-art algorithms for linked lists, skip lists, hash tables, and queues. Our results show that concurrent data structures built using OPTIK are more scalable than the state of the art. Rachid Guerraoui, Vasileios Trigonakis |
PPoPP | 1 |
| 2016 | Right on Time Distributed Shared MemoryabstractThe demand for real-time data storage in distributed control systems (DCSs) is growing. Yet, providing real-time DCS guarantees is challenging, especially when more and more sensor and actuator devices are connected to industrial plants and message loss needs to be taken into account.In this paper, we investigate how to build a shared memory abstraction for DCSs as a first step towards implementing different shared storage systems in a DCS context. We first prove that, in the presence of host crashes and message losses, the necessary guarantees of such an abstraction are impossible to implement using a traditional approach that has no access to the internals of existing DCS services, e.g., a modular approach where algorithms are built on top of existing software blocks like failure detectors. We propose a white-box approach that utilizes messages of existing services in any DCS as the sole means of communication. More precisely, we present TapeWorm, an algorithm that attaches itself to the heartbeat messages of the failure detector component in DCSs. We prove that TapeWorm implements the desired shared memory guarantees for applications running on a DCS. We also analyze the performance of TapeWorm and we showcase ways of adapting TapeWorm to various application needs and workloads. Rachid Guerraoui, David Kozhaya, Yvonne-Anne Pignolet |
RTSS | 1 |
| 2016 | Fast and Robust Memory Reclamation for Concurrent Data StructuresabstractIn concurrent systems without automatic garbage collection, it is challenging to determine when it is safe to reclaim memory, especially for lock-free data structures. Existing concurrent memory reclamation schemes are either fast but do not tolerate process delays, robust to delays but with high overhead, or both robust and fast but narrowly applicable. This paper proposes QSense, a novel concurrent memory reclamation technique. QSense is a hybrid technique with a fast path and a fallback path. In the common case (without process delays), a high-performing memory reclamation scheme is used (fast path). If process delays block memory reclamation through the fast path, a robust fallback path is used to guarantee progress. The fallback path uses hazard pointers, but avoids their notorious need for frequent and expensive memory fences. Oana Balmau, Rachid Guerraoui, Maurice Herlihy, Igor Zablotchi |
SPAA | 2 |
| 2016 | Concurrent Search Data Structures Can Be Blocking and Practically Wait-FreeabstractWe argue that there is virtually no practical situation in which one should seek a "theoretically wait-free" algorithm at the expense of a state-of-the-art blocking algorithm in the case of search data structures: blocking algorithms are simple, fast, and can be made "practically wait-free". We draw this conclusion based on the most exhaustive study of blocking search data structures to date. We consider (a) different search data structures of different sizes, (b) numerous uniform and non-uniform workloads, representative of a wide range of practical scenarios, with different percentages of update operations, (c) with and without delayed threads, (d) on different hardware technologies, including processors providing HTM instructions. We explain our claim that blocking search data structures are practically wait-free through an analogy with the birthday paradox, revealing that, in state-of-the-art algorithms implementing such data structures, the probability of conflicts is extremely small. When conflicts occur as a result of context switches and interrupts, we show that HTM-based locks enable blocking algorithms to cope with them. Tudor David, Rachid Guerraoui |
SPAA | 2 |
| 2016 | Who's On Board?: Probabilistic Membership for Real-Time Distributed Control SystemsabstractTo increase their dependability, distributed control systems (DCSs) need to agree in real time about which hosts have crashed, i.e., they need a real-time membership service. In this paper, we prove that such a service cannot be implemented deterministically if, besides host crashes, communication can also fail. We define implementable probabilistic variants of membership properties, which constitute what we call a synchronous membership service (SYMS). We present an algorithm, ViewSnoop, that implements SYMS with high-probability. We implement, deploy and evaluate ViewSnoop analytically as well as experimentally, within an industrial DCS framework. We show that ViewSnoop significantly improves the dependability of DCSs compared to membership schemes based on classic heartbeats, at low additional cost. Moreover, ViewSnoop distinguishes, with high probability, host crashes from message losses, enabling DCSs to counteract losses better than existing approaches. Rachid Guerraoui, David Kozhaya, Manuel Oriol, Yvonne-Anne Pignolet |
SRDS | 1 |
| 2016 | Unlocking Energy
Babak Falsafi, Rachid Guerraoui, Javier Picorel, Vasileios Trigonakis |
USENIX ATC | 2 |
| 2016 | Optimal Fair Computation
Rachid Guerraoui, Jingjing Wang 0007 |
DISC | 1 |
| 2015 | Asynchronized Concurrency: The Secret to Scaling Concurrent Search Data StructuresabstractWe introduce "asynchronized concurrency (ASCY)," a paradigm consisting of four complementary programming patterns. ASCY calls for the design of concurrent search data structures (CSDSs) to resemble that of their sequential counterparts. We argue that ASCY leads to implementations which are portably scalable: they scale across different types of hardware platforms, including single and multi-socket ones, for various classes of workloads, such as read-only and read-write, and according to different performance metrics, including throughput, latency, and energy. We substantiate our thesis through the most exhaustive evaluation of CSDSs to date, involving 6 platforms, 22 state-of-the-art CSDS algorithms, 10 re-engineered state-of-the-art CSDS algorithms following the ASCY patterns, and 2 new CSDS algorithms designed with ASCY in mind. We observe up to 30% improvements in throughput in the re-engineered algorithms, while our new algorithms out-perform the state-of-the-art alternatives. Tudor David, Rachid Guerraoui, Vasileios Trigonakis |
ASPLOS | 2 |
| 2015 | Hide & Share: Landmark-Based Similarity for Private KNN ComputationabstractComputing k-nearest-neighbor graphs constitutes a fundamental operation in a variety of data-mining applications. As a prominent example, user-based collaborative-filtering provides recommendations by identifying the items appreciated by the closest neighbors of a target user. As this kind of applications evolve, they will require KNN algorithms to operate on more and more sensitive data. This has prompted researchers to propose decentralized peer-to-peer KNN solutions that avoid concentrating all information in the hands of one central organization. Unfortunately, such decentralized solutions remain vulnerable to malicious peers that attempt to collect and exploit information on participating users. In this paper, we seek to overcome this limitation by proposing H&S (Hide & Share), a novel landmark-based similarity mechanism for decentralized KNN computation. Landmarks allow users (and the associated peers) to estimate how close they lay to one another without disclosing their individual profiles. We evaluate H&S in the context of a user-based collaborative-filtering recommender with publicly available traces from existing recommendation systems. We show that although landmark-based similarity does disturb similarity values (to ensure privacy), the quality of the recommendations is not as significantly hampered. We also show that the mere fact of disturbing similarity values turns out to be an asset because it prevents a malicious user from performing a profile reconstruction attack against other users, thus reinforcing users' privacy. Finally, we provide a formal privacy guarantee by computing an upper bound on the amount of information revealed by H&S about a user's profile. Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Antoine Rault, François Taïani, Jingjing Wang 0007 |
DSN | 2 |
| 2015 | Making BFT Protocols Really AdaptiveabstractMany state-machine Byzantine Fault Tolerant (BFT) protocols have been introduced so far. Each protocol addressed a different subset of conditions and use-cases. However, if the underlying conditions of a service span different subsets, choosing a single protocol will likely not be a best fit. This yields robustness and performance issues which may be even worse in services that exhibit fluctuating conditions and workloads. In this paper, we reconcile existing state-machine BFT protocols in a single adaptive BFT system, called ADAPT, aiming at covering a larger set of conditions and use-cases, probably the union of individual subsets of these protocols. At anytime, a launched protocol in ADAPT can be aborted and replaced by another protocol according to a potential change (an event) in the underlying system conditions. The launched protocol is chosen according to an "evaluation process" that takes into consideration both: protocol characteristics and its performance. This is achieved by applying some mathematical formulas that match the profiles of protocols to given user (e.g., service owner) preferences. ADAPT can assess the profiles of protocols (e.g., throughput) at run-time using Machine Learning prediction mechanisms to get accurate evaluations. We compare ADAPT with well known BFT protocols showing that it outperforms others as system conditions change and under dynamic workloads. Jean Paul Bahsoun, Rachid Guerraoui, Ali Shoker |
IPDPS | 2 |
| 2015 | The Benefits of Entropy in Population ProtocolsabstractA distributed computing system can be viewed as the result of the interplay between a distributed algorithm specifying the effects of a local event (e.g. reception of a message), and an adversary choosing the interleaving (schedule) of these events in the execution. In the context of large networks of mobile pairwise interacting agents (population protocols), the adversary models the mobility of the agents by choosing the successive pairs of interacting agents. For some problems, assuming that the adversary selects the schedule according to some probability distribution greatly helps to devise (almost) correct solutions. But how much randomness is really necessary? To what extent does a problem admit implementations that are robust against a "not so random" schedule? This paper takes a first step in addressing this question by borrowing the concept of T-randomness, 0 <= T <= 1, from algorithmic information theory. Roughly speaking, the value T fixes the entropy rate of the considered schedules. For instance, the case T = 1 corresponds, in a specific sense, to schedules in which the pairs of interacting agents are chosen independently and uniformly (perfect randomness). The holy grail question can then be precisely stated as determining the optimal entropy rate to solve a given problem. We first show that perfect randomness is never required. Precisely, if a finite-state algorithm solves a problem with 1-randomness, then this algorithm still solves the same problem with T-randomness for some T < 1. Second, we illustrate how to compute bounds on the optimal entropy rate of a specific problem, namely the leader election problem. Joffroy Beauquier, Peva Blanchard, Janna Burman, Rachid Guerraoui |
OPODIS | 4 |
| 2015 | Safety-Liveness Exclusion in Distributed ComputingabstractThe history of distributed computing is full of trade-offs between safety and liveness. For instance, one of the most celebrated results in the field, namely the impossibility of consensus in an asynchronous system basically says that we cannot devise an algorithm that deterministically ensures consensus agreement and validity (i.e., safety) on the one hand, and consensus wait-freedom (i.e., liveness) on the other hand. The motivation of this work is to study the extent to which safety and liveness properties inherently exclude each other. More specifically, we ask, given any safety property S, whether we can determine the strongest (resp. weakest) liveness property that can (resp. cannot) be achieved with S. We show that, maybe surprisingly, the answers to these safety-liveness exclusion questions are in general negative. This has several ramifications in various distributed computing contexts. In the context of consensus for example, this means that it is impossible to determine the strongest (resp. the weakest) liveness property that can (resp. cannot) be ensured with linearizability. Victor Bushkov, Rachid Guerraoui |
PODC | 2 |
| 2015 | The Weakest Failure Detector for Eventual ConsistencyabstractIn its classical form, a consistent replicated service requires all replicas to witness the same evolution of the service state. Assuming a message-passing environment with a majority of correct processes, the necessary and sufficient information about failures for implementing a general state machine replication scheme ensuring consistency is captured by the Ω failure detector. Swan Dubois, Rachid Guerraoui, Petr Kuznetsov, Franck Petit, Pierre Sens 0001 |
PODC | 2 |
| 2015 | To Transmit Now or Not to Transmit NowabstractGiven an unreliable communication link, this paper studies how to build, in an energy-efficient manner, a reliable communication service that is synchronous with high probability. We consider a Partially Observable Markov Decision Process (POMDP) setting in which a communication link's transmission quality: (i) changes according to a classic Markovian model and (ii) can be only partially observed, through feedback relative to previous transmissions. We perform a thorough analysis under several variations of Ack/Nack feedback mechanisms. Despite the general intractability of POMDPs, we prove that our communication service, under reliable feedback, can be inexpensively implemented. We obtain closed form solutions specifying when to transmit over the link, which allows to derive an energy-optimal implementation. We also analyse the impact of lossy feedback on implementing our communication service. Considering multiple lossy feedback mechanisms, we show that an easily implementable structure for our communication service can also be obtained, depending on the feedback mechanism itself. Dacfey Dzung, Rachid Guerraoui, David Kozhaya, Yvonne-Anne Pignolet |
SRDS | 2 |
| 2015 | Privacy-Conscious Information Diffusion in Social NetworksabstractWe present Riposte , a distributed algorithm for disseminating information (ideas, news, opinions, or trends) in a social network. Riposte ensures that information spreads widely if and only if a large fraction of users find it interesting, and this is done in a “privacy-conscious” manner, namely without revealing the opinion of any individual user. Whenever an information item is received by a user, Riposte decides to either forward the item to all the user’s neighbors, or not to forward it to anyone. The decision is randomized and is based on the user’s (private) opinion on the item, as well as on an upper bound s on the number of user’s neighbors that have not received the item yet. In short, if the user likes the item, Riposte forwards it with probability slightly larger than 1 / s , and if not, the item is forwarded with probability slightly smaller than 1 / s . Using a comparison to branching processes, we show for a general family of random directed graphs with arbitrary out-degree sequences, that if the information item appeals to a sufficiently large (constant) fraction of users, then the item spreads to a constant fraction of the network; while if fewer users like it, the dissemination process dies out quickly. In addition, we provide extensive experimental evaluation of Riposte on topologies taken from online social networks, including Twitter and Facebook. These keywords were added by machine and not by the authors. This process is experimental and the keywords may be updated as the learning algorithm improves. George Giakkoupis, Rachid Guerraoui, Arnaud Jégou, Anne-Marie Kermarrec, Nupur Mittal |
DISC | 2 |
| 2015 | Byzantine Fireflies
Rachid Guerraoui, Alexandre Maurer |
DISC | 1 |
| 2015 | D2P: Distance-Based Differential Privacy in RecommendersabstractThe upsurge in the number of web users over the last two decades has resulted in a significant growth of online information. This information growth calls for recommenders that personalize the information proposed to each individual user. Nevertheless, personalization also opens major privacy concerns. This paper presents D 2 P , a novel protocol that ensures a strong form of differential privacy, which we call distance-based differential privacy, and which is particularly well suited to recommenders. D 2 P avoids revealing exact user profiles by creating altered profiles where each item is replaced with another one at some distance. We evaluate D 2 P analytically and experimentally on MovieLens and Jester datasets and compare it with other private and non-private recommenders. Rachid Guerraoui, Anne-Marie Kermarrec, Rhicheek Patra, Mahsa Taziki |
Proc. VLDB Endow. | 1 |
| 2015 | The Next 700 BFT ProtocolsabstractWe present Abstract (ABortable STate mAChine replicaTion), a new abstraction for designing and reconfiguring generalized replicated state machines that are, unlike traditional state machines, allowed to abort executing a client’s request if “something goes wrong.” Abstract can be used to considerably simplify the incremental development of efficient Byzantine fault-tolerant state machine replication ( BFT ) protocols that are notorious for being difficult to develop. In short, we treat a BFT protocol as a composition of Abstract instances. Each instance is developed and analyzed independently and optimized for specific system conditions. We illustrate the power of Abstract through several interesting examples. We first show how Abstract can yield benefits of a state-of-the-art BFT protocol in a less painful and error-prone manner. Namely, we develop AZyzzyva , a new protocol that mimics the celebrated best-case behavior of Zyzzyva using less than 35% of the Zyzzyva code. To cover worst-case situations, our abstraction enables one to use in AZyzzyva any existing BFT protocol. We then present Aliph , a new BFT protocol that outperforms previous BFT protocols in terms of both latency (by up to 360%) and throughput (by up to 30%). Finally, we present R-Aliph , an implementation of Aliph that is robust , that is, whose performance degrades gracefully in the presence of Byzantine replicas and Byzantine clients. Pierre-Louis Aublin, Rachid Guerraoui, Nikola Knezevic, Vivien Quéma, Marko Vukolic |
ACM Trans. Comput. Syst. | 2 |
| 2014 | Finding trojan message vulnerabilities in distributed systemsabstractTrojan messages are messages that seem correct to the receiver but cannot be generated by any correct sender. Such messages constitute major vulnerability points of a distributed system---they constitute ideal targets for a malicious actor and facilitate failure propagation across nodes. We describe Achilles, a tool that searches for Trojan messages in a distributed system. Achilles uses dynamic white-box analysis on the distributed system binaries in order to infer the predicate that defines messages parsed by receiver nodes and generated by sender nodes, respectively, and then computes Trojan messages as the difference between the two. Radu Banabic, George Candea, Rachid Guerraoui |
ASPLOS | 3 |
| 2014 | Reusable Concurrent Data Types
Vincent Gramoli, Rachid Guerraoui |
ECOOP | 2 |
| 2014 | HyRec: leveraging browsers for scalable recommendersabstractThe ever-growing amount of data available on the Internet calls for personalization. Yet, the most effective personalization schemes, such as those based on collaborative filtering (CF), are notoriously resource greedy. This paper presents HyRec, an online cost-effective scalable system for user-based CF personalization. HyRec offloads recommendation tasks onto the web browsers of users, while a server orchestrates the process and manages the relationships between user profiles. Antoine Boutet, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Rhicheek Patra |
Middleware | 3 |
| 2014 | Consensus insideabstractScaling to a large number of cores with non-uniform communication latency and unpredictable response time may call for viewing a modern many-core architecture as a distributed system. In this view, the cores replicate shared data and ensure consistency among replicas through a message-passing based agreement protocol. Tudor David, Rachid Guerraoui, Maysam Yabandeh |
Middleware | 2 |
| 2014 | A paradox of eventual linearizability in shared memoryabstractThis paper compares, for the first time, the computational power of linearizable objects with that of eventually linearizable ones. We present the following paradox. We show that, unsurprisingly, no set of eventually linearizable objects can (1) implement any non-trivial linearizable object, nor (2) boost the consensus power of simple objects like linearizable registers. We also show, perhaps surprisingly, that any implementation of an eventually linearizable complex object like a fetch&increment counter (from linearizable base objects), can itself be viewed as a fully linearizable implementation of the same fetch&increment counter (using the exact same set of base objects). Rachid Guerraoui, Eric Ruppert |
PODC | 1 |
| 2014 | The PCL theorem: transactions cannot be parallel, consistent and liveabstractWe show that it is impossible to design a transactional memory system which ensures parallelism, i.e. transactions do not need to synchronize unless they access the same application objects, while ensuring very little consistency, i.e. a consistency condition, called weak adaptive consistency, introduced here and which is weaker than snapshot isolation, processor consistency, and any other consistency condition stronger than them (such as opacity, serializability, causal serializability, etc.), and very little liveness, i.e. that transactions eventually commit if they run solo. Victor Bushkov, Dmytro Dziuma, Panagiota Fatourou, Rachid Guerraoui |
SPAA | 4 |
| 2014 | Tracking freeriders in gossip-based content dissemination systemsabstractGossip-based protocols have proven very efficient for disseminating high-bandwidth content such as video streams in a peer-to-peer fashion. However, for the protocols to work, nodes are required to collaborate by devoting a fraction of their upload bandwidth, a scarce resource for some of them, to forward the content they receive to other nodes. Consequently, such protocols suffer from freeriding, a common phenomenon on the Internet, which consists in selfishly benefiting from the system without contributing its fair share. Due to the dynamic nature and the inherent randomness of gossip protocols and to the high scalability requirements of video streaming systems, detecting freeriders is a difficult challenge. This paper presents LiFTinG, the first protocol for detecting freeriders, including colluding ones, in gossip-based content dissemination systems with asymmetric data exchanges. In addition, LiFTinG is still able to detect freeriders when network coding, a widely used technique to improve the efficiency of content dissemination, is used. LiFTinG relies on nodes to track abnormal behavior by cross-checking the history of their previous interactions and exploits the fact that nodes pick neighbors at random to prevent colluding nodes from mutually covering up their bad actions. We present a methodology for setting the parameters of LiFTinG to their optimal value, based on a theoretical analysis and we quantify theoretically the performance of LiFTinG. We derive, based on simulations, the optimal strategy of freeriders by taking into account, through a utility function, the benefit of freeriding and the probability of being detected. In addition to these simulations, we report on the deployment of LiFTinG on PlanetLab. In a 300-node system, where a stream of 674 kbps is broadcasted, LiFTinG incurs a maximum overhead of only 8% and provides good detection results: For instance, with 10% of freeriders decreasing their contribution by up to 30%, LiFTinG detects 86% of the freeriders after only 30 s and wrongfully expels only a few honest nodes (most of them actually being buggy). Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod, Swagatika Prusty, Aline Roumy |
Comput. Networks | 1 |
| 2014 | Computing in social networks
Andrei Giurgiu, Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec |
Inf. Comput. | 2 |
| 2014 | Tight Bounds for Asynchronous RenamingabstractThis article presents the first tight bounds on the time complexity of shared-memory renaming, a fundamental problem in distributed computing in which a set of processes need to pick distinct identifiers from a small namespace. We first prove an individual lower bound of Ω( k ) process steps for deterministic renaming into any namespace of size subexponential in k , where k is the number of participants. The bound is tight: it draws an exponential separation between deterministic and randomized solutions, and implies new tight bounds for deterministic concurrent fetch-and-increment counters, queues, and stacks. The proof is based on a new reduction from renaming to another fundamental problem in distributed computing: mutual exclusion. We complement this individual bound with a global lower bound of Ω( k log ( k / c )) on the total step complexity of renaming into a namespace of size ck , for any c ≥ 1. This result applies to randomized algorithms against a strong adversary, and helps derive new global lower bounds for randomized approximate counter implementations, that are tight within logarithmic factors. On the algorithmic side, we give a protocol that transforms any sorting network into a randomized strong adaptive renaming algorithm, with expected cost equal to the depth of the sorting network. This gives a tight adaptive renaming algorithm with expected step complexity O (log k ), where k is the contention in the current execution. This algorithm is the first to achieve sublinear time, and it is time-optimal as per our randomized lower bound. Finally, we use this renaming protocol to build monotone-consistent counters with logarithmic step complexity and linearizable fetch-and-increment registers with polylogarithmic cost. Dan Alistarh, James Aspnes, Keren Censor-Hillel, Seth Gilbert, Rachid Guerraoui |
J. ACM | 5 |
| 2014 | Personalizing Top-k Processing Online in a Peer-to-Peer Social Tagging NetworkabstractThe rapidly increasing amount of user-generated content in social tagging systems provides a huge source of information. Yet, performing effective search in these systems is very challenging, especially when we seek the most appropriate items that match a potentially ambiguous query. Collaborative filtering-based personalization is appealing in this context, as it limits the search within a small network of participants with similar preferences. Offline personalization, which consists in maintaining, for every user, a network of similar participants based on their tagging behaviors, is effective for queries that are close to the querying user’s tagging profile but performs poorly when the queries, reflecting emerging interests, have little correlation with the querying user’s profile. We present P 2 TK 2 , the first protocol to personalize query processing in social tagging systems online. P 2 TK 2 is completely decentralized, and this design choice stems from the observation that the evolving social tagging systems naturally resemble P2P systems where users are both producers and consumers. This design exploits the power of the crowd and prevents any central authority from controlling personal information. P 2 TK 2 is gossip-based and probabilistic. It dynamically associates each user with social acquaintances sharing similar tagging behaviors. Appropriate users for answering a query are discovered at query time with the help of social acquaintances. This is achieved according to the hybrid interest of the querying user, taking into account both her tagging behavior and her query. Results are iteratively refined and returned to the querying user. We evaluate P 2 TK 2 on CiteULike and Delicious traces involving up to 50,000 users. We highlight the advantages of online personalization compared to offline personalization, as well as its efficiency, scalability, and inherent ability to cope with user departure and interest evolution in P2P systems. Xiao Bai 0002, Rachid Guerraoui, Anne-Marie Kermarrec |
ACM Trans. Internet Techn. | 2 |
| 2013 | WHATSUP: A Decentralized Instant News RecommenderabstractWe present WHATSUP, a collaborative filtering system for disseminating news items in a large-scale dynamic setting with no central authority. WHATSUP constructs an implicit social network based on user profiles that express the opinions of users about the news items they receive (like-dislike). Users with similar tastes are clustered using a similarity metric reflecting long-standing and emerging (dis)interests. News items are disseminated through a novel heterogeneous gossip protocol that (1) biases the orientation of its targets towards those with similar interests, and (2) amplifies dissemination based on the level of interest in every news item. We report on an extensive evaluation of WHATSUP through (a) simulations, (b) a ModelNet emulation on a cluster, and (c) a PlanetLab deployment based on real datasets. We show that WHATSUP outperforms various alternatives in terms of accurate and complete delivery of relevant news items while preserving the fundamental advantages of standard gossip: namely, simplicity of deployment and robustness. Antoine Boutet, Davide Frey, Rachid Guerraoui, Arnaud Jégou, Anne-Marie Kermarrec |
IPDPS | 3 |
| 2013 | Composing Relaxed TransactionsabstractAs the classic transactional abstraction is sometimes considered too restrictive in leveraging parallelism, a lot of work has been devoted to devising relaxed transactional models with the goal of improving concurrency. Nevertheless, the quest for improving concurrency has somehow led to neglect one of the most appealing aspects of transactions: software composition, namely, the ability to develop pieces of software independently and compose them into applications that behave correctly in the face of concurrency. Indeed, a closer look at relaxed transactional models reveals that they do jeopardize composition, raising the fundamental question whether it is at all possible to devise such models while preserving composition. This paper shows that the answer is positive. We present outheritance, a necessary and sufficient condition for a (potentially relaxed) transactional memory to support composition. Basically, outheritance requires child transactions to pass their conflict information to their parent transaction, which in turn maintains this information until commit time. Concrete instantiations of this idea have been used before, classic transactions being the most prevalent example, but we believe to be the first to capture this as a general principle as well as to prove that it is, strictly speaking, equivalent to ensuring composition. We illustrate the benefits of outheritance using elastic transactions and show how they can satisfy outheritance and provide composition without hampering concurrency. We leverage this to present a new (transactional) Java package, a composable alternative to the concurrency package of the JDK, and evaluate efficiency through an implementation that speeds up state of the art software transactional memory implementations (TL2, LSA, SwissTM) by almost a factor of 3. Vincent Gramoli, Rachid Guerraoui, Mihai Letia |
IPDPS | 2 |
| 2013 | Fast byzantine agreementabstractThis paper presents the first probabilistic Byzantine Agreement algorithm whose communication and time complexities are poly-logarithmic. So far, the most effective probabilistic Byzantine Agreement algorithm had communication complexity Õ(√n) and time complexity Õ(1). Nicolas Braud-Santoni, Rachid Guerraoui, Florian Huc |
PODC | 2 |
| 2013 | Introducing speculation in self-stabilization: an application to mutual exclusionabstractSelf-stabilization ensures that, after any transient fault, the system recovers in a finite time and eventually exhibits correct behavior. Speculation consists in guaranteeing that the system satisfies its requirements for any execution but exhibits significantly better performances for a subset of executions that are more probable. A speculative protocol is in this sense supposed to be both robust and efficient in practice. Swan Dubois, Rachid Guerraoui |
PODC | 2 |
| 2013 | Highly dynamic distributed computing with byzantine failuresabstractThis paper shows for the first time that distributed computing can be both reliable and efficient in an environment that is both highly dynamic and hostile. More specifically, we show how to maintain clusters of size O(log N), each containing more than two thirds of honest nodes with high probability, within a system whose size can vary polynomially with respect to its initial size. Furthermore, the communication cost induced by each node arrival or departure is polylogarithmic with respect to N, the maximal size of the system. Our clustering can be achieved despite the presence of a Byzantine adversary controlling a fraction τ ≤ 1⁄3-ε of the nodes, for some fixed constant ε > 0, independent of N. So far, such a clustering could only be performed for systems whose size can vary constantly and it was not clear whether that was at all possible for polynomial variances. Rachid Guerraoui, Florian Huc, Anne-Marie Kermarrec |
PODC | 1 |
| 2013 | Everything you always wanted to know about synchronization but were afraid to askabstractThis paper presents the most exhaustive study of synchronization to date. We span multiple layers, from hardware cache-coherence protocols up to high-level concurrent software. We do so on different types of architectures, from single-socket -- uniform and non-uniform -- to multi-socket -- directory and broadcast-based -- many-cores. We draw a set of observations that, roughly speaking, imply that scalability of synchronization is mainly a property of the hardware. Tudor David, Rachid Guerraoui, Vasileios Trigonakis |
SOSP | 2 |
| 2013 | A Distributed Polling with Probabilistic PrivacyabstractIn this paper, we present PDP, a distributed polling protocol that enables a set of participants to gather their opinion on a common interest without revealing their point of view. PDP does not rely on any centralized authority or on heavyweight cryptography. PDP is an overlay-based protocol where a subset of participants may use a simple sharing scheme to express their votes. In a system of M participants arranged in groups of size N where at least 2k-1 participants are honest, PDP bounds the probability for a given participant to have its vote recovered with certainty by a coalition of B dishonest participants by π(B/N)(k+1), where π is the proportion of participants splitting their votes, and k a privacy parameter. PDP bounds the impact of dishonest participants on the global outcome by 2(kα + BN), where represents the number of dishonest nodes using the sharing scheme. Yahya Benkaouz, Rachid Guerraoui, Mohammed Erradi, Florian Huc |
SRDS | 2 |
| 2013 | Byzantine agreement with homonyms
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Anne-Marie Kermarrec, Eric Ruppert, Hung Tran-The |
Distributed Comput. | 3 |
| 2013 | Asynchronous gossipabstractWe study the complexity of gossip in an asynchronous, message-passing fault-prone distributed system. We show that an adaptive adversary can significantly hamper the spreading of a rumor, while an oblivious adversary cannot. The algorithmic techniques proposed in this article can be used for improving the message complexity of distributed algorithms that rely on an all-to-all message exchange paradigm and are designed for an asynchronous environment. As an example, we show how to improve the message complexity of asynchronous randomized consensus. Chryssis Georgiou, Seth Gilbert, Rachid Guerraoui, Dariusz R. Kowalski |
J. ACM | 3 |
| 2012 | TM2C: a software transactional memory for many-coresabstractTransactional memory is an appealing paradigm for concurrent programming. Many software implementations of the paradigm were proposed in the last decades for both shared memory multi-core systems and clusters of distributed machines. However, chip manufacturers have started producing many-core architectures, with low network-on-chip communication latency and limited support for cache-coherence, rendering existing transactional memory implementations inapplicable. Vincent Gramoli, Rachid Guerraoui, Vasileios Trigonakis |
EuroSys | 2 |
| 2012 | How to Allocate Tasks AsynchronouslyabstractAsynchronous task allocation is a fundamental problem in distributed computing in which p asynchronous processes must execute a set of m tasks. Also known as write-all or do-all, this problem been studied extensively, both independently and as a key building block for various distributed algorithms. In this paper, we break new ground on this classic problem: we introduce the To-Do Tree concurrent data structure, which improves on the best known randomized and deterministic upper bounds. In the presence of an adaptive adversary, the randomized To-Do Tree algorithm has O(m+p log p log2m) work complexity. We then show that there exists a deterministic variant of the To-Do Tree algorithm with work complexity O(m+p log5m log2max(m, p)). For all values of m and p, our algorithms are within log factors of the O(m + p log p) lower bound for this problem. The key technical ingredient in our results is a new approach for analyzing concurrent executions against a strong adaptive scheduler. This technique allows us to handle the complex dependencies between the processes' coin flips and their scheduling, and to tightly bound the work needed to perform subsets of the tasks. Dan Alistarh, Michael A. Bender, Seth Gilbert, Rachid Guerraoui |
FOCS | 4 |
| 2012 | Unifying Thread-Level Speculation and Transactional Memory
João Barreto 0001, Aleksandar Dragojevic, Paulo Ferreira 0001, Ricardo Filipe, Rachid Guerraoui |
Middleware | 5 |
| 2012 | Speculative linearizabilityabstractLinearizability is a key design methodology for reasoning about implementations of concurrent abstract data types in both shared memory and message passing systems. It provides the illusion that operations execute sequentially and fault-free, despite the asynchrony and faults inherent to a concurrent system, especially a distributed one. A key property of linearizability is inter-object composability: a system composed of linearizable objects is itself linearizable. However, devising linearizable objects is very difficult, requiring complex algorithms to work correctly under general circumstances, and often resulting in bad average-case behavior. Concurrent algorithm designers therefore resort to speculation: optimizing algorithms to handle common scenarios more efficiently. The outcome are even more complex protocols, for which it is no longer tractable to prove their correctness. Rachid Guerraoui, Viktor Kuncak, Giuliano Losa |
PLDI | 1 |
| 2012 | On the liveness of transactional memoryabstractDespite the large amount of work on Transactional Memory (TM), little is known about how much liveness it could provide. This paper presents the first formal treatment of the question. We prove that no TM implementation can ensure local progress, the analogous of wait-freedom in the TM context, and we highlight different ways to circumvent the impossibility. Victor Bushkov, Rachid Guerraoui, Michal Kapalka |
PODC | 2 |
| 2012 | Early Deciding Synchronous Renaming in O( logf ) Rounds or Less
Dan Alistarh, Hagit Attiya, Rachid Guerraoui, Corentin Travers |
SIROCCO | 3 |
| 2012 | On the cost of composing shared-memory algorithmsabstractDecades of research in distributed computing have led to a variety of perspectives on what it means for a concurrent algorithm to be efficient, depending on model assumptions, progress guarantees, and complexity metrics. It is therefore natural to ask whether one could compose algorithms that perform efficiently under different conditions, so that the composition preserves the performance of the original components when their conditions are met. Dan Alistarh, Rachid Guerraoui, Petr Kuznetsov, Giuliano Losa |
SPAA | 2 |
| 2012 | Scalable and Secure Polling in Dynamic Distributed NetworksabstractWe consider the problem of securely conducting a poll in synchronous dynamic networks equipped with a Public Key Infrastructure (PKI). Whereas previous distributed solutions had a communication cost of O(n2) in an n nodes system, we present SPP (Secure and Private Polling), the first distributed polling protocol requiring only a communication complexity of O(n log3n), which we prove is near-optimal. Our protocol ensures perfect security against a computationally-bounded adversary, tolerates (1/2 - ϵ)n Byzantine nodes for any constant 1/2 >; ϵ >; 0 (not depending on n), and outputs the exact value of the poll with high probability. SPP is composed of two sub-protocols, which we believe to be interesting on their own: SPP-Overlay maintains a structured overlay when nodes leave or join the network, and SPP-Computation conducts the actual poll. We validate the practicality of our approach through experimental evaluations and describe briefly two possible applications of SPP: (1) an optimal Byzantine Agreement protocol whose communication complexity is Θ(n log n) and (2) a protocol solving an open question of King and Saia in the context of aggregation functions, namely on the feasibility of performing multiparty secure aggregations with a communication complexity of o(n2). Sébastien Gambs, Rachid Guerraoui, Hamza Harkous, Florian Huc, Anne-Marie Kermarrec |
SRDS | 2 |
| 2012 | Of Choices, Failures and Asynchrony: The Many Faces of Set Agreement
Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Corentin Travers |
Algorithmica | 3 |
| 2012 | Special section with selected papers from PODC 2010
Rachid Guerraoui |
Distributed Comput. | 1 |
| 2012 | Decentralized polling with respectable participants
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod, Ymir Vigfusson |
J. Parallel Distributed Comput. | 1 |
| 2012 | Generating Fast Indulgent Algorithms
Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Corentin Travers |
Theory Comput. Syst. | 3 |
| 2012 | The Weakest Failure Detectors to Solve Quittable Consensus and Nonblocking Atomic CommitabstractWe define quittable consensus, a natural variation of the consensus problem, where processes have the option to agree on “quit” if failures occur, and we relate this problem to the well-known problem of nonblocking atomic commit. We then determine the weakest failure detectors for these two problems in all environments, regardless of the number of faulty processes. Rachid Guerraoui, Vassos Hadzilacos, Petr Kuznetsov, Sam Toueg |
SIAM J. Comput. | 1 |
| 2011 | Generalized Universality
Eli Gafni, Rachid Guerraoui |
CONCUR | 2 |
| 2011 | The Complexity of RenamingabstractWe study the complexity of renaming, a fundamental problem in distributed computing in which a set of processes need to pick distinct names from a given namespace. We prove an individual lower bound of Ω( k ) process steps for deterministic renaming into any namespace of size sub-exponential in k, where k is the number of participants. This bound is tight: it draws an exponential separation between deterministic and randomized solutions, and implies new tight bounds for deterministic fetch-and-increment registers, queues and stacks. The proof of the bound is interesting in its own right, for it relies on the first reduction from renaming to another fundamental problem in distributed computing: mutual exclusion. We complement our individual bound with a global lower bound of Ω(k log (k/c)) on the total step complexity of renaming into a namespace of size ck, for any c ≥ 1. This applies to randomized algorithms against a strong adversary, and helps derive new global lower bounds for randomized approximate counter and fetch-and-increment implementations, all tight within logarithmic factors. Dan Alistarh, James Aspnes, Seth Gilbert, Rachid Guerraoui |
FOCS | 4 |
| 2011 | Democratizing Transactional Programming
Vincent Gramoli, Rachid Guerraoui |
Middleware | 2 |
| 2011 | Model Checking a Networked System Without the Network
Rachid Guerraoui, Maysam Yabandeh |
NSDI | 1 |
| 2011 | Byzantine agreement with homonymsabstractInternational audience Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Anne-Marie Kermarrec, Eric Ruppert, Hung Tran-The |
PODC | 3 |
| 2011 | The complexity of robust atomic storageabstractWe study the time-complexity of robust atomic read/write storage from fault-prone storage components in asynchronous message-passing systems. Robustness here means wait-free tolerating the largest possible number t of Byzantine storage component failures (optimal resilience) without relying on data authentication. We show that no single-writer multiple-reader (SWMR) robust atomic storage implementation exists if (a) read operations complete in less than four communication round-trips (rounds), and (b) the time complexity of write operations is constant. More precisely, we present two lower bounds. The first is a read lower bound stating that three rounds of communication are necessary to read from a SWMR robust atomic storage. The second is a write lower bound, showing that Ω(log(t)) write rounds are necessary to read in three rounds from such a storage. Applied to known results, our lower bounds close a fundamental gap: we show that time-optimal robust atomic storage can be obtained using well-known transformations from regular to atomic storage and existing time-optimal regular storage implementations. © 2011 ACM. Dan Dobre, Rachid Guerraoui, Matthias Majuntke, Neeraj Suri, Marko Vukolic |
PODC | 2 |
| 2011 | Laws of order: expensive synchronization in concurrent algorithms cannot be eliminatedabstractBuilding correct and efficient concurrent algorithms is known to be a difficult problem of fundamental importance. To achieve efficiency, designers try to remove unnecessary and costly synchronization. However, not only is this manual trial-and-error process ad-hoc, time consuming and error-prone, but it often leaves designers pondering the question of: is it inherently impossible to eliminate certain synchronization, or is it that I was unable to eliminate it on this attempt and I should keep trying? Hagit Attiya, Rachid Guerraoui, Danny Hendler, Petr Kuznetsov, Maged M. Michael, Martin T. Vechev |
POPL | 2 |
| 2011 | Brief announcement: transaction polymorphismabstractIn this work, we present transaction polymorphism, a synchronization technique that consists of providing more control to the programmer than traditional (i.e., monomorphic) transactions to achieve comparable performance to generic lock-based and lock-free solutions. Vincent Gramoli, Rachid Guerraoui |
SPAA | 2 |
| 2011 | The disagreement power of an adversary
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Andreas Tielmann |
Distributed Comput. | 3 |
| 2011 | Verification of STM on relaxed memory models
Rachid Guerraoui, Thomas A. Henzinger, Vasu Singh |
Formal Methods Syst. Des. | 1 |
| 2011 | The impossibility of boosting distributed service resilience
Paul C. Attie, Rachid Guerraoui, Petr Kuznetsov, Nancy A. Lynch, Sergio Rajsbaum |
Inf. Comput. | 2 |
| 2011 | The Complexity of Early Deciding Set AgreementabstractIn the k-set agreement problem, each processor starts with a private input value and eventually decides on an output value. At most k distinct output values may be chosen, and every processor's output value must be one of the proposed values. We consider a synchronous message passing system, and we prove a tight bound of $\lfloor f/k\rfloor+2$ rounds of communication for all processors to decide in every run in which at most f processors fail. The lower bound proof proceeds through a simulation of a synchronous solution to k-set agreement in message passing, in an asynchronous shared memory system in which $k-1$ processors may fail, and which was proven to be impossible using topological approaches. In contrast to past complexity results on set agreement, our lower bound proof is purely algorithmic. It does not use any direct topological argument but uses instead the impossibility of asynchronous set agreement to encapsulate the needed topology. We thus derive an adaptive complexity lower bound for a message passing system from a static impossibility in a shared memory system. Eli Gafni, Rachid Guerraoui, Bastian Pochon |
SIAM J. Comput. | 2 |
| 2011 | Stabilization, Safety, and Security of Distributed Systems (SSS 2009)
Ajoy K. Datta, Franck Petit, Rachid Guerraoui |
Theor. Comput. Sci. | 3 |
| 2011 | Collaborative personalized top-k processingabstractThis article presents P4Q, a fully decentralized gossip-based protocol to personalize query processing in social tagging systems. P4Q dynamically associates each user with social acquaintances sharing similar tagging behaviors. Queries are gossiped among such acquaintances, computed on-the-fly in a collaborative, yet partitioned manner, and results are iteratively refined and returned to the querier. Analytical and experimental evaluations convey the scalability of P4Q for top- k query processing, as well its inherent ability to cope with users updating profiles and departing. Xiao Bai 0002, Rachid Guerraoui, Anne-Marie Kermarrec, Vincent Leroy 0001 |
ACM Trans. Database Syst. | 2 |
| 2010 | Gossiping personalized queriesabstractInternational audience Xiao Bai 0002, Marin Bertier, Rachid Guerraoui, Anne-Marie Kermarrec, Vincent Leroy 0001 |
EDBT | 3 |
| 2010 | The next 700 BFT protocolsabstractModern Byzantine fault-tolerant state machine replication (BFT) protocols involve about 20,000 lines of challenging C++ code encompassing synchronization, networking and cryptography. They are notoriously difficult to develop, test and prove. We present a new abstraction to simplify these tasks. We treat a BFT protocol as a composition of instances of our abstraction. Each instance is developed and analyzed independently. Rachid Guerraoui, Nikola Knezevic, Vivien Quéma, Marko Vukolic |
EuroSys | 1 |
| 2010 | How Efficient Can Gossip Be? (On the Cost of Resilient Information Exchange)
Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Morteza Zadimoghaddam |
ICALP (2) | 3 |
| 2010 | The Gossple Anonymous Social Network
Marin Bertier, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Vincent Leroy 0001 |
Middleware | 3 |
| 2010 | LiFTinG: Lightweight Freerider-Tracking in Gossip
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod, Swagatika Prusty |
Middleware | 1 |
| 2010 | WhatsUp: News, From, For, Through, EveryoneabstractWhatsUp (WUP) is a new form of electronic news. It is personalized and decentralized. Users receive news and have the ability to express their interest in it. This opinion, in turn, is used as an implicit and dynamic subscription scheme to filter and personalize future information. The system is peer-to-peer: no big brother company controls the news, and no central server makes it vulnerable to failures, censorship or attacks. At the heart of WUP lies the idea of collaborative filtering applied to the dissemination of news: people who liked the same news in the past might as well like the same news in the future: irrelevant news disappear by themselves. The idea is put to work through Beep: a biased epidemic dissemination (gossip) protocol that delivers news to interested users in a timely manner, despite jamming and churn. Beep is dynamically parameterized on a per- user, per-news, and per-dissemination-hop basis. When compared to a classical epidemic dissemination protocol, Beep has two key characteristics: orientation and amplification. Every user forwards the news of interest to a randomly selected set of users largely constituted by those who have similar interests (orientation). Moreover, the size of this set of users depends on the level of interest in the news itself (amplification). Antoine Boutet, Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec |
Peer-to-Peer Computing | 3 |
| 2010 | Boosting Gossip for Live StreamingabstractGossip protocols are considered very effective to disseminate information in a large scale dynamic distributed system. Their inherent simplicity makes them easy to implement and deploy. However, whereas their probabilistic guarantees are often enough to disseminate data in the context of low- bandwidth applications, they typically do not suffice for high-bandwidth content dissemination: missing 1% is unacceptable for live streaming. In this paper, we show how the combination of two simple mechanisms copes with this seemingly inherent deficiency of gossip: (i) codec, an erasure coding scheme, and (ii) claim2, a content- request scheme that leverages gossip duplication to diversify the retransmission sources of missing information. We show how these mechanisms can effectively complement each other in a new gossip protocol, gossip++, which retains the simplicity of deployment of plain gossip. In a realistic setting with an average bandwidth capability (800 kbps) close to the stream rate (680 kbps) and 1% message loss, plain gossip can provide at most 99% of the stream. Using gossip++, on the other hand, all nodes can view a perfectly clear stream. Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Maxime Monod |
Peer-to-Peer Computing | 2 |
| 2010 | Leveraging parallel nesting in transactional memoryabstractExploiting the emerging reality of affordable multi-core architectures goes through providing programmers with simple abstractions that would enable them to easily turn their sequential programs into concurrent ones that expose as much parallelism as possible. While transactional memory promises to make concurrent programming easy to a wide programmer community, current implementations either disallow nested transactions to run in parallel or do not scale to arbitrary parallel nesting depths. This is an important obstacle to the central goal of transactional memory, as programmers can only start parallel threads in restricted parts of their code. João Barreto 0001, Aleksandar Dragojevic, Paulo Ferreira 0001, Rachid Guerraoui, Michal Kapalka |
PPoPP | 4 |
| 2010 | Securing every bit: authenticated broadcast in radio networksabstractThis paper studies non-cryptographic authenticated broadcast in radio networks subject to malicious failures. We introduce two protocols that address this problem. The first, NeighborWatchRB, makes use of a novel strategy in which honest devices monitor their neighbors for malicious behavior. Second, we present a more robust variant, MultiPathRB, that tolerates the maximum possible density of malicious devices per region, using an elaborate voting strategy. We also introduce a new proof technique to show that both protocols ensure asymptotically optimal running time. Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Zarko Milosevic 0001, Calvin C. Newport |
SPAA | 3 |
| 2010 | Brief announcement: byzantine agreement with homonymsabstractIn this work, we address Byzantine agreement in a message passing system with homonyms, i.e. a system with a number l of authenticated identities that is independent of the total number of processes n, in the presence of t < n Byzantine processes. Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Anne-Marie Kermarrec |
SPAA | 3 |
| 2010 | Collaborative scoring with dishonest participantsabstractConsider a set of players that are interested in collectively evaluating a set of objects. We develop a collaborative scoring protocol in which each player evaluates a subset of the objects, after which we can accurately predict each players' individual opinion of the remaining objects. The accuracy of the predictions is near optimal, depending on the number of objects evaluated by each player and the correlation among the players' preferences. Seth Gilbert, Rachid Guerraoui, Faezeh Malakouti Rad, Morteza Zadimoghaddam |
SPAA | 2 |
| 2010 | Transactions in the jungleabstractTransactional memory (TM) has shown potential to simplify the task of writing concurrent programs. Inspired by classical work on databases, formal definitions of the semantics of TM executions have been proposed. Many of these definitions assumed that accesses to shared data are solely performed through transactions. In practice, due to legacy code and concurrency libraries, transactions in a TM have to share data with non-transactional operations. The semantics of such interaction, while widely discussed by practitioners, lacks a clear formal specification. Those interactions can vary, sometimes in subtle ways, between TM implementations and underlying memory models. Rachid Guerraoui, Thomas A. Henzinger, Michal Kapalka, Vasu Singh |
SPAA | 1 |
| 2010 | Computing in Social Networks
Andrei Giurgiu, Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec |
SSS | 2 |
| 2010 | Fast Randomized Test-and-Set and Renaming
Dan Alistarh, Hagit Attiya, Seth Gilbert, Andrei Giurgiu, Rachid Guerraoui |
DISC | 5 |
| 2010 | Brief Announcement: New Bounds for Partially Synchronous Set Agreement
Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Corentin Travers |
DISC | 3 |
| 2010 | Foundations of Speculative Distributed Computing - (Invited Lecture Extended Abstract)
Rachid Guerraoui |
DISC | 1 |
| 2010 | Model checking transactional memories
Rachid Guerraoui, Thomas A. Henzinger, Vasu Singh |
Distributed Comput. | 1 |
| 2010 | Refined quorum systems
Rachid Guerraoui, Marko Vukolic |
Distributed Comput. | 1 |
| 2010 | Tight failure detection bounds on atomic object implementationsabstractThis article determines the weakest failure detectors to implement shared atomic objects in a distributed system with crash-prone processes. We first determine the weakest failure detector for the basic register object. We then use that to determine the weakest failure detector for all popular atomic objects including test-and-set, fetch-and-add, queue, consensus and compare-and-swap, which we show is the same. Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui |
J. ACM | 3 |
| 2010 | Fast Access to Distributed Atomic MemoryabstractWe study efficient and robust implementations of an atomic read-write data structure over an asynchronous distributed message-passing system made of reader and writer processes, as well as a number of servers implementing the data structure. We determine the exact conditions under which every read and write involves one round of communication with the servers. These conditions relate the number of readers to the tolerated number of faulty servers and the nature of these failures. Partha Dutta, Rachid Guerraoui, Ron R. Levy, Marko Vukolic |
SIAM J. Comput. | 2 |
| 2010 | Reflexes: Abstractions for integrating highly responsive tasks into Java applicationsabstractAchieving submillisecond response times in a managed language environment such as Java or C# requires overcoming significant challenges. In this article, we propose Reflexes, a programming model and runtime system infrastructure that lets developers seamlessly mix highly responsive tasks and timing-oblivious Java applications. Thus enabling gradual addition of real-time features, to a non-real-time application without having to resort to recoding the real-time parts in a different language such as C or Ada. Experiments with the Reflex prototype implementation show that it is possible to run a real-time task with a period of 45 μs with an accuracy of 99.996% (only 0.001% worse than the corresponding C implementation) in the presence of garbage collection and heavy load ordinary Java threads. Jesper Honig Spring, Filip Pizlo, Jean Privat, Rachid Guerraoui, Jan Vitek |
ACM Trans. Embed. Comput. Syst. | 4 |
| 2010 | Throughput optimal total order broadcast for cluster environmentsabstractTotal order broadcast is a fundamental communication primitive that plays a central role in bringing cheap software-based high availability to a wide range of services. This article studies the practical performance of such a primitive on a cluster of homogeneous machines. We present LCR, the first throughput optimal uniform total order broadcast protocol. LCR is based on a ring topology. It only relies on point-to-point inter-process communication and has a linear latency with respect to the number of processes. LCR is also fair in the sense that each process has an equal opportunity of having its messages delivered by all processes. We benchmark a C implementation of LCR against Spread and JGroups, two of the most widely used group communication packages. LCR provides higher throughput than the alternatives, over a large number of scenarios. Rachid Guerraoui, Ron R. Levy, Bastian Pochon, Vivien Quéma |
ACM Trans. Comput. Syst. | 1 |
| 2009 | Software Transactional Memory on Relaxed Memory Models
Rachid Guerraoui, Thomas A. Henzinger, Vasu Singh |
CAV | 1 |
| 2009 | Transactional Memory: Glimmer of a Theory
Rachid Guerraoui, Michal Kapalka |
CAV | 1 |
| 2009 | High-Performance Transactional Event Processing
Antonio Cunei, Rachid Guerraoui, Jesper Honig Spring, Jean Privat, Jan Vitek |
COORDINATION | 2 |
| 2009 | Stretching gossip with live streamingabstractGossip-based information dissemination protocols are considered easy to deploy, scalable and resilient to network dynamics. They are also considered highly flexible, namely tunable at will to increase their robustness and adapt to churn. So far however, they have mainly been evaluated through simulation, very often assuming ideal settings. Instead, in this paper, we report on an extensive study of gossip protocols, deployed on a 230 Planetlab node testbed, in the context of a challenging video streaming application in environments with constrained bandwidths. More precisely, we assess the impact of varying the well known knobs of gossip, fanout and refresh rate, in various upload-bandwidth distributions and churn. Our results show that in such challenging contexts, the performance of gossip protocols may be hampered by high fanout values. We also show that the more proactive a gossip protocol, the better it copes with churn. For instance, when 20% of the nodes simultaneously crash, 70% of the remaining nodes do not suffer any loss in stream quality, while the others only experience a performance decrease for an average of 5 seconds around the churn event. Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Maxime Monod, Vivien Quéma |
DSN | 2 |
| 2009 | Names Trump Malice: Tiny Mobile Agents Can Tolerate Byzantine Failures
Rachid Guerraoui, Eric Ruppert |
ICALP (2) | 1 |
| 2009 | Interference-Resilient Information ExchangeabstractThis paper presents an efficient protocol for reliably exchanging information in a single-hop, multi-channel radio network subject to unpredictable interference. We model the interference by an adversary that can simultaneously disrupt up to t of the C available channels. We assume no shared secret keys or third-party infrastructure. The running time of our protocol depends on the gap between C and t: when the number of channels C = Q,(t2), the running time is linear; when only C = t +1 channels are available, the running time is exponential. We prove that exponential-time is unavoidable in the latter case. At the core of our protocol lies a combinatorial function, possibly of independent interest, described for the first time in this paper: the multi-selector. A multi-selector generates a sequence of channel assignments for each device such that every sufficiently large subset of devices is partitioned onto distinct channels by at least one of these assignments. Seth Gilbert, Rachid Guerraoui, Dariusz R. Kowalski, Calvin C. Newport |
INFOCOM | 2 |
| 2009 | Of Choices, Failures and Asynchrony: The Many Faces of Set Agreement
Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Corentin Travers |
ISAAC | 3 |
| 2009 | Heterogeneous Gossip
Davide Frey, Rachid Guerraoui, Anne-Marie Kermarrec, Boris Koldehofe, Martin Mogensen, Maxime Monod, Vivien Quéma |
Middleware | 2 |
| 2009 | Decentralized Polling with Respectable Participants
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod |
OPODIS | 1 |
| 2009 | On Tracking Freeriders in Gossip ProtocolsabstractPeer-to-peer content dissemination applications suffer immensely from freeriders, i.e., nodes that do not provide their fair share. The Tit-for-Tat (TfT) incentives have received much attention as they help make such systems more robust against freeriding. However, these rely on an asymmetric component, namely opportunistic pushes, that let peers receive content without sending anything in return. Opportunistic push constitutes the Achilles' heel of TfT-based protocols as illustrated by the fact that all known attacks against them exploit it. This problem becomes even more serious when used by colluding freeriders. In this paper, we discuss the possibility of using accountability to secure gossip-based dissemination protocols based on asymmetric exchanges. The fact that gossip protocols are dynamic and randomized makes our approach robust against collusion and alleviates the need for cryptography. We present the challenges raised by an auditing approach and give insights into how to build a freerider-tracking protocol for gossip-based content dissemination. Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod |
Peer-to-Peer Computing | 1 |
| 2009 | Stretching transactional memoryabstractTransactional memory (TM) is an appealing abstraction for programming multi-core systems. Potential target applications for TM, such as business software and video games, are likely to involve complex data structures and large transactions, requiring specific software solutions (STM). So far, however, STMs have been mainly evaluated and optimized for smaller scale benchmarks. Aleksandar Dragojevic, Rachid Guerraoui, Michal Kapalka |
PLDI | 2 |
| 2009 | The disagreement power of an adversary: extended abstractabstractAt the heart of distributed computing lies the fundamental result that the level of agreement that can be obtained in an asynchronous shared memory model where t processes can crash is exactly t+1. In other words, an adversary that can crash any subset of size at most t can prevent the processes from agreeing on t values. But what about the rest (22n − n) adversaries that might crash certain combination of processes and not others? Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Andreas Tielmann |
PODC | 3 |
| 2009 | The wireless synchronization problemabstractIn this paper, we study the wireless synchronization problem which requires devices activated at different times on a congested single-hop radio network to synchronize their round numbering. We assume a collection of n synchronous devices with access to a shared band of the radio spectrum, divided into F narrowband frequencies. We assume that the communication medium suffers from unpredictable, perhaps even malicious interference, which we model by an adversary that can disrupt up to t frequencies per round. Devices begin executing in different rounds and the exact number of participants is not known in advance. Shlomi Dolev, Seth Gilbert, Rachid Guerraoui, Fabian Kuhn, Calvin C. Newport |
PODC | 3 |
| 2009 | Preventing versus curing: avoiding conflicts in transactional memoriesabstractTransactional memories are typically speculative and rely on contention managers to cure conflicts. This paper explores a complementary approach that prevents conflicts by scheduling transactions according to predictions on their access sets. Aleksandar Dragojevic, Rachid Guerraoui, Anmol V. Singh, Vasu Singh |
PODC | 2 |
| 2009 | The semantics of progress in lock-based transactional memoryabstractTransactional memory (TM) is a promising paradigm for concurrent programming. Whereas the number of TM implementations is growing, however, little research has been conducted to precisely define TM semantics, especially their progress guarantees. This paper is the first to formally define the progress semantics of lockbased TMs, which are considered the most effective in practice. Rachid Guerraoui, Michal Kapalka |
POPL | 1 |
| 2009 | The 2009 Edsger W. Dijkstra Prize in Distributed Computing
Lorenzo Alvisi, Rachid Guerraoui, Prasad Jayanti, Idit Keidar, Shay Kutten, Jennifer L. Welch |
DISC | 2 |
| 2009 | The Disagreement Power of an Adversary
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Andreas Tielmann |
DISC | 3 |
| 2009 | Elastic Transactions
Pascal Felber, Vincent Gramoli, Rachid Guerraoui |
DISC | 3 |
| 2009 | Brief Announcement: Towards Secured Distributed Polling in Social Networks
Rachid Guerraoui, Kévin Huguenin, Anne-Marie Kermarrec, Maxime Monod |
DISC | 1 |
| 2009 | On the weakest failure detector ever
Rachid Guerraoui, Maurice Herlihy, Petr Kuznetsov, Nancy A. Lynch, Calvin C. Newport |
Distributed Comput. | 1 |
| 2009 | The complexity of obstruction-free implementationsabstractObstruction-free implementations of concurrent objects are optimized for the common case where there is no step contention , and were recently advocated as a solution to the costs associated with synchronization without locks. In this article, we study this claim and this goes through precisely defining the notions of obstruction-freedom and step contention. We consider several classes of obstruction-free implementations, present corresponding generic object implementations, and prove lower bounds on their complexity. Viewed collectively, our results establish that the worst-case operation time complexity of obstruction-free implementations is high, even in the absence of step contention. We also show that lock-based implementations are not subject to some of the time-complexity lower bounds we present. Hagit Attiya, Rachid Guerraoui, Danny Hendler, Petr Kuznetsov |
J. ACM | 2 |
| 2009 | Of malicious motes and suspicious sensors: On the efficiency of malicious interference in wireless networks
Seth Gilbert, Rachid Guerraoui, Calvin C. Newport |
Theor. Comput. Sci. | 2 |
| 2009 | A topological treatment of early-deciding set-agreement
Rachid Guerraoui, Maurice Herlihy, Bastian Pochon |
Theor. Comput. Sci. | 1 |
| 2008 | Completeness and Nondeterminism in Model Checking Transactional Memories
Rachid Guerraoui, Thomas A. Henzinger, Vasu Singh |
CONCUR | 1 |
| 2008 | A Scalable and Oblivious Atomicity Assertion
Rachid Guerraoui, Marko Vukolic |
CONCUR | 1 |
| 2008 | The Return of Transactions
Rachid Guerraoui |
ECOOP | 1 |
| 2008 | An Arbitrary Tree-Structured Replica Control ProtocolabstractTraditional replication protocols that arrange logically the replicas into a tree structure have reasonable availability, low communication costs but induce high system load. We propose in this paper the arbitrary protocol: a tree-based replica control protocol that can be configured based on the frequencies of read and write operations in order to provide lower system load than existing tree replication protocols, yet with comparable cost and availability. Our protocol enables the shifting from one configuration into another by just modifying the structure of the tree. There is no need to implement a new protocol whenever the frequencies of read and write operations change. At the heart of our protocol lies the new idea of logical and physical levels in a tree. In short, read operations are carried out on any physical node of every physical level of the tree whereas the write operation is performed on all physical nodes of a single physical level of the tree. We discuss optimal configurations, proving in particular a new lower bound, of independent interest, for the case of a binary tree. Jean Paul Bahsoun, Robert Basmadjian, Rachid Guerraoui |
ICDCS | 3 |
| 2008 | Flexible task graphs: a unified restricted thread programming model for javaabstractThe disadvantages of unconstrained shared-memory multi-threading in Java, especially with regard to latency and determinism in realtime systems, have given rise to a variety of language extensions that place restrictions on how threads allocate, share, and communicate memory, leading to order-of-magnitude reductions in latency and jitter. However, each model makes different trade-offs with respect to expressiveness, efficiency, enforcement, and latency, and no one model is best for all applications. Joshua S. Auerbach, David F. Bacon, Rachid Guerraoui, Jesper Honig Spring, Jan Vitek |
LCTES | 3 |
| 2008 | The Next 700 BFT Protocols
Rachid Guerraoui |
OPODIS | 1 |
| 2008 | Model checking transactional memoriesabstractModel checking software transactional memories (STMs) is difficult because of the unbounded number, length, and delay of concurrent transactions and the unbounded size of the memory. We show that, under certain conditions, the verification problem can be reduced to a finite-state problem, and we illustrate the use of the method by proving the correctness of several STMs, including two-phase locking, DSTM, TL2, and optimistic concurrency control. The safety properties we consider include strict serializability and opacity; the liveness properties include obstruction freedom, livelock freedom, and wait freedom. Rachid Guerraoui, Thomas A. Henzinger, Barbara Jobstmann, Vasu Singh |
PLDI | 1 |
| 2008 | Sharing is harder than agreeingabstractOne of the most celebrated results of the theory of distributed computing is the impossibility, in an asynchronous system of n processes that communicate through shared memory registers, to solve the set agreement problem where the processes need to decide on up to n-1 among their n initial values. In short, the result indicates that the register abstraction is too weak to implement the set agreement one. Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui |
PODC | 3 |
| 2008 | Secure communication over radio channelsabstractWe study the problem of secure communication in a multi-channel, single-hop radio network with a malicious adversary that can cause collisions and spoof messages. We assume no pre-shared secrets or trusted-third-party infrastructure. The main contribution of this paper is f-AME: a randomized (f)ast-(A)uthenticated (M)essage (E)xchange protocol that enables nodes to exchange messages in a reliable and authenticated manner. It runs in O(|E|t2 log n) time and has optimal resilience to disruption, where E is the set of pairs of nodes that need to swap messages, n is the total number of nodes, C the number of channels, and t < C the number of channels on which the adversary can participate in each round. We show how to use f-AME to establish a shared secret group key, which can be used to implement a secure, reliable and authenticated long-lived communication service. The resulting service requires O(nt3 log n) rounds for the setup phase, and O(t log n) rounds for an arbitrary pair to communicate. By contrast, existing solutions rely on pre-shared secrets, trusted third-party infrastructure, and/or the assumption that all interference is non-malicious. Shlomi Dolev, Seth Gilbert, Rachid Guerraoui, Calvin C. Newport |
PODC | 3 |
| 2008 | On the complexity of asynchronous gossipabstractIn this paper, we study the complexity of gossip in an asynchronous, message-passing fault-prone distributed system. In short, we show that an adaptive adversary can significantly hamper the spreading of a rumor, while an oblivious adversary cannot. This latter fact implies that there exist message-efficient asynchronous (randomized) consensus protocols, in the context of an oblivious adversary. Chryssis Georgiou, Seth Gilbert, Rachid Guerraoui, Dariusz R. Kowalski |
PODC | 3 |
| 2008 | Extensible encoding of type hierarchiesabstractThe subtyping test consists of checking whether a type t is a descendant of a type r (Agrawal et al. 1989). We study how to perform such a test efficiently, assuming a dynamic hierarchy when new types are inserted at run-time. The goal is to achieve time and space efficiency, even as new types are inserted. We propose an extensible scheme, named ESE, that ensures (1) efficient insertion of new types, (2) efficient subtyping tests, and (3) small space usage. On the one hand ESE provides comparable test times to the most efficient existing static schemes (e.g.,Zibin et al. (2001)). On the other hand, ESE has comparable insertion times to the most efficient existing dynamic scheme (Baehni et al. 2007), while ESE outperforms it by a factor of 2-3 times in terms of space usage. Hamed S. Alavi, Seth Gilbert, Rachid Guerraoui |
POPL | 3 |
| 2008 | On the correctness of transactional memoryabstractTransactional memory (TM) is perceived as an appealing alternative to critical sections for general purpose concurrent programming. Despite the large amount of recent work on TM implementations, however, very little effort has been devoted to precisely defining what guarantees these implementations should provide. A formal description of such guarantees is necessary in order to check the correctness of TM systems, as well as to establish TM optimality results and inherent trade-offs. Rachid Guerraoui, Michal Kapalka |
PPoPP | 1 |
| 2008 | Partial snapshot objectsabstractWe introduce a generalization of the atomic snapshot object, which we call the partial snapshot object. This object stores a vector of values. Processes may write components of the vector individually or atomically scan any subset of the components. We investigate implementations of the latter partial scan operation that are more efficient than the complete scans of traditional snapshot objects. We present an algorithm that is based on a new implementation of the active set abstraction, which may be of independent interest. Hagit Attiya, Rachid Guerraoui, Eric Ruppert |
SPAA | 2 |
| 2008 | On obstruction-free transactionsabstractThis paper studies obstruction-free software transactional memory systems (OFTMs). These systems are appealing, for they combine the atomicity property of transactions with a liveness property that ensures the commitment of every transaction that eventually encounters no contention. Rachid Guerraoui, Michal Kapalka |
SPAA | 1 |
| 2008 | How to Solve Consensus in the Smallest Window of Synchrony
Dan Alistarh, Seth Gilbert, Rachid Guerraoui, Corentin Travers |
DISC | 3 |
| 2008 | The Weakest Failure Detector for Message Passing Set-Agreement
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Andreas Tielmann |
DISC | 3 |
| 2008 | Optimistic Erasure-Coded Distributed Storage
Partha Dutta, Rachid Guerraoui, Ron R. Levy |
DISC | 2 |
| 2008 | Permissiveness in Transactional Memories
Rachid Guerraoui, Thomas A. Henzinger, Vasu Singh |
DISC | 1 |
| 2008 | Failure detectors as type boostersabstractThe power of an object type T can be measured as the maximum number n of processes that can solve consensus using only objects of T and registers. This number, denoted cons(T), is called the consensus power of T. This paper addresses the question of the weakest failure detector to solve consensus among a number k > n of processes that communicate using shared objects of a type T with consensus power n. In other words, we seek for a failure detector that is sufficient and necessary to “boost” the consensus power of a type T from n to k. It was shown in Neiger (Proceedings of the 14th annual ACM symposium on principles of distributed computing (PODC), pp. 100–109, 1995) that a certain failure detector, denoted Ω n , is sufficient to boost the power of a type T from n to k, and it was conjectured that Ω n was also necessary. In this paper, we prove this conjecture for one-shot deterministic types. We first show that, for any one-shot deterministic type T with cons(T) ≤ n, Ω n is necessary to boost the power of T from n to n + 1. Then we go a step further and show that Ω n is also the weakest to boost the power of (n + 1)-ported one-shot deterministic types from n to any k > n. Our result generalizes, in a precise sense, the result of the weakest failure detector to solve consensus in asynchronous message-passing systems (Chandra et al. in J ACM 43(4):685–722, 1996). As a corollary, we show that Ω t is the weakest failure detector to boost the resilience level of a distributed shared memory system, i.e., to solve consensus among n > t processes using (t − 1)-resilient objects of consensus power t. Rachid Guerraoui, Petr Kuznetsov |
Distributed Comput. | 1 |
| 2008 | The weakest failure detectors to boost obstruction-freedomabstractIt is considered good practice in concurrent computing to devise shared object implementations that ensure a minimal obstruction-free progress property and delegate the task of boosting liveness to independent generic oracles called contention managers. This paper determines necessary and sufficient conditions to implement wait-free and non-blocking contention managers, i.e., contention managers that ensure wait-freedom (resp. non-blockingness) of any associated obstruction-free object implementation. The necessary conditions hold even when universal objects (like compare-and-swap) or random oracles are available in the implementation of the contention manager. On the other hand, the sufficient conditions assume only basic read/write objects, i.e., registers. We show that failure detector $$\lozenge{\fancyscript{P}}$$ is the weakest to convert any obstruction-free algorithm into a wait-free one, and Ω *, a new failure detector which we introduce in this paper, and which is strictly weaker than $$\lozenge\fancyscript{P}$$ but strictly stronger than Ω, is the weakest to convert any obstruction-free algorithm into a non-blocking one. We also address the issue of minimizing the overhead imposed by contention management in low contention scenarios. We propose two intermittent failure detectors $$I_{\Omega^*}$$ and $$I_{\lozenge\fancyscript{P}}$$ that are in a precise sense equivalent to, respectively, Ω * and $$\lozenge\fancyscript{P}$$ , but allow for reducing the cost of failure detection in eventually synchronous systems when there is little contention. We present two contention managers: a non-blocking one and a wait-free one, that use, respectively, $$I_{\Omega^*}$$ and $$I_{\lozenge\fancyscript{P}}$$ . When there is no contention, the first induces very little overhead whereas the second induces some non-trivial overhead. We show that wait-free contention managers, unlike their non-blocking counterparts, impose an inherent non-trivial overhead even in contention-free executions. Rachid Guerraoui, Michal Kapalka, Petr Kuznetsov |
Distributed Comput. | 1 |
| 2008 | The gap in circumventing the impossibility of consensus
Rachid Guerraoui, Petr Kuznetsov |
J. Comput. Syst. Sci. | 1 |
| 2008 | A general characterization of indulgenceabstractAn indulgent algorithm is a distributed algorithm that, besides tolerating process failures, also tolerates unreliable information about the interleaving of the processes. This article presents a general characterization of indulgence in an abstract computing model that encompasses various communication and resilience schemes. We use our characterization to establish several results about the inherent power and limitations of indulgent algorithms. Rachid Guerraoui, Nancy A. Lynch |
ACM Trans. Auton. Adapt. Syst. | 1 |
| 2008 | The collective memory of amnesic processesabstractThis article considers the problem of robustly emulating a shared atomic memory over a distributed message-passing system where processes can fail by crashing and possibly recover. We revisit the notion of atomicity in the crash-recovery context and introduce a generic algorithm that emulates an atomic memory. The algorithm is instantiated for various settings according to whether processes have access to local stable storage, and whether, in every execution of the algorithm, a sufficient number of processes are assumed not to crash. We establish the optimality of specific instances of our algorithm in terms of resilience , log complexity (number of stable storage accesses needed in every read or write operation), as well as time complexity (number of communication steps needed in every read or write operation). The article also discusses the impact of considering a multiwriter versus a single-writer memory, as well as the impact of weakening the consistency of the memory by providing safe or regular semantics instead of atomicity. Rachid Guerraoui, Ron R. Levy, Bastian Pochon, Jim Pugh |
ACM Trans. Algorithms | 1 |
| 2007 | A Universal Construction for Concurrent ObjectsabstractA concurrent object is an object that can be concurrently accessed by several processes. A wait-free implementation of an object is such that any operation issued by a non-faulty process terminates in a finite number of its own steps, whatever the behavior of the other processes (that can be very slow or even have crashed). An object type is universal if objects of that type, together with atomic registers, allows implementing any concurrent object defined by a sequential specification. A universal construction is a wait-free algorithm, based only on atomic registers and universal objects, that, given any sequential object type T, provides the processes with a wait-free concurrent object of the type T. In a famous paper (titled "Wait-free synchronization") Herlihy has shown that consensus objects are universal, and has presented a consensus-based universal construction. We present here a new universal construction. That construction, that is built incrementally, is particularly simple. While, in addition to consensus objects, Herlihy's universal construction uses low-level objects such as pointers, the design of the construction presented here is based on the simple and well-known state machine replication paradigm. Its proof is also simple and consequently allows to better understand not only the power of consensus objects but also the subtleties of wait-free computations and the way the consensus objects allow coping with both process failures and non-determinism. In that sense, this paper has a pedagogical flavor. Rachid Guerraoui, Michel Raynal |
ARES | 1 |
| 2007 | A comparison of optimistic approaches to collaborative editing of Wiki pagesabstractWikis, a popular tool for sharing knowledge, are basically collaborative editing systems. However, existing Wiki systems offer limited support for co-operative authoring, and they do not scale well, because they are based on a centralised architecture. This paper compares the well-known centralised MediaWiki system with several peer-to-peer approaches to editing of wiki pages: an operational transformation approach (MOT2), a commutativity-oriented approach (WOOTO) and a conflict resolution approach (ACF). We evaluate and compare them, according to a number of qualitative and quantitative metrics. Claudia-Lavinia Ignat, Gérald Oster, Pascal Molli, Michèle Cart, Jean Ferrié, Anne-Marie Kermarrec, Pierre Sutra, Marc Shapiro 0001, Lamia Benmouffok, Jean-Michel Busca, Rachid Guerraoui |
CollaborateCom | 11 |
| 2007 | STMBench7: a benchmark for software transactional memoryabstractSoftware transactional memory (STM) is a promising technique for controlling concurrency in modern multi-processor architectures. STM aims to be more scalable than explicit coarse-grained locking and easier to use than fine-grained locking. However, STM implementations have yet to demonstrate that their runtime overheads are acceptable. To date, empiric evaluations of these implementations have suffered from the lack of realistic benchmarks. Measuring performance of an STM in an overly simplified setting can be at best uninformative and at worst misleading as it may steer researchers to try to optimize irrelevant aspects of their implementations. Rachid Guerraoui, Michal Kapalka, Jan Vitek |
EuroSys | 1 |
| 2007 | A High Throughput Atomic Storage AlgorithmabstractThis paper presents an algorithm to ensure the atomicity of a distributed storage that can be read and written by any number of clients. In failure-free and synchronous situations, and even if there is contention, our algorithm has a high write throughput and a read throughput that grows linearly with the number of available servers. The algorithm is devised with a homogeneous cluster of servers in mind. It organizes servers around a ring and assumes point-to-point communication. It is resilient to the crash failure of any number of readers and writers as well as to the crash failure of all but one server. We evaluated our algorithm on a cluster of 24 nodes with dual fast ethernet network interfaces (100 Mbps). We achieve 81 Mbps of write throughput and 8×90 Mbps of read throughput (with up to 8 servers) which conveys the linear scalability with the number of servers. Rachid Guerraoui, Dejan Kostic, Ron R. Levy, Vivien Quéma |
ICDCS | 1 |
| 2007 | Streamflex: high-throughput stream programming in javaabstractThe stream programming paradigm aims to expose coarsegrained parallelism in applications that must process continuous sequences of events.The appeal of stream programming comes from its conceptual simplicity.A program is a collection of independent filters which communicate by the means of uni-directional data channels.This model lends itself naturally to concurrent and efficient implementations on modern multiprocessors.As the output behavior of filters is determined by the state of their input channels, stream programs have fewer opportunities for the errors (such as data races and deadlocks) that plague shared memory concurrent programming.This paper introduces STREAMFLEX, an extension to Java which marries streams with objects and thus enables to combine, in the same Java virtual machine, stream processing code with traditional object-oriented components.STREAMFLEX targets high-throughput low-latency applications with stringent quality-of-service requirements.To achieve these goals, it must, at the same time, extend and restrict Java.To allow for program optimization and provide latency guarantees, the STREAMFLEX compiler restricts Java by imposing a stricter typing discipline on filters.On the other hand, STREAMFLEX extends the Java virtual machine with real-time capabilities, transactional memory and type-safe region-based allocation.The result is a rich and expressive language that can be implemented efficiently. Jesper Honig Spring, Jean Privat, Rachid Guerraoui, Jan Vitek |
OOPSLA | 3 |
| 2007 | Secretive Birds: Privacy in Population Protocols
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Eric Ruppert |
OPODIS | 3 |
| 2007 | On the weakest failure detector everabstractMany problems in distributed computing are impossible when no information about process failures is available. It is common to ask what information about failures is necessary and sufficient to circumvent some specific impossibility, e.g., consensus, atomic commit, mutual exclusion, etc. This paper asks what information about failures is needed to circumvent any impossibility and sufficient to circumvent some impossibility. In other words, what is the minimal yet non-trivial failure informatio. Rachid Guerraoui, Maurice Herlihy, Petr Kuznetsov, Nancy A. Lynch, Calvin C. Newport |
PODC | 1 |
| 2007 | Refined quorum systemsabstractIt is considered good distributed computing practice to devise object implementations that tolerate contention, periods of asynchrony and a large number of failures, but perform fast if few failures occur, the system is synchronous and there is no contention. This paper initiates the first study of quorum systems that help design such implementations by encompassing, at the same time, optimal resilience (just like traditional quorum systems), as well as optimal best-case complexity (unlike traditional quorum systems). Rachid Guerraoui, Marko Vukolic |
PODC | 1 |
| 2007 | Reflexes: abstractions for highly responsive systemsabstractCommercial Java virtual machines are designed to maximize the performance of applications at the expense of predictability. High throughput garbage collection algorithms, for example, can introduce pauses of 100 milliseconds or more. We are interested in supporting applications with response times in the tens of microseconds and their integration with larger timing-oblivious applications in the same Java virtual machine. We propose Reflexes, a new abstraction for writing highly responsive systems in Java and investigate the virtual machine support needed to add Reflexes to a Java environment. Our implementation of Reflexes was evaluated on several programs including an audio-processing application running at 22.05KHz. The number of missed deadlines, less than 0.2% for 10 million observations, compares favorably to a native C implementation. Jesper Honig Spring, Filip Pizlo, Rachid Guerraoui, Jan Vitek |
VEE | 3 |
| 2007 | Amnesic Distributed Storage
Gregory V. Chockler, Rachid Guerraoui, Idit Keidar |
DISC | 2 |
| 2007 | Gossiping in a Multi-channel Radio Network
Shlomi Dolev, Seth Gilbert, Rachid Guerraoui, Calvin C. Newport |
DISC | 3 |
| 2007 | On the Message Complexity of Indulgent Consensus
Seth Gilbert, Rachid Guerraoui, Dariusz R. Kowalski |
DISC | 2 |
| 2007 | The Alpha of Indulgent ConsensusabstractThis paper presents a simple framework unifying a family of consensus algorithms that can tolerate process crash failures and asynchronous periods of the network, also called indulgent consensus algorithms. Key to the framework is a new abstraction we introduce here, called Alpha, and which precisely captures consensus safety. Implementations of Alpha in shared memory, storage area network, message passing and active disk systems are presented, leading to directly derived consensus algorithms suited to these communication media. The paper also considers the case where the number of processes is unknown and can be arbitrarily large. Rachid Guerraoui, Michel Raynal |
Comput. J. | 1 |
| 2007 | The overhead of consensus failure recovery
Partha Dutta, Rachid Guerraoui, Idit Keidar |
Distributed Comput. | 2 |
| 2007 | Anonymous and fault-tolerant shared-memory computing
Rachid Guerraoui, Eric Ruppert |
Distributed Comput. | 1 |
| 2007 | The perfectly synchronized round-based model of distributed computing
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Bastian Pochon |
Inf. Comput. | 3 |
| 2007 | The Time-Complexity of Local Decision in Distributed AgreementabstractAgreement is at the heart of distributed computing. In its simple form, it requires a set of processes to decide on a common value out of the values they propose. The time-complexity of distributed agreement problems is generally measured in terms of the number of communication rounds needed to achieve a global decision, i.e., for all nonfaulty (correct) processes to reach a decision. This paper studies the time-complexity of local decisions in agreement problems, which we define as the number of communication rounds needed for at least one correct process to decide. We explore bounds for early local decision that depend on the number f of actual failures (that occur in a given run of an algorithm), out of the maximum number t of failures tolerated (by the algorithm). We first consider the synchronous message-passing model where we give tight local decision bounds for three variants of agreement: consensus, uniform consensus, and (nonblocking) atomic commit. We use these results to (1) show that, for consensus, local decision bounds are not compatible with global decision bounds (roughly speaking, they cannot be reached by the same algorithm), and (2) draw the first sharp line between the time-complexity of uniform consensus and atomic commit. Then we consider the eventually synchronous model, where we give tight local decision bounds for synchronous runs of uniform consensus. (In this model, consensus and uniform consensus are similar, atomic commit is impossible, and one cannot bound the number of rounds to reach a decision in nonsynchronous runs of consensus algorithms.) We prove a counterintuitive result that the early local decision bound is the same as the early global decision bound. We also give a matching early deciding consensus algorithm that is significantly better than previous eventually synchronous consensus algorithms. Partha Dutta, Rachid Guerraoui, Bastian Pochon |
SIAM J. Comput. | 2 |
| 2007 | Gossip-based peer samplingabstractGossip-based communication protocols are appealing in large-scale distributed applications such as information dissemination, aggregation, and overlay topology management. This paper factors out a fundamental mechanism at the heart of all these protocols: the peer-sampling service. In short, this service provides every node with peers to gossip with. We promote this service to the level of a first-class abstraction of a large-scale distributed system, similar to a name service being a first-class abstraction of a local-area system. We present a generic framework to implement a peer-sampling service in a decentralized manner by constructing and maintaining dynamic unstructured overlays through gossiping membership information itself. Our framework generalizes existing approaches and makes it easy to discover new ones. We use this framework to empirically explore and compare several implementations of the peer-sampling service. Through extensive simulation experiments we show that---although all protocols provide a good quality uniform random stream of peers to each node locally---traditional theoretical assumptions about the randomness of the unstructured overlays as a whole do not hold in any of the instances. We also show that different design decisions result in severe differences from the point of view of two crucial aspects: load balancing and fault tolerance. Our simulations are validated by means of a wide-area implementation. Márk Jelasity, Spyros Voulgaris, Rachid Guerraoui, Anne-Marie Kermarrec, Maarten van Steen |
ACM Trans. Comput. Syst. | 3 |
| 2006 | When Birds Die: Making Population Protocols Fault-Tolerant
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Eric Ruppert |
DCOSS | 3 |
| 2006 | High Throughput Total Order Broadcast for Cluster EnvironmentsabstractTotal order broadcast is a fundamental communication primitive that plays a central role in bringing cheap software-based high availability to a wide array of services. This paper studies the practical performance of such a primitive on a cluster of homogeneous machines. We present FSR, a (uniform) total order broadcast protocol that provides high throughput, regardless of message broadcast patterns. FSR is based on a ring topology, only relies on point-to-point inter-process communication, and has a linear latency with respect to the total number of processes in the system. Moreover, it is fair in the sense that each process has an equal opportunity of having its messages delivered by all processes. On a cluster of Itanium based machines, FSR achieves a throughput of 79 Mbit/s on a 100 Mbit/s switched Ethernet network Rachid Guerraoui, Ron R. Levy, Bastian Pochon, Vivien Quéma |
DSN | 1 |
| 2006 | Lucky Read/Write Access to Robust Atomic StorageabstractThis paper establishes tight bounds on the best-case time-complexity of distributed atomic read/write storage implementations that tolerate worst-case conditions. We study asynchronous robust implementations where a writer and a set of reader processes (clients) access an atomic storage implemented over a set of 2t+b+1 server processes of which t can fail: b of these can be malicious and the rest can crash. We define a lucky operation (read or write) as one that runs synchronously and without contention. It is often argued in practice that lucky operations are the most frequent. We determine the exact conditions under which a lucky operation can be fast, namely expedited in onecommunication round-trip with no data authentication. We show that every lucky write (resp., read) can be fast despite fw(resp., fr) actual failures, if and only if fw + fr \lt t-b. Rachid Guerraoui, Ron R. Levy, Marko Vukolic |
DSN | 1 |
| 2006 | Of Malicious Motes and Suspicious Sensors: On the Efficiency of Malicious Interference in Wireless Networks
Seth Gilbert, Rachid Guerraoui, Calvin C. Newport |
OPODIS | 2 |
| 2006 | A Topological Treatment of Early-Deciding Set-Agreement
Rachid Guerraoui, Maurice Herlihy, Bastian Pochon |
OPODIS | 1 |
| 2006 | GosSkip, an Efficient, Fault-Tolerant and Self Organizing Overlay Using Gossip-based Construction and Skip-Lists PrinciplesabstractThis paper presents GosSkip, a self organizing and fully distributed overlay that provides a scalable support to data storage and retrieval in dynamic environments. The structure of GosSkip, while initially possibly chaotic, eventually matches a perfect set of Skip-list-like structures, where no hash is used on data attributes, thus preserving semantic locality and permitting range queries. The use of epidemic-based protocols is the key to scalability, fairness and good behavior of the protocol under churn, while preserving the simplicity of the approach and maintaining O(log(N)) state per peer and O(log(N)) routing costs. In addition, we propose a simple and efficient mechanism to exploit the presence of multiple data items on a single physical node. GosSkip's behavior in both a static and a dynamic scenario is further conveyed by experiments with an actual implementation and real traces of a peer to peer workload Rachid Guerraoui, Sidath B. Handurukande, Kévin Huguenin, Anne-Marie Kermarrec, Fabrice Le Fessant, Etienne Rivière |
Peer-to-Peer Computing | 1 |
| 2006 | Synchronizing without locks is inherently expensiveabstractIt has been considered bon ton to blame locks for their fragility, especially since researchers identified obstruction-freedom: a progress condition that precludes locking while being weak enough to raise the hope for good performance. This paper attenuates this hope by establishing lower bounds on the complexity of obstructionfree implementations in contention-free executions: those where obstruction-freedom was precisely claimed to be effective. Through our lower bounds, we argue for an inherent cost of concurrent computing without locks. We first prove that obstruction-free implementations of a large class of objects, using only overwriting or trivial primitives in contention-free executions, have Omega(n) space complexity and Omega(log^2 n) (obstruction-free) step complexity. These bounds apply to implementations of many popular objects, including variants of fetch&add, counter, compare&swap, and LL/SC. When arbitrary primitives can be applied in contention-free executions, we show that, in any implementation of binary consensus, or any perturbable object, the number of distinct base objects accessed and memory stalls incurred by some process in a contention free execution is Omega(sqrt{n}). All these results hold regardless of the behavior of processes after they become aware of contention. We also prove that, in any obstruction-free implementation of a perturbable object in which processes are not allowed to fail their operations, the number of memory stalls incurred by some process that is unaware of contention is Omega(n). Hagit Attiya, Rachid Guerraoui, Danny Hendler, Petr Kuznetsov |
PODC | 2 |
| 2006 | Towards a theory of transactional contention managersabstractNo abstract available. Rachid Guerraoui, Maurice Herlihy, Bastian Pochon |
PODC | 1 |
| 2006 | How fast can a very robust read be?abstractThis paper studies the time complexity of reading unauthenticated data from a distributed storage made of a set of failure-prone base objects. More specifically, we consider the abstraction of a robust read/write storage that provides wait-free access to unauthenticated data over a set of base storage objects with t possible failures, out of which at most b are arbitrary and the rest are simple crash failures.We prove a 2 communication round-trip lower bound for reading from a safe storage that uses at most 2t+2b base objects, independently of the number or round-trips needed by the writer. We then prove the lower bound tight by exhibiting a regular storage that uses 2t+b+1 base objects (optimal resilience) and features 2 communication round-trips for both read and write operations. Rachid Guerraoui, Marko Vukolic |
PODC | 1 |
| 2006 | Unconscious Eventual Consistency with Gossips
Roberto Baldoni, Rachid Guerraoui, Ron R. Levy, Vivien Quéma, Sara Tucci Piergiovanni |
SSS | 2 |
| 2006 | A General Characterization of Indulgence
Rachid Guerraoui, Nancy A. Lynch |
SSS | 1 |
| 2006 | The Weakest Failure Detectors to Boost Obstruction-Freedom
Rachid Guerraoui, Michal Kapalka, Petr Kuznetsov |
DISC | 1 |
| 2006 | Editorial: Introduction
Rachid Guerraoui |
Distributed Comput. | 1 |
| 2005 | How Fast Can Eventual Synchrony Lead to Consensus?abstractIt is well known that the consensus problem can be solved in a distributed system if, after some time T/sub S/, no process fails and there is some upper bound /spl delta/ on how long it takes to deliver a message. We know of no existing algorithm that guarantees consensus among N processes before time T/sub S/+O(N/spl delta/). We show that consensus can be achieved by time T/sub S/+O(/spl delta/). Partha Dutta, Rachid Guerraoui, Leslie Lamport |
DSN | 2 |
| 2005 | The Impossibility of Boosting Distributed Service ResilienceabstractWe prove two theorems saying that no distributed system in which processes coordinate using reliable registers and f-resilient services can solve the consensus problem in the presence of f + 1 undetectable process stopping failures. (A service is f-resilient if it is guaranteed to operate as long as no more than f of the processes connected to it fail.) Our first theorem assumes that the given services are atomic objects, and allows any connection pattern between processes and services. In contrast, we show that it is possible to boost the resilience of systems solving problems easier than consensus: the k-set consensus problem is solvable for 2k - 1 failures using 1-resilient consensus services. The first theorem and its proof generalize to the larger class of failure-oblivious services. Our second theorem allows the system to contain failure-aware services, such as failure detectors, in addition to failure-oblivious services; however, it requires that each failure-aware service be connected to all processes. Thus, f + 1 process failures overall can disable all the failure-aware services. In contrast, it is possible to boost the resilience of a system solving consensus if arbitrary patterns of connectivity are allowed between processes and failure-aware services: consensus is solvable for any number of failures using only 1-resilient 2-process perfect failure detectors Paul C. Attie, Rachid Guerraoui, Petr Kuznetsov, Nancy A. Lynch, Sergio Rajsbaum |
ICDCS | 2 |
| 2005 | Frugal Event Dissemination in a Mobile Environment
Sébastien Baehni, Chirdeep Singh Chhabra, Rachid Guerraoui |
Middleware | 3 |
| 2005 | Toward a theory of transactional contention managersabstractIn recent software transactional memory proposals, a contention manager module is responsible for ensuring that the system as a whole makes progress. A number of contention manager algorithms have been proposed and empirically evaluated.In this paper we lay some foundations for a theory of contention management. We present the greedy contention manager, the first to combine non-trivial provable properties with good practical performance.In a model where transaction delays are finite, the greedy manager guarantees that every transaction commits within a bounded time, and the time to complete n concurrent transactions that share s objects is within a factor of s(s+1)/2 of the time that would have been taken by an optimal off-line list scheduler. No contention manager reviewed in the literature satisfies both the properties. Benchmark results convey our claim of the practicality of the greedy manager. Rachid Guerraoui, Maurice Herlihy, Bastian Pochon |
PODC | 1 |
| 2005 | From a static impossibility to an adaptive lower bound: the complexity of early deciding set agreementabstractSet agreement, where processors decisions constitute a set of outputs, is notoriously harder to analyze than consensus where the decisions are restricted to a single output. This is because the topological questions that underly set agreement are not about simple connectivity as in consensus. Analyzing set agreement inspired the discovery of the relation between topology and distributed algorithms, and consequently the impossibility of asynchronous set agreement.Yet, the application of topological reasoning has been to the static case, that of asynchronous and synchronous tasks. It is not known yet for example, how to characterize starvation-free solvability of non-terminating tasks. Non-terminating tasks are dynamic entities with no defined end. In a similar vain, early deciding synchronous set agreement, in which the number of rounds it takes a processor to decide adapts to the actual number of failures, falls in this category of dynamic entities.This paper develops a simulation technique that brings to bear topological results to deal with the dynamic situation that arises with early decisions. The novelty of the new simulation is the ability of simulators to look back at the transcript of past rounds of the simulation to influence their current behavior.Using our new technique, we not only re-derive past results, but we propose and prove a lower bound to synchronous early stopping set agreement. We then provide an algorithm to match the lower bound. Our technique uses the BG simulation, in the most creative way it was used to-date, to obtain a rather simple reduction from a static asynchronous impossibility. This reduction is a simple alternative to yet unknown topological argument, and in fact may suggest the way of finding such an argument. Eli Gafni, Rachid Guerraoui, Bastian Pochon |
STOC | 2 |
| 2005 | Computing with Reads and Writes in the Absence of Step Contention
Hagit Attiya, Rachid Guerraoui, Petr Kuznetsov |
DISC | 2 |
| 2005 | (Almost) All Objects Are Universal in Message Passing Systems
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui |
DISC | 3 |
| 2005 | Polymorphic Contention Management
Rachid Guerraoui, Maurice Herlihy, Bastian Pochon |
DISC | 1 |
| 2005 | What Can Be Implemented Anonymously?
Rachid Guerraoui, Eric Ruppert |
DISC | 1 |
| 2005 | The inherent price of indulgence
Partha Dutta, Rachid Guerraoui |
Distributed Comput. | 2 |
| 2005 | Reliable and total order broadcast in the crash-recovery model
Romain Boichat, Rachid Guerraoui |
J. Parallel Distributed Comput. | 2 |
| 2005 | Mutual exclusion in asynchronous systems with failure detectors
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Petr Kuznetsov |
J. Parallel Distributed Comput. | 3 |
| 2004 | Data-Aware MulticastabstractThis paper presents a multicast algorithm for peer-to-peer dissemination of events in a distributed topic-based publish-subscribe system, where processes publish events of certain topics, organized in a hierarchy, and expect events of topics they subscribed to. Our algorithm is "data-aware" in the sense that it exploits information about process subscriptions and topic inclusion relationships to build dynamic groups of processes and efficiently manage the flow of information within and between these process groups. This "data-awareness" helps limit the membership information that each process needs to maintain and preserves processes from receiving messages related to topics they have not subscribed to. It also provides the application with means to control, for each topic in a hierarchy, the trade-off between the message complexity and the reliability of event dissemination. We convey this trade-off through both analysis and simulation. Sébastien Baehni, Patrick Eugster, Rachid Guerraoui |
DSN | 3 |
| 2004 | Linguistic Support for Distributed Programming AbstractionsabstractWe contribute to addressing context of Java and the type-based publish/subscribe (TPS) abstraction, an object-oriented variant of the publish/subscribe paradigm. We present an experience that compares implementations of TPS in (1) a variant of Java we designed to inherently support TPS, (2) standard Java, and (3) Java augmented with genericity. We derive from our implementation experience general observations on what features a programming language should support in order to enable a satisfactory library implementation of TPS, and finally, also alternative abstractions. In particular, we (re-) insist here on the importance of providing genericity and reflective features in the language, and point out the very fact that current efforts towards providing such features are still insufficient. Christian Heide Damm, Patrick Eugster, Rachid Guerraoui |
ICDCS | 3 |
| 2004 | D-Reliable Broadcast: A Probabilistic Measure of Broadcast ReliabilityabstractWe introduce a new probabilistic specification of reliable broadcast communication primitives, called /spl Delta/ - reliable broadcast. This specification captures in a precise way the reliability of practical broadcast algorithms that, on the one hand, were devised with some form of reliability in mind but, on the other hand, are not considered reliable according to "traditional" reliability specifications. We illustrate the use of our specification by precisely measuring and comparing the reliability of two popular broadcast algorithms, namely bimodal multicast and IP multicast. In particular, we quantify how the reliability of each algorithm scales with the size of the system. Patrick Eugster, Rachid Guerraoui, Petr Kuznetsov |
ICDCS | 2 |
| 2004 | Robust Emulations of Shared Memory in a Crash-Recovery ModelabstractA shared memory abstraction can be robustly emulated over an asynchronous message passing system where any process can fail by crashing and possibly recover (crash-recovery model), by having (a) the processes exchange messages to synchronize their read and write operations and (b) log key information on their local stable storage. This paper extends the existing atomicity consistency criterion defined for multiwriter/multireader shared memory in a crash-stop model, by providing two new criteria for the crash-recovery model. We introduce lower bounds on the log-complexity for each of the two corresponding types of robust shared memory emulations. We demonstrate that our lower bounds are tight by providing algorithms that match them. Besides being optimal, these algorithms have the same message and time complexity as their most efficient counterpart we know of in the crash-stop model. Rachid Guerraoui, Ron R. Levy |
ICDCS | 1 |
| 2004 | Towards Safe Distributed Application DevelopmentabstractDistributed application development is overly tedious, as the dynamic composition of distributed components is hard to combine with static safety with respect to types (type safety) and data (encapsulation). Achieving such safety usually goes through specific compilation to generate the glue between components, or making use of a single programming language for all individual components with a hardwired abstraction for the distributed interaction. In this paper, we investigate general-purpose programming language features for supporting third-party implementations of programming abstractions for distributed interaction among components. We report from our experiences in developing a stock market application based on type-based publish/subscribe (TPS) implemented (1) as a library in standard Java as well as with (2) a homegrown extension of the Java language augmented with specific primitives for TPS, motivated by the lacks of former implementation. We then revisit the library approach, investigating the impact of genericity, reflective features, and the type system, on the implementation of a satisfactory TPS library. We then discuss the impact of these features also on other distributed programming abstractions, and hence on the engineering of distributed applications in general, pointing out lacks of mainstream programming environments such as Java as well as .NET. Patrick Eugster, Christian Heide Damm, Rachid Guerraoui |
ICSE | 3 |
| 2004 | The Peer Sampling Service: Experimental Evaluation of Unstructured Gossip-Based Implementations
Márk Jelasity, Rachid Guerraoui, Anne-Marie Kermarrec, Maarten van Steen |
Middleware | 2 |
| 2004 | The weakest failure detectors to solve certain fundamental problems in distributed computingabstractWe determine the weakest failure detectors to solve several fundamental problems in distributed message-passing systems, for all environments -- i.e., regardless of the number and timing of crashes. The problems that we consider are: implementing an atomic register, solving consensus, solving quittable consensus (a variant of consensus in which processes have the option to decide 'quit' if a failure occurs), and solving non-blocking atomic commit. Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Vassos Hadzilacos, Petr Kuznetsov, Sam Toueg |
PODC | 3 |
| 2004 | How fast can a distributed atomic read be?abstractThis paper addresses the problem of designing an efficient implementation of a basic atomic read-write data structure over an asynchronous message-passing system. In particular, we consider time-efficient implementations of this abstraction in the case of a single writer, multiple readers (also called a SWMR atomic register) and S servers: the writer, the readers, and t out of the S servers may fail by crashing. Previous implementations tolerate the failure of any minority of servers (i.e., t < S/2) and require one communication round-trip for every write, and two round-trips for every read.We investigate the possibility of fast implementations, namely, implementations that complete both reads and writes in one round-trip. We show that, interestingly, the existence of a fast implementation depends on the maximum number of readers considered. More precisely, we show that a fast implementation is possible if and only if the number of readers is less that S t-2. We also show that a fast implementation is impossible in a multiple writers setting when t ≥ 1.Our results draw sharp lines between the time-complexity of regular and atomic register implementations, as well as between single-writer and multi-writer implementations. The results lead also to revisit, in a message-passing context, the folklore theorem that "atomic reads must write". Partha Dutta, Rachid Guerraoui, Ron R. Levy, Arindam Chakraborty |
PODC | 2 |
| 2004 | Fast non-blocking atomic commit: an inherent trade-off
Partha Dutta, Rachid Guerraoui, Bastian Pochon |
Inf. Process. Lett. | 2 |
| 2004 | The Information Structure of Indulgent ConsensusabstractTo solve consensus, distributed systems have to be equipped with oracles such as a failure detector, a leader capability, or a random number generator. For each oracle, various consensus algorithms have been devised. Some of these algorithms are indulgent toward their oracle in the sense that they never violate consensus safety, no matter how the underlying oracle behaves. We present a simple and generic indulgent consensus algorithm that can be instantiated with any specific oracle and be as efficient as any ad hoc consensus algorithm initially devised with that oracle in mind. The key to combining genericity and efficiency is to factor out the information structure of indulgent consensus executions within a new distributed abstraction, which we call "Lambda". Interestingly, identifying this information structure also promotes a fine-grained study of the inherent complexity of indulgent consensus. We show that instantiations of our generic algorithm with specific oracles, or combinations of them, match lower bounds on oracle-efficiency, zero-degradation, and one-step-decision. We show, however, that no leader or failure detector-based consensus algorithm can be, at the same time, zero-degrading and configuration-efficient. Moreover, we show that leader-based consensus algorithms that are oracle-efficient are inherently zero-degrading, but some failure detector-based consensus algorithms can be both oracle-efficient and configuration-efficient. These results highlight some of the fundamental trade offs underlying each oracle. Rachid Guerraoui, Michel Raynal |
IEEE Trans. Computers | 1 |
| 2003 | Adaptive Gossip-Based BroadcastabstractThis paper presents a novel adaptation mechanism that allows every node of a gossip-based broadcast algorithm to adjust the rate of message emission 1) to the amount of resources available to the nodes within the same broadcast group and 2) to the global level of congestion in the system. The adaptation mechanism can be applied to all gossip-based broadcast algorithms we know of and makes their use more realistic in practical situations where nodes have limited resources whose quantity changes dynamically with time without decreasing the reliability. 1 Luís E. T. Rodrigues, Sidath B. Handurukande, José Pereira 0001, Rachid Guerraoui, Anne-Marie Kermarrec |
DSN | 4 |
| 2003 | An Equational Theory for Transactions
Andrew P. Black, Vincent Cremet, Rachid Guerraoui, Martin Odersky |
FSTTCS | 3 |
| 2003 | Pragmatic Type InteroperabilityabstractProviding type interoperability consists in ensuring that, even if written by different programmers, possibly in different languages and running on different platforms, types that are supposed to represent the same software module are indeed treated as one single type. This form of interoperability is crucial in modern distributed programming. We present a pragmatic approach to deal with type interoperability in a dynamic and distributed environment. Our approach is based on an optimistic transport protocol, specific serialization mechanisms and a set of implicit type conformance rules. We experiment the approach over the .NET platform which we indirectly evaluate. Sébastien Baehni, Patrick Eugster, Rachid Guerraoui, Philippe Altherr |
ICDCS | 3 |
| 2003 | A Generic Framework for Indulgent ConsensusabstractConsensus is a fundamental distributed agreement problem that has to be solved when one has to design or implement reliable applications. As consensus cannot be solved in pure asynchronous distributed systems, those systems have to be equipped with appropriate oracles to circumvent the impossibility. Several oracles (unreliable failure detector leader capability, random number generator) have been proposed, and consensus protocols based on such ad hoc oracles have been designed This paper presents a generic consensus framework that can be instantiated with any oracle, or combination of oracles, that satisfies a set of properties. This generic framework provides indulgent consensus protocols that are particularly simple and efficient both in well-behaved runs (i.e., when there are no failures), and in stable runs (i.e., when there is no failure during the execution although some processes can be initially crashed). In those runs, the protocols terminate in two communication steps (which is optimal). Indulgence means that the resulting protocol never violates its safety property even when the underlying oracle behaves arbitrarily. Interestingly, the protocol can also allow processes to decide in one communication step in some specific configurations. Rachid Guerraoui, Michel Raynal |
ICDCS | 1 |
| 2003 | Distributed Programming for Dummies: A Shifting Transformation TechniqueabstractThe perfectly synchronized round model provides the powerful abstraction of crash-stop failures with atomic message delivery. This abstraction makes distributed programming very easy. We present an implementation of this abstraction in a distributed system with general message omissions. Protocols devised using our abstraction (i.e., in the perfectly synchronized round model) are automatically transformed into protocols for the omission model. The transformation is achieved using a round shifting technique with a constant time complexity overhead. This transformation is in a precise sense optimal. Furthermore, and rather surprisingly, no automatic transformation from a weaker model, say the traditional crash-stop model (with no atomic message delivery), onto an even stronger model than the general-omission one, say the send-omission model, can provide better time complexity performance. Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui, Bastian Pochon |
SRDS | 3 |
| 2003 | Tight Bounds on Early Local Decisions in Uniform Consensus
Partha Dutta, Rachid Guerraoui, Bastian Pochon |
DISC | 2 |
| 2003 | On Failure Detectors and Type Boosters
Rachid Guerraoui, Petr Kuznetsov |
DISC | 1 |
| 2003 | The Database State Machine Approach
Fernando Pedone, Rachid Guerraoui, André Schiper |
Distributed Parallel Databases | 2 |
| 2003 | Editorial: MiddlewareabstractMiddlewareIn theory, the term middleware denotes any software that can be used for more than one application and more than one piece of hardware.In practice, the term middleware tends to denote software abstractions for distributed computing, including communication abstractions such as RPC, and reliability abstractions such as transactions.Devising and implementing such abstractions has constituted an active area of research in the last decade, bridging the gap between various fields, including programming languages, distributed systems, networking and databases.This special issue gathers four middleware papers.The first paper is on the customization of quality of service requirements in distributed object systems, by J. He, M. Hiltunen, M. Rajagopalan and R. Schichting.The paper presents an interceptor-based architecture to enable transparent quality of service enhancements of a distributed application.The second paper is about thread transparency in information flow middleware, by R. Koster, A. Black, J. Huang, J. Walpole and C. Pu.It describes Infopipes, a high-level abstraction for writing information flow applications, which are hard to write with RPC-like abstractions.The third paper is about a CORBA activity service framework for supporting extended transactions, by I. Houston, M. Little, S. Shrivastava and S. Wheather.The service aims at structuring long-lived applications using an underlying signalling mechanism.The fourth paper is about access control and trust in the use of widely distributed services by J. Bacon, K. Moody and W. Yao.It describes OASIS, a role-based access control architecture for achieving secure interoperation of independently managed services.These papers are extended and revised versions of four papers selected out of the proceedings of Middleware 2001 (Lecture Notes in Computer Science, vol.2218), which was organized in November 2001 in Heidelberg.The proceedings gathered 20 papers, themselves selected out of 116 submissions.Every paper was reviewed by at least three members of the program committee.The selection of the papers was performed during a program committee meeting, held in Lausanne on 6 July 2001.The papers were judged according to their originality, presentation quality and relevance to the conference topics. Rachid Guerraoui |
Softw. Pract. Exp. | 1 |
| 2003 | Lightweight probabilistic broadcastabstractGossip-based broadcast algorithms, a family of probabilistic broadcast algorithms, trade reliability guarantees against "scalability" properties. Scalability in this context has usually been expressed in terms of message throughput and delivery latency, but there has been little work on how to reduce the memory consumption for membership management and message buffering at large scale.This paper presents lightweight probabilistic broadcast ( lpbcast ), a novel gossip-based broadcast algorithm, which complements the inherent throughput scalability of traditional probabilistic broadcast algorithms with a scalable memory management technique. Our algorithm is completely decentralized and based only on local information: in particular, every process only knows a fixed subset of processes in the system and only buffers fixed "most suitable" subsets of messages. We analyze our broadcast algorithm stochastically and compare the analytical results both with simulations and concrete implementation measurements. Patrick Eugster, Rachid Guerraoui, Sidath B. Handurukande, Petr Kuznetsov, Anne-Marie Kermarrec |
ACM Trans. Comput. Syst. | 2 |
| 2003 | Guest Editorial: Special Section on Middleware InfrastructuresabstractTHIS special section of Transactions on Parallel and Distributed Systems is devoted to middleware infrastructures and gathered nine papers. The first paper is about refactoring middleware with aspects and is authored by C. Zhang and H. Jacobson. The paper is a case for the introduction of aspect-oriented programming techniques within object request brokers. The second paper describes an energy-efficient object discovery protocol for context-sensitive middleware for ubiquitous computing, and is authored by S. Yau and F. Karim. The paper introduces a technique for discovering objects in a distributed environment that is efficient in terms of energy consumption. The third paper describes a middleware platform called OBIWAN, and is authored by P. Ferreira, L. Veiga, and C. Ribeiro. The platform performs automatic creation of object replicas (e.g., incremental on-demand replication) as well as garbage collection of useless objects. The fourth paper presents a middleware infrastructure for parallel and distributed programming models on heterogeneous systems, and is authored by J. Al-Jaroodi, N. Mohamed, H. Jiang, and D. Swanson. The infrastructure handles class loading and distributed deployment in a transparent manner. The fifth paper is about an adaptive quality-of-service aware middleware for replicated services, and is authored by S. Krishnamurthy, W. Sanders, and M. Cukier. The idea is to provide the clients of a replicated service the ability to specify temporal and consistency requirements and have the server adjust its replication strategy according to these requirements. The sixth paper describes an OCI-based group communication support for CORBA, and is authored by D. Lee, D. Nam, H. Youn, and C. Yu. The paper presents a way to transparently enhance CORBA with group communication and object group management primitives. The seventh paper introduces a cluster programming middleware for streamorientedapplications, and is authored by U. Ramachandra, R. Nikhil, J. Rehg, Y. Angelov, A. Paul, S.Adhikari,K.MacKenzie,N.Harel, andK.Knobe.Thepaper presents a middleware infrastructure with adequate support for data abstractions, dynamic cluster-wide threads, data parallelism, and multiple address spaces. Theeighthpaper focuseson thedesignandperformanceof real-time Javamiddleware and is authored byA.Corsaro and D. Schmidt. The paper describes an open-source implementation of the real-time specification for Java middleware and corresponding performance measures. The ninth paper describes clustering support and replication management for scalable network services, and is authored by K. Shen, T. Yang, and L. Chu. The presented middleware infrastructure, named Neptune, employs a loosely connected and functionally symmetric clustering to achieve scalability and robustness. We are extremely grateful to all the reviewers who provided very useful feedback to select the papers and improve their presentation, as well as to all of the authors of the submitted papers for their interest in this special section. Rachid Guerraoui, Willy Zwaenepoel |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2002 | A Realistic Look At Failure DetectorsabstractThis paper shows that, in an environment where we do not bound the number of faulty processes, the class P of perfect failure detectors is the weakest (among realistic failure detectors) to solve fundamental agreement problems like uniform consensus, atomic broadcast, and terminating reliable broadcast (also called Byzantine generals). Roughly speaking, in this environment, we collapse the Chandra-Toueg failure detector hierarchy, by showing that P ends up being the only class to solve those agreement problems. This contributes in explaining why most reliable distributed systems we know of do rely on some group membership service that precisely aims at emulating P. As an interesting side effect of our work, we show that, in our general environment, uniform consensus is strictly harder than consensus, and we revisit the view that uniform consensus and atomic broadcast are strictly weaker than terminating reliable broadcast. Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui |
DSN | 3 |
| 2002 | Probabilistic MulticastabstractGossip-based broadcast algorithms have been considered as a viable alternative to traditional deterministic reliable broadcast algorithms in large scale environments. However, these algorithms focus on broadcasting events inside a large group of processes, while the multicasting of events to a subset of processes in a group only, potentially varying for every event, has not been considered. We propose a scalable gossip-based multicast algorithm which ensures, with a high probability, that (1) a process interested in a multicast event delivers that event (just like in typical gossip-based broadcast algorithms), and that (2) a process not interested in that event does not receive it (unlike in broadcast algorithms). Patrick Eugster, Rachid Guerraoui |
DSN | 2 |
| 2002 | AOP: Does It Make Sense? The Case of Concurrency and Failures
Jörg Kienzle, Rachid Guerraoui |
ECOOP | 2 |
| 2002 | The LEAF Platform: Incremental Enhancements for the J2EEabstractLeaf, the Lean and Extensible Architectural Framework is an enhancement wrapper for J2EE implementations. Basically, LEAF fixes some identified J2EE issues and extends, as well as simplifies. the use of the J2EE by providing several incremental improvements. These improvements are seamlessly integrated, include an additional component type, allow the same interfaces for local and remote service implementations offer better J2EE implementation compatibility and ORB interceptors, and encompass several new technical services. This paper explains the need for LEAF through a diagnosis of the J2EE, presents the fundamental concepts underlying LEAF, overviews its implementation, reports on field experiences from using it in a number of commercial projects, and points out some interesting tradeoffs in using the J2EE with and without LEAF. Philipp Oser, Christian Gasser, Daniel Gorostidi, Rachid Guerraoui |
EDOC | 4 |
| 2002 | OS Support for P2P Programming: a Case for TPSabstractJust as the remote procedure call (RPC) turned out to be a very effective OS abstraction in building client-server applications over LANs, type-based publish-subscribe (TPS) can be viewed as a high-level candidate abstraction for building peer-to-peer (P2P) applications over WANs. This paper relates our preliminary, though positive, experience of implementing and using TPS over JXTA, which can be viewed as the P2P counterpart to sockets. We show that, at least for P2P applications with the Java type model, TPS provides a high-level programming support that ensures type safety and encapsulation, without hampering the decoupled nature of these applications. Furthermore, the loss of flexibility (inherent to the use of any high level abstraction) and the performance overhead, are negligible with respect to the simplicity gained by using TPS. Sébastien Baehni, Patrick Eugster, Rachid Guerraoui |
ICDCS | 3 |
| 2002 | The inherent price of indulgenceabstractThis paper presents a tight lower bound on the time complexity of indulgent consensus algorithms, i.e., consensus algorithms that use unreliable failure detectors. We state and prove our tight lower bound in the unifying framework of round-by-round fault detectors.We show that any ⋄P-based t-resilient consensus algorithm requires at least t + 2 rounds for a global decision even in runs that are synchronous. We then prove the bound to be tight by exhibiting a new ⋄P-based t-resilient consensus algorithm that reaches a global decision at round t + 2 in every synchronous run. Our new algorithm is in this sense significantly faster than the most efficient indulgent algorithm we knew of (which requires 2t + 2 rounds).We contrast our lower bound with the well-known t + 1 round tight lower bound on consensus for the synchronous model, pointing out the price of indulgence. Partha Dutta, Rachid Guerraoui |
PODC | 2 |
| 2002 | Failure Detection Lower Bounds on Registers and Consensus
Carole Delporte-Gallet, Hugues Fauconnier, Rachid Guerraoui |
DISC | 3 |
| 2002 | An Efficient Universal Construction for Message-Passing Systems
Partha Dutta, Svend Frølund, Rachid Guerraoui, Bastian Pochon |
DISC | 3 |
| 2002 | Non-blocking atomic commit in asynchronous distributed systems with failure detectors
Rachid Guerraoui |
Distributed Comput. | 1 |
| 2002 | Dictatorial Transaction Processing: Atomic Commitment Without Veto Right
Maha Abdallah, Rachid Guerraoui, Philippe Pucheral |
Distributed Parallel Databases | 2 |
| 2002 | e-Transactions: End-to-End Reliability for Three-Tier ArchitecturesabstractA three-tier application is organized as three layers: human users interact with front-end clients (e.g., browsers), middle-tier application servers (e.g., Web servers) contain the business logic of the application, and perform transactions against back-end databases. Although three-tier applications are becoming mainstream, they usually fail to provide sufficient reliability guarantees to the users. Usually, replication and transaction-processing techniques are applied to specific parts of the application, but their combination does not provide end-to-end reliability. The aim of this paper is to provide a precise specification of a desirable, yet realistic, end-to-end reliability contract in three-tier applications. The paper presents the specification in the form of the Exactly-Once Transaction (e-Transaction) abstraction: an abstraction that encompasses both safety and liveness properties in three-tier environments. It gives an example implementation of that abstraction and points out alternative implementations and tradeoffs. Svend Frølund, Rachid Guerraoui |
IEEE Trans. Software Eng. | 2 |
| 2001 | Lightweight Probabilistic BroadcastabstractThe growing interest in peer-to-peer applications has underlined the importance of scalability in modern distributed systems. Not surprisingly, much research effort has been invested in gossip-based broadcast protocols. These trade the traditional strong reliability guarantees against very good "scalability" properties. Scalability is in that context usually expressed in terms of throughput and delivery latency, but there is only little work on how to reduce the overhead of membership management on a large scale. The paper presents Lightweight Probabilistic Broadcast (lpbcast), a novel gossip-based broadcast algorithm which preserves the inherent throughput scalability of traditional gossip-based algorithms and adds a notion of membership management scalability: every process only knows a random subset of fixed size of the processes in the system. We formally analyze our broadcast algorithm in terms of scalability with respect to the size of individual views, and compare the analytical results both with simulations and concrete measurements. Patrick Eugster, Rachid Guerraoui, Sidath B. Handurukande, Petr Kuznetsov, Anne-Marie Kermarrec |
DSN | 2 |
| 2001 | On Objects and EventsabstractThis paper presents linguistic primitives for publish/subscribe programming using events and objects. We integrate our primitives into a strongly typed object-oriented language through four mechnisms: (1) serialization, (2) multiple subtyping, (3)closures, and (4) deferred code evaluation. We illustrate our primitives through Java, showing how we have overcome its respective lacks. A precompiler transforms statements based on our publish/subscribe primitives into calls to specifically generated typed adapters, which resemble the typed stubs and skeletons by the rmic precompiler for remote method invocations in Java Patrick Eugster, Rachid Guerraoui, Christian Heide Damm |
OOPSLA | 2 |
| 2001 | Reducing Noise in Gossip-Based Reliable BroadcastabstractWe present in this paper a general garbage collection scheme that reduces the "noise" in gossip-based broadcast algorithms. In short, our garbage collection scheme uses a simple heuristic to trade "useless" messages with "useful" ones. Used with a given gossip-based broadcast algorithm, a given size of buffers, and a given number of disseminated messages (e.g., per gossip round), our garbage collection scheme provides higher overall reliability than more conventional schemes. We illustrate our approach through two algorithms: bimodal multicast (pbcast) and lightweight probabilistic broadcast (lpbcast). Our scheme is based on the intuitive idea of discarding messages according to their "age". The "age" of a message represents the number of times the message has been retransmitted. Petr Kuznetsov, Rachid Guerraoui, Sidath B. Handurukande, Anne-Marie Kermarrec |
SRDS | 2 |
| 2001 | Open consensusabstractAbstract This paper presents the abstraction of open consensus and argues for its use as an effective component for building reliable agreement protocols in practical asynchronous systems where processes and links can crash and recover. The specification of open consensus has a decoupled, on‐demand and re‐entrant flavour that make its use very efficient, especially in terms of forced logs, which are known to be major sources of overhead in distributed systems. We illustrate the use of open consensus as a basic building block to develop a modular, yet efficient, total‐order broadcast protocol. Finally, we describe our Java implementation of our open‐consensus abstraction and we convey our efficiency claims through some practical performance measures. Copyright © 2001 John Wiley & Sons, Ltd. Romain Boichat, Svend Frølund, Rachid Guerraoui |
Concurr. Comput. Pract. Exp. | 3 |
| 2001 | Effective multicast programming in large scale distributed systemsabstractAbstract Many distributed applications have a strong requirement for efficient dissemination of large amounts of information to widely spread consumers in large networks. These include applications in e‐commerce and telecommunication. Publish/subscribe is considered one of the most important interaction styles with which to model communication on a large scale. Producers publish information on a topic and consumers subscribe to the topics they wish to be informed of. The decoupling of producers and consumers in time, space, and flow makes the publish/subscribe paradigm very attractive for large scale distribution, especially in environments like the Internet. This paper describes the architecture and implementation of DACE (Distributed Asynchronous Computing Environment), a framework for publish/subscribe communication based on an object‐oriented programming abstraction in the form of Distributed Asynchronous Collection (DAC). DACs capture the variants of publish/subscribe, without blurring their respective advantages. The architecture we present is tolerant of network partitions and crash failures. The underlying model is based on the notion of Topic Membership: a weak membership for the parties involved in a topic. We present how Topic Membership enables the realization of a robust and efficient reliable multicast on a large scale. The protocol ensures that, inside a topic, even a subscriber who is temporarily partitioned away eventually receives a published message. Copyright © 2001 John Wiley & Sons, Ltd. Patrick Eugster, Romain Boichat, Rachid Guerraoui, Joseph S. Sventek |
Concurr. Comput. Pract. Exp. | 3 |
| 2001 | X-Ability: a theory of replication
Svend Frølund, Rachid Guerraoui |
Distributed Comput. | 2 |
| 2001 | On the hardness of failure-sensitive agreement problems
Rachid Guerraoui |
Inf. Process. Lett. | 1 |
| 2001 | Genuine atomic multicast in asynchronous distributed systems
Rachid Guerraoui, André Schiper |
Theor. Comput. Sci. | 1 |
| 2001 | Implementing E-Transactions with Asynchronous ReplicationabstractThis paper describes a distributed algorithm that implements the abstraction of e-Transaction: a transaction that executes exactly-once despite failures. Our algorithm is based on an asynchronous replication scheme that generalizes well-known active-replication and primary-backup schemes. We devised the algorithm with a three-tier architecture in mind: the end-user interacts with front-end clients (e.g., browsers) that invoke middle-tier application servers (e.g., web servers) to access back-end databases. The algorithm preserves the three-tier nature of the architecture and introduces a very acceptable overhead with respect to unreliable solutions. Svend Frølund, Rachid Guerraoui |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2001 | The Generic Consensus ServiceabstractThis paper describes a modular approach for the construction of fault-tolerant agreement protocols. The approach is based on a generic consensus service. Fault-tolerant agreement protocols are built using a client-server interaction, where the clients are the processes that must solve the agreement problem and the servers implement the consensus service. This service is accessed through a generic consensus filter, customized for each specific agreement problem. We illustrate our approach on the construction of various fault-tolerant agreement protocols, such as nonblocking atomic commitment, group membership, view synchronous communication, and total order multicast. Through a systematic reduction to consensus, we provide a simple way to solve agreement problems. In addition to its modularity, our approach enables efficient implementations of agreement protocols and precise characterization of the assumptions underlying their liveness and safety properties. Rachid Guerraoui, André Schiper |
IEEE Trans. Software Eng. | 1 |
| 2000 | Synchronous System and Perfect Failure Detector: Solvability and Efficiency IssueabstractWe compare, in terms of solvability and efficiency, the synchronous model, noted Ss, with the asynchronous model augmented with a perfect failure detector, noted S/sub P/. We first exhibit a problem that, although time-free, is solvable in S/sub S/ but not in S/sub P/. We then examine whether one of these two models allows more efficient solutions for designing fault-tolerant applications. In particular, we concentrate on the uniform consensus problem which is solvable in both models, and we design a uniform consensus algorithm for the S/sub S/ model that is more efficient than any algorithm solving uniform consensus in S/sub P/ with respect to some significant time complexity measure. From a practical viewpoint, the synchronous model thus seems better than the asynchronous model augmented with a perfect failure detector. Bernadette Charron-Bost, Rachid Guerraoui, André Schiper |
DSN | 2 |
| 2000 | Implementing e-Transactions with Asynchronous ReplicationabstractAn e-Transaction is one that executes exactly-once despite failures. This paper describes a distributed protocol that implements the abstraction of e-Transaction in three-tier architectures. Three-tier architectures are typically Internet-oriented architectures, where the end-user interacts with front-end clients (e.g., browsers) that invoke middle-tier application servers (e.g., web servers) to access back-end databases. We implement the e-Transaction abstraction using an asynchronous replication scheme that preserves the three-tier nature of the architecture and introduces a very acceptable overhead with respect to unreliable solutions. Svend Frølund, Rachid Guerraoui |
DSN | 2 |
| 2000 | Distributed Asynchronous Collections: Abstractions for Publish/Subscribe Interaction
Patrick Eugster, Rachid Guerraoui, Joseph S. Sventek |
ECOOP | 2 |
| 2000 | X-ability: a theory of replicationabstractDifferent replication mechanisms provide different solutions to the same basic problem. However, there is no precise specification of the problem itself, only of particular classes of solutions, such as active replication and primary-backup. Having a precise specification of the problem would help us better understand the space of possible solutions. Svend Frølund, Rachid Guerraoui |
PODC | 2 |
| 2000 | Indulgent algorithms (preliminary version)abstractInformally, an indulgent algorithm is a distributed algorithm that tolerates unreliable failure detection: the algorithm is indulgent towards its failure detector. This paper formally characterises such algorithms and states some of their interesting features. We show that indulgent algorithms are inherently safe and uniform. We also state impossibility results for indulgent solutions to divergent problems like consensus, and failure-sensitive problems like non-blocking atomic commit and terminating reliable broadcast. Rachid Guerraoui |
PODC | 1 |
| 2000 | Reliable Broadcast in the Crash-Recovery ModelabstractThe paper addresses the problem of broadcasting messages in a reliable manner within a practical asynchronous system where processes and channels may crash and recover. In this crash-recovery model, we present meaningful specifications of reliable broadcast and we describe algorithms that implement those specifications. Our approach is modular and incremental. It is modular in the sense that we give the properties of reliable broadcast separately, and then consider their composition. It is incremental in the sense that we show how to automatically transform any reliable broadcast algorithm that implements a given specification into one that implements a stronger specification. In particular we show how to reuse, in a crash-recovery model, reliable broadcast algorithms that were initially designed in a simpler crash-stop model. Romain Boichat, Rachid Guerraoui |
SRDS | 2 |
| 2000 | Abstractions for Devising Byzantine-Resilient State Machine ReplicationabstractDCL Assia Doudou, Rachid Guerraoui, Benoît Garbinato |
SRDS | 2 |
| 2000 | A Pragmatic Implementation of e-TransactionsabstractThree-tier applications have nice properties, which make them scalable and manageable: clients are thin and servers are stateless. However, it is challenging to implement or even define end-to-end reliability for such applications. Furthermore, it is especially hard to make these applications reliable without violating their nice properties. In previous work, we identified e-transactions as a desirable and practical end-to-end reliability guarantee for three-tier applications (S. Frolund and R. Guerraoui, 1999). Essentially, an e-transaction guarantees that the server-side transactional side-effect happens exactly once, and that the client receives the result of the server-side computation. Thus, e-transactions mask server and database failures relative to the client. We present a pragmatic implementation of e-transactions that maintains the nice properties of three-tier applications in the special, but very common case of a single back-end database. Svend Frølund, Rachid Guerraoui |
SRDS | 2 |
| 2000 | Special issue: European Conference on Object-oriented Programming 1999
Rachid Guerraoui |
Concurr. Pract. Exp. | 1 |
| 2000 | Experiences with object group systemsabstractThe GARF, Bast, and OGS systems represent the resulting efforts of a multi-year ‘object group’ program at the Swiss Federal Institute of Technology in Lausanne. The intent of the program was to understand the extent to which one could build flexible and performance system supports to encapsulate object plurality. That is, we experimented with various ways to build libraries and services to support object groups in a distributed setting. This paper summarizes the main steps of the efforts and draws some conclusions about the successes and failures of our approaches. Copyright © 2000 John Wiley & Sons, Ltd. Rachid Guerraoui, Patrick Eugster, Pascal Felber, Benoît Garbinato, Karim Mazouni |
Softw. Pract. Exp. | 1 |
| 1999 | Workshop on Reliable Middleware - ForewordabstractPresents the introductory welcome message from the conference proceedings. Pascal Felber, Rachid Guerraoui |
SRDS | 2 |
| 1998 | Exploiting Atomic Broadcast in Replicated Databases
Fernando Pedone, Rachid Guerraoui, André Schiper |
Euro-Par | 2 |
| 1998 | Scalable Atomic MulticastabstractWe present a new scalable fault-tolerant algorithm which ensures total order delivery of messages sent to multiple groups of processes. Our algorithm is particularly well suited for large scale systems because: (1) any process can multicast a message to one or more groups of processes without being forced to join those groups; (2) inter-group total order is ensured system-wide but, for each individual multicast, the number and size of messages exchanged depends only on the number of addressees; (3) process failure detection does not need to be reliable. Our algorithm also exhibits a modular design. It uses two companion protocols, namely a reliable multicast protocol and a consensus protocol, and these protocols are not required to use the same communication channels or to share common variables with the total order protocol. This approach follows a design methodology based on the composition of (encapsulated) micro-protocols. Luís E. T. Rodrigues, Rachid Guerraoui, André Schiper |
ICCCN | 2 |
| 1998 | Flexible Protocol Composition in BastabstractThe paper presents BAST, an object-oriented library of reliable distributed protocols. The authors show how BAST can be used to build fault-tolerant distributed applications, and how new protocols can be added to it. They discuss some distributed protocol design issues and the way these issues are circumvented in BAST. They briefly describe their Smalltalk and Java implementations of BAST, together with some performance results. Benoît Garbinato, Rachid Guerraoui |
ICDCS | 2 |
| 1998 | One-Phase Commit: Does it make Sense?abstractAlthough widely used in distributed transactional systems, the so-called Two-Phase Commit (2PC) protocol introduces a substantial delay in transaction processing, even in the absence of failures. This has led several researchers to look for alternative commit protocols that minimize the time cost associated with coordination messages and forced log writes in 2PC. In particular, variations of a One-Phase Commit (1PC) protocol have recently been proposed. Although efficient, 1PC is however rarely considered in practice because of the strong assumptions it requires from the distributed transactional system. The aim of the paper is to better identify and understand those assumptions. Through a careful look into the intrinsic characteristics of 1PC, we dissect the assumptions underlying it and we present simple techniques that minimize them. We believe that these techniques constitute a first step towards a serious reconsideration of 1PC in the transactional world. Maha Abdallah, Rachid Guerraoui, Philippe Pucheral |
ICPADS | 2 |
| 1998 | System Support for Object GroupsabstractThis paper draws several observations from our experiences in building support for object groups. These observations actually go beyond our experiences and may apply to many other developments of object based distributed systems.Our first experience aimed at building support for Smalltalk object replication using the Isis process group toolkit. It was quite easy to achieve group transparency but we were confronted with a strong mismatch between the rigidity of the process group model and the flexible nature of object interactions. Consequently, we decided to build our own object oriented protocol framework, specifically dedicated to support object groups (instead of using a process group toolkit). We built our framework in such a way that basic distributed protocols, such as failure detection and multicasts, are considered as first class entities, directly accessible to the programmers. To achieve flexible and dynamic protocol composition, we had to go beyond inheritance and objectify distributed algorithms.Our second experience consisted in building a CORBA service aimed at managing group of objects written on different languages and running on different platforms. This experience revealed a mismatch between the asynchrony of group protocols and the synchrony of standard CORBA interaction mechanisms, which limited the portability of our CORBA object group service. We restricted the impact of this mismatch by encapsulating asynchrony issues inside a specific messaging sub-service.We dissect the cost of object group transparency in our various implementations, and we point out the recurrent sources of overheads, namely message indirection, marshaling/unmarshaling and strong consistency. Rachid Guerraoui, Pascal Felber, Benoît Garbinato, Karim Mazouni |
OOPSLA | 1 |
| 1997 | Total Order Multicast to Multiple GroupsabstractWe present a fault tolerant algorithm that ensures total order delivery of messages sent to multiple groups of processes. Our algorithm is a multiple group "genuine" multicast algorithm in the sense that: (1) any process can send a message to any set of process groups; and (2) only the sender and the receivers of a message take part in the algorithm needed to deliver the message. The correctness of our algorithm does not require reliable failure detectors, but requires causal order delivery of messages. This establishes a new and interesting link between causal order delivery and fault tolerance with unreliable failure detectors. Rachid Guerraoui, André Schiper |
ICDCS | 1 |
| 1997 | Transaction Reordering in Replicated DatabasesabstractThe paper presents a fault tolerant lazy replication protocol that ensures 1-copy serializability at a relatively low cost. Unlike eager replication approaches, our protocol enables local transaction execution and does not lead to any deadlock situation. Compared to previous lazy replication approaches, we significantly reduce the abort rate of transactions and we do not require any reconciliation procedure. Our protocol first executes transactions locally, then broadcasts a transaction certification message to all replica managers, and finally employs a certification procedure to ensure 1-copy serializability. Certification messages are broadcast using a non blocking atomic broadcast primitive, which alleviates the need for a more expensive non blocking atomic commitment algorithm. The certification procedure uses a reordering technique to reduce the probability of transaction aborts. Fernando Pedone, Rachid Guerraoui, André Schiper |
SRDS | 2 |
| 1996 | Protocol Classes for Designing Reliable Distributed Environments
Benoît Garbinato, Pascal Felber, Rachid Guerraoui |
ECOOP | 3 |
| 1996 | Reducing the Cost for Non-Blocking in Atomic CommitmentabstractNon-blocking atomic commitment protocols enable a decision (commit or abort) to be reached at every correct participant, despite the failure of others. The cost for non-blocking implies however (1) a high number of messages and communication steps required to reach commit, and (2) a complicated termination protocol needed in the case of failure suspicions. In this paper, we present a non-blocking protocol, called MDSPC (Modular and Decentralized Three Phase Commit), which enables to trade resiliency against efficiency. As conveyed by our performance measures, MDSPC is faster than existing non-blocking protocols, and in the case of a broadcast network and a reasonable resiliency rate (e.g. 2 or 3) is almost as efficient as the classical (blocking) 2PC. The termination protocol of MDSPC is encapsulated inside a majority consensus protocol. This modularity leads to a simple structure of MDSPC and enables a precise characterization of its liveness in an asynchronous system with an unreliable failure detector. Rachid Guerraoui, Mikel Larrea, André Schiper |
ICDCS | 1 |
| 1996 | The Design of a CORBA Group Communication ServiceabstractThe common object request broker architecture (CORBA) is becoming a standard for distributed application middleware, and there are increasing needs for enriching the basic functionalities of CORBA. While mechanisms for persistence, transactions, event channels, etc. have been designed and specified for CORBA, no standard support is provided to handle object replication. In this paper we discuss the issue of augmenting CORBA with group communication, which is considered an adequate paradigm to handle replication. We distinguish two main approaches: the integration approach and the service approach. We argue that the service approach is more appropriate to CORBA as it preserves the modularity of the architecture. We describe a proposal for a group communication service and discuss some implementation issues. Pascal Felber, Benoît Garbinato, Rachid Guerraoui |
SRDS | 3 |
| 1995 | Non-Blocking Atomic Commitment with an Unreliable Failure DetectorabstractIn a transactional system, an atomic commitment protocol ensures that for any transaction, all data manager processes agree on the same outcome (commit or abort). A non-blocking atomic commitment protocol enables an outcome to be decided at every correct process despite the failure of others. In this paper we apply, for the first time, the fundamental result of T. Chandra and S. Toueg (1991) on solving the abstract consensus problem, to non-blocking atomic commitment. More precisely, we present a non-blocking atomic commitment protocol in an asynchronous system augmented with an unreliable failure detector that can make an infinity of false failure suspicions. If no process is suspected to have failed, then our protocol is similar to a three phase commit protocol. In the case where processes are suspected, our protocol does not require any additional termination protocol: failure scenarios are handled within our regular protocol and are thus much simpler to manage. Rachid Guerraoui, Mikel Larrea, André Schiper |
SRDS | 1 |
| 1995 | Nested Transactions: Reviewing the Coherence Contract
Rachid Guerraoui |
Inf. Sci. | 1 |
| 1994 | Atomic Object Composition
Rachid Guerraoui |
ECOOP | 1 |
| 1992 | Nesting Actions through Asynchronous Message Passing: the ACS Protocol
Rachid Guerraoui, Riccardo Capobianchi, Agnes Lanusse |
ECOOP | 1 |