Lorenzo Alvisi

dblp:a/LAlvisi · DBLP profile ↗
← Back
96ranked-venue papers
17as first author
11since 2021 · last 2026
0000-0002-9857-5528ORCID · conflict

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

Systems, architecture and hardware · 36 · 10 first-author · 1 since 2021Security and privacy · 19 · 2 first-author · 2 since 2021Software engineering, systems software and programming languages · 19 · 1 first-author · 4 since 2021Computer networks · 12 · 1 first-author · 1 since 2021Databases, data management, data science and information retrieval · 6 · 1 first-author · 1 since 2021Theory of computation · 2 · 1 first-authorArtificial intelligence and machine learning · 1 · 1 first-author · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 1 · 1 first-authorHuman-computer interaction and ubiquitous computing · 1 · 1 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 1
YearPublicationVenuePosition
2026 Fast Deterministically Safe Proof-of-Work Consensus
Ali Farahbakhsh, Giuliano Losa, Youer Pu, Lorenzo Alvisi
SP4
2025 Modeling Metastability
abstract
Recently, there has been increasing concern about a new failure mode in data-center systems: when there is an external shock, such as a sudden load spike or some machine failures, systems will sometimes respond with reduced throughput - but, in contrast to a traditional overload situation, the throughput does not recover once the external shock disappears, and remains permanently degraded. This phenomenon has been called a metastable failure.
Ali Farahbakhsh, Andreas Haeberlen, Qingjie Lu, Lorenzo Alvisi, Robbert van Renesse, Shir Cohen
HotNets4
2025 Pesto: Cooking up High Performance BFT Queries
abstract
This paper presents Pesto, a high-performance Byzantine Fault Tolerant (BFT) database that offers full SQL compatibility. Pesto intentionally forgoes the use of State Machine Replication (SMR); SMR-based designs offer poor performance due to the several round trips required to order transactions. Pesto, instead, allows for replicas to remain inconsistent, and only synchronizes on demand to ensure that the database remain serializable in the presence of concurrent transactions and malicious actors. On TPC-C, Pesto matches the throughput of Peloton [20] and Postgres [21], two unreplicated SQL database systems, while increasing throughput by 2.3x compared to classic SMR-based BFT-architectures, and reducing latency by 2.7x to 3.9x. Pesto's leaderless design minimizes the impact of replica failures and ensures robust performance.
Florian Suri-Payer, Neil Giridharan, Liam Arzola, Shir Cohen, Lorenzo Alvisi, Natacha Crooks
SOSP5
2024 Unraveling the Italian and English Telegram Conspiracy Spheres Through Message Forwarding
Lorenzo Alvisi, Serena Tardelli, Maurizio Tesconi
ASONAM (2)1
2024 Autobahn: Seamless high speed BFT
abstract
Today's practical, high performance Byzantine Fault Tolerant (BFT) consensus protocols operate in the partial synchrony model. However, existing protocols are inefficient when deployments are indeed partially synchronous. They deliver either low latency during fault-free, synchronous periods (good intervals) or robust recovery from events that interrupt progress (blips). At one end, traditional, view-based BFT protocols optimize for latency during good intervals, but, when blips occur, can suffer from performance degradation (hangovers) that can last beyond the return of a good interval. At the other end, modern DAG-based BFT protocols recover more gracefully from blips, but exhibit lackluster latency during good intervals. To close the gap, this work presents Autobahn, a novel high-throughput BFT protocol that offers both low latency and seamless recovery from blips. By combining a highly parallel asynchronous data dissemination layer with a low-latency, partially synchronous consensus mechanism, Autobahn (i) avoids the hangovers incurred by traditional BFT protocols and (ii) matches the throughput of state of the art DAG-based BFT protocols while cutting their latency in half, matching the latency of traditional BFT protocols.
Neil Giridharan, Florian Suri-Payer, Ittai Abraham, Lorenzo Alvisi, Natacha Crooks
SOSP4
2023 Morty: Scaling Concurrency Control with Re-Execution
abstract
Serializable systems often perform poorly under high contention. In this work, we analyze this performance limitation through a novel take on conflict windows. Through the lens of these windows, we develop a new concurrency control technique that leverages transaction re-execution to improve throughput scalability under high contention. Our system, Morty, achieves up to 1.7x-96x the throughput of state-of-the-art systems, with similar or better latency.
Matthew Burke 0001, Florian Suri-Payer, Jeffrey Helt, Lorenzo Alvisi, Natacha Crooks
EuroSys4
2023 ScaleDB: A Scalable, Asynchronous In-Memory Database
Syed Akbar Mehdi, Deukyeon Hwang, Simon Peter 0001, Lorenzo Alvisi
OSDI4
2023 Gorilla: Safe Permissionless Byzantine Consensus
abstract
Nakamoto's consensus protocol works in a permissionless model and tolerates Byzantine failures, but only offers probabilistic agreement. Recently, the Sandglass protocol has shown such weaker guarantees are not a necessary consequence of a permissionless model; yet, Sandglass only tolerates benign failures, and operates in an unconventional partially synchronous model. We present Gorilla Sandglass, the first Byzantine tolerant consensus protocol to guarantee, in the same synchronous model adopted by Nakamoto, deterministic agreement and termination with probability 1 in a permissionless setting. We prove the correctness of Gorilla by mapping executions that would violate agreement or termination in Gorilla to executions in Sandglass, where we know such violations are impossible. Establishing termination proves particularly interesting, as the mapping requires reasoning about infinite executions and their probabilities.
Youer Pu, Ali Farahbakhsh, Lorenzo Alvisi, Ittay Eyal
DISC3
2022 Safe Permissionless Consensus
abstract
Consensus protocols have traditionally been studied in a setting where all participants are known to each other from the start of the protocol execution. In the parlance of the 'blockchain' literature, this is referred to as the permissioned setting. What differentiates Bitcoin from these previously studied protocols is that it operates in a permissionless setting, i.e. it is a protocol for establishing consensus over an unknown network of participants that anybody can join, with as many identities as they like in any role. The arrival of this new form of protocol brings with it many questions. Beyond Bitcoin, what can we prove about permissionless protocols in a general sense? How does recent work on permissionless protocols in the blockchain literature relate to the well-developed history of research on permissioned protocols in distributed computing? To answer these questions, we describe a formal framework for the analysis of both permissioned and permissionless systems. Our framework allows for "apples-to-apples" comparisons between different categories of protocols and, in turn, the development of theory to formally discuss their relative merits. A major benefit of the framework is that it facilitates the application of a rich history of proofs and techniques in distributed computing to problems in blockchain and the study of permissionless systems. Within our framework, we then address the questions above. We consider the Byzantine Generals Problem as a formalisation of the problem of reaching consensus, and address a programme of research that asks, "Under what adversarial conditions, and for what types of permissionless protocol, is consensus possible?" We prove a number of results for this programme, our main result being that deterministic consensus is not possible for decentralised permissionless protocols. To close, we give a list of eight open questions.
Youer Pu, Lorenzo Alvisi, Ittay Eyal
DISC2
2021 Basil: Breaking up BFT with ACID (transactions)
abstract
This paper presents Basil, the first transactional, leaderless Byzantine Fault Tolerant key-value store. Basil leverages ACID transactions to scalably implement the abstraction of a trusted shared log in the presence of Byzantine actors. Unlike traditional BFT approaches, Basil executes non-conflicting operations in parallel and commits transactions in a single round-trip during fault-free executions. Basil improves throughput over traditional BFT systems by four to five times, and is only four times slower than TAPIR, a non-Byzantine replicated system. Basil's novel recovery mechanism further minimizes the impact of failures: with 30% Byzantine clients, throughput drops by less than 25% in the worst-case.
Florian Suri-Payer, Matthew Burke 0001, Zheng Wang 0078, Lorenzo Alvisi, Natacha Crooks
SOSP5
2021 Building Systems of Systems with Escher
Burcu Canakci, Lorenzo Alvisi, Robbert van Renesse
SSS2
2020 Scalog: Seamless Reconfiguration and Total Order in a Scalable Shared Log
Cong Ding 0001, David Chu, Evan Zhao, Lorenzo Alvisi, Robbert van Renesse
NSDI5
2020 Byzantine Ordered Consensus without Byzantine Oligarchy
Srinath Setty, Qi Chen 0009, Lidong Zhou, Lorenzo Alvisi
OSDI5
2019 2019 Edsger W. Dijkstra Prize in Distributed Computing
abstract
The committee decided to award the 2019 Edsger W. Dijkstra Prize in Distributed Computing to Alessandro Panconesi and Aravind Srinivasan for their paper Randomized Distributed Edge Coloring via an Extension of the Chernoff-Hoeffding Bounds, SIAM Journal on Computing, volume 26, number 2, 1997, pages 350-368. A preliminary version of this paper appeared as Fast Randomized Algorithms for Distributed Edge Coloring, Proceedings of the Eleventh Annual ACM Symposium Principles of Distributed Computing (PODC), 1992, pages 251-262.
Lorenzo Alvisi, Shlomi Dolev, Faith Ellen, Idit Keidar, Fabian Kuhn, Jukka Suomela
PODC1
2018 Obladi: Oblivious Serializable Transactions in the Cloud
Natacha Crooks, Matthew Burke 0001, Ethan Cecchetti, Sitar Harel, Rachit Agarwal 0001, Lorenzo Alvisi
OSDI6
2018 2018 Doctoral Dissertation Award
abstract
The winner of the 2018 Principles of Distributed Computing Doctoral Dissertation Award is Dr. Rati Gelashvili, for his dissertation titled "On the Complexity of Synchronization," written under the supervision of Prof. Nir Shavit at the Massachusetts Institute of Technology.
Lorenzo Alvisi, Idit Keidar, Andréa W. Richa, Alexander A. Schwarzmann
PODC1
2017 I Can't Believe It's Not Causal! Scalable Causal Consistency with No Slowdown Cascades
Syed Akbar Mehdi, Cody Littley, Natacha Crooks, Lorenzo Alvisi, Nathan Bronson, Wyatt Lloyd
NSDI4
2017 Seeing is Believing: A Client-Centric Specification of Database Isolation
abstract
This paper introduces the first state-based formalization of isolation guarantees. Our approach is premised on a simple observation: applications view storage systems as black-boxes that transition through a series of states, a subset of which are observed by applications. Defining isolation guarantees in terms of these states frees definitions from implementation-specific assumptions. It makes immediately clear what anomalies, if any, applications can expect to observe, thus bridging the gap that exists today between how isolation guarantees are defined and how they are perceived. The clarity that results from definitions based on client-observable states brings forth several benefits. First, it allows us to easily compare the guarantees of distinct, but semantically close, isolation guarantees. We find that several well-known guarantees, previously thought to be distinct, are in fact equivalent, and that many previously incomparable flavors of snapshot isolation can be organized in a clean hierarchy. Second, freeing definitions from implementation-specific artefacts can suggest more efficient implementations of the same isolation guarantee. We show how a client-centric implementation of parallel snapshot isolation can be more resilient to slowdown cascades, a common phenomenon in large-scale datacenters.
Natacha Crooks, Youer Pu, Lorenzo Alvisi, Allen Clement
PODC3
2017 Pretzel: Email encryption and provider-supplied functions are compatible
abstract
Emails today are often encrypted, but only between mail servers---the vast majority of emails are exposed in plaintext to the mail servers that handle them. While better than no encryption, this arrangement leaves open the possibility of attacks, privacy violations, and other disclosures. Publicly, email providers have stated that default end-to-end encryption would conflict with essential functions (spam filtering, etc.), because the latter requires analyzing email text. The goal of this paper is to demonstrate that there is no conflict. We do so by designing, implementing, and evaluating Pretzel. Starting from a cryptographic protocol that enables two parties to jointly perform a classification task without revealing their inputs to each other, Pretzel refines and adapts this protocol to the email context. Our experimental evaluation of a prototype demonstrates that email can be encrypted end-to-end and providers can compute over it, at tolerable cost: clients must devote some storage and processing, and provider overhead is roughly 5x versus the status quo.
Trinabh Gupta, Henrique Fingler, Lorenzo Alvisi, Michael Walfish
SIGCOMM3
2017 Bringing Modular Concurrency Control to the Next Level
abstract
This paper presents Tebaldi, a distributed key-value store that explores new ways to harness the performance opportunity of combining different specialized concurrency control mechanisms (CCs) within the same database. Tebaldi partitions conflicts at a fine granularity and matches them to specialized CCs within a hierarchical framework that is modular, extensible, and able to support a wide variety of concurrency control techniques, from single-version to multiversion and from lock-based to timestamp-based. When running the TPC-C benchmark, Tebaldi yields more than 20× the throughput of the basic two-phase locking protocol, and over 3.7× the throughput of Callas, a recent system that, like Tebaldi, aims to combine different CCs.
Chunzhi Su, Natacha Crooks, Cong Ding 0001, Lorenzo Alvisi
SIGMOD Conference4
2016 Scalable and Private Media Consumption with Popcorn
Trinabh Gupta, Natacha Crooks, Whitney Mulhern, Srinath Setty, Lorenzo Alvisi, Michael Walfish
NSDI5
2016 TARDiS: A Branch-and-Merge Approach To Weak Consistency
abstract
This paper presents the design, implementation, and evaluation of TARDiS (Transactional Asynchronously Replicated Divergent Store), a transactional key-value store explicitly designed for weakly-consistent systems. Reasoning about these systems is hard, as neither causal consistency nor per-object eventual convergence allow applications to deal satisfactorily with write-write conflicts. TARDiS instead exposes as its fundamental abstraction the set of conflicting branches that arise in weakly-consistent systems. To this end, TARDiS introduces a new concurrency control mechanism: branch-on-conflict. On the one hand, TARDiS guarantees that storage will appear sequential to any thread of execution that extends a branch, keeping application logic simple. On the other, TARDiS provides applications, when needed, with the tools and context necessary to merge branches atomically, when and how applications want. Since branch-on-conflict in TARDiS is fast, weakly-consistent applications can benefit from adopting this paradigm not only for operations issued by different sites, but also, when appropriate, for conflicting local operations. We find that TARDiS reduces coding complexity for these applications and that judicious branch-on-conflict can improve their local throughput at each site by two to eight times.
Natacha Crooks, Youer Pu, Nancy Estrada, Trinabh Gupta, Lorenzo Alvisi, Allen Clement
SIGMOD Conference5
2015 High-performance ACID via modular concurrency control
abstract
This paper describes the design, implementation, and evaluation of Callas, a distributed database system that offers to unmodified, transactional ACID applications the opportunity to achieve a level of performance that can currently only be reached by rewriting all or part of the application in a BASE/NoSQL Style. The key to combining performance and ease of programming is to decouple the ACID abstraction---which Callas offers identically for all transactions---from the mechanism used to support it. MCC, the new Modular approach to Concurrency Control at the core of Callas, makes it possible to partition transactions in groups with the guarantee that, as long as the concurrency control mechanism within each group upholds a given isolation property, that property will also hold among transactions in different groups. Because of their limited and specialized scope, these group-specific mechanisms can be customized for concurrency with unprecedented aggressiveness. In our MySQL Cluster-based prototype, Callas yields an 8.2x throughput gain for TPC-C with no programming effort.
Chunzhi Su, Cody Littley, Lorenzo Alvisi, Manos Kapritsos, Yang Wang 0009
SOSP4
2014 Exalt: Empowering Researchers to Evaluate Large-Scale Storage Systems
Yang Wang 0009, Manos Kapritsos, Lara Schmidt, Lorenzo Alvisi, Michael Dahlin
NSDI4
2014 Salt: Combining ACID and BASE in a Distributed Database
Chunzhi Su, Manos Kapritsos, Yang Wang 0009, Navid Yaghmazadeh, Lorenzo Alvisi, Prince Mahajan
OSDI6
2014 Lazy Means Smart: Reducing Repair Bandwidth Costs in Erasure-coded Distributed Storage
abstract
Erasure coding schemes provide higher durability at lower storage cost, and thus constitute an attractive alternative to replication in distributed storage systems, in particular for storing rarely accessed "cold" data. These schemes, however, require an order of magnitude higher recovery bandwidth for maintaining a constant level of durability in the face of node failures. In this paper we propose lazy recovery, a technique to reduce recovery bandwidth demands down to the level of replicated storage. The key insight is that a careful adjustment of recovery rate substantially reduces recovery bandwidth, while keeping the impact on read performance and data durability low. We demonstrate the benefits of lazy recovery via extensive simulation using a realistic distributed storage configuration and published component failure parameters. For example, when applied to the commonly used RS(14, 10) code, lazy recovery reduces repair bandwidth by up to 76% even below replication, while increasing the amount of degraded stripes by 0.1 percentage points. Lazy recovery works well with a variety of erasure coding schemes, including the recently introduced bandwidth efficient codes, achieving up to a factor of 2 additional bandwidth savings.
Mark Silberstein, Lakshmi Ganesh, Yang Wang 0009, Lorenzo Alvisi, Michael Dahlin
SYSTOR4
2013 Reasoning with MAD Distributed Systems
Lorenzo Alvisi, Edmund L. Wong
CONCUR1
2013 Robustness in the Salus Scalable Block Store
Yang Wang 0009, Manos Kapritsos, Zuocheng Ren, Prince Mahajan, Jeevitha Kirubanandam, Lorenzo Alvisi, Michael Dahlin
NSDI6
2013 What's a little collusion between friends?
abstract
This paper proposes a fundamentally different approach to addressing the challenge posed by colluding nodes to the sustainability of cooperative services. Departing from previous work that tries to address the threat by disincentivizing collusion or by modeling colluding nodes as faulty, this paper describes two new notions of equilibrium, k-indistinguishability and k-stability, that allow coalitions to leverage their associations without harming the stability of the service.
Edmund L. Wong, Lorenzo Alvisi
PODC2
2013 SoK: The Evolution of Sybil Defense via Social Networks
abstract
Sybil attacks in which an adversary forges a potentially unbounded number of identities are a danger to distributed systems and online social networks. The goal of sybil defense is to accurately identify sybil identities. This paper surveys the evolution of sybil defense protocols that leverage the structural properties of the social graph underlying a distributed system to identify sybil identities. We make two main contributions. First, we clarify the deep connection between sybil defense and the theory of random walks. This leads us to identify a community detection algorithm that, for the first time, offers provable guarantees in the context of sybil defense. Second, we advocate a new goal for sybil defense that addresses the more limited, but practically useful, goal of securely white-listing a local region of the graph.
Lorenzo Alvisi, Allen Clement, Alessandro Epasto, Silvio Lattanzi, Alessandro Panconesi
IEEE Symposium on Security and Privacy1
2012 All about Eve: Execute-Verify Replication for Multi-Core Servers
Manos Kapritsos, Yang Wang 0009, Vivien Quéma, Allen Clement, Lorenzo Alvisi, Michael Dahlin
OSDI5
2012 Gnothi: Separating Data and Metadata for Efficient and Available Storage Replication
Yang Wang 0009, Lorenzo Alvisi, Michael Dahlin
USENIX ATC2
2011 Regret Freedom Isn't Free
Edmund L. Wong, Isaac Levy, Lorenzo Alvisi, Allen Clement, Michael Dahlin
OPODIS3
2011 Special issue on PODC 2009
Lorenzo Alvisi
Distributed Comput.1
2011 Depot: Cloud Storage with Minimal Trust
abstract
This article describes the design, implementation, and evaluation of Depot, a cloud storage system that minimizes trust assumptions. Depot tolerates buggy or malicious behavior byany numberof clients or servers, yet it provides safety and liveness guarantees to correct clients. Depot provides these guarantees using a two-layer architecture. First, Depot ensures that the updates observed by correct nodes are consistently ordered under Fork-Join-Causal consistency (FJC). FJC is a slight weakening of causal consistency that can be both safe and live despite faulty nodes. Second, Depot implements protocols that use this consistent ordering of updates to provide other desirable consistency, staleness, durability, and recovery properties. Our evaluation suggests that the costs of these guarantees are modest and that Depot can tolerate faults and maintain good availability, latency, overhead, and staleness even when significant faults occur.
Prince Mahajan, Srinath Setty, Allen Clement, Lorenzo Alvisi, Michael Dahlin, Michael Walfish
ACM Trans. Comput. Syst.5
2010 Depot: Cloud Storage with Minimal Trust
Prince Mahajan, Srinath Setty, Allen Clement, Lorenzo Alvisi, Michael Dahlin, Michael Walfish
OSDI5
2010 It's on Me! The Benefit of Altruism in BAR Environments
Edmund L. Wong, Joshua B. Leners, Lorenzo Alvisi
DISC3
2010 Dual-Quorum: A Highly Available and Consistent Replication System for Edge Services
abstract
This paper introduces dual-quorum replication, a novel data replication algorithm designed to support Internet edge services. Edge services allow clients to access Internet services via distributed edge servers that operate on a shared collection of underlying data. Although it is generally difficult to share data while providing high availability, good performance, and strong consistency, replication algorithms designed for specific access patterns can offer nearly ideal trade-offs among these metrics. In this paper, we focus on the key problem of sharing read/write data objects across a collection of edge servers when the references to each object 1) tend not to exhibit high concurrency across multiple nodes and 2) tend to exhibit bursts of read-dominated or write-dominated behavior. Dual-quorum replication combines volume leases and quorum-based techniques to achieve excellent availability, response time, and consistency for such workloads. In particular, through both analytical and experimental evaluations, we show that the dual-quorum protocol can (for the workloads of interest) approach the optimal performance and availability of Read-One/Write-All-Asynchronously (ROWA-A) epidemic algorithms without suffering the weak consistency guarantees and resulting design complexity inherent in ROWA-A systems.
Michael Dahlin, Jiandan Zheng, Lorenzo Alvisi, Arun Iyengar
IEEE Trans. Dependable Secur. Comput.4
2009 Making Byzantine Fault Tolerant Systems Tolerate Byzantine Faults
Allen Clement, Edmund L. Wong, Lorenzo Alvisi, Michael Dahlin, Mirco Marchetti
NSDI3
2009 Upright cluster services
abstract
The UpRight library seeks to make Byzantine fault tolerance (BFT) a simple and viable alternative to crash fault tolerance for a range of cluster services. We demonstrate UpRight by producing BFT versions of the Zookeeper lock service and the Hadoop Distributed File System (HDFS). Our design choices in UpRight favor simplifying adoption by existing applications; performance is a secondary concern. Despite these priorities, our BFT Zookeeper and BFT HDFS implementations have performance comparable with the originals while providing additional robustness.
Allen Clement, Manos Kapritsos, Yang Wang 0009, Lorenzo Alvisi, Michael Dahlin, Taylor L. Riché
SOSP5
2009 Model Checking Coalition Nash Equilibria in MAD Distributed Systems
Federico Mari, Igor Melatti, Ivano Salvo, Enrico Tronci, Lorenzo Alvisi, Allen Clement, Harry C. Li
SSS5
2009 The 2009 Edsger W. Dijkstra Prize in Distributed Computing
Lorenzo Alvisi, Rachid Guerraoui, Prasad Jayanti, Idit Keidar, Shay Kutten, Jennifer L. Welch
DISC1
2009 Zyzzyva: Speculative Byzantine fault tolerance
abstract
A longstanding vision in distributed systems is to build reliable systems from unreliable components. An enticing formulation of this vision is Byzantine Fault-Tolerant (BFT) state machine replication, in which a group of servers collectively act as a correct server even if some of the servers misbehave or malfunction in arbitrary (“Byzantine”) ways. Despite this promise, practitioners hesitate to deploy BFT systems, at least partly because of the perception that BFT must impose high overheads. In this article, we present Zyzzyva, a protocol that uses speculation to reduce the cost of BFT replication. In Zyzzyva, replicas reply to a client's request without first running an expensive three-phase commit protocol to agree on the order to process requests. Instead, they optimistically adopt the order proposed by a primary server, process the request, and reply immediately to the client. If the primary is faulty, replicas can become temporarily inconsistent with one another, but clients detect inconsistencies, help correct replicas converge on a single total ordering of requests, and only rely on responses that are consistent with this total order. This approach allows Zyzzyva to reduce replication overheads to near their theoretical minima and to achieve throughputs of tens of thousands of requests per second, making BFT replication practical for a broad range of demanding services.
Ramakrishna Kotla, Lorenzo Alvisi, Michael Dahlin, Allen Clement, Edmund L. Wong
ACM Trans. Comput. Syst.2
2009 Practical and low-overhead masking of failures of TCP-based servers
abstract
This article describes an architecture that allows a replicated service to survive crashes without breaking its TCP connections. Our approach does not require modifications to the TCP protocol, to the operating system on the server, or to any of the software running on the clients. Furthermore, it runs on commodity hardware. We compare two implementations of this architecture (one based on primary/backup replication and another based on message logging) focusing on scalability, failover time, and application transparency. We evaluate three types of services: a file server, a Web server, and a multimedia streaming server. Our experiments suggest that the approach incurs low overhead on throughput, scales well as the number of clients increases, and allows recovery of the service in near-optimal time.
Dmitrii Zagorodnov, Keith Marzullo, Lorenzo Alvisi, Thomas C. Bressoud
ACM Trans. Comput. Syst.3
2008 BAR primer
abstract
Byzantine and rational behaviors are increasingly recognized as unavoidable realities in todaypsilas cooperative services. Yet, how to design BAR-tolerant protocols and rigorously prove them strategy proof remains somewhat of a mystery: existing examples tend either to focus on unrealistically simple problems or to want in rigor. The goal of this paper is to demystify the process by presenting the full algorithmic development cycle that, starting from the classic synchronous repeated terminating reliable broadcast (R-TRB) problem statement, leads to a provably BAR-tolerant solution. We show i) how to express R-TRB as a game; ii) why the strategy corresponding to the optimal Byzantine fault tolerant algorithm of Dolev and strong does not guarantee safety when non-Byzantine players behave rationally; iii) how to derive a BAR-tolerant R-TRB protocol: iv) how to prove rigorously that the protocol ensures safety in the presence of non-Byzantine rational players.
Allen Clement, Harry C. Li, Jeff Napper, Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin
DSN5
2008 Model Checking Nash Equilibria in MAD Distributed Systems
abstract
We present a symbolic model checking algorithm for verification of Nash equilibria in finite state mechanisms modeling multiple administrative domains (MAD) distributed systems. Given a finite state mechanism, a proposed protocol for each agent and an indifference threshold for rewards, our model checker returns PASS if the proposed protocol is a Nash equilibrium (up to the given indifference threshold) for the given mechanism, FAIL otherwise. We implemented our model checking algorithm inside the NuSMV model checker and present experimental results showing its effectiveness for moderate size mechanisms.
Federico Mari, Igor Melatti, Ivano Salvo, Enrico Tronci, Lorenzo Alvisi, Allen Clement, Harry C. Li
FMCAD5
2008 FlightPath: Obedience vs. Choice in Cooperative Services
Harry C. Li, Allen Clement, Mirco Marchetti, Manos Kapritsos, Luke Robison, Lorenzo Alvisi, Michael Dahlin
OSDI6
2008 Matrix Signatures: From MACs to Digital Signatures in Distributed Systems
Amitanand S. Aiyer, Lorenzo Alvisi, Rida A. Bazzi, Allen Clement
DISC2
2007 Bounded wait-free implementation of optimally resilient byzantine storage without (unproven) cryptographic assumptions
abstract
No abstract available.
Amitanand S. Aiyer, Lorenzo Alvisi, Rida A. Bazzi
PODC2
2007 Theory of BAR games
abstract
No abstract available.
Allen Clement, Jeff Napper, Harry C. Li, Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin
PODC5
2007 Truth in advertising: lightweight verification of route integrity
abstract
We design and evaluate a lightweight route verification mechanism that enables a router to discover route failures and inconsistencies between advertised Internet routes and actual paths taken by the data packets. Our mechanism is accurate, incrementally deployable, and secure against malicious intermediary routers. By carefully avoiding any cryptographic operations in the data path, our prototype implementation achieves the overhead of less than 1% on a 1 Gbps link, demonstrating that our method is suitable even for high-performance networks.
Edmund L. Wong, Praveen Balasubramanian, Lorenzo Alvisi, Mohamed G. Gouda, Vitaly Shmatikov
PODC3
2007 Zyzzyva: speculative byzantine fault tolerance
abstract
We present Zyzzyva, a protocol that uses speculation to reduce the cost and simplify the design of Byzantine fault tolerant state machine replication. In Zyzzyva, replicas respond to a client's request without first running an expensive three-phase commit protocol to reach agreement on the order in which the request must be processed. Instead, they optimistically adopt the order proposed by the primary and respond immediately to the client. Replicas can thus become temporarily inconsistent with one another, but clients detect inconsistencies, help correct replicas converge on a single total ordering of requests, and only rely on responses that are consistent with this total order. This approach allows Zyzzyva to reduce replication overheads to near their theoretical minimal.
Ramakrishna Kotla, Lorenzo Alvisi, Michael Dahlin, Allen Clement, Edmund L. Wong
SOSP2
2007 The Paxos Register
abstract
We introduce the Paxos register to simplify and unify the presentation of Paxos-style consensusprotocols. We use our register to show how Lamport7.s Classic Paxos and Castro and Liskov's Byzantine Paxos are the same consensusprotocol, but for different failure models. We also use our register to compare and contrast Byzantine Paxos with Martin and Alvisi's Fast Byzantine Consensus. The Paxos register is a write-once register that exposes two important abstractions for reaching consensus: ( i ) read and write operations that capture how processes in Pams protocols progose and decide values and (ii) tokens that capture how these protocols guarantee agreement despite partial failures. We encapsulate the difSerences of several Paxos-style protocols in the implementation details of these abstractions.
Harry C. Li, Allen Clement, Amitanand S. Aiyer, Lorenzo Alvisi
SRDS4
2007 SafeStore: A Durable and Practical Storage System
Ramakrishna Kotla, Lorenzo Alvisi, Michael Dahlin
USENIX ATC2
2007 Bounded Wait-Free Implementation of Optimally Resilient Byzantine Storage Without (Unproven) Cryptographic Assumptions
Amitanand S. Aiyer, Lorenzo Alvisi, Rida A. Bazzi
DISC2
2006 Key Grids: A Protocol Family for Assigning Symmetric Keys
abstract
We describe a family of log n protocols for assigning symmetric keys to n processes in a network so that each process can use its assigned keys to communicate securely with every other process. The k-th protocol in our protocol family, where 1 les k les log n, assigns O(k2kradicn) symmetric keys to each process in the network. (Thus, our (log n)-th protocol assigns O(log2n) symmetric keys to each process. This is not far from the lower bound of O(log n) symmetric keyswhichweshowis needed for each process to communicate securely with every other process in the network.) The protocols in our protocol family can be used to assign symmetric keys to the processes in a sensor network, or ad-hoc or mobile network, where each process has a small memory to store its assigned keys. We also discuss the vulnerability of our protocols to "collusion". In particular, we show thatkradicn colluding processes can compromise the security of the k-th protocol in our protocol family.
Amitanand S. Aiyer, Lorenzo Alvisi, Mohamed G. Gouda
ICNP2
2006 BAR Gossip
Harry C. Li, Allen Clement, Edmund L. Wong, Jeff Napper, Indrajit Roy 0001, Lorenzo Alvisi, Michael Dahlin
OSDI6
2006 Byzantine and Multi-writer K-Quorums
Amitanand S. Aiyer, Lorenzo Alvisi, Rida A. Bazzi
DISC2
2006 Fast Byzantine Consensus
abstract
We present the first protocol that reaches asynchronous Byzantine consensus in two communication steps in the common case. We prove that our protocol is optimal in terms of both number of communication steps and number of processes for two-step consensus. The protocol can be used to build a replicated state machine that requires only three communication steps per request in the common case. Further, we show a parameterized version of the protocol that is safe despite f Byzantine failures and, in the common case, guarantees two-step execution despite some number t of failures (t les f). We show that this parameterized two-step consensus protocol is also optimal in terms of both number of communication steps and number of processes
Jean-Philippe Martin, Lorenzo Alvisi
IEEE Trans. Dependable Secur. Comput.2
2006 Correction to "Fast Byzantine Consensus"
Jean-Philippe Martin, Lorenzo Alvisi
IEEE Trans. Dependable Secur. Comput.2
2005 Fast Byzantine Consensus
abstract
We present the first consensus protocol that reaches asynchronous Byzantine consensus in two communication steps in the common case. We prove that our protocol is optimal in terms of both number of communication step, and number of processes for 2-step consensus. The protocol can be used to build a replicated state machine that requires only three communication steps per request in the common case.
Jean-Philippe Martin, Lorenzo Alvisi
DSN2
2005 Dual-Quorum Replication for Edge Services
Michael Dahlin, Jiandan Zheng, Lorenzo Alvisi, Arun Iyengar
Middleware4
2005 BAR fault tolerance for cooperative services
abstract
This paper describes a general approach to constructing cooperative services that span multiple administrative domains. In such environments, protocols must tolerate both Byzantine behaviors when broken, misconfigured, or malicious nodes arbitrarily deviate from their specification and rational behaviors when selfish nodes deviate from their specification to increase their local benefit. The paper makes three contributions: (1) It introduces the BAR (Byzantine, Altruistic, Rational) model as a foundation for reasoning about cooperative services; (2) It proposes a general three-level architecture to reduce the complexity of building services under the BAR model; and (3) It describes an implementation of BAR-B the first cooperative backup service to tolerate both Byzantine users and an unbounded number of rational users. At the core of BAR-B is an asynchronous replicated state machine that provides the customary safety and liveness guarantees despite nodes exhibiting both Byzantine and rational behaviors. Our prototype provides acceptable performance for our application: our BAR-tolerant state machine executes 15 requests per second, and our BAR-B backup service can back up 100MB of data in under 4 minutes.
Amitanand S. Aiyer, Lorenzo Alvisi, Allen Clement, Michael Dahlin, Jean-Philippe Martin, Carl Porth
SOSP2
2005 On the Availability of Non-strict Quorum Systems
Amitanand S. Aiyer, Lorenzo Alvisi, Rida A. Bazzi
DISC2
2005 Improving the Performance of Software Distributed Shared Memory with Speculation
abstract
We study the performance benefits of speculation in a release consistent software distributed shared memory system. We propose a new protocol, speculative home-based release consistency (SHRC) that speculatively updates data at remote nodes to reduce the latency of remote memory accesses. Our protocol employs a predictor that uses patterns in past accesses to shared memory to predict future accesses. We have implemented our protocol in a release consistent software distributed shared memory system that runs on commodity hardware. We evaluate our protocol implementation using eight software distributed shared memory benchmarks and show that it can result in significant performance improvements.
Michael Kistler, Lorenzo Alvisi
IEEE Trans. Parallel Distributed Syst.2
2004 A Framework for Dynamic Byzantine Storage
abstract
We present a framework for transforming several quorum-based protocols so that they can dynamically adapt their failure threshold and server count, allowing them to be reconfigured in anticipation of possible failures or to replace servers as desired. We demonstrate this transformation on the dissemination quorum protocol. The resulting system provides confirmable wait-free atomic semantics while tolerating Byzantine failures from the clients or servers. The system can grow without bound to tolerate as many failures as desired. Finally, the protocol is optimal and fast: only the minimal number of servers - 3f + 1 - is needed to tolerate any f failures and, in the common case, reads require only one message round-trip.
Jean-Philippe Martin, Lorenzo Alvisi
DSN2
2003 A Fault-Tolerant Java Virtual Machine
abstract
The Java programming language was designed for portability and safe code distribution, not for fault-tolerance. We modify the Sun JDK1.2 to provide transparent fault-tolerance for many Java applications under the crash failure model. Our approach is to log non-deterministic events at the JVM interface using a primary-backup architecture. In particular, we identify the sources of non-determinism in the JVM due to asynchronous exceptions and multi-threaded access to shared data, as well as the non-determinism present at the native method interface. We analyze the overhead introduced in our system by each of these sources of non-determinism and compare the performance of dierent techniques for handling multi-threading.
Jeff Napper, Lorenzo Alvisi, Harrick M. Vin
DSN2
2003 Engineering Fault-Tolerant TCP/IP Servers Using FT-TCP
abstract
In a recent paper we have proposed FT-TCP: an architecture that allows a replicated service to survive crashes without breaking its TCP connections. FT-TCP is attractive in principle because it does not require modifications to the TCP protocol and does not affect any of the software running on the clients; however, its practicality for real-world applications remains to be proven. In this paper, we report on our experience in engineering FT-TCP for two such applications---the Samba file server and a multimedia streaming server from Apple. We compare two implementations of FT-TCP, one based on primary-backup and another based on message logging, focusing on scalability, failover time, and application transparency. Our experiments suggest that FT-TCP is a practicable approach for replicating TCP/IP-based services that incurs low overhead on throughput, scales well as the number of clients increases, and allows recovery of the service in near-optimal time.
Dmitrii Zagorodnov, Keith Marzullo, Lorenzo Alvisi, Thomas C. Bressoud
DSN3
2003 Separating agreement from execution for byzantine fault tolerant services
abstract
We describe a new architecture for Byzantine fault tolerant state machine replication that separates agreement that orders requests from execution that processes requests. This separation yields two fundamental and practically significant advantages over previous architectures. First, it reduces replication costs because the new architecture can tolerate faults in up to half of the state machine replicas that execute requests. Previous systems can tolerate faults in at most a third of the combined agreement/state machine replicas. Second, separating agreement from execution allows a general privacy firewall architecture to protect confidentiality through replication. In contrast, replication in previous systems hurts confidentiality because exploiting the weakest replica can be sufficient to compromise the system. We have constructed a prototype and evaluated it running both microbenchmarks and an NFS server. Overall, we find that the architecture adds modest latencies to unreplicated systems and that its performance is competitive with existing Byzantine fault tolerant systems.
Jian Yin 0002, Jean-Philippe Martin, Arun Venkataramani, Lorenzo Alvisi, Michael Dahlin
SOSP4
2003 Scalable causal message logging for wide-area environments
abstract
Abstract Wide‐area systems are gaining in popularity as an infrastructure for running scientific applications. From a fault tolerance perspective, these environments are challenging because of their scale and their variability. Causal message logging protocols have attractive properties that make them suitable for these environments. They spread fault tolerance information around in the system providing high availability. This information can also be used to replicate objects that are otherwise inaccessible because of network partitions. However, current causal message logging protocols do not scale to thousands or millions of processes. We describe the Hierarchical Causal Message Logging Protocol (HCML) that uses a hierarchy of shared logging sites, or proxies, to significantly reduce the space requirements as compared with existing protocols. These proxies also act as caches for fault tolerance information and reduce the overall message overhead by as much as 50%. HCML also leverages differences in bandwidth between processes that reduces overall message latency by as much as 97%. Copyright © 2003 John Wiley & Sons, Ltd.
Karan Bhatia, Keith Marzullo, Lorenzo Alvisi
Concurr. Comput. Pract. Exp.3
2002 Small Byzantine Quorum Systems
abstract
In this paper we present two protocols for asynchronous Byzantine quorum systems (BQS) built on top of reliable channels-one for self-verifying data and the other for any data. Our protocols tolerate f Byzantine failures with f fewer servers than existing solutions by eliminating nonessential work in the write protocol and by using read and write quorums of different sizes. Since engineering a reliable network layer on an unreliable network is difficult, two other possibilities must be explored. The first is to strengthen the model by allowing synchronous networks that use time-outs to identify failed links or machines. We consider running synchronous and asynchronous Byzantine quorum protocols over synchronous networks and conclude that, surprisingly, "self-timing" asynchronous Byzantine protocols may offer significant advantages for many synchronous networks when network time-outs are long. We show how to extend an existing Byzantine quorum protocol to eliminate its dependency on reliable networking and to handle message loss and retransmission explicitly.
Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin
DSN2
2002 Modeling the Effect of Technology Trends on the Soft Error Rate of Combinational Logic
abstract
This paper examines the effect of technology scaling and microarchitectural trends on the rate of soft errors in CMOS memory and logic circuits. We describe and validate an end-to-end model that enables us to compute the soft error rates (SER) for existing and future microprocessor-style designs. The model captures the effects of two important masking phenomena, electrical masking and latching-window masking, which inhibit soft errors in combinational logic. We quantify the SER due to high-energy neutrons in SRAM cells, latches, and logic circuits for feature sizes from 600 nm to 50 nm and clock periods from 16 to 6 fan-out-of-4 inverter delays. Our model predicts that the SER per chip of logic circuits will increase nine orders of magnitude from 1992 to 2011 and at that point will be comparable to the SER per chip of unprotected memory elements. Our result emphasizes that computer system designers must address the risks of soft errors in logic circuits for future designs.
Premkishore Shivakumar, Michael Kistler, Stephen W. Keckler, Doug Burger, Lorenzo Alvisi
DSN5
2002 Half-Pipe Anchoring: An Efficient Technique for Multiple Connection Handoff
abstract
We present half-pipe anchoring, a technique that enables multiple connection handoff mechanisms that are efficient and easy to implement. In a server cluster, these mechanisms result in better resource utilization and improved scalability. More importantly, half-pipe anchoring supports the only connection handoff mechanism that operates efficiently in heterogeneous clusters composed of specialized nodes. The key idea behind our approach is to decouple the two unidirectional half-pipes that make up a TCP connection between a client and a cluster. We anchor the unidirectional half-pipe from the client to the cluster at a designated server while allowing the half-pipe from the cluster to the client to migrate on a per-request basis to an optimal server where the request is best serviced We describe the design and implementation of a multiple connection handoff mechanism in the Linux kernel that demonstrates the benefits of our technique.
Ravi Kokku, Ramakrishnan Rajamony, Lorenzo Alvisi, Harrick M. Vin
ICNP3
2002 Minimal Byzantine Storage
Jean-Philippe Martin, Lorenzo Alvisi, Michael Dahlin
DISC2
2002 Causality tracking in causal message-logging protocols
Lorenzo Alvisi, Karan Bhatia, Keith Marzullo
Distributed Comput.1
2002 Engineering web cache consistency
abstract
Server-driven consistency protocols can reduce read latency and improve data freshness for a given network and server overhead, compared to the traditional consistency protocols that rely on client polling. Server-driven consistency protocols appear particularly attractive for large-scale dynamic Web workloads because dynamically generated data can change rapidly and unpredictably. However, there have been few reports on engineering server-driven consistency for such workloads. This article reports our experience in engineering server-driven consistency for a sporting and event Web site hosted by IBM, one of the most popular sites on the Internet for the duration of the event. We also examine an e-commerce site for a national retail store. Our study focuses on scalability and cachability of dynamic content. To assess scalability, we measure both the amount of state that a server needs to maintain to ensure consistency and the bursts of load in sending out invalidation messages when a popular object is modified. We find that server-driven protocols can cap the size of the server's state to a given amount without significant performance costs, and can smooth the bursts of load with minimal impact on the consistency guarantees. To improve performance, we systematically investigate several design issues for which prior research has suggested widely different solutions, including whether servers should send invalidations to idle clients. Finally, we quantify the performance impact of caching dynamic data with server-driven consistency protocols and the benefits of server-driven consistency protocols for large-scale dynamic Web services. We find that (i) caching dynamically generated data can increase cache hit rates by up to 10%, compared to the systems that do not cache dynamically generated data; and (ii) server-driven consistency protocols can increase cache hit rates by a factor of 1.5-3 for large-scale dynamic Web services, compared to client polling protocols. We have implemented a prototype of a server-driven consistency protocol based on our findings by augmenting the popular Squid cache.
Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Arun Iyengar
ACM Trans. Internet Techn.2
2001 Scalable Causal Message Logging for Wide-Area Environments
Karan Bhatia, Keith Marzullo, Lorenzo Alvisi
Euro-Par3
2001 Wrapping Server-Side TCP to Mask Connection Failures
abstract
We present an implementation of a fault-tolerant TCP (FT-TCP) that allows a faulty server to keep its TCP connections open until it either recovers or it is failed over to a backup. The failure and recovery of the server process are completely transparent to client processes connected with it via TCP. FT-TCP does not affect the software running on a client, does not require to change the server's TCP implementation, and does not use a proxy.
Lorenzo Alvisi, Thomas C. Bressoud, Ayman El-Khashab, Keith Marzullo, Dmitrii Zagorodnov
INFOCOM1
2001 Heterogeneous networking: a new survivability paradigm
abstract
We believe that a network, to be survivable, must be heterogeneous. Just like a species that draws on a small gene pool can succumb to a single environmental threat, so a homogeneous network is vulnerable to a malicious attack that exploits a single weakness common to all of its components. In contrast, in a network in which each critical functionality is provided by a diverse set of protocols and implementations, attacks that focus on a weakness of one such protocol or implementation will not be able to bring down the entire network, even though all elements are not be bulletproof and even if some of components are compromised.Following this survivability through heterogeneity philosophy, we propose a new survivability paradigm, called heterogeneous networking, for improving a network's defense capabilities. Rather than following the current trend of converging towards single solutions to provide the desired functionality at every element of the network architecture, this methodology calls for systematically increasing the network's heterogeneity without sacrificing its interoperability.
Yongguang Zhang, Harrick M. Vin, Lorenzo Alvisi, Wenke Lee, Son K. Dao
NSPW3
2001 A framework for semantic reasoning about Byzantine quorum systems
abstract
We present a set of definitions and theorems that allow us to reason about the semantics of quorum system variables, including Byzantine quorum system variables, as a class. Using these tools, we present a formal proof that the problem of atomic semantics for such variables can be reduced to the simpler problem of regular semantics for such systems. Specifically, any regular masking quorum system protocol can be combined with a writeback mechanism to produce an atomic protocol. We then describe a subclass of TS-variables for which the latter problem is not solvable by traditional approaches in an asynchronous environment. Finally, for such variables we define the notion of pseudoregular and pseudoatomic semantics, and show briey that the same reduction holds for these concepts.
Evelyn Tumlin Pierce, Lorenzo Alvisi
PODC2
2001 Engineering server-driven consistency for large scale dynamic Web services
abstract
Many researchers have shown that server-driven consistency protocols can potentially reduce read latency. Server-driven consistency protocols are particularly attractive for large-scale dynamic web workloads because dynamically generated data can change rapidly and unpredictably. However, there have been no reports on engineering server-driven consistency for such a workload. This paper reports our experience in engineering server-driven consistency for a Sporting and Eventweb site hosted by IBM, one of the most popular web sites on the Internet for the duration of the event. Our study focuses on scalability and cachability of dynamic content. To assess scalability, we measure both the amount of state that a server needs to maintain to ensure consistency and the bursts of load that a server sustains to send out invalidation messages when a popular object is modified. We find that it is possible to limit the size of the server's state without significant performance costs and that bursts of load can be smoothed out with minimal impact on the consistency guarantees. To improve performance, we systematically investigate several design issues for which prior research has suggested widely different solutions, including how long servers should send invalidations to idle clients. Finally, we quantify the performance impact of caching dynamic data with server-driven consistency protocols and find that it can reduce read latency by more than 10%. We have implemented a prototype of a server-driven consistency protocol based on our findings on top of the popular Squid cache.
Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Arun Iyengar
WWW2
2001 Fault Detection for Byzantine Quorum Systems
abstract
In this paper, we explore techniques to detect Byzantine server failures in asynchronous replicated data services. Our goal is to detect arbitrary failures of data servers in a system where each client accesses the replicated data at only a subset (quorum) of servers in each operation. In such a system, some correct servers can be out-of-date after a write and can therefore, return values other than the most up-to-date value in response to a client's read request, thus complicating the task of determining the number of faulty servers in the system at any point in time. We initiate the study of detecting server failures in this context, and propose two statistical approaches for estimating the risk posed by faulty servers based on responses to read requests.
Lorenzo Alvisi, Dahlia Malkhi, Evelyn Tumlin Pierce, Michael K. Reiter
IEEE Trans. Parallel Distributed Syst.1
2000 Dynamic Byzantine Quorum Systems
abstract
Byzantine quorum systems enhance the availability and efficiency of fault-tolerant replicated services when servers may suffer Byzantine failures. An important limitation of Byzantine quorum systems is their dependence on a static threshold limit on the number of server faults. The correctness of the system is only guaranteed if the actual number of faults is lower than the the threshold at all times. However, a threshold chosen for the worst case wastes expensive replication in the common situation where the number of faults averages well below the worst case. In this paper, we present protocols for dynamically raising and lowering the resilience threshold of a quorum-based Byzantine fault-tolerant data service in response to current information on the number of server failures. Using such protocols, a system can operate in an efficient low-threshold mode with relatively small quorums in the absence of faults, increasing and decreasing the quorum size (and thus the tolerance) as faults appear and are dealt with, respectively.
Lorenzo Alvisi, Evelyn Tumlin Pierce, Dahlia Malkhi, Michael K. Reiter, Rebecca N. Wright
DSN1
2000 The Cost of Recovery in Message Logging Protocols
abstract
Past research in message logging has focused on studying the relative overhead imposed by pessimistic, optimistic and causal protocols during failure-free executions. In this paper, we give the first experimental evaluation of the performance of these protocols during recovery. Our results suggest that applications face a complex tradeoff when choosing a message logging protocol for fault tolerance. On the one hand, optimistic protocols can provide fast failure-free execution and good performance during recovery, but are complex to implement and can create orphan processes. On the other hand, orphan-free protocols either risk being slow during recovery (e.g. sender-based pessimistic and causal protocols) or incur a substantial overhead during failure-free execution (e.g. receiver-based pessimistic protocols). To address this tradeoff, we propose hybrid logging protocols, which are a new class of orphan-free protocols. We show that hybrid protocols perform within 2% of causal logging during failure-free execution and within 2% of receiver-based logging during recovery.
Sriram Rao, Lorenzo Alvisi, Harrick M. Vin
IEEE Trans. Knowl. Data Eng.2
1999 Volume Leases for Consistency in Large-Scale Systems
abstract
This article introduces volume leases as a mechanism for providing server-driven cache consistency for large-scale, geographically distributed networks. Volume leases retain the good performance, fault tolerance, and server scalability of the semantically weaker client-driven protocols that are now used on the Web. Volume leases are a variation of object leases, which were originally designed for distributed file systems. However, whereas traditional object leases amortize overheads over long lease periods, volume leases exploit spatial locality to amortize overheads across multiple objects in a volume. This approach allows systems to maintain good write performance even in the presence of failures. Using trace-driven simulation, we compare three volume lease algorithms against four existing cache consistency algorithms and show that our new algorithms provide strong consistency while maintaining scalability and fault-tolerance. For a trace-based workload of Web accesses, we find that volumes can reduce message traffic at servers by 40 percent compared to a standard lease algorithm, and that volumes can considerably reduce the peak load at servers when popular objects are modified.
Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Calvin Lin
IEEE Trans. Knowl. Data Eng.2
1998 Low-Overhead Protocols for Fault-Tolerant File Sharing
abstract
We quantify the adverse effect of file sharing on the performance of reliable distributed applications. We demonstrate that file sharing incurs significant overhead, which is likely to triple over the next five years. We present a novel approach that eliminates this overhead. Our approach: tracks causal dependencies resulting from file sharing using determinants; efficiently replicates the determinants in the volatile memory of agents to ensure their availability during recovery; and reproduces during recovery the interactions with the file server as well as the file data lost in a failure. Our approach allows agents to exchange files directly without first saving the files on disks at the server. As a consequence, the costs of supporting file sharing and message passing in a reliable distributed application become virtually identical. The result is a simple, uniform approach, which can provide low-overhead fault tolerance to applications in which communication is performed through message passing, file sharing, or a combination of the two.
Lorenzo Alvisi, Sriram Rao, Harrick M. Vin
ICDCS1
1998 Using Leases to Support Server-Driven Consistency in Large-Scale Systems
abstract
The paper introduces volume leases as a mechanism for providing cache consistency for large scale, geographically distributed networks. Volume leases are a variation of leases, which were originally designed for distributed file systems. Using trace driven simulation, we compare two new algorithms against four existing cache consistency algorithms and show that our new algorithms provide strong consistency while maintaining scalability and fault tolerance. For a trace based workload of Web accesses, we find that volumes can reduce message traffic at servers by 40% compared to a standard lease algorithm, and that volumes can considerably reduce the peak load at servers when popular objects are modified.
Jian Yin 0002, Lorenzo Alvisi, Michael Dahlin, Calvin Lin
ICDCS2
1998 The Relative Overhead of Piggybacking in Causal Message Logging Protocols
abstract
Message logging protocols ensure that crashed processes make the same choices when re-executing nondeterministic events during recovery. Causal message logging protocols achieve this by piggybacking the results of these choices (called determinants) on the ambient message traffic. By doing so, these protocols do not create orphan processes nor introduce blocking in failure-free executions. To survive f failures, they ensure that determinants are stored by at least f+1 processes. Causal logging protocols differ in the kind of information they piggyback to other processes. The more information they send, the better each process is able to estimate global properties of the determinants, which in turn results in fewer needless piggybacking of determinants. This paper quantifies the tradeoff between the cost of sending more information and the benefit of doing so.
Karan Bhatia, Keith Marzullo, Lorenzo Alvisi
SRDS3
1998 The Cost of Recovery in Message Logging Protocols
abstract
Past research in message logging has focused on studying the relative overhead imposed by pessimistic, optimistic, and causal protocols during failure-free executions. We give the first experimental evaluation of the performance of these protocols during recovery. We discover that, if a single failure is to be tolerated, pessimistic and causal protocols perform best, because they avoid rollbacks of correct processes. For multiple failures, however, the dominant factor in determining performance becomes where the recovery information is logged (i.e. at the sender, at the receiver, or replicated at a subset of the processes in the system) rather than when this information is logged (i.e. if logging is synchronous or asynchronous).
Sriram Rao, Lorenzo Alvisi, Harrick M. Vin
SRDS2
1998 Message Logging: Pessimistic, Optimistic, Causal, and Optimal
abstract
Message-logging protocols are an integral part of a popular technique for implementing processes that can recover from crash failures. All message-logging protocols require that, when recovery is complete, there be no orphan processes, which are surviving processes whose states are inconsistent with the recovered state of a crashed process. We give a precise specification of the consistency property "no orphan processes". From this specification, we describe how different existing classes of message-logging protocols (namely optimistic, pessimistic, and a class that we call causal) implement this property. We then propose a set of metrics to evaluate the performance of message-logging protocols, and characterize the protocols that are optimal with respect to these metrics. Finally, starting from a protocol that relies on causal delivery order, we show how to derive optimal causal protocols that tolerate f overlapping failures and recoveries for a parameter f (1/spl les/f/spl les/n).
Lorenzo Alvisi, Keith Marzullo
IEEE Trans. Software Eng.1
1996 Trade-Offs in Implementing Optimal Message Logging Protocols
abstract
Casual message logging protocols[3] have several attractive properties: they introduce no blocking, send no additional messages over those sent by the application, and can never cause orphans to be created by crashes.Causal message logging, however, does require additional data to be piggybacked on application messages.The amount of such piggybacked data can become large.In this paper, we present five different implementations of casual message logging.All of the corresponding protocols are parameterized by ~, the maximum number of processes that can fail concurrently.We also explore how the application's communication structure can be exploited to limit the amount of piggybacked data.ing recorded that state.When a process crashes, a new process is created in its place: the new process is given
Lorenzo Alvisi, Keith Marzullo
PODC1
1996 Parallel Computing in Networks of Workstations with Paralex
abstract
Modern distributed systems consisting of powerful workstations and high-speed interconnection networks are an economical alternative to special-purpose supercomputers. The technical issues that need to be addressed in exploiting the parallelism inherent in a distributed system include heterogeneity, high-latency communication, fault tolerance and dynamic load balancing. Current software systems for parallel programming provide little or no automatic support towards these issues and require users to be experts in fault-tolerant distributed computing. The Paralex system is aimed at exploring the extent to which the parallel application programmer can be liberated from the complexities of distributed systems. Paralex is a complete programming environment and makes extensive use of graphics to define, edit, execute, and debug parallel scientific applications. All of the necessary code for distributing the computation across a network and replicating it to achieve fault tolerance and dynamic load balancing is automatically generated by the system. In this paper we give an overview of Paralex and present our experiences with a prototype implementation.
Renzo Davoli, Luigi-Alberto Giachini, Özalp Babaoglu, Alessandro Amoroso, Lorenzo Alvisi
IEEE Trans. Parallel Distributed Syst.5
1995 Message Logging: Pessimistic, Optimistic, and Causal
abstract
Message logging protocols are an integral part of a technique for implementing processes that can recover from crash failures. All message logging protocols require that, when recovery is complete, there be no orphan processes, which are surviving processes whose states are inconsistent with the recovered state of a crashed process. We give a precise specification of the consistency property "no orphan processes". From this specification, we describe how different existing classes of message logging protocols (namely optimistic, pessimistic, and a class that we call causal) implement this property. We then propose a set of metrics to evaluate the performance of message logging protocols, and characterize the protocols that are optimal with respect to these metrics. Finally, starting from a protocol that relies on causal delivery order, we show how to derive optimal causal protocols that tolerate f overlapping failures and recoveries for a parameter f:1/spl les/f/spl les/n.
Lorenzo Alvisi, Keith Marzullo
ICDCS1
1995 Deriving Optimal Checkpoint Protocols for Distributed Shared Memory Architectures (Abstract)
abstract
No abstract available.
Lorenzo Alvisi, Keith Marzullo
PODC1
1992 Paralex: an environment for parallel programming in distributed systems
abstract
Modern distributed systems consisting of powerful workstations and high-speed interconnection networks are an economical alternative to special-purpose super computers. The technical issues that need to be addressed in exploiting the parallelism inherent in a distributed system include heterogeneity, high-latency communication, fault tolerance and dynamic load balancing. Current software systems for parallel programming provide little or no automatic support towards these issues and require users to be experts in fault-tolerant distributed computing. The Paralex system is aimed at exploring the extent to which the parallel application programmer can be liberated from the complexities of distributed systems. Paralex is a complete programming environment and makes extensive use of graphics to define, edit, execute and debug parallel scientific applications. All of the necessary code for distributing the computation across a network and replicating it to achieve fault tolerance and dynamic load balancing is automatically generated by the system. In this paper we give an overview of Paralex and present our experiences with a prototype implementation.
Özalp Babaoglu, Lorenzo Alvisi, Alessandro Amoroso, Renzo Davoli, Luigi-Alberto Giachini
ICS2
1989 On the two array mask hidden-line algorithm
Lorenzo Alvisi, Giulio Casciola
Comput. Graph.1