This commit is contained in:
Silas Brack 2026-03-07 13:09:53 +01:00
parent 1461b41a36
commit 7f3ec69cf6
7 changed files with 236 additions and 158 deletions

194
src/db.rs
View file

@ -1,7 +1,6 @@
use rusqlite::{params, Connection, OpenFlags};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use tokio::sync::{mpsc, oneshot};
use crate::error::AppError;
@ -151,174 +150,73 @@ pub fn all_records(conn: &Connection) -> Result<Vec<Record>, AppError> {
Ok(records)
}
// --- Write commands ---
pub enum WriteCmd {
Put {
key: String,
volumes: Vec<String>,
size: Option<i64>,
reply: oneshot::Sender<Result<(), AppError>>,
},
Delete {
key: String,
reply: oneshot::Sender<Result<(), AppError>>,
},
BulkPut {
records: Vec<(String, Vec<String>, Option<i64>)>,
reply: oneshot::Sender<Result<(), AppError>>,
},
}
fn execute_cmd(
conn: &Connection,
cmd: WriteCmd,
) -> (Result<(), AppError>, oneshot::Sender<Result<(), AppError>>) {
match cmd {
WriteCmd::Put {
key,
volumes,
size,
reply,
} => {
let volumes_json = encode_volumes(&volumes);
let result = conn
.prepare_cached(
"INSERT INTO kv (key, volumes, size) VALUES (?1, ?2, ?3)
ON CONFLICT(key) DO UPDATE SET volumes = ?2, size = ?3",
)
.and_then(|mut s| s.execute(params![key, volumes_json, size]))
.map(|_| ())
.map_err(AppError::from);
(result, reply)
}
WriteCmd::Delete { key, reply } => {
let result = conn
.prepare_cached("DELETE FROM kv WHERE key = ?1")
.and_then(|mut s| s.execute(params![key]))
.map(|_| ())
.map_err(AppError::from);
(result, reply)
}
WriteCmd::BulkPut { records, reply } => {
let result = (|| -> Result<(), AppError> {
let mut stmt = conn.prepare_cached(
"INSERT INTO kv (key, volumes, size) VALUES (?1, ?2, ?3)
ON CONFLICT(key) DO UPDATE SET volumes = ?2, size = ?3",
)?;
for (key, volumes, size) in &records {
let volumes_json = encode_volumes(volumes);
stmt.execute(params![key, volumes_json, size])?;
}
Ok(())
})();
(result, reply)
}
}
}
// --- WriterHandle ---
#[derive(Clone)]
pub struct WriterHandle {
tx: mpsc::Sender<WriteCmd>,
conn: Arc<Mutex<Connection>>,
}
impl WriterHandle {
pub fn new(path: &str) -> Self {
let conn = open_readwrite(path);
create_tables(&conn);
Self {
conn: Arc::new(Mutex::new(conn)),
}
}
pub async fn put(
&self,
key: String,
volumes: Vec<String>,
size: Option<i64>,
) -> Result<(), AppError> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCmd::Put {
key,
volumes,
size,
reply: reply_tx,
})
.await
.map_err(|_| AppError::WriterDead)?;
reply_rx.await.map_err(|_| AppError::WriterDroppedReply)?
let conn = self.conn.clone();
tokio::task::spawn_blocking(move || {
let conn = conn.lock().unwrap();
let volumes_json = encode_volumes(&volumes);
conn.prepare_cached(
"INSERT INTO kv (key, volumes, size) VALUES (?1, ?2, ?3)
ON CONFLICT(key) DO UPDATE SET volumes = ?2, size = ?3",
)?
.execute(params![key, volumes_json, size])?;
Ok(())
})
.await
.unwrap()
}
pub async fn delete(&self, key: String) -> Result<(), AppError> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCmd::Delete {
key,
reply: reply_tx,
})
.await
.map_err(|_| AppError::WriterDead)?;
reply_rx.await.map_err(|_| AppError::WriterDroppedReply)?
let conn = self.conn.clone();
tokio::task::spawn_blocking(move || {
let conn = conn.lock().unwrap();
conn.prepare_cached("DELETE FROM kv WHERE key = ?1")?
.execute(params![key])?;
Ok(())
})
.await
.unwrap()
}
pub async fn bulk_put(
&self,
records: Vec<(String, Vec<String>, Option<i64>)>,
) -> Result<(), AppError> {
let (reply_tx, reply_rx) = oneshot::channel();
self.tx
.send(WriteCmd::BulkPut {
records,
reply: reply_tx,
})
.await
.map_err(|_| AppError::WriterDead)?;
reply_rx.await.map_err(|_| AppError::WriterDroppedReply)?
let conn = self.conn.clone();
tokio::task::spawn_blocking(move || {
let conn = conn.lock().unwrap();
let mut stmt = conn.prepare_cached(
"INSERT INTO kv (key, volumes, size) VALUES (?1, ?2, ?3)
ON CONFLICT(key) DO UPDATE SET volumes = ?2, size = ?3",
)?;
for (key, volumes, size) in &records {
let volumes_json = encode_volumes(volumes);
stmt.execute(params![key, volumes_json, size])?;
}
Ok(())
})
.await
.unwrap()
}
}
// --- spawn_writer ---
pub fn spawn_writer(path: String) -> (WriterHandle, oneshot::Receiver<()>) {
let (tx, mut rx) = mpsc::channel::<WriteCmd>(4096);
let (ready_tx, ready_rx) = oneshot::channel();
std::thread::spawn(move || {
let conn = open_readwrite(&path);
create_tables(&conn);
let _ = ready_tx.send(());
loop {
let Some(first) = rx.blocking_recv() else {
break;
};
let mut batch = vec![first];
while batch.len() < 512 {
match rx.try_recv() {
Ok(cmd) => batch.push(cmd),
Err(_) => break,
}
}
let _ = conn.execute_batch("BEGIN");
let mut replies: Vec<(Result<(), AppError>, oneshot::Sender<Result<(), AppError>>)> =
Vec::with_capacity(batch.len());
for (i, cmd) in batch.into_iter().enumerate() {
let sp = format!("sp{i}");
let _ = conn.execute(&format!("SAVEPOINT {sp}"), []);
let (result, reply) = execute_cmd(&conn, cmd);
if result.is_ok() {
let _ = conn.execute(&format!("RELEASE {sp}"), []);
} else {
let _ = conn.execute(&format!("ROLLBACK TO {sp}"), []);
let _ = conn.execute(&format!("RELEASE {sp}"), []);
}
replies.push((result, reply));
}
let _ = conn.execute_batch("COMMIT");
for (result, reply) in replies {
let _ = reply.send(result);
}
}
});
(WriterHandle { tx }, ready_rx)
}

View file

@ -5,8 +5,6 @@ use axum::response::{IntoResponse, Response};
pub enum AppError {
NotFound,
Db(rusqlite::Error),
WriterDead,
WriterDroppedReply,
VolumeError(String),
NoHealthyVolume,
}
@ -25,8 +23,6 @@ impl std::fmt::Display for AppError {
match self {
AppError::NotFound => write!(f, "not found"),
AppError::Db(e) => write!(f, "database error: {e}"),
AppError::WriterDead => write!(f, "writer dead"),
AppError::WriterDroppedReply => write!(f, "writer dropped reply"),
AppError::VolumeError(msg) => write!(f, "volume error: {msg}"),
AppError::NoHealthyVolume => write!(f, "no healthy volume available"),
}

View file

@ -24,8 +24,7 @@ pub async fn build_app(config: config::Config) -> axum::Router {
});
}
let (writer, ready_rx) = db::spawn_writer(db_path.to_string());
ready_rx.await.expect("writer failed to initialize");
let writer = db::WriterHandle::new(db_path);
let num_readers = std::thread::available_parallelism()
.map(|n| n.get())

View file

@ -88,8 +88,7 @@ pub async fn run(config: &Config, dry_run: bool) {
}
// Open writer for updates
let (writer, ready_rx) = db::spawn_writer(db_path.to_string());
ready_rx.await.expect("writer failed to initialize");
let writer = db::WriterHandle::new(db_path);
let client = VolumeClient::new();
let mut moved = 0;

View file

@ -66,8 +66,7 @@ pub async fn run(config: &Config) {
let _ = std::fs::remove_file(format!("{db_path}-wal"));
let _ = std::fs::remove_file(format!("{db_path}-shm"));
let (writer, ready_rx) = db::spawn_writer(db_path.to_string());
ready_rx.await.expect("writer failed to initialize");
let writer = db::WriterHandle::new(db_path);
let volume_urls = config.volume_urls();