use alloc::string::{String, ToString};
use alloc::vec::Vec;
use std::path::{Path, PathBuf};
use crate::bytes::Bytes;
use crate::future::block_on;
use crate::persistence::NamespaceSummary;
use crate::persistence::turso::{self as database, SCHEMA_VERSION, SCHEMA_VERSION_KEY, connect};
use super::export::storage_error;
use super::{Bundle, BundleError, BundleManifest};
const MAGIC: &[u8; 16] = b"SQLite format 3\0";
#[derive(Debug)]
pub struct SqliteBundle {
database: turso::Database,
manifest: BundleManifest,
path: PathBuf,
}
impl SqliteBundle {
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self, BundleError> {
let path = path.as_ref();
if !is_sqlite(path)? {
return Err(BundleError::NotABundle);
}
let location = path.to_str().ok_or_else(|| {
BundleError::Storage(std::format!("bundle path {path:?} is not valid UTF-8"))
})?;
let database = block_on(turso::Builder::new_local(location).read_only(true).build())
.map_err(storage_error)?;
let connection = connect(&database).map_err(storage_error)?;
let expected = SCHEMA_VERSION.to_string();
match database::meta_get(&connection, SCHEMA_VERSION_KEY).map_err(missing_meta)? {
Some(found) if found == expected => {}
Some(found) => return Err(BundleError::UnsupportedDatabase(found)),
None => return Err(BundleError::NotABundle),
}
let manifest = BundleManifest::read(&connection)?;
manifest.warn_on_version_mismatch();
Ok(Self {
database,
manifest,
path: path.to_path_buf(),
})
}
pub fn manifest(&self) -> &BundleManifest {
&self.manifest
}
pub fn summary(&self) -> Vec<NamespaceSummary> {
self.read("summarize", database::summarize)
.unwrap_or_default()
}
fn read<T>(
&self,
name: &str,
read: impl FnOnce(&turso::Connection) -> Result<T, turso::Error>,
) -> Option<T> {
connect(&self.database)
.and_then(|connection| read(&connection))
.inspect_err(|err| log::warn!("Unable to {name} {}: {err}", self.describe()))
.ok()
}
}
impl Bundle for SqliteBundle {
fn get(&self, namespace: &str, key: &[u8]) -> Option<Bytes> {
self.read("read", |connection| {
let mut rows = block_on(connection.query(
"SELECT value FROM entries WHERE namespace = ?1 AND key = ?2",
(namespace, key.to_vec()),
))?;
let Some(row) = block_on(rows.next())? else {
return Ok(None);
};
let value: Vec<u8> = row.get(0)?;
Ok(Some(Bytes::from_bytes_vec(value)))
})
.flatten()
}
fn scan(&self, namespace: &str, visit: &mut dyn FnMut(&[u8], &[u8])) {
self.read("scan", |connection| {
let mut rows = block_on(connection.query(
"SELECT key, value FROM entries WHERE namespace = ?1",
(namespace,),
))?;
while let Some(row) = block_on(rows.next())? {
let key: Vec<u8> = row.get(0)?;
let value: Vec<u8> = row.get(1)?;
visit(&key, &value);
}
Ok(())
});
}
fn namespaces(&self) -> Vec<String> {
self.read("list", |connection| {
let mut rows = block_on(connection.query(
"SELECT DISTINCT namespace FROM entries ORDER BY namespace",
(),
))?;
let mut namespaces = Vec::new();
while let Some(row) = block_on(rows.next())? {
namespaces.push(row.get::<String>(0)?);
}
Ok(namespaces)
})
.unwrap_or_default()
}
fn describe(&self) -> String {
alloc::format!("bundle '{}' at {:?}", self.manifest.name, self.path)
}
}
pub(super) fn missing_meta(err: turso::Error) -> BundleError {
let message = err.to_string();
if message.contains("no such table") {
BundleError::NotABundle
} else {
BundleError::Storage(message)
}
}
fn is_sqlite(path: &Path) -> Result<bool, BundleError> {
use std::io::Read;
let mut header = [0u8; 16];
match std::fs::File::open(path)?.read_exact(&mut header) {
Ok(()) => Ok(&header == MAGIC),
Err(err) if err.kind() == std::io::ErrorKind::UnexpectedEof => Ok(false),
Err(err) => Err(err.into()),
}
}