use std::{ path::Path, sync::{Arc, RwLock}, }; use anyhow::{Context, Result, anyhow, ensure}; use deadpool_sqlite::{Config as DatabaseConfig, Hook, HookError, Pool, Runtime}; use rusqlite::{Connection, OpenFlags, OptionalExtension, params}; use crate::{Resolution, opcodec, resolution_from_operation}; const LATEST_OPERATION: &str = "SELECT operation FROM ops WHERE did = ?1 ORDER BY seq DESC LIMIT 1"; pub struct Db { pool: Pool, strings: RwLock>, } pub fn open_database(path: &Path, max_size: usize) -> Result> { // read-only open refuses to create an empty database where plox's should be let connection = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY) .with_context(|| format!("opening {}", path.display()))?; let version: i64 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; ensure!( version == 1, "unsupported plox database version {version}; expected compact format version 1" ); let compacting: bool = connection.query_row( "SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'compaction')", [], |row| row.get(0), )?; if compacting { let complete: bool = connection.query_row( "SELECT complete FROM compaction WHERE singleton = 1", [], |row| row.get(0), )?; ensure!( complete, "database compaction is unfinished; resume plox compact before serving it" ); } connection .prepare(LATEST_OPERATION) .context("reading plox ops schema")?; connection .prepare("SELECT id, value FROM strings") .context("reading plox string dictionary schema")?; let pool = DatabaseConfig::new(path) .builder(Runtime::Tokio1)? .max_size(max_size) .post_create(Hook::async_fn(|connection, _| { Box::pin(async move { match connection .interact(|connection| { connection.busy_timeout(std::time::Duration::from_secs(5))?; connection.pragma_update(None, "query_only", true) }) .await { Ok(Ok(())) => Ok(()), Ok(Err(error)) => Err(HookError::Backend(error)), Err(error) => Err(HookError::message(error.to_string())), } }) })) .build() .context("creating database connection pool")?; Ok(Arc::new(Db { pool, strings: RwLock::new(Vec::new()), })) } pub async fn resolve(db: &Arc, did: &str) -> Result> { let connection = db .pool .get() .await .context("checking out database connection")?; let query_did = opcodec::encode_did(did); let db = Arc::clone(db); let operation = connection .interact(move |connection| -> Result<_> { let tx = connection.transaction()?; let Some(program) = tx .prepare_cached(LATEST_OPERATION)? .query_row(params![query_did], |row| row.get::<_, Vec>(0)) .optional()? else { return Ok(None); }; let strings = db .strings .read() .map_err(|_| anyhow!("string dictionary lock poisoned"))?; if let Ok(operation) = opcodec::decode(&program, |id| strings.get(id).map(String::as_str)) { return Ok(Some(operation)); } drop(strings); // plox's dictionary is append-only; retry with additions from the operation's snapshot. let mut strings = db .strings .write() .map_err(|_| anyhow!("string dictionary lock poisoned"))?; let mut stmt = tx.prepare_cached("SELECT id, value FROM strings WHERE id >= ?1 ORDER BY id")?; let mut rows = stmt.query([strings.len() as i64])?; while let Some(row) = rows.next()? { let id: i64 = row.get(0)?; ensure!( id == strings.len() as i64, "non-contiguous string dictionary" ); strings.push(row.get(1)?); } opcodec::decode(&program, |id| strings.get(id).map(String::as_str)) .context("decoding PLC operation") .map(Some) }) .await .map_err(|error| anyhow!("querying database: {error}"))??; operation .map(|operation| resolution_from_operation(did, operation)) .transpose() }