Robbert van Renesse

dblp:r/RvRenesse · DBLP profile ↗
← Back
99ranked-venue papers
17as first author
8since 2021 · last 2026
0000-0003-3598-0283ORCID · verified

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

Systems, architecture and hardware · 46 · 6 first-author · 5 since 2021Security and privacy · 22 · 3 first-author · 1 since 2021Software engineering, systems software and programming languages · 14 · 3 first-authorComputer networks · 13 · 1 first-author · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 2Human-computer interaction and ubiquitous computing · 1 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 1
YearPublicationVenuePosition
2026 FicusDB: Scalable Multi-Versioned Authenticated Archival Storage
abstract
As consensus protocols scale, storage has become the dominant bottleneck in modern blockchains. Systems must maintain historical versions, generate integrity proofs, and sustain high throughput as state grows. Prior work often redesigns the authenticated data structure (ADS), but such changes sacrifice compatibility and require disruptive hard forks. We take the complementary approach: redesigning the storage layer while preserving the existing ADS interface.
Maofan Yin, Robbert van Renesse
EuroSys3
2026 Editor-in-Chief's Message
Sam H. Noh, Robbert van Renesse
ACM Trans. Comput. Syst.2
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
HotNets5
2023 Charlotte: Reformulating Blockchains into a Web of Composable Attested Data Structures for Cross-Domain Applications
abstract
Cross-domain applications are rapidly adopting blockchain techniques for immutability, availability, integrity, and interoperability. However, for most applications, global consensus is unnecessary and may not even provide sufficient guarantees. We propose a new distributed data structure: Attested Data Structures (ADS), which generalize not only blockchains but also many other structures used by distributed applications. As in blockchains, data in ADSs is immutable and self-authenticating. ADSs go further by supporting application-defined proofs ( attestations ). Attestations enable applications to plug in their own mechanisms to ensure availability and integrity. We present Charlotte , a framework for composable ADSs. Charlotte deconstructs conventional blockchains into more primitive mechanisms. Charlotte can be used to construct blockchains but does not impose the usual global-ordering overhead. Charlotte offers a flexible foundation for interacting applications that define their own policies for availability and integrity. Unlike traditional distributed systems, Charlotte supports heterogeneous trust: different observers have their own beliefs about who might fail, and how. Nevertheless, each observer has a consistent, available view of data. Charlotte’s data structures are interoperable and composable : applications and data structures can operate fully independently or can share data when desired. Charlotte defines a language-independent format for data blocks and a network API for servers. To demonstrate Charlotte’s flexibility, we implement several integrity mechanisms, including consensus and proof of work. We explore the power of disentangling availability and integrity mechanisms in prototype applications. The results suggest that Charlotte can be used to build flexible, fast, composable applications with strong guarantees.
Isaac C. Sheff, Xinwen Wang, Kushal Babel, Haobin Ni, Robbert van Renesse, Andrew C. Myers
ACM Trans. Comput. Syst.5
2022 BloomBox: Improving Availability and Efficiency in Geographic Hash Tables
abstract
Mobile Ad Hoc Networks are important today for scenarios in which centralized cloud infrastructure is missing, has broken down, or imposes censure or undesirable monitoring of storage and communication. Unfortunately, existing peer-to-peer storage systems such as a Geographic Hash Table (GHT) can consume a significant amount of network bandwidth just to maintain a certain required number of replicas of the data due to the churn present in the network. The replicas have to continuously exchange heartbeat messages in order to detect failures of replicas. If heartbeat messages get lost, an unnecessary but expensive recovery protocol ends up wasting significant bandwidth. To avoid this, replicas are placed close to one another, but that makes them vulnerable to dependent failures.Based on Mergeable Bloom Filters, a new data structure proposed for peer-to-peer distributed systems, we build BloomBox, a failure detection protocol for a GHT. Our simulations show that BloomBox can significantly reduce bandwidth usage needed for regenerated blocks compared to heartbeat-based failure detection. Moreover, BloomBox can provide significantly better availability than protocols based on heartbeats by placing replicas in geographically diverse locations.
Xinwen Wang, Robbert van Renesse
ICDCS2
2021 A Fresh Look at the Design and Implementation of Communication Paradigms (Invited Talk)
Robbert van Renesse
OPODIS1
2021 Building Systems of Systems with Escher
Burcu Canakci, Lorenzo Alvisi, Robbert van Renesse
SSS3
2021 CacheInspector: Reverse Engineering Cache Resources in Public Clouds
abstract
Infrastructure-as-a-Service cloud providers sell virtual machines that are only specified in terms of number of CPU cores, amount of memory, and I/O throughput. Performance-critical aspects such as cache sizes and memory latency are missing or reported in ways that make them hard to compare across cloud providers. It is difficult for users to adapt their application’s behavior to the available resources. In this work, we aim to increase the visibility that cloud users have into shared resources on public clouds. Specifically, we present CacheInspector , a lightweight runtime that determines the performance and allocated capacity of shared caches on multi-tenant public clouds. We validate CacheInspector ’s accuracy in a controlled environment, and use it to study the characteristics and variability of cache resources in the cloud, across time, instances, availability regions, and cloud providers. We show that CacheInspector ’s output allows cloud users to tailor their application’s behavior, including their output quality, to avoid suboptimal performance when resources are scarce.
Weijia Song, Christina Delimitrou, Zhiming Shen, Robbert van Renesse, Hakim Weatherspoon, Lotfi Benmohamed, Frederic J. de Vaulx, Charif Mahmoudi
ACM Trans. Archit. Code Optim.4
2020 Scalog: Seamless Reconfiguration and Total Order in a Scalable Shared Log
Cong Ding 0001, David Chu, Evan Zhao, Lorenzo Alvisi, Robbert van Renesse
NSDI6
2020 Heterogeneous Paxos
abstract
In distributed systems, a group of learners achieve consensus when, by observing the output of some acceptors, they all arrive at the same value. Consensus is crucial for ordering transactions in failure-tolerant systems. Traditional consensus algorithms are homogeneous in three ways: - all learners are treated equally, - all acceptors are treated equally, and - all failures are treated equally. These assumptions, however, are unsuitable for cross-domain applications, including blockchains, where not all acceptors are equally trustworthy, and not all learners have the same assumptions and priorities. We present the first consensus algorithm to be heterogeneous in all three respects. Learners set their own mixed failure tolerances over differently trusted sets of acceptors. We express these assumptions in a novel Learner Graph, and demonstrate sufficient conditions for consensus. We present Heterogeneous Paxos, an extension of Byzantine Paxos. Heterogeneous Paxos achieves consensus for any viable Learner Graph in best-case three message sends, which is optimal. We present a proof-of-concept implementation and demonstrate how tailoring for heterogeneous scenarios can save resources and reduce latency.
Isaac C. Sheff, Xinwen Wang, Robbert van Renesse, Andrew C. Myers
OPODIS3
2020 Scaling Membership of Byzantine Consensus
abstract
Scaling Byzantine Fault Tolerant (BFT) systems in terms of membership is important for secure applications with large participation such as blockchains. While traditional protocols have low latency, they cannot handle many processors. Conversely, blockchains often have hundreds to thousands of processors to increase robustness, but they typically have high latency or energy costs. We describe various sources of unscalability in BFT consensus protocols. To improve performance, many BFT protocols optimize the “normal case,” where there are no failures. This can be done in a modular fashion by wrapping existing BFT protocols with a building block that we call alliance . In normal case executions, alliance can scalably determine if the initial conditions of a BFT consensus protocol predetermine the outcome, obviating running the consensus protocol. We give examples of existing protocols that solve alliance. We show that a solution based on hypercubes and MAC s has desirable scalability and performance in normal case executions, with only a modest overhead otherwise. We provide important optimizations. Finally, we evaluate our solution using the ns3 simulator and show that it scales up to thousands of processors and compare with prior work in various network topologies.
Burcu Canakci, Robbert van Renesse
ACM Trans. Comput. Syst.2
2019 X-Containers: Breaking Down Barriers to Improve Performance and Isolation of Cloud-Native Containers
abstract
"Cloud-native" container platforms, such as Kubernetes, have become an integral part of production cloud environments. One of the principles in designing cloud-native applications is called Single Concern Principle, which suggests that each container should handle a single responsibility well. In this paper, we propose X-Containers as a new security paradigm for isolating single-concerned cloud-native containers. Each container is run with a Library OS (LibOS) that supports multi-processing for concurrency and compatibility. A minimal exokernel ensures strong isolation with small kernel attack surface. We show an implementation of the X-Containers architecture that leverages Xen paravirtualization (PV) to turn Linux kernel into a LibOS. Doing so results in a highly efficient LibOS platform that does not require hardware-assisted virtualization, improves inter-container isolation, and supports binary compatibility and multi-processing. By eliminating some security barriers such as seccomp and Meltdown patch, X-Containers have up to 27X higher raw system call throughput compared to Docker containers, while also achieving competitive or superior performance on various benchmarks compared to recent container platforms such as Google's gVisor and Intel's Clear Containers.
Zhiming Shen, Gur-Eyal Sela, Eugene Bagdasarian, Christina Delimitrou, Robbert van Renesse, Hakim Weatherspoon
ASPLOS6
2018 Untethered: Deployable Blockchains for IoT Environments
abstract
No abstract available.
Kolbeinn Karlsson, Danny Adams, Gloire Rubambiza, Zangyueyang Xian, Robbert van Renesse, Hakim Weatherspoon, Stephen B. Wicker
SoCC5
2018 Vegvisir: A Partition-Tolerant Blockchain for the Internet-of-Things
abstract
While the intersection of blockchains and the Internet of Things (IoT) have received considerable research interest lately, Nakamoto-style blockchains possess a number of qualities that make them poorly suited for many IoT scenarios. Specifically, they require high network connectivity and are power-intensive. This is a drawback in IoT environments where battery-constrained nodes form an unreliable ad hoc network such as in digital agriculture. In this paper we present Vegvisir, a partition-tolerant blockchain for use in power-constrained IoT environments with limited network connectivity. It is a permissioned, directed acyclic graph (DAG)-structured blockchain that can be used to create a shared, tamperproof data repository that keeps track of data provenance. We discuss the use cases, architecture, and challenges of such a blockchain.
Kolbeinn Karlsson, Weitao Jiang, Stephen B. Wicker, Danny Adams, Edwin Ma, Robbert van Renesse, Hakim Weatherspoon
ICDCS6
2018 Derecho: Fast State Machine Replication for Cloud Services
abstract
Cloud computing services often replicate data and may require ways to coordinate distributed actions. Here we present Derecho, a library for such tasks. The API provides interfaces for structuring applications into patterns of subgroups and shards, supports state machine replication within them, and includes mechanisms that assist in restart after failures. Running over 100Gbps RDMA, Derecho can send millions of events per second in each subgroup or shard and throughput peaks at 16GB/s, substantially outperforming prior solutions. Configured to run purely on TCP, Derecho is still substantially faster than comparable widely used, highly-tuned, standard tools. The key insight is that on modern hardware (including non-RDMA networks), data-intensive protocols should be built from non-blocking data-flow components.
Sagar Jha, Jonathan Behrens, Theo Gkountouvas, Mae Milano, Weijia Song, Edward Tremel, Robbert van Renesse, Sydney Zink, Kenneth P. Birman
ACM Trans. Comput. Syst.7
2017 Building smart memories and high-speed cloud services for the internet of things with derecho
abstract
The coming generation of Internet-of-Things (IoT) applications will process massive amounts of incoming data while supporting data mining and online learning. In cases with demanding real-time requirements, such systems behave as smart memories: a high-bandwidth service that captures sensor input, processes it using machine-learning tools, replicates and stores "interesting" data (discarding uninteresting content), updates knowledge models, and triggers urgently-needed responses.
Sagar Jha, Jonathan Behrens, Theo Gkountouvas, Mae Milano, Weijia Song, Edward Tremel, Sydney Zink, Kenneth P. Birman, Robbert van Renesse
SoCC9
2017 Towards an emergency edge supercloud
abstract
The "cloud paradigm" can provide a wealth of sophisticated emergency communication services that are gamechangers in emergency response, but its current implementation is not suitable to the challenging environments in which these responses often take place. The networking infrastructure may be all but unavailable, and access to centralized datacenters may be impossible.
Kolbeinn Karlsson, Zhiming Shen, Weijia Song, Hakim Weatherspoon, Robbert van Renesse, Stephen B. Wicker
SoCC5
2017 REM: Resource-Efficient Mining for Blockchains
Fan Zhang 0022, Ittay Eyal, Robert Escriva, Ari Juels, Robbert van Renesse
USENIX Security Symposium5
2017 Supercloud: A Library Cloud for Exploiting Cloud Diversity
abstract
Infrastructure-as-a-Service (IaaS) cloud providers hide available interfaces for virtual machine (VM) placement and migration, CPU capping, memory ballooning, page sharing, and I/O throttling, limiting the ways in which applications can optimally configure resources or respond to dynamically shifting workloads. Given these interfaces, applications could migrate VMs in response to diurnal workloads or changing prices, adjust resources in response to load changes, and so on. This article proposes a new abstraction that we call a Library Cloud and that allows users to customize the diverse available cloud resources to best serve their applications. We built a prototype of a Library Cloud that we call the Supercloud . The Supercloud encapsulates applications in a virtual cloud under users’ full control and can incorporate one or more availability zones within a cloud provider or across different providers. The Supercloud provides virtual machine, storage, and networking complete with a full set of management operations, allowing applications to optimize performance. In this article, we demonstrate various innovations enabled by the Library Cloud.
Zhiming Shen, Qin Jia, Gur-Eyal Sela, Weijia Song, Hakim Weatherspoon, Robbert van Renesse
ACM Trans. Comput. Syst.6
2016 Safe Serializable Secure Scheduling: Transactions and the Trade-Off Between Security and Consistency
abstract
Modern applications often operate on data in multiple administrative domains. In this federated setting, participants may not fully trust each other. These distributed applications use transactions as a core mechanism for ensuring reliability and consistency with persistent data. However, the coordination mechanisms needed for transactions can both leak confidential information and allow unauthorized influence. By implementing a simple attack, we show these side channels can be exploited. However, our focus is on preventing such attacks. We explore secure scheduling of atomic, serializable transactions in a federated setting. While we prove that no protocol can guarantee security and liveness in all settings, we establish conditions for sets of transactions that can safely complete under secure scheduling. Based on these conditions, we introduce \ti{staged commit}, a secure scheduling protocol for federated transactions. This protocol avoids insecure information channels by dividing transactions into distinct stages. We implement a compiler that statically checks code to ensure it meets our conditions, and a system that schedules these transactions using the staged commit protocol. Experiments on this implementation demonstrate that realistic federated transactions can be scheduled securely, atomically, and efficiently.
Isaac C. Sheff, Tom Magrino, Jed Liu, Andrew C. Myers, Robbert van Renesse
CCS5
2016 Follow the Sun through the Clouds: Application Migration for Geographically Shifting Workloads
abstract
Global cloud services have to respond to workloads that shift geographically as a function of time-of-day or in response to special events. While many such services have support for adding nodes in one region and removing nodes in another, we demonstrate that such mechanisms can lead to significant performance degradation. Yet other services do not support application-level migration at all. Live VM migration between availability zones or even across cloud providers would be ideal, but cloud providers do not support this flexible mechanism.
Zhiming Shen, Qin Jia, Gur-Eyal Sela, Ben Rainero, Weijia Song, Robbert van Renesse, Hakim Weatherspoon
SoCC6
2016 Bitcoin-NG: A Scalable Blockchain Protocol
Ittay Eyal, Adem Efe Gencer, Emin Gün Sirer, Robbert van Renesse
NSDI4
2016 Moving Participants Turtle Consensus
abstract
We present Moving Participants Turtle Consensus (MPTC), an asynchronous consensus protocol for crash and Byzantine-tolerant distributed systems. MPTC uses various moving target defense strategies to tolerate certain Denial-of-Service (DoS) attacks issued by an adversary capable of compromising a bounded portion of the system. MPTC supports on the fly reconfiguration of the consensus strategy as well as of the processes executing this strategy when solving the problem of agreement. It uses existing cryptographic techniques to ensure that reconfiguration takes place in an unpredictable fashion thus eliminating the adversary’s advantage on predicting protocol and execution-specific information that can be used against the protocol. We implement MPTC as well as a State Machine Replication protocol and evaluate our design under different attack scenarios. Our evaluation shows that MPTC approximates best case scenario performance even under a well-coordinated DoS attack.
Stavros Nikolaou, Robbert van Renesse
OPODIS2
2016 Proactive Cache Placement on Cooperative Client Caches for Online Social Networks
abstract
This paper investigates cache placement on a cooperative cache built from individual client caches in an online social network or web service. We use a service that maintains a mapping between content and the clients that cache it, and propose cache placement schemes that leverage relationships between clients (for example, social links) and workload statistics, proactively placing content on clients that are likely to access it. We evaluate efficacy through simulation, comparing our schemes against commonly used cache placement algorithms as well as optimal placement. We synthesize a workload to match characteristics of online social networks. Simulation results of our proposed caching schemes impose moderate network overhead and show considerable improvement to the client's cache hit ratio, even under churn.
Stavros Nikolaou, Robbert van Renesse, Nicolas Schiper
IEEE Trans. Parallel Distributed Syst.2
2015 Cache Serializability: Reducing Inconsistency in Edge Transactions
abstract
Read-only caches are widely used in cloud infrastructures to reduce access latency and load on backend databases. Operators view coherent caches as impractical at genuinely large scale and many client-facing caches are updated asynchronously with best-effort pipelines. Existing solutions that support cache consistency are inapplicable to this scenario since they require a round trip to the database on every cache transaction. Existing incoherent cache technologies are oblivious to transactional data access, even if the backend database supports transactions. We propose T-Cache, a novel caching policy for read-only transactions in which inconsistency is tolerable (won't cause safety violations) but undesirable (has a cost). T-Cache improves cache consistency despite asynchronous and unreliable communication between the cache and the database. We define cache-serializability, a variant of serializability that is suitable for incoherent caches, and prove that with unbounded resources T-Cache implements this new specification. With limited resources, T-Cache allows the system manager to choose a trade-off between performance and consistency. Our evaluation shows that T-Cache detects many inconsistencies with only nominal overhead. We use synthetic workloads to demonstrate the efficacy of T-Cache when data accesses are clustered and its adaptive reaction to workload changes. With workloads based on the real-world topologies, T-Cache detects 43 -- 70% of the inconsistencies and increases the rate of consistent transactions by 33 -- 58%.
Ittay Eyal, Kenneth P. Birman, Robbert van Renesse
ICDCS3
2015 Configuring Distributed Computations Using Response Surfaces
abstract
Configuring large distributed computations is a challenging task. Efficiently executing distributed computations requires configuration tuning based on careful examination of application and hardware properties. Considering the large number of parameters and impracticality of using trial and error in a production environment, programmers tend to make these decisions based on their experience and rules of thumb. Such configurations can lead to underutilized and costly clusters, and missed deadlines.
Adem Efe Gencer, David Bindel, Emin Gün Sirer, Robbert van Renesse
Middleware4
2015 Turtle Consensus: Moving Target Defense for Consensus
abstract
Consensus is a basic building block in middleware configuration services [4, 18]. While such services are designed to tolerate crash failures in asynchronous settings, they may not stand up well to Denial-of-Service (DoS) attacks. Specifically, malicious clients can carefully craft workloads that substantially degrade the performance of many state-of-the-art consensus protocols. By exploiting protocol-specific vulnerabilities, attackers can constantly force the protocol participants to slow execution paths [8]. In this paper, we investigate designing consensus protocols that provide acceptable performance under DoS attacks that aim to saturate the bandwidth of protocol participants.
Stavros Nikolaou, Robbert van Renesse
Middleware2
2015 Wormhole: Reliable Pub-Sub to Support Geo-replicated Internet Services
Yogeshwer Sharma, Philippe Ajoux, Petchean Ang, David Callies, Abhishek Choudhary, Laurent Demailly, Thomas Fersch, Liat Atsmon Guz, Andrzej Kotulski, Sachin Kulkarni, Harry C. Li, Evgeniy Makeev, Kowshik Prakasam, Robbert van Renesse, Sabyasachi Roy, Pratyush Seth, Yee Jiun Song, Benjamin Wester, Kaushik Veeraraghavan, Peter Xie
NSDI16
2015 Vive La Différence: Paxos vs. Viewstamped Replication vs. Zab
abstract
Paxos, Viewstamped Replication, and Zab are replication protocols for high-availability in asynchronous environments with crash failures. Claims have been made about their similarities and differences. But how does one determine whether two protocols are the same, and if not, how significant are the differences? We address these questions using refinement mappings. Protocols are expressed as succinct specifications that are progressively refined to executable implementations. Doing so enables a principled understanding of the correctness of design decisions for implementing the protocols. Additionally, differences that have a significant impact on performance are surfaced by this exercise.
Robbert van Renesse, Nicolas Schiper, Fred B. Schneider
IEEE Trans. Dependable Secur. Comput.1
2015 Fireflies: A Secure and Scalable Membership and Gossip Service
abstract
An attacker who controls a computer in an overlay network can effectively control the entire overlay network if the mechanism managing membership information can successfully be targeted. This article describes Fireflies, an overlay network protocol that fights such attacks by organizing members in a verifiable pseudorandom structure so that an intruder cannot incorrectly modify the membership views of correct members. Fireflies provides each member with a view of the entire membership, and supports networks with moderate total churn. We evaluate Fireflies using both simulations and PlanetLab to show that Fireflies is a practical approach for secure membership maintenance in such networks.
Håvard D. Johansen, Robbert van Renesse, Ymir Vigfusson, Dag Johansen
ACM Trans. Comput. Syst.2
2015 Omni-Kernel: An Operating System Architecture for Pervasive Monitoring and Scheduling
abstract
Theomni-kernelarchitecture is designed around pervasive monitoring and scheduling. Motivated by new requirements in virtualized environments, this architecture ensures that all resource consumption is measured, that resource consumption resulting from a scheduling decision is attributable to an activity, and that scheduling decisions are fine-grained.Vortex, implemented for multi-core x86-64 platforms, instantiates the omni-kernel architecture, providing a wide range of operating system functionality and abstractions. With Vortex, we experimentally demonstrated the efficacy of the omni-kernel architecture to provide accurate scheduler control over resource allocation despite competing workloads. Experiments involving Apache, MySQL, and Hadoop quantify the cost of pervasive monitoring and scheduling in Vortex to be below$6$percent ofcpuconsumption.
Åge Kvalnes, Dag Johansen, Robbert van Renesse, Fred B. Schneider, Steffen Viken Valvåg
IEEE Trans. Parallel Distributed Syst.3
2014 Scalable State-Machine Replication
abstract
State machine replication (SMR) is a well-known technique able to provide fault-tolerance. SMR consists of sequencing client requests and executing them against replicas in the same order, thanks to deterministic execution, every replica will reach the same state after the execution of each request. However, SMR is not scalable since any replica added to the system will execute all requests, and so throughput does not increase with the number of replicas. Scalable SMR (S-SMR) addresses this issue in two ways: (i) by partitioning the application state, while allowing every command to access any combination of partitions, and (ii) by using a caching algorithm to reduce the communication across partitions. We describe Eyrie, a library in Java that implements S-SMR, and Volery, an application that implements Zookeeper's API. We assess the performance of Volery and compare the results against Zookeeper. Our experiments show that Volery scales throughput with the number of partitions.
Carlos Eduardo Benevides Bezerra, Fernando Pedone, Robbert van Renesse
DSN3
2014 The Energy Efficiency of Database Replication Protocols
abstract
Replication is a widely used technique to provide high-availability to online services. While being an effective way to mask failures, replication comes at a price: at least twice as much hardware and energy are required to mask a single failure. In a context where the electricity drawn by data centers worldwide is increasing each year, there is a need to maximize the amount of useful work done per Joule, a metric denoted as energy efficiency. In this paper, we review commonly-used database replication protocols and experimentally measure their energy efficiency. We observe that the most efficient replication protocol achieves less than 60% of the energy efficiency of a stand-alone server on the TPC-C benchmark. We identify algorithmic techniques that can be used by any protocol to improve its efficiency. Some approaches improve performance, others lower power consumption. Of particular interest is a technique derived from primary-backup replication that implements a transaction log on low-power backups. We demonstrate how this approach can lead to an energy efficiency that is 79% of the one of a stand-alone server. This constitutes an important step towards reconciling replication with energy efficiency.
Nicolas Schiper, Fernando Pedone, Robbert van Renesse
DSN3
2014 Developing Correctly Replicated Databases Using Formal Tools
abstract
Fault-tolerant distributed systems often contain complex error handling code. Such code is hard to test or model-check because there are often too many possible failure scenarios to consider. As we will demonstrate in this paper, formal methods have evolved to a state in which it is possible to generate this code along with correctness guarantees. This paper describes our experience with building highly-available databases using replication protocols that were generated with the help of correct-by-construction formal methods. The goal of our project is to obtain databases with unsurpassed reliability while providing good performance. We report on our experience using a total order broadcast protocol based on Paxos and specified using a new formal language called Event ML. We compile Event ML specifications into a form that can be formally verified while simultaneously obtaining code that can be executed. We have developed two replicated databases based on this code and show that they have performance that is competitive with popular databases in one of the two considered benchmarks.
Nicolas Schiper, Vincent Rahli, Robbert van Renesse, Mark Bickford, Robert L. Constable
DSN3
2014 Ironstack: Performance, Stability and Security for Power Grid Data Networks
abstract
Operators of the nationwide power grid use proprietary data networks to monitor and manage their power distribution systems. These purpose-built, wide area communication networks connect a complex array of equipment ranging from PMUs and synchrophasers to SCADA systems. Collectively, these equipment form part of an intricate feedback system that ensures the stability of the power grid. In support of this mission, the operational requirements of these networks mandates high performance, reliability, and security. We designed Iron Stack, a system to address these concerns. By using cutting-edge software defined networking technology, Iron Stack is able to use multiple network paths to improve communications bandwidth and latency, provide seamless failure recovery, and ensure signals security. Additionally, Iron Stack is incrementally deployable and backward-compatible with existing switching infrastructure.
Zhiyuan Teo, Vera Kutsenko, Kenneth P. Birman, Robbert van Renesse
DSN4
2014 Characterizing Load Imbalance in Real-World Networked Caches
abstract
Modern Web services rely extensively upon a tier of in-memory caches to reduce request latencies and alleviate load on backend servers. Within a given cache, items are typically partitioned across cache servers via consistent hashing, with the goal of balancing the number of items maintained by each cache server. Effects of consistent hashing vary by associated hashing function and partitioning ratio. Most real-world workloads are also skewed, with some items significantly more popular than others. Inefficiency in addressing both issues can create an imbalance in cache-server loads.
Helga Gudmundsdottir, Ymir Vigfusson, Daniel A. Freedman, Kenneth P. Birman, Robbert van Renesse
HotNets6
2013 Leveraging sharding in the design of scalable replication protocols
abstract
Most if not all datacenter services use sharding and replication for scalability and reliability. Shards are more-or-less independent of one another and individually replicated. In this paper, we challenge this design philosophy and present a replication protocol where the shards interact with one another: A protocol running within shards ensures linearizable consistency, while the shards interact in order to improve availability. We provide a specification for the protocol, prove its safety, analyze its liveness and availability properties, and evaluate a working implementation.
Hussam Abu-Libdeh, Robbert van Renesse, Ymir Vigfusson
SoCC2
2013 SuperCloud: economical cloud service on multiple vendors
abstract
Today, Infrastructure-as-a-Service (IaaS) cloud providers such as Amazon's Elastic Compute Engine (EC2), Google's Compute Engine, and Microsoft's Azure offer elastic and isolated compute resources via virtualization and users often choose one of these providers based on price, locality, performance, and features. Typically, a user will choose the same provider for computation and storage to minimize latency and networking costs. Unfortunately, it can be difficult to switch providers once one is selected due to vendor lock-in [2].
Qin Jia, Robbert van Renesse, Hakim Weatherspoon
SoCC2
2013 Application-driven TCP recovery and non-stop BGP
abstract
Some network protocols tie application state to underlying TCP connections, leading to unacceptable service outages when an endpoint loses TCP state during fail-over or migration. For example, BGP ties forwarding tables to its control plane connections so that the failure of a BGP endpoint can lead to widespread routing disruption, even if it recovers all of its state but what was encapsulated by its TCP implementation. Although techniques exist for recovering TCP state transparently, they make assumptions that do not hold for applications such as BGP. We introduce application-driven TCP recovery, a technique that separates application recovery from TCP recovery. We evaluate our prototype, TCPR, and show that it outperforms existing BGP recovery techniques.
Robert Surton, Kenneth P. Birman, Robbert van Renesse
DSN3
2013 Sprinkler - Reliable Broadcast for Geographically Dispersed Datacenters
Haoyan Geng, Robbert van Renesse
Middleware2
2013 Secure Abstraction with Code Capabilities
abstract
We propose embedding executable code fragments in cryptographically protected capabilities to enable flexible discretionary access control in cloud-like computing infrastructures. We demonstrate how such a code capability mechanism can be implemented completely in user space. Using a novel combination of X.509 certificates and JavaScript code, code capabilities support restricted delegation, confinement, revocation, and rights amplification for secure abstraction.
Robbert van Renesse, Håvard D. Johansen, Nihar Naigaonkar, Dag Johansen
PDP1
2013 An analysis of Facebook photo caching
abstract
This paper examines the workload of Facebook's photo-serving stack and the effectiveness of the many layers of caching it employs. Facebook's image-management infrastructure is complex and geographically distributed. It includes browser caches on end-user systems, Edge Caches at ~20 PoPs, an Origin Cache, and for some kinds of images, additional caching via Akamai. The underlying image storage layer is widely distributed, and includes multiple data centers.
Kenneth P. Birman, Robbert van Renesse, Wyatt Lloyd, Harry C. Li
SOSP3
2012 A diversified and correct-by-construction broadcast service
abstract
We present a fault-tolerant ordered broadcast service that is correct-by-construction. Our broadcast service allows for diversity in space, whereby the participants in the broadcast protocol run different code, as well as in time, whereby the protocol itself is changed periodically. We use the Nuprl proof assistant to specify the service, prove correctness, and synthesize the code. The paper includes initial performance results.
Vincent Rahli, Nicolas Schiper, Robbert van Renesse, Mark Bickford, Robert L. Constable
ICNP3
2012 Byzantine Chain Replication
Robbert van Renesse, Chi Ho, Nicolas Schiper
OPODIS1
2009 Refining the way to consensus
abstract
No abstract available.
Robbert van Renesse
PODC1
2009 Slicing Distributed Systems
abstract
Peer-to-peer (P2P) architectures are popular for tasks such as collaborative download, VoIP telephony, and backup. To maximize performance in the face of widely variable storage capacities and bandwidths, such systems typically need to shift work from poor nodes to richer ones. Similar requirements are seen in today's large data centers, where machines may have widely variable configurations, loads, and performance. In this paper, we consider the slicing problem, which involves partitioning the participating nodes into k subsets using a one-dimensional attribute, and updating the partition as the set of nodes and their associated attributes change. The mechanism thus facilitates the development of adaptive systems. We begin by motivating this problem statement and reviewing prior work. Existing algorithms are shown to have problems with convergence, manifesting as inaccurate slice assignments, and to adapt slowly as conditions change. Our protocol, Sliver, has provably rapid convergence, is robust under stress and is simple to implement. We present both theoretical and experimental evaluations of the protocol.
Vincent Gramoli, Ymir Vigfusson, Kenneth P. Birman, Anne-Marie Kermarrec, Robbert van Renesse
IEEE Trans. Computers5
2008 Tempest: Soft state replication in the service tier
abstract
Soft state in the middle tier is key to enabling scalable and responsive three tier service architectures. While soft-state can be reconstructed upon failure, replicating it across multiple service instances is critical for rapid fail-over and high availability. Current techniques for storing and managing replicated soft state require mapping data structures to different abstractions such as database records, which can be difficult and introduce inefficiencies. Tempest is a system that provides programmers with data structures that look very similar to conventional Java Collections but are automatically replicated. We evaluate Tempest against alternatives such as in-memory databases and we show that Tempest does scale well in real world service architectures.
Tudor Marian, Mahesh Balakrishnan 0001, Kenneth P. Birman, Robbert van Renesse
DSN4
2008 Nysiad: Practical Protocol Transformation to Tolerate Byzantine Failures
Chi Ho, Robbert van Renesse, Mark Bickford, Danny Dolev
NSDI2
2008 A fast distributed slicing algorithm
abstract
No abstract available.
Vincent Gramoli, Ymir Vigfusson, Kenneth P. Birman, Anne-Marie Kermarrec, Robbert van Renesse
PODC5
2008 Bosco: One-Step Byzantine Asynchronous Consensus
Yee Jiun Song, Robbert van Renesse
DISC2
2008 SecureStream: An intrusion-tolerant protocol for live-streaming dissemination
Maya Haridasan, Robbert van Renesse
Comput. Commun.2
2007 Self-stabilizing and Byzantine-Tolerant Overlay Network
Danny Dolev, Ezra N. Hoch, Robbert van Renesse
OPODIS3
2007 Making Distributed Applications Robust
Chi Ho, Danny Dolev, Robbert van Renesse
OPODIS3
2007 FirePatch: Secure and Time-Critical Dissemination of Software Patches
Håvard D. Johansen, Dag Johansen, Robbert van Renesse
SEC3
2006 Fireflies: scalable support for intrusion-tolerant network overlays
abstract
This paper describes and evaluates Fireflies, a scalable protocol for supporting intrusion-tolerant network overlays. While such a protocol cannot distinguish Byzantine nodes from correct nodes in general, Fireflies provides correct nodes with a reasonably current view of which nodes are live, as well as a pseudo-random mesh for communication. The amount of data sent by correct nodes grows linearly with the aggregate rate of failures and recoveries, even if provoked by Byzantine nodes. The set of correct nodes form a connected submesh; correct nodes cannot be eclipsed by Byzantine nodes. Fireflies is deployed and evaluated on PlanetLab.
Håvard D. Johansen, André Allavena, Robbert van Renesse
EuroSys3
2006 MISTRAL: : efficient flooding in mobile ad-hoc networks
abstract
Flooding is an important communication primitive in mobile ad-hoc networks and also serves as a building block for more complex protocols such as routing protocols. In this paper, we propose a novel approach to flooding, which relies on proactive compensation packets periodically broadcast by every node. The compensation packets are constructed from dropped data packets, based on techniques borrowed from forward error correction. Since our approach does not rely on proactive neighbor discovery and network overlays it is resilient to mobilit.We evaluate the implementation of Mistral through simulation and compare its performance and overhead to purely probabilistic flooding. Our results show that Mistral achieves a significantly higher node coverage with comparable overhead.
Stefan Pleisch, Mahesh Balakrishnan 0001, Kenneth P. Birman, Robbert van Renesse
MobiHoc4
2006 Defense against Intrusion in a Live Streaming Multicast System
abstract
Application-level multicast systems are vulnerable to attacks that impede nodes from receiving desired data. Live streaming protocols are especially susceptible to packet loss induced by malicious behavior. We describe SecureStream, an application-level live streaming system built using a pull-based architecture that results in improved tolerance of malicious behavior. SecureStream is implemented as a layer running over Fireflies, an intrusion-tolerant membership protocol. Our paper describes the SecureStream system and offers simulation and experimental results confirming its resilience to attack
Maya Haridasan, Robbert van Renesse
Peer-to-Peer Computing2
2006 Cognitive Adaptive Radio Teams
abstract
Cognitive adaptive radio teams (CART) is a new platform developed by our group in support of collaborative mapping of complex communications-challenged environments, for example in support of search and rescue operations in environments lacking an adequate communications infrastructure. Experience during the 9/11 terrorist attacks, Asian tsunami, Kashmir earthquake, and post-Katrina Gulf Coast make it clear that rescue workers cannot count upon computer networks or even cell telephone support in the immediate aftermath of such events. Similarly, military urban warfare operations must also be conducted in locations lacking communication infrastructure. CART combines state-of-the-art ad-hoc networking technology with machine learning and prediction algorithms to offer new capabilities under these very difficult conditions
Richard Lau, Stephanie Demers, Yibei Ling, Bruce Siegell, Einar Vollset, Kenneth P. Birman, Robbert van Renesse, Howard E. Shrobe, Jonathan Bachrach, Lester Foster
SECON7
2006 A Scalable Services Architecture
abstract
Data centers constructed as clusters of inexpensive machines have compelling cost-performance benefits, but developing services to run on them can be challenging. This paper reports on a new framework, the scalable services architecture (SSA), which helps developers develop scalable clustered applications. The work is focused on non-transactional high-performance applications; these are poorly supported in existing platforms. A primary goal was to keep the SSA as small and simple as possible. Key elements include a TCP-based "chain replication" mechanism and a gossip-based subsystem for managing configuration data and repairing inconsistencies after faults. Our experimental results confirm the effectiveness of the approach
Tudor Marian, Kenneth P. Birman, Robbert van Renesse
SRDS3
2005 Using randomized techniques to build scalable intrusion-tolerant overlay networks (Keynote)
abstract
Overlay networks provide important routing functionality not easily supported directly by the Internet. Distributed Hash Tables (DHTs) have been proposed to support such overlay networks. While it is often straightforward to support overlay networks on DHTs, this choice can be questioned. DHTs dictate routes that are not optimal, and DHTs are hard to secure. As overlay networks are beginning to be deployed for critical applications, efficiency and security are becoming important attributes. We present an alternative support structure called Fireflies. Fireflies provides each of its members with a complete view of its live peers. A small subset of these peers is marked as neighbors. With high probability, the mesh formed by the members and their neighbor links has a diameter logarithmic in the number of live members, and connects all the reachable members that are not Byzantine. Fireflies uses several randomized algorithms, which are briefly described. We also discuss how Fireflies may be used to build intrusion-tolerant overlay networks
Robbert van Renesse
CollaborateCom1
2005 dcOvercoming Communications Challenges in Software for Monitoring and Controlling Power Systems
abstract
The restructuring of the electric power grid has created new control and monitoring requirements for which classical technologies may be inadequate. The most obvious way of building such systems, using TCP connections to link monitoring systems with data sources, gives poor scalability and exhibits instability precisely when information is most urgently required. Astrolabe, Bimodal Multicast, and Gravitational Gossip, technologies of our own design, seek to overcome these problems using what are called "epidemic" communication protocols. This paper evaluates a hypothetical power monitoring scenario involving the New York State grid, and concludes that the technology is well matched to the need.
Kenneth P. Birman, Jie Chen 0017, E. M. Hopkinson, Robert J. Thomas, James S. Thorp, Robbert van Renesse, Werner Vogels
Proc. IEEE6
2005 JiST: an efficient approach to simulation using virtual machines
abstract
Discrete event simulators are important scientific tools and their efficient design and execution is the subject of much research. In this paper, we propose a new approach for constructing simulators that leverages virtual machines and combines advantages from the traditional systems-based and language-based simulator designs. We introduce JiST, a Java-based simulation system that executes discrete event simulations both efficiently and transparently by embedding simulation semantics directly into the Java execution model. The system provides standard benefits that the modern Java runtime affords. In addition, JiST is efficient, out-performing existing highly optimized simulation runtimes. As a case study, we illustrate the practicality of the JiST framework by applying it to the construction of SWANS, a scalable wireless ad hoc network simulator. We simulate million node wireless networks, which represents two orders of magnitude increase in scale over what existing simulators can achieve on equivalent hardware and at the same level of detail. Copyright © 2005 John Wiley & Sons, Ltd.
Rimon Barr, Zygmunt J. Haas, Robbert van Renesse
Softw. Pract. Exp.3
2005 APSS: proactive secret sharing in asynchronous systems
abstract
APSS, a proactive secret sharing (PSS) protocol for asynchronous systems, is explained and proved correct. The protocol enables a set of secret shares to be periodically refreshed with a new, independent set, thereby thwarting mobile-adversary attacks. Protocols for asynchronous systems are inherently less vulnerable to denial-of-service attacks, which slow processor execution or delay message delivery. So APSS tolerates certain attacks that PSS protocols for synchronous systems cannot.
Lidong Zhou, Fred B. Schneider, Robbert van Renesse
ACM Trans. Inf. Syst. Secur.3
2004 Adding High Availability and Autonomic Behavior to Web Services
abstract
Rapid acceptance of the Web Services architecture promises to make it the most widely supported and popular object-oriented architecture to date. One consequence is that a wave of mission-critical Web Services applications will certainly be deployed in coming years. Yet the reliability options available within Web Services are limited in important ways. To use a term proposed by IBM, Web Services systems need to become far more autonomic, configuring themselves, diagnosing faults, and managing themselves. High availability applications need more attention. Moreover, the scenarios in which such issues arise often entail very large deployments, raising questions of scalability. In this paper we propose a path by which the architecture could be extended in these respects.
Kenneth P. Birman, Robbert van Renesse, Werner Vogels
ICSE2
2004 Chain Replication for Supporting High Throughput and Availability
Robbert van Renesse, Fred B. Schneider
OSDI1
2003 Astrolabe: A robust and scalable technology for distributed system monitoring, management, and data mining
abstract
Scalable management and self-organizational capabilities are emerging as central requirements for a generation of large-scale, highly dynamic, distributed applications. We have developed an entirely new distributed information management system called Astrolabe. Astrolabe collects large-scale system state, permitting rapid updates and providing on-the-fly attribute aggregation. This latter capability permits an application to locate a resource, and also offers a scalable way to track system state as it evolves over time. The combination of features makes it possible to solve a wide variety of management and self-configuration problems. This paper describes the design of the system with a focus upon its scalability. After describing the Astrolabe service, we present examples of the use of Astrolabe for locating resources, publish-subscribe, and distributed synchronization in large systems. Astrolabe is implemented using a peer-to-peer protocol, and uses a restricted form of mobile code based on the SQL query language for aggregation. This protocol gives rise to a novel consistency model. Astrolabe addresses several security considerations using a built-in PKI. The scalability of the system is evaluated using both simulation and experiments; these confirm that Astrolabe could scale to thousands and perhaps millions of nodes, with information propagation delays in the tens of seconds.
Robbert van Renesse, Kenneth P. Birman, Werner Vogels
ACM Trans. Comput. Syst.1
2002 Optimizing Buffer Management for Reliable Multicast
abstract
Reliable multicast delivery requires that a multicast message be received by all members in a group. Hence certain or all members need to buffer messages for possible retransmissions. Designing an efficient buffer management algorithm is challenging in large multicast groups where no member has complete group membership information and the delivery latency to different members could differ by orders of magnitude. We propose an innovative two-phase buffering algorithm, which explicitly addresses variations in delivery latency seen in large multicast groups. The algorithm effectively reduces buffer requirements by adaptively allocating buffer space to messages most needed in the system and by spreading the load of buffering among all members in the group. Simulation and experimental results demonstrate that the algorithm has good performance.
Kenneth P. Birman, Robbert van Renesse
DSN3
2002 Workshop on Reliable Peer-to-Peer Distributed Systems
Özalp Babaoglu, Anne-Marie Kermarrec, Robbert van Renesse, Luís E. T. Rodrigues, Maarten van Steen, Amin Vadhat
SRDS3
2002 Power-Aware Epidemics
abstract
Epidemic protocols have been heralded as appropriate for wireless sensor networks. The nodes in such networks have limited battery resources. In this paper we investigate the use of power in three styles of epidemic protocols: basic epidemics, neighborhood flooding epidemics, and hierarchical epidemics. Basic epidemics turn out to be highly power hungry, and are not appropriate for power-aware applications. Both neighborhood and hierarchical epidemics can be made to use power judiciously, but a trade-off exists between scalability and latency.
Robbert van Renesse
SRDS1
2002 Collaborative Networking in an Uncooperative Internet
abstract
Collaborative applications often require peer-to-peer interaction and peer discovery mechanisms. In today's Internet, Firewall and NAT technology, and a lack of support of IP multicast, have made it very difficult to support such applications. Application Level Gateways and Directory Services can solve these problems to some extent, but have scalability problems and should be used as a last resort. This paper describes our experience with implementing a service called Astrolabe which uses a peer-to-peer epidemic protocol. We show how we solved peer-to-peer communication, auto-configuration, and peer discovery. The resulting Astrolabe service can be used to support the development of other peer-to-peer protocols and applications.
Robbert van Renesse, Dan Dumitriu
SRDS1
2002 Implementing IPv6 as a Peer-to-Peer Overlay Network
abstract
This paper proposes to implement an IPv6 routing infrastructure as a self-organizing overlay network on top of the current IPv4 infrastructure. The overlay network builds upon a distributed IPv6 edge router with a master/slave architecture. We show how different slaves can be constructed to tunnel through NATs and firewalls, as well as to improve robustness of the routing infrastructure and to provide efficient and resilient implementations for features such as multicast, anycast, and mobile IP using currently available peer-to-peer (P2P) protocols. The resulting IPv6 overlay network would restore the end-to-end property of the original Internet, support evolution and dynamic updating of the protocols running on the overlay network, make available IPv6 and the associated features to network applications immediately, and provide an ideal underlying infrastructure for P2P applications, without changing networking hardware and software in the core Internet.
Lidong Zhou, Robbert van Renesse, Michael A. Marsh
SRDS2
2002 A TACOMA retrospective
abstract
Abstract For seven years, the TACOMA project has investigated the design and implementation of software support for mobile agents. A series of prototypes has been developed, with experiences in distributed applications driving the effort. This paper describes the evolution of these TACOMA prototypes, what primitives each supports, and how the primitives are used in building distributed applications. Copyright © 2002 John Wiley & Sons, Ltd.
Dag Johansen, Kåre J. Lauvset, Robbert van Renesse, Fred B. Schneider, Nils P. Sudmann, Kjetil Jacobsen
Softw. Pract. Exp.3
2002 COCA: A secure distributed online certification authority
abstract
COCA is a fault-tolerant and secure online certification authority that has been built and deployed both in a local area network and in the Internet. Extremely weak assumptions characterize environments in which COCA's protocols execute correctly: no assumption is made about execution speed and message delivery delays; channels are expected to exhibit only intermittent reliability; and with 3t+ 1 COCA servers up totmay be faulty or compromised. COCA is the first system to integrate a Byzantine quorum system (used to achieve availability) with proactive recovery (used to defend against mobile adversaries which attack, compromise, and control one replica for a limited period of time before moving on to another). In addition to tackling problems associated with combining fault-tolerance and security, new proactive recovery protocols had to be developed. Experimental results give a quantitative evaluation for the cost and effectiveness of the protocols.
Lidong Zhou, Fred B. Schneider, Robbert van Renesse
ACM Trans. Comput. Syst.3
2001 Scalable Fault-Tolerant Aggregation in Large Process Groups
abstract
The paper discusses fault-tolerant, scalable solutions to the problem of accurately and scalably calculating global aggregate functions in large process groups communicating over unreliable networks. These groups could represent sensors or processes communicating over a network that is either fixed (e.g., the Internet) or dynamic (e.g., multihop ad-hoc). Group members are prone to failures. The ability to evaluate global aggregate properties (e.g., the average of sensor temperature readings) is important for higher-level coordination activities in such large groups. We first define the setting and problem, laying down metrics to evaluate different algorithms for the same. We discuss why the usual approaches to solve this problem are unviable and unscalable over an unreliable network prone to message delivery failures and crash failures. We then propose a technique to impose an abstract hierarchy on such large groups, describing how this hierarchy can be made to mirror the network topology. We discuss several alternatives to use this technique to solve the global aggregate function evaluation problem. Finally, we present a protocol based on gossiping that uses this hierarchical technique. We present mathematical analysis and performance results to validate the robustness, efficiency and accuracy of the Hierarchical Gossiping algorithm.
Indranil Gupta, Robbert van Renesse, Kenneth P. Birman
DSN2
2001 Distributing media transformation over multiple media gateways
abstract
Media gateways have been proposed as a solution to the network heterogeneity problem in media multicasting. Services on the gateways transform media streams as they flow through the gateways. In this paper, we present our work on composable services in media gateways. A user can request a computation to be performed on a set of media streams. The system then distributes the computation over multiple gateways for execution. We present an algorithm for decomposing the computation into sub-computations, and an application-level protocol that locates appropriate media gateways to run these sub-computations.
Wei Tsang Ooi, Robbert van Renesse
ACM Multimedia2
2000 An adaptive protocol for localing programmable media gateways
abstract
We describe a new control protocol called Adaptive Gateway Location Protocol (AGLP). In this protocol, a client requests a computation on a multimedia stream. AGLP discovers programmable Internet servers that process multimedia streams, and assigns the computation to one of these so-called gateways. AGLP continuously searches for alternate gateways, and, transparent to users, migrates computations between them to improve efficiency. The AGLP protocol uses soft-states for robustness and scale. Simulation results support that our protocol quickly locates gateways and migrates computations while keeping the load on the network low. We also outline planned enhancements to AGLP.
Wei Tsang Ooi, Robbert van Renesse
ACM Multimedia2
2000 Fast protocol transition in a distributed environment (brief announcement)
abstract
Adaptivity is a desired feature of the distributed systems. Because many characteristics of the environment (network topology, active process distribution, etc.) may change from time to time, a good system should be able to adapt itself and perform sufficiently well under different conditions.
Xiaoming Liu 0003, Robbert van Renesse
PODC2
2000 A Probabilistically Correct Leader Election Protocol for Large Groups
Indranil Gupta, Robbert van Renesse, Kenneth P. Birman
DISC2
1999 Six Misconceptions about Reliable Distributed Computing
abstract
This paper describes how experiences with building industrial strength distributed applications have dramatically changed the assumptions about what tools are needed to build these systems.
Werner Vogels, Robbert van Renesse, Kenneth P. Birman
HPDC2
1999 Building reliable, high-performance communication systems from components
abstract
Although building systems from components has attractions, this approach also has problems. Can we be sure that a certain configuration of components is correct? Can it perform as well as a monolithic system? Our paper answers these questions for the Ensemble communication architecture by showing how, with help of the Nuprl formal system, configurations may be checked against specifications, and how optimized code can be synthesized from these configurations. The performance results show that we can substantially reduce end-to-end latency in the already optimized Ensemble system. Finally, we discuss whether the techniques we used are general enough for systems other than communication systems.
Xiaoming Liu 0003, Christoph Kreitz, Robbert van Renesse, Jason Hickey, Mark Hayden, Kenneth P. Birman, Robert L. Constable
SOSP3
1999 Specifications and Proofs for Ensemble Layers
Jason Hickey, Nancy A. Lynch, Robbert van Renesse
TACAS3
1998 Building Adaptive Systems Using Ensemble
abstract
Trends in networking and distributed computing are creating a new generation of applications that must adapt as the environment within which they execute changes. Examples of adaptation include switching protocols to overcome a security exposure or failure mode seen only in certain settings, changing data rates to accommodate a slow link, or adapting the behavior of a high level application to match the set of participants using the application. We describe the Ensemble system, a tool for building adaptive distributed programs. © 1998 John Wiley & Sons, Ltd.
Robbert van Renesse, Kenneth P. Birman, Mark Hayden, Alexey Vaysburd, David A. Karr
Softw. Pract. Exp.1
1997 Packing Messages as a Tool for Boosting the Performance of Total Ordering Protocols
abstract
This paper compares the throughput and latency of four protocols that provide total ordering. Two of these protocols are measured with and without message packing. We used a technique that buffers application messages for a short period of time before sending them, so more messages are packed together. The main conclusion of this comparison is that message packing influences the performance of total ordering protocols under high load overwhelmingly more than any other optimization that was checked in this paper, both in terms of throughput and latency. This improved performance is attributed to the fact that packing messages reduces the header overhead for messages, the contention on the network, and the load on the receiving CPUs.
Roy Friedman 0001, Robbert van Renesse
HPDC2
1997 Optimizing Layered Communication Protocols
abstract
Layering of communication protocols offers many well-known advantages but typically leads to performance inefficiencies. We present a model for layering, and point out where the performance problems occur in stacks of layers using this model. We then investigate the common execution paths in these stacks and how to identify them. These paths are optimized using three techniques: optimizing the computation, compressing protocol headers, and delaying processing. All of the optimizations can be automated in a compiler with the help of minor annotations by the protocol designer. We describe the performance that we obtain after implementing the optimizations by hand on a full-scale system.
Mark Hayden, Robbert van Renesse
HPDC2
1997 The Hierarchical Daisy Architecture for Causal Delivery
abstract
We propose the hierarchical daisy architecture, which provides causal delivery of messages sent to any subset of processes. The architecture provides fault tolerance and maintains the amount of control information within a reasonable size. It divides processes into logical groups. Messages inside a logical group are sent directly, while messages that need to cross logical group boundaries are forwarded by servers. We prove the correctness of the daisy architecture and discuss possible optimizations.
Roberto Baldoni, Roy Friedman 0001, Robbert van Renesse
ICDCS3
1996 Masking the Overhead of Protocol Layering
abstract
Protocol layering has been advocated as a way of dealing with the complexity of computer communication. It has also been criticized for its performance overhead. In this paper, we present some insights in the design of protocols, and how these insights can be used to mask the overhead of layering, in a way similar to client caching in a file system. With our techniques, we achieve an order of magnitude improvement in end-to-end message latency in the Horus communication framework. Over an ATM network, we are able to do a round-trip message exchange, of varying levels of semantics, in about 170 µseconds, using a protocol stack of four layers that were written in ML, a high-level functional language.
Robbert van Renesse
SIGCOMM1
1996 Strong and Weak Virtual Synchrony in Horus
abstract
This paper presents two variants of virtual synchrony, which are supported by Horus. The first variant, called strong virtual synchrony, includes the property that every message is delivered within the view in which it is sent. This property is very useful in developing applications, since it helps in minimizing the amount of context information that needs to be sent on messages, and the amount of computation which is required in order to process a message. However, it is shown that in order to support this property, the application program has to block messages during view changes. An alternative definition, called weak virtual synchrony, which can be implemented without blocking messages, is then presented. This definition still guarantees that messages will be delivered within the view in which they were sent, only that it uses a slightly weaker notion of what the view in which a message was sent is. An implementation of weak virtual synchrony that does not block messages during view changes as also developed in this paper.
Roy Friedman 0001, Robbert van Renesse
SRDS2
1996 A Transparent Light-Weight Group Service
abstract
The virtual synchrony model for group communication has proven to be a powerful paradigm for building distributed applications. Implementations of virtual synchrony usually require the use of failure detectors and failure recovery protocols. In applications that require the use of a large number of groups, significant performance gains can be attained if these groups share the resources required to provide virtual synchrony. A service that maps user groups onto instances of a virtually synchronous implementation is called a light-weight group service. This paper proposes a new design for the light-weight group protocols that enables the usage of this service in a transparent manner as a test case, the new design was implemented in the Horus system, although the underlying principles can be applied to other architectures as well. The paper also presents performance results from this implementation.
Luís E. T. Rodrigues, Katherine Guo, Antonio Sargento, Robbert van Renesse, Bradford B. Glade, Paulo Veríssimo, Kenneth P. Birman
SRDS4
1995 Operating system support for mobile agents
abstract
The TACOMA project is concerned with implementing operating system support for agents, processes that migrate through a network. Two TACOMA prototypes have been completed; this paper outlines our experiences in building and using them. A mechanism for exchanging electronic cash was explored, as well as agent-based schemes for scheduling and fault-tolerance.
Dag Johansen, Robbert van Renesse, Fred B. Schneider
HotOS2
1995 A Framework for Protocol Composition in Horus
abstract
The Horus system supports a communication architecture that treats protocols as instances of an abstract data type.This approach encourages developers to partition complex protocols into simple microprotocols, each of which is implemented by a protocol layer.Protocol layers can be stacked on top of each other in a variety of ways, at run-time.First, we describe the classes of protocols that can be supported this way.Next, we present the Horus object model that we designed for this technology, and the interface between the layers that makes it all work.We then present an example layer that implements a group membership protocol.Next, we show how, given a set of required properties, an appropriate stack can be constructed.We look at an example stack of protocols, which provides fault-tolerant, totally ordered communication between a group of processes.The work contributes a standard framework for protocol development and experimentation, provides a high performance implementation of the virtual synchrony model, and introduces a methodology for increasing the robustness of the protocol development process.
Robbert van Renesse, Kenneth P. Birman, Roy Friedman 0001, Mark Hayden, David A. Karr
PODC1
1994 A Security Architecture for Fault-Toerant Systems
abstract
Process groups are a common abstraction for fault-tolerant computing in distributed systems. We present a security architecture that extends the process group into a security abstraction. Integral parts of this architecture are services that securely and fault tolerantly support cryptographic key distribution. Using replication only when necessary, and introducing novel replication techniques when it was necessary, we have constructed these services both to be easily defensible against attack and to permit key distribution despite the transient unavailability of a substantial number of servers. We detail the design and implementation of these services and the secure process group abstraction they support. We also give preliminary performance figures for some common group operations.
Michael K. Reiter, Kenneth P. Birman, Robbert van Renesse
ACM Trans. Comput. Syst.3
1993 FLIP: An Internetwork Protocol for Supporting Distributed Systems
abstract
Most modern network protocols give adequate support for traditional applications such as file transfer and remote login. Distributed applications, however, have different requirements (e.g., efficient at-most-once remote procedure call even in the face of processor failures). Instead of using ad hoc protocols to meet each of the new requirements, we have designed a new protocol, called the Fast Local Internet Protocol (FLIP), that provides a clean and simple integrated approach to these new requirements. FLIP is an unreliable message protocol that provides both point-to-point communication and multicast communication, and requires almost no network management. Furthermore, by using FLIP we have simplified higher-level protocols such as remote procedure call and group communication, and enhanced support for process migration and security. A prototype implementation of FLIP has been built as part of the new kernel for the Amoeba distributed operating system, and is in daily use. Measurements of its performance are presented.
M. Frans Kaashoek, Robbert van Renesse, Hans van Staveren, Andrew S. Tanenbaum
ACM Trans. Comput. Syst.2
1991 The Amoeba distributed operating system - A status report
Andrew S. Tanenbaum, M. Frans Kaashoek, Robbert van Renesse, Henri E. Bal
Comput. Commun.3
1989 The Design of a High-Performance File Server
abstract
The Bullet server is a file server that outperforms traditional file servers by more than a factor of three. It achieves high throughput and low delay by a software design radically different from that of file servers currently in use. Whereas files are normally stored as a sequence of disk blocks, each Bullet server file is stored contiguously, both on disk and in the server's random access memory cache. Furthermore, it uses the concept of an immutable file to improve performance, to enable caching, and to provide a clean semantic model to the user. The authors describe the design and implementation of the Bullet server in detail, present measurements of its performance, and compare this performance to that of the SUN file server running on the same hardware.>
Robbert van Renesse, Andrew S. Tanenbaum, Annita N. Wilschut
ICDCS1
1989 The Performance of the Amoeba Distributed Operating system
abstract
Abstract Amoeba is a capability‐based distributed operating system designed for high‐performance interactions between clients and servers using the well‐known RPC model. The paper starts out by describing the architecture of the Amoeba system, which is typified by specialized components such as workstations, several services, a processor pool, and gateways that connect other Amoeba systems transparently over wide‐area networks. Next the RPC interface is described. The paper presents performance measurements of the Amoeba RPC on unloaded and loaded systems. The time to perform the simplest RPC between two user processes has been measured to be 1‐4 ms. Compared to SUN 3/50's RPC, Amoeba has one ninth of the delay, and over three times the throughput. Finally we describe the Amoeba file server. The Amoeba file server is so fast that it is limited by the communication bandwidth. To the best of our knowledge this is the fastest file server yet reported in the literature for this class of hardware.
Robbert van Renesse, Hans van Staveren, Andrew S. Tanenbaum
Softw. Pract. Exp.1
1988 Voting with Ghosts
abstract
A mechanism called voting with ghosts (VWG) is proposed for maintaining consistency of replicated data. VWG is an improvement of the weighted voting (WV) algorithm; it performs as well as the available copies (AC) algorithm, but unlike AC, works correctly in the face of network partitioning. A detailed description of the VWG method is given, and it is analyzed in the presence of node crashes and network partitions. Its performance is compared with that of WV and AC.>
Robbert van Renesse, Andrew S. Tanenbaum
ICDCS1
1987 Connecting RPC-Based Distributed Systems Using Wide-Area Networks
Robbert van Renesse, Andrew S. Tanenbaum, Hans van Staveren, J. Hall
ICDCS1
1987 Reliability Issues in Distributed Operating Systems
Andrew S. Tanenbaum, Robbert van Renesse
SRDS2
1986 Using Sparse Capabilities in a Distributed Operating System
Andrew S. Tanenbaum, Sape J. Mullender, Robbert van Renesse
ICDCS3