Skip to content

Commit 8937f5b

Browse files
committed
fix #296
1 parent 41bc213 commit 8937f5b

7 files changed

Lines changed: 35 additions & 29 deletions

File tree

src/batch/mod.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -105,7 +105,7 @@ impl WriteBatch {
105105
}
106106

107107
log::trace!("batch: Acquiring journal writer");
108-
let mut journal_writer = self.db.supervisor.journal.get_writer();
108+
let mut journal_writer = self.db.supervisor.journal.get_writer()?;
109109

110110
// IMPORTANT: Check the poisoned flag after getting journal mutex, otherwise TOCTOU
111111
if self.db.is_poisoned.load(Ordering::Relaxed) {

src/db.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -288,7 +288,7 @@ impl Database {
288288
/// Returns the disk space usage of the journal.
289289
#[doc(hidden)]
290290
pub fn journal_disk_space(&self) -> crate::Result<u64> {
291-
Ok(self.supervisor.journal.get_writer().len()?
291+
Ok(self.supervisor.journal.get_writer()?.len()?
292292
+ self
293293
.supervisor
294294
.journal_manager
@@ -587,7 +587,7 @@ impl Database {
587587
log::debug!("journal recovery result: {journal_recovery:#?}");
588588

589589
let active_journal = Arc::new(journal_recovery.active);
590-
active_journal.get_writer().persist(PersistMode::SyncAll)?;
590+
active_journal.get_writer()?.persist(PersistMode::SyncAll)?;
591591

592592
let sealed_journals = journal_recovery.sealed;
593593

src/db_test.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ fn clear_recover_sealed() -> crate::Result<()> {
2121

2222
tree.rotate_memtable_and_wait()?;
2323
assert!(tree.is_empty()?);
24-
db.supervisor.journal.get_writer().rotate()?;
24+
db.supervisor.journal.get_writer()?.rotate()?;
2525

2626
tree.insert("b", "a")?;
2727
assert!(tree.contains_key("b")?);

src/journal/mod.rs

Lines changed: 13 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,13 @@ pub struct Journal {
3131

3232
impl std::fmt::Debug for Journal {
3333
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
34-
write!(f, "{}", self.path().display())
34+
write!(
35+
f,
36+
"{}",
37+
self.path()
38+
.map(|p| p.display().to_string())
39+
.unwrap_or_else(|_| String::from("<failed to read path>"))
40+
)
3541
}
3642
}
3743

@@ -95,23 +101,23 @@ impl Journal {
95101
}
96102

97103
/// Hands out write access for the journal.
98-
pub(crate) fn get_writer(&self) -> MutexGuard<'_, Writer> {
104+
pub(crate) fn get_writer(&self) -> crate::Result<MutexGuard<'_, Writer>> {
99105
#[expect(clippy::expect_used)]
100-
self.writer.lock().expect("lock is poisoned")
106+
self.writer.lock().map_err(|_| crate::Error::Poisoned)
101107
}
102108

103-
pub fn path(&self) -> PathBuf {
104-
self.get_writer().path.clone()
109+
pub fn path(&self) -> crate::Result<PathBuf> {
110+
Ok(self.get_writer()?.path.clone())
105111
}
106112

107113
pub fn get_reader(&self) -> crate::Result<JournalBatchReader> {
108-
let raw_reader = JournalReader::new(self.path())?;
114+
let raw_reader = JournalReader::new(self.path()?)?;
109115
Ok(JournalBatchReader::new(raw_reader))
110116
}
111117

112118
/// Persists the journal.
113119
pub fn persist(&self, mode: PersistMode) -> crate::Result<()> {
114-
let mut journal_writer = self.get_writer();
120+
let mut journal_writer = self.get_writer()?;
115121
journal_writer.persist(mode).map_err(Into::into)
116122
}
117123

src/journal/test.rs

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ fn journal_rotation() -> crate::Result<()> {
3434

3535
{
3636
let journal = Journal::create_new(&path)?;
37-
let mut writer = journal.get_writer();
37+
let mut writer = journal.get_writer()?;
3838

3939
writer.write_batch(
4040
[
@@ -70,7 +70,7 @@ fn journal_recovery_active() -> crate::Result<()> {
7070

7171
{
7272
let journal = Journal::create_new(&path)?;
73-
let mut writer = journal.get_writer();
73+
let mut writer = journal.get_writer()?;
7474

7575
writer.write_batch(
7676
[
@@ -110,7 +110,7 @@ fn journal_recovery_active() -> crate::Result<()> {
110110
assert!(next_next_path.try_exists()?);
111111

112112
let journal_recovered = Journal::recover(dir2, CompressionType::None, 0)?;
113-
assert_eq!(journal_recovered.active.path(), next_next_path);
113+
assert_eq!(journal_recovered.active.path()?, next_next_path);
114114
assert_eq!(journal_recovered.sealed, &[(0, path), (1, next_path)]);
115115

116116
Ok(())
@@ -133,7 +133,7 @@ fn journal_recovery_active_lz4() -> crate::Result<()> {
133133

134134
{
135135
let journal = Journal::create_new(&path)?.with_compression(CompressionType::Lz4, 1);
136-
let mut writer = journal.get_writer();
136+
let mut writer = journal.get_writer()?;
137137

138138
writer.write_batch(
139139
[
@@ -173,7 +173,7 @@ fn journal_recovery_active_lz4() -> crate::Result<()> {
173173
assert!(next_next_path.try_exists()?);
174174

175175
let journal_recovered = Journal::recover(dir2, CompressionType::None, 0)?;
176-
assert_eq!(journal_recovered.active.path(), next_next_path);
176+
assert_eq!(journal_recovered.active.path()?, next_next_path);
177177
assert_eq!(journal_recovered.sealed, &[(0, path), (1, next_path)]);
178178

179179
Ok(())
@@ -194,7 +194,7 @@ fn journal_recovery_no_active() -> crate::Result<()> {
194194
let journal = Journal::create_new(&path)?;
195195

196196
{
197-
let mut writer = journal.get_writer();
197+
let mut writer = journal.get_writer()?;
198198

199199
writer.write_batch(
200200
[
@@ -217,7 +217,7 @@ fn journal_recovery_no_active() -> crate::Result<()> {
217217
assert!(!next_path.try_exists()?);
218218

219219
let journal_recovered = Journal::recover(dir2, CompressionType::None, 0)?;
220-
assert_eq!(journal_recovered.active.path(), path);
220+
assert_eq!(journal_recovered.active.path()?, path);
221221
assert_eq!(journal_recovered.sealed, &[]);
222222

223223
Ok(())
@@ -241,7 +241,7 @@ fn journal_truncation_corrupt_bytes() -> crate::Result<()> {
241241
{
242242
let journal = Journal::create_new(&path)?;
243243
journal
244-
.get_writer()
244+
.get_writer()?
245245
.write_batch(values.iter(), values.len(), 0)?;
246246
}
247247

@@ -301,7 +301,7 @@ fn journal_truncation_repeating_start_marker() -> crate::Result<()> {
301301
{
302302
let journal = Journal::create_new(&path)?;
303303
journal
304-
.get_writer()
304+
.get_writer()?
305305
.write_batch(values.iter(), values.len(), 0)?;
306306
}
307307

@@ -369,7 +369,7 @@ fn journal_truncation_repeating_end_marker() -> crate::Result<()> {
369369
{
370370
let journal = Journal::create_new(&path)?;
371371
journal
372-
.get_writer()
372+
.get_writer()?
373373
.write_batch(values.iter(), values.len(), 0)?;
374374
}
375375

@@ -429,7 +429,7 @@ fn journal_truncation_repeating_item_marker() -> crate::Result<()> {
429429
{
430430
let journal = Journal::create_new(&path)?;
431431
journal
432-
.get_writer()
432+
.get_writer()?
433433
.write_batch(values.iter(), values.len(), 0)?;
434434
}
435435

src/keyspace/mod.rs

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -236,7 +236,7 @@ impl Keyspace {
236236
pub fn clear(&self) -> crate::Result<()> {
237237
use std::sync::atomic::Ordering;
238238

239-
let mut journal_writer = self.supervisor.journal.get_writer();
239+
let mut journal_writer = self.supervisor.journal.get_writer()?;
240240

241241
// IMPORTANT: Check the poisoned flag after getting journal mutex, otherwise TOCTOU
242242
if self.is_poisoned.load(Ordering::Relaxed) {
@@ -719,7 +719,7 @@ impl Keyspace {
719719
#[doc(hidden)]
720720
pub fn rotate_memtable(&self) -> crate::Result<bool> {
721721
log::trace!("acquiring journal lock");
722-
let journal_writer = self.supervisor.journal.get_writer();
722+
let journal_writer = self.supervisor.journal.get_writer()?;
723723
let active_memtable_id = self.tree.active_memtable().id();
724724
self.inner_rotate_memtable(journal_writer, active_memtable_id)
725725
}
@@ -916,7 +916,7 @@ impl Keyspace {
916916
let key = key.into();
917917
let value = value.into();
918918

919-
let mut journal_writer = self.supervisor.journal.get_writer();
919+
let mut journal_writer = self.supervisor.journal.get_writer()?;
920920

921921
// IMPORTANT: Check the poisoned flag after getting journal mutex, otherwise TOCTOU
922922
if self.is_poisoned.load(Ordering::Relaxed) {
@@ -987,7 +987,7 @@ impl Keyspace {
987987

988988
let key = key.into();
989989

990-
let mut journal_writer = self.supervisor.journal.get_writer();
990+
let mut journal_writer = self.supervisor.journal.get_writer()?;
991991

992992
// IMPORTANT: Check the poisoned flag after getting journal mutex, otherwise TOCTOU
993993
if self.is_poisoned.load(Ordering::Relaxed) {
@@ -1070,7 +1070,7 @@ impl Keyspace {
10701070

10711071
let key = key.into();
10721072

1073-
let mut journal_writer = self.supervisor.journal.get_writer();
1073+
let mut journal_writer = self.supervisor.journal.get_writer()?;
10741074

10751075
// IMPORTANT: Check the poisoned flag after getting journal mutex, otherwise TOCTOU
10761076
if self.is_poisoned.load(Ordering::Relaxed) {

src/worker_pool.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ fn worker_tick(ctx: &WorkerState) -> crate::Result<bool> {
140140
}
141141
WorkerMessage::RotateMemtable(keyspace, memtable_id) => {
142142
log::trace!("acquiring journal lock");
143-
let journal_writer = keyspace.supervisor.journal.get_writer();
143+
let journal_writer = keyspace.supervisor.journal.get_writer()?;
144144
keyspace.inner_rotate_memtable(journal_writer, memtable_id)?;
145145
}
146146
WorkerMessage::Flush => {
@@ -150,7 +150,7 @@ fn worker_tick(ctx: &WorkerState) -> crate::Result<bool> {
150150

151151
{
152152
log::trace!("acquiring journal lock to maybe rotate journal");
153-
let mut journal_writer = ctx.supervisor.journal.get_writer();
153+
let mut journal_writer = ctx.supervisor.journal.get_writer()?;
154154

155155
if journal_writer.pos()? > 64_000_000 {
156156
#[expect(clippy::expect_used)]

0 commit comments

Comments
 (0)