Skip to content
Merged
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
187 changes: 176 additions & 11 deletions src/io/sqlite_store/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ use lightning_types::string::PrintableString;
use rusqlite::{named_params, Connection};

use crate::io::utils::{check_namespace_key_validity, create_dir_all_private};
use crate::runtime::StoreRuntime;

mod migrations;

Expand All @@ -51,6 +52,8 @@ const PAGE_SIZE: usize = 50;
/// [SQLite]: https://sqlite.org
pub struct SqliteStore {
inner: Arc<SqliteStoreInner>,
// Drive I/O independently of synchronous callers that can exhaust their blocking pool.
internal_runtime: Arc<StoreRuntime>,

// Version counter to ensure that writes are applied in the correct order. It is assumed that read and list
// operations aren't sensitive to the order of execution.
Expand All @@ -70,9 +73,10 @@ impl SqliteStore {
data_dir: PathBuf, db_file_name: Option<String>, kv_table_name: Option<String>,
) -> io::Result<Self> {
let inner = Arc::new(SqliteStoreInner::new(data_dir, db_file_name, kv_table_name)?);
let internal_runtime = Arc::new(StoreRuntime::new("ldk-node-sqlite-runtime", 2, "SQLite")?);

let next_write_version = AtomicU64::new(1);
Ok(Self { inner, next_write_version })
Ok(Self { inner, internal_runtime, next_write_version })
}

fn build_locking_key(
Expand Down Expand Up @@ -108,10 +112,12 @@ impl KVStore for SqliteStore {
let secondary_namespace = secondary_namespace.to_string();
let key = key.to_string();
let inner = Arc::clone(&self.inner);
let fut = tokio::task::spawn_blocking(move || {
let runtime = Arc::clone(&self.internal_runtime);
let fut = runtime.handle().spawn_blocking(move || {
inner.read_internal(&primary_namespace, &secondary_namespace, &key)
});
async move {
let _runtime = runtime;
fut.await.unwrap_or_else(|e| {
let msg = format!("Failed to IO operation due join error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, msg))
Expand All @@ -128,7 +134,8 @@ impl KVStore for SqliteStore {
let secondary_namespace = secondary_namespace.to_string();
let key = key.to_string();
let inner = Arc::clone(&self.inner);
let fut = tokio::task::spawn_blocking(move || {
let runtime = Arc::clone(&self.internal_runtime);
let fut = runtime.handle().spawn_blocking(move || {
inner.write_internal(
inner_lock_ref,
locking_key,
Expand All @@ -140,6 +147,7 @@ impl KVStore for SqliteStore {
)
});
async move {
let _runtime = runtime;
fut.await.unwrap_or_else(|e| {
let msg = format!("Failed to IO operation due join error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, msg))
Expand All @@ -156,7 +164,8 @@ impl KVStore for SqliteStore {
let secondary_namespace = secondary_namespace.to_string();
let key = key.to_string();
let inner = Arc::clone(&self.inner);
let fut = tokio::task::spawn_blocking(move || {
let runtime = Arc::clone(&self.internal_runtime);
let fut = runtime.handle().spawn_blocking(move || {
inner.remove_internal(
inner_lock_ref,
locking_key,
Expand All @@ -167,6 +176,7 @@ impl KVStore for SqliteStore {
)
});
async move {
let _runtime = runtime;
fut.await.unwrap_or_else(|e| {
let msg = format!("Failed to IO operation due join error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, msg))
Expand All @@ -180,10 +190,12 @@ impl KVStore for SqliteStore {
let primary_namespace = primary_namespace.to_string();
let secondary_namespace = secondary_namespace.to_string();
let inner = Arc::clone(&self.inner);
let fut = tokio::task::spawn_blocking(move || {
inner.list_internal(&primary_namespace, &secondary_namespace)
});
let runtime = Arc::clone(&self.internal_runtime);
let fut = runtime
.handle()
.spawn_blocking(move || inner.list_internal(&primary_namespace, &secondary_namespace));
async move {
let _runtime = runtime;
fut.await.unwrap_or_else(|e| {
let msg = format!("Failed to IO operation due join error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, msg))
Expand All @@ -199,10 +211,12 @@ impl PaginatedKVStore for SqliteStore {
let primary_namespace = primary_namespace.to_string();
let secondary_namespace = secondary_namespace.to_string();
let inner = Arc::clone(&self.inner);
let fut = tokio::task::spawn_blocking(move || {
let runtime = Arc::clone(&self.internal_runtime);
let fut = runtime.handle().spawn_blocking(move || {
inner.list_paginated_internal(&primary_namespace, &secondary_namespace, page_token)
});
async move {
let _runtime = runtime;
fut.await.unwrap_or_else(|e| {
let msg = format!("Failed to IO operation due join error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, msg))
Expand All @@ -216,8 +230,10 @@ impl MigratableKVStore for SqliteStore {
&self,
) -> impl Future<Output = Result<Vec<(String, String, String)>, io::Error>> + 'static + Send {
let inner = Arc::clone(&self.inner);
let fut = tokio::task::spawn_blocking(move || inner.list_all_keys_internal());
let runtime = Arc::clone(&self.internal_runtime);
let fut = runtime.handle().spawn_blocking(move || inner.list_all_keys_internal());
async move {
let _runtime = runtime;
fut.await.unwrap_or_else(|e| {
let msg = format!("Failed to IO operation due join error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, msg))
Expand Down Expand Up @@ -717,15 +733,164 @@ mod tests {
do_read_write_remove_list_persist, do_test_store, random_storage_path,
};

impl Drop for SqliteStore {
impl Drop for SqliteStoreInner {
fn drop(&mut self) {
match fs::remove_dir_all(&self.inner.data_dir) {
match fs::remove_dir_all(&self.data_dir) {
Err(e) => println!("Failed to remove test store directory: {}", e),
_ => {},
}
}
}

#[tokio::test]
async fn pending_operations_keep_runtime_alive_after_store_drop() {
async fn check_operation(
operation: impl FnOnce(
&SqliteStore,
Arc<SqliteStoreInner>,
) -> std::pin::Pin<Box<dyn Future<Output = ()> + Send>>,
) {
let store = SqliteStore::new(random_storage_path(), None, None).unwrap();
KVStore::write(&store, "test", "", "key", vec![42]).await.unwrap();
let runtime_handle = store.internal_runtime.handle().clone();

// Occupy both blocking threads so the operation cannot start before the store is
// dropped. Channels make this deterministic without relying on scheduling delays.
let mut blockers = Vec::new();
let mut releases = Vec::new();
for _ in 0..2 {
let (started_tx, started_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
blockers.push(store.internal_runtime.handle().spawn_blocking(move || {
started_tx.send(()).unwrap();
let _ = release_rx.recv();
}));
started_rx.recv_timeout(std::time::Duration::from_secs(5)).unwrap();
releases.push(release_tx);
}

let future = operation(&store, Arc::clone(&store.inner));
drop(store);
// A handle alone does not keep a runtime alive. Probe it while the operation is
// still pending, since eagerly spawned blocking work can survive runtime shutdown.
tokio::time::timeout(std::time::Duration::from_secs(5), runtime_handle.spawn(async {}))
.await
.unwrap()
.expect("SQLite runtime shut down while an operation future was still outstanding");
for release in releases {
release.send(()).unwrap();
}
tokio::time::timeout(std::time::Duration::from_secs(5), future)
.await
.expect("SQLite operation did not complete after the store was dropped");
for blocker in blockers {
blocker.await.unwrap();
}
assert!(runtime_handle.spawn(async {}).await.unwrap_err().is_cancelled());
}

check_operation(|store, _| {
let future = KVStore::read(store, "test", "", "key");
Box::pin(async move {
assert_eq!(future.await.expect("read was cancelled after store drop"), vec![42]);
})
})
.await;
check_operation(|store, inner| {
let future = KVStore::write(store, "test", "", "key", vec![43]);
Box::pin(async move {
future.await.expect("write was cancelled after store drop");
assert_eq!(inner.read_internal("test", "", "key").unwrap(), vec![43]);
})
})
.await;
check_operation(|store, inner| {
let future = KVStore::remove(store, "test", "", "key", false);
Box::pin(async move {
future.await.expect("remove was cancelled after store drop");
assert_eq!(
inner.read_internal("test", "", "key").unwrap_err().kind(),
io::ErrorKind::NotFound
);
})
})
.await;
check_operation(|store, _| {
let future = KVStore::list(store, "test", "");
Box::pin(async move {
assert_eq!(future.await.expect("list was cancelled after store drop"), vec!["key"]);
})
})
.await;
check_operation(|store, _| {
let future = PaginatedKVStore::list_paginated(store, "test", "", None);
Box::pin(async move {
let page = future.await.expect("paginated list was cancelled after store drop");
assert_eq!(page.keys, vec!["key"]);
assert!(page.next_page_token.is_none());
})
})
.await;
check_operation(|store, _| {
let future = MigratableKVStore::list_all_keys(store);
Box::pin(async move {
assert_eq!(
future.await.expect("list_all_keys was cancelled after store drop"),
vec![("test".to_string(), "".to_string(), "key".to_string())]
);
})
})
.await;
}

#[test]
fn io_completes_with_a_saturated_caller_blocking_pool() {
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.max_blocking_threads(1)
.enable_all()
.build()
.unwrap();
let store = SqliteStore::new(random_storage_path(), None, None).unwrap();

let result = runtime.block_on(async move {
// This synchronous caller occupies the only blocking thread. SQLite must be able
// to make progress without submitting its I/O to that same pool.
tokio::task::spawn_blocking(move || {
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(async {
// The async worker remains free to drive the timeout even if I/O deadlocks.
tokio::time::timeout(std::time::Duration::from_secs(5), async {
KVStore::write(&store, "test", "", "key", vec![42]).await.unwrap();
assert_eq!(
KVStore::read(&store, "test", "", "key").await.unwrap(),
vec![42]
);
assert_eq!(
KVStore::list(&store, "test", "").await.unwrap(),
vec!["key"]
);
let page = PaginatedKVStore::list_paginated(&store, "test", "", None)
.await
.unwrap();
assert_eq!(page.keys, vec!["key"]);
assert_eq!(
MigratableKVStore::list_all_keys(&store).await.unwrap(),
vec![("test".to_string(), "".to_string(), "key".to_string())]
);
KVStore::remove(&store, "test", "", "key", false).await.unwrap();
assert!(KVStore::list(&store, "test", "").await.unwrap().is_empty());
})
.await
})
})
})
.await
.unwrap()
});
assert!(result.is_ok(), "SQLite I/O deadlocked on the caller's blocking pool");
}

#[cfg(unix)]
#[test]
fn creates_private_database_storage() {
Expand Down
Loading