EDBT 2026 Demo / reviewers in the wild / expert
Lewis Tseng
dblp:85/10827
· DBLP profile ↗
51ranked-venue papers
21as first author
21since 2021 · last 2025
0000-0002-4717-4038ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 17 · 7 first-author · 9 since 2021Computer networks · 12 · 3 first-author · 6 since 2021Security and privacy · 6 · 5 first-author · 1 since 2021Software engineering, systems software and programming languages · 6 · 4 first-author · 2 since 2021Theory of computation · 2 · 2 first-authorDatabases, data management, data science and information retrieval · 1 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 1 · 1 since 2021
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2025 | Revisiting State Machine Replication in Practice: Lessons from Building an etcd-inspired SystemabstractState Machine Replication (SMR) is a foundational technique for building fault-tolerant distributed systems. It underpins infrastructure across cloud platforms, databases, and microservices, yet remains surprisingly difficult to implement efficiently in real world. While prior works in both academia and industry technical blogs have explored individual components, such as consensus protocols or deployment techniques, there is still no clear, integrated guide for building high-performance SMR systems end-to-end. Lucas Lebow, Mason Dunkle, Christopher Siems, Jonathan Zarnstorff, Lewis Tseng |
SoCC | 5 |
| 2025 | Content-Aware Gossip Protocol for Improving Access Latency in Mobile Device CloudabstractMobile Edge Computing (MEC) has been proposed to bring content and computation closer to users, which can reduce latency and network load. However, using only MEC is neither scalable nor cost-effective for large-scale services. An alternative solution, namely Mobile Device Clouds (MDCs), was proposed. MDC is a cluster of mobile devices that can share idle computation and storage resources, using device-to-device (D2D) communication. While promising, performance for accessing cached content remains a key challenge in MDC.We propose CAG (Content-Aware Gossip), a lightweight, fully distributed protocol that improves content dissemination through context-aware gossiping, without relying on fixed node identifiers. By decoupling content from specific devices, CAG enables adaptive replication tailored to the dynamic and infrastructure-free nature of in MDC.To guide caching under mobility and resource constraints, we develop a fluid-based system model and formulate an optimization problem balancing replication cost and access latency. Extensive simulations in MATLAB show that CAG reduces average content access latency by over 20% compared to state-of-the-art protocols, demonstrating its effectiveness for scalable and low-latency caching in MDC. Venkatraman Balasubramanian 0002, Lewis Tseng |
GLOBECOM | 2 |
| 2025 | SEED: A Distributed Framework for Multi-Drone Search via Satellite-Edge-Enabled Drones
Lewis Tseng, Vina Dang, Layann Shaban, Wen-Ping Tsai, Moayad Aloqaily |
GLOBECOM | 1 |
| 2025 | Pineapple: Unifying Multi-Paxos and Atomic Shared Registers
Tigran Bantikyan, Jonathan Zarnstorff, Te-Yen Chou, Lewis Tseng, Roberto Palmieri |
NSDI | 4 |
| 2024 | Racos: Improving Erasure Coding State Machine Replication using Leaderless ConsensusabstractCloud storage systems often adopt state machine replication (SMR) to ensure reliability and availability. Most SMR systems use "full-copy" replication across all nodes, which leads to degraded performance for data-intensive workloads, due to high disk and network I/O costs. Erasure coding has recently been integrated with leader-based SMR systems to reduce the costs, e.g., RS-Paxos, CRaft, HRaft, and FRaft. However, these systems still have bottlenecks at the leader, limiting their performance when handling large datasets. Jonathan Zarnstorff, Lucas Lebow, Christopher Siems, Dillon Remuck, Colin Ruiz, Lewis Tseng |
SoCC | 6 |
| 2024 | Fault-tolerant Consensus in Anonymous Dynamic NetworkabstractThis paper studies the feasibility of reaching consensus in an anonymous dynamic network. In our model, n anonymous nodes proceed in synchronous rounds. We adopt a hybrid fault model in which up to$f$nodes may suffer crash or Byzantine faults, and the dynamic message adversary chooses a communication graph for each round. We introduce a stability property of the dynamic network - (T, D)-dynaDegree for$T$≥ 1 and$n$- 1 ≥$D$≥1 - which requires that for every$T$consecutive rounds, any fault-free node must have incoming directed links from at least$D$distinct neighbors. These links might occur in different rounds during a$T$- round interval. (1,$n$-1) -dynaDegree means that the graph is a complete graph in every round. (1, 1) -dynaDegree means that each node has at least one incoming neighbor in every round, but the set of incoming neighbor(s) at each node may change arbitrarily between rounds. We show that exact consensus is impossible even with (1,$n$- 2) -dynaDegree. For an arbitrary T, we show that for crash-tolerant approximate consensus, (T, ⌊ n/2 ⌋)-dynaDegree and$n$> 2$f$are together necessary and sufficient, whereas for Byzantine approximate consensus, (T, ⌊(n + 3f)/2⌋)-dynaDegree and$n$> 5f are together necessary and sufficient. Qinzi Zhang, Lewis Tseng |
ICDCS | 2 |
| 2024 | ALock: Asymmetric Lock Primitive for RDMA SystemsabstractRemote direct memory access (RDMA) networks are being rapidly adopted into industry for their high speed, low latency, and reduced CPU overheads compared to traditional kernel-based TCP/IP networks. RDMA enables threads to access remote memory without interacting with another process. However, atomicity between local accesses and remote accesses is not guaranteed by the technology, hence complicating synchronization significantly. The current solution is to require threads wanting to access local memory in an RDMA-accessible region to pass through the RDMA card using a mechanism known as loopback, but this can quickly degrade performance. In this paper, we introduce ALock, a novel locking primitive designed for RDMA-based systems. ALock allows programmers to synchronize local and remote accesses without using loopback or remote procedure calls (RPCs). We draw inspiration from the classic Peterson's algorithm to create a hierarchical design that includes embedded MCS locks for two cohorts, remote and local. To evaluate the ALock we implement a distributed lock table, measuring throughput and latency in various cluster configurations and workloads. In workloads with a majority of local operations, the ALock outperforms competitors up to 29x and achieves a latency up to 20x faster. Amanda Baran, Jacob Nelson-Slivon, Lewis Tseng, Roberto Palmieri |
SPAA | 3 |
| 2024 | Brief Announcement: Racos: A Leaderless Erasure Coding State Machine ReplicationabstractCloud storage systems often use state machine replication (SMR) to ensure reliability and availability. Erasure coding has recently been integrated with SMR to reduce disk and network I/O costs. This brief announcement shares our experience in developing a leaderless erasure coding SMR system. We integrate our system Racos with etcd, a distributed key-value storage that powers Kubernetes. Racos outperforms competitors by up to 3.36x in throughput. Jonathan Zarnstorff, Lucas Lebow, Dillon Remuck, Colin Ruiz, Lewis Tseng |
SPAA | 5 |
| 2024 | The Power of Abstract MAC Layer: A Fault-Tolerance PerspectiveabstractThis paper studies the power of the "abstract MAC layer" model in a single-hop asynchronous network. The model captures primitive properties of modern wireless MAC protocols. In this model, Newport [PODC '14] proves that it is impossible to achieve deterministic consensus when nodes may crash. Subsequently, Newport and Robinson [DISC '18] present randomized consensus algorithms that terminate with O(n³ log n) expected broadcasts in a system of n nodes. We are not aware of any results on other fault-tolerant distributed tasks in this model. We first study the computability aspect of the abstract MAC layer. We present a wait-free algorithm that implements an atomic register. Furthermore, we show that in general, k-set consensus is impossible. Second, we aim to minimize storage complexity. Existing algorithms require Ω(n log n) bits. We propose two wait-free approximate consensus and two wait-free randomized binary consensus algorithms that only need constant storage complexity (except for the phase index). One randomized algorithm terminates with O(n log n) expected broadcasts. All our algorithms are anonymous, meaning that at the algorithm level, nodes do not need to have a unique identifier. Qinzi Zhang, Lewis Tseng |
DISC | 2 |
| 2024 | Iterative approximate Byzantine consensus in arbitrary directed graphs
Lewis Tseng, Guanfeng Liang, Nitin H. Vaidya |
Distributed Comput. | 1 |
| 2023 | Cryptocurrency meets CAP TheoremabstractGuerraoui et al. [PODC, 2019] implement a cryptocurrency (in a permissioned setting) using an asset transfer object. Their main result implies that consensus is not necessary for implementing a cryptocurrency, since computationally speaking, an asset transfer object is equivalent to an atomic read/write register. In this work, we take a step further to understand fundamental limitations of cryptocurrency under the CAP framework. Particularly, we show that no cryptocurrency that tolerates Byzantine adversary works in a partitioned network. We point out future directions that can circumvent the impossibility. Lewis Tseng, Moayad Aloqaily |
ICBC | 1 |
| 2023 | Poster: Timestamp Verifiability in Proof-of-WorkabstractVarious blockchain systems have been designed for dynamic networked systems. Due to the nature of the systems, the notion of "time" in such systems is somewhat subjective; hence, it is important to understand how the notion of time may impact these systems. This work focuses on an adversary who attacks a Proof-of-Work (POW) blockchain by selfishly constructing an alternative longest chain. We characterize optimal strategies employed by the adversary when a difficulty adjustment rule alà Bitcoin applies. Tzuo Hann Law, Selman Erol, Lewis Tseng |
MobiHoc | 3 |
| 2023 | Byzantine Consensus in Abstract MAC LayerabstractThis paper studies the design of Byzantine consensus algorithms in an \textit{asynchronous }single-hop network equipped with the "abstract MAC layer" [DISC09], which captures core properties of modern wireless MAC protocols. Newport [PODC14], Newport and Robinson [DISC18], and Tseng and Zhang [PODC22] study crash-tolerant consensus in the model. In our setting, a Byzantine faulty node may behave arbitrarily, but it cannot break the guarantees provided by the underlying abstract MAC layer. To our knowledge, we are the first to study Byzantine faults in this model. We harness the power of the abstract MAC layer to develop a Byzantine approximate consensus algorithm and a Byzantine randomized binary consensus algorithm. Both of our algorithms require \textit{only} the knowledge of the upper bound on the number of faulty nodes $f$, and do \textit{not} require the knowledge of the number of nodes $n$. This demonstrates the "power" of the abstract MAC layer, as consensus algorithms in traditional message-passing models require the knowledge of \textit{both} $n$ and $f$. Additionally, we show that it is necessary to know $f$ in order to reach consensus. Hence, from this perspective, our algorithms require the minimal knowledge. The lack of knowledge of $n$ brings the challenge of identifying a quorum explicitly, which is a common technique in traditional message-passing algorithms. A key technical novelty of our algorithms is to identify "implicit quorums" which have the necessary information for reaching consensus. The quorums are implicit because nodes do not know the identity of the quorums -- such notion is only used in the analysis. Lewis Tseng, Callie Sardina |
OPODIS | 1 |
| 2023 | Distributed Multi-writer Multi-reader Atomic Register with Optimistically Fast Read and WriteabstractA distributed multi-writer multi-reader (MWMR) atomic register is an important primitive that enables a wide range of distributed algorithms. Hence, improving its performance can have large-scale consequences. Since the seminal work of ABD emulation in the message-passing networks, many researchers study fast implementations of atomic registers under various conditions. "Fast'' means that a read or a write can be completed with 1 round-trip time (RTT), by contacting a simple majority. In this work, we explore an atomic register with optimal resilience and ''optimistically fast'' read and write operations. That is, both operations can be fast if there is no concurrent write. Lewis Tseng, Neo Zhou, Cole Dumas, Tigran Bantikyan, Roberto Palmieri |
SPAA | 1 |
| 2022 | Reliable Broadcast in Critical Applications: Asset Transfer and Smart HomeabstractAsynchronous Byzantine reliable broadcast receives renewed attention recently, as it is fundamental to many fault-tolerant critical applications. This paper focuses on the Byzantine Reliable Broadcast protocol, which was first proposed by Bracha in 1987. Several recent protocols have improved the round and bit complexity of these algorithms. Motivated by practical network constraints in modern applications, this paper revisits the problem and reduces both complexity in communication and local computation. State-of-the-arts protocols are evaluated using the developed framework that simulates realistic bandwidth constraints. The evaluation demonstrates that our protocols, which use cryptographic hash functions and erasure coding in a novel way, have superior performance in critical applications such as asset transfer and smart home. Yingjian Wu, Yicheng Shen, Haochen Pan, Lewis Tseng, Moayad Aloqaily |
ICC | 4 |
| 2022 | Fault-tolerant Snapshot Objects in Message Passing SystemsabstractThe atomic snapshot object (ASO) can be seen as a generalization of the atomic read/write register. ASO divides the object into$n$segments such that each node can update its own segment, and instantaneously scan all segments of the object. ASO is a powerful data structure that has many important applications, such as update-query state machines, linearizable conflict-free replicated data types, generalized lattice agreement, and cryptocurrency as in the form of an asset transfer object. This paper studies ASO in asynchronous message passing systems and proposes a framework for implementing efficient fault-tolerant snapshot objects. Denote by$D$the maximum message delay and$k$the actual number of failures in an execution. Our framework derives two ASO algorithms: •A crash-tolerant ASO algorithm that achieves O(√k. D) time complexity for both update and scan operations, and achieves amortized constant time operations if there are Ω(√k) operations. •A Byzantine ASO algorithm that achieves O(k.D) time complexity for both update and scan operations, and achieves amortized constant time operations if there is no Byzantine node in a given execution. The framework can also be adapted to implement sequentially consistent snapshot objects (SSO) that complete scan operations locally without any communication, and have the same time complexlty for update onerations as in our ASO algorithms. Vijay K. Garg, Saptaparni Kumar, Lewis Tseng, Xiong Zheng |
IPDPS | 3 |
| 2022 | Brief Announcement: Computability and Anonymous Storage-Efficient Consensus with an Abstract MAC LayerabstractThis paper explores fault-tolerant algorithms in the abstract MAC layer [7] in a single-hop network. The model captures the basic properties of modern wireless MAC protocols. Newport [11] proves that it is impossible to achieve deterministic fault-tolerant consensus, and Newport and Robinson [10] present randomized crash-tolerant consensus algorithms. We are not aware of any study on the computability in this model. Lewis Tseng, Qinzi Zhang |
PODC | 1 |
| 2022 | Brief Announcement: Asymmetric Mutual Exclusion for RDMAabstractCoordinating concurrent access to a shared resource using mutual exclusion is a fundamental problem in computation. In this paper, we present a novel approach to mutual exclusion designed specifically for distributed systems leveraging a popular network communication technology, remote direct memory access (RDMA). Our approach enables local processes to avoid using RDMA operations entirely, limits the number of RDMA operations required by remote processes, and guarantees both starvation-freedom and fairness. Jacob Nelson-Slivon, Lewis Tseng, Roberto Palmieri |
DISC | 2 |
| 2022 | CRACAU: Byzantine Machine Learning Meets Industrial Edge Computing in Industry 5.0abstractIndustry 5.0 is emerging as a result of the advancement in networking and communication technologies, artificial intelligence, distributed computing, and beyond 5G. Among the important enabling technologies, federated learning, industrial edge computing, and Byzantine-tolerant machine learning (ML) are key accelerators in Industry 5.0. We propose a framework to integrate these key components. Recent works have designed various Byzantine-tolerant ML algorithms for a datacenter or a cluster. However, these algorithms are difficult to be applied to industrial edge computing paradigms. In this article, a novel Byzantine-tolerant federated learning algorithm, CRACAU, is designed for the popular three-level edge computing architecture. In this algorithm, edge devices jointly learn an ML model using the data collected at each device, and their private data are never shared with others. Under standard assumptions, we formally prove that CRACAU converges to the optimal point, i.e., CRACAU finds the optimal parameters of the ML model. We also implement CRACAU in the MXNet framework and evaluate it on the popular benchmark MNIST and CIFAR-10 image classification datasets. Experimental results show that CRACAU achieves satisfying accuracy. Anran Du, Yicheng Shen, Qinzi Zhang, Lewis Tseng, Moayad Aloqaily |
IEEE Trans. Ind. Informatics | 4 |
| 2021 | Practical approximate consensus algorithms for small devices in lossy networksabstractThis paper studies a fundamental distributed primitive - approximate consensus - in connected things using wireless networks. It has been extensively studied in different disciplines, such as fault-tolerant computing, distributed computing, control, and robotics communities. To our surprise, we have not found any practical algorithm that is appropriate for our target scenario - a system of small things that have limited computation and storage capability, and use lossy wireless links to communicate with each other. Qinzi Zhang, Tigran Bantikyan, Lewis Tseng |
MobiCom | 3 |
| 2021 | Rabia: Simplifying State-Machine Replication Through RandomizationabstractWe introduce Rabia, a simple and high performance framework for implementing state-machine replication (SMR) within a datacenter. The main innovation of Rabia is in using randomization to simplify the design. Rabia provides the following two features: (i) It does not need any fail-over protocol and supports trivial auxiliary protocols like log compaction, snapshotting, and reconfiguration, components that are often considered the most challenging when developing SMR systems; and (ii) It provides high performance, up to 1.5x higher throughput than the closest competitor (i.e., EPaxos) in a favorable setup (same availability zone with three replicas) and is comparable with a larger number of replicas or when deployed in multiple availability zones. Haochen Pan, Jesse Tuglu, Neo Zhou, Yicheng Shen, Xiong Zheng, Joseph Tassarotti, Lewis Tseng, Roberto Palmieri |
SOSP | 8 |
| 2020 | BBB: A Lightweight Approach to Evaluate Private Blockchains in CloudsabstractEvaluating Blockchain performance is not an easy task. It is difficult to compare different systems, since the evaluation is often incomprehensible and conducted in different environments with distinct workloads. Only a handful of prior tools were proposed, e.g., BLOCKBENCH and HFBench. Unfortunately, these tools have several limitations. We first identify these limitations. Second, motivated by our observations, we then present a benchmarking tool, Boston Blockchain Benchmarking (BBB). BBB is configurable, extensible, and easy-touse. In particular, BBB can be used to test Blockchain from a networking perspective, a feature that we have not observed in prior tools. Similar to BLOCKBENCH, we focus on the private Blockchain. Concretely, we integrate our tool with Mininet, and provide a simple mechanism to test how network properties (e.g., latency, bandwidth, package loss rate) affect the performance of the chosen Blockchain. We present our preliminary result of evaluating Ethereum. We stress that the architecture of BBB is general, and could be extended to other Blockchain systems. BBB is extremely lightweight and can be used on your laptop to test a small network. Such a feature allows quick evaluation of the Blockchain and speeds up innovation and development. Haochen Pan, Xuheng Duan, Yingjian Wu, Lewis Tseng, Moayad Aloqaily, Azzedine Boukerche |
GLOBECOM | 4 |
| 2020 | Efficient and Robust Top-k Algorithms for Big Data IoTabstractTop-k considers as a technique to retrieve, from a hypothetically big data set, only the k (k ≥ 1) best (most relevant/important) candidates. Top-k query processing is a decisive necessity in various collaborative environments that comprise big data such as the Internet of Things (IoT) networks. Particularly, efficient top-k processing in large-scale distributed systems has shown a positively noticeable effect on their performance. This paper considers the distributed approximate top-k processing algorithms dedicated to the IoT-based networks and improve the accuracy of algorithms introduced previously. We then propose a safety-based fault-tolerance notation and contribute to improving a known algorithm in terms of accuracy. Our algorithms have been evaluated using simulation and real-world data and show superiority over conventional methods. Ruifan Yang, Lewis Tseng, Moayad Aloqaily, Azzedine Boukerche |
ICC | 3 |
| 2020 | Semi-Fast Byzantine-tolerant Shared Register without Reliable BroadcastabstractShared register emulations on top of message-passing systems provide an illusion of a simpler shared memory system which can make the task of a system designer easier. Numerous shared register applications have a considerably high read-to-write ratio. Thus, having algorithms that make reads more efficient than writes is a fair trade-off.Typically, such algorithms for reads and writes are asymmetric and sacrifice the stringent consistency condition atomicity, as it is impossible to have fast reads for multi-writer atomicity. Safety is a consistency condition that has has gathered interest from both the systems and theory community as it is weaker than atomicity yet provides strong enough guarantees like "strong consistency" or read-my-write consistency. One requirement that is assumed by many researchers is that of the reliable broadcast (RB) primitive, which ensures the "all or none" property during a broadcast. One drawback is that such a primitive takes 1.5 rounds to complete and requires server-to-server communication.This paper implements an efficient multi-writer multi-reader safe register without using a reliable broadcast primitive. Moreover, we provide fast reads or one-shot reads - our read operations can be completed in one round of client-to-server communication. Of course, this comes with the price of requiring more servers when compared to prior solutions assuming reliable broadcast. However, we show that this increased number of servers is indeed necessary as we prove a tight bound on the number of servers required to implement Byzantine-fault tolerant safe registers in a system without reliable broadcast.We extend our results to data stored using erasure coding as well. We present an emulation of single-writer multi-reader safe register based on MDS codes. The usage of MDS codes reduces storage and communication costs. On the negative side, we also show that to use MDS codes and at the same time achieve one-shot reads, we need even more servers. Kishori M. Konwar, Saptaparni Kumar, Lewis Tseng |
ICDCS | 3 |
| 2020 | Exact Consensus under Global Asymmetric Byzantine LinksabstractFault-tolerant distributed consensus is an important primitive in many large-scale distributed systems and applications. The consensus problem has been investigated under various fault models in the literature since the seminal work by Lamport et al. in 1982. In this paper, we study the exact consensus problem in a new faulty link model, namely global asymmetric Byzantine (GAB) link model. Our link-fault model is simple, yet to our surprise, not studied before. In our system, all the nodes are fault-free and each pair of nodes can communicate directly with each other. In the GAB link model, up to f directed links may become Byzantine, and have arbitrary behavior. Non-faulty links deliver messages reliably. In our model, it is possible that the link from node a to node b is faulty, but the link from node b to node a is fault-free. Unlike some prior models with a local constraint, which enforced a local upper bound on the number of failure links attached to each node, we adopt the global constraint, which allows any link to be corrupted in the GAB model. These global and asymmetric features distinguish our model from all prior faulty link models. In our GAB model, we study the consensus problem in both synchronous and asynchronous systems. We show that 2f + 1 nodes is both necessary and sufficient for solving synchronous consensus, whereas 2f+2 nodes is the tight condition on resilience for solving asynchronous consensus. We also study the models where faulty links are mobile (or transient), i.e., the set of faulty links might change from round to round. We show that 2f + 3 nodes is necessary and sufficient for a family of algorithms that update local state in an iterative fashion. Lewis Tseng, Qinzi Zhang, Saptaparni Kumar |
ICDCS | 1 |
| 2020 | Federated Vehicular Networks: Design, Applications, Routing, and EvaluationabstractIn this paper, we propose a new concept of vehicular networks, namely a Federated Vehicular Networks (FVN), which can be viewed as a stationary vehicular cloud. We first identify the motivation - namely the limits of traditional vehicular clouds due to their instability and rapidly changing topology - and then present the design and applications of an FVN. Finally, we model and describe the unique routing problem in FVN, and present and evaluate different routing algorithms. Jason Posner, Lewis Tseng, Moayad Aloqaily, Mohsen Guizani |
LCN | 2 |
| 2020 | CarML: distributed machine learning in vehicular cloudsabstractThis paper presents CarML, a distributed machine learning platform built on top of an emerging computing paradigm, vehicular clouds. We discuss our design and technical challenges, followed by our preliminary solutions. We verify the efficacy of our solutions using a customized simulator based on Python and SUMO. Anran Du, Yicheng Shen, Lewis Tseng |
MobiCom | 3 |
| 2020 | CassandrEAS: Highly Available and Storage-Efficient Distributed Key-Value Store with Erasure CodingabstractIn this work, we propose an erasure coding-based protocol that implements a key-value store with atomicity and near-optimal storage cost. Our protocol supports concurrent read and write operations while tolerating asynchronous communication and crash failures of any client and some fraction of servers. One novel feature is a tunable knob between the number of supported concurrent operations, availability, and storage cost. We implement our protocol into Cassandra, namely Cassan-drEAS (Cassandra + Erasure-coding Atomic Storage). Extensive evaluation using YCSB on Google Cloud Platform shows that CassandrEAS incurs moderate penalty on latency and throughput, yet saves significant amount of storage space. Viveck R. Cadambe, Kishori M. Konwar, Muriel Médard, Haochen Pan, Lewis Tseng, Yingjian Wu |
NCA | 5 |
| 2020 | Echo-CGC: A Communication-Efficient Byzantine-Tolerant Distributed Machine Learning Algorithm in Single-Hop Radio NetworkabstractIn the past few years, many Byzantine-tolerant distributed machine learning (DML) algorithms have been proposed in the point-to-point communication model. In this paper, we focus on a popular DML framework - the parameter server computation paradigm and iterative learning algorithms that proceed in rounds, e.g., [11, 8, 6]. One limitation of prior algorithms in this domain is the high communication complexity. All the Byzantine-tolerant DML algorithms that we are aware of need to send n d-dimensional vectors from worker nodes to the parameter server in each round, where n is the number of workers and d is the number of dimensions of the feature space (which may be in the order of millions). In a wireless network, power consumption is proportional to the number of bits transmitted. Consequently, it is extremely difficult, if not impossible, to deploy these algorithms in power-limited wireless devices. Motivated by this observation, we aim to reduce the communication complexity of Byzantine-tolerant DML algorithms in the single-hop radio network [1, 3, 14]. Qinzi Zhang, Lewis Tseng |
OPODIS | 2 |
| 2020 | Asynchronous Byzantine Approximate Consensus in Directed NetworksabstractThis paper considers the problem of approximate consensus in directed asynchronous message-passing networks where some nodes may become Byzantine faulty. We obtain a tight necessary and sufficient condition on the underlying directed communication network for asynchronous Byzantine approximate consensus to be achievable. Interestingly, this condition coincides with the tight condition for synchronous Byzantine exact consensus. Our consensus algorithm may be viewed as a non-trivial generalization of an algorithm previously proposed for the special case of complete networks. The tight condition and techniques identified in the paper shed light on the fundamental properties for solving approximate consensus in asynchronous directed networks. Dimitris Sakavalas, Lewis Tseng, Nitin H. Vaidya |
PODC | 2 |
| 2020 | Brief Announcement: Reaching Approximate Consensus When Everyone May CrashabstractFault-tolerant consensus is of great importance in distributed systems. This paper studies the asynchronous approximate consensus problem in the crash-recovery model with fair-loss links. In our model, up to f nodes may crash forever, while the rest may crash intermittently. Each node is equipped with a limited-size persistent storage that does not lose data when crashed. We present an algorithm that only stores three values in persistent storage - state, phase index, and a counter. Lewis Tseng, Qinzi Zhang |
DISC | 1 |
| 2020 | Reliable broadcast with trusted nodes: Energy reduction, resilience, and speed
Lewis Tseng, Yingjian Wu, Haochen Pan, Moayad Aloqaily, Azzedine Boukerche |
Comput. Networks | 1 |
| 2019 | Reliable Broadcast in Networks with Trusted NodesabstractBroadcast is one of the fundamental primitives to enable large-scale networks such as sensor networks and IoT. There is a rich study on achieving reliable broadcast under various kind of failures. In this paper, we use the notion of trust to improve the performance of reliable broadcast. We focus on Certified Propagation Algorithm (CPA), one of the simple algorithms that does not rely on a cryptographic infrastructure and has a proven guarantee on resilience (number of node failures tolerated). Specifically, the paper has two main contributions: (i) A new algorithm Trust-CPA which integrates CPA with trusted nodes has been proposed and shown to increase the resilience from the original CPA, and (ii) A natural optimization problem related to Trust-CPA (i.e., finding the location to place trusted nodes to reduce the broadcast latency) has been proposed as well. We first show that it is NP-hard to find an exact answer and even NP-hard to find a good approximation. A greedy heuristic algorithm has been used and its efficacy has been examined using simulation. We show that our algorithm performs relatively well in geometric random graphs, an appropriate model for large- scale wireless sensor networks. Lewis Tseng, Yingjian Wu, Haochen Pan, Moayad Aloqaily, Azzedine Boukerche |
GLOBECOM | 1 |
| 2019 | Distributed Causal Memory in the Presence of Byzantine ServersabstractWe study distributed causal shared memory (or distributed read/write objects) in the client-server model over asynchronous message-passing networks in which some servers may suffer Byzantine failures. Since Ahamad et al. proposed causal memory in 1994, there have been abundant research on causal storage. Lately, there is a renewed interest in enforcing causal consistency in large-scale distributed storage systems (e.g., COPS, Eiger, Bolt-on). However, to the best of our knowledge, the fault-tolerance aspect of causal memory is not well studied, especially on the tight resilience bound. In our prior work, we showed that 2 f+1 servers is the tight bound to emulate crash-tolerant causal shared memory when up to f servers may crash. In this paper, we adopt a typical model considered in many prior works on Byzantine-tolerant storage algorithms and quorum systems. In the system, up to f servers may suffer Byzantine failures and any number of clients may crash. We constructively present an emulation algorithm for Byzantine causal memory using 3 f+1 servers. We also prove that 3 f+1 is necessary for tolerating up to f Byzantine servers. In other words, we show that 3 f+1 is a tight bound. For evaluation, we implement our algorithm in Golang and compare their performance with two state-of-the-art fault-tolerant algorithms that ensure atomicity in the Google Cloud Platform. Lewis Tseng, Zezhi Wang, Haochen Pan |
NCA | 1 |
| 2019 | Exact Byzantine Consensus on Arbitrary Directed Graphs Under Local Broadcast ModelabstractWe consider Byzantine consensus in a synchronous system where nodes are connected by a network modeled as a directed graph, i.e., communication links between neighboring nodes are not necessarily bi-directional. The directed graph model is motivated by wireless networks wherein asymmetric communication links can occur. In the classical point-to-point communication model, a message sent on a communication link is private between the two nodes on the link. This allows a Byzantine faulty node to equivocate, i.e., send inconsistent information to its neighbors. This paper considers the local broadcast model of communication, wherein transmission by a node is received identically by all of its outgoing neighbors, effectively depriving the faulty nodes of the ability to equivocate. Prior work has obtained sufficient and necessary conditions on undirected graphs to be able to achieve Byzantine consensus under the local broadcast model. In this paper, we obtain tight conditions on directed graphs to be able to achieve Byzantine consensus with binary inputs under the local broadcast model. The results obtained in the paper provide insights into the trade-off between directionality of communication links and the ability to achieve consensus. Muhammad Samir Khan, Lewis Tseng, Nitin H. Vaidya |
OPODIS | 2 |
| 2019 | BBB: Make Benchmarking Blockchains Configurable and ExtensibleabstractInterest in Blockchain technology gradually developed after Bitcoin was introduced in 2008, and has grown exponentially as Blockchain can be applied in many scenarios that require coordination among parties that do not typically trust each other. Due to its popularity and wide applications, enormous number of Blockchain systems have been proposed in the past few years. One way to understand the performance is to build a benchmarking tool to evaluate different systems under different scenarios. A few research teams have build their benchmarking tools, such as BLOCKBENCH and HFBench. Unfortunately, there are several limitations of these tools. In this paper, we enumerate these limitations and challenges. We then present our tool - Boston Blockchain Benchmarking (BBB), which is more configurable and extensible than prior benchmarking tools. Moreover, BBB allows us to evaluate the impact of some common attacks against Blockchains. Finally, we present some preliminary results using BBB. Xuheng Duan, Haochen Pan, Lewis Tseng, Yingjian Wu |
PRDC | 3 |
| 2019 | Resilient Distributed Causal Memory in Client-Server ModelabstractWe study distributed causal shared memory (or key-value pairs) in an asynchronous network under crash failures. Causal memory, introduced by Ahamad et al. in the context of multi-processor environment in 1994, is an abstraction which ensures that nodes agree on the relative ordering of read and write operations that are causally related on key-value pairs. Inspired by the recent interests in geo-replicated causal storage systems (e.g., COPS, Eiger, Bolt-on), we systematically study the fault-tolerance property of the causal shared memory in the client-server model in this work. We identify that 2f + 1 servers is both necessary and sufficient to build a resilient causal memory in the presence of up to f crashed servers. We provide both the necessity proof and a new optimal algorithm that matches the bound. For evaluation, we implement our algorithm in Golang and compare the performance with state-of-the-art fault-tolerant algorithms that ensure strong consistency in the Google Cloud Platform. Lewis Tseng, Zezhi Wang |
PRDC | 1 |
| 2018 | Delivery Delay and Mobile FaultsabstractIn this work we address the problem of reaching approximate consensus in a complete network ofnnodes, where message deliveries can be delayed by at mostdtime-steps. We consider a mobile adversary, which corrupts at mostfnodes in any step, modeled as asynchronousround. We explicitly study howdaffects the feasibility of the problem. More precisely, we propose a framework to analyze mobile fault-tolerance in the presence of message delays. We prove that approximate consensus is feasible if and only ifn> 4df. We assume no knowledge of time (round index) by the nodes; instead, in our model, whenever a message is sent, it is timestamped by the communication channel. We propose the tightTimeStampsalgorithm, which utilizes timestamps to optimally bound the number of faulty messages. Dimitris Sakavalas, Lewis Tseng |
NCA | 2 |
| 2018 | Effects of Topology Knowledge and Relay Depth on Asynchronous Appoximate ConsensusabstractConsider a point-to-point message-passing network. We are interested in the asynchronous crash-tolerant consensus problem in incomplete networks. We study the feasibility and efficiency of approximate consensus under different restrictions on topology knowledge and the relay depth, i.e., the maximum number of hops any message can be relayed. These two constraints are common in large-scale networks, and are used to avoid memory overload and network congestion respectively. Specifically, for positive integer values k and k', we consider that each node knows all its neighbors of at most k-hop distance (k-hop topology knowledge), and the relay depth is k'. We consider both directed and undirected graphs. More concretely, we answer the following question in asynchronous systems: "What is a tight condition on the underlying communication graphs for achieving approximate consensus if each node has only a k-hop topology knowledge and relay depth k'?" To prove that the necessary conditions presented in the paper are also sufficient, we have developed algorithms that achieve consensus in graphs satisfying those conditions: - The first class of algorithms requires k-hop topology knowledge and relay depth k. Unlike prior algorithms, these algorithms do not flood the network, and each node does not need the full topology knowledge. We show how the convergence time and the message complexity of those algorithms is affected by k, providing the respective upper bounds. - The second set of algorithms requires only one-hop neighborhood knowledge, i.e., immediate incoming and outgoing neighbors, but needs to flood the network (i.e., relay depth is n, where n is the number of nodes). One result that may be of independent interest is a topology discovery mechanism to learn and "estimate" the topology in asynchronous directed networks with crash faults. Dimitris Sakavalas, Lewis Tseng, Nitin H. Vaidya |
OPODIS | 2 |
| 2018 | Erasure Coding in Object Stores: Challenges and Opportunities
Lewis Tseng |
PODC | 1 |
| 2018 | Brief Announcement: Effects of Topology Knowledge and Relay Depth on Asynchronous ConsensusabstractConsider an asynchronous incomplete directed network. We study the feasibility and efficiency of approximate crash-tolerant consensus under different restrictions on topology knowledge and relay depth, i.e., the maximum number of hops any message can be relayed. Dimitris Sakavalas, Lewis Tseng, Nitin H. Vaidya |
DISC | 2 |
| 2017 | Towards reliable broadcast in practical sensor networksabstractBroadcast is one of the fundamental primitives to enable sensor networks. In this paper, we address some practical concerns regarding reliable broadcast. Particularly, we consider the following issues: Hybrid fault model: We prove the tight necessary and sufficient condition for using Certified Propagation Algorithm (CPA) to achieve reliable broadcast in directed networks under hybrid fault model. The hybrid fault model considers both node and link failures; Geometric Random Graph: Geometric random graph is considered to be a suitable network model for sensor network deployment. We show that our condition has close relations with connectivity in geometric random graphs in the one-dimensional space. We prove that with a suitable choice of parameters, we can use CPA to achieve reliable broadcast in such graphs; Eventually reliable broadcast: We also study the reliable broadcast problem with relaxed properties. Especially, we explore the trade-off between latency of broadcast and validity (or correctness). For validity, we provide a lower bound on the probability of receiving a correct value in any graphs. For latency, we show that in the geometric random graphs in the one-dimensional space, the proposed algorithm achieves low latency (through simulations). Lewis Tseng |
NCA | 1 |
| 2017 | Voting in the Presence of Byzantine FaultsabstractVoting (or election) algorithms are used widely in many safety-critical systems to mask errors. Most systems only tolerate malicious (or Byzantine) voters - these systems assume the existence of a correct and centralized mechanism to collect the votes and propagate the voting output to each voter. However, in many realistic scenarios, such a centralized voting mechanism is not feasible. Thus, we study the Byzantine voting problem - no centralized mechanism exists in the system, and voters may become Byzantine faulty. We first present impossibility results in both synchronous and asynchronous systems. To circumvent the impossibility results presented in this paper, we propose two relaxed voting properties that are achievable and present optimal voting algorithms that satisfy the relaxed properties. Finally, we show that it is possible to design Byzantine voting algorithms that produce the voting output in one communication step under contention-free scenarios. Lewis Tseng |
PRDC | 1 |
| 2017 | Bitcoin's Consistency PropertyabstractBitcoin is the most popular cryptocurrency nowadays. Inspired by the success, both industry and academia seek to apply Bitcoin's core technique, Blockchain, to other fields like finance, healthcare and Internet-of-Things. One main application is to use Blockchain as the distributed transaction ledger system, a ledger (or a log) of all transactions that is maintained by anonymous participants in a distributed fashion. For the ledger system, consistency is an important property which specifies how the system orders the transactions. Intuitively, Blockchain is designed to maintain a single ground truth - the chain itself is the order of the transactions that all participants should respect. However, we show that under some circumstances, Bitcoin violates eventual consistency, i.e., participants would not converge to a single chain. Thus, we urge a more thorough study on Bitcoin's consistency properties. At the end of the paper, we propose related research directions. Lewis Tseng |
PRDC | 1 |
| 2017 | An Improved Approximate Consensus Algorithm in the Presence of Mobile Faults
Lewis Tseng |
SSS | 1 |
| 2017 | Characterizing and Adapting the Consistency-Latency Tradeoff in Distributed Key-Value StoresabstractThe CAP theorem is a fundamental result that applies to distributed storage systems. In this article, we first present and prove two CAP-like impossibility theorems. To state these theorems, we present probabilistic models to characterize the three important elements of the CAP theorem: consistency (C), availability or latency (A), and partition tolerance (P). The theorems show the un-achievable envelope, that is, which combinations of the parameters of the three models make them impossible to achieve together. Next, we present the design of a class of systems called Probabilistic CAP (PCAP) that perform close to the envelope described by our theorems. In addition, these systems allow applications running on a single data center to specify either a latency Service Level Agreement (SLA) or a consistency SLA. The PCAP systems automatically adapt, in real time and under changing network conditions, to meet the SLA while optimizing the other C/A metric. We incorporate PCAP into two popular key-value stores: Apache Cassandra and Riak. Our experiments with these two deployments, under realistic workloads, reveal that the PCAP systems satisfactorily meets SLAs and perform close to the achievable envelope. We also extend PCAP from a single data center to multiple geo-distributed data centers. Muntasir Raihan Rahman, Lewis Tseng, Indranil Gupta, Nitin H. Vaidya |
ACM Trans. Auton. Adapt. Syst. | 2 |
| 2016 | Recent Results on Fault-Tolerant Consensus in Message-Passing Networks
Lewis Tseng |
SIROCCO | 1 |
| 2015 | Fault-Tolerant Consensus in Directed GraphsabstractConsider a point-to-point network in which nodes are connected by directed links. This paper proves tight necessary and sufficient conditions on the underlying communication graphs for solving the following fault-tolerant consensus problems: Exact crash-tolerant consensus in synchronous systems, Approximate crash-tolerant consensus in asynchronous systems, and Exact Byzantine consensus in synchronous systems. Lewis Tseng, Nitin H. Vaidya |
PODC | 1 |
| 2015 | Broadcast using certified propagation algorithm in presence of Byzantine faults
Lewis Tseng, Nitin H. Vaidya, Vartika Bhandari |
Inf. Process. Lett. | 1 |
| 2014 | Asynchronous convex hull consensus in the presence of crash faultsabstractThis paper defines a new consensus problem, convex hull consensus. The input at each process is a d-dimensional vector of reals (or, equivalently, a point in the d-dimensional Euclidean space), and the output at each process is a convex polytope contained within the convex hull of the inputs at the fault-free processes. We explore the convex hull consensus problem under crash faults with incorrect inputs, and present an asynchronous approximate convex hull consensus algorithm with optimal fault tolerance that reaches consensus on an optimal output polytope. Lewis Tseng, Nitin H. Vaidya |
PODC | 1 |
| 2012 | Iterative approximate byzantine consensus in arbitrary directed graphsabstractThis paper proves a necessary and sufficient condition for the existence of iterative, algorithms that achieve approximate Byzantine consensus in arbitrary directed graphs, where each directed edge represents a communication channel between a pair of nodes. The class of iterative algorithms considered in this paper ensures that, after each iteration of the algorithm, the state of each fault-free node remains in the convex hull of the states of the fault-free nodes at the end of the previous iteration. The following convergence requirement is imposed: for any ε > 0, after a sufficiently large number of iterations, the states of the fault-free nodes are guaranteed to be within ε of each other. Nitin H. Vaidya, Lewis Tseng, Guanfeng Liang |
PODC | 2 |