Idit Keidar

dblp:k/IditKeidar · DBLP profile ↗
← Back
152ranked-venue papers
25as first author
16since 2021 · last 2024
0000-0002-6417-1250ORCID · verified

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

Systems, architecture and hardware · 76 · 15 first-author · 5 since 2021Databases, data management, data science and information retrieval · 14 · 2 first-author · 1 since 2021Theory of computation · 11 · 4 first-authorSecurity and privacy · 10 · 1 first-authorComputer networks · 8 · 1 first-author · 1 since 2021Software engineering, systems software and programming languages · 7 · 2 first-authorApplied, interdisciplinary, general and emerging computing · 5 · 1 since 2021Artificial intelligence and machine learning · 1Graphics, computer vision, multimedia, augmented reality and games · 1
YearPublicationVenuePosition
2024 Expected linear round synchronization: the missing link for linear Byzantine SMR
Oded Naor, Idit Keidar
Distributed Comput.2
2023 Nova: Safe Off-Heap Memory Allocation and Reclamation
abstract
In recent years, we begin to see Java-based systems embrace off-heap allocation for their big data demands. As of today, these system rely on simple ad-hoc garbage-collection solutions, which restrict the usage of off-heap data. This paper introduces the abstraction of safe off-heap memory allocation and reclamation (SOMAR), a thread-safe memory allocation and reclamation scheme for off-heap data in otherwise managed environments. SOMAR allows multi-threaded Java programs to use off-heap memory seamlessly. To realize this abstraction, we present Nova, Novel Off-heap Versioned Allocator, a lock-free SOMAR implementation. Our experiments show that Nova can be used to store off-heap data in Java data structures with better performance than ones managed by Java’s automatic GC. We further integrate Nova into the open-source Oak concurrent map library, which allows Oak to reclaim keys while the data structure is being accessed.
Ramy Fakhoury, Anastasia Braginsky, Idit Keidar, Yoav Zuriel
OPODIS3
2023 Quancurrent: A Concurrent Quantiles Sketch
abstract
Sketches are a family of streaming algorithms widely used in the world of big data to perform fast, real-time analytics. A popular sketch type is Quantiles, which estimates the data distribution of a large input stream. We present Quancurrent, a highly scalable concurrent Quantiles sketch. Quancurrent's throughput increases linearly with the number of available threads, and with 32 threads, it reaches an update speedup of 12x and a query speedup of 30x over a sequential sketch. Quancurrent allows queries to occur concurrently with updates and achieves an order of magnitude better query freshness than existing scalable solutions.
Shaked Elias-Zada, Arik Rinberg, Idit Keidar
SPAA3
2023 Brief Announcement: Subquadratic Multivalued Asynchronous Byzantine Agreement WHP
abstract
There have been several reductions from multivalued consensus to binary consensus over the past 20 years. To the best of our knowledge, none of them solved it for Byzantine asynchronous settings. In this paper, we close this gap. Moreover, we do so in subquadratic communication, using newly developed subquadratic binary Byzantine Agreement techniques.
Shir Cohen, Idit Keidar
DISC2
2023 Cordial Miners: Fast and Efficient Consensus for Every Eventuality
abstract
Cordial Miners are a family of efficient Byzantine Atomic Broadcast protocols, with instances for asynchrony and eventual synchrony. They improve the latency of state-of-the-art DAG-based protocols by almost 2X and achieve optimal good-case complexity of O(n) by forgoing Reliable Broadcast as a building block. Rather, Cordial Miners use the blocklace -- a partially-ordered counterpart of the totally-ordered blockchain data structure -- to implement the three algorithmic components of consensus: Dissemination, equivocation-exclusion, and ordering.
Idit Keidar, Oded Naor, Ouri Poupko, Ehud Shapiro
DISC1
2023 Intermediate Value Linearizability: A Quantitative Correctness Criterion
abstract
Big data processing systems often employ batched updates and data sketches to estimate certain properties of large data. For example, a CountMin sketch approximates the frequencies at which elements occur in a data stream, and a batched counter counts events in batches. This article focuses on correctness criteria for concurrent implementations of such objects. Specifically, we consider quantitative objects, whose return values are from an ordered domain, with a particular emphasis on (ε,δ)-bounded objects that estimate a numerical quantity with an error of at most ε with probability at least 1 - δ. The de facto correctness criterion for concurrent objects is linearizability. Intuitively, under linearizability, when a read overlaps an update, it must return the object’s value either before the update or after it. Consider, for example, a single batched increment operation that counts three new events, bumping a batched counter’s value from 7 to 10. In a linearizable implementation of the counter, a read overlapping this update must return either 7 or 10. We observe, however, that in typical use cases, any intermediate value between 7 and 10 would also be acceptable. To capture this additional degree of freedom, we propose Intermediate Value Linearizability (IVL) , a new correctness criterion that relaxes linearizability to allow returning intermediate values, for instance, 8 in the example above. Roughly speaking, IVL allows reads to return any value that is bounded between two return values that are legal under linearizability. A key feature of IVL is that we can prove that concurrent IVL implementations of (ε,δ)-bounded objects are themselves (ε,δ)-bounded. To illustrate the power of this result, we give a straightforward and efficient concurrent implementation of an (ε,δ)-bounded CountMin sketch, which is IVL (albeit not linearizable). We present four examples for IVL objects, each showcasing a different way of using IVL. The first is a simple wait-free IVL batched counter, with O (1) step complexity for update. The next considers an (ε,δ)-bounded CountMin sketch and further shows how to relax IVL using the notion of r -relaxation. Our third example is a non-atomic iterator over a data structure. In this example, we augment the data structure with an auxiliary history variable state that includes “tombstones” for items deleted from the data structure. Here, IVL semantics are required at the augmented level. Finally, using a priority queue , we show that some objects require IVL to be paired with other correctness criteria; indeed, a natural correctness notion for a concurrent priority queue is IVL coupled with sequential consistency. Last, we show that IVL allows for inherently cheaper implementations than linearizable ones. In particular, we show a lower bound of Ω ( n ) on the step complexity of the update operation of any wait-free linearizable batched counter from single-writer multi-reader registers, which is more expensive than our O (1) IVL implementation.
Arik Rinberg, Idit Keidar
J. ACM2
2022 SwiSh: Distributed Shared State Abstractions for Programmable Switches
Lior Zeno, Dan R. K. Ports, Jacob Nelson 0001, Daehyeok Kim, Shir Landau Feibish, Idit Keidar, Arik Rinberg, Alon Rashelbach, Igor Lima de Paula, Mark Silberstein
NSDI6
2022 Make Every Word Count: Adaptive Byzantine Agreement with Fewer Words
abstract
Byzantine Agreement (BA) is a key component in many distributed systems. While Dolev and Reischuk have proven a long time ago that quadratic communication complexity is necessary for worst-case runs, the question of what can be done in practically common runs with fewer failures remained open. In this paper we present the first Byzantine Broadcast algorithm with O(n(f+1)) communication complexity in a model with resilience of n = 2t+1, where 0 ≤ f ≤ t is the actual number of process failures in a run. And for BA with strong unanimity, we present the first optimal-resilience algorithm that has linear communication complexity in the failure-free case and a quadratic cost otherwise.
Shir Cohen, Idit Keidar, Alexander Spiegelman
OPODIS2
2022 Brief Announcement: Make Every Word Count: Adaptive Byzantine Agreement with Fewer Words
abstract
Byzantine Agreement (BA) is a key component in many distributed systems. While Dolev and Reischuk have proven a long time ago that quadratic word complexity is necessary for worst-case runs, the question of what can be done in practically common runs with fewer failures remained open. In this paper we present the first Byzantine Broadcast algorithm with O(n(f+1)) word complexity in a model with resilience of n=2t+1, where 0 ≤ f ≤ t is the actual number of process failures in a run.
Shir Cohen, Idit Keidar, Alexander Spiegelman
PODC2
2022 On Payment Channels in Asynchronous Money Transfer Systems
abstract
Money transfer is an abstraction that realizes the core of cryptocurrencies. It has been shown that, contrary to common belief, money transfer in the presence of Byzantine faults can be implemented in asynchronous networks and does not require consensus. Nonetheless, existing implementations of money transfer still require a quadratic message complexity per payment, making attempts to scale hard. In common blockchains, such as Bitcoin and Ethereum, this cost is mitigated by payment channels implemented as a second layer on top of the blockchain allowing to make many off-chain payments between two users who share a channel. Such channels require only on-chain transactions for channel opening and closing, while the intermediate payments are done off-chain with constant message complexity. But payment channels in-use today require synchrony; therefore, they are inadequate for asynchronous money transfer systems. In this paper, we provide a series of possibility and impossibility results for payment channels in asynchronous money transfer systems. We first prove a quadratic lower bound on the message complexity of on-chain transfers. Then, we explore two types of payment channels, unidirectional and bidirectional. We define them as shared memory abstractions and prove that in certain cases they can be implemented as a second layer on top of an asynchronous money transfer system whereas in other cases it is impossible.
Oded Naor, Idit Keidar
DISC2
2022 DSON: JSON CRDT Using Delta-Mutations For Document Stores
abstract
We propose DSON, a space efficient δ-based CRDT approach for distributed JSON document stores, enabling high availability at a global scale, while providing strong eventual consistency guarantees. We define the semantics of our CRDT based approach formally, and prove its correctness and convergence. Previous approaches optimize for collaborative document editing and store metadata proportional to the number of updates to a document, which is not acceptable for long lived document management. The metadata stored with our approach is bounded by O ( k 2 D + n log n ), where n is the number of replicas, D is the number of document elements, and k ≤ n is the number of concurrent document updates. We also implement our approach[37] and demonstrate its space efficiency empirically. Experimental analysis shows that the metadata stored is typically significantly less than the worst case. This provides the basis for robust highly available distributed document stores with well defined semantics and safety guarantees, relieving application developers from the burden of conflict resolution.
Arik Rinberg, Tomer Solomon, Roee Shlomo, Guy Khazma, Gal Lushi, Idit Keidar, Paula Ta-Shma
Proc. VLDB Endow.6
2021 Game of Coins
abstract
The cryptocurrency market is blooming. Tens of new coins emerge every year and their total market cap keeps growing. The research community is trying to keep up by proposing improved mining protocols and attacking existing ones. However, surprisingly as it may sound, most existing works overlook the real-life multi-coin market, by focusing on a system with a single coin. To the best of our knowledge, this paper was the first to consider a system with many coins and strategic miners that are free to choose where to mine. We first formalize the current practice of strategic mining in multi-coin markets as a singleton weighted congestion game and prove that any better-response dynamics in such a game converges to an equilibrium. Then, in our main result, we present a reward design attack that moves the system configuration from any initial equilibrium to a desired one. The attack is executed via temporary manipulation of coin rewards, which leads strategic miners to switch between coins. It applies to any better-response dynamics of the miners. To motivate our attack we show that in any equilibrium there is always at least one miner whose profit is higher in a different one, and thus may benefit from such an attack.
Alexander Spiegelman, Idit Keidar, Moshe Tennenholtz
ICDCS2
2021 Using Nesting to Push the Limits of Transactional Data Structure Libraries
Gal Assa, Hagar Meir, Guy Golan-Gueta, Idit Keidar, Alexander Spiegelman
OPODIS4
2021 All You Need is DAG
abstract
We present DAG-Rider, the first asynchronous Byzantine Atomic Broadcast protocol that achieves optimal resilience, optimal amortized communication complexity, and optimal time complexity. DAG-Rider is post-quantum safe and ensures that all values proposed by correct processes eventually get delivered. We construct DAG-Rider in two layers: In the first layer, processes reliably broadcast their proposals and build a structured Directed Acyclic Graph (DAG) of the communication among them. In the second layer, processes locally observe their DAGs and totally order all proposals with no extra communication.
Idit Keidar, Eleftherios Kokoris-Kogias, Oded Naor, Alexander Spiegelman
PODC1
2021 Brief Announcement: Using Nesting to Push the Limits of Transactional Data Structure Libraries
abstract
Transactional data structure libraries (TDSL) combine the ease-of-programming of transactions with the high performance and scalability of custom-tailored concurrent data structures. They can be very efficient thanks to their ability to exploit data structure semantics in order to reduce overhead, aborts, and wasted work compared to general-purpose software transactional memory. However, TDSLs were not previously used for complex use-cases involving long transactions and a variety of data structures. In this paper, we boost the performance and usability of a TDSL, towards allowing it to support complex applications. A key idea is nesting. Nested transactions create checkpoints within a longer transaction, so as to limit the scope of abort, without changing the semantics of the original transaction. We build a Java TDSL with built-in support for nested transactions over a number of data structures. We conduct a case study of a complex network intrusion detection system that invests a significant amount of work to process each packet. Our study shows that our library outperforms publicly available STMs twofold without nesting, and by up to 16x when nesting is used.
Gal Assa, Hagar Meir, Guy Golan-Gueta, Idit Keidar, Alexander Spiegelman
DISC4
2021 Tame the Wild with Byzantine Linearizability: Reliable Broadcast, Snapshots, and Asset Transfer
abstract
We formalize Byzantine linearizability, a correctness condition that specifies whether a concurrent object with a sequential specification is resilient against Byzantine failures. Using this definition, we systematically study Byzantine-tolerant emulations of various objects from registers. We focus on three useful objects- reliable broadcast, atomic snapshot, and asset transfer. We prove that there exist n-process f-resilient Byzantine linearizable implementations of such objects from registers if and only if f < n/2.
Shir Cohen, Idit Keidar
DISC2
2020 EvenDB: optimizing key-value storage for spatial locality
abstract
Applications of key-value (KV-)storage often exhibit high spatial locality, such as when many data items have identical composite key prefixes. This prevalent access pattern is underused by the ubiquitous LSM design underlying high-throughput KV-stores today.
Eran Gilad, Edward Bortnikov, Anastasia Braginsky, Yonatan Gottesman, Eshcar Hillel, Idit Keidar, Nurit Moscovici, Rana Shahout
EuroSys6
2020 Byzantine Agreement and SMR with Sub-Quadratic Message Complexity (Invited Talk)
abstract
Byzantine Agreement (BA) has been studied for four decades by now, but until recently, has been considered at a fairly small scale. In recent years, however, we begin to see practical use-cases of BA in large-scale systems, which motivates a push for reduced communication complexity. Dolev and Reischuk’s well-known lower bound stipulates that any deterministic algorithm requires Ω(n²) communication in the worst-case, and until fairly recently, almost all randomized algorithms have had at least quadratic complexity as well. This talk will present two new algorithms breaking this barrier. The first part of the talk will consider a fully asynchronous setting, focusing on randomized BA whose safety and liveness guarantees hold with high probability. It will present the first asynchronous Byzantine Agreement algorithm with sub-quadratic communication complexity. This algorithm exploits VRF-based committee sampling, which it adapts for the asynchronous model. The second part of the talk will consider the eventually synchronous model, where BA and State Machine Replication (SMR) can be solved with deterministic safety and liveness guarantees. In this context, randomization is used in order to reduce the expected communication complexity. The talk will present an algorithm for round synchronization, which is a building block for BA and SMR and constitutes the main performance bottleneck therein. It will present an algorithm that, for the first time, achieves round synchronization with expected linear message complexity and expected constant latency. Existing protocols can use this round synchronization algorithm to solve Byzantine SMR with the same asymptotic performance. The first part of the talk is based on joint work with Shir Cohen and Alexander Spiegelman, and the second part of the talk is based on joint work with Oded Naor.
Idit Keidar
OPODIS1
2020 Brief Announcement: Not a COINcidence: Sub-Quadratic Asynchronous Byzantine Agreement WHP
abstract
King and Saia were the first to break the quadratic word complexity bound for Byzantine Agreement in synchronous systems against an adaptive adversary, and Algorand broke this bound with near-optimal resilience (first in the synchronous model and then with eventual-synchrony). Yet the question of asynchronous sub-quadratic Byzantine Agreement remained open. To the best of our knowledge, we are the first to answer this question in the affirmative. A key component of our solution is a shared coin algorithm based on a VRF. A second essential ingredient is VRF-based committee sampling, which we formalize and utilize in the asynchronous model for the first time. Our algorithms work against a delayed-adaptive adversary, which cannot perform after-the-fact removals but has full control of Byzantine processes and full information about communication in earlier rounds. Using committee sampling and our shared coin, we solve Byzantine Agreement with high probability, with a word complexity of Õ(n) and O(1) expected time, breaking the O(n2) bit barrier for asynchronous Byzantine Agreement.
Shir Cohen, Idit Keidar, Alexander Spiegelman
PODC2
2020 Brief Announcement: Intermediate Value Linearizability: A Quantitative Correctness Criterion
abstract
A common correctness criterion for concurrent objects is linearizability. Intuitively, under linearizability, when a read overlaps an update, it must return either the object's value before the update or the value after it. Consider, for example, a batched counter supporting "batched" increments, and a single operation that bumps its value from 7 to 10. A read overlapping this update is allowed to return either 7 or 10. In this paper, we propose Intermediate Value Linearizability (IVL), a new correctness criterion that relaxes linearizability to allow returning intermediate values, for instance, 8 in the example above. IVL is applicable to objects whose return values are from a totally ordered set. Roughly speaking, it allows reads to return any value that is bounded between two return values that are legal under linearizability. We show that this added degree of freedom inherently allows for cheaper implementations than linearizability. In particular, we show a lower bound of Ω(n) on the step complexity of the update operation of a wait-free linearizable batched counter, and give a wait-free IVL implementation of the same object with an O(1) step complexity for update.
Arik Rinberg, Idit Keidar
PODC2
2020 Nesting and composition in transactional data structure libraries
abstract
Transactional data structure libraries (TDSL) combine the ease-of-programming of transactions with the high performance and scalability of custom-tailored concurrent data structures. They can be very efficient thanks to their ability to exploit data structure semantics in order to reduce overhead, aborts, and wasted work compared to general-purpose software transactional memory. However, TDSLs were not previously used for complex use-cases involving long transactions and a variety of data structures.
Gal Assa, Hagar Meir, Guy Golan-Gueta, Idit Keidar, Alexander Spiegelman
PPoPP4
2020 Oak: a scalable off-heap allocated key-value map
abstract
Efficient ordered in-memory key-value (KV-)maps are paramount for the scalability of modern data platforms. In managed languages like Java, KV-maps face unique challenges due to the high overhead of garbage collection (GC).
Hagar Meir, Dmitry Basin, Edward Bortnikov, Anastasia Braginsky, Yonatan Gottesman, Idit Keidar, Eran Meir, Gali Sheffi, Yoav Zuriel
PPoPP6
2020 Fast concurrent data sketches
abstract
Data sketches are approximate succinct summaries of long data streams. They are widely used for processing massive amounts of data and answering statistical queries about it. Existing libraries producing sketches are very fast, but do not allow parallelism for creating sketches using multiple threads or querying them while they are being built. We present a generic approach to parallelising data sketches efficiently and allowing them to be queried in real time, while bounding the error that such parallelism introduces. Utilising relaxed semantics and the notion of strong linearisability we prove our algorithm's correctness and analyse the error it induces in two specific sketches. Our implementation achieves high scalability while keeping the error small. We have contributed one of our concurrent sketches to the open-source data sketches library.
Arik Rinberg, Alexander Spiegelman, Edward Bortnikov, Eshcar Hillel, Idit Keidar, Lee Rhodes, Hadar Serviansky
PPoPP5
2020 Scalable top-k retrieval with Sparta
abstract
Many big data processing applications rely on a top-k retrieval building block, which selects (or approximates) the k highest-scoring data items based on an aggregation of features. In web search, for instance, a document's score is the sum of its scores for all query terms. Top-k retrieval is often used to sift through massive data and identify a smaller subset of it for further analysis. Because it filters out the bulk of the data, it often constitutes the main performance bottleneck.
Gali Sheffi, Dmitry Basin, Edward Bortnikov, David Carmel, Idit Keidar
PPoPP5
2020 Not a COINcidence: Sub-Quadratic Asynchronous Byzantine Agreement WHP
abstract
King and Saia were the first to break the quadratic word complexity bound for Byzantine Agreement in synchronous systems against an adaptive adversary, and Algorand broke this bound with near-optimal resilience (first in the synchronous model and then with eventual-synchrony). Yet the question of asynchronous sub-quadratic Byzantine Agreement remained open. To the best of our knowledge, we are the first to answer this question in the affirmative. A key component of our solution is a shared coin algorithm based on a VRF. A second essential ingredient is VRF-based committee sampling, which we formalize and utilize in the asynchronous model for the first time. Our algorithms work against a delayed-adaptive adversary, which cannot perform after-the-fact removals but has full control of Byzantine processes and full information about communication in earlier rounds. Using committee sampling and our shared coin, we solve Byzantine Agreement with high probability, with a word complexity of $\widetilde{O}(n)$ and $O(1)$ expected time, breaking the $O(n^2)$ bit barrier for asynchronous Byzantine Agreement.
Shir Cohen, Idit Keidar, Alexander Spiegelman
DISC2
2020 Expected Linear Round Synchronization: The Missing Link for Linear Byzantine SMR
abstract
State Machine Replication (SMR) solutions often divide time into rounds, with a designated leader driving decisions in each round. Progress is guaranteed once all correct processes synchronize to the same round, and the leader of that round is correct. Recently suggested Byzantine SMR solutions such as HotStuff, Tendermint, and LibraBFT achieve progress with a linear message complexity and a constant time complexity once such round synchronization occurs. But round synchronization itself incurs an additional cost. By Dolev and Reischuk’s lower bound, any deterministic solution must have Ω(n²) communication complexity. Yet the question of randomized round synchronization with an expected linear message complexity remained open. We present an algorithm that, for the first time, achieves round synchronization with expected linear message complexity and expected constant latency. Existing protocols can use our round synchronization algorithm to solve Byzantine SMR with the same asymptotic performance.
Oded Naor, Idit Keidar
DISC2
2020 Intermediate Value Linearizability: A Quantitative Correctness Criterion
abstract
Big data processing systems often employ batched updates and data sketches to estimate certain properties of large data. For example, a CountMin sketch approximates the frequencies at which elements occur in a data stream, and a batched counter counts events in batches. This paper focuses on correctness criteria for concurrent implementations of such objects. Specifically, we consider quantitative objects, whose return values are from a totally ordered domain, with a particular emphasis on (ε,δ)-bounded objects that estimate a numerical quantity with an error of at most ε with probability at least 1 - δ. The de facto correctness criterion for concurrent objects is linearizability. Intuitively, under linearizability, when a read overlaps an update, it must return the object’s value either before the update or after it. Consider, for example, a single batched increment operation that counts three new events, bumping a batched counter’s value from 7 to 10. In a linearizable implementation of the counter, a read overlapping this update must return either 7 or 10. We observe, however, that in typical use cases, any intermediate value between 7 and 10 would also be acceptable. To capture this additional degree of freedom, we propose Intermediate Value Linearizability (IVL), a new correctness criterion that relaxes linearizability to allow returning intermediate values, for instance 8 in the example above. Roughly speaking, IVL allows reads to return any value that is bounded between two return values that are legal under linearizability. A key feature of IVL is that we can prove that concurrent IVL implementations of (ε,δ)-bounded objects are themselves (ε,δ)-bounded. To illustrate the power of this result, we give a straightforward and efficient concurrent implementation of an (ε, δ)-bounded CountMin sketch, which is IVL (albeit not linearizable). Finally, we show that IVL allows for inherently cheaper implementations than linearizable ones. In particular, we show a lower bound of Ω(n) on the step complexity of the update operation of any wait-free linearizable batched counter from single-writer objects, and propose a wait-free IVL implementation of the same object with an O(1) step complexity for update.
Arik Rinberg, Idit Keidar
DISC2
2019 FairLedger: A Fair Blockchain Protocol for Financial Institutions
abstract
Financial institutions are currently looking into technologies for permissioned blockchains. A major effort in this direction is Hyperledger, an open source project hosted by the Linux Foundation and backed by a consortium of over a hundred companies. A key component in permissioned blockchain protocols is a byzantine fault tolerant (BFT) consensus engine that orders transactions. However, currently available BFT solutions in Hyperledger (as well as in the literature at large) are inadequate for financial settings; they are not designed to ensure fairness or to tolerate selfish behavior that arises when financial institutions strive to maximize their own profit. We present FairLedger, a permissioned blockchain BFT protocol, which is fair, designed to deal with rational behavior, and, no less important, easy to understand and implement. The secret sauce of our protocol is a new communication abstraction, called detectable all-to-all (DA2A), which allows us to detect participants (byzantine or rational) that deviate from the protocol, and punish them. We implement FairLedger in the Hyperledger open source project, using Iroha framework, one of the biggest projects therein. To evaluate FairLegder's performance, we also implement it in the PBFT framework and compare the two protocols. Our results show that in failure-free scenarios FairLedger achieves better throughput than both Iroha's implementation and PBFT in wide-area settings.
Kfir Lev-Ari, Alexander Spiegelman, Idit Keidar, Dahlia Malkhi
OPODIS3
2019 2019 Edsger W. Dijkstra Prize in Distributed Computing
abstract
The committee decided to award the 2019 Edsger W. Dijkstra Prize in Distributed Computing to Alessandro Panconesi and Aravind Srinivasan for their paper Randomized Distributed Edge Coloring via an Extension of the Chernoff-Hoeffding Bounds, SIAM Journal on Computing, volume 26, number 2, 1997, pages 350-368. A preliminary version of this paper appeared as Fast Randomized Algorithms for Distributed Edge Coloring, Proceedings of the Eleventh Annual ACM Symposium Principles of Distributed Computing (PODC), 1992, pages 251-262.
Lorenzo Alvisi, Shlomi Dolev, Faith Ellen, Idit Keidar, Fabian Kuhn, Jukka Suomela
PODC4
2019 Fast Concurrent Data Sketches
abstract
Data sketches are approximate succinct summaries of long data streams. They are widely used for processing massive amounts of data and answering statistical queries about it. Existing libraries producing sketches are very fast, but do not allow parallelism for creating sketches using multiple threads or querying them while they are being built. We present a generic approach to parallelising data sketches efficiently and allowing them to be queried in real time, while bounding the error that such parallelism introduces. Utilising relaxed semantics and the notion of strong linearisability we prove our algorithm's correctness and analyse the error it induces in two specific sketches. Our implementation achieves high scalability while keeping the error small. We have contributed one of our concurrent sketches to the open-source data sketches library.
Arik Rinberg, Alexander Spiegelman, Edward Bortnikov, Eshcar Hillel, Idit Keidar, Hadar Serviansky
PODC5
2018 2018 Edsger W. Dijkstra Prize in Distributed Computing
abstract
The Dijkstra Prize Committee has decided to grant the 2018 Edsger W. Dijkstra Prize in Distributed Computing to Bowen Alpern and Fred B. Schneider for their paper:
Yehuda Afek, Idit Keidar, Boaz Patt-Shamir, Sergio Rajsbaum, Ulrich Schmid 0001, Gadi Taubenfeld
PODC2
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
PODC2
2018 Session details: Session 1A: Persistent Memory
Idit Keidar
PODC1
2018 Session details: Session 2A: Approximation and Learning
Idit Keidar
PODC1
2018 Session details: Session 3A: Congest
Idit Keidar
PODC1
2018 Integrated Bounds for Disintegrated Storage
abstract
We point out a somewhat surprising similarity between non-authenticated Byzantine storage, coded storage, and certain emulations of shared registers from smaller ones. A common characteristic in all of these is the inability of reads to safely return a value obtained in a single atomic access to shared storage. We collectively refer to such systems as disintegrated storage, and show integrated space lower bounds for asynchronous regular wait-free emulations in all of them. In a nutshell, if readers are invisible, then the storage cost of such systems is inherently exponential in the size of written values; otherwise, it is at least linear in the number of readers. Our bounds are asymptotically tight to known algorithms, and thus justify their high costs.
Alon Berger, Idit Keidar, Alexander Spiegelman
DISC2
2018 Accordion: Better Memory Organization for LSM Key-Value Stores
abstract
Log-structured merge (LSM) stores have emerged as the technology of choice for building scalable write-intensive key-value storage systems. An LSM store replaces random I/O with sequential I/O by accumulating large batches of writes in a memory store prior to flushing them to log-structured disk storage; the latter is continuously re-organized in the background through a compaction process for efficiency of reads. Though inherent to the LSM design, frequent compactions are a major pain point because they slow down data store operations, primarily writes, and also increase disk wear. Another performance bottleneck in today's state-of-the-art LSM stores, in particular ones that use managed languages like Java, is the fragmented memory layout of their dynamic memory store. In this paper we show that these pain points may be mitigated via better organization of the memory store. We present Accordion - an algorithm that addresses these problems by re-applying the LSM design principles to memory management. Accordion is implemented in the production code of Apache HBase, where it was extensively evaluated. We demonstrate Accordion's double-digit performance gains versus the baseline HBase implementation and discuss some unexpected lessons learned in the process.
Edward Bortnikov, Anastasia Braginsky, Eshcar Hillel, Idit Keidar, Gali Sheffi
Proc. VLDB Endow.4
2018 Taking Omid to the Clouds: Fast, Scalable Transactions for Real-Time Cloud Analytics
abstract
We describe how we evolve Omid, a transaction processing system for Apache HBase, to power Apache Phoenix, a cloud-grade real-time SQL analytics engine. Omid was originally designed for data processing pipelines at Yahoo, which are, by and large, throughput-oriented monolithic NoSQL applications. Providing a platform to support converged real-time transaction processing and analytics applications - dubbed translytics - introduces new functional and performance requirements. For example, SQL support is key for developer productivity, multi-tenancy is essential for cloud deployment, and latency is cardinal for just-in-time data ingestion and analytics insights. We discuss our efforts to adapt Omid to these new domains, as part of the process of integrating it into Phoenix as the transaction processing backend. A central piece of our work is latency reduction in Omid's protocol, which also improves scalability. Under light load, the new protocol's latency is 4x to 5x smaller than the legacy Omid's, whereas under increased loads it is an order of magnitude faster. We further describe a fast path protocol for single-key transactions, which enables processing them almost as fast as native HBase operations.
Ohad Shacham, Yonatan Gottesman, Aran Bergman, Edward Bortnikov, Eshcar Hillel, Idit Keidar
Proc. VLDB Endow.6
2017 Fishing in the stream: Similarity search over endless data
abstract
Similarity search is the task of retrieving data items that are similar to a given query. In this paper, we introduce the time-sensitive notion of similarity search over endless data-streams (SSDS), which takes into account data quality and temporal characteristics in addition to similarity. SSDS is challenging as it needs to process unbounded data, while computation resources are bounded. We propose Stream-LSH, a randomized SSDS algorithm that bounds the index size by retaining items according to their freshness, quality, and dynamic popularity attributes. We show that Stream-LSH increases recall when searching for similar items compared to alternative approaches using the same space capacity.
Naama Kraus, David Carmel, Idit Keidar
IEEE BigData3
2017 Fragola: low-latency transactions in distributed data stores
abstract
As transaction processing services begin to be used in new application domains, low transaction latency becomes an important consideration. Motivated by such use cases we developed Fragola, a highly scalable low-latency and high-throughput transaction processing engine for Apache HBase. Similarly to other modern transaction managers, Fragola provides a variant of generalized snapshot isolation (SI), which scales better than traditional serializability implementations.
Yonatan Gottesman, Aran Bergman, Edward Bortnikov, Eshcar Hillel, Idit Keidar, Ohad Shacham
SoCC5
2017 Omid, Reloaded: Scalable and Highly-Available Transaction Processing
Edward Bortnikov, Eshcar Hillel, Idit Keidar, Ivan Kelly, Matthieu Morel, Sameer Paranjpye, Francisco Perez-Sorrosal, Ohad Shacham
FAST3
2017 KiWi: A Key-Value Map for Scalable Real-Time Analytics
abstract
Modern big data processing platforms employ huge in-memory key-value (KV) maps. Their applications simultaneously drive high-rate data ingestion and large-scale analytics. These two scenarios expect KV-map implementations that scale well with both real-time updates and large atomic scans triggered by range queries.
Dmitry Basin, Edward Bortnikov, Anastasia Braginsky, Guy Golan-Gueta, Eshcar Hillel, Idit Keidar, Moshe Sulamy
PPoPP6
2017 On Liveness of Dynamic Storage
Alexander Spiegelman, Idit Keidar
SIROCCO2
2017 WatchIT: Who Watches Your IT Guy?
abstract
System administrators have unlimited access to system resources. As the Snowden case highlighted, these permissions can be exploited to steal valuable personal, classified, or commercial data. This problem is exacerbated when a third party administers the system. For example, a bank outsourcing its IT would not want to allow administrators access to the actual data. We propose WatchIT: a strategy that constrains IT personnel's view of the system and monitors their actions. To this end, we introduce the abstraction of perforated containers -- while regular Linux containers are too restrictive to be used by system administrators, by "punching holes" in them, we strike a balance between information security and required administrative needs. Following the principle of least privilege, our system predicts which system resources should be accessible for handling each IT issue, creates a perforated container with the corresponding isolation, and deploys it as needed for fixing the problem.
Noam Shalev, Idit Keidar, Yaron Weinsberg, Yosef Moatti, Elad Ben-Yehuda
SOSP2
2017 Brief Announcement: Towards Reduced Instruction Sets for Synchronization
abstract
Contrary to common belief, a recent work by Ellen, Gelashvili, Shavit, and Zhu has shown that computability does not require multicore architectures to support "strong" synchronization instructions like compare-and-swap, as opposed to combinations of "weaker" instructions like decrement and multiply. However, this is the status quo, and in turn, most efficient concurrent data-structures heavily rely on compare-and-swap (e.g. for swinging pointers). We show that this need not be the case, by designing and implementing a concurrent linearizable Log data-structure (also known as a History object), supporting two operations: append(item), which appends the item to the log, and get-log(), which returns the appended items so far, in order. Readers are wait-free and writers are lock-free, hence this data-structure can be used in a lock-free universal construction to implement any concurrent object with a given sequential specification. Our implementation uses atomic read, xor, decrement, and fetch-and-increment instructions supported on X86 architectures, and provides similar performance to a compare-and-swap-based solution on today's hardware. This raises a fundamental question about minimal set of synchronization instructions that the architectures have to support.
Rati Gelashvili, Idit Keidar, Alexander Spiegelman, Roger Wattenhofer
DISC2
2017 Dynamic Reconfiguration: Abstraction and Optimal Asynchronous Solution
abstract
Providing clean and efficient foundations and tools for reconfiguration is a crucial enabler for distributed system management today. This work takes a step towards developing such foundations. It considers classic fault-tolerant atomic objects emulated on top of a static set of fault-prone servers, and turns them into dynamic ones. The specification of a dynamic object extends the corresponding static (non-dynamic) one with an API for changing the underlying set of fault-prone servers. Thus, in a dynamic model, an object can start in some configuration and continue in a different one. Its liveness is preserved through the reconfigurations it undergoes, tolerating a versatile set of faults as it shifts from one configuration to another. In this paper we present a general abstraction for asynchronous reconfiguration, and exemplify its usefulness for building two dynamic objects: a read/write register and a max-register. We first define a dynamic model with a clean failure condition that allows an administrator to reconfigure the system and switch off a server once the reconfiguration operation removing it completes. We then define the Reconfiguration abstraction and show how it can be used to build dynamic registers and max-registers. Finally, we give an optimal asynchronous algorithm implementing the Reconfiguration abstraction, which in turn leads to the first asynchronous (consensus-free) dynamic register emulation with optimal complexity. More concretely, faced with n requests for configuration changes, the number of configurations that the dynamic register is implemented over is n; and the complexity of each client operation is O(n).
Alexander Spiegelman, Idit Keidar, Dahlia Malkhi
DISC2
2017 Composing ordered sequential consistency
Kfir Lev-Ari, Edward Bortnikov, Idit Keidar, Alexander Shraer
Inf. Process. Lett.3
2016 CSR: Core Surprise Removal in Commodity Operating Systems
abstract
One of the adverse effects of shrinking transistor sizes is that processors have become increasingly prone to hardware faults. At the same time, the number of cores per die rises. Consequently, core failures can no longer be ruled out, and future operating systems for many-core machines will have to incorporate fault tolerance mechanisms.
Noam Shalev, Eran Harpaz, Hagar Meir, Idit Keidar, Yaron Weinsberg
ASPLOS4
2016 Dynamic Atomic Snapshots
abstract
Snapshots are useful tools for monitoring big distributed and parallel systems. In this paper, we adapt the well-known atomic snapshot abstraction to dynamic models with an unbounded number of participating processes. Our dynamic snapshot specification extends the API to allow changing the set of processes whose values should be returned from a scan operation. We introduce the ephemeral memory model, which consists of a dynamically changing set of nodes; when a node is removed, its memory can be immediately reclaimed. In this model, we present an algorithm for wait-free dynamic atomic snapshots.
Alexander Spiegelman, Idit Keidar
OPODIS2
2016 Transactional data structure libraries
abstract
We introduce transactions into libraries of concurrent data structures; such transactions can be used to ensure atomicity of sequences of data structure operations. By focusing on transactional access to a well-defined set of data structure operations, we strike a balance between the ease-of-programming of transactions and the efficiency of custom-tailored data structures. We exemplify this concept by designing and implementing a library supporting transactions on any number of maps, sets (implemented as skiplists), and queues. Our library offers efficient and scalable transactions, which are an order of magnitude faster than state-of-the-art transactional memory toolkits. Moreover, our approach treats stand-alone data structure operations (like put and enqueue) as first class citizens, and allows them to execute with virtually no overhead, at the speed of the original data structure library.
Alexander Spiegelman, Guy Golan-Gueta, Idit Keidar
PLDI3
2016 Brief Announcement: A Key-Value Map for Massive Real-Time Analytics
abstract
Modern big data processing platforms employ huge in-memory key-value (KV-) maps. Their applications simultaneously drive high-rate data ingestion and large-scale analytics. These two scenarios expect KV-map implementations that scale well with both real-time updates and massive atomic scans triggered by range queries. However, today's state-of-the art concurrent KV-maps fall short of satisfying these requirements -- they either provide only limited or non-atomic scans, or severely hamper updates when scans are ongoing. We present KiWi, the first atomic KV-map to efficiently support simultaneous massive data retrieval and real-time access. The key to achieving this is treating scans as first class citizens, whereas most existing concurrent KV-maps do not provide atomic scans, and some others add them to existing maps without rethinking the design anew.
Dmitry Basin, Edward Bortnikov, Anastasia Braginsky, Guy Golan-Gueta, Eshcar Hillel, Idit Keidar, Moshe Sulamy
PODC6
2016 Space Bounds for Reliable Storage: Fundamental Limits of Coding
abstract
We study the inherent space requirements of reliable storage algorithms in asynchronous distributed systems. A number of recent works have used codes in order to achieve a better storage cost than the well-known replication approach. However, a closer look reveals that they incur extra costs in certain scenarios. Specifically, if multiple clients access the storage concurrently, then existing asynchronous code-based algorithms may store a number of copies of the data that grows linearly with the number of concurrent clients. We prove here that this is inherent. Given three parameters, (1) the data size -- D bits, (2) the concurrency level -- c, and (3) the number of storage node failures that need to be tolerated -- f, we show a lower bound of Omega(min(f,c)D) bits on the space complexity of asynchronous distributed storage algorithms. Intuitively, this implies that the asymptotic storage cost is either as high as with replication, namely O(fD), or as high under concurrency as with the aforementioned code-based algorithms, i.e., O(cD).
Alexander Spiegelman, Yuval Cassuto, Gregory V. Chockler, Idit Keidar
PODC4
2016 NearBucket-LSH: Efficient Similarity Search in P2P Networks
Naama Kraus, David Carmel, Idit Keidar, Meni Orenbach
SISAP3
2016 Brief Announcement: Transactional Data Structure Libraries
abstract
We introduce transactions into libraries of concurrent data structures; such transactions can be used to ensure atomicity of sequences of data structure operations. By restricting transactional access to a well-defined set of data structure operations, we strike a balance between the ease-of-programming of transactions and the efficiency of custom-tailored data structures. We exemplify this concept by designing and implementing a library supporting transactions on any number of maps, sets (implemented as skiplists), and queues. Our library offers efficient and scalable transactions, which are an order of magnitude faster than state-of-the-art transactional memory toolkits. Moreover, our approach treats stand-alone data structure operations (like put and enqueue) as first class citizens, and allows them to execute with virtually no overhead, at the speed of the original data structure library.
Alexander Spiegelman, Guy Golan-Gueta, Idit Keidar
SPAA3
2016 Modular Composition of Coordination Services
Kfir Lev-Ari, Edward Bortnikov, Idit Keidar, Alexander Shraer
USENIX ATC3
2016 EFS: Energy-Friendly Scheduler for memory bandwidth constrained systems
Tomer Y. Morad, Noam Shalev, Idit Keidar, Avinoam Kolodny, Uri C. Weiser
J. Parallel Distributed Comput.3
2015 Scaling concurrent log-structured data stores
abstract
Log-structured data stores (LSM-DSs) are widely accepted as the state-of-the-art implementation of key-value stores. They replace random disk writes with sequential I/O, by accumulating large batches of updates in an in-memory data structure and merging it with the on-disk store in the background. While LSM-DS implementations proved to be highly successful at masking the I/O bottleneck, scaling them up on multicore CPUs remains a challenge. This is nontrivial due to their often rich APIs, as well as the need to coordinate the RAM access with the background I/O.
Guy Golan-Gueta, Edward Bortnikov, Eshcar Hillel, Idit Keidar
EuroSys4
2015 Space Bounds for Reliable Storage: Fundamental Limits of Coding (Keynote)
abstract
We present here a synopsis of a keynote presentation given by Idit Keidar at OPODIS 2015, the International Conference on Principles of Distributed Systems, which took place in Rennes, France, on December 14-17 2015.
Alexander Spiegelman, Yuval Cassuto, Gregory V. Chockler, Idit Keidar
OPODIS4
2015 Dynamic Reconfiguration: A Tutorial (Tutorial)
abstract
A key challenge for distributed systems is the problem of reconfiguration. Clearly, any production storage system that provides data reliability and availability for long periods must be able to reconfigure in order to remove failed or old servers and add healthy or new ones. This is far from trivial since we do not want the reconfiguration management to be centralized or cause a system shutdown. In this tutorial we look into existing reconfigurable storage algorithms. We propose a common model and failure condition capturing their guarantees. We define a reconfiguration problem around which dynamic object solutions may be designed. To demonstrate its strength, we use it to implement dynamic atomic storage. We present a generic framework for solving the reconfiguration problem, show how to recast existing algorithms in terms of this framework, and compare among them.
Alexander Spiegelman, Idit Keidar, Dahlia Malkhi
OPODIS2
2015 Towards Automatic Lock Removal for Scalable Synchronization
Maya Arbel-Raviv, Guy Golan-Gueta, Eshcar Hillel, Idit Keidar
DISC4
2015 A Constructive Approach for Proving Data Structures' Linearizability
Kfir Lev-Ari, Gregory V. Chockler, Idit Keidar
DISC3
2015 On Avoiding Spare Aborts in Transactional Memory
Idit Keidar, Dmitri Perelman
Theory Comput. Syst.1
2014 On Correctness of Data Structures under Reads-Write Concurrency
Kfir Lev-Ari, Gregory V. Chockler, Idit Keidar
DISC3
2014 LiMoSense: live monitoring in dynamic sensor networks
Ittay Eyal, Idit Keidar, Raphael Rom
Distributed Comput.2
2014 GPUfs: Integrating a file system with GPUs
abstract
As GPU hardware becomes increasingly general-purpose, it is quickly outgrowing the traditional, constrained GPU-as-coprocessor programming model. This article advocates for extending standard operating system services and abstractions to GPUs in order to facilitate program development and enable harmonious integration of GPUs in computing systems. As an example, we describe the design and implementation of GPUFs, a software layer which provides operating system support for accessing host files directly from GPU programs. GPUFs provides a POSIX-like API, exploits GPU parallelism for efficiency, and optimizes GPU file access by extending the host CPU's buffer cache into GPU memory. Our experiments, based on a set of real benchmarks adapted to use our file system, demonstrate the feasibility and benefits of the GPUFs approach. For example, a self-contained GPU program that searches for a set of strings throughout the Linux kernel source tree runs over seven times faster than on an eight-core CPU.
Mark Silberstein, Bryan Ford, Idit Keidar, Emmett Witchel
ACM Trans. Comput. Syst.3
2013 GPUfs: integrating a file system with GPUs
abstract
PU hardware is becoming increasingly general purpose, quickly outgrowing the traditional but constrained GPU-as-coprocessor programming model. To make GPUs easier to program and easier to integrate with existing systems, we propose making the host's file system directly accessible from GPU code. GPUfs provides a POSIX-like API for GPU programs, exploits GPU parallelism for efficiency, and optimizes GPU file access by extending the buffer cache into GPU memory. Our experiments, based on a set of real benchmarks adopted to use our file system, demonstrate the feasibility and benefits of our approach. For example, we demonstrate a simple self-contained GPU program which searches for a set of strings in the entire tree of Linux kernel source files over seven times faster than an eight-core CPU run.
Mark Silberstein, Bryan Ford, Idit Keidar, Emmett Witchel
ASPLOS3
2013 Distributed sparse signal recovery for sensor networks
abstract
We propose a distributed algorithm for sparse signal recovery in sensor networks based on Iterative Hard Thresholding (IHT). Every agent has a set of measurements of a signal x, and the objective is for the agents to recover x from their collective measurements at a minimal communication cost and with low computational complexity. A naïve distributed implementation of IHT would require global communication of every agent's full state in each iteration. We find that we can dramatically reduce this communication cost by leveraging solutions to the distributed top-K problem in the database literature. Evaluations show that our algorithm requires up to three orders of magnitude less total bandwidth than the best-known distributed basis pursuit method.
Stacy Patterson, Yonina C. Eldar, Idit Keidar
ICASSP3
2013 In-Network Analytics for Ubiquitous Sensing
Ittay Eyal, Idit Keidar, Stacy Patterson, Raphi Rom
DISC2
2012 SALSA: scalable and low synchronization NUMA-aware algorithm for producer-consumer pools
abstract
We present a highly-scalable non-blocking producer-consumer task pool, designed with a special emphasis on lightweight synchronization and data locality. The core building block of our pool is SALSA, Scalable And Low Synchronization Algorithm for a single-consumer container with task stealing support. Each consumer operates on its own SALSA container, stealing tasks from other containers if necessary. We implement an elegant self-tuning policy for task insertion, which does not push tasks to overloaded SALSA containers, thus decreasing the likelihood of stealing.
Elad Gidron, Idit Keidar, Dmitri Perelman, Yonathan Perez
SPAA2
2011 LiMoSense - Live Monitoring in Dynamic Sensor Networks
Ittay Eyal, Idit Keidar, Raphael Rom
ALGOSENSORS2
2011 Tolerant Value Speculation in Coarse-Grain Streaming Computations
abstract
Streaming applications are the subject of growing interest, as the need for fast access to data continues to grow. In this work, we present the design requirements and implementation of coarse-grain value speculation in streaming applications. We explain how this technique can be useful in cases where serial parts of applications constitute bottlenecks, and when slower I/O favors using available prefixes of the data. Contrary to previous work, we show how allowing some tolerance can justify early predictions on a scale of a large window of values. We suggest a methodology for runtime support of speculation, along with the mechanisms required for rollback. We present resource management issues consequent to our technique. We study how validation and speculation frequencies impact the performance of the program. Finally, we present our implementation in the context of the Huffman encoder benchmark, running it in different configurations and on different architectures.
Nathaniel Azuelos, Idit Keidar, Ayal Zaks
IPDPS2
2011 CAFÉ: Scalable Task Pools with Adjustable Fairness and Contention
Dmitry Basin, Rui Fan 0004, Idit Keidar, Ofer Kiselov, Dmitri Perelman
DISC3
2011 SMV: Selective Multi-Versioning STM
Dmitri Perelman, Anton Byshevsky, Oleg Litmanovich, Idit Keidar
DISC4
2011 Distributed data clustering in sensor networks
Ittay Eyal, Idit Keidar, Raphael Rom
Distributed Comput.2
2011 Special issue with selected papers from DISC 2009
Idit Keidar
Distributed Comput.1
2011 Dynamic atomic storage without consensus
abstract
This article deals with the emulation of atomic read/write (R/W) storage in dynamic asynchronous message passing systems. In static settings, it is well known that atomic R/W storage can be implemented in a fault-tolerant manner even if the system is completely asynchronous, whereas consensus is not solvable. In contrast, all existing emulations of atomic storage in dynamic systems rely on consensus or stronger primitives, leading to a popular belief that dynamic R/W storage is unattainable without consensus. In this article, we specify the problem of dynamic atomic read/write storage in terms of the interface available to the users of such storage. We discover that, perhaps surprisingly, dynamic R/W storage is solvable in a completely asynchronous system: we present DynaStore, an algorithm that solves this problem. Our result implies that atomic R/W storage is in fact easier than consensus, even in dynamic systems.
Marcos K. Aguilera, Idit Keidar, Dahlia Malkhi, Alexander Shraer
J. ACM2
2011 Fail-Aware Untrusted Storage
abstract
We consider a set of clients collaborating through an online service provider that is subject to attacks and hence not fully trusted by the clients. We introduce the abstraction of a fail-aware untrusted service, with meaningful semantics even when the provider is faulty. In the common case, when the provider is correct, such a service guarantees consistency (linearizability) and liveness (wait-freedom) of all operations. In addition, the service always provides accurate and complete consistency and failure detection. We illustrate our new abstraction by presenting a Fail-Aware Untrusted STorage service (FAUST). Existing storage protocols in this model guarantee so-called forking semantics. We observe, however, that none of the previously suggested protocols suffices for implementing fail-aware untrusted storage with the desired liveness and consistency properties (at least wait-freedom and linearizability when the server is correct). We present a new storage protocol, which does not suffer from this limitation, and implements a new consistency notion, called weak fork-linearizability. We show how to extend this protocol to provide eventual consistency and failure awareness in FAUST.
Christian Cachin, Idit Keidar, Alexander Shraer
SIAM J. Comput.2
2010 Threads vs. caches: Modeling the behavior of parallel workloads
abstract
A new generation of high-performance engines now combine graphics-oriented parallel processors with a cache architecture. In order to meet this new trend, new highly-parallel workloads are being developed. However, it is often difficult to predict how a given application would perform on a given architecture. This paper provides a new model capturing the behavior of such parallel workloads on different multi-core architectures. Specifically, we provide a simple analytical model, which, for a given application, describes its performance and power as a function of the number of threads it runs in parallel, on a range of architectures. We use our model (backed by simulations) to study both synthetic workloads and real ones from the PARSEC suite. Our findings recognize distinctly different behavior patterns for different application families and architectures.
Zvika Guz, Oved Itzhak, Idit Keidar, Avinoam Kolodny, Avi Mendelson, Uri C. Weiser
ICCD3
2010 Brief announcement: sources of instability in data center multicast
abstract
No abstract available.
Dmitry Basin, Kenneth P. Birman, Idit Keidar, Ymir Vigfusson
PODC3
2010 Distributed data classification in sensor networks
abstract
Low overhead analysis of large distributed data sets is necessary for current data centers and for future sensor networks. In such systems, each node holds some data value, e.g., a local sensor read, and a concise picture of the global system state needs to be obtained. In resource-constrained environments like sensor networks, this needs to be done without collecting all the data at any location, i.e., in a distributed, manner. To this end, we define the distributed classification problem, in which numerous interconnected nodes compute a classification of their data, i.e., partition these values into multiple collections, and describe each collection concisely.
Ittay Eyal, Idit Keidar, Raphael Rom
PODC2
2010 On maintaining multiple versions in STM
abstract
An effective way to reduce the number of aborts in software transactional memory (STM) is to keep multiple versions of transactional objects. In this paper, we study inherent properties of STMs that use multiple versions to guarantee successful commits of all read-only transactions.
Dmitri Perelman, Rui Fan 0004, Idit Keidar
PODC3
2010 Order is power: Selective Packet Interleaving for energy efficient Networks-on-Chip
abstract
Network-on-Chip (NoC) links consume a significant fraction of the total NoC power. We present Selective Packet Interleaving (SPI), a flit transmission scheme that reduces power consumption in NoC links. SPI decreases the number of bit transitions in the links by exploiting the multiplicity of virtual channels in a NoC router. SPI multiplexes flits to the router's output link so as to minimize the number of bit transitions from the previously transmitted flit. Analysis and simulations demonstrate a reduction of up to 55% in the number of bit transitions and up to 40% savings in power consumed on the link. SPI benefits grow with the number of virtual channels. SPI works better for links with a small number of bits in parallel. While SPI compares favorably against bus inversion, combining both schemes helps to further reduce bit transitions.
Amit Berman, Ran Ginosar, Idit Keidar
VLSI-SoC3
2010 Correctness of Gossip-Based Membership under Message Loss
abstract
Due to their simplicity and effectiveness, gossip-based membership protocols have become the method of choice for maintaining partial membership in large peer-to-peer systems. A variety of gossip-based membership protocols were proposed. Some were shown to be effective empirically, lacking analytic understanding of their properties. Others were analyzed under simplifying assumptions, such as lossless and delayless network. It is not clear whether the analysis results hold in dynamic networks, where both nodes and network links can fail. In this paper we try to bridge this gap. We first enumerate the desirable properties of a gossip-based membership protocol, such as view uniformity, independence, and load balance. We then propose a simple send & forget protocol, and show that even in the presence of message loss, it achieves the desirable properties.
Maxim Gurevich, Idit Keidar
SIAM J. Comput.2
2009 Fail-Aware Untrusted Storage
abstract
We consider a set of clients collaborating through an online service provider that is subject to attacks, and hence not fully trusted by the clients. We introduce the abstraction of a fail-aware untrusted service, with meaningful semantics even when the provider is faulty. In the common case, when the provider is correct, such a service guarantees consistency (linearizability) and liveness (wait-freedom) of all operations. In addition, the service always provides accurate and complete consistency and failure detection. We illustrate our new abstraction by presenting a Fail-Aware Untrusted STorage service (FAUST). Existing storage protocols in this model guarantee so-called forking semantics. We observe, however, that none of the previously suggested protocols suffice for implementing fail-aware untrusted storage with the desired liveness and consistency properties (at least wait-freedom and linearizability when the server is correct). We present a new storage protocol, which does not suffer from this limitation, and implements a new consistency notion, called weak fork-linearizability. We show how to extend this protocol to provide eventual consistency and failure awareness in FAUST.
Christian Cachin, Idit Keidar, Alexander Shraer
DSN2
2009 Low-overhead error detection for Networks-on-Chip
abstract
In the current deep sub-micron age, interconnect reliability is a subject of major concern, and is crucial for a successful product. Coding is a widely-used method to achieve communication reliability, which can be very useful in a network-on-chip (NoC). A key challenge for NoC error detection is to provide a defined detection level, while minimizing the number of redundant parity bits, using small encoder and decoder circuits, and ensuring shortest path routing. We present parity routing (PaR), a novel method to reduce the number of redundant bits transmitted. PaR exploits NoC path diversity to reduce the number of redundant parity bits. Our analysis shows that, for example, on a 4×4 NoC with a demand of one parity bit, PaR reduces the redundant information transmitted by 75%, and the savings increase asymptotically to 100% with the size of the NoC. In addition, we show that PaR can yield power savings due to the reduced number of bit transmissions and simple decoding process. Furthermore, PaR utilizes low complexity, small-area circuits.
Amit Berman, Idit Keidar
ICCD2
2009 Dynamic atomic storage without consensus
abstract
This paper deals with the emulation of atomic read/write (R/W) storage in dynamic asynchronous message passing systems. In static settings, it is well known that atomic R/W storage can be implemented in a fault-tolerant manner even if the system is completely asynchronous, whereas consensus is not solvable. In contrast, all existing emulations of atomic storage in dynamic systems rely on consensus or stronger primitives, leading to a popular belief that dynamic R/W storage is unattainable without consensus.
Marcos K. Aguilera, Idit Keidar, Dahlia Malkhi, Alexander Shraer
PODC2
2009 Correctness of gossip-based membership under message loss
abstract
Due to their simplicity and effectiveness, gossip-based membership protocols have become the method of choice for maintaining partial membership in large P2P systems. A variety of gossip-based membership protocols were proposed. Some were shown to be effective empirically, lacking analytic understanding of their properties. Others were analyzed under simplifying assumptions, such as lossless and delay-less network. It is not clear whether the analysis results hold in dynamic networks where both nodes and network links can fail.
Maxim Gurevich, Idit Keidar
PODC2
2009 On avoiding spare aborts in transactional memory
abstract
This paper takes a step toward developing a theory for understanding aborts in transactional memory systems (TMs). Existing TMs may abort many transactions that could, in fact, commit without violating correctness. We call such unnecessary aborts spare aborts. We classify what kinds of spare aborts can be eliminated, and which cannot. We further study what kinds of spare aborts can be avoided efficiently. Specifically, we show that some unnecessary aborts cannot be avoided, and that there is an inherent tradeoff between the overhead of a TM and the extent to which it reduces the number of spare aborts. We also present an efficient example TM algorithm that avoids certain kinds of spare aborts, and analyze its properties and performance.
Idit Keidar, Dmitri Perelman
SPAA1
2009 Transactifying Apache's cache module
abstract
Apache is a large-scale industrial multi-process and multithreaded application, which uses lock-based synchronization. We report on our experience in modifying Apache's cache module to employ transactional memory instead of locks, a process we refer to as transactification; we are not aware of any previous efforts to transactify legacy software of such a large scale. Along the way, we learned some valuable lessons about which tools one should use, which parts of the code one should transactify and which are better left untouched, as well as on the intricacy of commit handlers. We also stumbled across weaknesses of existing software transactional memory (STM) toolkits, leading us to identify desirable features they are currently lacking. Finally, we present performance results from running Apache on a 32-core machine, showing that, there are scenarios where the performance of the STM-based version is close to that of the lock-based version. These results suggest that there are applications for which the overhead of using a software-only implementation of transactional memory is insignificant.
Haggai Eran, Ohad Lutzky, Zvika Guz, Idit Keidar
SYSTOR4
2009 The 2009 Edsger W. Dijkstra Prize in Distributed Computing
Lorenzo Alvisi, Rachid Guerraoui, Prasad Jayanti, Idit Keidar, Shay Kutten, Jennifer L. Welch
DISC4
2009 Brahms: Byzantine resilient random membership sampling
Edward Bortnikov, Maxim Gurevich, Idit Keidar, Gabriel Kliot, Alexander Shraer
Comput. Networks3
2009 EquiCast: Scalable multicast with selfish users
Idit Keidar, Roie Melamed, Ariel Orda
Comput. Networks1
2009 Fork sequential consistency is blocking
Christian Cachin, Idit Keidar, Alexander Shraer
Inf. Process. Lett.2
2009 Deleting files in the Celeste peer-to-peer storage system
Gal Badishi, Germano Caronni, Idit Keidar, Raphael Rom, Glenn Scott
J. Parallel Distributed Comput.3
2009 Impossibility Results and Lower Bounds for Consensus under Link Failures
abstract
We provide a suite of impossibility results and lower bounds for the required number of processes and rounds for synchronous consensus under transient link failures. Our results show that consensus can be solved even in the presence of $O(n^2)$ moving omission and/or arbitrary link failures per round, provided that both the number of affected outgoing and incoming links of every process is bounded. Providing a step further toward the weakest conditions under which consensus is solvable, our findings are applicable to a variety of dynamic phenomena such as transient communication failures and end-to-end delay variations. We also prove that our model surpasses alternative link failure modeling approaches in terms of assumption coverage.
Ulrich Schmid 0001, Bettina Weiss, Idit Keidar
SIAM J. Comput.3
2009 Do not crawl in the DUST: Different URLs with similar text
abstract
We consider the problem of DUST: Different URLs with Similar Text. Such duplicate URLs are prevalent in Web sites, as Web server software often uses aliases and redirections, and dynamically generates the same page from various different URL requests. We present a novel algorithm, DustBuster , for uncovering DUST; that is, for discovering rules that transform a given URL to others that are likely to have similar content. DustBuster mines DUST effectively from previous crawl logs or Web server logs, without /examining page contents. Verifying these rules via sampling requires fetching few actual Web pages. Search engines can benefit from information about DUST to increase the effectiveness of crawling, reduce indexing overhead, and improve the quality of popularity statistics such as PageRank.
Ziv Bar-Yossef, Idit Keidar, Uri Schonfeld
ACM Trans. Web2
2008 Dynamic service assignment in mobile networks: the magma approach
Edward Bortnikov, Israel Cidon, Idit Keidar
PODC3
2008 Brahms: byzantine resilient random membership sampling
abstract
We present Brahms, an algorithm for sampling random nodes in a large dynamic system prone to malicious behavior. Brahms stores small membership views at each node, and yet overcomes Byzantine attacks by a linear portion of the system. Brahms is composed of two components. The first one is a resilient gossip-based membership protocol. The second one uses a novel memory-efficient approach for uniform sampling from a possibly biased stream of ids that traverse the node. We evaluate Brahms using rigorous analysis, backed by simulations, which show that our theoretical model captures the protocol's essentials. We study two representative attacks, and show that with high probability, an attacker cannot create a partition between correct nodes. We further prove that each node's sample converges to a uniform one over time. To our knowledge, no such properties were proven for gossip protocols in the past.
Edward Bortnikov, Maxim Gurevich, Idit Keidar, Gabriel Kliot, Alexander Shraer
PODC3
2008 Principles of untrusted storage: a new look at consistency conditions
abstract
No abstract available.
Christian Cachin, Idit Keidar, Alexander Shraer
PODC2
2008 Utilizing shared data in chip multiprocessors with the nahalal architecture
abstract
This paper addresses a new cache organization in a Chip Multiprocessors (CMP) environment. We introduce Nahalal, an architecture whose novel floorplan topology partitions cached data according to its usage (shared versus private data), and thus enables fast access to shared data for all processors while preserving the vicinity of private data to each processor. The Nahalal architecture combines the best of both shared caches and private caches, enabling fast accesses to data as in private caches while eliminating the need for inter-cache coherence transactions. Detailed simulations in Simics demonstrate that Nahalal decreases cache access latency by up to 41.1% compared to traditional CMP designs, yielding performance gains of up to 12.65% in run time.
Zvika Guz, Idit Keidar, Avinoam Kolodny, Uri C. Weiser
SPAA2
2008 An Empirical Study of Denial of Service Mitigation Techniques
abstract
We present an empirical study of the resistance of several protocols to denial of service (DoS) attacks on client-server communication. We show that protocols that use authentication alone, e.g., IPSec, provide protection to some extent, but are still susceptible to DoS attacks, even when the network is not congested. In contrast, a protocol that uses a changing filtering identifier (FI) is usually immune to DoS attacks, as long as the network itself is not congested. This approach is called FI hopping. We build and experiment with two prototype implementations of FI hopping. One implementation is a modification of IPSec in a Linux kernel, and a second implementation comes as an NDIS hook driver on a Windows machine. We present results of experiments in which client-server communication is subject to a DoS-attack. Our measurements illustrate that FI hopping withstands severe DoS attacks without hampering the client-server communication. Moreover, our implementations show that FI hopping is simple, practical, and easy to deploy.
Gal Badishi, Amir Herzberg, Idit Keidar, Oleg Romanov, Avital Yachin
SRDS3
2008 Araneola: A scalable reliable multicast system for dynamic environments
Roie Melamed, Idit Keidar
J. Parallel Distributed Comput.2
2008 How to Choose a Timing Model
abstract
When employing a consensus algorithm for state machine replication, should one optimize for the case that all communication links are usually timely or for fewer timely links? Does optimizing a protocol for better message complexity hamper the time complexity? In this paper, we investigate these types of questions using mathematical analysis as well as experiments over PlanetLab (WAN) and a LAN. We present a new and efficient leader-based consensus protocol that has O(n) stable-state message complexity (in a system with n processes) and requires only O(n) links to be timely at stable times. We compare this protocol with several previously suggested protocols. Our results show that a protocol that requires fewer timely links can achieve better performance, even if it sends fewer messages.
Idit Keidar, Alexander Shraer
IEEE Trans. Parallel Distributed Syst.1
2008 Octopus: A fault-tolerant and efficient ad-hoc routing protocol
Roie Melamed, Idit Keidar, Yoav Barel
Wirel. Networks2
2007 Scalable real-time gateway assignment in mobile mesh networks
abstract
The perception of future wireless mesh network (WMN) deployment and usage is rapidly evolving. WMNs are now being envisaged to provide citywide "last-mile" access for numerous mobile devices running media-rich applications with stringent quality of service (QoS) requirements. Consequently, some current-day conceptions underlying application support in WMNs need to be revisited. In particular, in a large WMN, the dynamic assignment of users to Internet gateways will become a complex traffic engineering problem that will need to consider load peaks, user mobility, and handoff penalties. We propose QMesh, a framework for user-gateway assignment that runs inside the WMN, and is oblivious to underlying routing protocols. It solves the handoff management problem in a scalable distributed manner. We evaluate QMesh through an extensive simulation (mostly of VoIP), in two settings: (1) a real campus network, with user mobility traces from the public CRAWDAD dataset, and (2) a large-scale urban WMN. Simulation results demonstrate that QMesh achieves significant QoS improvements and network capacity increases compared to traditional handoff policies, and illustrate the need for intelligent gateway assignment within the mesh.
Edward Bortnikov, Israel Cidon, Idit Keidar
CoNEXT3
2007 How to Choose a Timing Model?
abstract
When employing a consensus algorithm for state machine replication, should one optimize for the case that all communication links are usually timely, or for fewer timely links? Does optimizing a protocol for better message complexity hamper the time complexity? In this paper, we investigate these types of questions using mathematical analysis as well as experiments over PlanetLab (WAN) and a LAN. We present a new and efficient leader-based consensus protocol that has O(n) stable-state message complexity (in a system with n processes) and requires only O(n) links to be timely at stable times. We compare this protocol with several previously suggested protocols. Our results show that a protocol that requires fewer timely links can achieve better performance, even if it sends fewer messages.
Idit Keidar, Alexander Shraer
DSN1
2007 NoC-Based FPGA: Architecture and Routing
abstract
We present a novel network-on-chip-based architecture for future programmable chips (FPGAs). A key challenge for FPGA design is supporting numerous highly variable design instances with good performance and low cost. Our architecture minimizes the cost of supporting a wide range of design instances with given throughput requirements by balancing the amount of efficient hard-coded NoC infrastructure and the allocation of "soft" networking resources at configuration time. Although traffic patterns are design-specific, the physical link infrastructure is a performance bottleneck, and hence should be hard-coded. It is therefore important to employ routing schemes that allow for high flexibility to efficiently accommodate different traffic patterns during configuration. We examine the required capacity allocation for supporting a collection of typical traffic patterns on such chips under a number of routing schemes. We propose a new routing scheme, weighted ordered toggle (WOT), and show that it allows high design flexibility with low infrastructure cost. Moreover, WOT utilizes simple, small-area, on-chip routers, and has low memory demands
Roman Gindin, Israel Cidon, Idit Keidar
NOCS3
2007 Scalable Load-Distance Balancing
Edward Bortnikov, Israel Cidon, Idit Keidar
DISC3
2007 Amnesic Distributed Storage
Gregory V. Chockler, Rachid Guerraoui, Idit Keidar
DISC3
2007 Do not crawl in the dust: different urls with similar text
abstract
We consider the problem of DUST: Different URLs with Similar Text. Such duplicate URLs are prevalent in web sites, as web server software often uses aliases and redirections, and dynamically generates the same page from various different URLrequests. We present a novel algorithm, DustBuster, for uncovering DUST; that is, for discovering rules that transform a given URL to others that are likely to have similar content. DustBuster mines DUST effectively from previous crawl logs or web server logs, without examining page contents. Verifying these rules via sampling requires fetching few actual web pages. Search engines can benefit from information about DUST to increase the effectiveness of crawling, reduce indexing overhead, and improve the quality of popularity statistics such as PageRank.
Ziv Bar-Yossef, Idit Keidar, Uri Schonfeld
WWW2
2007 The overhead of consensus failure recovery
Partha Dutta, Rachid Guerraoui, Idit Keidar
Distributed Comput.3
2007 Wait-free regular storage from Byzantine components
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi
Inf. Process. Lett.3
2007 Keeping Denial-of-Service Attackers in the Dark
abstract
We consider the problem of overcoming (distributed) denial-of-service (DoS) attacks by realistic adversaries that have knowledge of their attack's successfulness, for example, by observing service performance degradation or by eavesdropping on messages or parts thereof. A solution for this problem in a high-speed network environment necessitates lightweight mechanisms for differentiating between valid traffic and the attacker's packets. The main challenge in presenting such a solution is to exploit existing packet-filtering mechanisms in a way that allows fast processing of packets but is complex enough so that the attacker cannot efficiently craft packets that pass the filters. We show a protocol that mitigates DoS attacks by adversaries that can eavesdrop and (with some delay) adapt their attacks accordingly. The protocol uses only available efficient packet-filtering mechanisms based mainly on addresses and port numbers. Our protocol avoids the use of fixed ports and instead performs "pseudorandom port hopping." We model the underlying packet-filtering services and define measures for the capabilities of the adversary and for the success rate of the protocol. Using these, we provide a novel rigorous analysis of the impact of DoS on an end-to-end protocol and show that our protocol provides effective DoS prevention for realistic attack and deployment scenarios.
Gal Badishi, Amir Herzberg, Idit Keidar
IEEE Trans. Dependable Secur. Comput.3
2007 Nomadic Service Assignment
abstract
We consider the problem of dynamically assigning application sessions of mobile users or user groups to service points. Such assignments must balance the trade-off between two conflicting goals. On the one hand, we would like to connect a user to the closest server in order to reduce network costs and service latencies. On the other hand, we would like to minimize the number of costly session migrations, or handoffs, between service points. We tackle this problem using two approaches. First, we employ algorithmic online optimization to obtain algorithms whose worst-case performance is within a factor of the optimal. Next, we extend them with opportunistic heuristics that achieve near-optimal practical average performance and scalability. We conduct case studies of two settings where such algorithms are required: wireless mesh networks with mobile users and wide-area groupware applications with or without mobility.
Edward Bortnikov, Israel Cidon, Idit Keidar
IEEE Trans. Mob. Comput.3
2006 Nomadic Service Points
abstract
Abstract — We consider the novel problem of dynamically assigning application sessions of mobile users or user groups to service points. Such assignments must balance the tradeoff between two conflicting goals. On the one hand, we would like to connect a user to the closest server, in order to reduce network costs and service latencies. On the other hand, we would like to minimize the number of costly session migrations, or handoffs, between service points. We tackle this problem using two approaches. First, we employ algorithmic online optimization to obtain algorithms whose worst-case performance is within a factor of the optimal. Next, we extend them with opportunistic versions that achieve excellent practical average performance and scalability. We conduct case studies of two settings where such algorithms are required: wireless mesh networks with mobile users, and wide-area groupware applications with or without mobility. I.
Edward Bortnikov, Israel Cidon, Idit Keidar
INFOCOM3
2006 Veracity radius: capturing the locality of distributed computations
abstract
This paper focuses on local computations of distributed aggregation problems on fixed graphs. We define a new metric on problem instances, Veracity Radius (VR), which captures the inherent possibility to compute them locally. We prove that VR yields a tight lower bound on output-stabilization time, i.e., the time until all nodes fix their outputs, as well as a lower bound on quiescence time. We present an efficient aggregation algorithm, I-LEAG, which reaches both output stabilization and quiescence within a time that is proportional to the VR of the problem instance, and is also efficient in terms of per-node communication and memory. We empirically show that the VR metric also effectively captures the performance of previously suggested efficient aggregation protocols, and that I-LEAG significantly outperforms these protocols in several respects.
Yitzhak Birk, Idit Keidar, Liran Liss, Assaf Schuster, Ran Wolff 0001
PODC2
2006 EquiCast: scalable multicast with selfish users
abstract
Peer-to-peer (P2P) networks suffer from the problem of "free-loaders", i.e., users who consume resources without contributing anything in return. In this paper, we tackle this problem taking a game theoretic perspective by modeling the system as a non-cooperative game. We introduce Equi-Cast, a wide-area P2P multicast protocol for large groups of selfish nodes. EquiCast is the first P2P multicast protocol that is formally proven to enforce cooperation in selfish environments. Additionally, we prove that EquiCast incurs a low constant load on each user.
Idit Keidar, Roie Melamed, Ariel Orda
PODC1
2006 Timeliness, failure-detectors, and consensus performance
abstract
We study the implication that various timeliness and failure detector assumptions have on the performance of consensus algorithms that exploit them. We present a general framework, GIRAF, for expressing such assumptions, and reasoning about the performance of indulgent algorithms.
Idit Keidar, Alexander Shraer
PODC1
2006 Deleting Files in the Celeste Peer-to-Peer Storage System
abstract
Celeste is a robust peer-to-peer object store built on top of a distributed hash table (DHT). Celeste is a working system, developed by Sun Microsystems Laboratories. During the development of Celeste, we faced the challenge of complete object deletion, and moreover, of deleting "files" composed of several different objects. This important problem is not solved by merely deleting meta-data, as there are scenarios in which all file contents must be deleted, e.g., due to a court order. Complete file deletion in a realistic peer-to-peer storage system has not been previously dealt with due to the intricacy of the problem - the system may experience high churn rates, nodes may crash or have intermittent connectivity, and the overlay network may become partitioned at times. We present an algorithm that eventually deletes all file content, data and meta-data, in the aforementioned complex scenarios. The algorithm is fully functional and has been successfully integrated into Celeste
Gal Badishi, Germano Caronni, Idit Keidar, Raphael Rom, Glenn Scott
SRDS3
2006 Efficient Dynamic Aggregation
Yitzhak Birk, Idit Keidar, Liran Liss, Assaf Schuster
DISC2
2006 Do not crawl in the DUST: different URLs with similar text
abstract
We consider the problem of dust: Different URLs with Similar Text. Such duplicate URLs are prevalent in web sites, as web server software often uses aliases and redirections, translates URLs to some canonical form, and dynamically generates the same page from various different URL requests. We present a novel algorithm, DustBuster, for uncovering dust; that is, for discovering rules for transforming a given URL to others that are likely to have similar content. DustBuster is able to detect dust effectively from previous crawl logs or web server logs, without examining page contents. Verifying these rules via sampling requires fetching few actual web pages. Search engines can benefit from this information to increase the effectiveness of crawling, reduce indexing overhead as well as improve the quality of popularity statistics such as PageRank.
Uri Schonfeld, Ziv Bar-Yossef, Idit Keidar
WWW3
2006 Byzantine disk paxos: optimal resilience with byzantine shared memory
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi
Distributed Comput.3
2006 An architecture for adaptive intrusion-tolerant applications
abstract
Abstract Applications that are part of a mission‐critical information system need to maintain a usable level of key services through ongoing cyber‐attacks. In addition to the well‐publicized denial of service (DoS) attacks, these networked and distributed applications are increasingly threatened by sophisticated attacks that attempt to corrupt system components and violate service integrity. While various approaches have been explored to deal with DoS attacks, corruption‐inducing attacks remain largely unaddressed. We have developed a collection of mechanisms based on redundancy, Byzantine fault tolerance, and adaptive middleware that help distributed, object‐based applications tolerate corruption‐inducing attacks. In this paper, we present the ITUA architecture, which integrates these mechanisms in a framework for auto‐adaptive intrusion‐tolerant systems, and we describe our experience in using the technology to defend a critical application that is part of a larger avionics system as an example. We also motivate the adaptive responses that are key to intrusion tolerance, and explain the use of the ITUA architecture to support them in an architectural framework. Copyright © 2006 John Wiley & Sons, Ltd.
Partha P. Pal, Paul Rubel, Michael Atighetchi, Franklin Webber, William H. Sanders, Mouna Seri, HariGovind V. Ramasamy, James Lyons, Tod Courtney, Adnan Agbaria, Michel Cukier, Jeanna M. Gossett, Idit Keidar
Softw. Pract. Exp.13
2006 Exposing and Eliminating Vulnerabilities to Denial of Service Attacks in Secure Gossip-Based Multicast
abstract
We propose a framework and methodology for quantifying the effect of denial of service (DoS) attacks on a distributed system. We present a systematic study of the resistance of gossip-based multicast protocols to DoS attacks. We show that even distributed and randomized gossip-based protocols, which eliminate single points of failure, do not necessarily eliminate vulnerabilities to DoS attacks. We propose Drum - a simple gossip-based multicast protocol that eliminates such vulnerabilities. Drum was implemented in Java and tested on a large cluster. We show, using closed-form mathematical analysis, simulations, and empirical tests, that Drum survives severe DoS attacks.
Gal Badishi, Idit Keidar, Amir Sasson
IEEE Trans. Dependable Secur. Comput.2
2005 Topic 8 - Distributed Systems and Algorithms
Marc Shapiro 0001, Idit Keidar, Felix C. Freiling, Luís E. T. Rodrigues
Euro-Par2
2005 Octopus: A Fault-Tolerant and Ef.cient Ad-hoc Routing Protocol
abstract
Mobile ad-hoc networks (MANETs) are failure-prone environments; it is common for mobile wireless nodes to intermittently disconnect from the network, e.g., due to signal blockage. This paper focuses on withstanding such failures in large MANETs: we present Octopus, a fault-tolerant and efficient position-based routing protocol. Fault-tolerance is achieved by employing redundancy, i.e., storing the location of each node at many other nodes, and by keeping frequently refreshed soft state. At the same time, Octopus achieves a low location update overhead by employing a novel aggregation technique, whereby a single packet updates the location of many nodes at many other nodes. Octopus is highly scalable: for a fixed node density, the number of location update packets sent does not grow with the network size. And when the density increases, the overhead drops. Thorough empirical evaluation using the ns2 simulator with up to 675 mobile nodes shows that Octopus achieves excellent fault-tolerance at a modest overhead: when all nodes intermittently disconnect and reconnect, Octopus achieves the same high reliability as when all nodes are constantly up.
Roie Melamed, Idit Keidar, Yoav Barel
SRDS2
2005 Keeping Denial-of-Service Attackers in the Dark
Gal Badishi, Amir Herzberg, Idit Keidar
DISC3
2005 MaGMA: mobility and group management architecture for real-time collaborative applications
abstract
We introduce MaGMA, a mobility and group management architecture, enabling real-time collaborative group applications such as push-to-talk (PTT) for mobile users. MaGMA provides, for the first time, a comprehensive and scalable solution for group management, seamless mobility, and quality-of-service (QoS). MaGMA is a distributed IP-based architecture consisting of an overlay server network deployed as part of the service infrastructure. MaGMA's architecture consists of a collection of mobile group managers (MGMs), which manage group membership and may also implement a multicast overlay for data delivery. The architecture is very flexible, and can co-exist with current as well as emerging wireless network technologies. We see such services as essential components in beyond-3G (B3G) networks. We propose two group management approaches in the context of MaGMA. We devise protocols for both approaches, evaluate both solutions using simulations, and validate the results through mathematical analysis. Finally, we present a proof-of-concept prototype implementation. Copyright © 2005 John Wiley & Sons, Ltd.
Nadav Lavi, Israel Cidon, Idit Keidar
Wirel. Commun. Mob. Comput.3
2004 Exposing and Eliminating Vulnerabilities to Denial of Service Attacks in Secure Gossip-Based Multicast
abstract
We propose a framework and methodology for quantifying the effect of denial of service (DoS) attacks on a distributed system. We present a systematic study of the resistance of gossip-based multicast protocols to DoS attacks. We show that even distributed and randomized gossip-based protocols, which eliminate single points of failure, do not necessarily eliminate vulnerabilities to DoS attacks. We propose Drum - a simple gossip-based multicast protocol that eliminates such vulnerabilities. Drum was implemented in Java and tested on a large cluster. We show, using closed-form mathematical analysis, simulations, and empirical tests, that Drum survives severe DoS attacks.
Gal Badishi, Idit Keidar, Amir Sasson
DSN2
2004 Caching-Enhanced Scalable Reliable Multicast
abstract
We present the caching-enhanced scalable reliable multicast (CESRM) protocol. CESRM augments the scalable reliable multicast (SRM) protocol (S. Floyd et al., 1995 and 1997) with a caching-based expedited recovery scheme. CESRM exploits the packet loss locality occurring in IP multicast transmissions in order to expeditiously recover from losses in the manner in which recent losses were recovered. Trace-driven simulations show that CESRM reduces the average recovery latency of SRM by roughly 50% and, moreover, drastically reduces the overhead in terms of recovery traffic and control messages.
Carolos Livadas, Idit Keidar
DSN2
2004 Araneola: A Scalable Reliable Multicast System for Dynamic Environments
abstract
We present Araneola, a scalable reliable application-level multicast system for highly dynamic wide-area environments. Araneola supports multi-point to multi-point reliable communication in a fully distributed manner while incurring constant load on each node. For a tunable parameter k /spl ges/ 3, Araneola constructs and dynamically maintains an overlay structure in which each node's degree is either k or k + 1, and roughly 90% of the nodes have degree k. Empirical evaluation shows that Araneola's overlay structure achieves three important mathematical properties of k-regular random graphs (i.e., random graphs in which each node has exactly k neighbors) with N nodes: (i) its diameter grows logarithmically with N; (ii) it is generally k-connected; and (iii) it remains highly connected following random removal of linear-size subsets of edges or nodes. The overlay is constructed at a very low cost: each join, leave, or failure is handled locally, and entails the sending of only about 3k messages in total. Given this overlay, Araneola disseminates multicast messages by gossiping over the overlay's links. We show that compared to a standard gossip-based multicast protocol, Araneola achieves substantial improvements in load, reliability, and latency. Finally, we present an extension to Araneola in which the basic overlay is enhanced with additional links chosen according to geographic proximity and available bandwidth. We show that this approach reduces the number of physical hops messages traverse without hurting the overlay's robustness.
Roie Melamed, Idit Keidar
NCA2
2004 Byzantine disk paxos: optimal resilience with byzantine shared memory
abstract
We present Byzantine Disk Paxos, an asynchronous shared-memory consensus protocol that uses a collection of n > 3t disks, t of which may fail by becoming non-responsive or arbitrarily corrupted. We give two constructions of this protocol; that is, we construct two different building blocks, each of which can be used, along with a leader oracle, to solve consensus. One building block is a shared wait-free safe register. The second building block is a regular register that satisfies a weaker termination (liveness) condition than wait freedom: its write operations are wait-free, whereas its read operations are guaranteed to return only in executions with a finite number of writes. We call this termination condition finite writes (FW), and show that consensus is solvable with FW-terminating registers and a leader oracle. We construct each of these reliable registers from n > 3t base registers, t of which can be non-responsive or Byzantine. All the previous wait-free constructions in this model used at least 4t+1 fault-prone registers, and we are not familiar with any prior FW-terminating constructions in this model.
Ittai Abraham, Gregory V. Chockler, Idit Keidar, Dahlia Malkhi
PODC3
2004 Brief announcement: exposing and eliminating vulnerabilities to denial of service attacks in secure gossip-based multicast
abstract
No abstract available.
Gal Badishi, Idit Keidar, Amir Sasson
PODC2
2004 Brief announcement: Trilix: a scalable unstructured lookup system for dynamic environments
Idit Keidar, Roie Melamed
PODC1
2003 A simple proof of the uniform consensus synchronous lower bound
Idit Keidar, Sergio Rajsbaum
Inf. Process. Lett.1
2002 Evaluating the running time of a communication round over the internet
abstract
We study the running time of distributed algorithms deployed in a widely distributed setting over the Internet using TCP. We consider a simple primitive that corresponds to a communication round in which every host sends information to every other host; this primitive occurs in numerous distributed algorithms. We experiment with four algorithms that typically implement this primitive. We run our experiments on ten hosts at geographically disperse locations over the Internet. We observe that message loss has a large impact on algorithm running times, which causes leader-based algorithms to usually outperform decentralized ones.
Omar Bakr, Idit Keidar
PODC2
2002 Early-Delivery Dynamic Atomic Broadcast
Ziv Bar-Joseph, Idit Keidar, Nancy A. Lynch
DISC2
2002 A Virtually Synchronous Group Multicast Algorithm for WANs: Formal Approach
abstract
This paper presents a formal design for a novel group communication service targeted for wide-area networks (WANs). The service provides virtual synchrony semantics. Such semantics facilitate the design of fault tolerant distributed applications. The presented design is more suitable for WANs than previously suggested ones. In particular, it features the first algorithm to achieve virtual synchrony semantics in a single communication round. The design also employs a scalable WAN-oriented architecture: it effectively decouples the main two components of virtually synchronous group communication---group membership and reliable group multicast. The design is carried out formally and rigorously. This paper includes formal specifications of both safety and liveness properties. The algorithm is formally modeled and assertionally verified.
Idit Keidar, Roger I. Khazan
SIAM J. Comput.1
2002 Moshe: A group membership service for WANs
abstract
We present Moshe, a novel scalable group membership algorithm built specifically for use in wide area networks (WANs), which can suffer partitions. Moshe is designed with three new significant features that are important in this setting: it avoids delivering views that reflect out-of-date memberships; it requires a single round of messages in the common case; and it employs a client-server design for scalability. Furthermore, Moshe's interface supplies the hooks needed to provide clients with full virtual synchrony semantics. We have implemented Moshe on top of a network event mechanism also designed specifically for use in a WAN. In addition to specifying the properties of the algorithm and proving that this specification is met, we provide empirical results of an implementation of Moshe running over the Internet. The empirical results justify the assumptions made by our design and exhibit good performance. In particular, Moshe terminates within a single communication round over 98% of the time. The experimental results also lead to interesting observations regarding the performance of membership algorithms over the Internet.
Idit Keidar, Jeremy B. Sussman, Keith Marzullo, Danny Dolev
ACM Trans. Comput. Syst.1
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.1
2001 Availability Study of Dynamic Voting Algorithms
abstract
Fault-tolerant distributed systems often select a primary component to allow a subset of the processes to function when failures occur. The dynamic voting paradigm defines rules for selecting the primary component adaptively: when a partition occurs, if a majority of the previous primary component is connected, a new and possibly smaller primary component is chosen. Several studies have shown that dynamic voting leads to more available solutions than other paradigms for maintaining a primary component. However, these studies have assumed that every attempt made by the algorithm to form a new primary component terminates successfully. Unfortunately, in real systems, this is not always the case: a change in connectivity can interrupt the algorithm while it is still attempting to form a new primary component; in such cases, algorithms may block until the processes can resolve the outcome of the interrupted attempt. This paper uses simulations to evaluate the effect of interruptions on the availability of dynamic voting algorithm. We study four dynamic voting algorithms and identify two important characteristics that impact an algorithm's availability in runs with frequent connectivity changes. First, we show that the number of processes that need to be present in order to resolve past attempts impacts the availability, especially during long runs with numerous connectivity changes. Second, we show that the number of communication rounds exchanged in an algorithm plays a significant role in the availability achieved, especially in the degradation of availability as connectivity changes become more frequent.
Kyle Ingols, Idit Keidar
ICDCS2
2000 A Client-Server Approach to Virtually Synchronous Group Multicast: Specifications and Algorithms
abstract
This paper presents a formal design for a novel group multicast service that provides virtually synchronous semantics in asynchronous fault-prone environments. The design employs a client-server architecture in which group membership is maintained not by every process but only by dedicated membership servers, while virtually synchronous group multicast is implemented by service end-points running at the clients. Specifically, the paper defines service semantics for the client-server interface, that is, for the group membership service. The paper then specifies virtually synchronous semantics for the new group multicast service, as a collection of commonly used safety and liveness properties. Finally, the paper presents new algorithms that use the defined group membership service to implement the specified properties. The algorithm that provides the complete virtually synchronous semantics executes in a single message round in parallel with the membership service's agreement on views, and is therefore more efficient than previously suggested algorithms providing such semantics.
Idit Keidar, Roger I. Khazan
ICDCS1
2000 A Client-Server Oriented Algorithm for Virtually Synchronous Group Membership in WANs
abstract
We describe a novel scalable group membership service designed explicitly for wide area networks. Our membership service is scalable in the number of groups supported, in the number of members in each group, and in the topology each group spans. Our service also supplies the hooks needed to provide clients with full virtual synchrony semantics. Our service attains, on average, a low message overhead by agreeing on membership within a single message round. Furthermore, our service avoids notifying the application of obsolete membership views when the network is unstable, yet it converges when the network has stabilized.
Idit Keidar, Jeremy B. Sussman, Keith Marzullo, Danny Dolev
ICDCS1
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
ICSE1
2000 Totally Ordered Multicast with Bounded Delays and Variable Rates
Ziv Bar-Joseph, Idit Keidar, Tal Anker, Nancy A. Lynch
OPODIS2
2000 Optimistic Virtual Synchrony
abstract
We present Optimistic Virtual Synchrony (OVS), a new form of group communication which provides the same capabilities as Virtual Synchrony with better performance. It does so by allowing applications to send messages during periods in which services implementing Virtual Synchrony block. OVS also allows applications to determine the policy as to when messages sent optimistically should be delivered and when they should be discarded. Thus, OVS gives applications fine grain control over the specific semantics they require, and does not impose costs for enforcing any semantics that they do not require. At the same time, OVS provides a single easy-to-use interface for all applications.
Jeremy B. Sussman, Idit Keidar, Keith Marzullo
SRDS2
1999 Fault Tolerant Video on Demand Services
abstract
This paper describes a highly available distributed video on demand (VoD) service which is inherently fault tolerant. The VoD service is provided by multiple servers that reside at different sites. New servers may be brought up "on the fly" to alleviate the load on other servers. When a server crashes it is replaced by another server in a transparent way; the clients are unaware of the change of service provider. In test runs of our VoD service prototype, such transitions are not noticeable to a human observer who uses the service. Our VoD service uses a sophisticated flow control mechanism and supports adjustment of the video quality to client capabilities. It does not assume any proprietary network technology: it uses commodity hardware and publicly available network technologies (e.g., TCP/IP, ATM). Our service may run on any machine connected to the Internet. The service exploits a group communication system as a building block for high availability. The utilization of group communication greatly simplifies the service design.
Tal Anker, Danny Dolev, Idit Keidar
ICDCS3
1998 Increasing the Resilience of Distributed and Replicated Database Systems
Idit Keidar, Danny Dolev
J. Comput. Syst. Sci.1
1997 Failure Detectors in Omission Failure Environments
abstract
No abstract available.
Danny Dolev, Roy Friedman 0001, Idit Keidar, Dahlia Malkhi
PODC3
1997 Dynamic Voting for Consistent Primary Components
abstract
Distributed applications often use quorums in order to guarantee consistency. With emerging world-wide communication technology, many new applications (e.g. conferencing applications and interactive games) wish to allow users to freely join and leave, without restarting the entire system. The dynamic voting paradigm allows such systems to define quorums adaptively, accounting for the changes in the set of participants. Furthermore, dynamic voting was proven to be the most available paradigm for maintaining quorums in unreliable networks. However, the subtleties of implementing dynamic voting were not well understood, in fact many of the suggested protocols may lead to inconsistencies in case of failures. Other protocols severely limit the availability in case failures occur during the protocol. In this paper we present a robust and efficient dynamic voting protocol for unreliable asynchronous networks. The protocol consistently maintains the primary component in a distributed system. O...
Esti Yeger Lotem, Idit Keidar, Danny Dolev
PODC2
1996 Efficient Message Ordering in Dynamic Networks
abstract
We present an algorithm for totally ordering messages in the face of network partitions and site failures.The algorithm aJways aJlows a majority of connected processors in the network to make progress (z.e. to order messages), if they remain connected for sufficiently long, regardless of past failures.Furthermore, our aJgorithm always allows processors to initiate messages, even when they are not members of a connected majority component in the network.Thus, messages can eventually become totally ordered even if their initiator is never a member of a majority component.The algorithm guarantees that when a majority is connected, each message is ordered within two communication rounds, if no failures occur during these rounds. 1 Introduction Consistent order is a powerful paradigm for the design of fault tolerant applications, e.g.consistent replication [Sch90, Kei94].We present an efficient algorithm for consistent message ordering in the face of network partitions and site failures, The network may partition into several components, and remerge.The algorithm is most adequate for dynamic networks where failures are transient.The algorithm uses an underlaying group communication service as a building block.Problem Definition Atomic broadcast deals with consistent message ordering.Informally, atomic broadcast requires that all the correct processors will deliver all the messages to the application in the same order and that they eventually deliver all messages sent by correct processors.In our
Idit Keidar, Danny Dolev
PODC1
1995 Increasing the Resilience of Atomic Commit at No Additional Cost
abstract
This paper presents a new atomic commitment protocol, Enhanced Three Phase Commit (E3PC ), that always allows a quorum in the system to make progress. Previously suggested quorum-based protocols (e.g. the quorum-based Three Phase Commit (3PC) [Ske82]) allow a quorum to make progress in case of one failure. If failures cascade, however, and the quorum in the system is "lost" (i.e. at a given time no quorum component exists, e.g. because of a total crash), a quorum can later become connected and still remain blocked. With our protocol, a connected quorum never blocks. E3PC is based on the quorumbased 3PC [Ske82], and it does not require more time or communication than 3PC. The principles demonstrated in this paper can be used to increase the resilience of a variety of distributed services, e.g. replicated database systems, by ensuring that a quorum will always be able to make progress. 1 Introduction Reliability and availability of loosely coupled distributed database systems is beco...
Idit Keidar, Danny Dolev
PODS1