Roy Friedman 0001

dblp:f/RoyFriedman · DBLP profile ↗
← Back
128ranked-venue papers
50as first author
19since 2021 · last 2025
0000-0001-6460-9665ORCID · verified

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

Systems, architecture and hardware · 36 · 15 first-author · 3 since 2021Computer networks · 29 · 10 first-author · 3 since 2021Security and privacy · 23 · 11 first-author · 6 since 2021Software engineering, systems software and programming languages · 13 · 4 first-author · 4 since 2021Theory of computation · 13 · 4 first-author · 1 since 2021Databases, data management, data science and information retrieval · 9 · 4 first-author · 2 since 2021Human-computer interaction and ubiquitous computing · 1
YearPublicationVenuePosition
2025 Mean Tail: Top-K and Frequency Estimation with Fewer Counters and More Keys
Dvir Biton, Roy Friedman 0001
AINA (2)2
2025 Geometric Sketch: The Inflatable-Shrinkable Sketch
Dvir Biton, Roy Friedman 0001, Rana Shahout
AINA (2)2
2025 Ethereum Conflicts Graphed
Dvir Biton, Roy Friedman 0001, Yaron Hay
ICBC2
2025 On Quorum Sizes in DAG-Based BFT Protocols
Razya Ladelsky, Roy Friedman 0001
ICBC2
2025 An Analysis of Sui's Transactions and Conflicts
abstract
Local execution of smart contracts related trans-actions is largely considered the main bottleneck in contem-porary blockchains' performance. To enable key performance optimizations, such as preloading frequently accessed objects and exploiting transaction-level parallelism, it is vital to understand the interaction patterns of such transactions and the structure of their object dependencies. Sui is a modern high-performance Layer 1 blockchain, which applies an object-centric data model. In this paper, we analyze over 1.5 million recent Sui checkpoints as well as 125,000 check-points related to 3 past high load events, by tracing transaction inputs and effects using the Sui full node API and Move call tracing tools. We report on the distribution of transactions per checkpoint, the object access patterns in Move function calls, the prevalence of shared vs. owned object usage, and a comprehensive characterization of the structure and sparsity of checkpoint-level conflict graphs. Our findings suggest that Sui's object model results in highly parallelizable conflict graphs exhibiting a complex hub and spokes structure, enabling efficient execution. We also identify distinct differences in errors, inputs and writes between the baseline usage and periods of heavy load. The data [1] and code [2] are available in open source.
Dvir Biton, Roy Friedman 0001
PRDC2
2025 Efficient Scheduling of Smart Contract Transactions via Conflict Graph Coloring
abstract
A smart contract is a special type of transaction designed for the execution of automated logic on blockchains. Alas, smart contracts transactions are one of the major hindrances to blockchain throughput. Hence, improving the execution time of smart contracts is a prime challenge for Blockchains at large. To that end, concurrent execution of smart contract is an appealing direction, which has been adopted by several contemporary Blockchains like Solana, Aptos, Sui, Sei, and Monad. Executing smart contracts in parallel requires applying deterministic concurrency controls based on ensuring consistent ordering of all conflicting transactions in all miners/validators. Existing implementations rely on the Block's total ordering to resolve this requirement. Recently, it has been suggested that relying on minimal coloring of the conflict graph corresponding to the Block's transactions can provide a better performance potential, yet without any evaluation. In this paper, we compare between approaches to smart contracts parallelization. Our study’ finds that in many situations, indeed the coloring-based ordering leads to significantly better performance than the Block order preserving approach. However, this gain has its limits, and it is not always guaranteed. In particular, the results are largely dependent on the conflict ratio in the conflict graph and the type of application.
Ankit Ravish, Yaron Hay, Manaswini Piduguralla, Roy Friedman 0001, Sathya Peri
PRDC5
2025 On the data persistency of replicated erasure codes in distributed storage systems
Roy Friedman 0001, Rafal Kapelko, Karol Marchwicki
Inf. Comput.1
2024 TrafficGrinder: A 0-RTT-Aware QUIC Load Balancer
abstract
QUIC is an emerging transport protocol, offering multiple advantages over TCP. We propose a novel 0-RTT-aware load balancing scheme for QUIC. The proposed scheme is scalable, resilient to 0 -RTT replay attacks, and guarantees perfect forward secrecy. It requires no modifications to the QUIC standard and it is QUIC version independent. 0-RTT is crucial for web performance, particularly on mobile networks. Our experiments show that it reduces the time-to-first-byte by half and the server load by 40% compared to 1-RTT, regardless of the network latency. Using both synthetic and real-world traffic traces, we show that the proposed load balancer guarantees near-optimal load balancing performance. It also guarantees faster time-to-first-byte and faster completion time compared to other load balancing schemes, such as least loaded, power-of-two-choices, and maximum session affinity.
Robert J. Shahla, Reuven Cohen, Roy Friedman 0001
ICNP3
2024 Distributed Recoverable Sketches
Diana Cohen, Roy Friedman 0001, Rana Shahout
OPODIS2
2024 Batch-Schedule-Execute: On Optimizing Concurrent Deterministic Scheduling for Blockchains
abstract
Executing smart contracts is a compute and storage-intensive task, which currently dominates modern blockchain's performance. Given that computers are becoming increasingly multicore, concurrency is an attractive approach to improve programs' execution runtime. A unique challenge of blockchains is that all replicas (miners or validators) must execute all smart contracts in the same logical order to maintain the semantics of State Machine Replication (SMR). In this work, we study the maximal level of parallelism attainable when focusing on the conflict graph between transactions packaged in the same block. This exposes a performance vulnerability that block creators may exploit against existing blockchain concurrency solutions, which rely on a total ordering phase for maintaining consistency amongst all replicas. To facilitate the formal aspects of our study, we develop a novel generic framework for Active State Machine Replication (ASMR) that is strictly serializable. We introduce the concept of graph scheduling and the definition of the minimal latency scheduling problem, which we prove to be NP-hard. We show that the restricted version of this problem for homogeneous transactions is equivalent to the classic Graph Vertex Coloring Problem, yet show that the heterogeneous case is more complex. We discuss the practical implications of these results.
Yaron Hay, Roy Friedman 0001
SRDS2
2023 Sketching the Path to Efficiency: Lightweight Learned Cache Replacement
Rana Shahout, Roy Friedman 0001
OPODIS2
2023 Together is Better: Heavy Hitters Quantile Estimation
abstract
Stream monitoring is fundamental in many data stream applications, such as financial data trackers, security, anomaly detection, and load balancing. In that respect, quantiles are of particular interest, as they often capture the user's utility. For example, if a video connection has high tail (e.g., 99'th percentile) latency, the perceived quality will suffer, even if the average and median latencies are low. In this work, we consider the problem of approximating the per-item quantiles. Elements in our stream are (ID, value) tuples, and we wish to track the quantiles for each ID. Existing quantile sketches are designed for a plain number stream (i.e., containing just a value). While one could allocate a separate sketch instance for each ID, this may require an infeasible amount of memory. Instead, we consider tracking the quantiles for the heavy hitters (most frequent items), which are often considered particularly important, without knowing them beforehand. We first present a couple of simple and effective algorithms that serve as baselines, a sampling approach and a sketching approach. Then, we present SQUAD, an algorithm that combines sampling and sketching while improving the asymptotic space complexity. Intuitively, SQUAD uses a background sampling process to capture the behaviour of the quantiles of an item before it is allocated with a sketch, thereby allowing us to use fewer samples and sketches. The algorithms are rigorously analyzed, and we demonstrate SQUAD's superiority using extensive~simulations on real-world traces.
Rana Shahout, Roy Friedman 0001, Ran Ben-Basat
Proc. ACM Manag. Data2
2022 SQUAD: combining sketching and sampling is better than either for per-item quantile estimation
abstract
Latency quantiles measurements are essential as they often capture the user's utility. For example, if a video connection has high tail latency, the perceived quality will suffer, even if the average and median latencies are low. In this work, we consider the problem of approximating the per-item quantiles. Elements in our stream are (ID, latency) tuples, and we wish to track the latency quantiles for each ID. Existing quantile sketches are designed for a single number stream (e.g., containing just the latency). While one could allocate a separate sketch instance for each ID, this may require an infeasible amount of memory. Instead, we consider tracking the quantiles for the heavy hitters (most frequent items), which are often considered particularly important, without knowing them beforehand.
Rana Shahout, Roy Friedman 0001, Ran Ben-Basat
SYSTOR2
2022 Box queries over multi-dimensional streams
Roy Friedman 0001, Rana Shahout
Inf. Syst.1
2022 Lightweight Robust Size Aware Cache Management
abstract
Modern key-value stores, object stores, Internet proxy caches, and Content Delivery Networks (CDN) often manage objects of diverse sizes, e.g., blobs, video files of different lengths, images with varying resolutions, and small documents. In such workloads, size-aware cache policies outperform size-oblivious algorithms. Unfortunately, existing size-aware algorithms tend to be overly complicated and computationally expensive. Our work follows a more approachable pattern; we extend the prevalent (size-oblivious) TinyLFU cache admission policy to handle variable-sized items. Implementing our approach inside two popular caching libraries only requires minor changes. We show that our algorithms yield competitive or better hit-ratios and byte hit-ratios compared to the state-of-the-art size-aware algorithms such as AdaptSize, LHD, LRB, and GDSF. Further, a runtime comparison indicates that our implementation is faster by up to 3× compared to the best alternative, i.e., it imposes a much lower CPU overhead.
Gil Einziger, Ohad Eytan, Roy Friedman 0001, Ben Manes
ACM Trans. Storage3
2021 CELL: Counter Estimation for Per-flow Traffic in Streams and Sliding Windows
abstract
Measurement capabilities are fundamental for a variety of network applications. Typically, recent data items are more relevant than old ones, a notion we can capture through a sliding window abstraction. These capabilities require a large number of counters in order to monitor the traffic of all network flows. However, SRAM memories are too small to contain these counters. Previous works suggested replacing counters with small estimators, trading accuracy for reduced space. But these estimators only focus on the counters’ size, whereas often flow ids consume more space than their respective counters. In this work, we present the CELL algorithm that combines estimators with efficient flow representation for superior memory reduction.We also extend CELL to the sliding window model, which prioritizes the recent data, by presenting two variants named RAND-CELL and SHIFT-CELL. We formally analyze the error and memory consumption of our algorithms and compare their performance against competing approaches using real-world Internet traces. These measurements exhibit the benefits of our work and show that CELL consumes at least 30% less space than the best-known alternative. The code is available in open source.
Rana Shahout, Roy Friedman 0001, Dolev Adas
ICNP2
2021 Sliding Window CRDT Sketches
abstract
Sketches maintain compact approximate statistics about streams of data, thereby enabling quickly answering queries regarding the data stream without having to reprocess it. Often, recent data is considered more important than older one, which is captured by the sliding window model. In distributed settings, where parts of the stream are seen by different, potentially geographically distributed components of the system, it makes sense to collect global statistics about the stream, but in a decentralized manner. Further, in order to ensure availability, scalability, and good performance, it is appealing to treat sketches as a CRDT data-type. In this work we introduce the notion of sliding window CRDT sketches. We then present the CRDT All Timestamps (aka CRDT-AT) and CRDT Last Timestamp (aka CRDT-LT) algorithms for implementing such sketches and analyze them. We also study the performance of CRDT-AT and CRDT-LT using real workloads, to establish their viability.
Dolev Adas, Roy Friedman 0001
SRDS2
2021 CELL: counter estimation for per-flow traffic over sliding windows
abstract
Estimators reduce the memory footprint of maintaining network statistics, while keeping the estimation error of each flow proportional to its size. This is unlike sketches and other approximate algorithms that only guarantee an error proportional to the entire stream size. In this work we present the CELL algorithm that combines estimators with efficient flow representation to obtain superior memory reduction compared to the state of the art.
Rana Shahout, Dolev Adas, Roy Friedman 0001
SYSTOR3
2021 Access Strategies for Network Caching
abstract
Having multiple data stores that can potentially serve content is common in modern networked applications. Data stores often publish approximate summaries of their content to enable effective utilization. Since these summaries are not entirely accurate, forming an efficient access strategy to multiple data stores becomes a complex risk management problem. This paper formally models this problem as a cost minimization problem, while taking into account both access costs, the inaccuracy of the approximate summaries, as well as the penalties incurred by failed requests. We introduce practical algorithms with guaranteed approximation ratios and further show that they are optimal in various settings. We also perform an extensive simulation study based on real data and show that our algorithms are more robust than existing heuristics. That is, they exhibit near-optimal performance in various settings, whereas the efficiency of existing approaches depends upon system parameters that may change over time, or be otherwise unknown.
Itamar Cohen, Gil Einziger, Roy Friedman 0001, Gabriel Scalosub
IEEE/ACM Trans. Netw.3
2020 It's Time to Revisit LRU vs. FIFO
Ohad Eytan, Danny Harnik, Effi Ofer, Roy Friedman 0001, Ronen I. Kat
HotStorage4
2020 Brief Announcement: Jiffy: A Fast, Memory Efficient, Wait-Free Multi-Producers Single-Consumer Queue
abstract
In applications such as sharded data processing systems, data flow programming and load sharing applications, multiple concurrent data producers are feeding requests into the same data consumer. This can be naturally realized through concurrent queues, where each consumer pulls its tasks from its dedicated queue. For scalability, wait-free queues are often preferred over lock based structures. The vast majority of wait-free queue implementations, and even lock-free ones, support the multi-producer multi-consumer model. Yet, this comes at a premium, since implementing wait-free multi-producer multi-consumer queues requires utilizing complex helper data structures. The latter increases the memory consumption of such queues and limits their performance and scalability. Additionally, many such designs employ (hardware) cache unfriendly memory access patterns. In this work we study the implementation of wait-free multi-producer single-consumer queues. Specifically, we propose Jiffy, an efficient memory frugal novel wait-free multi-producer single-consumer queue and formally prove its correctness. We then compare the performance and memory requirements of Jiffy with other state of the art lock-free and wait-free queues. We show that indeed Jiffy can maintain good performance with up to 128 threads, delivers better throughput than other constructions we compared against, and consumes less memory.
Dolev Adas, Roy Friedman 0001
DISC2
2020 FireLedger: A High Throughput Blockchain Consensus Protocol
abstract
Blockchains are distributed secure ledgers to which transactions are issued continuously and each block of transactions is tightly coupled to its predecessors. Permissioned blockchains place special emphasis on transactions throughput. In this paper we present FireLedger, which leverages the iterative nature of blockchains in order to improve their throughput in optimistic execution scenarios. FireLedger trades latency for throughput in the sense that in FireLedger the last f + 1 blocks of each node's blockchain are considered tentative, i.e., they may be rescinded in case one of the last f + 1 blocks proposers was Byzantine. Yet, when optimistic assumptions are met, a new block is decided in each communication step, which consists of a proposer that sends only its proposal and all other participants are sending a single bit each. In our performance study FireLedger obtained 20% -- 600% better throughput than state of the art protocols like HotStuff and BFT-SMaRt, depending on the configuration.
Yehonatan Buchnik, Roy Friedman 0001
Proc. VLDB Endow.2
2019 Access Strategies for Network Caching
abstract
Having multiple data stores that can potentially serve content is common in modern networked applications. Data stores often publish approximate summaries of their content to enable effective utilization. Since these summaries are not entirely accurate, forming an efficient access strategy to multiple data stores becomes a complex risk management problem.This paper formally models this problem, and introduces practical algorithms with guaranteed approximation ratios, and in particular we show that our algorithms are optimal in a variety of settings. We also perform an extensive simulation study based on real data, and show that our algorithms are more robust than existing heuristics. That is, they exhibit near optimal performance in various settings whereas the efficiency of existing approaches depends upon system parameters that may change over time, or be otherwise unknown.
Itamar Cohen, Gil Einziger, Roy Friedman 0001, Gabriel Scalosub
INFOCOM3
2019 Nitrosketch: robust and general sketch-based monitoring in software switches
abstract
Software switches are emerging as a vital measurement vantage point in many networked systems. Sketching algorithms or sketches, provide high-fidelity approximate measurements, and appear as a promising alternative to traditional approaches such as packet sampling. However, sketches incur significant computation overhead in software switches. Existing efforts in implementing sketches in virtual switches make sacrifices on one or more of the following dimensions: performance (handling 40 Gbps line-rate packet throughput with low CPU footprint), robustness (accuracy guarantees across diverse workloads), and generality (supporting various measurement tasks).
Zaoxing Liu, Ran Ben-Basat, Gil Einziger, Yaron Kassner, Vladimir Braverman, Roy Friedman 0001, Vyas Sekar
SIGCOMM6
2019 Succinct Summing over Sliding Windows
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Yaron Kassner
Algorithmica3
2019 Give me some slack: Efficient network measurements
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001
Theor. Comput. Sci.3
2019 Randomized Admission Policy for Efficient Top-k, Frequency, and Volume Estimation
abstract
Network management protocols often require timely and meaningful insight about per flow network traffic. This paper introduces Randomized Admission Policy (RAP) -a novel algorithm for the frequency, top-k, and byte volume estimation problems, which are fundamental in network monitoring. We demonstrate space reductions compared to the alternatives, for the frequency estimation problem, by a factor of up to 32 on real packet traces and up to 128 on heavy-tailed workloads. For top-k identification, RAP exhibits memory savings by a factor of between 4 and 64 depending on the workloads' skewness. These empirical results are backed by formal analysis, indicating the asymptotic space improvement of our probabilistic admission approach. In Addition, we present d-way RAP, a hardware friendly variant of RAP that empirically maintains its space and accuracy benefits.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Yaron Kassner
IEEE/ACM Trans. Netw.4
2018 Brief Announcement: Give Me Some Slack: Efficient Network Measurements
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001
ICALP3
2018 Pay for a Sliding Bloom Filter and Get Counting, Distinct Elements, and Entropy for Free
abstract
For many networking applications, recent data is more significant than older data, motivating the need for sliding window solutions. Various capabilities, such as DDoS detection and load balancing, require insights about multiple metrics including Bloom filters, per-flow counting, count distinct and entropy estimation. In this work, we present a unified construction that solves all the above problems in the sliding window model. Our single solution offers a better space to accuracy tradeoff than the state-of-the-art for each of these individual problems! We show this both analytically and by running multiple real Internet backbone and datacenter packet traces.
Eran Assaf, Ran Ben-Basat, Gil Einziger, Roy Friedman 0001
INFOCOM4
2018 Volumetric Hierarchical Heavy Hitters
abstract
Hierarchical heavy hitters (HHH) identification is useful for various network utilities such as anomaly detection, DDoS mitigation, and traffic analysis. However, the increasing support for jumbo frames enables attackers to overload the system with fewer packets, avoiding detection by packet counting techniques. This paper suggests an efficient algorithm for detecting HHH based on their traffic volume that asymptotically improves the runtime of previous works. We implement our algorithm in Open vSwitch (OVS) and incur a 4-6% overhead compared to a 42% throughput reduction experienced by the state-of-the-art.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Marcelo Caggiani Luizelli, Erez Waisbard
MASCOTS3
2018 Give Me Some Slack: Efficient Network Measurements
abstract
Many networking applications require timely access to recent network measurements, which can be captured using a sliding window model. Maintaining such measurements is a challenging task due to the fast line speed and scarcity of fast memory in routers. In this work, we study the impact of allowing slack in the window size on the asymptotic requirements of sliding window problems. That is, the algorithm can dynamically adjust the window size between W and W(1+tau) where tau is a small positive parameter. We demonstrate this model's attractiveness by showing that it enables efficient algorithms to problems such as Maximum and General-Summing that require Omega(W) bits even for constant factor approximations in the exact sliding window model. Additionally, for problems that admit sub-linear approximation algorithms such as Basic-Summing and Count-Distinct, the slack model enables a further asymptotic improvement. The main focus of the paper is on the widely studied Basic-Summing problem of computing the sum of the last W integers from {0,1 ...,R} in a stream. While it is known that Omega(W log R) bits are needed in the exact window model, we show that approximate windows allow an exponential space reduction for constant tau. Specifically, for tau=Theta(1), we present a space lower bound of Omega(log(RW)) bits. Additionally, we show an Omega(log (W/epsilon)) lower bound for RW epsilon additive approximations and a Omega(log (W/epsilon)+log log R) bits lower bound for (1+epsilon) multiplicative approximations. Our work is the first to study this problem in the exact and additive approximation settings. For all settings, we provide memory optimal algorithms that operate in worst case constant time. This strictly improves on the work of [Mayur Datar et al., 2002] for (1+epsilon)-multiplicative approximation that requires O(epsilon^(-1) log(RW)log log (RW)) space and performs updates in O(log (RW)) worst case time. Finally, we show asymptotic improvements for the Count-Distinct, General-Summing and Maximum problems.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001
MFCS3
2018 Adaptive Software Cache Management
abstract
Developing a silver bullet software cache management policy is a daunting task due to the variety of potential workloads. In this paper, we investigate an adaptivity mechanism for software cache management schemes which offer tuning parameters targeted at the frequency vs. recency bias in the workload. The goal is automatic tuning of the parameters for best performance based on the workload without any manual intervention. We study two approaches for this problem, a hill climbing solution and an indicator based solution. In hill climbing, we repeatedly reconfigure the system hoping to find its best setting. In the indicator approach, we estimate the workloads' frequency vs. recency bias and adjust the parameters accordingly in a single swoop.
Gil Einziger, Ohad Eytan, Roy Friedman 0001, Ben Manes
Middleware3
2018 Space Efficient Elephant Flow Detection
abstract
Identifying the large flows in terms of byte volume, known as elephant flows, is a fundamental capability that many network algorithms require. While optimal solutions that find the largest flows in terms of packet-count are known [5], constant update time algorithms for byte-volume were only recently discovered [1, 2]. Here, we propose an improved variant of the DIMSUM algorithm [2] that reduces the space requirement by 50% while allowing O(1) update time.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001
SYSTOR3
2018 Fast flow volume estimation
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001
Pervasive Mob. Comput.3
2018 Stream Frequency Over Interval Queries
abstract
Stream frequency measurements are fundamental in many data stream applications such as financial data trackers, intrusion-detection systems, and network monitoring. Typically, recent data items are more relevant than old ones, a notion we can capture through a sliding window abstraction. This paper considers a generalized sliding window model that supports stream frequency queries over an interval given at query time. This enables drill-down queries, in which we can examine the behavior of the system in finer and finer granularities. For this model, we asymptotically improve the space bounds of existing work, reduce the update and query time to a constant, and provide deterministic solutions. When evaluated over real Internet packet traces, our fastest algorithm processes items 90--250 times faster, serves queries at least 730 times quicker and consumes at least 40% less space than the best known method.
Ran Ben-Basat, Roy Friedman 0001, Rana Shahout
Proc. VLDB Endow.2
2018 ICE Buckets: Improved Counter Estimation for Network Measurement
Gil Einziger, Benny Fellman, Roy Friedman 0001, Yaron Kassner
IEEE/ACM Trans. Netw.3
2017 TinyCache - An Effective Cache Admission Filter
abstract
Effective management policies for datastore caches should provide good hit-ratios on a large number of workloads, operate in constant time, and maintain a small amount of metadata. In certain workloads, cache stability is an important metric, as limiting the number of cache updates can improve power consumption, increase the life expectancy of flash memories, and conserve network bandwidth in distributed settings. This paper introduces TinyCache, a compact table based management policy for datastore caches. TinyCache achieves similar hit ratio compared to the leading alternatives while operating in worst case constant time and only accessing a fixed sized memory word for each update. TinyCache encodes its metadata in a memory optimal manner and reduces the number of cache updates by up to X6 compared to state of the art.
Dolev Adas, Gil Einziger, Roy Friedman 0001
GLOBECOM3
2017 Randomized admission policy for efficient top-k and frequency estimation
abstract
Network management protocols often require timely and meaningful insight about per flow network traffic. This paper introduces Randomized Admission Policy (RAP) - a novel algorithm for the frequency and top-k estimation problems, which are fundamental in network monitoring. We demonstrate space reductions compared to the alternatives by a factor of up to 32 on real packet traces and up to 128 on heavy-tailed workloads. For top-k identification, RAP exhibits memory savings by a factor of between 4 and 64 depending on the workloads' skewness. These empirical results are backed by formal analysis, indicating the asymptotic space improvement of our probabilistic admission approach. Additionally, we present d-Way RAP, a hardware friendly variant of RAP that empirically maintains its space and accuracy benefits.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Yaron Kassner
INFOCOM3
2017 Optimal elephant flow detection
abstract
Monitoring the traffic volumes of elephant flows, including the total byte count per flow, is a fundamental capability for online network measurements. We present an asymptotically optimal algorithm for solving this problem in terms of both space and time complexity. This improves on previous approaches, which can only count the number of packets in constant time. We evaluate our work on real packet traces, demonstrating an up to X2.5 speedup compared to the best alternative.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Yaron Kassner
INFOCOM3
2017 Hardening Cassandra Against Byzantine Failures
abstract
Cassandra is one of the most widely used distributed data stores these days. Cassandra supports flexible consistency guarantees over a wide-column data access model and provides almost linear scale-out performance. This enables application developers to tailor the performance and availability of Cassandra to their exact application's needs and required semantics. Yet, Cassandra is designed to withstand benign failures, and cannot cope with most forms of Byzantine attacks. In this work, we present an analysis of Cassandra's vulnerabilities and propose protocols for hardening Cassandra against Byzantine failures. We examine several alternative design choices and compare between them both qualitatively and empirically by using the Yahoo! Cloud Serving Benchmark (YCSB) performance benchmark. We include incremental performance analysis for our algorithmic and cryptographic adjustments, supporting our design choices.
Roy Friedman 0001, Roni Licher
OPODIS1
2017 Constant Time Updates in Hierarchical Heavy Hitters
abstract
Monitoring tasks, such as anomaly and DDoS detection, require identifying frequent flow aggregates based on common IP prefixes. These are known as hierarchical heavy hitters (HHH), where the hierarchy is determined based on the type of prefixes of interest in a given application. The per packet complexity of existing HHH algorithms is proportional to the size of the hierarchy, imposing significant overheads.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Marcelo Caggiani Luizelli, Erez Waisbard
SIGCOMM3
2017 Counting distinct elements over sliding windows
abstract
In Distributed Denial of Service (DDoS) attacks, an attacker tries to disable a service with a flood of seemingly legitimate requests from multiple devices; this is usually accompanied by a sharp spike in the number of distinct IP addresses / flows accessing the system in a short time frame. Hence, the number of distinct elements over sliding windows is a fundamental signal in DDoS identification. Additionally, assessing whether a specific flow has recently accessed the system, known as the Set Membership problem, can help us identify the attacking parties. Here, we show how to extend the functionality of a state of the art algorithm for set membership over a W elements sliding window. We now also support estimation of the distinct flow count, using as little as log2 (W) additional bits.
Eran Assaf, Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Yaron Kassner
SYSTOR4
2017 TinySet - An Access Efficient Self Adjusting Bloom Filter Construction
abstract
Bloom filters are a very popular and efficient data structure for approximate set membership queries. However, Bloom filters have several key limitations as they require 44% more space than the lower bound, their operations access multiple memory words, and they do not support removals. This paper presents TinySet, an alternative Bloom filter construction that is more space efficient than Bloom filters for false positive rates smaller than 2.8%, accesses only a single memory word and partially supports removals. TinySet is mathematically analyzed and extensively tested and is shown to be fast and more space efficient than a variety of Bloom filter variants. TinySet also has low sensitivity to configuration parameters and is therefore more flexible than a Bloom filter.
Gil Einziger, Roy Friedman 0001
IEEE/ACM Trans. Netw.2
2017 TinyLFU: A Highly Efficient Cache Admission Policy
abstract
This article proposes to use a frequency-based cache admission policy in order to boost the effectiveness of caches subject to skewed access distributions. Given a newly accessed item and an eviction candidate from the cache, our scheme decides, based on the recent access history, whether it is worth admitting the new item into the cache at the expense of the eviction candidate. This concept is enabled through a novel approximate LFU structure called TinyLFU , which maintains an approximate representation of the access frequency of a large sample of recently accessed items. TinyLFU is very compact and lightweight as it builds upon Bloom filter theory. We study the properties of TinyLFU through simulations of both synthetic workloads and multiple real traces from several sources. These simulations demonstrate the performance boost obtained by enhancing various replacement policies with the TinyLFU admission policy. Also, a new combined replacement and eviction policy scheme nicknamed W-TinyLFU is presented. W-TinyLFU is demonstrated to obtain equal or better hit ratios than other state-of-the-art replacement policies on these traces. It is the only scheme to obtain such good results on all traces.
Gil Einziger, Roy Friedman 0001, Ben Manes
ACM Trans. Storage2
2016 Heavy hitters in streams and sliding windows
abstract
Identifying heavy hitter flows is a fundamental problem in various network domains. The well established method of using sketches to approximate flow statistics suffers from space inefficiencies. In addition, flow arrival rates are dynamic, thus keeping track of the most recent heavy hitters poses a challenge. Sliding window approximations address this problem, reducing space at the cost of increasing point query time. This paper presents two novel algorithms for identifying heavy hitters in streams and sliding windows. Both algorithms use statically allocated memory and support constant time point queries. For sliding windows, this is an asymptotic improvement over previous work. We also demonstrate reduced memory requirements of up to 85% in streams and 66% in sliding windows over synthetic and real Internet packet traces.
Ran Ben-Basat, Gil Einziger, Roy Friedman 0001, Yaron Kassner
INFOCOM3
2016 File System Usage in Android Mobile Phones
abstract
In this paper, we report on the analysis of data from Android mobile phones of 38 users, composed of access traces of the users' mobile file systems during 30 days. We shed new light on the file usage patterns and present the data in terms of file size distributions, file sessions, file lifetime, file access activity and read / write access patterns. We characterize different distributions and extract conclusions about usage patterns of Android file systems.
Roy Friedman 0001, David Sainz
SYSTOR1
2016 Shades: Expediting Kademlia's lookup process
Gil Einziger, Roy Friedman 0001, Yoav Kantor
Comput. Networks2
2015 TinySet - An Access Efficient Self Adjusting Bloom Filter Construction
abstract
Bloom filters are a very popular and efficient data structure for approximate set membership queries. However, Bloom filters have several key limitations as they require 44% more space than the lower bound, their operations access multiple memory words and they do not support removals. This work presents TinySet, an alternative Bloom filter construction that is more space efficient than Bloom filters for false positive rates smaller than 2.8%, accesses only a single memory word and partially supports removals. TinySet is mathematically analyzed and extensively tested and is shown to be fast and more space efficient than a variety of Bloom filter variants. TinySet also has low sensitivity to configuration parameters and is therefore more flexible than a Bloom filter.
Gil Einziger, Roy Friedman 0001
ICCCN2
2015 Distilling the ingredients of P2P live streaming systems
abstract
Peer-to-peer live streaming systems involve complex engineering and are difficult to test and to deploy. To cut through the complexity, we advocate such systems be designed by composing ingredients: a novel abstraction denoting the smallest interoperable units of code that each express a single design choice. We present a system, STREAMAID, that provides tools for designing protocols in terms of ingredients, systematically testing the impact of every design decision in a simulator, and deploying them in a wide-area testbed such as PlanetLab for evaluation. We show how to decompose popular P2P live streaming systems, such as CoolStreaming, BitTorrent Live and others, into ingredients and how STREAMAID can help optimize and adapt these protocols. By experimenting with the essential building blocks of which P2P live streaming protocols are comprised, we gain a unique vantage point of their relative quality, their bottlenecks and their potential for future improvement.
Roy Friedman 0001, Alexander Libov, Ymir Vigfusson
P2P1
2015 Probabilistic Byzantine Tolerance for Cloud Computing
abstract
Tolerating Byzantine failures in the context of cloud computing is costly. Traditional BFT protocols induce a fixed degree of replication for computations and are therefore wasteful. This paper explores probabilistic Byzantine tolerance, in which computation tasks are replicated on dynamic replication sets whose size is determined based on ensuring probabilistic thresholds of correctness. The probabilistic assessment of a trustworthy output by selecting reputable nodes allows a significant reduction in the number of nodes involved in each computation task. The paper further studies several reputation management policies, including the one used by BOINC as well as a couple of novel ones, in terms of their impact of the possible damage inflicted on the system by various Byzantine behavior strategies, and reports some encouraging insights.
Luciana Arantes, Roy Friedman 0001, Olivier Marin, Pierre Sens 0001
SRDS2
2015 Chameleon - a group communication framework for smartphones
abstract
Summary This paper reports about our experience in designing and developingChameleon, a highly portable and adaptable group communication framework for smartphones. Chameleon owes its level of portability to several design choices, including the following: (i) a layered architecture, where the headers of each layer have a standard XML‐based format, enabling automatic, error‐resistant generation of efficient serialization code in any platform; (ii) reliance only on the J2ME library, which serves as least common denominator for Java dialects and facilitates automatic translation to.NET; (iii) having flexible membership models; and (iv) supporting multiple concurrent protocol stacks.Through a single codebase,Chameleonis currently available as an open‐source project for J2ME, J2SE, Android,.NET CF, and.NET.Chameleonis easily extendable and is bundled with tools, configurations, and third‐party code tuned in a way that lifts some of the burden normally associated with multiplatform development for smartphones. Both the header generation from XML and automatic translation to.NET features of Chameleon are readily available to any application that is based on it.Chameleon's threading model separates between execution of internal layers and application's code and by that protects one from the other. As we describe in the paper, it simplifies layers' development and allows the protocol stack to easily block application calls when this is required by internal algorithms. Additionally, this model simplifies testing, and an extensive testing framework is supplied along withChameleon, which is also usable for testing of application‐specific layers. Copyright © 2014 John Wiley & Sons, Ltd.
Alex Dvinsky, Roy Friedman 0001
Softw. Pract. Exp.2
2015 A generic decentralized trust management framework
abstract
Summary This paper describes TRUSTPACK, a decentralized trust management framework that provides trust management as a generic service. TRUSTPACK is unique in that it does not provide a central service. Instead, it is run by many autonomous services. This design enables TRUSTPACK to alleviate privacy concerns, as well as potentially provide better personalization and scalability when compared with current centralized solutions. A major component of TRUSTPACK is a generic decentralized graph query processing framework called GRAPHPACK, which was also developed as part of this work. GRAPHPACK consists of a decentralized graph processing language as well as an execution engine, as elaborated in this paper. The paper also presents several examples and a case study showing how TRUSTPACK can be used to handle various trust management scenarios, as well as its incorporation in an existing third party P2P file sharing application. Prototypes of TRUSTPACK and GRAPHPACK are available as open source projects at http://code.google.com/p/trustpack/ and http://code.google.com/p/graphpack/ , respectively. Copyright © 2013 John Wiley & Sons, Ltd.
Roy Friedman 0001, Amit Portnoy
Softw. Pract. Exp.1
2014 Shades: Expediting Kademlia's Lookup Process
Gil Einziger, Roy Friedman 0001, Yoav Kantor
Euro-Par2
2014 MOLStream: A Modular Rapid Development and Evaluation Framework for Live P2P Streaming
abstract
We present MOL Stream, a modular framework for rapid development and evaluation of P2P live streaming systems. MOL Stream allows P2P streaming protocols to be decomposed into basic blocks, each associated with a standard functional specification. By exposing structural commonalities between these components, MOL Stream enables specific implementations of these building blocks to be combined in order to devise, refine and evaluate new P2P live streaming protocols. Our approach offers several benefits. First, block encapsulation entails that more advanced individual components, e.g., the overlay, can seamlessly replace existing ones without affecting the rest of the system. As a case study, we show how MOL Stream can seamlessly substitute the overlay used by DONet/Coolstreaming, a popular P2P live streaming implementation, for an improved version. Second, MOL Stream facilitates the comparison between various protocols over local clusters or wide-area test beds such as Planet Lab. The combination of rapid prototyping and minimum effort valuation enables researchers and students to faster understand how various design choices at different levels impact the performance and scalability of the protocol, as shown through several examples in this paper. MOL Stream is written in Java and is freely available as an open-source project at https://sourceforge.net/projects/molstream/.
Roy Friedman 0001, Alexander Libov, Ymir Vigfusson
ICDCS1
2014 Replicated erasure codes for storage and repair-traffic efficiency
abstract
This paper introduces a new family of redundancy schemes for distributed storage systems, called replicated erasure codes (REC), which combine the storage-space efficiency of erasure codes and the repair-traffic efficiency of replication. A formal model for analyzing the storage and repair-traffic costs under availability and persistency constraints is also developed. It is shown that under parameters that characterize common P2P environments, REC generally achieves better results than each of the two methods separately.
Roy Friedman 0001, Yoav Kantor, Amir Kantor
P2P1
2014 TinyLFU: A Highly Efficient Cache Admission Policy
abstract
This paper proposes to use a frequency based cache admission policy in order to boost the effectiveness of caches subject to skewed access distributions. Rather than deciding on which object to evict, TinyLFU decides, based on the recent access history, whether it is worth admitting an accessed object into the cache at the expense of the eviction candidate. Realizing this concept is enabled through a novel approximate LFU structure called TinyLFU, which maintains an approximate representation of the access frequency of recently accessed objects. TinyLFU is extremely compact and lightweight as it builds upon Bloom filter theory. The paper shows an analysis of the properties of TinyLFU including simulations of both synthetic workloads as well as YouTube and Wikipedia traces.
Gil Einziger, Roy Friedman 0001
PDP2
2014 Postman: An Elastic Highly Resilient Publish/Subscribe Framework for Self Sustained Service Independent P2P Networks
Gil Einziger, Roy Friedman 0001
SSS2
2013 Hybrid Distributed Consensus
Roy Friedman 0001, Gabriel Kliot, Alex Kogan
OPODIS1
2013 Kaleidoscope: Adding colors to Kademlia
abstract
Kademlia is considered to be one of the most effective key based routing protocols. It is nowadays implemented in many file sharing peer-to-peer networks such as BitTorrent, KAD, and Gnutella. This paper introduces Kaleidoscope, a novel routing/caching scheme designed to significantly reduce the cost of lookup operations in Kademlia by using a color-based distributed cache. Moreover, Kaleidoscope greatly improves load balancing among the nodes and reduces the well documented hot spots problem. The paper also includes an extensive performance study demonstrating the benefits of Kaleidoscope.
Gil Einziger, Roy Friedman 0001, Eyal Kibbar
P2P2
2013 An advertising mechanism for P2P networks
abstract
In order for P2P systems to be viable, users must be given incentives to donate resources. Such incentives can be in the form of tit-for-tat like mechanisms, in which a user is rewarded with better service for contributing resources to the system. Alternatively, such incentives can be economical, i.e., users get paid for their contribution. In particular, the latter can be achieved through a P2P advertisement mechanism. This paper investigates how to incorporate an advertisement mechanism into a P2P system to serve as an incentive to donate resources, especially for services in which users often interact with the system through mobile devices. First, the precise P2P advertisement dissemination model is presented. Second, the paper proposes and explores several advertisement dissemination schemes combined with a few payment models and compares between them through simulations. The reported results are encouraging for this direction and in particular identify payment models whose payment is super-linear with the availability of donated machines. This means that they serve as good incentives for owners of donated machines to keep them connected to the P2P network for long durations.
Roy Friedman 0001, Alexander Libov
P2P1
2013 A density-driven publish subscribe service for mobile ad-hoc networks
Roy Friedman 0001, Anna Kaplun Shulman
Ad Hoc Networks1
2013 On Power and Throughput Tradeoffs of WiFi and Bluetooth in Smartphones
abstract
This paper describes a combined power and throughput performance study of WiFi and Bluetooth usage in smartphones. The work measures the obtained throughput in various settings while employing each of these technologies, and the power consumption level associated with them. In addition, the power requirements of Bluetooth and WiFi in their respective noncommunicating modes are also compared. The study reveals several interesting phenomena and tradeoffs. In particular, the paper identifies many situations in which WiFi is superior to Bluetooth, countering previous reports. The study also identifies a couple of scenarios that are better handled by Bluetooth. The conclusions from this study suggest preferred usage patterns, as well as operative suggestions for researchers and smartphone developers. This includes a cross-layer optimization for TCP/IP that could greatly improve the throughput to power ratio whenever the transmitter is more capable than the receiver.
Roy Friedman 0001, Alex Kogan, Yevgeny Krivolapov
IEEE Trans. Mob. Comput.1
2012 Efficient and Reliable Multicast in Multi-radio Networks
abstract
This paper investigates a novel efficient approach to utilize multiple radio interfaces for enhancing the performance of reliable multicasts from a single sender to a group of receivers. In the proposed scheme, one radio channel (and interface) is dedicated only for recovery information transmissions. We apply this concept to both ARQ and hybrid ARQ+FEC protocols, formally analyzing the number of packets each receiver needs to process in both our approach and in the common single channel approach. We also present a corresponding efficient protocol, and study its performance by simulation. Both the formal analysis and the simulations demonstrate the benefits of our scheme.
Roy Friedman 0001, Alex Kogan
SRDS1
2011 On power and throughput tradeoffs of WiFi and Bluetooth in smartphones
abstract
This paper describes a combined power and throughput performance study of WiFi and Bluetooth usage in smartphones. The study reveals several interesting phenomena and tradeoffs. The conclusions from this study suggest preferred usage patterns, as well as operative suggestions for researchers and smartphone developers.
Roy Friedman 0001, Alex Kogan, Yevgeny Krivolapov
INFOCOM1
2011 On Reliable Dissemination in Wireless Ad Hoc Networks
abstract
Reliable broadcast is a basic service for many collaborative applications as it provides reliable dissemination of the same information to many recipients. This paper studies three common approaches for achieving scalable reliable broadcast in ad hoc networks, namely probabilistic flooding, counter-based broadcast, and lazy gossip. The strength and weaknesses of each scheme are analyzed, and a new protocol that combines these three techniques, called RAPID, is developed. Specifically, the analysis in this paper focuses on the trade-offs between reliability (percentage of nodes that receive each message), latency, and the message overhead of the protocol. Each of these methods excel in some of these parameters, but no single method wins in all of them. This motivates the need for a combined protocol that benefits from all of these methods and allows to trade between them smoothly. Interestingly, since the RAPID protocol only relies on local computations and probability, it is highly resilient to mobility and failures and even selfish behavior. By adding authentication, it can even be made malicious tolerant. Additionally, the paper includes a detailed performance evaluation by simulation. The simulations confirm that RAPID obtains higher reliability with low latency and good communication overhead compared with each of the individual methods.
Vadim Drabkin, Roy Friedman 0001, Gabriel Kliot, Marc Segal
IEEE Trans. Dependable Secur. Comput.2
2010 Brief announcement: deterministic dominating set construction in networks with bounded degree
abstract
This paper considers the problem of calculating dominating sets in bounded degree networks. In these networks, the maximal degree of any node is bounded by δ, which is usually significantly smaller than n, the total number of nodes in the system. Such networks arise in various settings of wireless and peer-to-peer communication. A trivial approach of choosing all nodes into the dominating set yields an algorithm with an approximation ratio of δ + 1. We show that any deterministic algorithm with a non-trivial approximation ratio requires Ω(log* n) rounds, meaning effectively that no local o(δ)-approximation deterministic algorithm may ever exist. On the positive side, we show two deterministic algorithms that achieve log δ and 2 log δ-approximation in O(δ3 + log* n) and O(δ2 logδ + log* n) time, respectively. These algorithms rely on coloring rather than node IDs to break symmetry.
Roy Friedman 0001, Alex Kogan
PODC1
2010 DEEP: Density-based proactive data dissemination protocol for wireless sensor networks with uncontrolled sink mobility
Massimo Vecchio, Aline Carneiro Viana, Artur Ziviani, Roy Friedman 0001
Comput. Commun.4
2010 Probabilistic quorum systems in wireless Ad Hoc networks
abstract
Quorums are a basic construct in solving many fundamental distributed computing problems. One of the known ways of making quorums scalable and efficient is by weakening their intersection guarantee to being probabilistic. This article explores several access strategies for implementing probabilistic quorums in ad hoc networks. In particular, we present the first detailed study of asymmetric probabilistic biquorum systems, that allow to mix different access strategies and different quorums sizes, while guaranteeing the desired intersection probability. We show the advantages of asymmetric probabilistic biquorum systems in ad hoc networks. Such an asymmetric construction is also useful for other types of networks with nonuniform access costs (e.g, peer-to-peer networks). The article includes a formal analysis of these approaches backed up by an extensive simulation-based study. The study explores the impact of various parameters such as network size, network density, mobility speed, and churn. In particular, we show that one of the strategies that uses random walks exhibits the smallest communication overhead, thus being very attractive for ad hoc networks.
Roy Friedman 0001, Gabriel Kliot, Chen Avin
ACM Trans. Comput. Syst.1
2009 Power Aware Management Middleware for Multiple Radio Interfaces
Roy Friedman 0001, Alex Kogan
Middleware1
2009 3DLS: density-driven data location service for mobile ad-hoc networks
abstract
Finding data items is one of the most basic services of any distributed system. It is particular challenging in ad-hoc networks, due to their inherent decentralized nature and lack of infrastructure. A data location service (DLS) provides this capability. This paper presents 3DLS, a novel density driven data location service. 3DLS is based on performing biased walks over a density based virtual topography. 3DLS also includes an autonomic dynamic configuration mechanism for adapting the lengths of the walks, in order to ensure good performance in varying circumstances and loads. This is without any explicit knowledge of the network characteristics, such as size, mobility speed, etc. Moreover, 3DLS does not rely on geographical knowledge, its decisions are based only on local information, it does not invoke multi-hop routing, and it avoids flooding the network. The paper includes a detailed performance study of 3DLS, carried by simulations, which compares 3DLS to other known approaches. The simulations results validate the viability of 3DLS.
Roy Friedman 0001, Noam Mori
MobiHoc1
2009 Efficient Power Utilization in Multi-radio Wireless Ad Hoc Networks
Roy Friedman 0001, Alex Kogan
OPODIS1
2009 Brief Announcement: Efficient Utilization of Multiple Interfaces in Wireless Ad Hoc Networks
Roy Friedman 0001, Alex Kogan
DISC1
2009 Efficient route discovery in hybrid networks
Roy Friedman 0001, Ari Shotland, Gwendal Simon
Ad Hoc Networks1
2008 Probabilistic quorum systems in wireless ad hoc networks
abstract
Quorums are a basic construct in solving many fundamental distributed computing problems. One of the known ways of making quorums scalable and efficient is by weakening their intersection guarantee to being probabilistic. This paper explores several access strategies for implementing probabilistic quorums in ad hoc networks. In particular, we present the first detailed study of asymmetric probabilistic bi-quorum systems and show its advantages in ad hoc networks. The paper includes both a formal analysis of these approaches backed by a simulation based study. In particular, we show that one of the strategies, based on random walks, exhibits the smallest communication overhead.
Roy Friedman 0001, Gabriel Kliot, Chen Avin
DSN1
2008 Model-based performance evaluation of distributed checkpointing protocols
Adnan Agbaria, Roy Friedman 0001
Perform. Evaluation2
2008 RaWMS - Random Walk Based Lightweight Membership Service for Wireless Ad Hoc Networks
abstract
This article presents RaWMS, a novel lightweight random membership service for ad hoc networks. The service provides each node with a partial uniformly chosen view of network nodes. Such a membership service is useful, for example, in data dissemination algorithms, lookup and discovery services, peer sampling services, and complete membership construction. The design of RaWMS is based on a novel reverse random walk (RW) sampling technique. The article includes a formal analysis of both the reverse RW sampling technique and RaWMS and verifies it through a detailed simulation study. In addition, RaWMS is compared both analytically and by simulations with a number of other known methods such as flooding and gossip-based techniques.
Ziv Bar-Yossef, Roy Friedman 0001, Gabriel Kliot
ACM Trans. Comput. Syst.2
2007 RAPID: Reliable Probabilistic Dissemination in Wireless Ad-Hoc Networks
abstract
In this paper, we propose a novel reliable probabilistic dissemination protocol, RAPID, for mobile wireless ad-hoc networks that tolerates message omissions, node crashes, and selfish behavior. The protocol employs a combination of probabilistic forwarding with deterministic corrective measures. The forwarding probability is set based on the observed number of nodes in each one-hop neighborhood, while the deterministic corrective measures include deterministic gossiping as well as timer based corrections of the probabilistic process. These aspects of the protocol are motivated by a theoretical analysis that is also presented in the paper, which explains why this unique protocol design is inherent to ad-hoc networks environments. Since the protocol only relies on local computations and probability, it is highly resilient to mobility and failures. The paper includes a detailed performance evaluation by simulation. We compare the performance and the overhead of RAPID with the performance of other probabilistic approaches. Our results show that RAPID achieves a significantly higher node coverage with a smaller overhead.
Vadim Drabkin, Roy Friedman 0001, Gabriel Kliot, Marc Segal
SRDS2
2007 Managed Agreement: Generalizing two fundamental distributed agreement problems
Emmanuelle Anceaume, Roy Friedman 0001, Maria Potop-Butucaru
Inf. Process. Lett.2
2007 Asynchronous Agreement and Its Relation with Error-Correcting Codes
abstract
The condition-based approach identifies sets of input vectors, called conditions, for which it is possible to design an asynchronous protocol solving a distributed problem despite process crashes. This paper establishes a direct correlation between distributed agreement problems and error-correcting codes. In particular, crash failures in distributed agreement problems correspond to erasure failures in error-correcting codes and Byzantine and value domain faults correspond to corruption errors. This correlation is exemplified by concentrating on two well-known agreement problems, namely, consensus and interactive consistency, in the context of the condition-based approach. Specifically, the paper presents the following results: first, it shows that the conditions that allow interactive consistency to be solved despite fccrashes and fcvalue domain faults correspond exactly to the set of error-correcting codes capable of recovering from fcerasures and fccorruptions. Second, the paper proves that consensus can be solved despite fccrash failures if the condition corresponds to a code whose Hamming distance is fc+ 1 and Byzantine consensus can be solved despite fbByzantine faults if the Hamming distance of the code is 2 fb+ 1. Finally, the paper uses the above relations to establish several results in distributed agreement that are derived from known results in error-correcting codes and vice versa.
Roy Friedman 0001, Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
IEEE Trans. Computers1
2007 On the Respective Power of *P and *S to Solve One-Shot Agreement Problems
abstract
Unreliable failure detectors are abstract devices that, when added to asynchronous distributed systems, enable solving distributed computing problems (e.g., consensus) that otherwise would be impossible to solve in these systems. This paper focuses on two classes of failure detectors defined by Chandra and Toueg, namely, the classes denoted diamP (eventually perfect) and diamS (eventually strong). Both classes include failure detectors that eventually detect permanently all process crashes, but while the failure detectors of diamP eventually make no erroneous suspicions, the failure detectors of diamS are only required to eventually not suspect a single correct process. Informally, in a one-shot agreement problem, a new problem instance is created each time the processes propose new values to be decided on (e.g., consensus is one-shot). In such a context, this paper addresses the following question related to the comparative power of these classes, namely: "Are there one-shot agreement problems that can be solved in asynchronous distributed systems with reliable links but prone to process crash failures augmented with op, but cannot be solved when those systems are augmented with diamS?" Surprisingly, the paper shows that the answer to this question is "no." An important consequence of this result is that diamP cannot be the weakest class of failure detectors that enables solving one-shot agreement problems in unreliable asynchronous distributed systems
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Parallel Distributed Syst.1
2006 Practical Byzantine Group Communication
abstract
This paper presents an adaptation of the JazzEnsemble group communication system that enables it to tolerate Byzantine failures. The work emphasizes scalability and good performance in the normal case, i.e., when there are no failures, while providing strong semantics to the application. The paper presents the main concepts and protocols that enable the Byzantine tolerant version of JazzEnsemble to obtain these goals. In particular, this includes fuzzy mute and fuzzy verbose failure detectors, an efficient Byzantine vector consensus protocol, and a novel Byzantine uniform broadcast protocol, as well as modifications at each layer of the system to overcome potential Byzantine attacks. Additionally, high-level protocols only rely on the oral messages model, and thus messages need to be signed only once at a low level of the system. Finally, the paper presents an extensive performance evaluation, which demonstrates the system’s scalability and efficiency, and is used to analyze the sources of performance degradation associated with various aspects of overcoming Byzantine failures.
Vadim Drabkin, Roy Friedman 0001, Alon Kama
ICDCS2
2006 RaWMS -: random walk based lightweight membership service for wireless ad hoc network
abstract
RaWMS is a novel lightweight random membership service for ad hoc networks. The service provides each node with a partial uniformly chosen view of network nodes. Such a membership service is useful, e.g., in data dissemination algorithms, lookup and discovery services, peer sampling services, and complete member-ship construction. The design of RaWMS is based on a novel re-verse random walk (RW) sampling technique. The paper includes a formal analysis of both the reverse RW sampling technique and RaWMS and verifies it through a detailed simulation study. In addition, RaWMS is compared with a number of other known methods such as flooding and gossip-based techniques.
Ziv Bar-Yossef, Roy Friedman 0001, Gabriel Kliot
MobiHoc2
2006 Self-stabilizing Wireless Connected Overlays
Vadim Drabkin, Roy Friedman 0001, Maria Potop-Butucaru
OPODIS2
2006 Adaptive Batching for Replicated Servers
abstract
This paper presents two novel generic adaptive batching schemes for replicated servers. Both schemes are oblivious to the underlying communication protocols. Our novel schemes adapt their batching levels automatically and immediately according to the current communication load. This is done without any explicit monitoring or calibration of the system. Additionally, the paper includes a detailed performance evaluation
Roy Friedman 0001, Erez Hadad
SRDS1
2005 Efficient Byzantine Broadcast in Wireless Ad-Hoc Networks
abstract
This paper presents an overlay based Byzantine tolerant broadcast protocol for wireless ad-hoc networks. The use of an overlay results in a significant reduction in the number of messages. The protocol overcomes Byzantine failures by combining digital signatures, gossiping of message signatures, and failure detectors. These ensure that messages dropped or modified by Byzantine nodes will be detected and retransmitted and that the overlay will eventually consist of enough correct processes to enable message dissemination. An appealing property of the protocol is that it only requires the existence of one correct node in each one-hop neighborhood. The paper also includes a detailed performance evaluation by simulation.
Vadim Drabkin, Roy Friedman 0001, Marc Segal
DSN2
2005 Timed grid routing (TIGR) bites off energy
abstract
Energy efficiency and collisions avoidance are both critical properties to increase the lifetime and effectiveness of wireless networks. This paper proposes a family of algorithms for reducing both energy consumption and packets collisions in ad-hoc networks. In particular, this family of protocols offers a tradeoff between bandwidth utilization and power consumption. The proposed algorithms are based on geographic knowledge to form a virtual grid and on synchronized clocks in order to achieve a collision free locally computable transmission schedule.As an added benefit to the proposed approach, a new efficient location service that is adjusted to the grid oriented geo-routing is also presented. The paper analyzes the proposed family of algorithms for their induced latency and the energy vs. latency tradeoffs they provide.
Roy Friedman 0001, Guy Korland
MobiHoc1
2005 Two Abstractions for Implementing Atomic Objects in Dynamic Systems
Roy Friedman 0001, Michel Raynal, Corentin Travers
OPODIS1
2005 Brief announcement: abstractions for implementing atomic objects in dynamic systems
abstract
No abstract available.
Roy Friedman 0001, Michel Raynal, Corentin Travers
PODC1
2005 Intersecting Sets: a Basic Abstraction for Asynchronous Agreement Problems
abstract
Defining good abstractions is a central issue when one wants to understand the deep structure and basic principles that underlie computing mechanisms. This paper introduces a basic and particularly simple distributed computing abstraction suited to asynchronous distributed agreement problems. This abstraction, called intersecting sets, requires each process to deposit a value and allows each non-faulty process to obtain a subset of these values such that any two such sets have a non-empty intersection. This simple abstraction captures an essential part of distributed agreement problems. After having introduced and motivated this abstraction, the paper investigates its properties, its power and its benefit when solving distributed agreement problems.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
PRDC1
2005 Asynchronous bounded lifetime failure detectors
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.1
2005 Simple and Efficient Oracle-Based Consensus Protocols for Asynchronous Byzantine Systems
abstract
This paper is on the consensus problem in asynchronous distributed systems where (up to f) processes (among n) can exhibit a Byzantine behavior, i.e., can deviate arbitrarily from their specification. One way to solve the consensus problem in such a context consists of enriching the system with additional oracles that are powerful enough to cope with the uncertainty and unpredictability created by the combined effect of Byzantine behavior and asynchrony. This paper presents two kinds of Byzantine asynchronous consensus protocols using two types of oracles, namely, a common coin that provides processes with random values and a failure detector oracle. Both allow the processes to decide in one communication step in favorable circumstances. The first is a randomized protocol for an oblivious scheduler model that assumes n > 6f. The second one is a failure detector-based protocol that assumes n > tif. These protocols are designed to be particularly simple and efficient in terms of communication steps, the number of messages they generate in each step, and the size of messages. So, although they are not optimal in the number of Byzantine processes that can be tolerated, they are particularly efficient when we consider the number of communication steps they require to decide and the number and size of the messages they use. In that sense, they are practically appealing.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
IEEE Trans. Dependable Secur. Comput.1
2004 Locating cache proxies in manets
abstract
Caching Internet based services is a potentially important application for MANETs, as it can improve mobile users' perceived quality of service, reduce their energy consumption, and lower their air-time costs. This paper considers the problem of locating cache proxies in MANETs using several search techniques. The paper first examines several existing and a few novel search techniques including flooding, constrained flooding, a novel dynamic variation of probabilistic flooding, and BFS. These are superimposed on a Maximal Independent Set (MIS), a Connected Dominating Set (DS), and a novel adaptation of BFS-tree based overlays, where each of these overlays is maintained in a self stabilizing manner. The paper also includes a comparison of the performance of these search techniques and overlays by extensive simulations.
Roy Friedman 0001, Maria Potop-Butucaru, Gwendal Simon
MobiHoc1
2004 Brief announcement: veto number and the respective power of eventual failure detectors
abstract
No abstract available.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
PODC1
2004 Simple and Efficient Oracle-Based Consensus Protocols for Asynchronous Byzantine Systems
abstract
This paper is on the consensus problem in asynchronous distributed systems where (up to f) processes (among n) can exhibit a Byzantine behavior, i.e., can deviate arbitrarily from their specification. A way to solve the consensus problem in such a context consists of enriching the system with additional oracles that are powerful enough to cope with the uncertainty and unpredictability created by the combined effect of Byzantine behavior and asynchrony. Considering two types of such oracles, namely, an oracle that provides processes with random values, and a failure detector oracle, the paper presents two families of Byzantine asynchronous consensus protocols. Two of these protocols are particularly noteworthy: they allow the processes to decide in one communication step in favorable circumstances. The first is a randomized protocol that assumes n > 5f. The second one is a failure detector-based protocol that assumes n > 6f. These protocols are designed to be particularly simple and efficient in terms of communication steps, the number of messages they generate in each step, and the size of messages. So, although they are not optimal in the number of Byzantine processes that can be tolerated, they are particularly efficient when we consider the number of communication steps they require to decide, and the number and size of the messages they use. In that sense, they are practically appealing.
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
SRDS1
2004 The Notion of Veto Number and the Respective Power of OP and OS to Solve One-Shot Agreement Problems
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
DISC1
2004 A weakest failure detector-based asynchronous consensus protocol for f<n
Roy Friedman 0001, Achour Mostéfaoui, Michel Raynal
Inf. Process. Lett.1
2004 Quantifying rollback propagation in distributed checkpointing
Adnan Agbaria, Hagit Attiya, Roy Friedman 0001, Roman Vitenberg
J. Parallel Distributed Comput.3
2003 Evaluating Distributed Checkpointing Protocol
abstract
This paper presents an objective measure, called overhead ratio, for evaluating distributed checkpointing protocols. This measure extends previous evaluation schemes by incorporating several additional parameters that are inherent in distributed environments. In particular, we take into account the rollback propagation of the protocol, which impacts the length of the recovery process, and therefore the expected program run-time in executions that involve failures and recoveries. The paper also analyzes several known protocols and compares their overhead ratio.
Adnan Agbaria, Ari Freund 0001, Roy Friedman 0001
ICDCS3
2003 Transparent Fault-Tolerant Java Virtual Machine
abstract
Replication is one of the prominent approaches for obtaining fault tolerance. Implementing replication on commodity hardware and in a transparent fashion, i.e., without changing the programming model, has many challenges. Deciding at what level to implement the replication has ramifications on development costs and portability of the programs. Other difficulties lie in the coordination of the copies in the face of non-determinism. We report on an implementation of transparent fault tolerance at the virtual machine level of Java. We describe the design of the system and present performance results that in certain cases are equivalent to those of non-replicated executions. We also discuss design decisions stemming from implementing replication at the virtual machine level, and the special considerations necessary in order to support symmetric multi-processors (SMP).
Roy Friedman 0001, Alon Kama
SRDS1
2003 On the Locality of Consistency Conditions
Roman Vitenberg, Roy Friedman 0001
DISC2
2003 On the composability of consistency conditions
Roy Friedman 0001, Roman Vitenberg, Gregory V. Chockler
Inf. Process. Lett.1
2002 Distributed Agreement and Its Relation with Error-Correcting Codes
Roy Friedman 0001, Achour Mostéfaoui, Sergio Rajsbaum, Michel Raynal
DISC1
2002 Virtual-machine-based heterogeneous checkpointing
abstract
Abstract Checkpointingan application is the act of saving the application's state during its execution on stable storage, so that if the application fails it can berestartedfrom the last saved state, thereby avoiding loss of the work that was already done. Aheterogeneous checkpoint/restartmechanism allows one to restart an application on a possibly different hardware architecture and/or operating system than those in which the application was saved. This paper explores how to construct such a mechanism at the virtual machine level. That is, rather than dumping the entire state of the application process, the mechanism reported here dumps the state of the application as maintained by a virtual machine. During restart, the saved state is loaded into a new copy of the virtual machine, which continues running from there. The heterogeneous checkpoint/restart mechanism reported here was developed for the OCaml variant of ML. The paper reports on the main issues encountered in building such a mechanism and the design choices made, presents performance evaluations, and discusses some lessons and ideas for extending the work to native code OCaml and Java. Copyright © 2002 John Wiley & Sons, Ltd.
Adnan Agbaria, Roy Friedman 0001
Softw. Pract. Exp.2
2002 Scalable Stability Detection Using Logical Hypercube
abstract
This paper proposes to use a logical hypercube structure for detecting message stability in distributed systems. In particular, a stability detection protocol that uses such a superimposed logical structure is presented, and its scalability is compared with other known stability detection protocols. The main benefits of the logical hypercube approach are scalability, fault-tolerance, and refraining from overloading a single node or link in the system. These benefits become evident both by an analytical comparison and by simulations. Another important feature of the logical hypercube approach is that the performance of the protocol is in general not sensitive to the topology of the underlying physical network.
Roy Friedman 0001, Shiri Manor, Katherine Guo
IEEE Trans. Parallel Distributed Syst.1
2001 Quantifying Rollback Propagation in Distributed Checkpointing
abstract
Proposes a new classification of executions with checkpoints that is based on the notion of k-rollback, indicating the maximal number of checkpoints that may need to be rolled back during recovery. The relation between known execution classes is explored, and it is shown that coordinated checkpointing, SZPF (strictly Z-path free) and ZPF (Z-path free) are 1-rollback mechanisms, while ZCF (Z-cycle free) is (n-1)-rollback, where n is the number of participants in an execution. A new class of executions, called d-BC (d-bounded cycles), is introduced, and is shown to be an [(n-1)/spl middot/d]-rollback mechanism (ZCF is a special case of d-BC for d=1). Finally, a d-BC protocol is presented. This protocol has the nice property that it does not impose any control information overhead on an application's messages, yet it only sends a few control messages of its own. Moreover, the protocol maintains information about recovery lines, which enables very efficient discovery of the most recent recovery line that existed a short time before the failure.
Adnan Agbaria, Hagit Attiya, Roy Friedman 0001, Roman Vitenberg
SRDS3
2000 Implementing a Caching Service for Distributed CORBA Objects
Gregory V. Chockler, Danny Dolev, Roy Friedman 0001, Roman Vitenberg
Middleware3
2000 Consistency Conditions for a CORBA Caching Service
Gregory V. Chockler, Roy Friedman 0001, Roman Vitenberg
DISC2
1999 Symphony: Managing Virtual Servers in the Global Village
Roy Friedman 0001, Assaf Schuster, Ayal Itzkovitz, Eli Biham, Erez Hadad, Vladislav Kalinovsky, Sergey Kleyman, Roman Vitenberg
Euro-Par1
1999 Starfish: Fault-Tolerant Dynamic MPI Programs on Clusters of Workstations
abstract
This paper reports on the architecture and design of Starfish, an environment for executing dynamic (and static) MPI-2 programs on a cluster of workstations. Starfish is unique in being efficient, fault-tolerant, highly available, and dynamic as a system internally, and in supporting fault-tolerance and dynamicity for its application programs as well. Starfish achieves these goals by combining group communication technology with checkpoint/restart, and uses a novel architecture that is both flexible and portable and keeps group communication outside the critical data path, for maximum performance.
Adnan Agbaria, Roy Friedman 0001
HPDC2
1999 Scalable Stability Detection using Logical Hypercube
abstract
This paper proposes to use a logical hypercube structure for detecting message stability in distributed systems. In particular, a stability detection protocol that uses such a superimposed logical structure is presented, and its scalability is compared with other known stability detection protocols. The main benefits of the logical hypercube approach are scalability, fault-tolerance, and refraining from overloading a single node or link in the system. These benefits become evident both by an analytical comparison and by simulations. Another important feature of the logical hypercube approach is that the performance of the protocol is in general not sensitive to the topology of the underlying physical network.
Roy Friedman 0001, Shiri Manor, Katherine Guo
SRDS1
1999 Load-Balancing Schemes for High-Throughput Distributed Fault-Tolerant Servers
abstract
Clusters of workstations, connected by a fast network, are emerging as a viable architecture for building high-throughput fault-tolerant servers. This type of architecture is more scalable and more cost-effective than a tightly coupled multiprocessor and may achieve as good a throughput. Two of the most important issues that a designer of such clustered servers must consider in order for the system to meet its fault-tolerance and throughput goals are the load-balancing scheme and the fault-tolerance scheme that the system will use. This paper explores several combinations of such fault-tolerance and load-balancing schemes and compares their impact on the maximum throughout achievable by the system, and on its survivability. In particular, we show that a fault-tolerance scheme may have an effect on the throughput of the system, while a load-balancing scheme may affect the ability of the system to override failures. We study the scalability of the different schemes under different loads and failure conditions. Our simulations take into consideration the overhead of each scheme, the network contention, and the resource loads.
Roy Friedman 0001, Daniel Mossé
J. Parallel Distributed Comput.1
1999 Middleware support for distributed multimedia and collaborative computing
abstract
Maestro is a middleware support tool for distributed multimedia and collaborative computing applications. These applications share a common need for managing multiple subgroups while providing possibly different quality-of-service guarantees for each of these groups. Maestro's functionality maps well into these requirements, and can significantly shorten the development time of such applications. In this paper, we report on Maestro, and demonstrate its utility in implementing several multimedia and collaborative computing applications. In particular, we provide a detailed description of the implementation of IMUX, a pseudo X-server (proxy) for collaborative computing applications that is based on Maestro. Copyright © 1999 John Wiley & Sons, Ltd.
Kenneth P. Birman, Roy Friedman 0001, Mark Hayden, Injong Rhee
Softw. Pract. Exp.2
1998 Shared Memory Consistency Conditions for Nonsequential Execution: Definitions and Programming Strategies
abstract
To enhance performance on shared memory multiprocessors, various techniques have been proposed to reduce the latency of memory accesses, including pipelining of accesses, out-of-order execution of accesses, and branch prediction with speculative execution. These optimizations can, however, complicate the user's model of memory. This paper attacks the problem of simplifying programming on two fronts. First, a general framework is presented for defining shared memory consistency conditions that allows nonsequential execution of memory accesses. The interface at which conditions are defined is between the program and the system and is architecture-independent. The framework is used to generalize three consistency conditions---sequential consistency, hybrid consistency, and weak consistency---for nonsequential execution. Thus, familiar consistency conditions can be precisely specified even in optimized architectures. Second, three techniques are described for structuring programs so that a shared memory that provides the weaker (and more efficient) condition of hybrid consistency appears to guarantee the stronger (and more costly) condition of sequential consistency. The benefit of these techniques is that sequentially consistent executions are easier to reason about. The first technique statically classifies accesses based on their type. This approach is extremely simple to use and leads to a general technique for writing efficient synchronization code. The third technique is to avoid data races in the program, which was previously studied in a somewhat different setting. Precise, yet short and comprehensible, proofs are provided for the correctness of the programming techniques. Such proofs shed light on the reasons these techniques work; we believe that the insight gained can lead to the development of other techniques.
Hagit Attiya, Soma Chaudhuri, Roy Friedman 0001, Jennifer L. Welch
SIAM J. Comput.3
1998 A Correctness Condition for High-Performance Multiprocessors
abstract
Hybrid consistency, a consistency condition for shared memory multiprocessors, attempts to capture the guarantees provided by contemporary high-performance architectures. It combines the expressiveness of strong consistency conditions (e.g., sequential consistency, linearizability) and the efficiency of weak consistency conditions (e.g., pipelined RAM, causal memory). Memory access operations are classified as either strong or weak. A global ordering of strong operations at different processes is guaranteed, but there is very little guarantee on the ordering of weak operations at different processes, except for what is implied by their interleaving with the strong operations. A formal and precise definition of this condition is given and an algorithm for providing hybrid consistency on distributed memory machines is presented. The response time of the algorithm is proved to be within a constant multiplicative factor of the (theoretical) optimal time bounds.
Hagit Attiya, Roy Friedman 0001
SIAM J. Comput.2
1997 Packing Messages as a Tool for Boosting the Performance of Total Ordering Protocols
abstract
This paper compares the throughput and latency of four protocols that provide total ordering. Two of these protocols are measured with and without message packing. We used a technique that buffers application messages for a short period of time before sending them, so more messages are packed together. The main conclusion of this comparison is that message packing influences the performance of total ordering protocols under high load overwhelmingly more than any other optimization that was checked in this paper, both in terms of throughput and latency. This improved performance is attributed to the fact that packing messages reduces the header overhead for messages, the contention on the network, and the load on the receiving CPUs.
Roy Friedman 0001, Robbert van Renesse
HPDC1
1997 The Hierarchical Daisy Architecture for Causal Delivery
abstract
We propose the hierarchical daisy architecture, which provides causal delivery of messages sent to any subset of processes. The architecture provides fault tolerance and maintains the amount of control information within a reasonable size. It divides processes into logical groups. Messages inside a logical group are sent directly, while messages that need to cross logical group boundaries are forwarded by servers. We prove the correctness of the daisy architecture and discuss possible optimizations.
Roberto Baldoni, Roy Friedman 0001, Robbert van Renesse
ICDCS2
1997 Failure Detectors in Omission Failure Environments
abstract
No abstract available.
Danny Dolev, Roy Friedman 0001, Idit Keidar, Dahlia Malkhi
PODC2
1997 Load Balancing Schemes for High-Throughput Distributed Fault-Tolerant Servers
abstract
Clusters of workstations, connected by a fast network, are emerging as a viable architecture for building high-throughput fault-tolerant servers. This type of architecture is more scalable and more cost-effective than a tightly coupled multiprocessor and may achieve as good a throughput. We explore several combinations of fault tolerance (FT) and load-balancing (LB) schemes, and compare their impact on the maximum throughput achievable by the system, and on its survivability. In particular, we show that the FT scheme has an effect on the throughput of the system, while the LB scheme affects the ability of the system to override failures. We study the scalability of the different schemes under different loads and failure conditions. Our simulations take into consideration the overhead of each scheme, the network contention, and the resource loads.
Roy Friedman 0001, Daniel Mossé
SRDS1
1997 Fast Replicated State Machines Over Partitionable Networks
abstract
The paper presents an implementation of replicated state machines in asynchronous distributed environments prone to node failures and network partitions. This implementation has several appealing properties: it guarantees that progress will be made whenever a majority of replicas can communicate with each other; it allows minority partitions to continue providing service for idempotent requests; it offers the application the choice between optimistic or safe message delivery. Performance measurements have shown that our implementation incurs low latency and achieves high throughput while providing globally consistent replicated state machine semantics.
Roy Friedman 0001, Alexey Vaysburd
SRDS1
1997 MILLIPEDE: Easy Parallel Programming in Available Distributed Environments
abstract
MILLIPEDE is a project aimed at developing a distributed shared memory environment for parallel programming. A major goal of this project is to support easy-to-grasp parallel programming languages that will also make it straightforward to parallelize existing code. Other targets are forward compatibility and availability of both the user programs (hence the shared memory support and the C-like parallel language PARC) and the system itself (which is thus implemented in user-level and using the operating system exported services). Locality of memory references, which implies efficiency and speedups, is maintained by MILLIPEDE} using page and thread migration, through which dynamic load-balancing and weak memory are implemented. ©1997 by John Wiley & Sons, Ltd.
Roy Friedman 0001, Maxim Goldin, Ayal Itzkovitz, Assaf Schuster
Softw. Pract. Exp.1
1996 Strong and Weak Virtual Synchrony in Horus
abstract
This paper presents two variants of virtual synchrony, which are supported by Horus. The first variant, called strong virtual synchrony, includes the property that every message is delivered within the view in which it is sent. This property is very useful in developing applications, since it helps in minimizing the amount of context information that needs to be sent on messages, and the amount of computation which is required in order to process a message. However, it is shown that in order to support this property, the application program has to block messages during view changes. An alternative definition, called weak virtual synchrony, which can be implemented without blocking messages, is then presented. This definition still guarantees that messages will be delivered within the view in which they were sent, only that it uses a slightly weaker notion of what the view in which a message was sent is. An implementation of weak virtual synchrony that does not block messages during view changes as also developed in this paper.
Roy Friedman 0001, Robbert van Renesse
SRDS1
1996 Limitations of Fast Consistency Conditions for Distributed Shared Memories
Hagit Attiya, Roy Friedman 0001
Inf. Process. Lett.2
1995 A Framework for Protocol Composition in Horus
abstract
The Horus system supports a communication architecture that treats protocols as instances of an abstract data type.This approach encourages developers to partition complex protocols into simple microprotocols, each of which is implemented by a protocol layer.Protocol layers can be stacked on top of each other in a variety of ways, at run-time.First, we describe the classes of protocols that can be supported this way.Next, we present the Horus object model that we designed for this technology, and the interface between the layers that makes it all work.We then present an example layer that implements a group membership protocol.Next, we show how, given a set of required properties, an appropriate stack can be constructed.We look at an example stack of protocols, which provides fault-tolerant, totally ordered communication between a group of processes.The work contributes a standard framework for protocol development and experimentation, provides a high performance implementation of the virtual synchrony model, and introduces a methodology for increasing the robustness of the protocol development process.
Robbert van Renesse, Kenneth P. Birman, Roy Friedman 0001, Mark Hayden, David A. Karr
PODC3
1995 Implementing Hybrid Consistency with High-Level Synchronization Operations
Roy Friedman 0001
Distributed Comput.1
1994 Programming DEC-Alpha Based Multiprocessors the Easy Way (Extended Abstract)
abstract
Article Free Access Share on Programming DEC-Alpha based multiprocessors the easy way (extended abstract) Authors: Hagit Attiya Department of Computer Science, The Technion, Haifa 32000, Israel Department of Computer Science, The Technion, Haifa 32000, IsraelView Profile , Roy Friedman Department of Computer Science, The Technion, Haifa 32000, Israel Department of Computer Science, The Technion, Haifa 32000, IsraelView Profile Authors Info & Claims SPAA '94: Proceedings of the sixth annual ACM symposium on Parallel algorithms and architecturesAugust 1994Pages 157–166https://doi.org/10.1145/181014.192323Published:01 August 1994Publication History 12citation49DownloadsMetricsTotal Citations12Total Downloads49Last 12 Months25Last 6 weeks1 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Publisher SiteeReaderPDF
Hagit Attiya, Roy Friedman 0001
SPAA2
1993 Implementing Hybrid Consistency with High-Level Synchronization Operations (Extenced Abstract)
abstract
In recent years, there is a growing tendency to support high-level synchronization operations, such as read-modifywrite, FIFO queues and stacks, as part of the programmer's shared memory model.This paper examines the problem of implementing hybrid consistency with high-level synchronization operations.Lower bounds on the time required to implement several common synchronization operations 12th ACM Symposfum
Roy Friedman 0001
PODC1
1993 Shared Memory Consistency Conditions for Non-Sequential Execution: Definitions and Programming Strategies
abstract
Article Free Access Share on Shared memory consistency conditions for non-sequential execution: definitions and programming strategies Authors: Hagit Attiya View Profile , Soma Chaudhuri View Profile , Roy Friedman View Profile , Jennifer L. Welch View Profile Authors Info & Claims SPAA '93: Proceedings of the fifth annual ACM symposium on Parallel Algorithms and ArchitecturesAugust 1993 Pages 241–250https://doi.org/10.1145/165231.165263Published:01 August 1993Publication History 10citation357DownloadsMetricsTotal Citations10Total Downloads357Last 12 Months7Last 6 weeks1 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteeReaderPDF
Hagit Attiya, Soma Chaudhuri, Roy Friedman 0001, Jennifer L. Welch
SPAA3
1992 A Correctness Condition for High-Performance Multiprocessors (Extended Abstract)
abstract
Hybrid consistency, a new consistency condition for shared memory multiprocessors, attempts to capture the guarantees provided by contemporary high-performance architectures. It combines the expressiveness of strong consistency conditions(e.g., sequential consistency, linearizability) and the efficiency of weak consistency conditions (e.g., Pipelined RAM, causal memory). Memory access operations are classified either strong or weak. A global ordering of strong operations at different processes is guaranteed, but there is very little guarantee on the ordering of weak operations at different processes, except for what is implied by their interleaving with the strong operations. A formal and precise definition of this condition is given. An efficient implementation of hybrid consistency on distributed memory machines is presented. In this implementation, weak opearations are executed instantaneously, while the response time for strong operations is linear in the network delay. (It is proven that this is within a constant factor of the optimal time bounds.)
Hagit Attiya, Roy Friedman 0001
STOC2