Rodrigo Rodrigues 0001

dblp:87/6689 · also Rodrigo Seromenho Miragaia Rodrigues · DBLP profile ↗
← Back
62ranked-venue papers
3as first author
9since 2021 · last 2026
0000-0001-8367-4024ORCID · verified

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

Systems, architecture and hardware · 30 · 1 first-author · 4 since 2021Security and privacy · 13 · 1 first-author · 2 since 2021Software engineering, systems software and programming languages · 13 · 1 first-author · 2 since 2021Databases, data management, data science and information retrieval · 6Computer networks · 5 · 1 since 2021Applied, interdisciplinary, general and emerging computing · 2
YearPublicationVenuePosition
2026 Vardalith: Hybrid Detection of Persistent Memory Concurrency Bugs
abstract
Persistent Memory offers byte-addressable persistence but exposes developers to new concurrency bugs - persistency-induced races - where a thread might read unpersisted data, potentially leading to inconsistencies after crashes. Existing tools face important practical limitations: they either require exhaustive exploration of thread interleavings, depend on application-specific semantics or specialized testing drivers, or report many interleavings that do not correspond to real persistency-induced races. This paper introduces a hybrid approach for detecting persistency-induced races that overcomes these limitations. Our method operates without application-specific knowledge and does not require observing the exact racy interleaving during testing. Instead, by precisely extending the detection window around persistent memory accesses, we can infer the existence of racy interleavings whenever conflicting executions are observed. Our evaluation across multiple applications found 26 bugs (7 new) demonstrating that our approach provides a principled and practical foundation for detecting persistency-induced races.
José Fragoso Santos, Rodrigo Rodrigues 0001, Miguel Matos
ECOOP3
2025 CoVault: Secure, Scalable Analytics of Personal Data
Roberta De Viti, Isaac Sheff, Noemi Glaeser, Baltasar Dinis, Rodrigo Rodrigues 0001, Bobby Bhattacharjee, Anwar Hithnawi, Deepak Garg 0001, Peter Druschel
USENIX Security Symposium5
2024 Alea-BFT: Practical Asynchronous Byzantine Fault Tolerance
Diogo S. Antunes, Afonso N. Oliveira, André Breda, Matheus Guilherme Franco, Henrique Moniz, Rodrigo Rodrigues 0001
NSDI6
2023 With Great Freedom Comes Great Opportunity: Rethinking Resource Allocation for Serverless Functions
abstract
Current serverless offerings give users limited flexibility for configuring the resources allocated to their function invocations. This simplifies the interface for users to deploy server-less computations but creates deployments that are resource inefficient. In this paper, we take a principled approach to the problem of resource allocation for serverless functions, analyzing the effects of automating this choice in a way that leads to the best combination of performance and cost. In particular, we systematically explore the opportunities that come with decoupling memory and CPU resource allocations and also enabling the use of different VM types, and we find a rich trade-off space between performance and cost. The provider can use this in a number of ways, e.g., exposing all these parameters to the user; eliding preferences for performance and cost from users and simply offer the same performance with lower cost; or exposing a small number of choices for users to trade performance for cost.
Muhammad Bilal 0007, Marco Canini, Rodrigo Fonseca, Rodrigo Rodrigues 0001
EuroSys4
2023 Mumak: Efficient and Black-Box Bug Detection for Persistent Memory
abstract
The advent of Persistent Memory (PM) opens the door to novel application designs that explore its performance and durability benefits. However, there is no free lunch, and to program PM applications, developers need to be aware of potential inconsistent application state upon machine or application crashes. To overcome this difficulty, several tools have been proposed to detect the presence of the so-called crash-consistency bugs. While these are effective in detecting a variety of bugs, they present several key limitations, namely relying on application-specific semantics, requiring the programmer to manually annotate the program or modify the PM library, and relying on techniques with poor scalability, making them impractical for production code.
Miguel Matos, Rodrigo Rodrigues 0001
EuroSys3
2023 RR: A Fault Model for Efficient TEE Replication
Baltasar Dinis, Peter Druschel, Rodrigo Rodrigues 0001
NDSS3
2023 Antipode: Enforcing Cross-Service Causal Consistency in Distributed Applications
abstract
Modern internet-scale applications suffer from cross-service inconsistencies, arising because applications combine multiple independent and mutually-oblivious datastores. The end-to-end execution flow of each user request spans many different services and datastores along the way, implicitly establishing ordering dependencies among operations at different datastores. Readers should observe this ordering and, in today's systems, they do not.
João Loff, Daniel Porto 0002, João Garcia 0001, Jonathan Mace, Rodrigo Rodrigues 0001
SOSP5
2022 SconeKV: A Scalable, Strongly Consistent Key-Value Store
abstract
For decades, relational databases provided a strong foundation for constructing applications due to their ACID properties. However, distributed applications reached a scale, both in terms of data volume and number of concurrent clients, that traditional databases cannot accommodate. NoSQL databases addressed this problem by trading consistency for scalability, namely through horizontal scalability schemes supported by optimistic replication protocols, which only guarantee eventual consistency. In this paper, we explore a novel design between the two extremes, which is able to scale to large deployments while still offering strong consistency guarantees in the form of serializable transactions. Our key insight is to leverage recent advances in membership services that provide strongly consistent views at scale. Those assurances from the membership layer simplify building efficient and consistent storage protocols. Our evaluation of the resulting system,SconeKV, in a realistic scenario shows that it scales and performs better than CockroachDB while being competitive with Cassandra.
Miguel Matos, Rodrigo Rodrigues 0001
IEEE Trans. Parallel Distributed Syst.3
2021 Particle-In-Cell Simulation Using Asynchronous Tasking
Nicolas L. Guidotti, Pedro Ceyrat, João Barreto 0001, José Monteiro 0001, Rodrigo Rodrigues 0001, Ricardo Fonseca, Xavier Martorell, Antonio J. Peña
Euro-Par5
2020 Finding the right cloud configuration for analytics clusters
abstract
Finding good cloud configurations for deploying a single distributed system is already a challenging task, and it becomes substantially harder when a data analytics cluster is formed by multiple distributed systems since the search space becomes exponentially larger. In particular, recent proposals for single system deployments rely on benchmarking runs that become prohibitively expensive as we shift to joint optimization of multiple systems, as users have to wait until the end of a long optimization run to start the production run of their job.
Muhammad Bilal 0007, Marco Canini, Rodrigo Rodrigues 0001
SoCC3
2020 Bandwidth-Aware Page Placement in NUMA
abstract
Page placement is a critical problem for memory-intensive applications running on a shared-memory multiprocessor with a non-uniform memory access (NUMA) architecture. State-of-the-art page placement mechanisms interleave pages evenly across NUMA nodes. However, this approach fails to maximize memory throughput in modern NUMA systems, characterized by asymmetric bandwidths and latencies, and sensitive to memory contention and interconnect congestion phenomena.We propose BWAP, a novel page placement mechanism based on asymmetric weighted page interleaving. BWAP combines an analytical performance model of the target NUMA system with on-line iterative tuning of page distribution for a given memory-intensive application. Our experimental evaluation with representative memory-intensive workloads shows that BWAP performs up to 66% better than state-of-the-art techniques. These gains are particularly relevant when multiple co-located applications run in disjoint partitions of a large NUMA machine or when applications do not scale up to the total number of cores.
David Gureya, João Neto 0001, Reza Karimi, João Barreto 0001, Pramod Bhatotia, Vivien Quéma, Rodrigo Rodrigues 0001, Paolo Romano 0002, Vladimir Vlassov
IPDPS7
2020 Do the Best Cloud Configurations Grow on Trees? An Experimental Evaluation of Black Box Algorithms for Optimizing Cloud Workloads Sub
Muhammad Bilal 0007, Marco Serafini, Marco Canini, Rodrigo Rodrigues 0001
Proc. VLDB Endow.4
2019 A More Consistent Understanding of Consistency
abstract
Recent storage systems trade strong consistency for performance, availability, and scalability. However, this makes it hard to understand the semantics that the storage system provides, and also makes the design and implementation of the storage system itself more error-prone. This paper proposes a comprehensive solution to these problems. In particular, we propose a specification language named ConSpec, which enables the formalization of different consistency semantics that a storage system may provide, using a uniform syntax that is independent of the design and implementation of the target storage system. We use ConSpec to revisit several existing models in light of a common way to define and compare them. Furthermore, we generalize the CAP theorem, whose original formulation only considered linearizability, to precisely define the class of consistency definitions that can and cannot be implemented in a highly-available, partition-tolerant way. Finally, we present the design and implementation of a new consistency checker that takes a trace from a storage system (e.g., the output of a test suite) and validates whether it meets any consistency semantics defined using ConSpec. The evaluation of our consistency checker shows that it is able to verify the correctness of long traces in a reasonable time.
Subhajit Sidhanta, Ricardo J. Dias, Rodrigo Rodrigues 0001
SRDS3
2018 The Tortoise and the Hare: Characterizing Synchrony in Distributed Environments (Practical Experience Report)
abstract
The design of distributed protocols that run in data centers and enterprise clusters is heavily dependent on synchrony assumptions regarding the timing behavior of the participating nodes and the network. However, little is known about the actual synchrony of real distributed systems, and how it varies across deployments. To better understand this timing behavior and how it impacts the design and implementation of distributed protocols, we conduct an extensive measurement study of the latency for transmitting and processing messages between nodes in four different environments. Our study determines how protocol characteristics affect the latency behavior. We also determine how different environmental factors can affect the measured latency and whether high latency events manifest globally or locally. Our results suggest several directions for reducing latency, and for leveraging recent distributed computing models in a more judicious way.
Daniel Porto 0002, João Leitão 0001, Flavio Paiva Junqueira, Rodrigo Rodrigues 0001
DSN4
2018 Fine-grained consistency for geo-replicated systems
Cheng Li 0001, Nuno M. Preguiça, Rodrigo Rodrigues 0001
USENIX ATC3
2018 IPA: Invariant-preserving Applications for Weakly consistent Replicated Databases
abstract
It is common to use weakly consistent replication to achieve high availability and low latency at a global scale. In this setting, concurrent updates may lead to states where application invariants do not hold. Some systems coordinate the execution of (conflicting) operations to avoid invariant violations, leading to high latency and reduced availability for those operations. This problem is worsened by the difficulty in identifying precisely which operations conflict. In this paper we propose a novel approach to preserve application invariants without coordinating the execution of operations. The approach consists of modifying operations in a way that application invariants are maintained in the presence of concurrent updates. When no conflicting updates occur, the modified operations present their original semantics. Otherwise, we use sensible and deterministic conflict resolution policies that preserve the invariants of the application. To implement this approach, we developed a static analysis, IPA, that identifies conflicting operations and proposes the necessary modifications to operations. Our analysis shows that IPA can avoid invariant violations in many applications, including typical database applications. Our evaluation reveals that the offline static analysis runs fast enough for being used with large applications. The overhead introduced in the modified operations is low and it leads to lower latency and higher throughput when compared with other approaches that enforce invariants.
Valter Balegas, Sérgio Duarte, Carla Ferreira 0001, Rodrigo Rodrigues 0001, Nuno M. Preguiça
Proc. VLDB Endow.4
2017 Fine-Grained Consistency Upgrades for Online Services
abstract
Online services such as Facebook or Twitter have public APIs to enable an easy integration of these services with third party applications. However, the developers who design these applications have no information about the consistency provided by these services, which exacerbates the complexity of reasoning about the semantics of the applications they are developing. In this paper, we show that is possible to deploy a transparent middleware between the application and the service, which enables a fine-grained control over the session guarantees that comprise the consistency semantics provided by these APIs, without having to gain access to the implementation of the underlying services. We evaluated our middleware using the Facebook public API and the Redis datastore, and our results show that we are able to provide fine-grained control of the consistency semantics incurring in a small local storage and modest latency overhead.
Filipe Freitas, João Leitão 0001, Nuno M. Preguiça, Rodrigo Rodrigues 0001
SRDS4
2017 Generalized Paxos Made Byzantine (and Less Complex)
Miguel Pires, Srivatsan Ravi, Rodrigo Rodrigues 0001
SSS3
2017 Blotter: Low Latency Transactions for Geo-Replicated Storage
abstract
Most geo-replicated storage systems use weak consistency to avoid the performance penalty of coordinating replicas in different data centers. This departure from strong semantics poses problems to application programmers, who need to address the anomalies enabled by weak consistency. In this paper we use a recently proposed isolation level, called Non-Monotonic Snapshot Isolation, to achieve ACID transactions with low latency. To this end, we present Blotter, a geo-replicated system that leverages these semantics in the design of a new concurrency control protocol that leaves a small amount of local state during reads to make commits more efficient, which is combined with a configuration of Paxos that is tailored for good performance in wide area settings. Read operations always run on the local data center, and update transactions complete in a small number of message steps to a subset of the replicas. We implemented Blotter as an extension to Cassandra. Our experimental evaluation shows that Blotter has a small overhead at the data center scale, and performs better across data centers when compared with our implementations of the core Spanner protocol and of Snapshot Isolation on the same codebase.
Henrique Moniz, João Leitão 0001, Ricardo J. Dias, Johannes Gehrke, Nuno M. Preguiça, Rodrigo Rodrigues 0001
WWW6
2016 Characterizing the Consistency of Online Services (Practical Experience Report)
abstract
While several proposals for the specification and implementation of various consistency models exist, little is known about what is the consistency currently offered by online services with millions of users. Such knowledge is important, not only because it allows for setting the right expectations and justifying the behavior observed by users, but also because it can be used for improving the process of developing applications that use APIs offered by such services. To fill this gap, this paper presents a measurement study of the consistency of the APIs exported by four widely used Internet services, the Facebook Feed, Facebook Groups, Blogger, and Google+. To conduct this study, our work (1) proposes definitions for a set of relevant consistency properties, (2) develops a simple, yet generic methodology comprising a small number of tests, which probe these services from a user perspective, and try to uncover consistency anomalies that are key to our definitions, and (3) reports on the analysis of the data obtained from running these tests for a period of several weeks. Our measurement study shows that some of these services do exhibit consistency anomalies, including some behaviors that may appear counter-intuitive for users, such as the lack of session guarantees for write monotonicity.
Filipe Freitas, João Leitão 0001, Nuno M. Preguiça, Rodrigo Rodrigues 0001
DSN4
2016 IncApprox: A Data Analytics System for Incremental Approximate Computing
abstract
Incremental and approximate computations are increasingly being adopted for data analytics to achieve low-latency execution and efficient utilization of computing resources. Incremental computation updates the output incrementally instead of re-computing everything from scratch for successive runs of a job with input changes. Approximate computation returns an approximate output for a job instead of the exact output. Both paradigms rely on computing over a subset of data items instead of computing over the entire dataset, but they differ in their means for skipping parts of the computation. Incremental computing relies on the memoization of intermediate results of sub-computations, and reusing these memoized results across jobs. Approximate computing relies on representative sampling of the entire dataset to compute over a subset of data items.
Dhanya R. Krishnan, Do Le Quoc, Pramod Bhatotia, Christof Fetzer, Rodrigo Rodrigues 0001
WWW5
2015 iThreads: A Threading Library for Parallel Incremental Computation
abstract
Incremental computation strives for efficient successive runs of applications by re-executing only those parts of the computation that are affected by a given input change instead of recomputing everything from scratch. To realize these benefits automatically, we describe iThreads, a threading library for parallel incremental computation. iThreads supports unmodified shared-memory multithreaded programs: it can be used as a replacement for pthreads by a simple exchange of dynamically linked libraries, without even recompiling the application code. To enable such an interface, we designed algorithms and an implementation to operate at the compiled binary code level by leveraging MMU-assisted memory access tracking and process-based thread isolation. Our evaluation on a multicore platform using applications from the PARSEC and Phoenix benchmarks and two case-studies shows significant performance gains.
Pramod Bhatotia, Pedro Fonseca 0001, Umut A. Acar, Björn B. Brandenburg, Rodrigo Rodrigues 0001
ASPLOS5
2015 Putting consistency back into eventual consistency
abstract
Geo-replicated storage systems are at the core of current Internet services. The designers of the replication protocols used by these systems must choose between either supporting low-latency, eventually-consistent operations, or ensuring strong consistency to ease application correctness. We propose an alternative consistency model, Explicit Consistency, that strengthens eventual consistency with a guarantee to preserve specific invariants defined by the applications. Given these application-specific invariants, a system that supports Explicit Consistency identifies which operations would be unsafe under concurrent execution, and allows programmers to select either violation-avoidance or invariant-repair techniques. We show how to achieve the former, while allowing operations to complete locally in the common case, by relying on a reservation system that moves coordination off the critical path of operation execution. The latter, in turn, allows operations to execute without restriction, and restore invariants by applying a repair operation to the database state. We present the design and evaluation of Indigo, a middleware that provides Explicit Consistency on top of a causally-consistent data store. Indigo guarantees strong application invariants while providing similar latency to an eventually-consistent system in the common case.
Valter Balegas, Sérgio Duarte, Carla Ferreira 0001, Rodrigo Rodrigues 0001, Nuno M. Preguiça, Mahsa Najafzadeh, Marc Shapiro 0001
EuroSys4
2015 Visigoth fault tolerance
abstract
We present a new technique for designing distributed protocols for building reliable stateful services called Visigoth Fault Tolerance (VFT). VFT introduces the Visigoth model, which makes it possible to calibrate the timing assumptions of a system using a threshold of slow processes or messages, and also to distinguish between non-malicious arbitrary faults and correlated attack scenarios. This enables solutions that leverage the characteristics of data center systems, namely their secure environment and predictable performance, in order to allow replicated systems to be more efficient with respect to the utilization of resources than those designed under asynchrony and Byzantine assumptions, while avoiding the need to make a system synchronous, or to restrict failure modes to silent crashes. We implemented a VFT protocol for a state machine replication library, and ran several benchmarks. Our evaluation shows that VFT has comparable performance to existing schemes and brings significant benefits in terms of the throughput per dollar, i.e., the server cost for sustaining a certain level of request execution.
Daniel Porto 0002, João Leitão 0001, Cheng Li 0001, Allen Clement, Aniket Kate, Flavio Paiva Junqueira, Rodrigo Rodrigues 0001
EuroSys7
2015 Guardat: enforcing data policies at the storage layer
abstract
In today's data processing systems, both the policies protecting stored data and the mechanisms for their enforcement are spread over many software components and configuration files, increasing the risk of policy violation due to bugs, vulnerabilities and misconfigurations. Guardat addresses this problem. Users, developers and administrators specify file protection policies declaratively, concisely and separate from code, and Guardat enforces these policies by mediating I/O in the storage layer. Policy enforcement relies only on the integrity of the Guardat controller and any external policy dependencies. The semantic gap between the storage layer enforcement and per-file policies is bridged using cryptographic attestations from Guardat. We present the design and prototype implementation of Guardat, enforce example policies in a Web server, and show experimentally that its overhead is low.
Anjo Vahldiek-Oberwagner, Eslam Elnikety, Aastha Mehta, Deepak Garg 0001, Peter Druschel, Rodrigo Rodrigues 0001, Johannes Gehrke, Ansley Post
EuroSys6
2015 Extending Eventually Consistent Cloud Databases for Enforcing Numeric Invariants
abstract
Geo-replicated databases often offer high availability and low latency by relying on weak consistency models. The inability to enforce invariants across all replicas remains a key shortcoming that prevents the adoption of such databases in several applications. In this paper we show how to extend an eventually consistent cloud database for enforcing numeric invariants. Our approach builds on ideas from escrow transactions, but our novel design overcomes the limitations of previous works. First, by relying on a new replicated data type, our design has no central authority and uses pairwise asynchronous communication only. Second, by layering our design on top of a fault-tolerant database, our approach exhibits better availability during network partitions and data center faults. The evaluation of our prototype, built on top of Riak, shows much lower latency and better scalability than the traditional approach of using strong consistency to enforce numeric invariants.
Valter Balegas, Diogo Serra, Sérgio Duarte, Carla Ferreira 0001, Marc Shapiro 0001, Rodrigo Rodrigues 0001, Nuno M. Preguiça
SRDS6
2015 PIXIDA: Optimizing Data Parallel Jobs in Wide-Area Data Analytics
abstract
In the era of global-scale services, big data analytical queries are often required to process datasets that span multiple data centers (DCs). In this setting, cross-DC bandwidth is often the scarcest, most volatile, and/or most expensive resource. However, current widely deployed big data analytics frameworks make no attempt to minimize the traffic traversing these links. In this paper, we present P ixida , a scheduler that aims to minimize data movement across resource constrained links. To achieve this, we introduce a new abstraction called S ilo , which is key to modeling P ixida 's scheduling goals as a graph partitioning problem. Furthermore, we show that existing graph partitioning problem formulations do not map to how big data jobs work, causing their solutions to miss opportunities for avoiding data movement. To address this, we formulate a new graph partitioning problem and propose a novel algorithm to solve it. We integrated P ixida in Spark and our experiments show that, when compared to existing schedulers, P ixida achieves a significant traffic reduction of up to ~ 9x on the aforementioned links.
Konstantinos Kloudas, Rodrigo Rodrigues 0001, Nuno M. Preguiça, Margarida Mamede
Proc. VLDB Endow.2
2014 Slider: incremental sliding window analytics
abstract
Sliding window analytics is often used in distributed data-parallel computing for analyzing large streams of continuously arriving data. When pairs of consecutive windows overlap, there is a potential to update the output incrementally, more efficiently than recomputing from scratch. However, in most systems, realizing this potential requires programmers to explicitly manage the intermediate state for overlapping windows, and devise an application-specific algorithm to incrementally update the output.
Pramod Bhatotia, Umut A. Acar, Flavio Paiva Junqueira, Rodrigo Rodrigues 0001
Middleware4
2014 SKI: Exposing Kernel Concurrency Bugs through Systematic Schedule Exploration
Pedro Fonseca 0001, Rodrigo Rodrigues 0001, Björn B. Brandenburg
OSDI2
2014 Automating the Choice of Consistency Levels in Replicated Systems
Cheng Li 0001, João Leitão 0001, Allen Clement, Nuno M. Preguiça, Rodrigo Rodrigues 0001, Viktor Vafeiadis
USENIX ATC5
2012 Scalable testing of file system checkers
abstract
File system checkers (like e2fsck) are critical, complex, and hard to develop, and developers today rely on hand-written tests to exercise this intricate code. Test suites for file system checkers take a lot of effort to develop and require careful reasoning to cover a sufficiently comprehensive set of inputs and recovery mechanisms. We present a tool and methodology for testing file system checkers that reduces the need for a specification of the recovery process and the development of a test suite. Our methodology splits the correctness of the checker into two objectives: consistency and completeness of recovery. For each objective, we leverage either the file system checker code itself or a comparison among the outputs of multiple checkers to extract an implicit specification of correct behavior. Our methodology is embodied in a testing tool called SWIFT, which uses a mix of symbolic and concrete execution; it introduces two new techniques: a specific concretization strategy and a corruption model that leverages test suites of file system checkers. We used SWIFT to test the file system checkers of ext2, ext3, ext4, ReiserFS, and Minix; we found bugs in all checkers, including cases leading to data loss. Additionally, we automatically generated test suites achieving code coverage on par with manually constructed test suites shipped with the checkers.
João Carlos Menezes Carreira, Rodrigo Rodrigues 0001, George Candea, Rupak Majumdar
EuroSys2
2012 Shredder: GPU-accelerated incremental storage and computation
Pramod Bhatotia, Rodrigo Rodrigues 0001, Akshat Verma
FAST2
2012 Enhancing the OS against Security Threats in System Administration
Nuno Santos 0001, Rodrigo Rodrigues 0001, Bryan Ford
Middleware2
2012 Orchestrating the Deployment of Computations in the Cloud with Conductor
Alexander Wieder, Pramod Bhatotia, Ansley Post, Rodrigo Rodrigues 0001
NSDI4
2012 Making Geo-Replicated Systems Fast as Possible, Consistent when Necessary
Cheng Li 0001, Daniel Porto 0002, Allen Clement, Johannes Gehrke, Nuno M. Preguiça, Rodrigo Rodrigues 0001
OSDI6
2012 On the (limited) power of non-equivocation
abstract
In recent years, there have been a few proposals to add a small amount of trusted hardware at each replica in a Byzantine fault tolerant system to cut back replication factors. These trusted components eliminate the ability for a Byzantine node to perform equivocation, which intuitively means making conflicting statements to different processes.
Allen Clement, Flavio Paiva Junqueira, Aniket Kate, Rodrigo Rodrigues 0001
PODC4
2012 Policy-Sealed Data: A New Abstraction for Building Trusted Cloud Services
Nuno Santos 0001, Rodrigo Rodrigues 0001, Krishna P. Gummadi, Stefan Saroiu
USENIX Security Symposium2
2012 Adaptive Search Radius - Using hop count to reduce P2P traffic
Ricardo J. F. Lopes Pereira, Teresa Vazão, Rodrigo Rodrigues 0001
Comput. Networks3
2012 Automatic Reconfiguration for Large-Scale Reliable Storage Systems
abstract
Byzantine-fault-tolerant replication enhances the availability and reliability of Internet services that store critical state and preserve it despite attacks or software errors. However, existing Byzantine-fault-tolerant storage systems either assume a static set of replicas, or have limitations in how they handle reconfigurations (e.g., in terms of the scalability of the solutions or the consistency levels they provide). This can be problematic in long-lived, large-scale systems where system membership is likely to change during the system lifetime. In this paper, we present a complete solution for dynamically changing system membership in a large-scale Byzantine-fault-tolerant system. We present a service that tracks system membership and periodically notifies other system nodes of membership changes. The membership service runs mostly automatically, to avoid human configuration errors; is itself Byzantine-fault-tolerant and reconfigurable; and provides applications with a sequence of consistent views of the system membership. We demonstrate the utility of this membership service by using it in a novel distributed hash table called dBQS that provides atomic semantics even across changes in replica sets. dBQS is interesting in its own right because its storage algorithms extend existing Byzantine quorum protocols to handle changes in the replica set, and because it differs from previous DHTs by providing Byzantine fault tolerance and offering strong semantics. We implemented the membership service and dBQS. Our results show that the approach works well, in practice: the membership service is able to manage a large system and the cost to change the system membership is low.
Rodrigo Rodrigues 0001, Barbara Liskov, Kathryn Chen, Moses D. Liskov, David A. Schultz
IEEE Trans. Dependable Secur. Comput.1
2011 Incoop: MapReduce for incremental computations
abstract
Many online data sets evolve over time as new entries are slowly added and existing entries are deleted or modified. Taking advantage of this, systems for incremental bulk data processing, such as Google's Percolator, can achieve efficient updates. To achieve this efficiency, however, these systems lose compatibility with the simple programming models offered by non-incremental systems, e.g., MapReduce, and more importantly, requires the programmer to implement application-specific dynamic algorithms, ultimately increasing algorithm and code complexity.
Pramod Bhatotia, Alexander Wieder, Rodrigo Rodrigues 0001, Umut A. Acar, Rafael Pasquini
SoCC3
2011 Finding complex concurrency bugs in large multi-threaded applications
abstract
Parallel software is increasingly necessary to take advantage of multi-core architectures, but it is also prone to concurrency bugs which are particularly hard to avoid, find, and fix, since their occurrence depends on specific thread interleavings. In this paper we propose a concurrency bug detector that automatically identifies when an execution of a program triggers a concurrency bug. Unlike previous concurrency bug detectors, we are able to find two particularly hard classes of bugs. The first are bugs that manifest themselves by subtle violation of application semantics, such as returning an incorrect result. The second are latent bugs, which silently corrupt internal data structures, and are especially hard to detect because when these bugs are triggered they do not become immediately visible. Pike detects these concurrency bugs by checking both the output and the internal state of the application for linearizability at the level of user requests. This paper presents this technique for finding concurrency bugs, its application in the context of a testing tool that systematically searches for such problems, and our experience in applying our approach to MySQL, a large-scale complex multi-threaded application. We were able to find several concurrency bugs in a stable version of the application, including subtle violations of application semantics, latent bugs, and incorrect error replies.
Pedro Fonseca 0001, Cheng Li 0001, Rodrigo Rodrigues 0001
EuroSys3
2011 Efficient middleware for byzantine fault tolerant database replication
abstract
Byzantine fault tolerance (BFT) enhances the reliability and availability of replicated systems subject to software bugs, malicious attacks, or other unexpected events. This paper presents Byzantium, a BFT database replication middleware that provides snapshot isolation semantics. It is the first BFT database system that allows for concurrent transaction execution without relying on a centralized component, which is essential for having both performance and robustness. Byzantium builds on an existing BFT library but extends it with a set of techniques for increasing concurrency in the execution of operations, for optimistically executing operations in a single replica, and for striping and load-balancing read operations across replicas. Experimental results show that our replication protocols introduce only a modest performance overhead for read-write dominated workloads and perform better than a non-replicated database system for read-only workloads.
Rui Garcia, Rodrigo Rodrigues 0001, Nuno M. Preguiça
EuroSys2
2010 A study of the internal and external effects of concurrency bugs
abstract
Concurrent programming is increasingly important for achieving performance gains in the multi-core era, but it is also a difficult and error-prone task. Concurrency bugs are particularly difficult to avoid and diagnose, and therefore in order to improve methods for handling such bugs, we need a better understanding of their characteristics. In this paper we present a study of concurrency bugs in MySQL, a widely used database server. While previous studies of real-world concurrency bugs exist, they have centered their attention on the causes of these bugs. In this paper we provide a complementary focus on their effects, which is important for understanding how to detect or tolerate such bugs at run-time. Our study uncovered several interesting facts, such as the existence of a significant number of latent concurrency bugs, which silently corrupt data structures and are exposed to the user potentially much later. We also highlight several implications of our findings for the design of reliable concurrent systems.
Pedro Fonseca 0001, Cheng Li 0001, Vishal Singhal, Rodrigo Rodrigues 0001
DSN4
2010 Accountable Virtual Machines
Andreas Haeberlen, Paarijaat Aditya, Rodrigo Rodrigues 0001, Peter Druschel
OSDI3
2010 Brief announcement: modelling MapReduce for optimal execution in the cloud
abstract
We describe a model for MapReduce computations that can be used to optimize the increasingly complex choice of resources that cloud customers purchase.
Alexander Wieder, Pramod Bhatotia, Ansley Post, Rodrigo Rodrigues 0001
PODC4
2009 Fifth Workshop on Hot Topics in System Dependability (HotDep 2009)
abstract
The FifthWorkshop on Hot Topics in System Dependability (HotDep'09) brings forth cutting-edge research ideas in fault tolerance, reliability and systems. This year's edition of the workshop will feature a total of 10 presentations of original research on wide array of topics, including cloud computing, storage, program analysis, operating systems, replication protocols, or failure prediction.
Christof Fetzer, Rodrigo Rodrigues 0001
DSN2
2009 Verme: Worm containment in overlay networks
abstract
Topological worms, such as those that propagate by following links in an overlay network, have the potential to spread faster than traditional random scanning worms because they have knowledge of a subset of the overlay nodes, and choose these nodes to propagate themselves; and also because they can avoid traditional detection mechanisms. Furthermore, this worm propagation strategy is likely to become prevalent as the deployment of networks with a sparse address space, such as IPv6, makes the traditional random scanning strategy futile. We present a novel approach for containing topological worms based on the fact that some overlay nodes may not have common vulnerabilities, due to their platform diversity. By reorganizing the overlay graph, it is possible to contain topological worms in small islands of nodes with common vulnerabilities that only have knowledge of themselves or nodes running on distinct platforms. We also present the design of Verme, a peer-to-peer overlay based on Chord that follows this approach, and VerDi, a DHT layer built on top of the Verme routing overlay. Simulations show that Verme and VerDi have a low overhead when compared to Chord's corresponding layers, and that our new overlay design helps containing, or at least slowing down the propagation of topological worms.
Filipe Freitas, Edgar Marques, Rodrigo Rodrigues 0001, Carlos Ribeiro, Paulo Ferreira 0001, Luís E. T. Rodrigues
DSN3
2009 Zeno: Eventually Consistent Byzantine-Fault Tolerance
Atul Singh, Pedro Fonseca 0001, Petr Kuznetsov, Rodrigo Rodrigues 0001, Petros Maniatis
NSDI4
2009 Full-Information Lookups for Peer-to-Peer Overlays
abstract
Most peer-to-peer lookup schemes keep a small amount of routing state per node, typically logarithmic in the number of overlay nodes. This design assumes that routing information at each member node must be kept small so that the bookkeeping required to respond to system membership changes is also small, given that aggressive membership dynamics are expected. As a consequence, lookups have high latency as each lookup requires contacting several nodes in sequence. In this paper, we question these assumptions by presenting a peer-to-peer routing algorithm with small lookup paths. Our algorithm, called ldquoOneHop,rdquo maintains full information about the system membership at each node, routing in a single hop whenever that information is up to date and in a small number of hops otherwise. We show how to disseminate information about membership changes quickly enough so that nodes maintain accurate complete membership information. We also present analytic bandwidth requirements for our scheme that demonstrate that it could be deployed in systems with hundreds of thousands of nodes and high churn. We validate our analytic model using a simulated environment and a real implementation. Our results confirm that OneHop is able to achieve high efficiency, usually reaching the correct node directly 99 percent of the time.
Pedro Fonseca 0001, Rodrigo Rodrigues 0001, Barbara Liskov
IEEE Trans. Parallel Distributed Syst.2
2008 Pastel: Bridging the Gap between Structured and Large-State Overlays
abstract
Peer-to-peer overlays envision a single overlay substrate that can be used (possibly simultaneously) by many applications, but current overlays either target fast, few-hop lookups for contacting directly the responsible nodes, or slower multi-hop lookups that can be used by applications that exploit the overlay topology (like multicast or anycast). In this paper we present Pastel, an extension to Pastry that bridges the gap between the two types of overlays. Pastel maintains both Pastry routing tables and a full information table, and we show how we can exploit synergies between the maintenance of the two. We also propose a novel API that is richer than the one offered by existing overlays, to give applications control over the type of lookups (structured, multi-hop routing, or attempt direct contact). We implemented Pastel in a discrete-event packet level simulator and our results show that Pastel has lookups that are usually more efficient than Pastry's. Furthermore, the bandwidth required by Pastel is modest, even for a system with thousands of nodes.
Nuno Cruces, Rodrigo Rodrigues 0001, Paulo Ferreira 0001
CCGRID2
2007 GiGi: An Ocean of Gridlets on a "Grid-for-the-Masses"
abstract
There have been a few proposals aiming at bridging the gap between institutional grid infrastructures (e.g., Globus-based), popular cycle-sharing applications (e.g., SETIQhome), and massively used decentralized P2P file-sharing applications. Nonetheless, no such infrastructure was ever successful in allowing, in a large-scale, home users to run popular desktop applications faster, by using spare cycles in other users' machines and, in return, donate their spare cycles to run other users' applications. We present a novel application and programming model that was designed to overcome some of the barriers to the deployment of a generic peer-to-peer grid infrastructure. In particular, we want to enable a trivial deployment in such infrastructures of existing applications that are in widespread use but do not currently exploit parallelism for improved performance. The model presented in this paper revolves around the concept of a Gridlet, a semantics-aware unit of workload division and computation off-load. A gridlet is a chunk of data associated with the operations to be performed on the data, and in many cases these operations consist of unmodified application binaries. Moreover, the concept of gridlet is also employed for resource management, and accounting of peer contribution. We believe this new concept, absent in other proposals, will significantly lower the barriers for exploiting parallel execution in popular applications, thus improving the chances of the gridlet model being widely adopted.
Luís Veiga, Rodrigo Rodrigues 0001, Paulo Ferreira 0001
CCGRID2
2007 Adaptive Search Radius - Lowering Internet P2P File-Sharing Traffic through Self-Restraint
abstract
Peer-to-peer (P2P) file sharing accounts for a very significant part of the Internet's traffic, translating into significant peering costs for ISPs. It has been noticed that, just like WWW traffic, P2P file sharing traffic shows locality properties, which are not exploited by current P2P file sharing protocols. We propose a novel peer selection algorithm, adaptive search radius (ASR), whose primary goal is to reduce ISPs' peering costs, where peers exploit locality by only downloading from those other peers which are nearest (in network hops). Simulation studies, using the eMule protocol, show that ASR benefits both ISPs, by globally reducing P2P file sharing traffic, and users, who experience faster downloads.
Ricardo J. F. Lopes Pereira, Teresa Vazão, Rodrigo Rodrigues 0001
NCA3
2006 Tolerating Byzantine Faulty Clients in a Quorum System
abstract
Byzantine quorum systems have been proposed that work properly even when up to f replicas fail arbitrarily. However, these systems are not so successful when confronted with Byzantine faulty clients. This paper presents novel protocols that provide atomic semantics despite Byzantine clients. Our protocols prevent Byzantine clients from interfering with good clients: bad clients cannot prevent good clients from completing reads and writes, and they cannot cause good clients to see inconsistencies. In addition we also prevent bad clients that have been removed from operation from leaving behind more than a bounded number of writes that could be done on their behalf by a colluder. Our protocols are designed to work in an asynchronous system like the Internet and they are highly efficient. We require 3f +1 replicas, and either two or three phases to do writes; reads normally complete in one phase and require no more than two phases, no matter what the bad clients are doing. We also present strong correctness conditions for systems with Byzantine clients that limit what can be done on behalf of bad clients once they leave the system. Furthermore we prove that our protocols are both safe (they meet those conditions) and live.
Barbara Liskov, Rodrigo Rodrigues 0001
ICDCS2
2006 HQ Replication: A Hybrid Quorum Protocol for Byzantine Fault Tolerance
James A. Cowling, Dan S. Myers, Barbara Liskov, Rodrigo Rodrigues 0001, Liuba Shrira
OSDI4
2005 Byzantine Clients Rendered Harmless
Barbara Liskov, Rodrigo Rodrigues 0001
DISC2
2004 Efficient Routing for Peer-to-Peer Overlays
Barbara Liskov, Rodrigo Rodrigues 0001
NSDI3
2004 Brief announcement: reconfigurable byzantine-fault-tolerant atomic memory
abstract
No abstract available.
Rodrigo Rodrigues 0001, Barbara Liskov
PODC1
2003 High Availability, Scalable Storage, Dynamic Peer Networks: Pick Two
Charles Blake 0001, Rodrigo Rodrigues 0001
HotOS2
2003 One Hop Lookups for Peer-to-Peer Overlays
Barbara Liskov, Rodrigo Rodrigues 0001
HotOS3
2003 BASE: Using abstraction to improve fault tolerance
abstract
Software errors are a major cause of outages and they are increasingly exploited in malicious attacks. Byzantine fault tolerance allows replicated systems to mask some software errors but it is expensive to deploy. This paper describes a replication technique, BASE, which uses abstraction to reduce the cost of Byzantine fault tolerance and to improve its ability to mask software errors. BASE reduces cost because it enables reuse of off-the-shelf service implementations. It improves availability because each replica can be repaired periodically using an abstract view of the state stored by correct replicas, and because each replica can run distinct or nondeterministic service implementations, which reduces the probability of common mode failures. We built an NFS service where each replica can run a different off-the-shelf file system implementation, and an object-oriented database where the replicas ran the same, nondeterministic implementation. These examples suggest that our technique can be used in practice---in both cases, the implementation required only a modest amount of new code, and our performance results indicate that the replicated services perform comparably to the implementations that they reuse.
Miguel Castro 0001, Rodrigo Rodrigues 0001, Barbara Liskov
ACM Trans. Comput. Syst.2
2001 Using Abstraction To Improve Fault Tolerance
abstract
Software errors are a major cause of outages and they are increasingly exploited in malicious attacks. Byzantine fault tolerance allows replicated systems to mask some software errors but it is expensive to deploy. The paper describes a replication technique, BFTA, which uses abstraction to reduce the cost of Byzantine fault tolerance and to improve its ability to mask software errors. BFTA reduces cost because it enables reuse of off-the-shelf service implementations. It improves availability because each replica can be repaired periodically using an abstract view of the state stored by correct replicas, and because each replica can run distinct or non-deterministic service implementations, which reduces the probability of common mode failures. We built an NFS service that allows each replica to run a different operating system. This example suggests that BFTA can be used in practice; the replicated file system required only a modest amount of new code, and preliminary performance results indicate that it performs comparably to the off-the-shelf implementations that it wraps.
Miguel Castro 0001, Rodrigo Rodrigues 0001, Barbara Liskov
HotOS2
2001 BASE: Using Abstraction to Improve Fault Tolerance
abstract
Software errors are a major cause of outages and they are increasingly exploited in malicious attacks. Byzantine fault tolerance allows replicated systems to mask some software errors but it is expensive to deploy. This paper describes a replication technique, BASE, which uses abstraction to reduce the cost of Byzantine fault tolerance and to improve its ability to mask software errors. BASE reduces cost because it enables reuse of off-the-shelf service implementations. It improves availability because each replica can be repaired periodically using an abstract view of the state stored by correct replicas, and because each replica can run distinct or non-deterministic service implementations, which reduces the probability of common mode failures. We built an NFS service where each replica can run a different off-the-shelf file system implementation, and an object-oriented database where the replicas ran the same, non-deterministic implementation. These examples suggest that our technique can be used in practice --- in both cases, the implementation required only a modest amount of new code, and our performance results indicate that the replicated services perform comparably to the implementations that they reuse.
Rodrigo Rodrigues 0001, Miguel Castro 0001, Barbara Liskov
SOSP1