From 43690ab73006eb7b1d1a844b9f0ee4bc801b039b Mon Sep 17 00:00:00 2001 From: tabokie Date: Thu, 9 Dec 2021 23:07:35 +0800 Subject: [PATCH 1/5] fix stress default values Signed-off-by: tabokie --- stress/src/main.rs | 51 +++++++++++++++++----------------------------- 1 file changed, 19 insertions(+), 32 deletions(-) diff --git a/stress/src/main.rs b/stress/src/main.rs index b3023414..44382e26 100644 --- a/stress/src/main.rs +++ b/stress/src/main.rs @@ -35,9 +35,8 @@ impl MessageExt for MessageExtTyped { const DEFAULT_TIME: Duration = Duration::from_secs(60); const DEFAULT_REGIONS: u64 = 5; const DEFAULT_PURGE_INTERVAL: Duration = Duration::from_millis(10 * 1000); -const DEFAULT_COMPACT_TTL: Duration = Duration::from_millis(0); const DEFAULT_COMPACT_COUNT: u64 = 0; -const DEFAULT_FORCE_COMPACT_FACTOR_STR: &str = "0.5"; +const DEFAULT_FORCE_COMPACT_FACTOR: f64 = 0.5; const DEFAULT_WRITE_THREADS: u64 = 1; const DEFAULT_WRITE_OPS_PER_THREAD: u64 = 0; const DEFAULT_READ_THREADS: u64 = 0; @@ -106,7 +105,7 @@ struct ControlOpt { takes_value = true, default_value = "0.5", validator = |s| { - let factor = s.parse::().unwrap(); + let factor = s.parse::().unwrap(); if factor >= 1.0 { Err(String::from("Factor must be smaller than 1.0")) } else if factor <= 0.0 { @@ -117,7 +116,7 @@ struct ControlOpt { }, help = "Factor to shrink raft log during force compact" )] - force_compact_factor: f32, + force_compact_factor: f64, #[clap( long = "write-threads", @@ -131,13 +130,13 @@ struct ControlOpt { long = "write-ops-per-thread", value_name = "ops", takes_value = true, - default_value = formatcp!("{}", DEFAULT_READ_OPS_PER_THREAD), + default_value = formatcp!("{}", DEFAULT_WRITE_OPS_PER_THREAD), help = "Set the per-thread OPS for read entry requests" )] write_ops_per_thread: u64, #[clap( - long = "read-thread", + long = "read-threads", value_name = "threads", takes_value = true, default_value = formatcp!("{}", DEFAULT_READ_THREADS), @@ -181,20 +180,10 @@ struct ControlOpt { )] write_region_count: u64, - #[clap( - long = "write-async", - value_name = "async", - takes_value = true, - help = "Whether to async write raft logs" - )] - write_async: bool, + #[clap(long = "write-without-sync", help = "Do not sync after write")] + write_without_sync: bool, - #[clap( - long = "reuse-data", - value_name = "reuse", - takes_value = true, - help = "Whether to reuse existing data in specified path" - )] + #[clap(long = "reuse-data", help = "Reuse existing data in specified path")] reuse_data: bool, #[clap( @@ -231,13 +220,13 @@ struct ControlOpt { default_value = "0.6", help = "Purge if rewrite log files garbage ratio is greater than this threshold" )] - purge_rewrite_garbage_ratio: String, + purge_rewrite_garbage_ratio: f64, #[clap( long = "compression-threshold", value_name = "size", takes_value = true, - default_value = "10GB", + default_value = "8KB", help = "Compress log batch bigger than this threshold" )] batch_compression_threshold: String, @@ -248,9 +237,8 @@ struct TestArgs { time: Duration, regions: u64, purge_interval: Duration, - compact_ttl: Duration, compact_count: u64, - force_compact_factor: f32, + force_compact_factor: f64, write_threads: u64, write_ops_per_thread: u64, read_threads: u64, @@ -258,7 +246,7 @@ struct TestArgs { entry_size: usize, write_entry_count: u64, write_region_count: u64, - write_async: bool, + write_without_sync: bool, } impl Default for TestArgs { @@ -267,9 +255,8 @@ impl Default for TestArgs { time: DEFAULT_TIME, regions: DEFAULT_REGIONS, purge_interval: DEFAULT_PURGE_INTERVAL, - compact_ttl: DEFAULT_COMPACT_TTL, compact_count: DEFAULT_COMPACT_COUNT, - force_compact_factor: DEFAULT_FORCE_COMPACT_FACTOR_STR.parse::().unwrap(), + force_compact_factor: DEFAULT_FORCE_COMPACT_FACTOR, write_threads: DEFAULT_WRITE_THREADS, write_ops_per_thread: DEFAULT_WRITE_OPS_PER_THREAD, read_threads: DEFAULT_READ_THREADS, @@ -277,7 +264,7 @@ impl Default for TestArgs { entry_size: DEFAULT_ENTRY_SIZE, write_entry_count: DEFAULT_WRITE_ENTRY_COUNT, write_region_count: DEFAULT_WRITE_REGION_COUNT, - write_async: DEFAULT_WRITE_SYNC, + write_without_sync: DEFAULT_WRITE_SYNC, } } } @@ -441,7 +428,7 @@ fn spawn_write( // TODO(tabokie): compensate for slow requests wait_til(&mut start, last + i); } - if let Err(e) = engine.write(&mut log_batch, !args.write_async) { + if let Err(e) = engine.write(&mut log_batch, !args.write_without_sync) { println!("write error {:?} in thread {}", e, index); } let end = Instant::now(); @@ -508,7 +495,7 @@ fn spawn_purge( let first = engine.first_index(region).unwrap_or(0); let last = engine.last_index(region).unwrap_or(0); let compact_to = last - - ((last - first + 1) as f32 * args.force_compact_factor) as u64 + - ((last - first + 1) as f64 * args.force_compact_factor) as u64 + 1; engine.compact_to(region, compact_to); } @@ -585,12 +572,12 @@ fn main() { config.purge_threshold = ReadableSize::from_str(&opts.purge_threshold).unwrap(); config.purge_rewrite_threshold = Some(ReadableSize::from_str(&opts.purge_rewrite_threshold).unwrap()); - config.purge_rewrite_garbage_ratio = opts.purge_rewrite_garbage_ratio.parse::().unwrap(); + config.purge_rewrite_garbage_ratio = opts.purge_rewrite_garbage_ratio; config.batch_compression_threshold = ReadableSize::from_str(&opts.batch_compression_threshold).unwrap(); args.time = Duration::from_secs(opts.time); args.regions = opts.regions; - args.purge_interval = Duration::from_millis(opts.purge_interval); + args.purge_interval = Duration::from_secs(opts.purge_interval); if let Some(count) = opts.compact_count { args.compact_count = count; @@ -604,7 +591,7 @@ fn main() { args.entry_size = opts.entry_size; args.write_entry_count = opts.write_entry_count; args.write_region_count = opts.write_region_count; - args.write_async = opts.write_async; + args.write_without_sync = opts.write_without_sync; if !opts.reuse_data { // clean up existing log files let _ = std::fs::remove_dir_all(&config.dir); From 98a46d2a67afab36da4765e380fe7927e873302f Mon Sep 17 00:00:00 2001 From: tabokie Date: Thu, 9 Dec 2021 23:15:20 +0800 Subject: [PATCH 2/5] bump rust-toolchain Signed-off-by: tabokie --- rust-toolchain | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rust-toolchain b/rust-toolchain index d86522f4..e996cc42 100644 --- a/rust-toolchain +++ b/rust-toolchain @@ -1 +1 @@ -nightly-2021-07-28 \ No newline at end of file +nightly-2021-12-07 \ No newline at end of file From 14517df9ae6114587bb43cdeb4ac83da5d8b4640 Mon Sep 17 00:00:00 2001 From: tabokie Date: Thu, 9 Dec 2021 23:15:36 +0800 Subject: [PATCH 3/5] use vecdeque for storing indexes Signed-off-by: tabokie --- src/consistency.rs | 4 ++-- src/log_batch.rs | 17 +++++++++-------- src/memtable.rs | 15 +++++++-------- src/purge.rs | 8 ++++++-- 4 files changed, 24 insertions(+), 20 deletions(-) diff --git a/src/consistency.rs b/src/consistency.rs index 6d28a701..8317b4af 100644 --- a/src/consistency.rs +++ b/src/consistency.rs @@ -28,8 +28,8 @@ impl ReplayMachine for ConsistencyChecker { for item in item_batch.iter() { if let LogItemContent::EntryIndexes(ents) = &item.content { if !ents.0.is_empty() { - let incoming_first_index = ents.0.first().unwrap().index; - let incoming_last_index = ents.0.last().unwrap().index; + let incoming_first_index = ents.0.front().unwrap().index; + let incoming_last_index = ents.0.back().unwrap().index; let index_range = self .raft_groups .entry(item.raft_group_id) diff --git a/src/log_batch.rs b/src/log_batch.rs index 5e8d1c84..3fdb9aa4 100644 --- a/src/log_batch.rs +++ b/src/log_batch.rs @@ -8,6 +8,7 @@ use byteorder::{BigEndian, ReadBytesExt, WriteBytesExt}; use log::error; use protobuf::parse_from_bytes; use protobuf::Message; +use std::collections::VecDeque; use crate::codec::{self, NumberEncoder}; use crate::memtable::EntryIndex; @@ -63,12 +64,12 @@ type SliceReader<'a> = &'a [u8]; // Format: // { count | first index | [ tail offsets ] } #[derive(Clone, Debug, PartialEq)] -pub struct EntryIndexes(pub Vec); +pub struct EntryIndexes(pub VecDeque); impl EntryIndexes { pub fn decode(buf: &mut SliceReader, entries_size: &mut usize) -> Result { let mut count = codec::decode_var_u64(buf)?; - let mut entry_indexes = Vec::with_capacity(count as usize); + let mut entry_indexes = VecDeque::with_capacity(count as usize); let mut index = 0; if count > 0 { index = codec::decode_var_u64(buf)?; @@ -83,7 +84,7 @@ impl EntryIndexes { ..Default::default() }; *entries_size += entry_len; - entry_indexes.push(entry_index); + entry_indexes.push_back(entry_index); index += 1; count -= 1; } @@ -247,7 +248,7 @@ pub enum LogItemContent { } impl LogItem { - pub fn new_entry_indexes(raft_group_id: u64, entry_indexes: Vec) -> LogItem { + pub fn new_entry_indexes(raft_group_id: u64, entry_indexes: VecDeque) -> LogItem { LogItem { raft_group_id, content: LogItemContent::EntryIndexes(EntryIndexes(entry_indexes)), @@ -423,7 +424,7 @@ impl LogItemBatch { } } - pub fn add_entry_indexes(&mut self, region_id: u64, mut entry_indexes: Vec) { + pub fn add_entry_indexes(&mut self, region_id: u64, mut entry_indexes: VecDeque) { for ei in entry_indexes.iter_mut() { ei.entry_offset = self.entries_size as u64; self.entries_size += ei.entry_len; @@ -574,12 +575,12 @@ impl LogBatch { ) -> Result<()> { debug_assert!(self.buf_state == BufState::Open); - let mut entry_indexes = Vec::with_capacity(entries.len()); + let mut entry_indexes = VecDeque::with_capacity(entries.len()); self.buf_state = BufState::Incomplete; for e in entries { let buf_offset = self.buf.len(); e.write_to_vec(&mut self.buf)?; - entry_indexes.push(EntryIndex { + entry_indexes.push_back(EntryIndex { index: M::index(e), entry_len: self.buf.len() - buf_offset, ..Default::default() @@ -593,7 +594,7 @@ impl LogBatch { pub(crate) fn add_raw_entries( &mut self, region_id: u64, - mut entry_indexes: Vec, + mut entry_indexes: VecDeque, entries: Vec>, ) -> Result<()> { debug_assert!(entry_indexes.len() == entries.len()); diff --git a/src/memtable.rs b/src/memtable.rs index 8fba093b..70d6a6f5 100644 --- a/src/memtable.rs +++ b/src/memtable.rs @@ -186,29 +186,28 @@ impl MemTable { } } - pub fn append(&mut self, entry_indexes: Vec) { + pub fn append(&mut self, mut entry_indexes: VecDeque) { let len = entry_indexes.len(); if len > 0 { self.prepare_append(entry_indexes[0].index, false, false); self.global_stats.add(LogQueue::Append, len); - // TODO: Optimize this. - self.entry_indexes.extend(entry_indexes); + self.entry_indexes.append(&mut entry_indexes); } } // This will only be called during recovery. - pub fn append_rewrite(&mut self, entry_indexes: Vec) { + pub fn append_rewrite(&mut self, mut entry_indexes: VecDeque) { let len = entry_indexes.len(); if len > 0 { debug_assert_eq!(self.rewrite_count, self.entry_indexes.len()); self.prepare_append(entry_indexes[0].index, true, true); self.global_stats.add(LogQueue::Rewrite, len); - self.entry_indexes.extend(entry_indexes); + self.entry_indexes.append(&mut entry_indexes); self.rewrite_count = self.entry_indexes.len(); } } - pub fn rewrite(&mut self, rewrite_indexes: Vec, gate: Option) { + pub fn rewrite(&mut self, rewrite_indexes: VecDeque, gate: Option) { if rewrite_indexes.is_empty() { return; } @@ -240,8 +239,8 @@ impl MemTable { ); let rewrite_pos = (rewrite_first - rewrite_indexes[0].index) as usize; - for (i, rindex) in rewrite_indexes[rewrite_pos..rewrite_pos + rewrite_len] - .iter() + for (i, rindex) in rewrite_indexes + .range(rewrite_pos..rewrite_pos + rewrite_len) .enumerate() { let index = &mut self.entry_indexes[i + pos]; diff --git a/src/purge.rs b/src/purge.rs index 86b5a0e8..43591457 100644 --- a/src/purge.rs +++ b/src/purge.rs @@ -291,7 +291,11 @@ where std::mem::swap(&mut take_entries, &mut entries); let mut take_entry_indexes = entry_indexes.split_off(cursor + 1); std::mem::swap(&mut take_entry_indexes, &mut entry_indexes); - log_batch.add_raw_entries(region_id, take_entry_indexes, take_entries)?; + log_batch.add_raw_entries( + region_id, + take_entry_indexes.into(), + take_entries, + )?; self.rewrite_impl(&mut log_batch, rewrite, false)?; total_size = 0; cursor = 0; @@ -300,7 +304,7 @@ where } } if !entries.is_empty() { - log_batch.add_raw_entries(region_id, entry_indexes, entries)?; + log_batch.add_raw_entries(region_id, entry_indexes.into(), entries)?; } for (k, v) in kvs { log_batch.put(region_id, k, v); From d2927a394597df99c5ce8231527f367f4b104532 Mon Sep 17 00:00:00 2001 From: tabokie Date: Thu, 9 Dec 2021 23:58:30 +0800 Subject: [PATCH 4/5] fix issues Signed-off-by: tabokie --- src/log_batch.rs | 12 ++++++------ src/memtable.rs | 17 ++++++++++++++++- src/test_util.rs | 11 ++++++++--- 3 files changed, 30 insertions(+), 10 deletions(-) diff --git a/src/log_batch.rs b/src/log_batch.rs index 3fdb9aa4..9c2f768b 100644 --- a/src/log_batch.rs +++ b/src/log_batch.rs @@ -475,7 +475,7 @@ impl LogItemBatch { compression_type: CompressionType, ) -> Result { verify_checksum(buf)?; - *buf = &mut &buf[..buf.len() - LOG_BATCH_CHECKSUM_LEN]; + *buf = &mut buf[..buf.len() - LOG_BATCH_CHECKSUM_LEN]; let count = codec::decode_var_u64(buf)?; let mut items = LogItemBatch::with_capacity(count as usize); let mut entries_size = 0; @@ -837,7 +837,7 @@ mod tests { fn decode_entries_from_bytes( buf: &[u8], - entry_indexes: &[EntryIndex], + entry_indexes: &VecDeque, _encoded: bool, ) -> Vec { let mut entries = Vec::with_capacity(entry_indexes.len()); @@ -849,7 +849,7 @@ mod tests { #[test] fn test_entry_indexes_enc_dec() { - fn encode_and_decode(entry_indexes: &mut Vec) -> EntryIndexes { + fn encode_and_decode(entry_indexes: &mut VecDeque) -> EntryIndexes { let mut entries_size = 0; for idx in entry_indexes.iter_mut() { idx.entry_offset = entries_size as u64; @@ -869,7 +869,7 @@ mod tests { decoded_indexes } - let entry_indexes = vec![Vec::new(), generate_entry_indexes_opt(7, 17, None)]; + let entry_indexes = vec![VecDeque::new(), generate_entry_indexes_opt(7, 17, None)]; for mut idxs in entry_indexes.into_iter() { let decoded = encode_and_decode(&mut idxs); assert_eq!(idxs, decoded.0); @@ -1063,8 +1063,8 @@ mod tests { if let LogItemContent::EntryIndexes(entry_indexes) = &item.content { if !entry_indexes.0.is_empty() { let (begin, end) = ( - entry_indexes.0.first().unwrap().index, - entry_indexes.0.last().unwrap().index + 1, + entry_indexes.0.front().unwrap().index, + entry_indexes.0.back().unwrap().index + 1, ); let origin_entries = generate_entries(begin, end, Some(entry_data)); let decoded_entries = diff --git a/src/memtable.rs b/src/memtable.rs index 70d6a6f5..7de9d504 100644 --- a/src/memtable.rs +++ b/src/memtable.rs @@ -835,7 +835,7 @@ mod tests { memtable.consistency_check(); // Empty. - memtable.append(Vec::new()); + memtable.append(VecDeque::new()); // Hole. assert!( @@ -1730,4 +1730,19 @@ mod tests { memtable.put(key2.clone(), value.clone(), FileId::dummy(LogQueue::Append)); }); } + + #[bench] + fn bench_memtable_append(b: &mut test::Bencher) { + let mut memtable = MemTable::new(0, Arc::new(GlobalStats::default())); + b.iter(move || { + let last_index = memtable.last_index().unwrap_or(0); + let entry_indexes = generate_entry_indexes( + last_index + 1, + last_index + 11, + FileId::dummy(LogQueue::Append), + ); + memtable.append(entry_indexes); + memtable.compact_to(last_index + 1); + }); + } } diff --git a/src/test_util.rs b/src/test_util.rs index c5310632..dd800bc9 100644 --- a/src/test_util.rs +++ b/src/test_util.rs @@ -1,5 +1,6 @@ // Copyright (c) 2017-present, PingCAP, Inc. Licensed under Apache-2.0. +use std::collections::VecDeque; use std::panic::{self, AssertUnwindSafe}; use raft::eraftpb::Entry; @@ -22,7 +23,11 @@ pub fn generate_entries(begin_index: u64, end_index: u64, data: Option<&[u8]>) - v } -pub fn generate_entry_indexes(begin_idx: u64, end_idx: u64, file_id: FileId) -> Vec { +pub fn generate_entry_indexes( + begin_idx: u64, + end_idx: u64, + file_id: FileId, +) -> VecDeque { generate_entry_indexes_opt(begin_idx, end_idx, Some(file_id)) } @@ -30,7 +35,7 @@ pub fn generate_entry_indexes_opt( begin_idx: u64, end_idx: u64, file_id: Option, -) -> Vec { +) -> VecDeque { assert!(end_idx >= begin_idx); let mut ents_idx = vec![]; for idx in begin_idx..end_idx { @@ -43,7 +48,7 @@ pub fn generate_entry_indexes_opt( ents_idx.push(ent_idx); } - ents_idx + ents_idx.into() } /// Catch panic while suppressing default panic hook. From efe23f7473f038d2146593f42949ae02bfc46ac0 Mon Sep 17 00:00:00 2001 From: tabokie Date: Fri, 10 Dec 2021 14:11:55 +0800 Subject: [PATCH 5/5] fix clippy issues Signed-off-by: tabokie --- src/lib.rs | 1 - src/log_batch.rs | 2 +- src/write_barrier.rs | 81 ++++++++++++++++++++++++++++++++++---------- 3 files changed, 65 insertions(+), 19 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index 33e9a773..cc6f8844 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -11,7 +11,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -#![feature(shrink_to)] #![feature(map_first_last)] #![feature(generic_associated_types)] #![feature(test)] diff --git a/src/log_batch.rs b/src/log_batch.rs index 9c2f768b..86502a49 100644 --- a/src/log_batch.rs +++ b/src/log_batch.rs @@ -475,7 +475,7 @@ impl LogItemBatch { compression_type: CompressionType, ) -> Result { verify_checksum(buf)?; - *buf = &mut buf[..buf.len() - LOG_BATCH_CHECKSUM_LEN]; + *buf = &buf[..buf.len() - LOG_BATCH_CHECKSUM_LEN]; let count = codec::decode_var_u64(buf)?; let mut items = LogItemBatch::with_capacity(count as usize); let mut entries_size = 0; diff --git a/src/write_barrier.rs b/src/write_barrier.rs index 3eb33957..742d2457 100644 --- a/src/write_barrier.rs +++ b/src/write_barrier.rs @@ -8,7 +8,52 @@ use std::time::Instant; use fail::fail_point; use parking_lot::{Condvar, Mutex}; -type Ptr = Option>; +struct Ptr(Option>); + +unsafe impl Send for Ptr {} + +impl std::clone::Clone for Ptr { + #[inline] + fn clone(&self) -> Self { + Ptr(self.0) + } +} + +impl std::marker::Copy for Ptr {} + +impl PartialEq for Ptr { + #[inline] + fn eq(&self, other: &Self) -> bool { + self.0 == other.0 + } +} + +impl Ptr { + #[inline] + fn null() -> Self { + Ptr(None) + } + + #[inline] + fn from_mut(t: &mut T) -> Self { + unsafe { Ptr(Some(NonNull::new_unchecked(t))) } + } + + #[inline] + fn as_mut<'a>(&self) -> &'a mut T { + unsafe { self.0.unwrap().as_mut() } + } + + #[inline] + fn is_null(&self) -> bool { + self.0.is_none() + } + + #[inline] + fn set_null(&mut self) { + self.0 = None; + } +} pub struct Writer { next: Cell>>, @@ -19,11 +64,13 @@ pub struct Writer { pub(crate) start_time: Instant, } +unsafe impl Send for Writer {} + impl Writer { // SAFETY: Data pointed by `payload` is owned by this writer during its lifetime. pub fn new(payload: &mut P, sync: bool, start_time: Instant) -> Self { Writer { - next: Cell::new(None), + next: Cell::new(Ptr::null()), payload: payload as *mut _, output: None, sync, @@ -91,12 +138,12 @@ impl<'a, 'b, 'c, P, O> Iterator for WriterIter<'a, 'b, 'c, P, O> { type Item = &'a mut Writer; fn next(&mut self) -> Option { - if self.start.is_none() { + if self.start.is_null() { None } else { - let writer = unsafe { self.start.unwrap().as_mut() }; + let writer = self.start.as_mut(); if self.start == self.back { - self.start = None; + self.start.set_null(); } else { self.start = writer.get_next(); } @@ -118,9 +165,9 @@ unsafe impl Send for WriteBarrierInner {} impl Default for WriteBarrierInner { fn default() -> Self { WriteBarrierInner { - head: Cell::new(None), - tail: Cell::new(None), - pending_leader: Cell::new(None), + head: Cell::new(Ptr::null()), + tail: Cell::new(Ptr::null()), + pending_leader: Cell::new(Ptr::null()), pending_index: Cell::new(0), } } @@ -147,15 +194,15 @@ impl WriteBarrier { /// become the leader of a set of writers, returns a `WriteGroup` that /// contains them, `writer` included. pub fn enter<'a>(&self, writer: &'a mut Writer) -> Option> { - let node = unsafe { Some(NonNull::new_unchecked(writer)) }; + let node = Ptr::from_mut(writer); let mut inner = self.inner.lock(); - if let Some(tail) = inner.tail.get() { + if let Ptr(Some(tail)) = inner.tail.get() { unsafe { tail.as_ref().set_next(node); } inner.tail.set(node); - if inner.pending_leader.get().is_some() { + if inner.pending_leader.get().is_null() { // follower of next write group. self.follower_cvs[inner.pending_index.get() % 2].wait(&mut inner); return None; @@ -167,11 +214,11 @@ impl WriteBarrier { .set(inner.pending_index.get().wrapping_add(1)); // self.leader_cv.wait(&mut inner); - inner.pending_leader.set(None); + inner.pending_leader.set(Ptr::null()); } } else { // leader of a empty write group. proceed directly. - debug_assert!(inner.pending_leader.get().is_none()); + debug_assert!(inner.pending_leader.get().is_null()); inner.head.set(node); inner.tail.set(node); } @@ -189,17 +236,17 @@ impl WriteBarrier { fn leader_exit(&self) { fail_point!("write_barrier::leader_exit", |_| {}); let inner = self.inner.lock(); - if let Some(leader) = inner.pending_leader.get() { + if let Ptr(Some(leader)) = inner.pending_leader.get() { // wake up leader of next write group. self.leader_cv.notify_one(); // wake up follower of current write group. self.follower_cvs[inner.pending_index.get().wrapping_sub(1) % 2].notify_all(); - inner.head.set(Some(leader)); + inner.head.set(Ptr(Some(leader))); } else { // wake up follower of current write group. self.follower_cvs[inner.pending_index.get() % 2].notify_all(); - inner.head.set(None); - inner.tail.set(None); + inner.head.set(Ptr::null()); + inner.tail.set(Ptr::null()); } } }