Skip to content

Commit b09cbaf

Browse files
committed
[#226] row append 경로의 storage lock 누락 수정
append_table_rows_with_buffer_limit가 row_storage_lock 없이 row_buffer_pool만 갱신해 update/delete의 snapshot 기반 replace_rows와 경합할 수 있었다. update/delete가 snapshot을 읽은 뒤 replace_rows로 pending append를 비우는 사이 INSERT가 끼면 새 행이 snapshot에도 pending에도 남지 않아 유실될 수 있었다. append 경로도 update/delete와 동일하게 row_storage_lock을 잡도록 바꾸고, 버퍼 한도 초과 시에는 이미 보유한 락 안에서 flush_row_buffers_locked(false)를 호출해 재진입 deadlock을 피했다. append가 row_storage_lock을 기다리는 회귀 테스트도 추가했다.
1 parent 4f57c7a commit b09cbaf

1 file changed

Lines changed: 42 additions & 1 deletion

File tree

src/engine/actions/dml/scan.rs

Lines changed: 42 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@ impl DBEngine {
6969
return Ok(());
7070
}
7171

72+
let _guard = self.row_storage_lock.lock().await;
7273
let segment_path = self.row_segment_path(table_name)?;
7374
let frame = encode_row_frames(rows)?;
7475

@@ -80,7 +81,7 @@ impl DBEngine {
8081
};
8182

8283
if buffered_bytes >= buffer_limit_bytes {
83-
self.flush_row_buffers().await?;
84+
self.flush_row_buffers_locked(false).await?;
8485
}
8586

8687
Ok(())
@@ -204,6 +205,7 @@ impl DBEngine {
204205
Ok(rows)
205206
}
206207

208+
#[cfg(test)]
207209
pub(crate) async fn flush_row_buffers(&self) -> errors::Result<()> {
208210
let _guard = self.row_storage_lock.lock().await;
209211
self.flush_row_buffers_locked(false).await
@@ -333,6 +335,7 @@ impl DBEngine {
333335
#[cfg(test)]
334336
mod tests {
335337
use std::path::PathBuf;
338+
use std::sync::Arc;
336339

337340
use crate::config::launch_config::LaunchConfig;
338341
use crate::engine::DBEngine;
@@ -610,4 +613,42 @@ mod tests {
610613

611614
assert!(engine.row_buffer_pool.lock().await.is_unsynced_empty());
612615
}
616+
617+
#[tokio::test]
618+
async fn append_table_rows_waits_for_row_storage_lock() {
619+
let base_path = PathBuf::from(format!(
620+
"target/test_row_segments/append_waits_for_storage_lock_{}",
621+
std::process::id()
622+
));
623+
if base_path.exists() {
624+
tokio::fs::remove_dir_all(&base_path).await.unwrap();
625+
}
626+
627+
let config = LaunchConfig::default_for_base_path(&base_path);
628+
let table_name = TableName::new(Some("rrdb".to_string()), "users".to_string());
629+
let rows_path = PathBuf::from(&config.data_directory)
630+
.join("rrdb")
631+
.join("tables")
632+
.join("users")
633+
.join("rows");
634+
tokio::fs::create_dir_all(&rows_path).await.unwrap();
635+
636+
let engine = Arc::new(DBEngine::new(config));
637+
let row = TableDataRow {
638+
fields: vec![TableDataField {
639+
table_name: table_name.clone(),
640+
column_name: "id".to_string(),
641+
data: TableDataFieldType::Integer(1),
642+
}],
643+
};
644+
let _guard = engine.row_storage_lock.lock().await;
645+
let append_task = {
646+
let engine = engine.clone();
647+
tokio::spawn(async move { engine.append_table_rows(&table_name, &[row]).await })
648+
};
649+
650+
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
651+
652+
assert!(!append_task.is_finished());
653+
}
613654
}

0 commit comments

Comments
 (0)