Alexander A. Schwarzmann

dblp:s/AlexanderAShvartsman · also Alexander A. Shvartsman · DBLP profile ↗
← Back
123ranked-venue papers
6as first author
5since 2021 · last 2022
0000-0003-4447-3267ORCID · verified

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

Systems, architecture and hardware · 51 · 3 since 2021Theory of computation · 25 · 4 first-authorArtificial intelligence and machine learning · 6 · 1 since 2021Security and privacy · 6 · 1 first-author · 1 since 2021Software engineering, systems software and programming languages · 5Databases, data management, data science and information retrieval · 3 · 2 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 3 · 1 since 2021Computer networks · 1 · 1 first-author
YearPublicationVenuePosition
2022 Triggerability of Backdoor Attacks in Multi-Source Transfer Learning-based Intrusion Detection
abstract
Network-based Intrusion Detection Systems (NIDSs) automate monitoring of events in networks and analyze them for signatures of cyberattacks. With the advancement of machine learning algorithms, more organizations started using machine learning based IDSs (ML-IDSs) to identify and mitigate cyberattacks. However, the lack of training datasets is a major challenge when implementing ML-IDSs. Therefore, using training data from external sources or transfer learning models are some solutions to overcome this challenge. However, using training data from external sources introduces the risk of backdoored datasets, specifically, when the adversaries also have background knowledge on data sources inside the target organization. This work investigates the role of backdoor attacks on intrusion detection techniques trained using multi-source data. The backdoor examples are injected into one or more training data sources. Transfer learning models are then created by projecting data from different sources into a new subspace containing all source data. The backdoor is then triggered in the target data. An anomaly-based intrusion detection classifier is applied to examine the effectiveness of the introduced backdoors. The results have shown that backdoor attacks on multis-source transfer learning models are feasible, although having less impact compared to backdoors on traditional machine learning models.
Nour Alhussien, Ahmed Aleroud, Reza Rahaeimehr, Alexander A. Schwarzmann
BDCAT4
2022 2022 Edsger W. Dijkstra Prize in Distributed Computing
abstract
The Edsger W. Dijkstra Prize in Distributed Computing is awarded for outstanding papers on the principles of distributed computing, whose significance and impact on the theory or practice of distributed computing have been evident for at least a decade. It is sponsored jointly by the ACM Symposium on Principles of Distributed Computing (PODC) and the EATCS Symposium on Distributed Computing (DISC). The prize is presented annually, with the presentation taking place alternately at PODC and DISC.
Marcos Aguiliera, Andréa W. Richa, Alexander A. Schwarzmann, Alessandro Panconesi, Christian Scheideler, Philipp Woelfel
PODC3
2022 Implementing three exchange read operations for distributed atomic storage
Chryssis Georgiou, Theophanis Hadjistasi, Nicolas C. Nicolaou, Alexander A. Schwarzmann
J. Parallel Distributed Comput.4
2021 Towards a Robust Distributed Framework for Election-Day Voter Check-In
Alexander A. Schwarzmann
SSS1
2021 Tractable low-delay atomic memory
Antonio Fernández 0001, Theophanis Hadjistasi, Nicolas C. Nicolaou, Alexandru Popa 0001, Alexander A. Schwarzmann
Distributed Comput.5
2018 Consistent Distributed Memory Services: Resilience and Efficiency (Invited Paper)
abstract
Reading, 'Riting, and 'Rithmetic, the three R's underlying much of human intellectual activity, not surprisingly, also stand as a venerable foundation of modern computing technology. Indeed, both the Turing machine and von Neumann machine models operate by reading, writing, and computing, and all practical uniprocessor implementations are based on performing activities structured in terms of the three R's. With the advance of networking technology, communication became an additional major systemic activity. However, at a high level of abstraction, it is apparently still more natural to think in terms of reading, writing, and computing. While it is hard to imagine distributed systems - such as those implementing the World-Wide Web - without communication, we often imagine browser-based applications that operate by retrieving (i.e., reading) data, performing computation, and storing (i.e., writing) the results. In this article, we deal with the storage of shared readable and writable data in distributed systems that are subject to perturbations in the underlying distributed platforms composed of computers and networks that interconnect them. The perturbations may include permanent failures (or crashes) of individual computers, transient failures, and delays in the communication medium. The focus of this paper is on the implementations of distributed atomic memory services. Atomicity is a venerable notion of consistency, introduced in 1979 by Lamport [Lamport, 1979]. To this day atomicity remains the most natural type of consistency because it provides an illusion of equivalence with the serial object type that software designers expect. We define the overall setting, models of computation, definition of atomic consistency, and measures of efficiency. We then present algorithms for single-writer settings in the static models. Then we move to presenting algorithms for multi-writer settings. For both static settings we discuss design issues, correctness, efficiency, and trade-offs. Lastly we survey the implementation issues in dynamic settings, where the universe of participants may completely change over time. Here the expectation is that solutions are found by integrating static algorithms with a reconfiguration framework so that during periods of relative stability one benefits from the efficiency of static algorithms, and where during the more turbulent times performance degrades gracefully when reconfigurations are needed. We describe the most important approaches and provide examples.
Theophanis Hadjistasi, Alexander A. Schwarzmann
ICALP2
2018 2018 Doctoral Dissertation Award
abstract
The winner of the 2018 Principles of Distributed Computing Doctoral Dissertation Award is Dr. Rati Gelashvili, for his dissertation titled "On the Complexity of Synchronization," written under the supervision of Prof. Nir Shavit at the Massachusetts Institute of Technology.
Lorenzo Alvisi, Idit Keidar, Andréa W. Richa, Alexander A. Schwarzmann
PODC4
2017 Doing-it-All with bounded work and communication
Bogdan S. Chlebus, Leszek Gasieniec, Dariusz R. Kowalski, Alexander A. Schwarzmann
Inf. Comput.4
2017 Special issue containing selected expanded papers from the 17th International Symposium on Stabilization, Safety and Security of Distributed Systems (SSS 2015)
Andrzej Pelc, Alexander A. Schwarzmann
Inf. Comput.2
2017 Coordinated cooperative task computing using crash-prone processors with unreliable multicast
Seda Davtyan, Roberto De Prisco, Chryssis Georgiou, Theophanis Hadjistasi, Alexander A. Schwarzmann
J. Parallel Distributed Comput.5
2016 Storage-Optimized Data-Atomic Algorithms for Handling Erasures and Errors in Distributed Storage Systems
abstract
Erasure codes are increasingly being studied in the context of implementing atomic memory objects in large scale asynchronous distributed storage systems. When compared with the traditional replication based schemes, erasure codes have the potential of significantly lowering storage and communication costs while simultaneously guaranteeing the desired resiliency levels. In this work, we propose the Storage-Optimized Data-Atomic (SODA) algorithm for implementing atomic memory objects in the multi-writer multi-reader setting. SODA uses Maximum Distance Separable (MDS) codes, and is specifically designed to optimize the total storage cost for a given fault-tolerance requirement. For tolerating f server crashes in an n-server system, SODA uses an [n, k] MDS code with k = n - f, and incurs a total storage cost of n/n-f. SODA is designed under the assumption of reliable point-to-point communication channels. The communication cost of a write and a read operation are respectively given by O(f2) and n/n-f(δw+1), where δwdenotes the number of writes that are concurrent with the particular read. In comparison with the recent CASGC algorithm [1], which also uses MDS codes, SODA offers lower storage cost while pays more on the communication cost. We also present a modification of SODA, called SODAerr, to handle the case where some of the servers can return erroneous coded elements during a read operation. Specifically, in order to tolerate f server failures and e error-prone coded elements, the SODAerr algorithm uses an [n, k] MDS code such that k = n - 2e - f. SODAerr also guarantees liveness and atomicity, while maintaining an optimized total storage cost of n/n-f-2e.
Kishori M. Konwar, N. Prakash 0001, Erez Kantor, Nancy A. Lynch, Muriel Médard, Alexander A. Schwarzmann
IPDPS6
2016 Brief Announcement: Oh-RAM! One and a Half Round Read/Write Atomic Memory
abstract
Emulating atomic read/write shared objects in a message-passing system is a fundamental problem in distributed computing. Considering that network communication is the most expensive resource, efficiency is measured first of all in terms of the communication needed to implement read and write operations. It is well known that two communication round-trip phases involving in total four message exchanges are sufficient to implemented atomic operations. In this work we present a comprehensive treatment of the question of when and how it is possible to implement atomic memory where read and write operations complete in three message exchanges, i.e., we aim for One and half Round Atomic Memory, hence the name Oh-RAM! We present algorithms that allow operations to complete in three communication exchanges without imposing any constraints on the number of readers and writers. We present an implementation for the {single-writer/multiple-reader} (SWMR) setting, where reads complete in three communication exchanges and writes complete in two exchanges. Then we pose the question of whether it is possible to implement multiple-writer/multiple-reader (MWMR) memory where operations complete in at most three communication exchanges. In light of our impossibility result these algorithms are optimal in terms of the number of communication exchanges.
Theophanis Hadjistasi, Nicolas C. Nicolaou, Alexander A. Schwarzmann
PODC3
2015 Robust network supercomputing with unreliable workers
Kishori M. Konwar, Sanguthevar Rajasekaran, Alexander A. Schwarzmann
J. Parallel Distributed Comput.3
2015 Dealing with undependable workers in decentralized network supercomputing
Seda Davtyan, Kishori M. Konwar, Alexander Russell, Alexander A. Schwarzmann
Theor. Comput. Sci.4
2014 Coordinated Cooperative Work Using Undependable Processors with Unreliable Broadcast
abstract
With the end of Moore's Law in sight, parallelism became the main means for speeding up computationally intensive applications, especially in the cases where large collections of tasks need to be performed. Network supercomputing -- taking advantage of very large numbers of computers in a distributed environment is an effective approach to massive parallelism that harnesses the processing power inherent in large networked settings. In such settings, processor failures are no longer an exception, but the norm. Any algorithm designed for realistic settings must be able to deal with failures. This paper presents a new message-passing algorithm for distributed cooperative work in synchronous settings where processors may crash, and where any broadcasts performed by crashing processors are unreliable. We specify the algorithm, prove that it is correct, and perform extensive simulations that show that its performance is close to similar algorithms that use reliable broadcast, and that its work compares favorably to the relevant lower bounds.
Seda Davtyan, Roberto De Prisco, Chryssis Georgiou, Alexander A. Schwarzmann
PDP4
2014 Dependable Decentralized Cooperation with the Help of Reliability Estimation
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
SSS3
2013 Estimating Reliability of Workers for Cooperative Distributed Computing
abstract
Internet supercomputing is an approach to solving partitionable, computation-intensive problems by harnessing the power of a vast number of interconnected computers. For the problem of using network supercomputing to perform a large collection of independent tasks, prior work introduced a decentralized approach and provided randomized synchronous algorithms that perform all tasks correctly with high probability, while dealing with misbehaving or crash-prone processors. The main weaknesses of existing algorithms is that they assume either that the average probability of a non-crashed processor returning incorrect results is inferior to 12, or that the probability of returning incorrect results is known to each processor. Here we present a randomized synchronous distributed algorithm that tightly estimates the probability of each processor returning correct results. Starting with the set P of n processors, let F be the set of processors that crash. Our algorithm estimates the probability pi of returning a correct result for each processor i ∈ P - F, making the estimates available to all these processors. The estimation is based on the (ε, δ)-approximation, where each estimated probability p̃iof piobeys the bound Pr[pi(1 - ε) ≤ p̃i≤ pi(1 + ε)] > 1 - δ, for any constants δ > 0 and ε > 0 chosen by the user. An important aspect of this algorithm is that each processor terminates without global coordination. We assess the efficiency of the algorithm in three adversarial models as follows. For the model where the number of non-crashed processors P - F is linearly bounded the time complexity T (n) of the algorithm is O(log n), work complexity W(n) is O(n log n), and message complexity M(n) is O(n log2n). For the model where P - F is bounded by a fractional polynomial we have T(n) = O(n1-alog n log log n), W(n) = O(n log n log log n), and M(n) = O(n log2n log log n). For the model where P - F is bounded by a poly-logarithm we have T(n) = O(n), W(n) = O(n poly log n), and M(n) = O(n log2n poly log n). All bounds are shown to hold with high probability.
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
ISPDC3
2013 Self-stabilizing Resource Discovery Algorithm
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
OPODIS3
2013 Brief announcement: self-stabilizing resource discovery algorithm
abstract
Distributed cooperative computing in networks involves marshaling collections of network nodes possessing the necessary computational resources. Before the willing nodes can act in a concerted way they must first discover one another. This is the general setting of the Resource Discovery Problem (RDP). This paper presents a self-stabilizing algorithm that solves RDP in a deterministic synchronous setting. The solution approach is formulated in terms of evolving knowledge graphs, where vertices represent the participating network nodes, and edges represent one node's knowledge about another. Ideally, the diameter of such a graph is one, i.e., each node knows all others. The algorithm works in rounds as it evolves the knowledge graph with the goal of reducing its diameter. This is accomplished by nodes sharing their knowledge through gossip messages. We prove that the algorithm is self-stabilizing, i.e., it tolerates arbitrary perturbations in the nodes' local states and is guaranteed to solve the problem once such failures subside. The algorithm has stabilization time of O(D), and it takes at most 4D + 4 complete round to stabilize, where D is the diameter of the initial knowledge graph, and the corresponding message complexity is O,(|V|j ⋅D), where V is the set of participating nodes.
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
PODC3
2013 Special issue on DISC 2010
Nancy A. Lynch, Alexander A. Schwarzmann
Distributed Comput.2
2012 Brief announcement: decentralized network supercomputing in the presence of malicious and crash-prone workers
abstract
Internet supercomputing is an approach to solving partitionable, computation-intensive problems by harnessing the power of a vast number of interconnected computers. For the problem of using network supercomputing to perform a large collection of independent tasks, our prior work introduced the decentralized approach, and provided a synchronous algorithm that is able to perform all tasks with high probability (whp), while dealing with malicious behaviors under a rather strong assumption that the average probability of live (non-crashed) processors returning bogus results remains inferior to 1/2 during the computation. There the adversary is severely limited in its ability to crash processors that normally return correct results. This work develops an efficient synchronous decentralized algorithm that is able to deal with a much stronger adversary. We consider a failure model with crashes, where given the initial set of processors P, an adversary is able to crash any subset F of processors, where |F| ≤ f•n, for a constant f (0<f<1), under the constraint that there exists a subset H ⊆ P - F, with |H| = Ω(n), called the hardened set, such that the average probability of a processor in H returning a bogus result is inferior to 1/2. Here any processor may return bogus results, and H may be much smaller than P-F, while the average probability of processors in P-F returning a bogus result may be greater than 1/2. We develop an efficient randomized algorithm for n processors and t tasks (n≤t), where each live processor is able to determine locally when all tasks are performed, and obtain the results of all tasks. We prove that in Θ(t⁄n logn) rounds all live workers know the results of all tasks whp, and that these results are correct whp. The work complexity of the algorithm is Θ(tlogn), the message complexity is Θ(nlogn ), and the bit complexity is O(tn log3n).
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
PODC3
2011 Towards Feasible Implementations of Low-Latency Multi-writer Atomic Registers
abstract
This work explores implementations of multiwriter/multi-reader (MWMR) atomic registers in asynchronous, crash-prone, message-passing systems with the focus on low latency and computational feasibility. The efficiency of atomic read/write register implementations is traditionally measured in terms of the latency of read and write operations. To reduce operation latency researchers focused on the communication costs, expressed as the number of communication round-trips (or rounds), often ignoring the computation costs. In this paper we consider efficiency of a register implementation in terms of both communication and computation costs. As of this writing, algorithm SFW is the sole known MWMR algorithm that allows single round read and write operations. The algorithm uses collections of intersecting sets (quorums), and to enable single round operations, SFW relies on the evaluation of certain predicates. We formulate a new combinatorial problem that captures the computational burden of evaluating the predicates in algorithm SFW and we show that it is NP-Complete. To make the evaluation of the predicates feasible, we present a polynomial log-approximation algorithm for this problem and we show how to use it with algorithm SFW. Then we present a new algorithm, called CWFR, that allows fast operations independently of the underlying quorum system construction. The algorithm implements two-round writes and allows reads to complete in a single round. We conclude with experimental evaluations of our algorithms obtained from simulations in NS2.
Chryssis Georgiou, Nicolas C. Nicolaou, Alexander Russell, Alexander A. Schwarzmann
NCA4
2011 Robust Network Supercomputing without Centralized Control
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
OPODIS3
2011 Robust network supercomputing without centralized control
abstract
Traditional approaches to network supercomputing employ a master process and a large number of potentially undependable worker processes that must perform a collection of tasks on behalf of the master. In such a centralized scheme, the master process is a performance bottleneck and a single point of failure. This work develops an original approach that eliminates the master and instead uses a decentralized algorithm, where each worker is able to determine locally that all tasks have been performed, and to collect locally the results of all tasks. The failure model assumes that the average probability of a worker returning a wrong result is inferior to 1/2. A randomized synchronous algorithm for n processes and n tasks is presented. The algorithm terminates in Θ(log n) rounds, and it is proved that upon termination the workers know the results of all tasks with high probability, and that these results are correct with high probability. The message complexity of the algorithm is Θ(n log n), and the bit complexity is O(n2 log3 n).
Seda Davtyan, Kishori M. Konwar, Alexander A. Schwarzmann
PODC3
2010 Load Balancing and Almost Symmetries for RAMBO Quorum Hosting
Laurent D. Michel, Alexander A. Schwarzmann, Elaine L. Sonderegger, Pascal Van Hentenryck
CP2
2010 Rambo: a robust, reconfigurable atomic memory service for dynamic networks
Seth Gilbert, Nancy A. Lynch, Alexander A. Schwarzmann
Distributed Comput.3
2010 Emulating shared-memory Do-All algorithms in asynchronous message-passing systems
Dariusz R. Kowalski, Mariam Momenzadeh, Alexander A. Schwarzmann
J. Parallel Distributed Comput.3
2010 Editors' preface
Alexander A. Schwarzmann, Pascal Felber
Theor. Comput. Sci.1
2009 Online Selection of Quorum Systems for RAMBO Reconfiguration
Laurent D. Michel, Martijn Moraal, Alexander A. Schwarzmann, Elaine L. Sonderegger, Pascal Van Hentenryck
CP3
2009 Bandwidth-Limited Optimal Deployment of Eventually-Serializable Data Services
Laurent D. Michel, Pascal Van Hentenryck, Elaine L. Sonderegger, Alexander A. Schwarzmann, Martijn Moraal
CPAIOR4
2009 On the Efficiency of Atomic Multi-reader, Multi-writer Distributed Memory
Burkhard Englert, Chryssis Georgiou, Peter M. Musial, Nicolas C. Nicolaou, Alexander A. Schwarzmann
OPODIS5
2009 At-most-once semantics in asynchronous shared memory
abstract
This paper investigates the feasibility of implementing at-most-once access semantics in a model where a collection of actions is to be performed by failure-prone, asynchronous shared-memory processes. We introduce the At-Most-Once problem for performing a set of n jobs using m processors, and we define the notion of efficiency for such protocols, called effectiveness, that allows the classification of algorithms solving the problem. The effectiveness for an at-most-once implementation is the number of jobs safely completed by the implementation, expressed as a function of the number of jobs n, the number of processes m, and the number of process crashes f. We prove a lower bound of n--f on the effectiveness of any algorithm. We then present two process solutions that offer a trade off between work and space complexity. Finally, we generalize a two-process solution for the multi-process setting using a hierarchical algorithm that achieves effectiveness of n--log m†o(n), coming reasonably close, asymptotically, to the corresponding lower bound.
Sotiris Kentros, Aggelos Kiayias, Nicolas C. Nicolaou, Alexander A. Schwarzmann
SPAA4
2009 At-Most-Once Semantics in Asynchronous Shared Memory
Sotiris Kentros, Aggelos Kiayias, Nicolas C. Nicolaou, Alexander A. Schwarzmann
DISC4
2009 Reconfigurable distributed storage for dynamic networks
Gregory V. Chockler, Seth Gilbert, Vincent Gramoli, Peter M. Musial, Alexander A. Schwarzmann
J. Parallel Distributed Comput.5
2009 Fault-tolerant semifast implementations of atomic read/write registers
Chryssis Georgiou, Nicolas C. Nicolaou, Alexander A. Schwarzmann
J. Parallel Distributed Comput.3
2009 Node discovery in networks
Kishori M. Konwar, Dariusz R. Kowalski, Alexander A. Schwarzmann
J. Parallel Distributed Comput.3
2009 Editors' preface
Ivan Lavallée, Alexander A. Schwarzmann
Theor. Comput. Sci.2
2009 State-wide elections, optical scan voting systems, and the pursuit of integrity
abstract
In recent years, two distinct electronic voting technologies have been introduced and extensively utilized in election procedures: direct recording electronic systems and optical scan (OS) systems. The latter are typically deemed safer, as they inherently provide a voter-verifiable paper trail that enables hand-counted audits and recounts that rely on direct voter input. For this reason, OS machines have been widely deployed in the United States. Despite the growing popularity of these machines, they are known to suffer from various security vulnerabilities that, if left unchecked, can compromise the integrity of elections in which the machines are used. This article studies general auditing procedures designed to enhance the integrity of elections conducted with optical scan equipment and, additionally, describes the specific auditing procedures currently in place in the State of Connecticut. We present an abstract view of a typical OS voting technology and its relationship to the general election process. With this in place, we lay down a ldquotemporal-resourcerdquo adversarial model, providing a simple language for describing the disruptive power of a potential adversary. Finally, we identify how audit procedures, injected at various critical stages before, during, and after an election, can frustrate such adversarial interference and so contribute to election integrity. We present the implementation of such auditing procedures for elections in the State of Connecticut utilizing the Premiere (Diebold) AccuVote OS; these audits were conducted by the UConn VoTeR Center, at the University of Connecticut, on request of the Office of the Secretary of the State. We discuss the effectiveness of such procedures in every stage of the process and we present results and observations gathered from the analysis of past election data.
Tigran Antonyan, Seda Davtyan, Sotiris Kentros, Aggelos Kiayias, Laurent D. Michel, Nicolas C. Nicolaou, Alexander Russell, Alexander A. Schwarzmann
IEEE Trans. Inf. Forensics Secur.8
2009 Developing a Consistent Domain-Oriented Distributed Object Service
abstract
This paper presents a new algorithm for a reconfigurable distributed domain-oriented atomic object service, called DO-RAMBO, which stands for Domain-Oriented Reconfigurable Atomic Memory for Basic Objects. This service is suitable for inclusion as a middleware system service for distributed applications requiring atomic read/write data. The implementation substantially extends and refines the abstract RAMBO algorithm of Lynch and Shvartsman that supports individual atomic objects. In this paper, domains are introduced to allow the users to group related atomic objects. The new implementation manages configurations on the basis of domains, significantly improving the utility and the performance of the resulting service. DO-RAMBO guarantees consistency under asynchrony, message loss, node crashes, new node arrivals, and node departures. We present the formal algorithm development for DO-RAMBO and give analytical and empirical results that illustrate the benefit of the new approach.
Chryssis Georgiou, Peter M. Musial, Alexander A. Schwarzmann
IEEE Trans. Parallel Distributed Syst.3
2008 Optimal Deployment of Eventually-Serializable Data Services
Laurent D. Michel, Alexander A. Schwarzmann, Elaine L. Sonderegger, Pascal Van Hentenryck
CPAIOR2
2008 Spontaneous, Self-Sampling Quorum Systems for Ad Hoc Networks
abstract
Quorum systems-collections of sets with pairwise nonempty intersections-are used in distributed settings to implement services such as consensus and consistent memory. Quorums have been substantially studied in static settings, however the design and analysis of quorum-based distributed services in resource-limited ad hoc networks is a relatively unexplored area. The pioneering work of Chockler, Gilbert, and Patt-Shamir considers such networks and proposes an implementation of probabilistic quorum systems with per-node communication bit complexity of O(log2n), where n is the number of nodes. The authors assumes a priori knowledge of node failure probability p, where 0 ¿ p2n). We demonstrate the utility of our construction by presenting a single-writer, multi-reader algorithm that uses our probabilistic quorums to implement atomic objects in ad hoc networks, where consistency is guaranteed with high probability. We include simulation results illustrating the high probability guarantee for our atomic memory service.
Kishori M. Konwar, Peter M. Musial, Alexander A. Schwarzmann
ISPDC3
2008 An Abstract Channel Specification and an Algorithm Implementing It Using Java Sockets
abstract
Abstract models and specifications can be used in the design of distributed applications to formally reason about their safety properties. However, the benefits of using formal methods are often negated by the ad hoc process of mapping the semantics of an abstract specification to algorithms designed to be executed on target distributed platforms. The challenge of formally specifying communication channels and correctly implementing them as algorithms that use realistic distributed system services is the focus of this paper. This work provides an original formal specification of an abstract asynchronous communication channel with support for dynamic creation and tear down of links between participating network nodes, and its implementation as an algorithm using Java sockets. The specification and the algorithm are expressed using the Input/Output Automata formalism, and it is proved that the algorithm correctly implements the specification, viz. that any externally observable behavior (trace) of the algorithm has a corresponding behavior of the specification. The approach presented here can be used to implement algorithms for dynamic systems, where communicating nodes may join, leave, and experience delays. The result is also of direct benefit to automated code generation, such as that implemented within the Input/Output Automata Toolkit at MIT.
Chryssis Georgiou, Peter M. Musial, Alexander A. Schwarzmann, Elaine L. Sonderegger
NCA3
2008 On the robustness of (semi) fast quorum-based implementations of atomic shared memory
abstract
Atomic (linearizable) read/write memory is a fundamental abstractions in distributed computing. Following a seminal implementation of atomic memory of Attiya et al. [6], a folklore belief developed that in messaging-passing atomic memory implementations "reads must write." However, work by Dutta et al. [4] established that if the number of readers R is constrained with respect to the number of replicas S and the maximum number of crash-failures t so that R < S/t - 2, then single communication round-trip reads are possible. Such an implementation given in [4] is called fast. Subsequently, Georgiou et al. [3] relaxed the constraint in [4], and proposed semifast implementations with unbounded number of readers, where under realistic conditions most reads need only a single communication round-trip to complete. Their approach groups collections of readers into virtual nodes. Semifast behavior of their algorithm is preserved as long as the number of virtual nodes V is constrained by V < S/t - 2.
Chryssis Georgiou, Nicolas C. Nicolaou, Alexander A. Schwarzmann
PODC3
2008 On the Robustness of (Semi) Fast Quorum-Based Implementations of Atomic Shared Memory
Chryssis Georgiou, Nicolas C. Nicolaou, Alexander A. Schwarzmann
DISC3
2008 Introduction to special issue dedicated to the DISC 20th anniversary
Shlomi Dolev, Alexander A. Schwarzmann
Distributed Comput.2
2008 Writing-all deterministically and optimally using a nontrivial number of asynchronous processors
Dariusz R. Kowalski, Alexander A. Schwarzmann
ACM Trans. Algorithms2
2007 Tampering with Special Purpose Trusted Computing Devices: A Case Study in Optical Scan E-Voting
abstract
Special purpose trusted computing devices are currently being deployed to offer many services for which the general purpose computing paradigm is unsuitable. The nature of the services offered by many of these devices demand high security and reliability, as well as low cost and low power consumption. Electronic Voting machines is a canonical example of this phenomenon. With electronic voting machines currently being used in much of the United States and several other countries, there is a strong need for thorough security evaluation of these devices and the procedures in place for their use. In this work, we first put forth a general framework for special purpose trusted computing devices. We then focus on Optical Scan (OS) electronic voting technology as a specific instance of this framework. OS terminals are a popular e-voting technology with the decided advantage of a user-verified paper trail: the ballot sheets themselves. Still election results are based on machine- generated totals as well as machine-generated audit reports to validate the voting process. In this paper we present a security assessment of the Diebold AccuVote Optical Scan voting terminal (AV-OS), a popular OS terminal currently in wide deployment anticipating the 2008 Presidential elections. The assessment is developed using exclusively reverse-engineering, without any technical specifications provided by the machine suppliers. We demonstrate a number of security issues that relate to the machine's proprietary language, called AccuBasic, that is used for reporting election results. While this language is thought to be benign, especially given that it is essentially sandboxed by the firmware to have only read access, we demonstrate that it is powerful enough to (i) strengthen known attacks against the AV-OS so that they become undetectable prior to elections (and thus significantly increasing their magnitude) or, (ii) to conditionally bias the election results to reach a desired outcome. Given the discovered vulnerabilities and attacks we proceed to discuss how random audits can be used to validate with high confidence that a procedure carried out by special purpose devices such as the AV-OS has not been manipulated. We end with a set of recommendations for the design and safe-use of OS voting systems.
Aggelos Kiayias, Laurent D. Michel, Alexander Russell, Narasimha K. Shashidhar, Andrew See, Alexander A. Schwarzmann, Seda Davtyan
ACSAC6
2007 Implementing Atomic Data through Indirect Learning in Dynamic Networks
abstract
Developing middleware services for dynamic distributed systems, e.g., ad-hoc networks, is a challenging task given that such services deal with dynamically changing membership and asynchronous communication. Algorithms developed for static settings are often not usable in such settings because they rely on (logical) all-to-all node connectivity through routing protocols, which may be unfeasible or prohibitively expensive to implement in highly dynamic settings. This paper explores the indirect learning, via periodic gossip, approach to information dissemination within a dynamic, distributed data service implementing atomic read/write memory service. The indirect learning scheme is used to improve the liveness of the service in the settings with uncertain connectivity. The service is formally proved to guarantee atomicity in all executions. Conditional performance analysis of the new service is presented, where this analysis has the potential of being generalized to other similar dynamic algorithms. Under the assumption that the network is connected, and assuming reasonable timing conditions, the bounds on the duration of read/write operations of the new service are calculated. Finally, the paper proposes a deployment strategy where indirect learning leads to an improvement in communication costs relative to a previous solution that assumes all-to-all connectivity.
Kishori M. Konwar, Peter M. Musial, Nicolas C. Nicolaou, Alexander A. Schwarzmann
NCA4
2007 A formal treatment of an abstract channel implementation using java sockets and TCP
abstract
Abstract models and specifications can be used in the design of distributed applications to formally reason about their safety properties. However, the benefits of using formal methods are offset by the challenging process of mapping the functionality of an abstract specification to the low-level executable code for target distributed platforms. Formal specification and practical implementation of communication channels is one such challenge. This work provides the first formal specification of an abstract asynchronous communication channel with support for dynamic creation and tear down of communication links between participating network nodes, and its implementation using Java sockets and TCP. The specifications are formulated using Input/Output Automata formalism, and it is proved that the resulting implementation preserves the safety properties of the abstract channel. The approach presented here can be used to implement algorithms for dynamic systems, where communicating nodes may join, leave, and experience arbitrary delays, and it can directly benefit automated code generation.
Chryssis Georgiou, Peter M. Musial, Alexander A. Schwarzmann, Elaine L. Sonderegger
PODC3
2007 Long-lived Rambo: Trading knowledge for communication
Chryssis Georgiou, Peter M. Musial, Alexander A. Schwarzmann
Theor. Comput. Sci.3
2006 Resource Discovery in Networks under Bandwidth Limitations
abstract
The resource discovery problem, where cooperating machines need to find one another in a network, was introduced by Harchol-Balter, Leighton, and Lewin (1999) in the context of Akamai Technologies with the goal of building an Internet-wide content-distribution system. In the solutions for the synchronous setting proposed so far in the papers by Harchol-Bartel et al. (1999), Kutten et al. (2001) and Law and Siu (2000), there is a possibility that during some time step many machines may contact a single machine, and this is not a realistic assumption. This work assumes a synchronous model, however at each step a machine can send and receive only a constant number of messages. It is shown that the conjectured poly-logarithmic upper bound (Harchol-Bartel et al., 1999) for such a setting is not possible. This is done by proving a lower bound on time of Omega(n), where n is the number of participating nodes. For this model a randomized algorithm is presented that solves the resource discovery problem in O(n log2n) time, i.e., within a poly-logarithmic factor of the corresponding lower bound. The algorithm has a O(n2log2n) message complexity and O(n3log3n) communication complexity. Simulation results for the algorithm illustrate the lower and upper bounds, and lead to interesting observations
Kishori M. Konwar, Alexander A. Schwarzmann
ISPDC2
2006 Fault-tolerant semifast implementations of atomic read/write registers
abstract
This paper investigates time-efficient implementations of atomic read-write registers in message-passing systems where the number of readers can be unbounded. In particular we study the case of a single writer, multiple readers, and S servers, such that the writer, any subset of the readers, and up to t servers may crash. A recent result of Dutta et al. [3] shows how to obtain fast implementations in which both reads and writes complete in one communication round-trip, under the constraint that the number of readers is less than S t - 2, where t < S 2 . In that same paper the authors pose a question of whether it is possible to relax the bound on readers, and at what cost, if semifast implementations are considered, i.e., implementations that have fast reads or fast writes.This paper provides an answer to this question. It is shown that one can obtain implementations where all writes are fast, i.e., involving a single round-trip communication, and where reads complete in one to two communication rounds under the assumption that no more than t < S 2 servers crash. Simulated scenarios included in this paper indicate that only a small fraction of reads require a second communication round. Interestingly the correctness of the implementation does not depend on the number of concurrent readers in the system. The solution is obtained with the help of non-unique virtual ids assigned to each reader, where the readers sharing a virtual id form a virtual node. For the proposed definition of semifast implementations it is shown that implementations satisfying certain assumptions are semifast if and only if the number of virtual ids in the system is less than S t - 2. This result is proved to be tight in terms of the required communication. It is shown that only a single complete two-round read operation may be necessary for each write operation. It is furthermore shown that no semifast implementation exists for the multi-reader, multi-writer model.
Chryssis Georgiou, Nicolas C. Nicolaou, Alexander A. Schwarzmann
SPAA3
2006 Brief Announcement: Fault-Tolerant SemiFast Implementations of Atomic Read/Write Registers
Chryssis Georgiou, Nicolas C. Nicolaou, Alexander A. Schwarzmann
DISC3
2006 Robust Network Supercomputing with Malicious Processes
Kishori M. Konwar, Sanguthevar Rajasekaran, Alexander A. Schwarzmann
DISC3
2006 Distributed scheduling for disconnected cooperation
Grzegorz Malewicz, Alexander Russell, Alexander A. Schwarzmann
Distributed Comput.3
2006 Dynamic load balancing with group communication
Shlomi Dolev, Roberto Segala, Alexander A. Schwarzmann
Theor. Comput. Sci.3
2005 Improved algorithms for multiplex PCR primer set selection with amplification length constraints
Kishori M. Konwar, Ion I. Mandoiu, Alexander Russell, Alexander A. Schwarzmann
APBC4
2005 Explicit Combinatorial Structures for Cooperative Distributed Algorithms
abstract
Cooperation in distributed settings often involves activities that must be performed at least once by the participating processors. When processor failures or delays occur, it becomes unavoidable that some tasks are done redundantly. To make efficient use of the available processors, several distributed algorithms schedule the activities of the processors in terms of permutations of tasks that need to be performed at least once. This paper presents the first explicit practical deterministic construction of sets of permutations with certain combinatorial properties that immediately make practical several deterministic distributed algorithms. These algorithms solve a variety of problems, for example, cooperation in shared-memory and message-passing settings, and the gossip problem. Prior to this work, the most efficient algorithms for some of these problems were primarily of theoretical interest - they relied on permutations that are known to exist, but very expensive to construct, with the cost of construction being at least exponential in the size of the permutations. In this paper, the explicitly constructed permutations are ultimately used directly to produce practical instances of several classes of efficient deterministic algorithms. Most importantly, for all of these algorithms, the schedule construction cost is reduced from exponential to polynomial, at the expense of slight detuning, at most polylogarithmic, of the efficiency of these algorithms
Dariusz R. Kowalski, Peter M. Musial, Alexander A. Schwarzmann
ICDCS3
2005 Developing a Consistent Domain-Oriented Distributed Object Service
abstract
This paper presents a new algorithm for a reconfigurable distributed domain-oriented atomic object service, called DO-RAMBO, which stands for domain-oriented reconfigurable atomic memory for basic objects. This service is suitable for inclusion as a middleware system service for distributed applications requiring atomic read/write data. The implementation substantially extends and refines the abstract RAMBO algorithm of Lynch and Shvartsman that supports individual atomic objects. In this paper domains are introduced to allow the users to group related atomic objects. The new implementation manages configurations on the basis of domains, significantly improving the utility and the performance of the resulting service. DO-RAMBO guarantees consistency under asynchrony, message loss, node crashes, new node arrivals, and node departures. We present the formal algorithm development for DO-RAMBO and give analytical and preliminary empirical results that illustrate the benefit of the new approach
Chryssis Georgiou, Peter M. Musial, Alexander A. Schwarzmann
NCA3
2005 Reconfigurable Distributed Storage for Dynamic Networks
Gregory V. Chockler, Seth Gilbert, Vincent Gramoli, Peter M. Musial, Alexander A. Schwarzmann
OPODIS5
2005 Node Discovery in Networks
Kishori M. Konwar, Dariusz R. Kowalski, Alexander A. Schwarzmann
OPODIS3
2005 Autonomous virtual mobile nodes
abstract
This paper presents a new abstraction for virtual infrastructure in mobile ad hoc networks. An AutonomousVirtual Mobile Node (AVMN) is a robust and reliable entity that is designed to cope with theinherent difficulties caused by processors arriving, leaving, and moving according to their own agendas,as well as with failures and energy limitations. There are many types of applications that may make useof the AVMN infrastructure: tracking, supporting mobile users, or searching for energy sources.The AVMN extends the focal point abstraction in [9] and the virtual mobile node abstraction in [10].The new abstraction is that of a virtual general-purpose computing entity, an automaton that can makeautonomous on-line decisions concerning its own movement. We describe a self-stabilizing implementationof this new abstraction that is resilient to the chaotic behavior of the physical processors and providesautomatic recovery from any corrupted state of the system.
Shlomi Dolev, Seth Gilbert, Elad Michael Schiller, Alexander A. Schwarzmann, Jennifer L. Welch
SPAA4
2005 DNA-BAR: distinguisher selection for DNA barcoding
abstract
Summary: DNA-BAR is a software package for selecting DNA probes (henceforth referred to as distinguishers) that can be used in genomic-based identification of microorganisms. Given the genomic sequences of the microorganisms, DNA-BAR finds a near-minimum number of distinguishers yielding a distinct hybridization pattern for each microorganism. Selected distinguishers satisfy user specified bounds on length, melting temperature and GC content, as well as redundancy and cross-hybridization constraints. Availability: DNA-BAR can be used online through the web interface provided at http://dna.engr.uconn.edu/~software/DNA-BAR/. The open source C code, released under the GNU General Public License, is also available at the above address. Contact: [email protected]
Bhaskar DasGupta, Kishori M. Konwar, Ion I. Mandoiu, Alexander A. Schwarzmann
Bioinform.4
2005 GeoQuorums: implementing atomic memory in mobile ad hoc networks
Shlomi Dolev, Seth Gilbert, Nancy A. Lynch, Alexander A. Schwarzmann, Jennifer L. Welch
Distributed Comput.4
2005 Performing work with asynchronous processors: Message-delay-sensitive bounds
Dariusz R. Kowalski, Alexander A. Schwarzmann
Inf. Comput.2
2005 Work-Competitive Scheduling for Cooperative Computing with Dynamic Groups
abstract
The problem of cooperatively performing a set of t tasks in a decentralized computing environment subject to failures is one of the fundamental problems in distributed computing. The setting with partitionable networks is especially challenging, as algorithmic solutions must accommodate the possibility that groups of processors become disconnected (and, perhaps, reconnected) during the computation. The efficiency of task-performing algorithms is often assessed in terms of work: the total number of tasks, counting multiplicities, performed by all of the processors during the computation. In general, the scenario where the processors are partitioned into g disconnected components causes any task-performing algorithm to have work $\Omega(t\cdot g)$ even if each group of processors performs no more than the optimal number of $\Theta(t)$ tasks. Given that such pessimistic lower bounds apply to any scheduling algorithm, we pursue a competitive analysis. Specifically, this paper studies a simple randomized scheduling algorithm for p asynchronous processors, connected by a dynamically changing communication medium, to complete t known tasks. The performance of this algorithm is compared against that of an omniscient off-line algorithm with full knowledge of the future changes in the communication medium. The paper describes a notion of computation width, which associates a natural number with a history of changes in the communication medium, and shows both upper and lower bounds on work-competitiveness in terms of this quantity. Specifically, it is shown that the simple randomized algorithm obtains the competitive ratio $(1+\mathbf{cw}/e)$, where $\mathbf{cw}$ is the computation width and e is the base of the natural logarithm ($e=2.7182\ldots$); this competitive ratio is then shown to be tight.
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
SIAM J. Comput.3
2005 The Do-All problem with Byzantine processor failures
Antonio Fernández 0001, Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
Theor. Comput. Sci.4
2005 Efficient gossip and robust distributed computation
Chryssis Georgiou, Dariusz R. Kowalski, Alexander A. Schwarzmann
Theor. Comput. Sci.3
2004 The Join Problem in Dynamic Network Algorithms
abstract
Distributed algorithms in dynamic networks often employ communication patterns whose purpose is to disseminate information among the participants. Gossiping is one form of such communication pattern. In dynamic settings, the set of participants can change substantially as new participants join, and as failures and voluntary departures remove those who have joined previously. A natural question for such settings is: how soon can newly joined nodes discover each other by means of gossiping? This paper abstracts and studies the join problem for dynamic systems that use all-to-all gossip. The problem is studied in terms of join-connectivity graphs where vertices represent the participants and where each edge represents one participant's knowledge about another. Ideally, such a graph has diameter one, i.e., all participants know each other. The diameter can grow as new participants join, and as failures remove edges from the graph. Gossip helps participants discover one another, decreasing the diameter. The results describe the lower and upper bounds on the number of communication rounds such that the participants who have previously joined discover one another, under a variety of assumptions about the joining and failures. For example, in the case when new participants join at multiple participants and participants may crash, the number of rounds cannot be bounded. In the more benign cases when the failures can be controlled or when new participants join at only one participant, the bound on rounds is shown to be logarithmic in the diameter of the initial configuration.
Kishori M. Konwar, Dariusz R. Kowalski, Alexander A. Schwarzmann
DSN3
2004 Implementing a Reconfigurable Atomic Memory Service for Dynamic Networks
abstract
Summary form only given. Transforming abstract algorithm specifications into executable code is an error-prone process in the absence of sophisticated compilers that can automatically translate such specifications into the target distributed system. We present a framework that was developed for translating algorithms specified as Input/Output Automata (IOA) to distributed programs. The framework consists of a methodology that guides the software development process and a core set of functions needed in target implementations that reduce unnecessary software development. The systems developed using this methodology preserve the modularity of the original specifications, making it easier to track refinements and effect optimizations. As a proof of concept, this work also presents a distributed implementation of a reconfigurable atomic memory service for dynamic networks (RAMBO). This service emulates atomic read/write shared objects in the dynamic setting where processors can arbitrarily crash, or join and leave the computation. The algorithm tolerates processor crashes and message loss and guarantees atomicity for arbitrary patterns of asynchrony and failure. The algorithm implementing the service is given in terms of IOA. An important consideration in formulating RAMBO was that it could be employed as a building block in real systems. Following a formal presentation of RAMBO algorithm, this work describes an optimized implementation that was developed using the methodology presented here. The system is implemented in Java and runs on a network of workstations. Empirical data illustrates the behavior of the system.
Peter M. Musial, Alexander A. Schwarzmann
IPDPS2
2004 Brief announcement: virtual mobile nodes for mobile ad hoc networks
abstract
No abstract available.
Shlomi Dolev, Seth Gilbert, Nancy A. Lynch, Elad Michael Schiller, Alexander A. Schwarzmann, Jennifer L. Welch
PODC5
2004 Long-Lived Rambo: Trading Knowledge for Communication
Chryssis Georgiou, Peter M. Musial, Alexander A. Schwarzmann
SIROCCO3
2004 Writing-all deterministically and optimally using a non-trivial number of asynchronous processors
abstract
The problem of performing n tasks on p asynchronous or undependable processors is a basic problem in distributed computing. This paper considers an abstraction of this problem called Write-All: using p processors write 1's into all locations of an array of size n. In this problem writing 1 abstracts the notion of performing a simple task. Despite substantial research, there is a dearth of efficient deterministic asynchronous algorithms for Write-All. Efficiency of algorithms is measured in terms of work that accounts for all local steps performed by the processors in solving the problem. Thus an optimal algorithm would have work Θ(n), however it is known that optimality cannot be achieved when p=Ω(n). The quest then is to obtain work-optimal solutions for this problem using a non-trivial, compared to n, number of processors p. Recently it was shown that optimality can be achieved using a non-trivial number M of processors, where M=4√n/log n. The new result in this paper significantly extends the range of processors for which optimality is achieved. The result shows that optimality can be achieved using close to M2 processors; more precisely, using (M log M)2-ε processors, for any ε > 0. Additionally, the new result uses only the atomic read/write memory, without resorting to using the test-and-set primitive that was necessary in the previous solution. This paper presents the algorithm and gives its analysis showing that the work complexity of the algorithm is O(n+p2+ε), which is optimal when p = O(n1/(2+ε)), while all prior deterministic algorithms require super-linear work when p=Ω(n1/4).
Dariusz R. Kowalski, Alexander A. Schwarzmann
SPAA2
2004 Collective asynchronous reading with polylogarithmic worst-case overhead
abstract
The Collect problem for an asynchronous shared-memory system has the objective for the processors to learn all values of a collection of shared registers, while minimizing the total number of read and write operations. First abstracted by Saks, Shavit, and Woll [37], Collect is among the standard problems in distributed computing, The model consists of $n$ asynchronous processes, each with a single-writer multi-reader register of a polynomial capacity. The best previously known deterministic solution performs O(n3/2log n) reads and writes, and it is due to Ajtai, Aspnes, Dwork, and Waarts [3]. This paper presents a new deterministic algorithm that performs O(n log7 n) read/write operations, thus substantially improving the best previous upper bound. Using an approach based on epidemic rumor-spreading, the novelty of the new algorithm is in using a family of expander graphs and ensuring that each of the successive groups of processes collect and propagate sufficiently many rumors to the next group. The algorithm is adapted to the Repeatable Collect problem, which is an on-line version. The competitive latency of the new algorithm is O(log7 n) vs. the much higher competitive latency O(√nlog n) given in [3]. A result of independent interest in this paper abstracts a gossiping game that is played on a graph and that gives its payoff in terms of expansion.
Bogdan S. Chlebus, Dariusz R. Kowalski, Alexander A. Schwarzmann
STOC3
2004 Virtual Mobile Nodes for Mobile Ad Hoc Networks
Shlomi Dolev, Seth Gilbert, Nancy A. Lynch, Elad Michael Schiller, Alexander A. Schwarzmann, Jennifer L. Welch
DISC5
2004 The complexity of synchronous iterative Do-All with crashes
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
Distributed Comput.3
2004 Editor's introduction
Alexander A. Schwarzmann
Inf. Comput.1
2003 RAMBO II: Rapidly Reconfigurable Atomic Memory for Dynamic Networks
abstract
Future civilian rescue and military operations will depend on a complex system of communicating devices that can operate in highly dynamic environments. In order to present a consistent view of a complex world, these devices will need to maintain data objects with atomic (linearizable) read/write semantics.
Seth Gilbert, Nancy A. Lynch, Alexander A. Schwarzmann
DSN3
2003 Emulating Shared-Memory Do-All Algorithms in Asynchronous Message-Passing Systems
Dariusz R. Kowalski, Mariam Momenzadeh, Alexander A. Schwarzmann
OPODIS3
2003 Performing work with asynchronous processors: message-delay-sensitive bounds
abstract
This paper considers the problem of performing tasks in asynchronous distributed settings. This problem, called Do-All, has been substantially studied in synchronous models, but there is a dearth of efficient algorithms for asynchronous message-passing processors. Do-All can be trivially solved without any communication by an algorithm where each processor performs all tasks. Assuming p processors and t tasks, this requires work Θ(p · t). Thus it is important to develop subquadratic solutions (when p and t are comparable) by trading computation for communication. Following the observation that it is not possible to obtain subquadratic work when the message delay d is substantial, e.g., d = Θ(t), this work pursues a message-delay-sensitive approach. Here the upper bounds on work and communication are given as functions of p, t, and d, the upper bound on message delays, however algorithms have no knowledge of d and they cannot rely on the existence of an upper bound on d. This paper presents two families of asynchronous algorithms achieving, for the first time, subquadratie work as long as d = o(t). The first family uses as its basis a shared-memory algorithm without having to emulate atomic registers assumed by that algorithm. The second family uses specific permutations of tasks, with certain combinatorial properties, to sequence the work of the processors. Another important contribution in this work is the first delay-sensitive lower bound for this problem that helps explain the behavior of our algorithms.
Dariusz R. Kowalski, Alexander A. Schwarzmann
PODC2
2003 Work-competitive scheduling for cooperative computing with dynamic groups
abstract
The problem of cooperatively performing a set of t tasks in a decentralized setting where the computing medium is subject to failures is one of the fundamental problems in distributed computing. The setting with partitionable networks is especially challenging, as algorithmic solutions must accommodate the possibility that groups of processors become disconnected (and, perhaps, reconnected) during the computation. The efficiency of task-performing algorithms is often assessed in terms of their work: the total number of tasks, counting multiplicities, performed by all of the processors during the computation. In general, an adversary that is able to partition the network into g components can cause any task-performing algorithm to have work Ω(t•g) even if each group of processors performs no more than the optimal number of Θ(t) tasks.Given such pessimistic lower bounds, and in order to understand better the practical implications of performing work in partitionable settings, we study distributed work-scheduling andpursue a competitiveanalysis. Specifically, we study asimple randomized scheduling algorithm for p asynchronous processors, connected by a dynamically changing communication medium, to complete t known tasks. We compare the performance of the algorithm against that of an "off-line" algorithm with full knowledge of the future changes in the communication medium. We describe a notion of computation width, which associates a natural number with a history of changes in the communication medium, and show both upper and lower bounds on competitiveness in terms of this quantity. Specifically, we show that a simple randomized algorithm obtains the competitive ratio (1+cw/e), where cw is computation width; we then show that this ratio is tight.
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
STOC3
2003 GeoQuorums: Implementing Atomic Memory in Mobile Ad Hoc Networks
Shlomi Dolev, Seth Gilbert, Nancy A. Lynch, Alexander A. Schwarzmann, Jennifer L. Welch
DISC4
2003 Efficient Gossip and Robust Distributed Computation
Chryssis Georgiou, Dariusz R. Kowalski, Alexander A. Schwarzmann
DISC3
2002 Failure sensitive analysis for parallel algorithm with controlled memory access concurrency
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
OPODIS3
2002 Optimally work-competitive scheduling for cooperative computing with merging groups
abstract
No abstract available.
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
PODC3
2002 Bounding Work and Communication in Robust Cooperative Computation
Bogdan S. Chlebus, Leszek Gasieniec, Dariusz R. Kowalski, Alexander A. Schwarzmann
DISC4
2002 RAMBO: A Reconfigurable Atomic Memory Service for Dynamic Networks
Nancy A. Lynch, Alexander A. Schwarzmann
DISC2
2002 An inheritance-based technique for building simulation proofs incrementally
abstract
This paper presents a formal technique for incremental construction of system specifications, algorithm descriptions, and simulation proofs showing that algorithms meet their specifications.The technique for building specifications and algorithms incrementally allows a child specification or algorithm to inherit from its parent by two forms of incremental modification: (a) signature extension , where new actions are added to the parent, and (b) specialization (subtyping), where the child's behavior is a specialization (restriction) of the parent's behavior. The combination of signature extension and specialization provides a powerful and expressive incremental modification mechanism for introducing new types of behavior without overriding behavior of the parent; this mechanism corresponds to the subclassing for extension form of inheritance.In the case when incremental modifications are applied to both a parent specification S and a parent algorithm A, the technique allows a simulation proof showing that the child algorithm A′ implements the child specification S′ to be constructed incrementally by extending a simulation proof that algorithm A implements specification S. The new proof involves reasoning about the modifications only, without repeating the reasoning done in the original simulation proof.The paper presents the technique mathematically, in terms of automata. The technique has been used to model and verify a complex middleware system; the methodology and results of that experiment are summarized in this paper.
Idit Keidar, Roger I. Khazan, Nancy A. Lynch, Alexander A. Schwarzmann
ACM Trans. Softw. Eng. Methodol.4
2001 Developing and Refining an Adaptive Token-Passing Strategy
abstract
Token rotation algorithms play an important role in distributed computing, to support such activities as mutual exclusion, round-robin scheduling, group membership and group communication protocols. Ring-based protocols maximize throughput in busy systems but can incur a linear (in the number of processors) delay when a processor needs to obtain a token to perform an operation. This paper synthesizes new algorithmic techniques for improving the performance (responsiveness) of logical ring protocols. The parameterized technique presents the safety properties of ring protocols and maintains high throughput in busy systems, while reducing the delay in lightly loaded systems from a linear to a logarithmic function in the number of processors. The development in this paper is done using term rewriting systems, where our parameterized protocol is developed in a series of safety-preserving refinements of a basic specification.
Burkhard Englert, Larry Rudolph, Alexander A. Schwarzmann
ICDCS3
2001 Local Scheduling for Distributed Cooperation
abstract
The emergence of mobile computing paradigms has created new dimensions for the problem of performing a collection of tasks in a distributed setting. Indeed, an intrinsic feature of mobile computing is that the communication topology changes over time, and some devices may not be able to communicate with others for prolonged periods of time. Efficient utilization of resources in such a setting requires tools for structuring computation with highly variable, or absent, processor connectivity. This article provides a family of efficient distributed scheduling building blocks for this purpose. Specifically, this paper presents new bounds for a fundamental distributed cooperation problem under the assumption that processors may need to schedule their work in isolation due to a prolonged absence of communication. The problem for n processors is defined in terms of t tasks that must be performed efficiently and that are known to all processors. This study gives tight bounds on the ability of the processors to schedule their work so that when some group of processors establish communication, the wasted (redundant) work these processors have collectively performed prior to that time is controlled.
Grzegorz Malewicz, Alexander Russell, Alexander A. Schwarzmann
NCA3
2001 Optimal scheduling for disconnected cooperation
abstract
We consider a distributed environment consisting of n processors that need to perform t tasks. We assume that communication is initially unavailable and that processors begin work in isolation. At some unknown point of time an unknown collection of processors may establish communication. Before processors begin communication they execute tasks in the order given by their schedules. Our goal is to schedule work of isolated processors so that when communication is established for the first time, the number of redundantly executed tasks is controlled. We quantify worst case redundancy as a function of processor advancements through their schedules.
Grzegorz Malewicz, Alexander Russell, Alexander A. Schwarzmann
PODC3
2001 Two Optimization Techniques for Component-Based Systems Deployment
M. Cecilia Bastarrica, Rodrigo E. Caballero, Steven A. Demurjian, Alexander A. Schwarzmann
SEKE4
2001 Optimal Scheduling for Distributed Cooperation Without Communication
Grzegorz Malewicz, Alexander Russell, Alexander A. Schwarzmann
SIROCCO3
2001 Towards practical deteministic write-all algorithms
abstract
The problem of performing t tasks on n asynchronous or undependable processors is a basic problem in parallel and distributed computing. We consider an abstraction of this problem called the Write-All problem— using n processors write 1's into all locations of an array of size t. The most efficient known deterministic asynchronous algorithms for this problem are due to Anderson and Woll. The first class of algorithms has work complexity of Ο(t . n ε), for n ≰ ty and any ε > 0, and they are the best known for the full range of processors (n = t). To schedule the work of the processors, the algorithms use sets of q permutations on [q] (q ≰ n) that have certain combinatorial properties. Instantiating such an algorithm for a specific ε either requires substantial pre-processing (exponential in 1/ε2) to find the requisite permutations, or imposes a prohibitive constant (exponential in 1/ε3) hidden by the asymptotic analysis. The second class deals with the specific case of t = nu, u ≰ 2, and these algorithms have work complexity of Ο(t log t). They also use sets of permutations with the same combinatorial properties. However instantiating these algorithms requires exponential in n preprocessing to find the permutations. To alleviate this costly instantiation Kanellakis and Shvartsman proposed a simple way of computing the permutation schedules. They conjectured that their construction has the desired properties but they provided no analysis.
Bogdan S. Chlebus, Stefan Dobrev, Dariusz R. Kowalski, Grzegorz Malewicz, Alexander A. Schwarzmann, Imrich Vrto
SPAA5
2001 The Complexity of Synchronous Iterative Do-All with Crashes
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
DISC3
2001 Performing tasks on synchronous restartable message-passing processors
Bogdan S. Chlebus, Roberto De Prisco, Alexander A. Schwarzmann
Distributed Comput.3
2001 Specifying and using a partitionable group communication service
abstract
Group communication services are becoming accepted as effective building blocks for the construction of fault-tolerant distributed applications. Many specifications for group communication services have been proposed. However, there is still no agreement about what these specifications should say, especially in cases where the services are partitionable , i.e., where communication failures may lead to simultaneous creation of groups with disjoint memberships, such that each group is unware of the existence of any other group. In this paper, we present a new, succinct specification for a view-oriented partitionable group communication service. The service associates each message with a particular view of the group membership. All send and receive events for a message occur within the associated view. The service provides a total order on the messages within each view, and each processor receives a prefix of this order. Our specification separates safety requirements from performance and fault-tolerance requirements. The safety requirements are expressed by an abstract, global state machine . To present the performance and fault-tolerance requirements, we include failure-status input actions in the specification; we then give properties saying that consensus on the view and timely message delivery are guaranteed in an execution provided that the execution stabilizes to a situation in which the failure-status stops changing and corresponds to consistently partioned system. Because consensus is not required in every execution, the specification is not subject to the existing impossibility results for partionable systems. Our specification has a simple implementation, based on the membership algorithm of Christian and Schmuck. We show the utility of the specification by constructing an ordered-broadcast application, using an algorithm (based on algorithms of Amir, Dolev, Keidar, and others) that reconciles information derived from different instantiations of the group. The application manages the view-change activity to build a shared sequence of messages, i.e., the per-view total orders of the group service are combined to give a universal total order. We prove the correctness and analyze the performance and fault-tolerance of the resulting application.
Alan D. Fekete, Nancy A. Lynch, Alexander A. Schwarzmann
ACM Trans. Comput. Syst.3
2000 Graceful Quorum Reconfiguration in a Robust Emulation of Shared Memory
abstract
Providing shared-memory abstraction in message-passing systems often simplifies the development of distributed algorithms and allows for the reuse of shared-memory algorithms in the message-passing setting. A robust emulation of atomic single-writer/multi-reader registers in message-passing systems was developed by Attiya, Bar-Noy and Dolev (1995). This emulation was extended by Lynch and Shvartsman (1997) to multi-writer/multi-reader registers using reconfigurable quorum systems. In this work we present a new atomic multi-writer/multi-reader register service that includes a fault-tolerant reconfiguration service. This new emulation has a substantially improved performance and fault-tolerance characteristics. We introduce the concept of intermediate quorum configurations and show how they can be used by readers/writers during reconfiguration. The result is that the quorum reconfigurations are graceful: readers and writers no longer "busy-wait" during reconfigurations, bur are able to complete their operations. An additional advance is that the reconfigurer is eliminated as the single point of failure. When the reconfigurer fails, readers and writers continue using intermediate configurations. In finite executions, read and write operations terminate in bounded time using a bounded number of messages (the bounds depend on the "currency" of the configuration at the invoker of the operation). Finally, the service places no restrictions on the installed quorum configuration: a previously installed quorum system can be replaced by an arbitrary new quorum system.
Burkhard Englert, Alexander A. Schwarzmann
ICDCS2
2000 An inheritance-based technique for building simulation proofs incrementally
abstract
This paper presents a technique for incrementally constructing safety specifications, abstract algorithm descriptions, and simulation proofs showing that algorithms meet their specifications.
Idit Keidar, Roger I. Khazan, Nancy A. Lynch, Alexander A. Schwarzmann
ICSE4
2000 The Complexity of Distributed Cooperation in the Presence of Failures
Chryssis Georgiou, Alexander Russell, Alexander A. Schwarzmann
OPODIS3
2000 Distributed cooperation in the absence of communication (brief announcement)
abstract
This work studies a distributed cooperation problem under an extreme assumption that no two processors may be able to communicate during a prolonged period of time. For problems where the quality of distributed decision-making depends on, and can be traded for, communication, the solution space needs to consider the possibility of no communication. Notably, this is the case in the load-balancing setting introduced by Papadimitriou and Yanakakis [PY] and studied by Georgiades, Mavronicolas and Spirakis [GMS]. The distributed cooperation problem that we consider here is defined for n processors in terms of t tasks that need to be performed efficiently and that are known to all processors. The fact that the tasks are initially known makes it possible for the problem to be solved in the absence of communication. The efficiency requirement and the possibility of eventual availability of communication make it desirable to structure the work of the processors so that when eventually some processors are able to communicate, the amount of wasted (redundant) work they have collectively performed prior to that time is controlled. We model solutions to the problem as sets of n lists of distinct tasks from {1,… ,t}. We call such lists schedules. We define and study the notion of k-waste that, for a set of n schedules, measures the maximum number of redundant task identifiers contained in any subset of k (≤ n) schedules. We are interested in expressing k-waste as a function of the length of schedules. Our goal is to construct n schedules of length t such that k-waste is controlled for any prefixes of the schedules.
Grzegorz Malewicz, Alexander Russell, Alexander A. Schwarzmann
PODC3
2000 Cooperative computing with fragmentable and mergeable groups
Chryssis Georgiou, Alexander A. Schwarzmann
SIROCCO2
2000 Distributed Cooperation During the Absence of Communication
Grzegorz Malewicz, Alexander Russell, Alexander A. Schwarzmann
DISC3
1999 An Auction-Based Flexible Pricing Scheme for Renegotiated QoS Connections and Its Evaluation
abstract
This work presents a new renegotiated flexibly-priced communication service. The service is based on a novel auction-based flexible pricing scheme, where the flexibility is expressed in terms of client-supplied budget functions that represent the upper bound on the bandwidth unit price the client is willing to pay, as a function of time. This flexibility allows for the admission control and the bandwidth maintenance policies of the scheme to be decentralized. The policies are implemented by the involved switches in a distributed fashion, such that the clients do not need to participate so long as their budget functions remain competitive. We define four measures of flow satisfaction in terms of the quality and the price. We have performed a comprehensive simulation study of the scheme that confirms its positive qualities. One of the main observations is that clients benefit financially when choosing flexible budget functions, when compared to clients that use fixed budget functions and that achieve similar quality of service. Flexible flows achieving 90% throughput satisfaction can save from 7% to 12% in bandwidth unit price, and flows that may be content with lower satisfaction can save up to 46%.
Grzegorz Malewicz, Alexander A. Schwarzmann
MASCOTS2
1999 A Framework for Architectural Specification of Distributed Object Systems
M. Cecilia Bastarrica, Steven A. Demurjian, Alexander A. Schwarzmann
OPODIS3
1999 Dynamic Load Balancing with Group Communication
Shlomi Dolev, Roberto Segala, Alexander A. Schwarzmann
SIROCCO3
1999 A Dynamic Primary Configuration Group Communication Service
Roberto De Prisco, Alan D. Fekete, Nancy A. Lynch, Alexander A. Schwarzmann
DISC4
1999 Eventually-Serializable Data Services
Alan D. Fekete, David Gupta, Victor Luchangco, Nancy A. Lynch, Alexander A. Schwarzmann
Theor. Comput. Sci.5
1999 Timing Conditions for Linearizability in Uniform Counting Networks
Nancy A. Lynch, Nir Shavit, Alexander A. Schwarzmann, Dan Touitou
Theor. Comput. Sci.3
1998 A Binary Integer Programming Model for Optimal Object Distribution
M. Cecilia Bastarrica, Alexander A. Schwarzmann, Steven A. Demurjian
OPODIS2
1998 Implementing an EventuallySerializable Data Service as a Distributed System Building Block
Oleg M. Cheiner, Alexander A. Schwarzmann
OPODIS2
1998 Implementing and Evaluating an Eventualy-Serializable Data Service
abstract
No abstract available.
Oleg M. Cheiner, Alexander A. Schwarzmann
PODC2
1998 A Dynamic View-Oriented Group Communication Service
abstract
View-oriented group communication services are widely used for fault-tolerant distributed computing.For applications involving coherent data, it is importaut to know when a process has a primary view of the current group membership, usually defined as a view containing a majority out of a static universe of processes.For high availability in a system where processes can join and leave routinely, some researchers have suggested def?.ning primary views dynamically, depending on having enough members in common with recent views.We present a new formal automaton specification, DVS, for the safety guarantees made by a practical group communication service providing a dynamic notion of primary view.We demonstrate the value of DVS by showing both how it can be implemented and how it can be used in an application.First, we present a distributed algorithm based on a group membership algorithm of Lotem, Keidar and Dolev; our version integrates communication with the membership service, uses iuformation from the application processes saying when a view has been prepared for computation by the application, and uses a static view-oriented service internally.We prove that this algorithm implements DVS.Second, we present an application algorithm that is a variant of an algorithm of Amir, Dolev, Keidar, Melliar-Smith and Moser, modified to use DVS instead of a static service.We prove that it implements a (non-group-oriented) totally-orderedbroadcast service.
Roberto De Prisco, Alan D. Fekete, Nancy A. Lynch, Alexander A. Schwarzmann
PODC4
1997 Specifying and Using a Partitionable Group Communication Service
abstract
Group communication services are becoming accepted as effective building blocks for the construction of fault-tolerant distributed applications. Many specifications for group communication services have been proposed. However, there is still no agreement about what these specifications should say, especially in cases where the services are partitionable, i.e., where communication failures may lead to simultaneous creation of groups with disjoint memberships, such that each group is unware of the existence of any other group. In this paper, we present a new, succinct specification for a view-oriented partitionable group communication service. The service associates each message with a particular view of the group membership. All send and receive events for a message occur within the associated view. The service provides a total order on the messages within each view, and each processor receives a prefix of this order. Our specification separates safety requirements from performance and fault-tolerance requirements. The safety requirements are expressed by an abstract, global state machine. To present the performance and fault-tolerance requirements, we include failure-status input actions in the specification; we then give properties saying that consensus on the view and timely message delivery are guaranteed in an execution provided that the execution stabilizes to a situation in which the failure-status stops changing and corresponds to consistently partioned system. Because consensus is not required in every execution, the specification is not subject to the existing impossibility results for partionable systems. Our specification has a simple implementation, based on the membership algorithm of Christian and Schmuck. We show the utility of the specification by constructing an ordered-broadcast application, using an algorithm (based on algorithms of Amir, Dolev, Keidar, and others) that reconciles information derived from different instantiations of the group. The application manages the view-change activity to build a shared sequence of messages, i.e., the per-view total orders of the group service are combined to give a universal total order. We prove the correctness and analyze the performance and fault-tolerance of the resulting application.
Alan D. Fekete, Nancy A. Lynch, Alexander A. Schwarzmann
PODC3
1996 Eventually-Serializable Data Services
abstract
We present a new specification for distributed data services that trade-off immediate consistency guarantees for improved system availability and efficiency, while ensuring the long-term consistency of the data.An eventually-serializable data service maintains the operations requested in a partial order that gravitates over time towards a total order.It provides clear and unambiguous guarantees about the immediate and long-term behavior of the system.To demonstrate its utility, we present an algorithm, based on one of Ladin, Liskov, Shrira, and Ghemawat [12], that implements this specification.Our algorithm provides the interface of the abstract service, and generalizes their algorithm by allowing general operations and greater flexibility in specifying consistency requirements.We also describe how to use this specification as a building block for applications such as directory services.1
Alan D. Fekete, David Gupta, Victor Luchangco, Nancy A. Lynch, Alexander A. Schwarzmann
PODC5
1996 Counting Networks are Practically Linearizable
abstract
Counting networks are a class of concurrent structures that allow the design of highly scalable concurrent data structures in a way that eliminates sequential bottlenecks and contention.Linearizable counting networks assure that the order of the values returned by the network reflects the real-time order in which they were requested.We argue that in many concurrent systems the worst case scenarios that violate linearizability require a form of timing anomaly that is uncommon in practice.The linear time cost of designing networks that achieve linearizability under all circumstances may thus prove an unnecessary burden on applications that are willing to trade-off occasional non-linearizability for speed and parallelism.This paper presents a very simple measure that is iocal to the individual links and nodes of the network, and that quantifies the extent to which a network can suffer from timing anomalies and still remain linearizable.Perhaps counter-intuitively, this measure is independent of network depth.We use our measure to mathematically support our experiment al results: that in a variety of normal situations tested on a simulated shared memory multiprocessor, the Monic counting networks of Aspnes, Herlihy, and Shavit are "for all practical purposes" Iinearizable.
Nancy A. Lynch, Nir Shavit, Alexander A. Schwarzmann, Dan Touitou
PODC3
1994 Efficient Parallelism vs Reliable Distribution: A Trade-off for Concurrent Computations
Paris C. Kanellakis, Dimitrios Michailidis, Alexander A. Schwarzmann
CONCUR3
1993 A Historical Object Base in an Enterprise Management Director
Alexander A. Schwarzmann
Integrated Network Management1
1992 Efficient Parallel Algorithms can be Made Robust
Paris C. Kanellakis, Alexander A. Schwarzmann
Distributed Comput.2
1992 An Efficient Write-All Algorithm for Fail-Stop PRAM Without Initialized Memory
Alexander A. Schwarzmann
Inf. Process. Lett.1
1991 Efficient Parallel Algorithms on Restartable Fail-Stop Processors
abstract
We study efficient deterministic executions of parallel algorithms on restartable fail-stop CRCW PRAMs.We allow the PRAM processors to be subject to arbitrary stop failures and restarts, that are determined by an on-lineThe lower bound also applies to the expected completed work of randomized algorithms that are subject to on-line adversaries.Finally, we desribe a simple on-line adversary that causes inefficiency in many randomized algorithms.
Paris C. Kanellakis, Alexander A. Schwarzmann
PODC2
1991 Achieving Optimal CRCW PRAM Fault-Tolerance
Alexander A. Schwarzmann
Inf. Process. Lett.1
1989 Efficient Parallel Algorithms Can Be Made Robust
abstract
The efficient parallel algorithms proposed for many fundamental problems, such as list ranking, computing preorder numberings and other functions on trees, or integer sorting, are very sensitive to processor failures.The requirement of efficiency (commonly formalized using Parallel-time x Processors as a cost measure) has led to the design of highly tuned PRAM algorithms which, given the additional constraint of simple processor failures, unfortunately become inefficient or even incorrect.We propose a new notion of robustness, that combines efficiency with fault tolerance.For the common case of fail-stop errors, we develop a general (and easy to implement) technique to make robust many efficient parallel algorithms, e.g., algorithms for all the problems listed above.More specifically, for any dynamic pattern of fail-stop errors with at least one surviving processor, our method increases the original algorithm cost by at most a multiplicative factor polylogarithmic in the input size.
Paris C. Kanellakis, Alexander A. Schwarzmann
PODC2