From e1532e262b7dae026de96f7d5089d247983967a8 Mon Sep 17 00:00:00 2001 From: tabokie Date: Mon, 19 Jun 2023 09:51:51 +0200 Subject: [PATCH] support switching disk Signed-off-by: tabokie --- src/engine.rs | 20 ++++--- src/file_pipe_log/pipe.rs | 40 +++++++++----- src/pipe_log.rs | 2 +- src/purge.rs | 8 +-- src/write_barrier.rs | 69 ++++++++++++++++++++++++ tests/failpoints/test_engine.rs | 94 +++++++++++++++++++++++++++++++++ 6 files changed, 210 insertions(+), 23 deletions(-) diff --git a/src/engine.rs b/src/engine.rs index 42e55821..294ff0cb 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -225,6 +225,14 @@ where Ok(len) } + pub fn switch_disk(&self) -> Result<()> { + if let Some(_token) = self.write_barrier.enter_as_admin() { + self.pipe_log.rotate(LogQueue::Append, true) + } else { + Err(Error::TryAgain(String::new())) + } + } + /// Synchronizes the Raft engine. pub fn sync(&self) -> Result<()> { self.write(&mut LogBatch::default(), true)?; @@ -2439,7 +2447,7 @@ pub(crate) mod tests { builder.begin(&mut log_batch); log_batch.put(rid, key.clone(), value.clone()).unwrap(); flush(&mut log_batch); - engine.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + engine.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); } { let mut builder = AtomicGroupBuilder::with_id(3); @@ -2458,7 +2466,7 @@ pub(crate) mod tests { log_batch.put(rid, key.clone(), value.clone()).unwrap(); data.insert(rid); flush(&mut log_batch); - engine.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + engine.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); } { let mut builder = AtomicGroupBuilder::with_id(3); @@ -2482,7 +2490,7 @@ pub(crate) mod tests { log_batch.put(rid, key.clone(), value.clone()).unwrap(); data.insert(rid); flush(&mut log_batch); - engine.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + engine.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); } { let mut builder = AtomicGroupBuilder::with_id(3); @@ -2501,7 +2509,7 @@ pub(crate) mod tests { log_batch.put(rid, key.clone(), value.clone()).unwrap(); data.insert(rid); flush(&mut log_batch); - engine.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + engine.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); } { // We must change id to avoid getting merged with last group. @@ -2522,7 +2530,7 @@ pub(crate) mod tests { rid += 1; log_batch.put(rid, key.clone(), value.clone()).unwrap(); flush(&mut log_batch); - engine.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + engine.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); } { let mut builder = AtomicGroupBuilder::with_id(5); @@ -2542,7 +2550,7 @@ pub(crate) mod tests { log_batch.put(rid, key.clone(), value.clone()).unwrap(); data.insert(rid); flush(&mut log_batch); - engine.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + engine.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); } engine.pipe_log.sync(LogQueue::Rewrite).unwrap(); diff --git a/src/file_pipe_log/pipe.rs b/src/file_pipe_log/pipe.rs index 5a5916ea..b8d41a8b 100644 --- a/src/file_pipe_log/pipe.rs +++ b/src/file_pipe_log/pipe.rs @@ -237,7 +237,11 @@ impl SinglePipe { /// Creates a new file for write, and rotates the active log file. /// /// This operation is atomic in face of errors. - fn rotate_imp(&self, writable_file: &mut MutexGuard>) -> Result<()> { + fn rotate_imp( + &self, + writable_file: &mut MutexGuard>, + switch_disk: bool, + ) -> Result<()> { let _t = StopWatch::new(( &*LOG_ROTATE_DURATION_HISTOGRAM, perf_context!(log_rotate_duration), @@ -247,9 +251,21 @@ impl SinglePipe { writable_file.writer.close()?; - let (path_id, handle) = self - .recycle_file(new_seq) - .unwrap_or_else(|| self.new_file(new_seq))?; + let (path_id, handle) = if switch_disk { + if self.paths.len() == 1 { + return Err(Error::InvalidArgument("no available path".to_owned())); + } + let file_id = FileId { + seq: new_seq, + queue: self.queue, + }; + let path_id = (self.active_files.read().back().unwrap().path_id + 1) % self.paths.len(); + let path = file_id.build_file_path(&self.paths[path_id]); + (path_id, self.file_system.create(path)?) + } else { + self.recycle_file(new_seq) + .unwrap_or_else(|| self.new_file(new_seq))? + }; let f = File:: { seq: new_seq, handle: handle.into(), @@ -318,7 +334,7 @@ impl SinglePipe { fail_point!("file_pipe_log::append"); let mut writable_file = self.writable_file.lock(); if writable_file.writer.offset() >= self.target_file_size { - if let Err(e) = self.rotate_imp(&mut writable_file) { + if let Err(e) = self.rotate_imp(&mut writable_file, false) { panic!( "error when rotate [{:?}:{}]: {e}", self.queue, writable_file.seq, @@ -369,7 +385,7 @@ impl SinglePipe { // - [3] Both main-dir and spill-dir have several recycled logs. // But as `bytes.len()` is always smaller than `target_file_size` in common // cases, this issue will be ignored temprorarily. - if let Err(e) = self.rotate_imp(&mut writable_file) { + if let Err(e) = self.rotate_imp(&mut writable_file, false) { panic!( "error when rotate [{:?}:{}]: {e}", self.queue, writable_file.seq @@ -422,8 +438,8 @@ impl SinglePipe { (last_seq - first_seq + 1) as usize * self.target_file_size } - fn rotate(&self) -> Result<()> { - self.rotate_imp(&mut self.writable_file.lock()) + fn rotate(&self, switch_disk: bool) -> Result<()> { + self.rotate_imp(&mut self.writable_file.lock(), switch_disk) } fn purge_to(&self, file_seq: FileSeq) -> Result { @@ -532,8 +548,8 @@ impl PipeLog for DualPipes { } #[inline] - fn rotate(&self, queue: LogQueue) -> Result<()> { - self.pipes[queue as usize].rotate() + fn rotate(&self, queue: LogQueue, switch_disk: bool) -> Result<()> { + self.pipes[queue as usize].rotate(switch_disk) } #[inline] @@ -645,7 +661,7 @@ mod tests { assert_eq!(file_handle.offset, header_size); assert_eq!(pipe_log.file_span(queue).1, 2); - pipe_log.rotate(queue).unwrap(); + pipe_log.rotate(queue, false).unwrap(); // purge file 1 assert_eq!(pipe_log.purge_to(FileId { queue, seq: 2 }).unwrap(), 1); @@ -715,7 +731,7 @@ mod tests { handles.push(pipe_log.append(&mut &content(i)).unwrap()); pipe_log.sync().unwrap(); } - pipe_log.rotate().unwrap(); + pipe_log.rotate(false).unwrap(); let (first, last) = pipe_log.file_span(); // Cannot purge already expired logs or not existsed logs. assert!(pipe_log.purge_to(first - 1).is_err()); diff --git a/src/pipe_log.rs b/src/pipe_log.rs index 57e94d1e..15629145 100644 --- a/src/pipe_log.rs +++ b/src/pipe_log.rs @@ -201,7 +201,7 @@ pub trait PipeLog: Sized { /// /// Implementation should be atomic under error conditions but not /// necessarily panic-safe. - fn rotate(&self, queue: LogQueue) -> Result<()>; + fn rotate(&self, queue: LogQueue, switch_disk: bool) -> Result<()>; /// Deletes all log files smaller than the specified file ID. The scope is /// limited to the log queue of `file_id`. diff --git a/src/purge.rs b/src/purge.rs index cb76f776..7c29df70 100644 --- a/src/purge.rs +++ b/src/purge.rs @@ -139,7 +139,7 @@ where let (_, last) = self.pipe_log.file_span(LogQueue::Append); let watermark = watermark.map_or(last, |w| std::cmp::min(w, last)); if watermark == last { - self.pipe_log.rotate(LogQueue::Append).unwrap(); + self.pipe_log.rotate(LogQueue::Append, false).unwrap(); } self.rewrite_append_queue_tombstones().unwrap(); if exit_after_step == Some(1) { @@ -169,13 +169,13 @@ where pub fn must_purge_all_stale(&self) { let _lk = self.force_rewrite_candidates.try_lock().unwrap(); - self.pipe_log.rotate(LogQueue::Rewrite).unwrap(); + self.pipe_log.rotate(LogQueue::Rewrite, false).unwrap(); self.rescan_memtables_and_purge_stale_files( LogQueue::Rewrite, self.pipe_log.file_span(LogQueue::Rewrite).1, ) .unwrap(); - self.pipe_log.rotate(LogQueue::Append).unwrap(); + self.pipe_log.rotate(LogQueue::Append, false).unwrap(); self.rescan_memtables_and_purge_stale_files( LogQueue::Append, self.pipe_log.file_span(LogQueue::Append).1, @@ -275,7 +275,7 @@ where // Rewrites the entire rewrite queue into new log files. fn rewrite_rewrite_queue(&self) -> Result> { let _t = StopWatch::new(&*ENGINE_REWRITE_REWRITE_DURATION_HISTOGRAM); - self.pipe_log.rotate(LogQueue::Rewrite)?; + self.pipe_log.rotate(LogQueue::Rewrite, false)?; let mut force_compact_regions = vec![]; let memtables = self.memtables.collect(|t| { diff --git a/src/write_barrier.rs b/src/write_barrier.rs index 3d365456..ec67a8b8 100644 --- a/src/write_barrier.rs +++ b/src/write_barrier.rs @@ -126,12 +126,30 @@ impl<'a, 'b, 'c, P, O> Iterator for WriterIter<'a, 'b, 'c, P, O> { } } +pub struct AdminToken<'a, P: 'a, O: 'a> { + ref_barrier: &'a WriteBarrier, +} + +impl<'a, P, O> Drop for AdminToken<'a, P, O> { + fn drop(&mut self) { + self.ref_barrier.leader_exit(); + } +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum AdminStatus { + None, + Pending, + Working, +} + struct WriteBarrierInner { head: Cell>>, tail: Cell>>, pending_leader: Cell>>, pending_index: Cell, + admin_status: Cell, } unsafe impl Send for WriteBarrierInner {} @@ -143,6 +161,7 @@ impl Default for WriteBarrierInner { tail: Cell::new(None), pending_leader: Cell::new(None), pending_index: Cell::new(0), + admin_status: Cell::new(AdminStatus::None), } } } @@ -151,6 +170,7 @@ impl Default for WriteBarrierInner { pub struct WriteBarrier { inner: Mutex>, leader_cv: Condvar, + admin_cv: Condvar, follower_cvs: [Condvar; 2], } @@ -158,6 +178,7 @@ impl Default for WriteBarrier { fn default() -> Self { WriteBarrier { leader_cv: Condvar::new(), + admin_cv: Condvar::new(), follower_cvs: [Condvar::new(), Condvar::new()], inner: Mutex::new(WriteBarrierInner::default()), } @@ -196,6 +217,16 @@ impl WriteBarrier { debug_assert!(inner.pending_leader.get().is_none()); inner.head.set(node); inner.tail.set(node); + // Unless there's an ongoing admin. + if inner.admin_status.get() == AdminStatus::Working { + inner.pending_leader.set(node); + inner + .pending_index + .set(inner.pending_index.get().wrapping_add(1)); + // + self.leader_cv.wait(&mut inner); + inner.pending_leader.set(None); + } } Some(WriteGroup { @@ -206,11 +237,48 @@ impl WriteBarrier { }) } + /// Waits until the caller has become the admin writer. Only admin writer + /// alone will be able to perform works. When necessary this call will be + /// reordered ahead of others already waiting in queue. + pub fn enter_as_admin(&self) -> Option> { + let mut inner = self.inner.lock(); + if inner.admin_status.get() != AdminStatus::None { + return None; + } + if inner.tail.get().is_some() { + // Jump ahead of the pending group (if any). + inner.admin_status.set(AdminStatus::Pending); + self.admin_cv.wait(&mut inner); + inner.admin_status.set(AdminStatus::Working); + } else { + // leader of a empty write group. proceed directly. + debug_assert!(inner.pending_leader.get().is_none()); + inner.admin_status.set(AdminStatus::Working); + } + Some(AdminToken { ref_barrier: self }) + } + /// Must called when write group leader finishes processing its responsible /// writers, and next write group should be formed. fn leader_exit(&self) { fail_point!("write_barrier::leader_exit", |_| {}); let inner = self.inner.lock(); + if inner.admin_status.get() == AdminStatus::Pending { + self.admin_cv.notify_one(); + // wake up follower of current write group. + if let Some(leader) = inner.pending_leader.get() { + self.follower_cvs[inner.pending_index.get().wrapping_sub(1) % 2].notify_all(); + inner.head.set(Some(leader)); + } else { + self.follower_cvs[inner.pending_index.get() % 2].notify_all(); + inner.head.set(None); + inner.tail.set(None); + } + return; + } + if inner.admin_status.get() == AdminStatus::Working { + inner.admin_status.set(AdminStatus::None); + } if let Some(leader) = inner.pending_leader.get() { // wake up leader of next write group. self.leader_cv.notify_one(); @@ -251,6 +319,7 @@ mod tests { } } assert_eq!(writer.finish(), 7); + let _token = barrier.enter_as_admin().unwrap(); } assert_eq!(processed_writers, 4); diff --git a/tests/failpoints/test_engine.rs b/tests/failpoints/test_engine.rs index 4a29e9d8..652448dd 100644 --- a/tests/failpoints/test_engine.rs +++ b/tests/failpoints/test_engine.rs @@ -1150,3 +1150,97 @@ fn test_build_engine_with_recycling_and_multi_dirs() { ); } } + +#[test] +fn test_concurrent_write_admin() { + let dir = tempfile::Builder::new() + .prefix("test_concurrent_write_admin") + .tempdir() + .unwrap(); + + let cfg = Config { + dir: dir.path().join("a").to_str().unwrap().to_owned(), + spill_dir: Some(dir.path().join("b").to_str().unwrap().to_owned()), + target_file_size: ReadableSize(1000000), + ..Default::default() + }; + + let some_entries = vec![ + Entry::new(), + Entry { + index: 1, + ..Default::default() + }, + ]; + + let engine = Arc::new(Engine::open(cfg).unwrap()); + let steps = Arc::new(AtomicU64::new(0)); + let mut threads = Vec::new(); + + let steps1 = steps.clone(); + fail::cfg_callback("write_barrier::leader_exit", move || { + while steps1.load(Ordering::Relaxed) != 4 { + std::thread::sleep(Duration::from_millis(50)); + } + }) + .unwrap(); + + // Leader + let engine1 = engine.clone(); + let some_entries1 = some_entries.clone(); + let steps1 = steps.clone(); + threads.push(std::thread::spawn(move || { + let mut log_batch = LogBatch::default(); + log_batch + .add_entries::(1, &some_entries1) + .unwrap(); + assert_eq!(steps1.fetch_add(1, Ordering::Relaxed), 0); + engine1.write(&mut log_batch, true).unwrap(); + })); + + while steps.load(Ordering::Relaxed) != 1 { + std::thread::sleep(Duration::from_millis(50)); + } + std::thread::sleep(Duration::from_millis(50)); + + // Follower. + let engine1 = engine.clone(); + let steps1 = steps.clone(); + threads.push(std::thread::spawn(move || { + let mut log_batch = LogBatch::default(); + log_batch + .add_entries::(2, &some_entries) + .unwrap(); + assert_eq!(steps1.fetch_add(1, Ordering::Relaxed), 1); + engine1.write(&mut log_batch, true).unwrap(); + assert_eq!(steps1.fetch_add(1, Ordering::Relaxed), 4); + })); + + while steps.load(Ordering::Relaxed) != 2 { + std::thread::sleep(Duration::from_millis(50)); + } + std::thread::sleep(Duration::from_millis(50)); + + // Admin + let engine1 = engine.clone(); + let steps1 = steps.clone(); + threads.push(std::thread::spawn(move || { + assert_eq!(steps1.fetch_add(1, Ordering::Relaxed), 2); + engine1.switch_disk().unwrap(); + })); + + while steps.load(Ordering::Relaxed) != 3 { + std::thread::sleep(Duration::from_millis(50)); + } + std::thread::sleep(Duration::from_millis(50)); + + // Unblock leader. + fail::remove("write_barrier::leader_exit"); + assert_eq!(steps.fetch_add(1, Ordering::Relaxed), 3); + + while steps.load(Ordering::Relaxed) != 5 { + std::thread::sleep(Duration::from_millis(50)); + } + let (a, b) = engine.file_span(LogQueue::Append); + assert_eq!(a + 1, b); +}