Danny Dolev

dblp:d/DannyDolev · DBLP profile ↗
← Back
172ranked-venue papers
71as first author
4since 2021 · last 2023
0000-0001-8853-0644ORCID · verified

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

Theory of computation · 52 · 35 first-authorSystems, architecture and hardware · 50 · 12 first-author · 2 since 2021Security and privacy · 18 · 6 first-authorApplied, interdisciplinary, general and emerging computing · 16 · 9 first-authorComputer networks · 14 · 4 first-authorSoftware engineering, systems software and programming languages · 6 · 1 first-authorDatabases, data management, data science and information retrieval · 3 · 1 first-authorArtificial intelligence and machine learning · 1Graphics, computer vision, multimedia, augmented reality and games · 1

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.

Theoretical computer science
51 papers
Distributed computing theory · 77% Computational complexity · 13% Algorithmic game theory and mechanism design · 10%
Computer architecture, parallel and distributed computing, and storage systems
52 papers
Distributed systems · 87% Integrated circuit design · 6% Reconfigurable computing and FPGAs · 2%
Network and information security
18 papers
Network security · 47% Cryptographic protocols and secure computation · 24% Blockchain and cryptocurrency security · 19%
Computer networks
7 papers
Internet architecture and protocols · 65% Routing and switching · 28% Network performance modeling · 4%

Topics — the 30 heaviest of 159, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Distributed computing theory › fault tolerance › byzantine fault tolerance
byzantine agreement
1.2142020
Revisiting Asynchronous Fault Tolerant Computation with Optimal Resilience · PODC 2020
Communication Complexity of Byzantine Agreement, Revisited · PODC 2019
Byzantine Agreement with Optimal Early Stopping, Optimal Resilience and Polynomial Complexity · STOC 2015
Distributed computing theory
fault tolerance
0.792020
Revisiting Asynchronous Fault Tolerant Computation with Optimal Resilience · PODC 2020
Byzantine Agreement with Optimal Early Stopping, Optimal Resilience and Polynomial Complexity · STOC 2015
Failure Detectors in Omission Failure Environments · PODC 1997
Distributed systems
fault tolerance
0.6252012
An Optimal Self-Stabilizing Firing Squad · SIAM J. Comput. 2012
Steward: Scaling Byzantine Fault-Tolerant Replication to Wide Area Networks · IEEE Trans. Dependable Secur. Comput. 2010
Nysiad: Practical Protocol Transformation to Tolerate Byzantine Failures · NSDI 2008
Computational complexity
communication complexity
0.652019
Communication Complexity of Byzantine Agreement, Revisited · PODC 2019
Early-deciding consensus is expensive · PODC 2013
Determinism vs. Nondeterminism in Multiparty Communication Complexity · SIAM J. Comput. 1992
Distributed systems
consensus
0.6112020
Early-deciding consensus is expensive · PODC 2013
An Optimal Self-Stabilizing Firing Squad · SIAM J. Comput. 2012
Revisiting Asynchronous Fault Tolerant Computation with Optimal Resilience · PODC 2020
Distributed systems › fault tolerance
byzantine fault tolerance
0.572014
Fault-tolerant algorithms for tick-generation in asynchronous logic: Robust pulse generation · J. ACM 2014
Steward: Scaling Byzantine Fault-Tolerant Replication to Wide Area Networks · IEEE Trans. Dependable Secur. Comput. 2010
OCD: obsessive consensus disorder (or repetitive consensus) · PODC 2008
Distributed computing theory
self-stabilization
0.422017
Stateless Computation · PODC 2017
Fast self-stabilizing byzantine tolerant digital clock synchronization · PODC 2008
Distributed computing theory
impossibility results
0.312017
Stateless Computation · PODC 2017
Distributed systems › consensus
byzantine agreement
0.382013
Early-deciding consensus is expensive · PODC 2013
Self-stabilizing byzantine agreement · PODC 2006
Atomic Broadcast: From Simple Message Diffusion to Byzantine Agreement · Inf. Comput. 1995
Distributed computing theory
consensus
0.272015
Byzantine Agreement with Optimal Early Stopping, Optimal Resilience and Polynomial Complexity · STOC 2015
The Distributed Firing Squad Problem · SIAM J. Comput. 1989
On the minimal synchronism needed for distributed consensus · J. ACM 1987
Distributed computing theory › fault tolerance
byzantine fault tolerance
0.232015
Byzantine Agreement with Optimal Early Stopping, Optimal Resilience and Polynomial Complexity · STOC 2015
Shifting Gears: Changing Algorithms on the Fly to Expedite Byzantine Agreement · Inf. Comput. 1992
On the Possibility and Impossibility of Achieving Clock Synchronization · STOC 1984
Distributed systems › clock synchronization
fault-tolerant clock synchronization
0.222014
Fault-tolerant algorithms for tick-generation in asynchronous logic: Robust pulse generation · J. ACM 2014
Dynamic Fault-Tolerant Clock Synchronization · J. ACM 1995
Integrated circuit design › digital circuit design
pulse generation
0.212014
Fault-tolerant algorithms for tick-generation in asynchronous logic: Robust pulse generation · J. ACM 2014
Distributed systems › fault tolerance › self-stabilization
self-stabilizing protocols
0.212014
Fault-tolerant algorithms for tick-generation in asynchronous logic: Robust pulse generation · J. ACM 2014
Distributed computing theory › distributed complexity
round complexity
0.232013
Early-deciding consensus is expensive · PODC 2013
Authenticated Algorithms for Byzantine Agreement · SIAM J. Comput. 1983
Polynomial Algorithms for Multiple Processor Agreement · STOC 1982
Distributed systems
clock synchronization
0.252008
Fast self-stabilizing byzantine tolerant digital clock synchronization · PODC 2008
Brief announcement: linear time byzantine self-stabilizing clock synchronization · PODC 2004
Dynamic Fault-Tolerant Clock Synchronization · J. ACM 1995
Distributed systems › fault tolerance
self-stabilization
0.112012
An Optimal Self-Stabilizing Firing Squad · SIAM J. Comput. 2012
Distributed computing theory
asynchronous systems
0.122019
Implementing Mediators with Asynchronous Cheap Talk · PODC 2019
On the Minimal Synchronism Needed for Distributed Consensus · FOCS 1983
Blockchain and cryptocurrency security › blockchain scalability
scalable consensus
0.112019
Communication Complexity of Byzantine Agreement, Revisited · PODC 2019
Distributed systems
distributed coordination and fault tolerance
0.142008
OCD: obsessive consensus disorder (or repetitive consensus) · PODC 2008
Failure Detectors in Omission Failure Environments · PODC 1997
Atomic Snapshots of Shared Memory · J. ACM 1993
Internet architecture and protocols › network synchronization
network time protocol
0.112018
Preventing (Network) Time Travel with Chronos · NDSS 2018
Distributed computing theory › fault tolerance › byzantine fault tolerance › byzantine agreement
asynchronous byzantine agreement
0.122008
An almost-surely terminating polynomial protocol forasynchronous byzantine agreement with optimal resilience · PODC 2008
Asynchronous Byzantine Consensus · PODC 1984
Reconfigurable computing and FPGAs
programmable devices
0.112008
Tapping into the fountain of CPUs: on operating system support for programmable devices · ASPLOS 2008
Distributed systems
distributed computing theory
0.122006
Century papers at the first quarter-century milestone · PODC 2006
Renaming in an Asynchronous Environment · J. ACM 1990
Distributed computing theory
distributed algorithms
0.152003
Asynchronous resource discovery · PODC 2003
Observable Clock Synchronization (Extended Abstract) · PODC 1994
On Distributed Algorithms in a Broadcast Domain · ICALP 1993
Internet architecture and protocols
multicast
0.112006
On multicast trees: structure and size estimation · IEEE/ACM Trans. Netw. 2006
Internet architecture and protocols › multicast
multicast tree
0.112006
On multicast trees: structure and size estimation · IEEE/ACM Trans. Netw. 2006
Distributed systems › fault tolerance › fault-tolerant protocols
reliable broadcast
0.112006
Self-stabilizing byzantine agreement · PODC 2006
Integrated circuit design
asynchronous circuit design
0.112014
Fault-tolerant algorithms for tick-generation in asynchronous logic: Robust pulse generation · J. ACM 2014
Distributed systems
group communication
0.131998
An Adaptive Totally Ordered Multicast Protocol That Tolerates Partitions · PODC 1998
Efficient Message Ordering in Dynamic Networks · PODC 1996
A Framework for Partitionable Membership Service (Abstract) · PODC 1996

Methods — techniques the papers use, named apart from their topics

cheap talk · 0.4synchronous distributed algorithms · 0.3hardness proof · 0.3transient fault modeling · 0.2probabilistic convergence analysis · 0.2early-stopping protocols · 0.2probabilistic self-stabilization · 0.2byzantine fault tolerance · 0.2lower bounds · 0.2lower bound · 0.2shunning verifiable secret sharing · 0.2operating system abstraction design · 0.2state machine replication · 0.1nash equilibrium · 0.1self-stabilization · 0.1prototype implementation · 0.1simulation · 0.1analytical modeling · 0.0
YearPublicationVenuePosition
2023 Colordag: An Incentive-Compatible Blockchain
Ittai Abraham, Danny Dolev, Ittay Eyal, Joseph Y. Halpern
DISC2
2023 Communication complexity of byzantine agreement, revisited
Ittai Abraham, T.-H. Hubert Chan, Danny Dolev, Kartik Nayak, Rafael Pass, Ling Ren 0001, Elaine Shi
Distributed Comput.3
2022 Brief Announcement: Authenticated Consensus in Synchronous Systems with Mixed Faults
Ittai Abraham, Danny Dolev, Alon Kagan, Gilad Stern
DISC2
2022 Revisiting asynchronous fault tolerant computation with optimal resilience
abstract
The celebrated result of Fischer, Lynch and Paterson is the fundamental lower bound for asynchronous fault tolerant computation: any 1-crash resilient asynchronous agreement protocol must have some (possibly measure zero) probability of not terminating. In 1994, Ben-Or, Kelmer and Rabin published a proof-sketch of a lesser known lower bound for asynchronous fault tolerant computation with optimal resilience in face of a Byzantine adversary: if $$n\le 4t$$ then any t-resilient asynchronous verifiable secret sharing protocol must have some non-zero probability of not terminating. Our main contribution is to revisit this lower bound and provide a rigorous and more general proof. Our second contribution is to show how to avoid this lower bound. We provide a protocol with optimal resilience that is almost surely terminating for a strong common coin functionality. Using this new primitive we provide an almost surely terminating protocol with optimal resilience for asynchronous Byzantine agreement that has a new fair validity property. To the best of our knowledge this is the first asynchronous Byzantine agreement with fair validity in the information theoretic setting.
Ittai Abraham, Danny Dolev, Gilad Stern
Distributed Comput.2
2020 Revisiting Asynchronous Fault Tolerant Computation with Optimal Resilience
Ittai Abraham, Danny Dolev, Gilad Stern
PODC2
2019 Communication Complexity of Byzantine Agreement, Revisited
abstract
As Byzantine Agreement (BA) protocols find application in large-scale decentralized cryptocurrencies, an increasingly important problem is to design BA protocols with improved communication complexity. A few existing works have shown how to achieve subquadratic BA under an adaptive adversary. Intriguingly, they all make a common relaxation about the adaptivity of the attacker, that is, if an honest node sends a message and then gets corrupted in some round, the adversary cannot erase the message that was already sent - henceforth we say that such an adversary cannot perform "after-the-fact removal". By contrast, many (super-)quadratic BA protocols in the literature can tolerate after-the-fact removal. In this paper, we first prove that disallowing after-the-fact removal is necessary for achieving subquadratic-communication BA.
Ittai Abraham, T.-H. Hubert Chan, Danny Dolev, Kartik Nayak, Rafael Pass, Ling Ren 0001, Elaine Shi
PODC3
2019 Implementing Mediators with Asynchronous Cheap Talk
abstract
A mediator can help non-cooperative agents obtain an equilibrium that may otherwise not be possible. We study the ability of players to obtain the same equilibrium without a mediator, using only cheap talk, that is, nonbinding pre-play communication. Previous work has considered this problem in a synchronous setting. Here we consider the effect of asynchrony on the problem, and provide upper bounds for implementing mediators. Considering asynchronous environments introduces new subtleties, including exactly what solution concept is most appropriate and determining what move is played if the cheap talk goes on forever. Different results are obtained depending on whether the move after such "infinite play'' is under the control of the players or part of the description of the game.
Ittai Abraham, Danny Dolev, Ivan Geffner, Joseph Y. Halpern
PODC2
2018 Preventing (Network) Time Travel with Chronos
Omer Deutsch, Neta Rozen Schiff, Danny Dolev, Michael Schapira
NDSS3
2018 Compact routing messages in self-healing trees
Armando Castañeda, Danny Dolev, Amitabh Trehan
Theor. Comput. Sci.2
2017 Stateless Computation
abstract
We present and explore a model of stateless and self-stabilizing distributed computation, inspired by real-world applications such as routing on today's Internet. Processors in our model do not have an internal state, but rather interact by repeatedly mapping incoming messages ("labels") to outgoing messages and output values. While seemingly too restrictive to be of interest, stateless computation encompasses both classical game-theoretic notions of strategic interaction and a broad range of practical applications (e.g., Internet protocols, circuits, diffusion of technologies in social networks). Our main technical contribution is a general impossibility result for stateless self-stabilization in our model, showing that even modest asynchrony (with wait times that are linear in the number of processors) can prevent a stateless protocol from reaching a stable global configuration. Furthermore, we present hardness results for verifying stateless self-stabilization. We also address several aspects of the computational power of stateless protocols. Most significantly, we show that short messages (of length that is logarithmic in the number of processors) yield substantial computational power, even on very poorly connected topologies.
Danny Dolev, Michael Erdmann, Neil Lutz, Michael Schapira, Adva Zair
PODC1
2016 HEX: Scaling honeycombs is easier than scaling clock trees
abstract
We argue that a hexagonal grid with simple intermediate nodes is a robust alternative to buffered clock trees typically used for clock distribution in VLSI circuits, multi-core processors, and other applications that require accurate synchronization: Our HEX grid is Byzantine fault-tolerant, self-stabilizing, and seamlessly integrates with multiple synchronized clock sources, as used in multi-synchronous Globally Synchronous Locally Asynchronous (GALS) architectures. Moreover, HEX guarantees a small clock skew between neighbors even for wire delays that are only moderately balanced. We provide both a theoretical analysis of the worst-case skew and simulation results that demonstrate a very small average skew.
Danny Dolev, Matthias Függer, Christoph Lenzen 0001, Martin Perner, Ulrich Schmid 0001
J. Comput. Syst. Sci.1
2016 Synchronous counting and computational algorithm design
Danny Dolev, Keijo Heljanko, Matti Järvisalo, Janne H. Korhonen, Christoph Lenzen 0001, Joel Rybicki, Jukka Suomela, Siert Wieringa
J. Comput. Syst. Sci.1
2015 Byzantine Agreement with Optimal Early Stopping, Optimal Resilience and Polynomial Complexity
abstract
We provide the first protocol that solves Byzantine agreement with optimal early stopping (min{f+2,t+1} rounds) and optimal resilience (n>3t) using polynomial message size and computation.
Ittai Abraham, Danny Dolev
STOC2
2014 Fault-tolerant algorithms for tick-generation in asynchronous logic: Robust pulse generation
abstract
Today’s hardware technology presents a new challenge in designing robust systems. Deep submicron VLSI technology introduces transient and permanent faults that were never considered in low-level system designs in the past. Still, robustness of that part of the system is crucial and needs to be guaranteed for any successful product. Distributed systems, on the other hand, have been dealing with similar issues for decades. However, neither the basic abstractions nor the complexity of contemporary fault-tolerant distributed algorithms match the peculiarities of hardware implementations. This article is intended to be part of an attempt striving to bridge over this gap between theory and practice for the clock synchronization problem. Solving this task sufficiently well will allow to build an ultra-robust high-precision clocking system for hardware designs like systems-on-chips in critical applications. As our first building block, we describe and prove correct a novel distributed, Byzantine fault-tolerant, probabilistically self-stabilizing pulse synchronization protocol, called FATAL, that can be implemented using standard asynchronous digital logic: Correct FATAL nodes are guaranteed to generate pulses (i.e., unnumbered clock ticks) in a synchronized way, despite a certain fraction of nodes being faulty. FATAL uses randomization only during stabilization and, despite the strict limitations introduced by hardware designs, offers optimal resilience and smaller complexity than all existing protocols. Finally, we show how to leverage FATAL to efficiently generate synchronized, self-stabilizing, high-frequency clocks.
Danny Dolev, Matthias Függer, Ulrich Schmid 0001, Christoph Lenzen 0001
J. ACM1
2014 Rigorously modeling self-stabilizing fault-tolerant circuits: An ultra-robust clocking scheme for systems-on-chip
abstract
We present the first implementation of a distributed clock generation scheme for Systems-on-Chip that recovers from an unbounded number of arbitrary transient faults despite a large number of arbitrary permanent faults. We devise self-stabilizing hardware building blocks and a hybrid synchronous/asynchronous state machine enabling metastability-free transitions of the algorithm's states. We provide a comprehensive modeling approach that permits to prove, given correctness of the constructed low-level building blocks, the high-level properties of the synchronization algorithm (which have been established in a more abstract model). We believe this approach to be of interest in its own right, since this is the first technique permitting to mathematically verify, at manageable complexity, high-level properties of a fault-prone system in terms of its very basic components. We evaluate a prototype implementation, which has been designed in VHDL, using the Petrify tool in conjunction with some extensions, and synthesized for an Altera Cyclone FPGA.
Danny Dolev, Matthias Függer, Markus Posch, Ulrich Schmid 0001, Andreas Steininger, Christoph Lenzen 0001
J. Comput. Syst. Sci.1
2013 Early-deciding consensus is expensive
abstract
In consensus, the n nodes of a distributed system seek to take a consistent decision on some output, despite up to t of them crashing or even failing maliciously, i.e., behaving "Byzantine''. It is known that it is impossible to guarantee that synchronous, deterministic algorithms consistently decide on an output in fewer than f+1 rounds in executions in which the actual number of faults is f ≤ t. This even holds if faults are crash-only, and in this case the bound can be matched precisely. However, the question of whether this can be done efficiently, i.e., with little communication, so far has not been addressed.
Danny Dolev, Christoph Lenzen 0001
PODC1
2013 HEX: scaling honeycombs is easier than scaling clock trees
abstract
We argue that grid structures are a very promising alternative to the standard approach for distributing a clock signal throughout VLSI circuits and other hardware devices. Traditionally, this is accomplished by a delay-balanced clock tree, which distributes the signal supplied by a single clock source via carefully engineered and buffered signal paths.
Danny Dolev, Matthias Függer, Christoph Lenzen 0001, Martin Perner, Ulrich Schmid 0001
SPAA1
2013 Synchronous Counting and Computational Algorithm Design
Danny Dolev, Janne H. Korhonen, Christoph Lenzen 0001, Joel Rybicki, Jukka Suomela
SSS1
2013 Distributed Protocols for Leader Election: A Game-Theoretic Perspective
Ittai Abraham, Danny Dolev, Joseph Y. Halpern
DISC2
2013 Enhanced calibration technique for RSSI-based ranging in body area networks
Gaddi Blumrosen, Bracha Hod, Tal Anker, Danny Dolev, Boris Rubinsky
Ad Hoc Networks4
2013 Enhancing RSSI-based tracking accuracy in wireless sensor networks
abstract
In recent years, the demand for high-precision tracking systems has significantly increased in the field of Wireless Sensor Network (WSN). A new tracking system based on exploitation of Received Signal Strength Indicator (RSSI) measurements in WSN is proposed. The proposed system is designed in particular for WSNs that are deployed in close proximity and can transmit data at a high transmission rate. The close proximity and an optimized transmit power level enable accurate conversion of RSSI measurements to range estimates. Having an adequate transmission rate enables spatial-temporal correlation between consecutive RSSI measurements. In addition, advanced statistical and signal processing methods are used to mitigate channel distortion and to compensate for packet loss. The system is evaluated in indoor conditions and achieves tracking resolution of a few centimeters which is compatible with theoretical bounds.
Gaddi Blumrosen, Bracha Hod, Tal Anker, Danny Dolev, Boris Rubinsky
ACM Trans. Sens. Networks4
2012 No justified complaints: on fair sharing of multiple resources
abstract
Fair allocation has been studied intensively in both economics and computer science. This subject has aroused renewed interest with the advent of virtualization and cloud computing. Prior work has typically focused on mechanisms for fair sharing of a single resource. We consider a variant where each user is entitled to a certain fraction of the system's resources, and has a fixed usage profile describing how much he would want from each resource. We provide a new definition for the simultaneous fair allocation of multiple continuously-divisible resources that we call bottleneck-based fairness (BBF). Roughly speaking, an allocation of resources is considered fair if every user either gets all the resources he wishes for, or else gets at least his entitlement on some bottleneck resource, and therefore cannot complain about not receiving more. We show that BBF has several desirable properties such as providing an incentive for sharing, and also promotes high overall utilization of resources; we also compare BBF carefully to another notion of fairness proposed recently, dominant resource fairness.
Danny Dolev, Dror G. Feitelson, Joseph Y. Halpern, Raz Kupferman, Nathan Linial
ITCS1
2012 "Tri, Tri Again": Finding Triangles and Small Subgraphs in a Distributed Setting - (Extended Abstract)
Danny Dolev, Christoph Lenzen 0001, Shir Peled
DISC1
2012 An Optimal Self-Stabilizing Firing Squad
abstract
Consider a fully connected network where up to t processes may crash and all processes start in an arbitrary memory state. The self-stabilizing firing squad problem consists of eventually guaranteeing simultaneous response to an external input. This is modeled by requiring that the noncrashed processes “fire” simultaneously if some correct process received an external “go” input, and that they only fire as a response to some process receiving such an input. This paper presents FireSquad, the first self-stabilizing firing squad algorithm. A firing squad algorithm facilitates the use of algorithms that need to start in the same round. It allows a smooth transition between algorithms whose executions need to be disjoint. The FireSquad algorithm combines two forms of fault-tolerance properties: self-stabilization to allow recovery from arbitrary transient errors and resilience to crash failures to handle permanent ones. The FireSquad algorithm is optimal in two respects: (a) once the algorithm is in a safe state, it fires in response to a go input as fast as any other algorithm does, and (b) starting from an arbitrary state, it converges to a safe state as fast as any other algorithm does.
Danny Dolev, Ezra N. Hoch, Yoram Moses
SIAM J. Comput.1
2011 Delay Fast Packets (DFP): Prevention of DNS Cache Poisoning
Shimrit Tzur-David, Kiril Lashchiver, Danny Dolev, Tal Anker
SecureComm3
2011 Fault-Tolerant Algorithms for Tick-Generation in Asynchronous Logic: Robust Pulse Generation - [Extended Abstract]
Danny Dolev, Matthias Függer, Christoph Lenzen 0001, Ulrich Schmid 0001
SSS1
2010 Continuous Close-Proximity RSSI-Based Tracking in Wireless Sensor Networks
abstract
In this paper we develop a continuous high-precision tracking system based on Received Signal Strength Indicator(RSSI) measurements for small ranges. The proposed system uses minimal number of sensor nodes with RSSI capabilities to track a moving object in close-proximity and high transmission rate. The close-proximity enables conversion of RSSI measurements to range estimates and the high transmission rate enables continuous tracking of the moving object. The RSSI-based tracking system includes calibration, range estimation, location estimation and refinement. We use advanced statistical and signal processing methods to mitigate channel distortion and packet loss. The system is evaluated in indoor settings and achieves tracking resolution of few centimeters. Therefore, it becomes the motion trackers of notice in many applications.
Gaddi Blumrosen, Bracha Hod, Tal Anker, Danny Dolev, Boris Rubinsky
BSN4
2010 SPADE: Statistical Packet Acceptance Defense Engine
abstract
A security engine should detect network traffic attacks at line-speed. "Learning" capabilities can help detecting new and unknown threats even before a vulnerability is exploited. The principal way for achieving this goal is to model anticipated network traffic behavior, and to use this model for identifying anomalies. This paper focuses on denial of service (DoS) attacks and distributed DoS (DDoS). Our goal is detecting and preventing of attacks. The main challenges include minimizing the false-positive rate and the memory consumption. SPADE: a Statistical Packet Acceptance Defense Engine is presented. SPADE is an accurate engine that uses an hierarchical adaptive structure to detect suspicious traffic using a relatively small memory footprint, therefore can be easily applied on hardware. SPADE is based on the assumption that during DoS/DDoS attacks, a significant portion of the traffic that is seen belongs to the attack, therefore, SPADE applies a statistical mechanism to primarily filter the attack's traffic.
Shimrit Tzur-David, Harel Avissar, Danny Dolev, Tal Anker
HPSR3
2010 A Fault-Resistant Asynchronous Clock Function
Ezra N. Hoch, Michael Ben-Or, Danny Dolev
SSS3
2010 Brief Announcement: Simple Gradecast Based Algorithms
Michael Ben-Or, Danny Dolev, Ezra N. Hoch
DISC2
2010 Peer-to-peer secure multi-party numerical computation facing malicious adversaries
Danny Bickson, Tzachy Reinman, Danny Dolev, Benny Pinkas
Peer-to-Peer Netw. Appl.3
2010 Steward: Scaling Byzantine Fault-Tolerant Replication to Wide Area Networks
abstract
This paper presents the first hierarchical byzantine fault-tolerant replication architecture suitable to systems that span multiple wide-area sites. The architecture confines the effects of any malicious replica to its local site, reduces message complexity of wide-area communication, and allows read-only queries to be performed locally within a site for the price of additional standard hardware. We present proofs that our algorithm provides safety and liveness properties. A prototype implementation is evaluated over several network topologies and is compared with a flat byzantine fault-tolerant approach. The experimental results show considerable improvement over flat byzantine replication algorithms, bringing the performance of byzantine replication closer to existing benign fault-tolerant replication techniques over wide area networks.
Yair Amir, Claudiu Danilov 0001, Danny Dolev, Jonathan Kirsch, John Lane, Cristina Nita-Rotaru, Josh Olsen, David Zage
IEEE Trans. Dependable Secur. Comput.3
2009 Fixing convergence of Gaussian belief propagation
abstract
Gaussian belief propagation (GaBP) is an iterative message-passing algorithm for inference in Gaussian graphical models. It is known that when GaBP converges it converges to the correct MAP estimate of the Gaussian random vector and simple sufficient conditions for its convergence have been established. In this paper we develop a double-loop algorithm for forcing convergence of GaBP. Our method computes the correct MAP estimate even in cases where standard GaBP would not have converged. We further extend this construction to compute least-squares solutions of over-constrained linear systems. We believe that our construction has numerous applications, since the GaBP algorithm is linked to solution of linear systems of equations, which is a fundamental problem in computer science and engineering. As a case study, we discuss the linear detection problem. We show that using our new construction, we are able to force convergence of Montanari's linear detection algorithm, in cases where it would originally fail. As a consequence, we are able to increase significantly the number of users that can transmit concurrently.
Danny Dolev, Danny Bickson, Jason K. Johnson
ISIT1
2009 Distributed large scale network utility maximization
abstract
Recent work by Zymnis et al. proposes an efficient primal-dual interior-point method, using a truncated Newton method, for solving the network utility maximization (NUM) problem. This method has shown superior performance relative to the traditional dual-decomposition approach. Other recent work by Bickson et al. shows how to compute efficiently and distributively the Newton step, which is the main computational bottleneck of the Newton method, utilizing the Gaussian belief propagation algorithm. In the current work, we combine both approaches to create an efficient distributed algorithm for solving the NUM problem. Unlike the work of Zymnis, which uses a centralized approach, our new algorithm is easily distributed. Using an empirical evaluation we show that our new method outperforms previous approaches, including the truncated Newton method and dual-decomposition methods. As an additional contribution, this is the first work that evaluates the performance of the Gaussian belief propagation algorithm vs. the preconditioned conjugate gradient method, for a large scale problem.
Danny Dolev, Argyris Zymnis, Stephen P. Boyd, Danny Bickson, Yoav Tock
ISIT1
2009 MULAN: Multi-Level Adaptive Network Filter
Shimrit Tzur-David, Danny Dolev, Tal Anker
SecureComm2
2009 Distributed data flow language for multi-party protocols
abstract
This paper presents a novel object-oriented approach to modeling the semantics of distributed multi-party protocols such as leader election, distributed locks or reliable multicast, and a programming language that supports it. The approach extends our live distributed objects (LO) model with the new concept of a distributed flow (DF), a stream of events that flow concurrently at multiple locations. DFs correspond to local variables, private fields, and method parameters in Java-like languages; they're means by which one stores and communicates state. Protocol instances correspond to Java objects; they consume and output flows; their internal states are encapsulated as internal flows, and their internal logic is represented as operations on flows. Our language provides a new type of concern separation: the semantic structure of protocols is decoupled from implementation details such as construction and maintenance of overlays, trees, and other structures used for scalability. These can be generated by the compiler or at deployment time. This can be done differently in different parts of the network, to match the local environment.
Krzysztof Ostrowski, Kenneth P. Birman, Danny Dolev
PLOS@SOSP3
2009 An Optimal Self-stabilizing Firing Squad
Danny Dolev, Ezra N. Hoch, Yoram Moses
SSS1
2008 Tapping into the fountain of CPUs: on operating system support for programmable devices
abstract
The constant race for faster and more powerful CPUs is drawing to a close. No longer is it feasible to significantly increase the speed of the CPU without paying a crushing penalty in power consumption and production costs. Instead of increasing single thread performance, the industry is turning to multiple CPU threads or cores (such as SMT and CMP) and heterogeneous CPU architectures (such as the Cell Broadband Engine). While this is a step in the right direction, in every modern PC there is a wealth of untapped compute resources. The NIC has a CPU; the disk controller is programmable; some high-end graphics adaptersare already more powerful than host CPUs. Some of these CPUs can perform some functions more efficiently than the host CPUs. Our operating systems and programming abstractions should be expanded to let applications tap into these computational resources and make the best use of them.
Yaron Weinsberg, Danny Dolev, Tal Anker, Muli Ben-Yehuda, Pete Wyckoff
ASPLOS2
2008 Programming with Live Distributed Objects
Krzysztof Ostrowski, Kenneth P. Birman, Danny Dolev, Jong Hoon Ahnn
ECOOP3
2008 Efficient Clustering for Improving Network Performance in Wireless Sensor Networks
Tal Anker, Danny Bickson, Danny Dolev, Bracha Hod
EWSN3
2008 LiteLoad: Content unaware routing for localizing P2P protocols
abstract
In today's extensive worldwide Internet traffic, some 60% of network congestion is caused by Peer to Peer sessions. Consequently ISPs are facing many challenges like: paying for the added traffic requirement, poor customer satisfaction due to degraded broadband experience, purchasing costly backbone links and upstream bandwidth and having difficulty to effectively control P2P traffic with conventional devices. Existing solutions such as caching and indexing of P2P content are controversial as their legality is uncertain due to copyright violation, and therefore hardly being installed by ISPs. In addition these solutions are not capable to handle existing encrypted protocols that are on the rise in popular P2P networks. Other solutions that employ traffic shaping and blocking degrade the downloading throughput and cause end users to switch ISPs for a better service. LiteLoad discerns patterns of user communications in Peer to Peer file sharing networks without identifying the content being requested or transferred and uses least-cost routing rules to push peer-to-peer transfers into confined network segments. This approach maintains the performance of file transfer as opposed to traffic shaping solutions and precludes internet provider involvement in caching, cataloguing or indexing of the shared content. Simulation results expresses the potential of the solution and a proof of concept of the key technology is demonstrated on popular protocols, including encrypted ones.
Shay Horovitz, Danny Dolev
IPDPS2
2008 Gaussian belief propagation based multiuser detection
abstract
In this work, we present a novel construction for solving the linear multiuser detection problem using the Gaussian Belief Propagation algorithm. Our algorithm yields an efficient, iterative and distributed implementation of the MMSE detector. Compared to our previous formulation, the new algorithm offers a reduction in memory requirements, the number of computational steps, and the number of messages passed. We prove that a detection method recently proposed by Montanari et al. is an instance of ours, and we provide new convergence results applicable to both.
Danny Bickson, Danny Dolev, Ori Shental, Paul H. Siegel, Jack K. Wolf
ISIT2
2008 Gaussian belief propagation solver for systems of linear equations
abstract
The canonical problem of solving a system of linear equations arises in numerous contexts in information theory, communication theory, and related fields. In this contribution, we develop a solution based upon Gaussian belief propagation (GaBP) that does not involve direct matrix inversion. The iterative nature of our approach allows for a distributed message-passing implementation of the solution algorithm. We also address some properties of the GaBP solver, including convergence, exactness, its max-product version and relation to classical solution methods. The application example of decorrelation in CDMA is used to demonstrate the faster convergence rate of the proposed solver in comparison to conventional linear-algebraic iterative solution methods.
Ori Shental, Paul H. Siegel, Jack K. Wolf, Danny Bickson, Danny Dolev
ISIT5
2008 Quicksilver Scalable Multicast (QSM)
abstract
QSM is a multicast engine designed to support a style of distributed programming in which application objects are replicated among clients and updated via multicast. The model requires platforms that scale in dimensions previously unexplored; in particular, to large numbers of multicast groups. Prior systems werenpsilat optimized for such scenarios and canpsilat take advantage of regular group overlap patterns, a key feature of our application domain. Furthermore, little is known about performance and scalability of such systems in modern managed environments. We shed light on these issues and offer architectural insights based on our experience building QSM.
Krzysztof Ostrowski, Kenneth P. Birman, Danny Dolev
NCA3
2008 Nysiad: Practical Protocol Transformation to Tolerate Byzantine Failures
Chi Ho, Robbert van Renesse, Mark Bickford, Danny Dolev
NSDI4
2008 Peer-to-Peer Secure Multi-party Numerical Computation
abstract
We propose an efficient framework for enabling secure multi-party numerical computations in a Peer-to-Peer network. This problem arises in a range of applications such as collaborative filtering, distributed computation of trust and reputation, monitoring and numerous other tasks, where the computing nodes would like to preserve the privacy of their inputs while performing a joint computation of a certain function. Although there is a rich literature in the field of distributed systems security concerning secure multi-party computation, in practice it is hard to deploy those methods in very large scalePeer-to-Peer networks. In this work, we examine several possible approaches and discuss their feasibility. Among the possible approaches, we identify a single approach which is both scalable and theoretically secure. An additional novel contribution is that we show how to compute the neighborhood based collaborative filtering, a state-of-the-art collaborative filtering algorithm, winner of the Netflix progress prize of the year 2007. Our solution computes this algorithm in a Peer-to-Peer network, using a privacy preserving computation, without loss of accuracy. Using extensive large scale simulations on top of real Internet topologies, we demonstrate the applicability of our approach. Asfar as we know, we are the first to implement such a large scale secure multi-party simulation of networks of millions of nodes and hundreds of millions of edges.
Danny Bickson, Danny Dolev, Genia Bezman, Benny Pinkas
Peer-to-Peer Computing2
2008 An almost-surely terminating polynomial protocol forasynchronous byzantine agreement with optimal resilience
abstract
Consider an asynchronous system with private channels and n processes, up to t of which may be faulty. We settle a longstanding open question by providing a Byzantine agreement protocol that simultaneously achieves three properties: (optimal) resilience: it works as long as n>3t;(almost-sure) termination: with probability one, all nonfaulty processes terminate;(polynomial) efficiency: the expected computation time, memory consumption, message size, and number of messages sent are all polynomial in n. Earlier protocols have achieved only two of these three properties. In particular, the protocol of Bracha is not polynomially efficient, the protocol of Feldman and Micali is not optimally resilient, and the protocol of Canetti and Rabin does not have almost-sure termination. Our protocol utilizes a new primitive called shunning (asynchronous) verifiable secret sharing (SVSS), which ensures, roughly speaking, that either a secret is successfully shared or a new faulty process is ignored from this point onwards by some nonfaulty process.
Ittai Abraham, Danny Dolev, Joseph Y. Halpern
PODC2
2008 Fast self-stabilizing byzantine tolerant digital clock synchronization
abstract
Consider a distributed network in which up to a third of the nodes may be Byzantine, and in which the non-faulty nodes may be subject to transient faults that alter their memory in an arbitrary fashion. Within the context of this model, we are interested in the digital clock synchronization problem; which consists of agreeing on bounded integer counters, and increasing these counters regularly. It has been postulated in the past that synchronization cannot be solved in a Byzantine tolerant and self-stabilizing manner. The first solution to this problem had an expected exponential convergence time. Later, a deterministic solution was published with linear convergence time, which is optimal for deterministic solutions. In the current paper we achieve an expected constant convergence time. We thus obtain the optimal probabilistic solution, both in terms of convergence time and in terms of resilience to Byzantine adversaries.
Michael Ben-Or, Danny Dolev, Ezra N. Hoch
PODC2
2008 OCD: obsessive consensus disorder (or repetitive consensus)
abstract
Consider a distributed system S of sensors, where the goal is to continuously output an agreed reading. The input readings of non-faulty sensors may change over time; and some of the sensors may be faulty (Byzantine). Thus, the system is required to repeatedly perform consensus on the input values.
Danny Dolev, Ezra N. Hoch
PODC1
2008 Self-stabilizing Numerical Iterative Computation
Ezra N. Hoch, Danny Bickson, Danny Dolev
SSS3
2008 Lower Bounds on Implementing Robust and Resilient Mediators
Ittai Abraham, Danny Dolev, Joseph Y. Halpern
TCC2
2008 Belief Propagation in Wireless Sensor Networks - A Practical Approach
Tal Anker, Danny Dolev, Bracha Hod
WASA2
2008 Constant-Space Localized Byzantine Consensus
Danny Dolev, Ezra N. Hoch
DISC1
2007 Accelerating Distributed Computing Applications Using a Network Offloading Framework
abstract
During the last two decades, a considerable amount of academic research has been conducted in the field of distributed computing. Typically, distributed applications require frequent network communication, which becomes a dominate factor in the overall runtime overhead. The recent proliferation of programmable peripheral devices for computer systems may be utilized in order to improve the performance of such applications. Offloading application-specific network functions to peripheral devices can improve performance and reduce host CPU utilization. Due to the peculiarities of each particular device and the difficulty of programming an outboard CPU, the need for an abstracted offloading framework is apparent. This paper proposes a novel offloading framework, called HYDRA that enables utilization of such devices. The framework enables an application developer to design the offloading aspects of the application by specifying an "offloading layout", which is enforced by the runtime during application deployment. The performance of a variety of distributed algorithms can be significantly improved by utilizing such a framework. We demonstrate this claim by evaluating several offloaded applications: a distributed total message ordering algorithm and a packet generator.
Yaron Weinsberg, Danny Dolev, Pete Wyckoff, Tal Anker
IPDPS2
2007 Self-stabilizing and Byzantine-Tolerant Overlay Network
Danny Dolev, Ezra N. Hoch, Robbert van Renesse
OPODIS1
2007 Making Distributed Applications Robust
Chi Ho, Danny Dolev, Robbert van Renesse
OPODIS2
2007 Byzantine Self-stabilizing Pulse in a Bounded-Delay Model
Danny Dolev, Ezra N. Hoch
SSS1
2007 On Self-stabilizing Synchronous Actions Despite Byzantine Attacks
Danny Dolev, Ezra N. Hoch
DISC1
2006 Scaling Byzantine Fault-Tolerant Replication toWide Area Networks
abstract
This paper presents the first hierarchical Byzantine fault-tolerant replication architecture suitable to systems that span multiple wide area sites. The architecture confines the effects of any malicious replica to its local site, reduces message complexity of wide area communication, and allows read-only queries to be performed locally within a site for the price of additional hardware. A prototype implementation is evaluated over several network topologies and is compared with a flat Byzantine fault-tolerant approach
Yair Amir, Claudiu Danilov 0001, Jonathan Kirsch, John Lane, Danny Dolev, Cristina Nita-Rotaru, Josh Olsen, David Zage
DSN5
2006 On a NIC's Operating System, Schedulers and High-Performance Networking Applications
Yaron Weinsberg, Tal Anker, Danny Dolev, Scott Kirkpatrick
HPCC3
2006 Wire-speed total order
abstract
Many distributed systems may be limited in their performance by the number of transactions they are able to support per unit of time. In order to achieve fault tolerance and to boost a system's performance, active state machine replication is frequently used. It employs total ordering service to keep the state of replicas synchronized. In this paper, we present an architecture that enables a drastic increase in the number of ordered transactions in a cluster, using off-the-shelf network equipment. Performance supporting nearly one million ordered transactions per second has been achieved, which substantiates our claim.
Tal Anker, Danny Dolev, Gregory Greenman, Ilya Shnaiderman
IPDPS2
2006 HYDRA: A Novel Framework for Making High-Performance Computing Offload Capable
abstract
The proliferation of programmable peripheral devices for computer systems opens new possibilities for academic research that will influence system designs in the near future. Programmability is a key feature that enables application-specific extensions to improve performance and offer new features. Increasing transistor density and decreasing cost provide excess computational power in devices such as disk controllers, network interfaces and video cards. This paper proposes an innovative programming model and runtime support that enables utilization of such devices by providing a generic code offloading framework. The framework enables an application developer to design the offloading aspects of the application by specifying an "offloading layout", which is enforced by the runtime during application deployment. The framework also provides the necessary development tools and programming constructs for developing such applications. We test our framework by implementing a packet generator on a programmable network card for network testing. The offloaded application produces traffic at five times the rate, and with inter-packet variability that is many orders of magnitude smaller than the non-offloaded version
Yaron Weinsberg, Danny Dolev, Tal Anker, Pete Wyckoff
LCN2
2006 Distributed computing meets game theory: robust mechanisms for rational secret sharing and multiparty computation
abstract
We study k-resilient Nash equilibria, joint strategies where no member of a coalition C of size up to k can do better, even if the whole coalition defects. We show that such k-resilient Nash equilibria exist for secret sharing and multiparty computation, provided that players prefer to get the information than not to get it. Our results hold even if there are only 2 players, so we can do multiparty computation with only two rational agents. We extend our results so that they hold even in the presence of up to t players with "unexpected" utilities. Finally, we show that our techniques can be used to simulate games with mediators by games without mediators.
Ittai Abraham, Danny Dolev, Rica Gonen, Joseph Y. Halpern
PODC2
2006 Self-stabilizing byzantine agreement
abstract
Byzantine agreement algorithms typically assume implicit initial state consistency and synchronization among the correct nodes and then operate in coordinated rounds of information exchange to reach agreement based on the input values. The implicit initial assumptions enable correct nodes to infer about the progression of the algorithm at other nodes from their local state. This paper considers a more severe fault model than permanent Byzantine failures, one in which the system can in addition be subject to severe transient failures that can temporarily throw the system out of its assumption boundaries. When the system eventually returns to behave according to the presumed assumptions it may be in an arbitrary state in which any synchronization among the nodes might be lost, and each node may be at an arbitrary state. We present a self-stabilizing Byzantine agreement algorithm that reaches agreement among the correct nodes in optimal time, by using only the assumption of bounded message transmission delay. In the process of solving the problem, two additional important and challenging building blocks were developed: a unique self-stabilizing protocol for assigning consistent relative times to protocol initialization and a Reliable Broadcast primitive that progresses at the speed of actual message delivery time.
Ariel Daliot, Danny Dolev
PODC2
2006 Century papers at the first quarter-century milestone
abstract
My talk intends to look back at the first 25 years of PODC to surface issues regarding its potential impact in light of the extensive spread out of Distributed Computing. Distributed computing expands to a level that no one anticipated even few years ago. In the near future distributed computing will be part of every aspect of our life -- from monitoring personal health to environment tracking, from communication to commuting and will be applied from within a chip through a multi-cpu machine to a virtual globally spanned machine. Society will be completely depended on the ability to continuously provide services at one level or another. The various systems will interact and will never be in a steady state. Are we ready to address such issues? Is the tool-box PODC have developed over the last quarter century will help the community to address the challenges of such systems.When PODC was founded there were skeptics that did not view Distributed Computing as a viable field of Computer Science. Few pioneers envisioned the role it will play and decided that it is an identified discipline with specific focus of interest that will benefit from having a separate conference.An impact is a multi facet notion and there is no single scale to measure it. In the talk I will discuss the impact of PODC from several angles. PODC went through several cycles over the years, from the first very few embryonic years, through a period of stability, some hard years and the emerging focus over the last few years.To get an independent perspective I chose to also look at PODC through the eyes of Google Scholar. I listed regular papers that have at least 100 references to the PODC version and to the subsequent paper in a journal (when I noticed one). The compiled list was different than any list I would have compiled myself1. In the talk I will discuss the list presented in the bibliography and will explore various objective measures that are reflected by it. I will compare it to various other topics and will try to raise issues that the PODC community needs to discuss toward its future role in the continuously expanding Distributed Computing discipline.I will also comment on various emerging fields in which Distributed Computing will play a major role and will point out challenges that the main-stream research in PODC is not capable of addressing.
Danny Dolev
PODC1
2006 Self-stabilizing Byzantine Digital Clock Synchronization
Ezra N. Hoch, Danny Dolev, Ariel Daliot
SSS2
2006 Cooperative and reliable packet-forwarding on top of AODV
abstract
Cooperative and reliable packet forwarding presents a formidable challenge in mobile ad hoc networks (MANET), due to special network characteristics; e.g., mobility, dynamic topology and absence of centralized management. Lack of cooperation, due to misbehavior caused by selfishness or malice, may severely degrade the performance of the network. Previous studies, relying on reputation systems, have demonstrated solutions designed for Dynamic Source Routing (DSR) protocol. This paper highlights various aspects of cooperation enforcement and reliability, when AODV is the underlying protocol. Furthermore, it presents a scalable protocol that combines a reputation system with AODV that addresses reputation fading, second-chance, robustness against liars and load balancing.
Tal Anker, Danny Dolev, Bracha Hod
WiOpt2
2006 Asynchronous resource discovery
Ittai Abraham, Danny Dolev
Comput. Networks2
2006 Internet resiliency to attacks and failures under BGP policy routing
Danny Dolev, Sugih Jamin, Osnat Mokryn, Yuval Shavitt
Comput. Networks1
2006 On multicast trees: structure and size estimation
Danny Dolev, Osnat Mokryn, Yuval Shavitt
IEEE/ACM Trans. Netw.1
2004 Off-Piste QoS-aware Routing Protocol
abstract
QoS-aware routing protocols have been the focus of much attention in the past decade. This is due to the interesting challenges posed by the problem of QoS routing, in the absence of precise information - or even partial and altogether imprecise information. We present a new QoS-aware routing protocol, OPsAR, which achieves very good performance results at a very reasonable cost, in terms of memory and messages, OPsAR is based on the multi-path routing approach, but constrains the number of paths used and leverages on previous resource reservation attempts. Every reservation attempt - successful or not - results in an update to the knowledge-state recorded at all the nodes that had participated in those attempts. This knowledge helps in avoiding network bottlenecks and coping with congested path when they are encountered. We present extensive simulation results and even compare ourselves favorably with a protocol that explores all possible paths.
Tal Anker, Danny Dolev, Yigal Eliaspur
ICNP2
2004 Optimal Resilience Asynchronous Approximate Agreement
Ittai Abraham, Yonatan Amit, Danny Dolev
OPODIS3
2004 Brief announcement: linear time byzantine self-stabilizing clock synchronization
abstract
We present the first self-stabilizing Byzantine clock synchronization algorithm that has linear convergence time.
Ariel Daliot, Danny Dolev, Hanna Parnas
PODC2
2003 Facilitating Efficient and Reliable Monitoring through HAMSA
David Breitgand, Danny Dolev, Danny Raz, Gleb Shaviner
Integrated Network Management2
2003 On Multicast Trees: Structure and Size Estimation
abstract
This work presents a thorough investigation of the structure of multicast trees cut from the Internet and power-law topologies. Based on both generated topologies and real Internet data, we characterize the structure of such trees and show that they obey the rank-degree power law; that most high degree tree nodes are concentrated in a low diameter neighborhood; and that the sub-tree size also obeys a power law. Our most surprising empirical finding suggests that there is a linear ratio between the number of high-degree network nodes, namely nodes whose tree degree is higher than some constant, and the number of leaf nodes in the multicast tree (clients). We also derive this ratio analytically. Based on this finding, we develop the fast algorithm, that estimates the number of clients, and show that it converges faster than one round trip delay from the root to a randomly selected client.
Danny Dolev, Osnat Mokryn, Yuval Shavitt
INFOCOM1
2003 Linear Time Byzantine Self-Stabilizing Clock Synchronization
Ariel Daliot, Danny Dolev, Hanna Parnas
OPODIS2
2003 Asynchronous resource discovery
abstract
Consider a dynamic, large-scale communication infrastructure (e.g., the Internet) where nodes (e.g., in a peer to peer system) can communicate only with nodes whose id (e.g., IP address) are known to them. One of the basic building blocks of such a distributed system is resource discovery - efficiently discovering the ids of the nodes that currently exist in the system. We present both upper and lower bounds for the resource discovery problem. For the original problem raised by Harchol-Balter, Leighton, and Lewin [3] we present an Ω2(n log n) message complexity lower bound for asynchronous networks whose size is unknown. For this model, we give an asymptotically message optimal algorithm that improves the bit complexity of Kutten and Peleg [4]. When each node knows the size of its connected component, we provide a novel and highly efficient algorithm with near linear O(nα(n, n)) message complexity (where α is the inverse of Ackerman's function). In addition, we define and study the Ad-hoc Resource Discovery Problem, which is a practical relaxation of the original problem. Our algorithm for ad-hoc resource discovery has near linear O(nα(n, n)) message complexity. The algorithm efficiently deals with dynamic node additions to the system, thus addressing an open question of [3]. We present a Ω(nα(n, n)) lower bound for the Ad-hoc Resource Discovery Problem, showing that our algorithm is asymptotically message optimal.
Ittai Abraham, Danny Dolev
PODC2
2003 TCP-Friendly Many-to-Many End-to-End Congestion Control
abstract
The paper addresses the issue of TCP-friendly congestion control mechanism for many-to-many communication environment. Lack of congestion control inhibits deployment of WAN applications that involve collaboration of groups of processes in the Internet environment. Recent efforts targeted unicast WAN congestion control (TFRC). We extend that approach to multicast many-to-many applications that operate using a middleware framework. Our congestion control mechanism was implemented within a group communication middleware and tested in a multi-continent environment. The measurements have proved the proposed approach to be robust, efficient and TCP-friendly, as well as to provide fairness among processes that compete for shared resources.
Tal Anker, Danny Dolev, Ilya Shnayderman, Innocenty Sukhov
SRDS2
2002 An integrated architecture for the scalable delivery of semi-dynamic Web content
abstract
The competition on clients attention requires sites to update their content frequently. As a result, a large percentage of Web pages are semi-dynamic, i.e., change quite often and stay static between changes. The cost of maintaining consistency for such pages discourages caching solutions. We suggest here an integrated architecture for the scalable delivery of frequently changing hot pages. Our scheme enables sites to dynamically select whether to cyclically multicast a hot page or to unicast it, and to switch between multicast and unicast mechanisms in a transparent way. Our scheme defines a new protocol, called h.t.t.p.m. In addition, it uses currently deployed protocols, and dynamically directs browsers seeking for a URL to multicast channels, while using existing DNS mechanisms. Thus, we enable sites to deliver content to a growing number of users at less cost and during denial of service attacks, while reducing load on core links. We report simulation results that demonstrate the advantages of the integrated architecture, and its significant impact on server and network load, as well as clients delay.
Danny Dolev, Osnat Mokryn, Yuval Shavitt, Innocenty Sukhov
ISCC1
2002 Ad Hoc Membership for Scalable Applications
Tal Anker, Danny Dolev, Ilya Shnayderman
DISC2
2002 Moshe: A group membership service for WANs
abstract
We present Moshe, a novel scalable group membership algorithm built specifically for use in wide area networks (WANs), which can suffer partitions. Moshe is designed with three new significant features that are important in this setting: it avoids delivering views that reflect out-of-date memberships; it requires a single round of messages in the common case; and it employs a client-server design for scalability. Furthermore, Moshe's interface supplies the hooks needed to provide clients with full virtual synchrony semantics. We have implemented Moshe on top of a network event mechanism also designed specifically for use in a WAN. In addition to specifying the properties of the algorithm and proving that this specification is met, we provide empirical results of an implementation of Moshe running over the Internet. The empirical results justify the assumptions made by our design and exhibit good performance. In particular, Moshe terminates within a single communication round over 98% of the time. The experimental results also lead to interesting observations regarding the performance of membership algorithms over the Internet.
Idit Keidar, Jeremy B. Sussman, Keith Marzullo, Danny Dolev
ACM Trans. Comput. Syst.4
2001 Neighborhood Preserving Hashing and Approximate Queries
abstract
Let $D \subseteq \Sigma^n$ be a dictionary. We look for efficient data structures and algorithms to solve the following approximate query problem: Given a query $u \in \Sigma^n$ list all words $v \in D$ that are close to u in Hamming distance. The problem reduces to the following combinatorial problem: Hash the vertices of the n-dimensional hypercube into buckets so that (1) the c-neighborhood of each vertex is mapped into at most k buckets and (2) no bucket is too large. Lower and upper bounds are given for the tradeoff between k and the size of the largest bucket. These results are used to derive bounds for the approximate query problem.
Danny Dolev, Yuval Harari, Nathan Linial, Noam Nisan, Michal Parnas
SIAM J. Discret. Math.1
2001 The architecture and performance of security protocols in the ensemble group communication system: Using diamonds to guard the castle
abstract
Ensemble is a Group Communication System built at Cornell and the Hebrew universities. It allows processes to create process groups within which scalable reliable fifo-ordered multicast and point-to-point communication are supported. The system also supports other communication properties, such as causal and total multicast ordering, flow control, and the like. This article describes the security protocols and infrastructure of Ensemble. Applications using Ensemble with the extensions described here benefit from strong security properties. Under the assumption that trusted processes will not be corrupted, all communication is secured from tampering by outsiders. Our work extends previous work performed in the Horus system (Ensemble's predecessor) by adding support for multiple partitions, efficient rekeying, and application-defined security policies. Unlike Horus, which used its own security infrastructure with nonstandard key distribution and timing services, Ensemble's security mechanism is based on off-the shelf authentication systems, such as PGP and Kerberos. We extend previous results on group rekeying, with a novel protocol that makes use of diamondlike data structures. Our Diamond protocol allows the removal of untrusted members within milliseconds. In this work we are considering configurations of hundreds of members, and further assume that member trust policies are symmetric and transitive. These assumptions dictate some of our design decisions.
Ohad Rodeh, Kenneth P. Birman, Danny Dolev
ACM Trans. Inf. Syst. Secur.3
2000 A Client-Server Oriented Algorithm for Virtually Synchronous Group Membership in WANs
abstract
We describe a novel scalable group membership service designed explicitly for wide area networks. Our membership service is scalable in the number of groups supported, in the number of members in each group, and in the topology each group spans. Our service also supplies the hooks needed to provide clients with full virtual synchrony semantics. Our service attains, on average, a low message overhead by agreeing on membership within a single message round. Furthermore, our service avoids notifying the application of obsolete membership views when the network is unstable, yet it converges when the network has stabilized.
Idit Keidar, Jeremy B. Sussman, Keith Marzullo, Danny Dolev
ICDCS4
2000 Implementing a Caching Service for Distributed CORBA Objects
Gregory V. Chockler, Danny Dolev, Roy Friedman 0001, Roman Vitenberg
Middleware2
2000 Optimized Rekey for Group Communication Systems
Ohad Rodeh, Kenneth P. Birman, Danny Dolev
NDSS3
2000 Nonmalleable Cryptography
abstract
The notion of nonmalleable cryptography, an extension of semantically secure cryptography, is defined. Informally, in the context of encryption the additional requirement is that given the ciphertext it is impossible to generate a different ciphertext so that the respective plaintexts are related. The same concept makes sense in the contexts of string commitment and zero-knowledge proofs of possession of knowledge. Nonmalleable schemes for each of these three problems are presented. The schemes do not assume a trusted center; a user need not know anything about the number or identity of other system users. Our cryptosystem is the first proven to be secure against a strong type of chosen ciphertext attack proposed by Rackoff and Simon, in which the attacker knows the ciphertext she wishes to break and can query the decryption oracle on any ciphertext other than the target.
Danny Dolev, Cynthia Dwork, Moni Naor
SIAM J. Comput.1
1999 Fault Tolerant Video on Demand Services
abstract
This paper describes a highly available distributed video on demand (VoD) service which is inherently fault tolerant. The VoD service is provided by multiple servers that reside at different sites. New servers may be brought up "on the fly" to alleviate the load on other servers. When a server crashes it is replaced by another server in a transparent way; the clients are unaware of the change of service provider. In test runs of our VoD service prototype, such transitions are not noticeable to a human observer who uses the service. Our VoD service uses a sophisticated flow control mechanism and supports adjustment of the video quality to client capabilities. It does not assume any proprietary network technology: it uses commodity hardware and publicly available network technologies (e.g., TCP/IP, ATM). Our service may run on any machine connected to the Internet. The service exploits a group communication system as a building block for high availability. The utilization of group communication greatly simplifies the service design.
Tal Anker, Danny Dolev, Idit Keidar
ICDCS2
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
PODC3
1998 A Safe and Scalable Payment Infrastructure for Trade of Electronic Content
abstract
BARTER (a Backbone ARchitecture for Trade of ElectRonic content) is a payment infrastructure that facilitates digital content trade over an open network. BARTER is designed to operate over a large-scale, global and heterogeneous communication network. The BARTER protocols address two vital requirements, neglected from existing electronic commerce systems: scalability and transactional efficiency. These protocols possess strong properties such as delivery atomicity, agreement validation and the ability to resolve several classes of disputes. BARTER's novelty is twofold: First, BARTER servers are not required to perform expensive cryptographic operations such as commitment verification; commitments are cross-verified by the parties themselves, thus reducing the overhead of online transaction processing by orders of magnitude. Consequently, BARTER can serve as an efficient online/offline clearing infrastructure. Second, BARTER integrates scalability considerations into several system components (the authentication subsystem, the account management subsystem, and the maintenance of global data) that are likely to suffer service degradation in a world-wide setting. In addressing these issues, BARTER takes into account the inherent asynchronous, unreliable, insecure and failure-prone environment assumptions. We contend that by employing service distribution, BARTER is expected to scale well, meeting the demands of a world-wide setting, over which it is intended to operate.
Gadi Shamir, Michael Ben-Or, Danny Dolev
Int. J. Cooperative Inf. Syst.3
1998 Increasing the Resilience of Distributed and Replicated Database Systems
Idit Keidar, Danny Dolev
J. Comput. Syst. Sci.2
1997 Failure Detectors in Omission Failure Environments
abstract
No abstract available.
Danny Dolev, Roy Friedman 0001, Idit Keidar, Dahlia Malkhi
PODC1
1997 Dynamic Voting for Consistent Primary Components
abstract
Distributed applications often use quorums in order to guarantee consistency. With emerging world-wide communication technology, many new applications (e.g. conferencing applications and interactive games) wish to allow users to freely join and leave, without restarting the entire system. The dynamic voting paradigm allows such systems to define quorums adaptively, accounting for the changes in the set of participants. Furthermore, dynamic voting was proven to be the most available paradigm for maintaining quorums in unreliable networks. However, the subtleties of implementing dynamic voting were not well understood, in fact many of the suggested protocols may lead to inconsistencies in case of failures. Other protocols severely limit the availability in case failures occur during the protocol. In this paper we present a robust and efficient dynamic voting protocol for unreliable asynchronous networks. The protocol consistently maintains the primary component in a distributed system. O...
Esti Yeger Lotem, Idit Keidar, Danny Dolev
PODC3
1997 Efficient Message Passing Interface (MPI) for Parallel Computing on Clusters of Workstations
Jehoshua Bruck, Danny Dolev, C. T. Howard Ho, Marcel-Catalin Rosu, Ray Strong
J. Parallel Distributed Comput.2
1997 Report Dagstuhl Seminar on Time Services, Schloß Dagstuhl, March 11-15, 1996
Danny Dolev, Rüdiger Reischuk, Fred B. Schneider, Ray Strong
Real Time Syst.1
1997 Bounded Concurrent Time-Stamping
abstract
We introduce concurrent time-stamping, a paradigm that allows processes to temporally order concurrent events in an asynchronous shared-memory system. Concurrent time-stamp systems are powerful tools for concurrency control, serving as the basis for solutions to coordination problems such as mutual exclusion, $\ell$-exclusion, randomized consensus, and multiwriter multireader atomic registers. Unfortunately, all previously known methods for implementing concurrent time-stamp systems have been theoretically unsatisfying since they require unbounded-size time-stamps---in other words, unbounded-size memory. This work presents the first bounded implementation of a concurrent time-stamp system, providing a modular unbounded-to-bounded transformation of the simple unbounded solutions to problems such as those mentioned above. It allows solutions to two formerly open problems, the bounded-probabilistic-consensus problem of Abrahamson and the fifo-$\ell$-exclusion problem of Fischer, Lynch, Burns and Borodin, and a more efficient construction of multireader multiwriter atomic registers.
Danny Dolev, Nir Shavit
SIAM J. Comput.1
1996 A Framework for Partitionable Membership Service (Abstract)
abstract
No abstract available.
Danny Dolev, Dahlia Malkhi, Ray Strong
PODC1
1996 Efficient Message Ordering in Dynamic Networks
abstract
We present an algorithm for totally ordering messages in the face of network partitions and site failures.The algorithm aJways aJlows a majority of connected processors in the network to make progress (z.e. to order messages), if they remain connected for sufficiently long, regardless of past failures.Furthermore, our aJgorithm always allows processors to initiate messages, even when they are not members of a connected majority component in the network.Thus, messages can eventually become totally ordered even if their initiator is never a member of a majority component.The algorithm guarantees that when a majority is connected, each message is ordered within two communication rounds, if no failures occur during these rounds. 1 Introduction Consistent order is a powerful paradigm for the design of fault tolerant applications, e.g.consistent replication [Sch90, Kei94].We present an efficient algorithm for consistent message ordering in the face of network partitions and site failures, The network may partition into several components, and remerge.The algorithm is most adequate for dynamic networks where failures are transient.The algorithm uses an underlaying group communication service as a building block.Problem Definition Atomic broadcast deals with consistent message ordering.Informally, atomic broadcast requires that all the correct processors will deliver all the messages to the application in the same order and that they eventually deliver all messages sent by correct processors.In our
Idit Keidar, Danny Dolev
PODC2
1995 Increasing the Resilience of Atomic Commit at No Additional Cost
abstract
This paper presents a new atomic commitment protocol, Enhanced Three Phase Commit (E3PC ), that always allows a quorum in the system to make progress. Previously suggested quorum-based protocols (e.g. the quorum-based Three Phase Commit (3PC) [Ske82]) allow a quorum to make progress in case of one failure. If failures cascade, however, and the quorum in the system is "lost" (i.e. at a given time no quorum component exists, e.g. because of a total crash), a quorum can later become connected and still remain blocked. With our protocol, a connected quorum never blocks. E3PC is based on the quorumbased 3PC [Ske82], and it does not require more time or communication than 3PC. The principles demonstrated in this paper can be used to increase the resilience of a variety of distributed services, e.g. replicated database systems, by ensuring that a quorum will always be able to make progress. 1 Introduction Reliability and availability of loosely coupled distributed database systems is beco...
Idit Keidar, Danny Dolev
PODS2
1995 Efficient Message Passing Interface (MPI) for Parallel Computing on Clusters of Workstations
abstract
Parallel computing on clusters of workstations and personal computers has very high \npotential, since it leverages existing hardware and software. Parallel programming \nenvironments offer the user a convenient way to express parallel computation and communication. \nIn fact, recently, a Message Passing Interface (MPI) has been proposed as an industrial \nstandard for writing "portable" message-passing parallel programs. The communication \npart of MPI consists of the usual point-to-point communication as well as collective \ncommunication. However, existing implementations of programming environments for clusters \nare built on top of a point-to-point communication layer (send and receive) over local \narea networks (LANs) and, as a result, suffer from poor performance in the collective \ncommunication part. \nIn this paper, we present an efficient design and implementation of the collective \ncommunication part in MPI that is optimized for clusters of workstations. Our system consists \nof two main components: the MPI-CCL layer that includes the collective communication \nfunctionality of MPI and a User-level Reliable Transport Protocol (URTP) that interfaces \nwith the LAN Data-link layer and leverages the fact that the LAN is a broadcast medium. \nOur system is integrated with the operating system via an efficient kernel extension \nmechanism that we developed. The kernel extension significantly improves the performance of \nour implementation as it can handle part of the communication overhead without involving \nuser space. \nWe have implemented our system on a collection of IBM RS/6000 workstations con- \nnected via a lOMbit Ethernet LAN. Our performance measurements are taken from typical \nscientific programs that run in a parallel mode by means of the MPI. The hypothesis behind \nour design is that system's performance will be bounded by interactions between the kernel \nand user space rather than by the bandwidth delivered by the LAN Data-Link Layer. Our \nresults indicate that the performance of our MPI Broadcast (on top of Ethernet) is about \ntwice as fast as a recently published software implementation of broadcast on top of ATM.
Jehoshua Bruck, Danny Dolev, C. T. Howard Ho, Marcel-Catalin Rosu, Ray Strong
SPAA2
1995 Atomic Broadcast: From Simple Message Diffusion to Byzantine Agreement
Flaviu Cristian, Houtan Aghili, Ray Strong, Danny Dolev
Inf. Comput.4
1995 Sharing Memory Robustly in Message-Passing Systems
abstract
Emulators that translate algorithms from the shared-memory model to two different message-passing models are presented. Both are achieved by implementing a wait-free, atomic, single-writer multi-reader register in unreliable, asynchronous networks. The two message-passing models considered are a complete network with processor failures and an arbitrary network with dynamic link failures. These results make it possible to view the shared-memory model as a higher-level language for designing algorithms in asynchronous distributed systems. Any wait-free algorithm based on atomic, single-writer multi-reader registers can be automatically emulated in message-passing systems, provided that at least a majority of the processors are not faulty and remain connected. The overhead introduced by these emulations is polynomial in the number of processors in the system. Immediate new results are obtained by applying the emulators to known shared-memory algorithms. These include, among others, protocols to solve the following problems in the message-passing model in the presence of processor or link failures: multi-writer multi-reader registers, concurrent time-stamp systems,l-exclusion, atomic snapshots, randomized consensus, and implementation of data structures.
Hagit Attiya, Amotz Bar-Noy, Danny Dolev
J. ACM3
1995 Dynamic Fault-Tolerant Clock Synchronization
abstract
This paper gives two simple efficient distributed algorithms: one for keeping clocks in a network synchronized and one for allowing new processors to join the network with their clocks synchronized. Assuming a fault-tolerant authentication protocol, the algorithms tolerate both link and processor failures of any type. The algorithm for maintaining synchronization works for arbitrary networks (rather than just completely connected networks) and tolerates any number of processor or communication link faults as long as the correct processors remain connected by fault-free paths. It thus represents an improvement over other clock synchronization algorithms such as those of Lamport and Melliar Smith and Welch and Lynch, although, unlike them, it does require an authentication protocol to handle Byzantine faults. Our algorithm for allowing new processors to join requires that more than half the processors be correct, a requirement that is provably necessary.
Danny Dolev, Joseph Y. Halpern, Barbara B. Simons, Ray Strong
J. ACM1
1994 PCODE: Efficient Parallel Computing over Distributed Environments
Jehoshua Bruck, Danny Dolev, C. T. Howard Ho, Rimon Orni, Ray Strong
PODC2
1994 Observable Clock Synchronization (Extended Abstract)
abstract
While the synchronization of time-o!-day clocks ordinarily requires information f70w in both directions between the clocks, this information need not j70w directlp via messages.However, to take advantage of indirect information fiow, we have to make a number of
Danny Dolev, Rüdiger Reischuk, Ray Strong
PODC1
1994 Experience with RAPID prototypes
abstract
The goals of the RAPID environment are: firstly to make the programming of distributed protocols simple without restricting the protocol relevant choices of the programmer; secondly to provide encapsulation and reusability that are at least as powerful as those offered by object oriented programming; and thirdly to provide for different styles of programming that make RAPID an easy transitional programming environment between older and lower level languages and C. The environment provides and is programmed in the RAPID-FL subset of the functional language FL. Although the full power of FL is available to the programmer, a very small number of concepts need to be learned to program in RAPID-FL. Moreover, restriction to RAPID-FL means that one can have the safety of a functional language combined with reasonable uses of assignment. RAPID makes storage management trivial and reduces the complexity of communication management to handling a few simple commands. We describe our experience using RAPID to perform clock synchronization experiments and to serve as scaffolding for high performance C code that implements a collective communication protocol for parallel machines.>
Danny Dolev, Ray Strong, Ed Wimmers
RSP1
1994 Neighborhood Preserving Hashing and Approximate Queries
Danny Dolev, Yuval Harari, Nathan Linial, Noam Nisan, Michal Parnas
SODA1
1994 A Bounded First-In, First-Enabled Solution to the l-Exclusion Problem
abstract
This article presents a solution to the first-come, first-enabled ℓ-exclusion problem of Fischer et al. [1979]. Unlike their solution, this solution does not use powerful read-modify-write synchronization primitives and requires only bounded shared memory. Use of the concurrent timestamp system of Dolev and Shavir [1989] is key in solving the problem within bounded shared memory.
Yehuda Afek, Danny Dolev, Eli Gafni, Michael Merritt, Nir Shavit
ACM Trans. Program. Lang. Syst.2
1993 On Distributed Algorithms in a Broadcast Domain
Danny Dolev, Dahlia Malkhi
ICALP1
1993 Atomic Snapshots of Shared Memory
abstract
This paper introduces a general formulation of atomic snapshot memory , a shared memory partitioned into words written ( updated ) by individual processes, or instantaneously read ( scanned ) in its entirety. This paper presents three wait-free implementations of atomic snapshot memory. The first implementation in this paper uses unbounded (integer) fields in these registers, and is particularly easy to understand. The second implementation uses bounded registers. Its correctness proof follows the ideas of the unbounded implementation. Both constructions implement a single-writer snapshot memory, in which each word may be updated by only one process, from single-writer, n -reader registers. The third algorithm implements a multi-writer snapshot memory from atomic n -writer, n -reader registers, again echoing key ideas from the earlier constructions. All operations require Θ( n 2 ) reads and writes to the component shared registers in the worst case. — Authors' Abstract
Yehuda Afek, Hagit Attiya, Danny Dolev, Eli Gafni, Michael Merritt, Nir Shavit
J. ACM3
1993 Perfectly Secure Message Transmission
abstract
This paper studies the problem of perfectly secure communication in general network in which processors and communication lines may be faulty. Lower bounds are obtained on the connectivity required for successful secure communication. Efficient algorithms are obtained that operate with this connectivity and rely on no complexity-theoretic assumptions. These are the first algorithms for secure communication in a general network to simultaneously achieve the three goals of perfect secrecy, perfect resiliency, and worst-case time linear in the diameter of the network.
Danny Dolev, Cynthia Dwork, Orli Waarts, Moti Yung
J. ACM1
1993 A Partial Equivalence Between Shared-Memory and Message-Passing in an Asynchronous Fail-Stop Distributed Environment
Amotz Bar-Noy, Danny Dolev
Math. Syst. Theory2
1992 Shifting Gears: Changing Algorithms on the Fly to Expedite Byzantine Agreement
Amotz Bar-Noy, Danny Dolev, Cynthia Dwork, Ray Strong
Inf. Comput.2
1992 Determinism vs. Nondeterminism in Multiparty Communication Complexity
abstract
A given Boolean function has its input distributed among many parties. The aim is to determine which parties to talk to and what information to exchange in order to evaluate the function while minimizing the total communication. This paper shows that it is possible to evaluate the Boolean function deterministically with only a polynomial increase in communication and number of parties accessed with respect to the information lower bound given by the nondeterministic communication complexity of the function.
Danny Dolev, Tomás Feder
SIAM J. Comput.1
1991 Non-Malleable Cryptography (Extended Abstract)
abstract
The notion of non-malleable cryptography, an extension of semantically secure cryptography, is defined. Informally, the additional requirement is that given the ciphertext it is impossible to generate a different ciphertext so that the respective plaintexts are related. The same concept makes sense in the contexts of string commitment and zero-knowledge proofs of possession of knowledge. Non-malleable schemes for each of these three problems are presented. The schemes do not assume a trusted center; a user need not know anything about the number or identity of other system users. Keywords: cryptography, cryptanalysis, randomized algorithms, nonmalleability AMS subject classifications: 68M10, 68Q20, 68Q22, 68R05, 68R10 A preliminary version of this work appeared in STOC '91 Hebrew University Jerusalem, Israel y IBM Research Division, Almaden Research Center, 650 Harry Road, San Jose, CA 95120. E-mail: [email protected]. z Incumbent of the Morris and Rose Goldman Career Devel...
Danny Dolev, Cynthia Dwork, Moni Naor
STOC1
1991 Consensus Algorithms with One-Bit Messages
Amotz Bar-Noy, Danny Dolev
Distributed Comput.2
1991 Fault-Tolerant Critical Section Management in Asynchronous Environments
Amotz Bar-Noy, Danny Dolev, Daphne Koller, David Peleg
Inf. Comput.2
1990 Perfectly Secure Message Transmission
abstract
The problem of perfectly secure communication in a general network in which processors and communication lines may be faulty is studied. Lower bounds are obtained on the connectivity required for successful secure communication. Efficient algorithms that operate with this connectivity and rely on no complexity theoretic assumptions are derived. These are the first algorithms for secure communication in a general network to achieve simultaneously the goals of perfect secrecy, perfect resiliency, and a worst case time which is linear in the diameter of the network.>
Danny Dolev, Cynthia Dwork, Orli Waarts, Moti Yung
FOCS1
1990 Atomic Snapshots of Shared Memory
abstract
An atomic snapshot memory is a shared data structure allowing concurrent processes to store information in a collection of shared registers, all of which may be read in a single atomic scan operation.This paper presents three wait-free implementations of atomic snapshot memory.Two constructions implement wait-free single-writer atomic snapshot memory from wait-free atomic single-writer, n-reader registers.A third construction implements a wait-free n-writer atomic snapshot memory from n-writer, n-reader registers.The first implementation uses unbounded
Yehuda Afek, Danny Dolev, Hagit Attiya, Eli Gafni, Michael Merritt, Nir Shavit
PODC2
1990 Sharing Memory Robustly in Message-Passing Systems
abstract
Emulators that translate algorithms from the sharedmemory model to two different message-passing models are presented.Both are achieved by implementing a wait-free, atomic, single-writer multi-reader register in unreliable, asynchronous networks.The two message-passing models considered are a complete network with processor failures and an arbitrary network with dynamic link failures.These results make it possible to view the sharedmemory model as a higher-level language for designing algorithms in asynchronous distributed systems.Any wait-free algorithm based on atomic, single-writer multi-reader registers can be automatically emulated in message-passing systems.The overhead introduced by these emulations is polynomial in the number of processors in the systems.Immediate new results are obtained by applying the emulators to known shared-memory algorithms.
Hagit Attiya, Amotz Bar-Noy, Danny Dolev
PODC3
1990 New Latency Bounds for Atomic Broadcast
abstract
Tighter bounds are provided on the time required to reach agreement in a distributed system as a function of the failure model. After describing the model of a distributed system that is a context for this work the authors define several failure classes. They define a partial order on classes of failures that involves whether there is a latency penalty in converting from tolerance of one failure class to another. In this setting they distinguish clock and timing failures, showing that there can be a penalty in converting from timing failure tolerance to clock failure tolerance. The authors leave open the exact expression for the optimal latency for timing failure tolerant atomic broadcast, though it is conjectured that there is some penalty in converting from omission failure tolerance to timing failure tolerance.>
Ray Strong, Danny Dolev, Flaviu Cristian
RTSS2
1990 Renaming in an Asynchronous Environment
abstract
This paper is concerned with the solvability of the problem of processor renaming in unreliable, completely asynchronous distributed systems. Fischer et al. prove in [8] that “nontrivial consensus” cannot be attained in such systems, even when only a single, benign processor failure is possible. In contrast, this paper shows that problems of processor renaming can be solved even in the presence of up tot
Hagit Attiya, Amotz Bar-Noy, Danny Dolev, David Peleg, Rüdiger Reischuk
J. ACM3
1990 Early Stopping in Byzantine Agreement
abstract
Two different kinds of Byzantine Agreement for distributed systems with processor faults are defined and compared. The first is required when coordinated actions may be performed by each participant at different times. This kind is called Simultaneous Byzantine Agreement (SBA). This paper deals with the number of rounds of message exchange required to reach Byzantine Agreement of either kind (BA). If an algorithm allows its participants to reach Byzantine agreement in every execution in which at most t participants are faulty, then the algorithm is said to tolerate t faults. It is well known that any BA algorithm that tolerates t faults (with t < n - 1 where n denotes the total number of processors) must run at least t + 1 rounds in some execution. However, it might be supposed that in executions where the number f of actual faults is small compared to t , the number of rounds could be correspondingly small. A corollary of our first result states that (when t < n - 1) any algorithm for SBA must run t + 1 rounds in some execution where there are no faults. For EBA (with t < n - 1), a lower bound of min( t + 1, f + 2) rounds is proved. Finally, an algorithm for EBA is presented that achieves the lower bound, provided that t is on the order of the square root of the total number of processors.
Danny Dolev, Rüdiger Reischuk, Ray Strong
J. ACM1
1989 Multiparty Communication Complexity
abstract
A given Boolean function has its input distributed among many parties. The aim is to determine which parties to talk to and what information to exchange with each of them in order to evaluate the function while minimizing the total communication. It is shown that it is possible to obtain the Boolean answer deterministically with only a polynomial increase in communication with respect to the information lower bound given by the nondeterministic communication complexity of the function.>
Danny Dolev, Tomás Feder
FOCS1
1989 Bounded Polynomial Randomized Consensus
abstract
In [A&3], Abrahamson presented a solution to the randomized consensus problem of Chor, Israeli and Li [CIL87], without assuming the existence of an atomic coin flip operation.This elegant algorithm uses unbounded memory, and has expected exponential running time.In [AH89], Aspens and Herlihy provide a breakthrough polynomial-time algorithm.However, it too is based on the use of unbounded memory.In this paper, we present a solution to the randomized consensus problem, that is bounded in space and runs in polynomial expected time.
Hagit Attiya, Danny Dolev, Nir Shavit
PODC2
1989 Shared-Memory vs. Message-Passing in an Asynchronous Distributed Environment
abstract
No abstract available.
Amotz Bar-Noy, Danny Dolev
PODC2
1989 Bounded Concurrent Time-Stamp Systems Are Constructible
abstract
Concurrent time stamping is at the heart of solutions to some of the most fundamental problems in distributed computing. Based on concurrent-time-stamp-systems, elegant and simple solutions to core problems such as ƒcƒs-mutual-exclusion, construction of a multi-reader-multi-writer atomic register, probabilistic consensus,… were developed. Unfortunately, the only known implementation of a concurrent time stamp system has been theoretically unsatisfying, since it requires unbounded size time-stamps, in other words, unbounded memory. Not knowing if bounded concurrent-time-stamp-systems are at all constructible, researchers were led to constructing complicated problem-specific solutions to replace the simple unbounded ones. In this work, for the first time, a bounded implementation of a concurrent-time-stamp-system is presented. It provides a modular unbounded-to-bounded transformation of the simple unbounded solutions to problems such as above. It allows solutions to two formerly open problems, the bounded-probabilistic-consensus problem of Abrahamson [A88] and the ƒiƒo[email protected]@@@-exclusion problem of [FLBB85], and a more efficient construction of mrmw atomic registers.
Danny Dolev, Nir Shavit
STOC1
1989 Choice Coordination with Limited Failure
Amotz Bar-Noy, Michael Ben-Or, Danny Dolev
Distributed Comput.3
1989 The Distributed Firing Squad Problem
abstract
The distributed firing squad problem is defined in the context of a synchronous distributed system where the correct processors operate in lock-step synchrony but do not share a global clock. If one or more correct processors receive a command to start a firing squad synchronization, then at some future time all correct processors must “fire” (formally, enter a special state) at exactly the same step. For various fault models, upper and lower bounds are proved on the number of faulty processors that can be tolerated and on the number of rounds of communication required between the reception of the start command and firing. For example, if a firing squad protocol is resilient to t fail-stop faults, then at least $t + 1$ rounds are necessary and sufficient. For the case of Byzantine faults with authentication where the faulty processors can take steps in between the synchronous steps of the correct processors, the firing squad problem can be solved in $t + 5$ rounds, provided that $n > 3t$, where n is the number of processors and t is the number of faults, and the problem cannot be solved at all if $n \leqq 3t$. Moreover, in the case that $n \leqq 3t$, the impossibility of a firing squad protocol holds even for a weaker “timing fault model” where all processors generate messages correctly according to the protocol, but the faulty processors can affect the system by slightly slowing down or speeding up messages.
Brian A. Coan, Danny Dolev, Cynthia Dwork, Larry J. Stockmeyer
SIAM J. Comput.2
1988 Robust multi-agent decision making in faulty environment
abstract
A brief overview is given of current state of the art in robust multiagent decision making in a faulty environment (distributed consensus); and also a family of novel algorithms is presented to achieve consensus that are deterministic, simple, require single-bit messages, and can be implemented in hardware. For such systems to work properly, the issues of reaching common decision (consensus) in the presence of faults have to be addressed. The authors offer a structural solution to this problem. The approach is illustrated on a hypothetical example of a number of autonomous robots in manufacturing that all have to agree on a common decision determined by the values of some central controllers.>
Amotz Bar-Noy, Danny Dolev, Dragutin Petkovic
ICPR2
1988 Toward a Non-Atomic Era: \ell-Exclusion as a Test Case
abstract
Most of the research in concurrency control has been based on the existence of strong synchronization primitives such as test and set. Following Lamport, recent research promoting the use of weaker primitives, “safe” rather than “atomic,” has resulted in construction of atomic registers from safe ones, in the belief that they would be useful tools for process synchronization. We argue that the properties provided by atomic operations may be too powerful, masking core difficulties of problems and leading to inefficiency. We therefore advocate a different approach, to skip the intermediate step of achieving atomicity, and solve problems directly from safe registers. Though it has been shown that “test and set” cannot be implemented from safe registers, we show how to achieve a fair solution to l-exclusion, a classical concurrency control problem previously solved assuming a very powerful form of atomic “test and set”. We do so using safe registers alone and without introducing atomicity. The solution is based on the construction of a simple novel non-atomic synchronization primitive.
Danny Dolev, Eli Gafni, Nir Shavit
STOC1
1988 Some Geometry for General River Routing
abstract
Efficient solutions are given to compute the optimal placement for a pair of VLSI modules interconnected by river routing. Specifically, let the (perpendicular) distance between the two modules be the separation, and call the (transverse) displacement the offset. This paper principally considers the separation problem: Given an offset and a wiring rule, find the minimum separation permitting a legal wiring. The design rules might use wires which are exclusively rectilinear, polygonal with a finite number of slopes, or possibly restricted to some other class of shapes such as circular arcs plus linear pieces. Techniques are developed which unify a variety of different placement problems, and give efficient solutions under extremely general conditions. The advantage of these generalizations is not only their theoretical framework; the results extend naturally to more precise models of real river routing, and the theory is applicable to placement problems for collections of modules.
Alan R. Siegel, Danny Dolev
SIAM J. Comput.2
1987 Achievable Cases in an Asynchronous Environment (Extended Abstract)
abstract
The paper deals with achievability of fault tolerant goals in a completely asynchronous distributed system. Fischer, Lynch, and Paterson [FLP] proved that in such a system "nontrivial agreement" cannot be achieved even in the (possible) presence of a single "benign" fault. In contrast, we exhibit two pairs of goals that are achievable even in the presence of up to t ≪ n/2 faulty processors, contradicting the widely held assumption that no nontrivial goals are attainable in such a system. The first pair deals with renaming processors so as to reduce the size of the initial name space. When only uniqueness is required of the new names, we present a lower bound of n + 1 on the size of the new name space, and a renaming algorithm which establishes an upper bound of n + t. In case the new names are required also to preserve the original order, a tight bound of 2t(n- t + 1) - 1 is obtained. The second pair of goals deals with the multi-slot critical section problem. We present algorithms for controlled access to a critical section. As for the number of slots required, a tight bound of t + 1 is proved in case the slots are identical. In the case of distinct slots the upper bound is 2t + 1.
Hagit Attiya, Amotz Bar-Noy, Danny Dolev, Daphne Koller, David Peleg, Rüdiger Reischuk
FOCS3
1987 Shifting Gears: Changing Algorithms on the Fly To Expedite Byzantine Agreement
abstract
All in-text\treferences\tunderlined\tin\tblue\tare\tlinked\tto\tpublications\ton\tResearchGate, letting you\taccess\tand\tread\tthem\timmediately.
Amotz Bar-Noy, Danny Dolev, Cynthia Dwork, Ray Strong
PODC2
1987 Efficient Fault-Tolerant Routings in Networks
Andrei Z. Broder, Danny Dolev, Michael J. Fischer, Barbara B. Simons
Inf. Comput.2
1987 A New Look at Fault-Tolerant Network Routing
Danny Dolev, Joseph Y. Halpern, Barbara B. Simons, Ray Strong
Inf. Comput.1
1987 On the minimal synchronism needed for distributed consensus
abstract
Reaching agreement is a primitive of distributed computing. Whereas this poses no problem in an ideal, failure-free environment, it imposes certain constraints on the capabilities of an actual system: A system is viable only if it permits the existence of consensus protocols tolerant to some number of failures. Fischer et al. have shown that in a completely asynchronous model, even one failure cannot be tolerated. In this paper their work is extended: Several critical system parameters, including various synchrony conditions, are identified and how varying these affects the number of faults that can be tolerated is examined. The proofs expose general heuristic principles that explain why consensus is possible in certain models but not possible in others.
Danny Dolev, Cynthia Dwork, Larry J. Stockmeyer
J. ACM1
1986 Cheating Husbands and other Stories: A Case Study of Knowledge, Action, and Communication
Yoram Moses, Danny Dolev, Joseph Y. Halpern
Distributed Comput.2
1986 Reaching approximate agreement in the presence of faults
abstract
This paper considers a variant of the Byzantine Generals problem, in which processes start with arbitrary real values rather than Boolean values or values from some bounded range, and in which approximate, rather than exact, agreement is the desired goal. Algorithms are presented to reach approximate agreement in asynchronous, as well as synchronous systems. The asynchronous agreement algorithm is an interesting contrast to a result of Fischer et al, who show that exact agreement with guaranteed termination is not attainable in an asynchronous system with as few as one faulty process. The algorithms work by successive approximation, with a provable convergence rate that depends on the ratio between the number of faulty processes and the total number of processes. Lower bounds on the convergence rate for algorithms of this form are proved, and the algorithms presented are shown to be optimal.
Danny Dolev, Nancy A. Lynch, Shlomit S. Pinter, Eugene W. Stark, William E. Weihl
J. ACM1
1986 On the Possibility and Impossibility of Achieving Clock Synchronization
Danny Dolev, Joseph Y. Halpern, Ray Strong
J. Comput. Syst. Sci.1
1986 The Parallel Complexity of Scheduling with Precedence Constraints
Danny Dolev, Eli Upfal, Manfred K. Warmuth
J. Parallel Distributed Comput.1
1986 Bounds for Width Two Branching Programs
abstract
Branching programs have been studied as a fundamental model for space bounded computations and, in particular, as a model in which to try to establish nontrivial space lower bounds and time-space trade-offs. At present, there still do not exist any results for single output functions. We consider a class of severely constrained programs (those having width 2) and establish characterizations as well as lower bounds for some Boolean functions computable within this model.
Allan Borodin, Danny Dolev, Faith Ellen, Wolfgang J. Paul
SIAM J. Comput.2
1985 Choice Coordination with Bounded Failure (a Preliminary Version)
abstract
No abstract available.
Amotz Bar-Noy, Michael Ben-Or, Danny Dolev
PODC3
1985 Cheating Husbands and Other Stories: A Case Study of Knowledge, Action, and Communication (Preliminary Version)
abstract
By looking at a number of variants of the clteatir~g Itusbcads puzzle, we illustrate the subtle relationship between knowledge, communication, and action in a distributed environment.
Yoram Moses, Danny Dolev, Joseph Y. Halpern
PODC2
1985 The Distributed Firing Squad Problem (Preliminary Version)
abstract
this paper we justify the design assumption of simultaneous starts. Specifically, we provide algorithms to solve the associated synchronization problem, which we call the distributed firing squad problem (abbreviated DFS). A distributed algorithm for the DFS problem has two properties: (I) if any correct processor receives a .message to start a DFS synchronization, then at some future time all cor- rect processors will "fire" (formally, enter a special state), and (2) the correct processors all fire at exactly the same step
Brian A. Coan, Danny Dolev, Cynthia Dwork, Larry J. Stockmeyer
STOC2
1985 Bounds on Information Exchange for Byzantine Agreement
abstract
Byzantine Agreement has become increasingly important in establishing distributed properties when errors may exist in the systems. Recent polynomial algorithms for reaching Byzantine Agreement provide us with feasible solutions for obtaining coordination and synchronization in distributed systems. In this paper the amount of information exchange necessary to ensure Byzantine Agreement is studied. This is measured by the total number of messages the participating processors have to send in the worst case. In algorithms that use a signature scheme, the number of signatures appended to messages are also counted. First it is shown that Ω( nt ) is a lower bound for the number of signatures for any algorithm using authentication, where n denotes the number of processors and t the upper bound on the number of faults the algorithm is supposed to handle. For algorithms that reach Byzantine Agreement without using authentication this is even a lower bound for the total number of messages. If n is large compared to t , these bounds match the upper bounds from previously known algorithms. For the number of messages in the authenticated case we prove the lower bound Ω( n + t 2 ). Finally algorithms that achieve this bound are presented.
Danny Dolev, Rüdiger Reischuk
J. ACM1
1985 Scheduling Flat Graphs
abstract
The problem of scheduling a partially ordered set of unit length tasks on m identical processors is known to be NP-complete. There are efficient algorithms for only a few special cases of this problem. In this paper we analyze the effect of the structure of the precedence graph and the availability of the processors on the construction of optimal schedules. We prove that to find an optimal schedule it suffices to consider at each step only initial tasks which belong to the $m - 1$ highest components of the precedence graph. This result reduces the number of cases we have to check during the construction of an optimal schedule. Our method leads to polynomial algorithms if the number of processors is fixed and the precedence graph has a certain form. In particular, if the precedence graph contains only intrees and outtrees, this result leads to linear algorithms for finding an optimal schedule on two or three processors.
Danny Dolev, Manfred K. Warmuth
SIAM J. Comput.1
1984 Flipping coins in many pockets (Byzantine agreement on uniformly random values)
abstract
It was recently shown by Michael Rabin that a sequence of random 0-1 values, prepared and distributed by a trusted "dealer," can be used to achieve Byzantine agreement in constant expected time in a network of processors. A natural question is whether it is possible to generate these values uniformly at random within the network. In this paper we present a cryptography based protocol for agreernent on a 0-1 randona value, if less than half of the processors are faulty. In fact the protocol allows uniform sampling from any finite set, and thus solves the problem of choosing a network leader uniformly at random. The protocol is usable both when all the communication is via "broadcast," in which case it needs three rounds of information exchange, and when each pair of processors communicate on a private line, in which case it needs 3t + 3 rounds, where t is the number of faulty proccssors. The protocol remains valid even if passive eavesdropping is allowed. On the other hand we show that no (probabilistic) protocol can achieve agreement on a fair coin in fewer phases then necessary for Byzantine agreement, and hence the "pre-dealt" nature of the random sequence required for Rabin's algorithm is crucial.
Andrei Z. Broder, Danny Dolev
FOCS2
1984 Asynchronous Byzantine Consensus
abstract
Reaching agreement in an asynchronous environment is essential to guarantee consistency in distributed data processing. All previous asynchronous protocols were either probabilistic or they assumed a fail-stop mode of failure. The deterministic protocol presented in this paper reaches a Strong Byzantine Agreement in a system of asynchronous processors; and therefore can sustain arbitrary faults. In our model, processors can be completely asynchronous, though the communication network has the property that a message being sent by a correctly operating processor to a set of processors will reach its destinations within a predetermined period Δ. Additional results presented in the paper prove that in the above model one cannot reach a consensus within a bounded time. A correctly operating processor should wait to receive messages from other processors before making a decision. This result holds also for Weak Byzantine Agreement, but not for nontrivial consensus. We present a trivial protocol to reach a nontrivial consensus in bounded time.
Hagit Attiya, Danny Dolev, Joseph Gil
PODC2
1984 Fault-Tolerant Clock Synchronization
abstract
This paper gives two simple efficient distributed algorithms: one for keeping clocks in a network synchronized and one for allowing new processors to join the network with their clocks synchronized. The algorithms tolerate both link and node failures of any type. The algorithm for maintaining synchronization will work for arbitrary networks (rather than just completely connected networks) and tolerates any number of processor or communication link faults as long as the correct processors remain connected by fault-free paths. It thus represents an improvement over other clock synchronization algorithms such as [LM1,LM2,LL1]. Our algorithm for allowing new processors to join requires that more than half the processors be correct, a requirement which is provably necessary.
Joseph Y. Halpern, Barbara B. Simons, Ray Strong, Danny Dolev
PODC4
1984 Efficient Fault Tolerant Routings in Networks
abstract
We analyze the problem of constructing a network which will have a fixed routing and which will be highly fault tolerant. A construction is presented which forms a “product route graph” from two or more constituent “route graphs.” The analysis involves the surviving route graph, which consists of all non-faulty nodes in the network with two nodes being connected by a directed edge iff the route from the first to the second is still intact after a set of component failures. The diameter of the surviving route graph, that is, the maximum distance between any pair of nodes, is a measure of the worst-case performance degradation caused by the faults. The number of faults tolerated, the diameter, and the degree of the product graph are related in a simple way to the corresponding parameters of the constituent graphs. In addition, there is a “padding theorem” which allows one to add nodes to a graph and to extend a previous routing.
Andrei Z. Broder, Danny Dolev, Michael J. Fischer, Barbara B. Simons
STOC2
1984 On the Possibility and Impossibility of Achieving Clock Synchronization
abstract
It is known that clock synchronization can be achieved in the presence of faulty clocks numbering more than one-third of the total number of participating clocks provided that some authentication technique is used. Without authentication the number of faults that can be tolerated has been an open question. Here we show that if we restrict logical clocks to running within some linear function of real time, then clock synchronization is impossible, without authentication, when one-third or more of the processors are faulty. However, if there is a bound on the rate at which a processor can generate messages, then we show that clock synchronization is achievable, without authentication, as long as the faults do not disconnect the network. Finally, we provide a lower bound on the closeness to which simultaneity can be achieved in the network as a function of the transmission and processing delay properties of the network.
Danny Dolev, Joseph Y. Halpern, Ray Strong
STOC1
1984 A New Look at Fault Tolerant Network Routing
abstract
Consider a communication network G in which a limited number of link and/or node faults F might occur. A routing ρ for the network (a fixed path between each pair of nodes) must be chosen without any knowledge of which components might become faulty. Choosing a good routing corresponds to bounding the diameter of the surviving route graph R(G,ρ)/F, where two nonfaulty nodes are joined by an edge if there are no faults on the route between them. We prove a number of results concerning the diameter of surviving route graphs. We show that if ρ is a minimal length routing, then the diameter of R(G,ρ)/F can be on the order of the number of nodes of G, even if F consists of only a single node. However, if G is the n-dimensional cube, the diameter of R(G,ρ)/F≤3 for any minimal length routing ρ and any set of faults F with |F|
Danny Dolev, Joseph Y. Halpern, Barbara B. Simons, Ray Strong
STOC1
1984 Correcting Faults in Write-Once Memory
abstract
Article Correcting faults in write-once memory Share on Authors: Danny Dolev View Profile , David Maier View Profile , Ilarry Mairson View Profile , Jeffrey Ullman View Profile Authors Info & Claims STOC '84: Proceedings of the sixteenth annual ACM symposium on Theory of computingDecember 1984 Pages 225–229https://doi.org/10.1145/800057.808685Online:01 December 1984Publication History 2citation230DownloadsMetricsTotal Citations2Total Downloads230Last 12 Months2Last 6 weeks0 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteGet Access
Danny Dolev, David Maier 0001, Harry G. Mairson, Jeffrey D. Ullman
STOC1
1983 On the Minimal Synchronism Needed for Distributed Consensus
abstract
Reaching agreement is a primitive of distributed computing. While this poses no problem in an ideal, failure-free environment, it imposes certain constraints on the capabilities of an actual system: a system is viable only if it permits the existence of consensus protocols tolerant to some number of failures. Fischer, Lynch and Paterson [FLP] have shown that in a completely asynchronous model, even one failure cannot be tolerated. In this paper we extend their work, identifying several critical system parameters, including various synchronicity conditions, and examine how varying these affects the number of faults which can be tolerated. Our proofs expose general heuristic principles that explain why consensus is possible in certain models but not possible in others.
Danny Dolev, Cynthia Dwork, Larry J. Stockmeyer
FOCS1
1983 Bounds for Width Two Branching Programs
abstract
Branching programs for the computation of Boolean functions were first studied in the Master's thesis of Masek.7 In a rather straightforward manner they generalize the concept of a decision tree to a decision graph.
Allan Borodin, Danny Dolev, Faith Ellen, Wolfgang J. Paul
STOC2
1983 Superconcentrators, Generalizers and Generalized Connectors with Limited Depth (Preliminary Version)
abstract
We show that the minimum possible size of an n-superconcentrator with depth 2k≥4 is θ(nλ(k, n)), where λ(k, .) is the inverse of a certain function at the k-th level of the primitive recursive hierarchy. It follows that the minimum possible depth of an n-superconcentrator with linear size is θ(β(n)), where β is the inverse of a function growing more rapidly than any primitive recursive function. Similar results hold for generalizers. We give a simple explicit construction for a (d1...dk)-generalizer with depth k and size (d1+...+dk)d1...dk. This is applied to give a simple explicit construction for a generalized n-connector with depth 2k−3 and size (2d1+3d2+...+3dk−1+2dk) d1...dk. These are the best explicit constructions currently available. We also show that, for each fixed k≥2, the minimum possible size of a generalized n-connector with depth k is Ω(n1+1/k) and 0((n log n)1+1/k).
Danny Dolev, Cynthia Dwork, Nicholas Pippenger, Avi Wigderson
STOC1
1983 Authenticated Algorithms for Byzantine Agreement
abstract
Reaching agreement in a distributed system in the presence of faulty processors is a central issue for reliable computer systems. Using an authentication protocol, one can limit the undetected behavior of faulty processors to a simple failure to relay messages to all intended targets. In this paper we show that, in spite of such an ability to limit faulty behavior, and no matter what message types or protocols are allowed, reaching (Byzantine) agreement requires at least $t + 1$ phases or rounds of information exchange, where t is an upper bound on the number of faulty processors. We present algorithms for reaching agreement based on authentication that require a total number of messages sent by correctly operating processors that is polynomial in both t and the number of processors, n. The best algorithm uses only $t + 1$ phases and $O(nt)$ messages.
Danny Dolev, Ray Strong
SIAM J. Comput.1
1983 On the security of public key protocols
abstract
Recently the use of public key encryption to provide secure network communication has received considerable attention. Such public key systems are usually effective against passive eavesdroppers, who merely tap the lines and try to decipher the message. It has been pointed out, however, that an improperly designed protocol could be vulnerable to an active saboteur, one who may impersonate another user or alter the message being transmitted. Several models are formulated in which the security of protocols can be discussed precisely. Algorithms and characterizations that can be used to determine protocol security in these models are given.
Danny Dolev, Andrew Chi-Chih Yao
IEEE Trans. Inf. Theory1
1982 On the Security of Ping-Pong Protocols
Danny Dolev, Shimon Even, Richard M. Karp
CRYPTO1
1982 On the Security of Multi-Party Protocols in Distributed Systems
Danny Dolev, Avi Wigderson
CRYPTO1
1982 'Eventual' Is Earlier than 'Immediate'
abstract
Two different notions of Byzantine Agreement - immediate and eventually - are defined depending on whether the agreement involves an action to be performed synchronously or not. The lower bounds for time complexity depend on what kind of agreement has to be achieved. All previous algorithms to reach Byzantine Agreement ensure immediate agreement. We present two algorithms that in many cases reach the second type of agreement faster than previously known algorithms showing that there actually is a difference between the two notions: Eventual Byzantine Agreement can be reached earlier than Immediate.
Danny Dolev, Rüdiger Reischuk, Ray Strong
FOCS1
1982 Finding Safe Paths in a Faulty Environment
abstract
This paper addresses the problem of finding safe paths through a network, some of whose nodes may be faulty. By a safe path we mean one between two nodes that does not contain any faulty node. The kinds of faults that concern us are not limited to those that may cause a failure of a node or link, but include those that may cause a node to distort messages in arbitrary ways. Furthermore, we want a distributed algorithm to allow the network itself to discover suitable paths without depending on a central controller for the analysis. More broadly, we assume that each node has only local knowledge of the network structure.
Danny Dolev, José Meseguer 0001, Marshall C. Pease
PODC1
1982 Bounds on Information Exchange for Byzantine Agreement
abstract
Byzantine Agreement has become increasingly important in establishing distributed properties when there may exist errors in the systems. Recent polynomial algorithms for reaching Byzantine Agreement provide us with feasible solutions for obtaining coordination and synchronization in distributed systems. In this paper we study the amount of information exchange necessary to ensure Byzantine Agreement. This is measured by the number of messages and the number of signatures appended to messages (in case of authenticated algorithms) the participating processors need to send, in the worse case, in order to reach Byzantine Agreement. The lower bound for the number of signatures in the authenticated case is Ω(nt), where n is the number of participating processors and t is the upper bound on the number of faults. If n is large compared to t, it matches the upper bounds from previously known algorithms. The lower bound for the number of messages is Ω(n+t2). We present an algorithm that achieves this bound and for which the number of phases does not exceed the minimum t+1 by more than a constant factor.
Danny Dolev, Rüdiger Reischuk
PODC1
1982 Polynomial Algorithms for Multiple Processor Agreement
abstract
Reaching agreement in a distributed system while handling malfunctioning behavior is a central issue for reliable computer systems. All previous algorithms for reaching the agreement required an exponential number of messages to be sent, with or without authentication. We give polynomial algorithms for reaching (Byzantine) agreement, both with and without the use of authentication protocols. We also prove that no matter what kind of information is exchanged, there is no way to reach agreement with fewer than t+1 rounds of exchange, where t is the upper bound on the number of faults.
Danny Dolev, Ray Strong
STOC1
1982 On the Security of Ping-Pong Protocols
Danny Dolev, Shimon Even, Richard M. Karp
Inf. Control.1
1982 An Efficient Algorithm for Byzantine Agreement without Authentication
Danny Dolev, Michael J. Fischer, Robert J. Fowler, Nancy A. Lynch, Ray Strong
Inf. Control.1
1981 Unanimity in an Unknown and Unreliable Environment
abstract
Can unanimity be achieved in an unknown and unreliable distributed system? We analyze two extreme models of networks: one in which all the routes of communication are known, and the other in which not even the topology of the network is known. We prove that independently of the model, unanimity is achievable if and only if the number of faulty processors in the system is 1. less than one half of the connectivity of the system's network, and 2. less than one third of the total number of processors. In cases where unanimity is achievable, an algorithm to obtain it is given.
Danny Dolev
FOCS1
1981 On the Security of Public Key Protocols (Extended Abstract)
abstract
Recently the use of public key encryption to provide secure network communication has received considerable attention. Such public key systems are usually effective against passive eavesdroppers, who merely tap the lines and try to decipher the message. It has been pointed out, however, that an improperly designed protocol could be vulnerable to an active saboteur, one who may impersonate another user or alter the message being transmitted. Several models are formulated in which the security of protocols can be discussed precisely. Algorithms and characteri-zations that can be used to determine protocol security in these models are given.
Danny Dolev, Andrew Chi-Chih Yao
FOCS1
1981 Optimal Wiring between Rectangles
abstract
We consider the problem of wiring together two parallel rows of points under a variety of conditions. The options include whether we allow the rows to slide relative to one another, whether we use only rectilinear wires or arbitrary wires, and whether we can use wires in one layer or several layers. In almost all of these combinations of conditions, we can provide a polynomial-time algorithm to minimize the distance between the parallel rows of points. We also compare two fundamentally different wiring approaches, where one and two layers are used. We show that although the theoretical model implies that there can be great gains for the two-layer strategy, even in cases where no crossovers are required, when we consider typical design rules for laying out VLSI circuits there is no substantial advantage to the two-layer approach over the one-layer approach.
Danny Dolev, Kevin Karplus, Alan R. Siegel, Alex Strong, Jeffrey D. Ullman
STOC1
1979 Commutation Preperties and Generating Sets Characterize Slices of Various Synchronization Primitives
Danny Dolev
Theor. Comput. Sci.1
1978 Commutation Relations of Slices Characterize Some Synchronization Primitives
Danny Dolev, Eli Shamir 0001
Inf. Process. Lett.1