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 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/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 5e8d1c84..86502a49 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; @@ -474,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; @@ -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()); @@ -836,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()); @@ -848,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; @@ -868,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); @@ -1062,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 0fec169b..7474b8fe 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]; @@ -844,7 +843,7 @@ mod tests { memtable.consistency_check(); // Empty. - memtable.append(Vec::new()); + memtable.append(VecDeque::new()); // Hole. assert!( @@ -1739,4 +1738,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/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); 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. 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()); } } }