Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 14 additions & 6 deletions src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -225,6 +225,14 @@
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()))

Check warning on line 232 in src/engine.rs

View check run for this annotation

Codecov / codecov/patch

src/engine.rs#L232

Added line #L232 was not covered by tests
}
}

/// Synchronizes the Raft engine.
pub fn sync(&self) -> Result<()> {
self.write(&mut LogBatch::default(), true)?;
Expand Down Expand Up @@ -2439,7 +2447,7 @@
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);
Expand All @@ -2458,7 +2466,7 @@
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);
Expand All @@ -2482,7 +2490,7 @@
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);
Expand All @@ -2501,7 +2509,7 @@
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.
Expand All @@ -2522,7 +2530,7 @@
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);
Expand All @@ -2542,7 +2550,7 @@
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();

Expand Down
40 changes: 28 additions & 12 deletions src/file_pipe_log/pipe.rs
Original file line number Diff line number Diff line change
Expand Up @@ -237,7 +237,11 @@
/// 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<WritableFile<F>>) -> Result<()> {
fn rotate_imp(
&self,
writable_file: &mut MutexGuard<WritableFile<F>>,
switch_disk: bool,
) -> Result<()> {
let _t = StopWatch::new((
&*LOG_ROTATE_DURATION_HISTOGRAM,
perf_context!(log_rotate_duration),
Expand All @@ -247,9 +251,21 @@

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()));

Check warning on line 256 in src/file_pipe_log/pipe.rs

View check run for this annotation

Codecov / codecov/patch

src/file_pipe_log/pipe.rs#L256

Added line #L256 was not covered by tests
}
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::<F> {
seq: new_seq,
handle: handle.into(),
Expand Down Expand Up @@ -318,7 +334,7 @@
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,
Expand Down Expand Up @@ -369,7 +385,7 @@
// - [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
Expand Down Expand Up @@ -422,8 +438,8 @@
(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<usize> {
Expand Down Expand Up @@ -532,8 +548,8 @@
}

#[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]
Expand Down Expand Up @@ -645,7 +661,7 @@
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);
Expand Down Expand Up @@ -715,7 +731,7 @@
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());
Expand Down
2 changes: 1 addition & 1 deletion src/pipe_log.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand Down
8 changes: 4 additions & 4 deletions src/purge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -275,7 +275,7 @@ where
// Rewrites the entire rewrite queue into new log files.
fn rewrite_rewrite_queue(&self) -> Result<Vec<u64>> {
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| {
Expand Down
69 changes: 69 additions & 0 deletions src/write_barrier.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,12 +126,30 @@
}
}

pub struct AdminToken<'a, P: 'a, O: 'a> {
ref_barrier: &'a WriteBarrier<P, O>,
}

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<P, O> {
head: Cell<Ptr<Writer<P, O>>>,
tail: Cell<Ptr<Writer<P, O>>>,

pending_leader: Cell<Ptr<Writer<P, O>>>,
pending_index: Cell<usize>,
admin_status: Cell<AdminStatus>,
}

unsafe impl<P: Send, O: Send> Send for WriteBarrierInner<P, O> {}
Expand All @@ -143,6 +161,7 @@
tail: Cell::new(None),
pending_leader: Cell::new(None),
pending_index: Cell::new(0),
admin_status: Cell::new(AdminStatus::None),
}
}
}
Expand All @@ -151,13 +170,15 @@
pub struct WriteBarrier<P, O> {
inner: Mutex<WriteBarrierInner<P, O>>,
leader_cv: Condvar,
admin_cv: Condvar,
follower_cvs: [Condvar; 2],
}

impl<P, O> Default for WriteBarrier<P, O> {
fn default() -> Self {
WriteBarrier {
leader_cv: Condvar::new(),
admin_cv: Condvar::new(),
follower_cvs: [Condvar::new(), Condvar::new()],
inner: Mutex::new(WriteBarrierInner::default()),
}
Expand Down Expand Up @@ -196,6 +217,16 @@
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);

Check warning on line 228 in src/write_barrier.rs

View check run for this annotation

Codecov / codecov/patch

src/write_barrier.rs#L222-L228

Added lines #L222 - L228 were not covered by tests
}
}

Some(WriteGroup {
Expand All @@ -206,11 +237,48 @@
})
}

/// 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<AdminToken<P, O>> {
let mut inner = self.inner.lock();
if inner.admin_status.get() != AdminStatus::None {
return None;

Check warning on line 246 in src/write_barrier.rs

View check run for this annotation

Codecov / codecov/patch

src/write_barrier.rs#L246

Added line #L246 was not covered by tests
}
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);
}

Check warning on line 276 in src/write_barrier.rs

View check run for this annotation

Codecov / codecov/patch

src/write_barrier.rs#L273-L276

Added lines #L273 - L276 were not covered by tests
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();
Expand Down Expand Up @@ -251,6 +319,7 @@
}
}
assert_eq!(writer.finish(), 7);
let _token = barrier.enter_as_admin().unwrap();
}

assert_eq!(processed_writers, 4);
Expand Down
Loading