Gregory V. Chockler

dblp:c/GregoryChockler · DBLP profile ↗
← Back
46ranked-venue papers
20as first author
12since 2021 · last 2026
0000-0001-6700-9235ORCID · verified

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

Systems, architecture and hardware · 28 · 14 first-author · 6 since 2021Software engineering, systems software and programming languages · 2 · 1 first-author · 1 since 2021Databases, data management, data science and information retrieval · 2Theory of computation · 2Security and privacy · 1
YearPublicationVenuePosition
2026 A Verified High-Performance Composable Object Library for Remote Direct Memory Access
abstract
Remote Direct Memory Access (RDMA) is a memory technology that allows remote devices to directly write to and read from each other’s memory, bypassing components such as the CPU and operating system. This enables low-latency high-throughput networking, as required for many modern data centres, HPC applications and AI/ML workloads. However, baseline RDMA comprises a highly permissive weak memory model that is difficult to use in practice and has only recently been formalised. In this paper, we introduce the Library of Composable Objects (LOCO), a formally verified library for building multi-node objects on RDMA, filling the gap between shared memory and distributed system programming. LOCO objects are well-encapsulated and take advantage of the strong locality and the weak consistency characteristics of RDMA. They have performance comparable to custom RDMA systems (e.g. distributed maps), but with a far simpler programming model amenable to formal proofs of correctness. To support verification, we develop a novel modular declarative verification framework, called Mowgli , that is flexible enough to model multinode objects and is independent of a memory consistency model. We instantiate Mowgli with the RDMA memory model, and use it to verify correctness of LOCO libraries.
Guillaume Ambal, George Hodgkins, Mark Madler, Gregory V. Chockler, Brijesh Dongol, Joseph Izraelevitz, Azalea Raad, Viktor Vafeiadis
Proc. ACM Program. Lang.4
2025 Tight Bounds on Channel Reliability via Generalized Quorum Systems
abstract
Communication channel failures are a major concern for the developers of modern fault-tolerant systems. However, while tight bounds for process failures are well-established, extending them to include channel failures has remained an open problem. We introduce generalized quorum systems -- a framework that characterizes the necessary and sufficient conditions for implementing atomic registers, atomic snapshots, lattice agreement and consensus under arbitrary patterns of process-channel failures. Generalized quorum systems relax the connectivity constraints of classical quorum systems: instead of requiring bidirectional reachability for every pair of write and read quorums, they only require some write quorum to be unidirectionally reachable from some read quorum. This weak connectivity makes implementing registers particularly challenging, because it precludes the traditional request/response pattern of quorum access, making classical solutions like ABD inapplicable. To address this, we introduce novel logical clocks that allow write and read quorums to reliably track state updates without relying on bidirectional connectivity.
Alejandro Naser-Pastoriza, Gregory V. Chockler, Alexey Gotsman, Fedor Ryabinin
PODC2
2025 TEE Is Not a Healer: Rollback-Resistant Reliable Storage
abstract
Recent advances in secure hardware technologies, such as Intel SGX or ARM TrustZone, offer an opportunity to substantially reduce the costs of Byzantine fault-tolerance by placing the program code and state within a secure enclave known as a Trusted Execution Environment (TEE). However, the protection offered by a TEE only applies during program execution. Once power is switched off, the non-volatile portion of the program state becomes vulnerable to rollback attacks wherein it is undetectably reverted to an older version. In this paper we consider the problem of implementing reliable read/write registers out of failure-prone replicas subject to state rollbacks. To this end, we introduce a new unified model that captures multiple failure types that can affect a TEE-based system and establish tight bounds on the fault-tolerance of register constructions in this model. We consider both the static case, where failure thresholds hold throughout the entire execution, and the dynamic case, where any number of replicas can roll back, provided these failures do not occur too often. Our dynamic register emulation algorithm, TEE-Rex, provides the first correct implementation of a distributed state recovery procedure that requires neither durable storage nor specialized hardware, such as trusted monotonic counters.
Sadegh Keshavarzi, Gregory V. Chockler, Alexey Gotsman
DISC2
2024 Mangosteen: Fast Transparent Durability for Linearizable Applications using NVM
Sergey Egorov, Gregory V. Chockler, Brijesh Dongol, Dan O'Keeffe, Sadegh Keshavarzi
USENIX ATC2
2024 Vertical Atomic Broadcast and Passive Replication
abstract
Atomic broadcast is a reliable communication abstraction ensuring that all processes deliver the same set of messages in a common global order. It is a fundamental building block for implementing fault-tolerant services using either active (aka state-machine) or passive (aka primary-backup) replication. We consider the problem of implementing reconfigurable atomic broadcast, which further allows users to dynamically alter the set of participating processes, e.g., in response to failures or changes in the load. We give a complete safety and liveness specification of this communication abstraction and propose a new protocol implementing it, called Vertical Atomic Broadcast, which uses an auxiliary service to facilitate reconfiguration. In contrast to prior proposals, our protocol significantly reduces system downtime when reconfiguring from a functional configuration by allowing it to continue processing messages while agreement on the next configuration is in progress. Furthermore, we show that this advantage can be maintained even when our protocol is modified to support a stronger variant of atomic broadcast required for passive replication.
Manuel Bravo, Gregory V. Chockler, Alexey Gotsman, Alejandro Naser-Pastoriza, Christian Roldán
DISC2
2024 What Cannot Be Implemented on Weak Memory?
abstract
We present a general methodology for establishing the impossibility of implementing certain concurrent objects on different (weak) memory models. The key idea behind our approach lies in characterizing memory models by their mergeability properties, identifying restrictions under which independent memory traces can be merged into a single valid memory trace. In turn, we show that the mergeability properties of the underlying memory model entail similar mergeability requirements on the specifications of objects that can be implemented on that memory model. We demonstrate the applicability of our approach to establish the impossibility of implementing standard distributed objects with different restrictions on memory traces on three memory models: strictly consistent memory, total store order, and release-acquire. These impossibility results allow us to identify tight and almost tight bounds for some objects, as well as new separation results between weak memory models, and between well-studied objects based on their implementability on weak memory models.
Armando Castañeda, Gregory V. Chockler, Brijesh Dongol, Ori Lahav 0001
DISC2
2024 Liveness and latency of Byzantine state-machine replication
abstract
Byzantine state-machine replication (SMR) ensures the consistency of replicated state in the presence of malicious replicas and lies at the heart of the modern blockchain technology. Byzantine SMR protocols often guarantee safety under all circumstances and liveness only under synchrony. However, guaranteeing liveness even under this assumption is nontrivial. So far we have lacked systematic ways of incorporating liveness mechanisms into Byzantine SMR protocols, which often led to subtle bugs. To close this gap, we introduce a modular framework to facilitate the design of provably live and efficient Byzantine SMR protocols. Our framework relies on a view abstraction generated by a special SMR synchronizer primitive to drive the agreement on command ordering. We present a simple formal specification of an SMR synchronizer and its bounded-space implementation under partial synchrony. We also apply our specification to prove liveness and analyze the latency of three Byzantine SMR protocols via a uniform methodology. In particular, one of these results yields what we believe is the first rigorous liveness proof for the algorithmic core of the seminal PBFT protocol.
Manuel Bravo, Gregory V. Chockler, Alexey Gotsman
Distributed Comput.2
2023 Fault-Tolerant Computing with Unreliable Channels
abstract
We study implementations of basic fault-tolerant primitives, such as consensus and registers, in message-passing systems subject to process crashes and a broad range of communication failures. Our results characterize the necessary and sufficient conditions for implementing these primitives as a function of the connectivity constraints and synchrony assumptions. Our main contribution is a new algorithm for partially synchronous consensus that is resilient to process crashes and channel failures and is optimal in its connectivity requirements. In contrast to prior work, our algorithm assumes the most general model of message loss where faulty channels are flaky, i.e., can lose messages without any guarantee of fairness. This failure model is particularly challenging for consensus algorithms, as it rules out standard solutions based on leader oracles and failure detectors. To circumvent this limitation, we construct our solution using a new variant of the recently proposed view synchronizer abstraction, which we adapt to the crash-prone setting with flaky channels.
Alejandro Naser-Pastoriza, Gregory V. Chockler, Alexey Gotsman
OPODIS2
2022 Acuerdo: Fast Atomic Broadcast over RDMA
abstract
Atomic broadcast protocols ensure that messages are delivered to a group of machines in some total order, even when some of these machines can fail. These protocols are key to making distributed services fault-tolerant, as their total order guarantee allows keeping multiple service replicas in sync. But, unfortunately, atomic broadcast protocols are also notoriously expensive.
Joseph Izraelevitz, Gaukas Wang, Rhett Hanscom, Kayli Silvers, Tamara Silbergleit Lehman, Gregory V. Chockler, Alexey Gotsman
ICPP6
2022 Liveness and Latency of Byzantine State-Machine Replication
Manuel Bravo, Gregory V. Chockler, Alexey Gotsman
DISC2
2022 Making Byzantine consensus live
Manuel Bravo, Gregory V. Chockler, Alexey Gotsman
Distributed Comput.2
2021 Multi-shot distributed transaction commit
abstract
Atomic Commit Problem (ACP) is a single-shot agreement problem similar to consensus, meant to model the properties of transaction commit protocols in fault-prone distributed systems. We argue that ACP is too restrictive to capture the complexities of modern transactional data stores, where commit protocols are integrated with concurrency control, and their executions for different transactions are interdependent. As an alternative, we introduce Transaction Certification Service (TCS), a new formal problem that captures safety guarantees of multi-shot transaction commit protocols with integrated concurrency control. TCS is parameterized by a certification function that can be instantiated to support common isolation levels, such as serializability and snapshot isolation. We then derive a provably correct crash-resilient protocol for implementing TCS through successive refinement. Our protocol achieves a better time complexity than mainstream approaches that layer two-phase commit on top of Paxos-style replication.
Gregory V. Chockler, Alexey Gotsman
Distributed Comput.1
2020 Making Byzantine Consensus Live
abstract
Partially synchronous Byzantine consensus protocols typically structure their execution into a sequence of views, each with a designated leader process. The key to guaranteeing liveness in these protocols is to ensure that all correct processes eventually overlap in a view with a correct leader for long enough to reach a decision. We propose a simple view synchronizer abstraction that encapsulates the corresponding functionality for Byzantine consensus protocols, thus simplifying their design. We present a formal specification of a view synchronizer and its implementation under partial synchrony, which runs in bounded space despite tolerating message loss during asynchronous periods. We show that our synchronizer specification is strong enough to guarantee liveness for single-shot versions of several well-known Byzantine consensus protocols, including HotStuff, Tendermint, PBFT and SBFT. We furthermore give precise latency bounds for these protocols when using our synchronizer. By factoring out the functionality of view synchronization we are able to specify and analyze the protocols in a uniform framework, which allows comparing them and highlights trade-offs.
Manuel Bravo, Gregory V. Chockler, Alexey Gotsman
DISC2
2019 White-Box Atomic Multicast
abstract
Atomic multicast is a communication primitive that delivers messages to multiple groups of processes according to some total order, with each group receiving the projection of the total order onto messages addressed to it. To be scalable, atomic multicast needs to be genuine, meaning that only the destination processes of a message should participate in ordering it. In this paper we propose a novel genuine atomic multicast protocol that in the absence of failures takes as low as 3 message delays to deliver a message when no other messages are multicast concurrently to its destination groups, and 5 message delays in the presence of concurrency. This improves the latencies of both the fault-tolerant version of classical Skeen's multicast protocol (6 or 12 message delays, depending on concurrency) and its recent improvement by Coelho et al. (4 or 8 message delays). To achieve such low latencies, we depart from the typical way of guaranteeing fault-tolerance by replicating each group with Paxos. Instead, we weave Paxos and Skeen's protocol together into a single coherent protocol, exploiting opportunities for white-box optimisations. We experimentally demonstrate that the superior theoretical characteristics of our protocol are reflected in practical performance pay-offs.
Alexey Gotsman, Anatole Lefort, Gregory V. Chockler
DSN3
2018 Session details: Session 2C: Security, Blockchains, and Replication
Gregory V. Chockler
PODC1
2018 Multi-Shot Distributed Transaction Commit
Gregory V. Chockler, Alexey Gotsman
DISC1
2017 Space Complexity of Fault-Tolerant Register Emulations
abstract
Driven by the rising popularity of cloud storage, the costs associated with implementing reliable storage services from a collection of fault-prone servers have recently become an actively studied question. The well-known ABD result shows that an f-tolerant register can be emulated using a collection of 2f+1 fault-prone servers each storing a single read-modify-write object, which is known to be optimal. In this paper we generalize this bound: we investigate the inherent space complexity of emulating reliable multi-writer registers as a function of the type of the base objects exposed by the underlying servers, the number of writers to the emulated register, the number of available servers, and the failure threshold.
Gregory V. Chockler, Alexander Spiegelman
PODC1
2017 Scalable communication middleware for permissioned distributed ledgers
abstract
Distributed Ledger Technology (DLT) is rapidly emerging as a new paradigm for automating complex business processes in secure and decentralised fashion. Currently, however, its wider adoption is hampered by scalability problems [3] rooted in an inherent tension between stringent consistency, security, and robustness requirements on one hand, and growing application demand coupled with high performance expectations on the other. For example, popular peer-to-peer DLTs based on proof-of-work consensus [4] can only improve the transaction throughput by degrading their security and consistency guarantees, which is unacceptable in the enterprise and mission-critical settings.
Artem Barger, Yacov Manevich, Benjamin Mandler, Vita Bortnikov, Gennady Laventman, Gregory V. Chockler
SYSTOR6
2016 Space Bounds for Reliable Storage: Fundamental Limits of Coding
abstract
We study the inherent space requirements of reliable storage algorithms in asynchronous distributed systems. A number of recent works have used codes in order to achieve a better storage cost than the well-known replication approach. However, a closer look reveals that they incur extra costs in certain scenarios. Specifically, if multiple clients access the storage concurrently, then existing asynchronous code-based algorithms may store a number of copies of the data that grows linearly with the number of concurrent clients. We prove here that this is inherent. Given three parameters, (1) the data size -- D bits, (2) the concurrency level -- c, and (3) the number of storage node failures that need to be tolerated -- f, we show a lower bound of Omega(min(f,c)D) bits on the space complexity of asynchronous distributed storage algorithms. Intuitively, this implies that the asymptotic storage cost is either as high as with replication, namely O(fD), or as high under concurrency as with the aforementioned code-based algorithms, i.e., O(cD).
Alexander Spiegelman, Yuval Cassuto, Gregory V. Chockler, Idit Keidar
PODC3
2015 Space Bounds for Reliable Storage: Fundamental Limits of Coding (Keynote)
abstract
We present here a synopsis of a keynote presentation given by Idit Keidar at OPODIS 2015, the International Conference on Principles of Distributed Systems, which took place in Rennes, France, on December 14-17 2015.
Alexander Spiegelman, Yuval Cassuto, Gregory V. Chockler, Idit Keidar
OPODIS3
2015 A Constructive Approach for Proving Data Structures' Linearizability
Kfir Lev-Ari, Gregory V. Chockler, Idit Keidar
DISC2
2014 Dynamic Performance Profiling of Cloud Caches
abstract
Large-scale in-memory object caches such as memcached are widely used to accelerate popular web sites and to reduce burden on backend databases. Yet current cache systems give cache operators limited information on what resources are required to optimally accommodate the present workload. This paper focuses on a key question for cache operators: how much total memory should be allocated to the in-memory cache tier to achieve desired performance?
Trausti Saemundsson, Hjörtur Björnsson, Gregory V. Chockler, Ymir Vigfusson
SoCC3
2014 On Correctness of Data Structures under Reads-Write Concurrency
Kfir Lev-Ari, Gregory V. Chockler, Idit Keidar
DISC2
2013 Dynamic performance profiling of cloud caches
abstract
In-memory object caches, such as memcached, are critical to the success of popular web sites, such as Facebook [3], by reducing database load and improving scalability [2]. The prominence of caches implies that configuring their ideal memory size has the potential for significant savings on computation resources and energy costs, but unfortunately cache configuration is poorly understood. The modern practice of manually tweaking live caching systems takes significant effort and may both increase the variance for client request latencies and impose high load on the database backend.
Hjörtur Björnsson, Gregory V. Chockler, Trausti Saemundsson, Ymir Vigfusson
SoCC2
2012 Brief announcement: reconfigurable state machine replication from non-reconfigurable building blocks
abstract
Reconfigurable state machine replication is an important enabler of elasticity for replicated cloud services, which must be able to dynamically adjust their size as a function of changing load and resource availability. We introduce a new generic framework to allow the reconfigurable state machine implementation to be derived from a collection of arbitrary non-reconfigurable state machines. Our reduction framework follows the black box approach, and does not make any assumptions with respect to its execution environment apart from reliable channels. It allows higher-level services to leverage speculative command execution to ensure uninterrupted progress during the reconfiguration periods as well as in situations where failures prevent the reconfiguration agreement from being reached in a timely fashion. We apply our framework to obtain a reconfigurable speculative state machine from the non-reconfigurable Paxos implementation, and analyze its performance on a realistic distributed testbed. Our results show that our framework incurs negligible overheads in the absence of reconfiguration, and allows steady throughput to be maintained throughout the reconfiguration periods.
Vita Bortnikov, Gregory V. Chockler, Dmitri Perelman, Alexey Roytman, Shlomit Shachor, Ilya Shnayderman
PODC2
2011 Special Issue on Cloud Computing
Gregory V. Chockler, Eliezer Dekel, Joseph F. JáJá, Jimmy Lin
J. Parallel Distributed Comput.1
2010 Dr. multicast: Rx for data center communication scalability
abstract
IP Multicast (IPMC) in data centers becomes disruptive when the technology is used by a large number of groups, a capability desired by event notification systems. We trace the problem to root causes, and introduce Dr. Multicast (MCMD), a system that eliminates the issue by mapping IPMC operations to a combination of point-to-point unicast and traditional IPMC transmissions guaranteed to be safe. MCMD optimizes the use of IPMC addresses within a data center by merging similar multicast groups in a principled fashion, while simultaneously respecting hardware limits expressed through administrator-controlled policies. The system is fully transparent, making it backward-compatible with commodity hardware and software found in modern data centers. Experimental evaluation shows that MCMD allows a large number of IPMC groups to be used without disruption, restoring a powerful group communication primitive to its traditional role.
Ymir Vigfusson, Hussam Abu-Libdeh, Mahesh Balakrishnan 0001, Kenneth P. Birman, Robert Burgess, Gregory V. Chockler, Haoyuan Li 0001, Yoav Tock
EuroSys6
2009 Special Issue of the Journal of Parallel and Distributed Computing: Cloud Computing
Gregory V. Chockler, Eliezer Dekel, Joseph F. JáJá, Jimmy Lin
J. Parallel Distributed Comput.1
2009 Reconfigurable distributed storage for dynamic networks
Gregory V. Chockler, Seth Gilbert, Vincent Gramoli, Peter M. Musial, Alexander A. Schwarzmann
J. Parallel Distributed Comput.1
2008 Virtual infrastructure for collision-prone wireless networks
abstract
Wireless ad hoc networks pose several significant challenges: devices are unreliable; deployments are unpredictable; and communication is erratic. One proposed solution is Virtual Infrastructure, an abstraction in which unpredictable and unreliable devices are used to emulate reliable and predictable infrastructure. In this paper, we present a new protocol for emulating virtual infrastructure in collision-prone wireless networks. At the heart of our emulation is a convergent history agreement protocol that tolerates lost messages and crash failures. It is designed specifically for ad hoc deployments, for example, the set of participants a priori unknown. The convergent history agreement protocol is quite efficient, as each agreement instance completes in a constant number of communication rounds, and the size of the messages is constant, independent of the length of the execution. Building on the convergent history agreement protocol, our virtual infrastructure emulation introduces only constant overhead per virtual round emulated. We believe that the techniques developed in this paper help to bring virtual infrastructure one step closer to a reality.
Gregory V. Chockler, Seth Gilbert, Nancy A. Lynch
PODC1
2008 Consensus and collision detectors in radio networks
Gregory V. Chockler, Murat Demirbas, Seth Gilbert, Nancy A. Lynch, Calvin C. Newport, Tina Nolte
Distributed Comput.1
2007 Constructing scalable overlays for pub-sub with many topics
abstract
We investigate the problem of designing a scalable overlay network to support decentralized topic-based pub/sub communication. We introduce a new optimization problem, called Minimum Topic-Connected Overlay (Min-TCO), that captures the tradeoff between the scalability of the overlay (in terms of the nodes' fanout) and the message forwarding overhead incurred by the communicating parties. Roughly, the Min-TCO problem is as follows: Given a collection of nodes and their subscriptions, connect the nodes using the minimum possible number of edges so that for each topic t, a message published on t could reach all the nodes interested in t by being forwarded by onlythe nodes interested in t.
Gregory V. Chockler, Roie Melamed, Yoav Tock, Roman Vitenberg
PODC1
2007 Amnesic Distributed Storage
Gregory V. Chockler, Rachid Guerraoui, Idit Keidar
DISC1
2007 Wait-free regular storage from Byzantine components
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi
Inf. Process. Lett.2
2006 Byzantine disk paxos: optimal resilience with byzantine shared memory
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi
Distributed Comput.2
2005 Reconfigurable Distributed Storage for Dynamic Networks
Gregory V. Chockler, Seth Gilbert, Vincent Gramoli, Peter M. Musial, Alexander A. Schwarzmann
OPODIS1
2005 Consensus and collision detectors in wireless Ad Hoc networks
abstract
We consider the fault-tolerant consensus problem in wireless ad hoc networks with crash-prone nodes. We develop consensus algorithms for single-hop environments where the nodes are located within broadcast range of each other. Our algorithms tolerate highly unpredictable wireless communication, in which messages may be lost due to collisions, electromagnetic interference, or other anomalies. Accordingly, each node may receive a different set of messages in the same round. In order to minimize collisions, we design adaptive algorithms that attempt to minimize the broadcast contention. To cope with unreliable communication, we augment the nodes with collision detectors and present a new classification of collision detectors in terms of accuracy and completeness, based on practical realities. We show exactly in which cases consensus can be solved, and thus determine the requirements for a useful collision detector.We validate the feasibility of our algorithms, and the underlying wireless model, with simulations based on a realistic 802.11 MAC layer implementation and a detailed radio propagation model. We analyze the performance of our algorithms under varying sizes and densities of deployment and varying MAC layer parameters. We use our single-hop consensus algorithms as the basis for solving consensus in a multi-hop network, demonstrating the resilience of our algorithms to a challenging and noisy environment.
Gregory V. Chockler, Murat Demirbas, Seth Gilbert, Calvin C. Newport, Tina Nolte
PODC1
2005 Proving Atomicity: An Assertional Approach
Gregory V. Chockler, Nancy A. Lynch, Sayan Mitra 0001, Joshua A. Tauber
DISC1
2005 Active Disk Paxos with infinitely many processes
Gregory V. Chockler, Dahlia Malkhi
Distributed Comput.1
2004 Byzantine disk paxos: optimal resilience with byzantine shared memory
abstract
We present Byzantine Disk Paxos, an asynchronous shared-memory consensus protocol that uses a collection of n > 3t disks, t of which may fail by becoming non-responsive or arbitrarily corrupted. We give two constructions of this protocol; that is, we construct two different building blocks, each of which can be used, along with a leader oracle, to solve consensus. One building block is a shared wait-free safe register. The second building block is a regular register that satisfies a weaker termination (liveness) condition than wait freedom: its write operations are wait-free, whereas its read operations are guaranteed to return only in executions with a finite number of writes. We call this termination condition finite writes (FW), and show that consensus is solvable with FW-terminating registers and a leader oracle. We construct each of these reliable registers from n > 3t base registers, t of which can be non-responsive or Byzantine. All the previous wait-free constructions in this model used at least 4t+1 fault-prone registers, and we are not familiar with any prior FW-terminating constructions in this model.
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi
PODC2
2003 On the composability of consistency conditions
Roy Friedman 0001, Roman Vitenberg, Gregory V. Chockler
Inf. Process. Lett.3
2002 Active disk paxos with infinitely many processes
abstract
We present an improvement to the Disk Paxos protocol by Gafni and Lamport which utilizes extended functionality and flexibility provided by Active Disks and supports unmediated concurrent data access by an unlimited number of processes. The solution facilitates coordination by an infinite number of clients using finite shared memory. It is based on a collection of read-modify-write objects with faults, that emulate a new, reliable shared memory abstraction called a ranked register. The required read-modify-write objects are readily available in Active Disks and in Object Storage Device controllers, making our solution suitable for state-of-the-art Storage Area Network (SAN) environments.
Gregory V. Chockler, Dahlia Malkhi
PODC1
2001 Backoff Protocols for Distributed Mutual Exclusion and Ordering
abstract
Presents a simple and efficient protocol for mutual exclusion in synchronous message-passing distributed systems subject to failures. Our protocol borrows design principles from prior work in backoff protocols for multiple access channels such as the Ethernet. Our protocol is adaptive in that the expected amortized system response time - informally, the average time a process waits before entering the critical section - is a function only of the number of clients currently contending and is independent of the maximum number of processes that might contend. In particular, in the contention-free case, a process can enter the critical section after only one round-trip message delay. We use this protocol to derive a protocol for ordering operations on a replicated object in an asynchronous distributed system subject to failures. This protocol is always safe, is probabilistically live during periods of stability and is suitable for deployment in practical systems.
Gregory V. Chockler, Dahlia Malkhi, Michael K. Reiter
ICDCS1
2000 Implementing a Caching Service for Distributed CORBA Objects
Gregory V. Chockler, Danny Dolev, Roy Friedman 0001, Roman Vitenberg
Middleware1
2000 Consistency Conditions for a CORBA Caching Service
Gregory V. Chockler, Roy Friedman 0001, Roman Vitenberg
DISC1
1998 An Adaptive Totally Ordered Multicast Protocol That Tolerates Partitions
abstract
In this work we present a novel protocol for total ordering of messages in asynchronous distributed environments prone to machine and communication link failures. Using the protocol as a building block, we constructed a Totally Ordered Group Communication (TOGC) system, i.e., a group communication service with a totally ordered multicast primitive. TOGC is a powerful infrastructure for building distributed fault-tolerant applications such as totally ordered broadcast, consistent object replication, distributed shared memory, Computer Supported Cooperative Work (CSCW) applications and distributed monitoring and display applications. An important contribution of the total ordering protocol described in this work is its ability to dynamically adjust the message delivery flow to changes in the transmission rates of the participating processes. The adaptation is accomplished by assigning delivery priorities (weights) to messages according to sender transmission rates. The priorities are det...
Gregory V. Chockler, N. Huleihel, Danny Dolev
PODC1