EDBT 2026 Demo / reviewers in the wild / expert
Miguel Castro 0001
dblp:c/MiguelCastro
· DBLP profile ↗
42ranked-venue papers
17as first author
3since 2021 · last 2024
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Software engineering, systems software and programming languages · 18 · 9 first-authorSystems, architecture and hardware · 13 · 6 first-author · 2 since 2021Computer networks · 7 · 3 first-author · 1 since 2021Security and privacy · 5 · 3 first-authorDatabases, data management, data science and information retrieval · 3
Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.
| Computer architecture, parallel and distributed computing, and storage systems
21 papers |
Distributed systems · 53% Interconnection networks and networks-on-chip · 19% Storage systems · 16% | |
| Network and information security
12 papers |
Systems and software security · 46% Blockchain and cryptocurrency security · 29% Cryptographic protocols and secure computation · 9% | |
| Databases, data mining, and information retrieval
3 papers |
Transaction processing and concurrency control · 45% Distributed and cloud data management · 27% Graph data management · 26% | |
| Software engineering, system software, and programming languages
9 papers |
Operating systems · 20% Software testing · 18% Program analysis · 18% | |
| Computer networks
7 papers |
Internet architecture and protocols · 78% Routing and switching · 14% Cellular and mobile networks · 6% |
Topics — the 30 heaviest of 85, each with the papers that count most for it
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Storage systems
key-value storage |
0.8 | 1 | 2024 | Honeycomb: Ordered Key-Value Store Acceleration on an FPGA-Based SmartNIC · IEEE Trans. Computers 2024 |
Internet architecture and protocols
network topology |
0.6 | 1 | 2022 | HammingMesh: A Network Topology for Large-Scale Deep Learning · SC 2022 |
Interconnection networks and networks-on-chip › network topology
network topology design |
0.6 | 1 | 2022 | HammingMesh: A Network Topology for Large-Scale Deep Learning · SC 2022 |
Distributed systems
fault tolerance |
0.5 | 6 | 2019 | Fast General Distributed Transactions with Opacity · SIGMOD Conference 2019 FARSITE: Federated, Available, and Reliable Storage for an Incompletely Trusted Environment · OSDI 2002 BASE: Using Abstraction to Improve Fault Tolerance · SOSP 2001 |
Distributed systems
replication |
0.5 | 5 | 2019 | Fast General Distributed Transactions with Opacity · SIGMOD Conference 2019 BASE: Using abstraction to improve fault tolerance · ACM Trans. Comput. Syst. 2003 Practical byzantine fault tolerance and proactive recovery · ACM Trans. Comput. Syst. 2002 |
Graph data management › graph database
distributed graph database |
0.4 | 1 | 2020 | A1: A Distributed In-Memory Graph Database · SIGMOD Conference 2020 |
Distributed and cloud data management
distributed query processing |
0.4 | 1 | 2020 | A1: A Distributed In-Memory Graph Database · SIGMOD Conference 2020 |
Distributed systems › fault tolerance
high availability |
0.4 | 2 | 2019 | Fast General Distributed Transactions with Opacity · SIGMOD Conference 2019 Safe and Efficient Sharing of Persistent Objects in Thor · SIGMOD Conference 1996 |
Transaction processing and concurrency control
distributed transaction processing |
0.4 | 1 | 2019 | Fast General Distributed Transactions with Opacity · SIGMOD Conference 2019 |
Transaction processing and concurrency control › serializability
strict serializability |
0.4 | 1 | 2019 | Fast General Distributed Transactions with Opacity · SIGMOD Conference 2019 |
Distributed systems › fault tolerance
transparent fault tolerance |
0.4 | 1 | 2019 | Fast General Distributed Transactions with Opacity · SIGMOD Conference 2019 |
Systems and software security
memory safety |
0.3 | 3 | 2010 | Dynamically checking ownership policies in concurrent c/c++ programs · POPL 2010 Baggy Bounds Checking: An Efficient and Backwards-Compatible Defense against Out-of-Bounds Errors · USENIX Security Symposium 2009 Preventing Memory Error Exploits with WIT · SP 2008 |
Hardware accelerators and domain-specific architectures › network accelerator
SmartNIC |
0.2 | 1 | 2024 | Honeycomb: Ordered Key-Value Store Acceleration on an FPGA-Based SmartNIC · IEEE Trans. Computers 2024 |
Distributed systems › distributed database
distributed transactions |
0.2 | 1 | 2015 | No compromises: distributed transactions with consistency, availability, and performance · SOSP 2015 |
Interconnection networks and networks-on-chip › remote direct memory access
RDMA-based replication |
0.2 | 1 | 2015 | No compromises: distributed transactions with consistency, availability, and performance · SOSP 2015 |
Distributed systems › replication › replication and fault tolerance
replication and recovery |
0.2 | 1 | 2015 | No compromises: distributed transactions with consistency, availability, and performance · SOSP 2015 |
Software testing
software reliability |
0.2 | 1 | 2014 | Docovery: toward generic automatic document recovery · ASE 2014 |
Memory systems
remote memory |
0.2 | 1 | 2014 | FaRM: Fast Remote Memory · NSDI 2014 |
Machine learning › Efficient and distributed learning
distributed training |
0.2 | 1 | 2022 | HammingMesh: A Network Topology for Large-Scale Deep Learning · SC 2022 |
Distributed systems › fault tolerance
byzantine fault tolerance |
0.2 | 5 | 2003 | BASE: Using abstraction to improve fault tolerance · ACM Trans. Comput. Syst. 2003 Practical byzantine fault tolerance and proactive recovery · ACM Trans. Comput. Syst. 2002 BASE: Using Abstraction to Improve Fault Tolerance · SOSP 2001 |
Malware analysis › malware defense
worm containment |
0.1 | 2 | 2008 | Vigilante: End-to-end containment of Internet worm epidemics · ACM Trans. Comput. Syst. 2008 Vigilante: end-to-end containment of internet worms · SOSP 2005 |
Interconnection networks and networks-on-chip
remote direct memory access |
0.1 | 1 | 2020 | A1: A Distributed In-Memory Graph Database · SIGMOD Conference 2020 |
Concurrent programming › concurrency bugs
data races |
0.1 | 1 | 2010 | Dynamically checking ownership policies in concurrent c/c++ programs · POPL 2010 |
Program verification
dynamic verification |
0.1 | 1 | 2010 | Dynamically checking ownership policies in concurrent c/c++ programs · POPL 2010 |
Systems and software security
vulnerability discovery |
0.1 | 2 | 2008 | Vigilante: End-to-end containment of Internet worm epidemics · ACM Trans. Comput. Syst. 2008 Better bug reporting with better privacy · ASPLOS 2008 |
Distributed systems › peer-to-peer systems
overlay networks |
0.1 | 2 | 2005 | Debunking Some Myths About Structured and Unstructured Overlays · NSDI 2005 An Evaluation of Scalable Application-Level Multicast Built Using Peer-To-Peer Overlays · INFOCOM 2003 |
Systems and software security › isolation
kernel extension isolation |
0.1 | 1 | 2009 | Fast byte-granularity software fault isolation · SOSP 2009 |
Systems and software security
operating system security |
0.1 | 1 | 2009 | Fast byte-granularity software fault isolation · SOSP 2009 |
Systems and software security › isolation
software fault isolation |
0.1 | 1 | 2009 | Fast byte-granularity software fault isolation · SOSP 2009 |
Operating systems › extensible operating systems › kernel extensibility
kernel extension isolation |
0.1 | 1 | 2009 | Fast byte-granularity software fault isolation · SOSP 2009 |
Methods — techniques the papers use, named apart from their topics
workload analysis · 1.7FaRM · 1.2RDMA · 1.1wait-free read · 0.8timestamp ordering · 0.8out-of-order execution · 0.8failover protocol · 0.8clock synchronization · 0.8caching · 0.8dynamic analysis · 0.3type safety enforcement · 0.2byte-granularity memory protection · 0.2bounds checking · 0.2automatic document recovery · 0.2information leakage measurement · 0.2execution path preservation · 0.2code instrumentation · 0.2symbolic execution · 0.1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2024 | Honeycomb: Ordered Key-Value Store Acceleration on an FPGA-Based SmartNICabstractIn-memory ordered key-value stores are an important building block in modern distributed applications. We present Honeycomb, a hybrid software-hardware system for accelerating read-dominated workloads on ordered key-value stores that provides linearizability for all operations including scans. Honeycomb stores a B-Tree in host memory. It executesput,updateanddeleteon a CPU. At the same time, it offloadsscanandgetonto an FPGA-based SmartNIC. This approach enables large stores and simplifies the FPGA implementation but raises the challenge of data access and synchronization across the slow PCIe bus. We describe how Honeycomb overcomes this challenge with careful data structure design, caching, request parallelism with out-of-order execution, wait-free read operations, and fast synchronization between the CPU and the FPGA. For read-heavy YCSB workloads, Honeycomb increases the throughput of a state-of-the-art ordered key-value store by at least$1.8\times$. For scan-heavy workloads inspired by cloud storage, Honeycomb increases the throughput by more than$2\times$. The cost-performance, which is more important for large-scale deployments, is improved by at least$1.5\times$on these workloads. Aleksandar Dragojevic, Shane T. Fleming, Antonios Katsarakis, Dario Korolija, Igor Zablotchi, Ho-Cheung Ng, Anuj Kalia, Miguel Castro 0001 |
IEEE Trans. Computers | 9 |
| 2022 | IA-CCF: Individual Accountability for Permissioned Ledgers
Alex Shamis, Peter R. Pietzuch, Burcu Canakci, Miguel Castro 0001, Cédric Fournet, Edward Ashton, Amaury Chamayou, Sylvan Clebsch, Antoine Delignat-Lavaud, Matthew Kerner, Julien Maffre, Olga Vrousgou, Christoph M. Wintersteiger, Manuel Costa, Mark Russinovich |
NSDI | 4 |
| 2022 | HammingMesh: A Network Topology for Large-Scale Deep LearningabstractNumerous microarchitectural optimizations unlocked tremendous processing power for deep neural networks that in turn fueled the AI revolution. With the exhaustion of such optimizations, the growth of modern AI is now gated by the performance of training systems, especially their data movement. Instead of focusing on single accelerators, we investigate data-movement characteristics of large-scale training at full system scale. Based on our workload analysis, we design HammingMesh, a novel network topology that provides high bandwidth at low cost with high job scheduling flexibility. Specifically, HammingMesh can support full bandwidth and isolation to deep learning training jobs with two dimensions of parallelism. Furthermore, it also supports high global bandwidth for generic traffic. Thus, HammingMesh will power future large-scale deep learning systems with extreme bandwidth requirements. Torsten Hoefler, Tommaso Bonato, Daniele De Sensi, Salvatore Di Girolamo, Shigang Li 0002, Marco Heddes, Jon Belk, Deepak Goel, Miguel Castro 0001, Steve Scott |
SC | 9 |
| 2020 | A1: A Distributed In-Memory Graph DatabaseabstractA1 is an in-memory distributed database used by the Bing search engine to support complex queries over structured data. The key enablers for A1 are availability of cheap DRAM and high speed RDMA (Remote Direct Memory Access) networking in commodity hardware. A1 uses FaRM [11,12] as its underlying storage layer and builds the graph abstraction and query engine on top. The combination of in-memory storage and RDMA access requires rethinking how data is allocated, organized and queried in a large distributed system. A single A1 cluster can store tens of billions of vertices and edges and support a throughput of 350+ million of vertex reads per second with end to end query latency in single digit milliseconds. In this paper we describe the A1 data model, RDMA optimized data structures and query execution. Chiranjeeb Buragohain, Knut Magne Risvik, Paul Brett, Miguel Castro 0001, Wonhee Cho 0004, Joshua Cowhig, Nikolas Gloy, Karthik Kalyanaraman, Richendra Khanna, John Pao, Matthew Renzelmann, Alex Shamis, Timothy Tan, Shuheng Zheng |
SIGMOD Conference | 4 |
| 2019 | Fast General Distributed Transactions with OpacityabstractTransactions can simplify distributed applications by hiding data distribution, concurrency, and failures from the application developer. Ideally the developer would see the abstraction of a single large machine that runs transactions sequentially and never fails. This requires the transactional subsystem to provide opacity (strict serializability for both committed and aborted transactions), as well as transparent fault tolerance with high availability. As even the best abstractions are unlikely to be used if they perform poorly, the system must also provide high performance. Existing distributed transactional designs either weaken this abstraction or are not designed for the best performance within a data center. This paper extends the design of FaRM --- which provides strict serializability only for committed transactions --- to provide opacity while maintaining FaRM's high throughput, low latency, and high availability within a modern data center. It uses timestamp ordering based on real time with clocks synchronized to within tens of microseconds across a cluster, and a failover protocol to ensure correctness across clock master failures. FaRM with opacity can commit 5.4 million neworder transactions per second when running the TPC-C transaction mix on 90 machines with 3-way replication. Alex Shamis, Matthew Renzelmann, Stanko Novakovic, Georgios Chatzopoulos, Aleksandar Dragojevic, Dushyanth Narayanan, Miguel Castro 0001 |
SIGMOD Conference | 7 |
| 2015 | No compromises: distributed transactions with consistency, availability, and performanceabstractTransactions with strong consistency and high availability simplify building and reasoning about distributed systems. However, previous implementations performed poorly. This forced system designers to avoid transactions completely, to weaken consistency guarantees, or to provide single-machine transactions that require programmers to partition their data. In this paper, we show that there is no need to compromise in modern data centers. We show that a main memory distributed computing platform called FaRM can provide distributed transactions with strict serializability, high performance, durability, and high availability. FaRM achieves a peak throughput of 140 million TATP transactions per second on 90 machines with a 4.9 TB database, and it recovers from a failure in less than 50 ms. Key to achieving these results was the design of new transaction, replication, and recovery protocols from first principles to leverage commodity networks with RDMA and a new, inexpensive approach to providing non-volatile DRAM. Aleksandar Dragojevic, Dushyanth Narayanan, Ed Nightingale, Matthew Renzelmann, Alex Shamis, Anirudh Badam, Miguel Castro 0001 |
SOSP | 7 |
| 2014 | Docovery: toward generic automatic document recoveryabstractApplication crashes and errors that occur while loading a document are one of the most visible defects of consumer software. While documents become corrupted in various ways---from storage media failures to incompatibility across applications to malicious modifications---the underlying reason they fail to load in a certain application is that their contents cause the application logic to exercise an uncommon execution path which the software was not designed to handle, or which was not properly tested. Tomasz Kuchta, Cristian Cadar, Miguel Castro 0001, Manuel Costa |
ASE | 3 |
| 2014 | FaRM: Fast Remote Memory
Aleksandar Dragojevic, Dushyanth Narayanan, Miguel Castro 0001, Orion Hodson |
NSDI | 3 |
| 2010 | Dynamically checking ownership policies in concurrent c/c++ programsabstractConcurrent programming errors arise when threads share data incorrectly. Programmers often avoid these errors by using synchronization to enforce a simple ownership policy: data is either owned exclusively by a thread that can read or write the data, or it is read owned by a set of threads that can read but not write the data. Unfortunately, incorrect synchronization often fails to enforce these policies and memory errors in languages like C and C++ can violate these policies even when synchronization is correct. Jean-Phillipe Martin, Michael Hicks 0001, Manuel Costa, Periklis Akritidis, Miguel Castro 0001 |
POPL | 5 |
| 2009 | Fast byte-granularity software fault isolationabstractBugs in kernel extensions remain one of the main causes of poor operating system reliability despite proposed techniques that isolate extensions in separate protection domains to contain faults. We believe that previous fault isolation techniques are not widely used because they cannot isolate existing kernel extensions with low overhead on standard hardware. This is a hard problem because these extensions communicate with the kernel using a complex interface and they communicate frequently. We present BGI (Byte-Granularity Isolation), a new software fault isolation technique that addresses this problem. BGI uses efficient byte-granularity memory protection to isolate kernel extensions in separate protection domains that share the same address space. BGI ensures type safety for kernel objects and it can detect common types of errors inside domains. Our results show that BGI is practical: it can isolate Windows drivers without requiring changes to the source code and it introduces a CPU overhead between 0 and 16%. BGI can also find bugs during driver testing. We found 28 new bugs in widely used Windows drivers. Miguel Castro 0001, Manuel Costa, Jean-Philippe Martin, Marcus Peinado, Periklis Akritidis, Austin Donnelly, Paul Barham 0001, Richard Black |
SOSP | 1 |
| 2009 | Baggy Bounds Checking: An Efficient and Backwards-Compatible Defense against Out-of-Bounds Errors
Periklis Akritidis, Manuel Costa, Miguel Castro 0001, Steven Hand 0001 |
USENIX Security Symposium | 3 |
| 2008 | Better bug reporting with better privacyabstractSoftware vendors collect bug reports from customers to improve the quality of their software. These reports should include the inputs that make the software fail, to enable vendors to reproduce the bug. However, vendors rarely include these inputs in reports because they may contain private user data. We describe a solution to this problem that provides software vendors with new input values that satisfy the conditions required to make the software follow the same execution path until it fails, but are otherwise unrelated with the original inputs. These new inputs allow vendors to reproduce the bug while revealing less private information than existing approaches. Additionally, we provide a mechanism to measure the amount of information revealed in an error report. This mechanism allows users to perform informed decisions on whether or not to submit reports. We implemented a prototype of our solution and evaluated it with real errors in real programs. The results show that we can produce error reports that allow software vendors to reproduce bugs while revealing almost no private information. Miguel Castro 0001, Manuel Costa, Jean-Philippe Martin |
ASPLOS | 1 |
| 2008 | Preventing Memory Error Exploits with WITabstractAttacks often exploit memory errors to gain control over the execution of vulnerable programs. These attacks remain a serious problem despite previous research on techniques to prevent them. We present write integrity testing (WIT), a new technique that provides practical protection from these attacks. WIT uses points-to analysis at compile time to compute the control-flow graph and the set of objects that can be written by each instruction in the program. Then it generates code instrumented to prevent instructions from modifying objects that are not in the set computed by the static analysis, and to ensure that indirect control transfers are allowed by the control-flow graph. To improve coverage where the analysis is not precise enough, WIT inserts small guards between the original program objects. We describe an efficient implementation with optimizations to reduce space and time overhead. This implementation can be used in practice because it compiles C and C++ programs without modifications, it has high coverage with no false positives, and it has low overhead. WIT's average runtime overhead is only 7% across a set of CPU intensive benchmarks and it is negligible when IO is the bottleneck. Periklis Akritidis, Cristian Cadar, Costin Raiciu, Manuel Costa, Miguel Castro 0001 |
SP | 5 |
| 2008 | Vigilante: End-to-end containment of Internet worm epidemicsabstractWorm containment must be automatic because worms can spread too fast for humans to respond. Recent work proposed network-level techniques to automate worm containment; these techniques have limitations because there is no information about the vulnerabilities exploited by worms at the network level. We propose Vigilante, a new end-to-end architecture to contain worms automatically that addresses these limitations. In Vigilante, hosts detect worms by instrumenting vulnerable programs to analyze infection attempts. We introduce dynamic data-flow analysis : a broad-coverage host-based algorithm that can detect unknown worms by tracking the flow of data from network messages and disallowing unsafe uses of this data. We also show how to integrate other host-based detection mechanisms into the Vigilante architecture. Upon detection, hosts generate self-certifying alerts (SCAs), a new type of security alert that can be inexpensively verified by any vulnerable host. Using SCAs, hosts can cooperate to contain an outbreak, without having to trust each other. Vigilante broadcasts SCAs over an overlay network that propagates alerts rapidly and resiliently. Hosts receiving an SCA protect themselves by generating filters with vulnerability condition slicing : an algorithm that performs dynamic analysis of the vulnerable program to identify control-flow conditions that lead to successful attacks. These filters block the worm attack and all its polymorphic mutations that follow the execution path identified by the SCA. Our results show that Vigilante can contain fast-spreading worms that exploit unknown vulnerabilities, and that Vigilante's filters introduce a negligible performance overhead. Vigilante does not require any changes to hardware, compilers, operating systems, or the source code of vulnerable programs; therefore, it can be used to protect current software binaries. Manuel Costa, Jon Crowcroft, Miguel Castro 0001, Antony I. T. Rowstron, Lidong Zhou, Paul Barham 0001 |
ACM Trans. Comput. Syst. | 3 |
| 2007 | Third Workshop on Hot Topics in System Dependability HotDep'07abstractThe goals of HotDep are to bring forth cuttingedge research ideas spanning the domains of fault tolerance/ reliability and systems, and to build linkages between the two communities (e.g., between people who attend traditional "dependability" conferences such as DSN and ISSRE, and those who attend "systems" conferences such as OSDI, SOSP, and EuroSys). Miguel Castro 0001, John Wilkes |
DSN | 1 |
| 2007 | Bouncer: securing software by blocking bad inputabstractAttackers exploit software vulnerabilities to control or crash programs. Bouncer uses existing software instrumentation techniques to detect attacks and it generates filters automatically to block exploits of the target vulnerabilities. The filters are deployed automatically by instrumenting system calls to drop exploit messages. These filters introduce low overhead and they allow programs to keep running correctly under attack. Previous work computes filters using symbolic execution along the path taken by a sample exploit, but attackers can bypass these filters by generating exploits that follow a different execution path. Bouncer introduces three techniques to generalize filters so that they are harder to bypass: a new form of program slicing that uses a combination of static and dynamic analysis to remove unnecessary conditions from the filter; symbolic summaries for common library functions that characterize their behavior succinctly as a set of conditions on the input; and generation of alternative exploits guided by symbolic execution. Bouncer filters have low overhead, they do not have false positives by design, and our results show that Bouncer can generate filters that block all exploits of some real-world vulnerabilities. Manuel Costa, Miguel Castro 0001, Lidong Zhou, Marcus Peinado |
SOSP | 2 |
| 2006 | Network coding with traffic engineeringabstractIn network coding, a router in the network mixes information from different flows. In the seminal work by Ahlswede et al [1], network coding is established as a technique to potentially increase the network capacity. Miguel Castro 0001, Jon Crowcroft, Greg O'Shea, Antony I. T. Rowstron |
CoNEXT | 2 |
| 2006 | POS: A Practical Order Statistics Service forWireless Sensor NetworksabstractWe present the design and implementation of POS, an in-network service that computes accurate order statistics energy-efficiently. POS returns a stream of periodic samples from any order statistic. It initially computes the value of the order statistic and then periodically runs a validation protocol to determine whether the value is still valid. If not, it uses an optimized binary search to determine the new value and then resumes periodic validation. POS uses in-network aggregation and transmission suppression to reduce communication complexity. Results from both experiments on a mote testbed and simulations show that POS can compute order statistics accurately while consuming less energy than the best techniques to compute averages in common cases. Landon P. Cox, Miguel Castro 0001, Antony I. T. Rowstron |
ICDCS | 2 |
| 2006 | Securing Software by Enforcing Data-flow Integrity
Miguel Castro 0001, Manuel Costa, Tim Harris 0001 |
OSDI | 1 |
| 2006 | Virtual ring routing: network routing inspired by DHTsabstractThis paper presents Virtual Ring Routing (VRR), a new network routing protocol that occupies a unique point in the design space. VRR is inspired by overlay routing algorithms in Distributed Hash Tables (DHTs) but it does not rely on an underlying network routing protocol. It is implemented directly on top of the link layer. VRR provides both raditional point-to-point network routing and DHT routing to the node responsible for a hash table key.VRR can be used with any link layer technology but this paper describes a design and several implementations of VRR that are tuned for wireless networks. We evaluate the performance of VRR using simulations and measurements from a sensor network and an 802.11a testbed. The experimental results show that VRR provides robust performance across a wide range of environments and workloads. It performs comparably to, or better than, the best wireless routing protocol in each experiment. VRR performs well because of its unique features: it does not require network flooding or trans-lation between fixed identifiers and location-dependent addresses. Matthew Caesar 0001, Miguel Castro 0001, Ed Nightingale, Greg O'Shea, Antony I. T. Rowstron |
SIGCOMM | 2 |
| 2005 | Debunking Some Myths About Structured and Unstructured Overlays
Miguel Castro 0001, Manuel Costa, Antony I. T. Rowstron |
NSDI | 1 |
| 2005 | Vigilante: end-to-end containment of internet wormsabstractWorm containment must be automatic because worms can spread too fast for humans to respond. Recent work has proposed network-level techniques to automate worm containment; these techniques have limitations because there is no information about the vulnerabilities exploited by worms at the network level. We propose Vigilante, a new end-to-end approach to contain worms automatically that addresses these limitations. Vigilante relies on collaborative worm detection at end hosts, but does not require hosts to trust each other. Hosts run instrumented software to detect worms and broadcast self-certifying alerts (SCAs) upon worm detection. SCAs are proofs of vulnerability that can be inexpensively verified by any vulnerable host. When hosts receive an SCA, they generate filters that block infection by analysing the SCA-guided execution of the vulnerable software. We show that Vigilante can automatically contain fast-spreading worms that exploit unknown vulnerabilities without blocking innocuous traffic. Manuel Costa, Jon Crowcroft, Miguel Castro 0001, Antony I. T. Rowstron, Lidong Zhou, Paul Barham 0001 |
SOSP | 3 |
| 2004 | Performance and Dependability of Structured Peer-to-Peer OverlaysabstractStructured peer-to-peer (P2P) overlay networks provide a useful substrate for building distributed applications. They map object keys to overlay nodes and offer a primitive to send a message to the node responsible for a key. They can implement, for example, distributed hash tables and multicast trees. However, there are concerns about the performance and dependability of these overlays in realistic environments. Several studies have shown that current P2P environments have high churn rates: nodes join and leave the overlay continuously. This paper presents techniques that continuously detect faults and repair the overlay to achieve high dependability and good performance in realistic environments. The techniques are evaluated using large-scale network simulation experiments with fault injection guided by real traces of node arrivals and departures. The results show that previous concerns are unfounded; our techniques can achieve dependable routing in realistic environments with an average delay stretch below two and a maintenance overhead of less than half a message per second per node. Miguel Castro 0001, Manuel Costa, Antony I. T. Rowstron |
DSN | 1 |
| 2004 | PIC: Practical Internet Coordinates for Distance EstimationabstractWe introduce PIC, a practical coordinate-based mechanism to estimate Internet network distance (i.e., round-trip delay or network hops). Network distance estimation is important in many applications; for example, network-aware overlay construction and server selection. There are several proposals for distance estimation in the Internet but they all suffer from problems that limit their benefit. Most rely on a small set of infrastructure nodes that are a single point of failure and limit scalability. Others use sets of peers to compute coordinates but these coordinates can be arbitrarily wrong if one of these peers is malicious. While it may be reasonable to secure a small set of infrastructure nodes, it is unreasonable to secure all peers. PIC addresses these problems: it does not rely on infrastructure nodes and it can compute accurate coordinates even when some peers are malicious. We present PIC's design, experimental evaluation, and an application to network-aware overlay construction and maintenance. Manuel Costa, Miguel Castro 0001, Antony I. T. Rowstron, Peter B. Key |
ICDCS | 2 |
| 2003 | An Evaluation of Scalable Application-Level Multicast Built Using Peer-To-Peer OverlaysabstractStructured peer-to-peer overlay networks such as CAN, Chord, Pastry, and Tapestry can be used to implement Internet-scale application-level multicast. There are two general approaches to accomplishing this: tree building and flooding. This paper evaluates these two approaches using two different types of structured overlay: 1) overlays which use a form of generalized hypercube routing, e.g., Chord, Pastry and Tapestry, and 2) overlays which use a numerical distance metric to route through a Cartesian hyperspace, e.g., CAN. Pastry and CAN are chosen as the representatives of each type of overlay. To the best of our knowledge, this paper reports the first head-to-head comparison of CAN-style versus Pastry-style overlay networks, using multicast communication workloads running on an identical simulation infrastructure. The two approaches to multicast are independent of overlay network choice, and we provide a comparison of flooding versus tree-based multicast on both overlays. Results show that the tree-based approach consistently outperforms the flooding approach. Finally, for tree-based multicast, we show that Pastry provides better performance than CAN. Miguel Castro 0001, Michael B. Jones, Anne-Marie Kermarrec, Antony I. T. Rowstron, Marvin Theimer, Helen J. Wang, Alec Wolman |
INFOCOM | 1 |
| 2003 | SplitStream: high-bandwidth multicast in cooperative environmentsabstractIn tree-based multicast systems, a relatively small number of interior nodes carry the load of forwarding multicast messages. This works well when the interior nodes are highly-available, dedicated infrastructure routers but it poses a problem for application-level multicast in peer-to-peer systems. SplitStream addresses this problem by striping the content across a forest of interior-node-disjoint multicast trees that distributes the forwarding load among all participating peers. For example, it is possible to construct efficient SplitStream forests in which each peer contributes only as much forwarding bandwidth as it receives. Furthermore, with appropriate content encodings, SplitStream is highly robust to failures because a node failure causes the loss of a single stripe on average. We present the design and implementation of SplitStream and show experimental results obtained on an Internet testbed and via large-scale network simulation. The results show that SplitStream distributes the forwarding load among all peers and can accommodate peers with different bandwidth capacities while imposing low overhead for forest construction and maintenance. Miguel Castro 0001, Peter Druschel, Anne-Marie Kermarrec, Animesh Nandi, Antony I. T. Rowstron, Atul Singh |
SOSP | 1 |
| 2003 | BASE: Using abstraction to improve fault toleranceabstractSoftware 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. | 1 |
| 2002 | FARSITE: Federated, Available, and Reliable Storage for an Incompletely Trusted Environment
Atul Adya, William J. Bolosky, Miguel Castro 0001, Gerald Cermak, Ronnie Chaiken, John R. Douceur, Jon Howell, Jacob R. Lorch, Marvin Theimer, Roger Wattenhofer |
OSDI | 3 |
| 2002 | Secure Routing for Structured Peer-to-Peer Overlay Networks
Miguel Castro 0001, Peter Druschel, Ayalvadi J. Ganesh, Antony I. T. Rowstron, Dan S. Wallach |
OSDI | 1 |
| 2002 | Scribe: a large-scale and decentralized application-level multicast infrastructureabstractThis paper presents Scribe, a scalable application-level multicast infrastructure. Scribe supports large numbers of groups, with a potentially large number of members per group. Scribe is built on top of Pastry, a generic peer-to-peer object location and routing substrate overlayed on the Internet, and leverages Pastry's reliability, self-organization, and locality properties. Pastry is used to create and manage groups and to build efficient multicast trees for the dissemination of messages to each group. Scribe provides best-effort reliability guarantees, and we outline how an application can extend Scribe to provide stronger reliability. Simulation results, based on a realistic network topology model, show that Scribe scales across a wide range of groups and group sizes. Also, it balances the load on the nodes while achieving acceptable delay and link stress when compared with Internet protocol multicast. Miguel Castro 0001, Peter Druschel, Anne-Marie Kermarrec, Antony I. T. Rowstron |
IEEE J. Sel. Areas Commun. | 1 |
| 2002 | Practical byzantine fault tolerance and proactive recoveryabstractOur growing reliance on online services accessible on the Internet demands highly available systems that provide correct service without interruptions. Software bugs, operator mistakes, and malicious attacks are a major cause of service interruptions and they can cause arbitrary behavior, that is, Byzantine faults. This article describes a new replication algorithm, BFT, that can be used to build highly available systems that tolerate Byzantine faults. BFT can be used in practice to implement real services: it performs well, it is safe in asynchronous environments such as the Internet, it incorporates mechanisms to defend against Byzantine-faulty clients, and it recovers replicas proactively. The recovery mechanism allows the algorithm to tolerate any number of faults over the lifetime of the system provided fewer than 1/3 of the replicas become faulty within a small window of vulnerability. BFT has been implemented as a generic program library with a simple interface. We used the library to implement the first Byzantine-fault-tolerant NFS file system, BFS. The BFT library and BFS perform well because the library incorporates several important optimizations, the most important of which is the use of symmetric cryptography to authenticate messages. The performance results show that BFS performs 2% faster to 24% slower than production implementations of the NFS protocol that are not replicated. This supports our claim that the BFT library can be used to build practical systems that tolerate Byzantine faults. Miguel Castro 0001, Barbara Liskov |
ACM Trans. Comput. Syst. | 1 |
| 2001 | Byzantine Fault Tolerance Can Be FastabstractByzantine fault tolerance is important because it can be used to implement highly-available systems that tolerate arbitrary behavior from faulty components. We present a detailed performance evaluation of BFT, a state-machine replication algorithm that tolerates Byzantine faults in asynchronous systems. Our results contradict the common belief that Byzantine fault tolerance is too slow to be used in practice, BFT performs well so that it can be used to implement real systems. We implemented a replicated NFS file system using BFT that performs 2% faster to 24% slower than production implementations of the NFS protocol that are not fault-tolerant. Miguel Castro 0001, Barbara Liskov |
DSN | 1 |
| 2001 | Using Abstraction To Improve Fault ToleranceabstractSoftware 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 |
HotOS | 1 |
| 2001 | BASE: Using Abstraction to Improve Fault ToleranceabstractSoftware 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 |
SOSP | 2 |
| 2000 | Proactive Recovery in a Byzantine-Fault-Tolerant System
Miguel Castro 0001, Barbara Liskov |
OSDI | 1 |
| 1999 | Providing Persistent Objects in Distributed Systems
Barbara Liskov, Miguel Castro 0001, Liuba Shrira, Atul Adya |
ECOOP | 2 |
| 1999 | Practical Byzantine Fault Tolerance
Miguel Castro 0001, Barbara Liskov |
OSDI | 1 |
| 1997 | Fragment Reconstruction: Providing Global Cache Coherence in a Transactional Storage SystemabstractCooperative caching is a promising technique to avoid the increasingly formidable disk bottleneck problem in distributed storage systems; it reduces the number of disk accesses by servicing client cache misses from the caches of other clients. However, existing cooperative caching techniques do not provide adequate support for fine grained sharing. We describe a new storage system architecture, split caching, and a new cache coherence protocol, fragment reconstruction, that combine cooperative caching with efficient support for fine grained sharing and transactions. We also present the results of performance studies that show that our scheme introduces little overhead over the basic cooperative caching mechanism and provides better performance when there is fine grained sharing. Atul Adya, Miguel Castro 0001, Barbara Liskov, Umesh Maheshwari, Liuba Shrira |
ICDCS | 2 |
| 1997 | HAC: Hybrid Adaptive Caching for Distributed Storage SystemsabstractThis paper presents HAC, a novel technique for managing the client cache in a distributed, persistent object storage system.I-k2 is a hybrid between page and object caching that combines the virtues of both while avoiding their disadvantages.,It achieves the low miss penalties of a page-caching system, but is able to perform well even when locality is poor, since it can discard pages while retaining their hot objects.It realizes the potentially lower miss rates of object-caching systems, yet avoids their problems of fragmentation and high overheads.Furthermore, HAC is adaptive: when locality is good it behaves like a page-caching system, while if locality is poor it behaves like an object-caching system.It is able to adjust the amount of cache space devoted to pages dynamically so that space in the cache can be used in the way that Fcst matches tbe needs of the application.The paper also presents results of experiments that indicate that HAC outperforms other object storage systems across a wide range of cache sizes and workloads; it performs substantially better on the expected workloads, which have low to moderate locality.Thus we show that our hybrid, adaptive approach is the cache management technique of choice for distributed, persistent object systems.This research was supported in Miguel Castro 0001, Atul Adya, Barbara Liskov, Andrew C. Myers |
SOSP | 1 |
| 1996 | Lightweight Logging for Lazy Release Consistent Distributed Shared MemoryabstractNo abstract available. Manuel Costa, Paulo Guedes, Manuel Sequeira, Nuno Neves 0001, Miguel Castro 0001 |
OSDI | 5 |
| 1996 | Safe and Efficient Sharing of Persistent Objects in ThorabstractThor is an object-oriented database system designed for use in a heterogeneous distributed environment. It provides highly-reliable and highly-available persistent storage for objects, and supports safe sharing of these objects by applications written in different programming languages.Safe heterogeneous sharing of long-lived objects requires encapsulation: the system must guarantee that applications interact with objects only by invoking methods. Although safety concerns are important, most object-oriented databases forgo safety to avoid paying the associated performance costs.This paper gives an overview of Thor's design and implementation. We focus on two areas that set Thor apart from other object-oriented databases. First, we discuss safe sharing and techniques for ensuring it; we also discuss ways of improving application performance without sacrificing safety. Second, we describe our approach to cache management at client machines, including a novel adaptive prefetching strategy.The paper presents performance results for Thor, on several OO7 benchmark traversals. The results show that adaptive prefetching is very effective, improving both the elapsed time of traversals and the amount of space used in the client cache. The results also show that the cost of safe sharing can be negligible; thus it is possible to have both safety and high performance. Barbara Liskov, Atul Adya, Miguel Castro 0001, Mark Day, Sanjay Ghemawat, Robert Gruber, Umesh Maheshwari, Andrew C. Myers, Liuba Shrira |
SIGMOD Conference | 3 |
| 1994 | A Checkpoint Protocol for an Entry Consistent Shared Memory SystemabstractWorkstation clusters are becoming an interesting alter-native to dedicated multiprocessors. In this environment, the probability of a failure, during an application’s exeeution, increases with the execution time and the number of work-stations used. If no provision is made for handling failures, it is unlikely that long running applications will terminate successfully. One solution to this problem is process check-pointing. This paper presents a checkpoint protocol for a multi-threaded distributed shared memory system based on the en-try consistency memory model. The protocol allows trans-parent recovery from single node failures and, in some cases, from multiple node failures. A simple mechanism is used to determine if the system can be brought to a consistent state in the event of multiple machine crashes. The protocol keeps a distributed log of shared data ac-cessesin the volatile memory of the processes, taking advan-tage of the independent failure characteristics of workstation clusters. Periodically, or whenever the log reaches a high-water mark, each process checkpoints its state, independently from the others. The protocol needs no extra messagesdur-ing the failure-fke period, since atl checkpoint control in-formation is piggybacked on the memory coherence protocol messages. 1 Nuno Neves 0001, Miguel Castro 0001, Paulo Guedes |
PODC | 2 |