use std::path::Path;
use std::sync::{Arc, Mutex};
use oxigraph::model::{
GraphName, GraphNameRef, NamedOrBlankNodeRef, Quad, TermRef, Triple, TripleRef,
};
use oxigraph::store::{QuadIter, Store};
use crate::persist::{self, DiskStore};
use crate::{Error, Result};
#[derive(Clone, Copy, Debug, Default)]
pub struct StatementPattern<'a> {
pub subject: Option<NamedOrBlankNodeRef<'a>>,
pub predicate: Option<oxigraph::model::NamedNodeRef<'a>>,
pub object: Option<TermRef<'a>>,
pub graph_name: Option<GraphNameRef<'a>>,
}
#[must_use]
pub struct StatementMatches {
inner: QuadIter<'static>,
}
impl Iterator for StatementMatches {
type Item = Result<Quad>;
fn next(&mut self) -> Option<Self::Item> {
self.inner
.next()
.map(|item| item.map_err(|error| Error::Storage(error.to_string())))
}
}
#[derive(Clone)]
pub struct Model {
store: Store,
disk: Option<DiskStore>,
write_lock: Arc<Mutex<()>>,
}
impl Model {
pub fn new() -> Result<Self> {
Store::new()
.map(|store| Self {
store,
disk: None,
write_lock: Arc::new(Mutex::new(())),
})
.map_err(|error| Error::Storage(error.to_string()))
}
pub fn open(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref();
let disk = DiskStore::open(path)?;
let store = Store::new().map_err(|error| Error::OpenStore {
path: path.to_owned(),
message: error.to_string(),
})?;
disk.load_into(&store)?;
Ok(Self {
store,
disk: Some(disk),
write_lock: Arc::new(Mutex::new(())),
})
}
#[must_use]
pub fn store(&self) -> &Store {
&self.store
}
pub fn storage_backend_available(name: &str) -> Result<bool> {
match name {
"memory" | "fjall" => Ok(true),
"rocksdb" | "redb" => Err(Error::Unsupported(
"storage backend was replaced by fjall; use Model::open".into(),
)),
other => Err(Error::Unsupported(format!(
"storage backend '{other}' is not recognized"
))),
}
}
pub fn add(&self, statement: impl Into<Triple>) -> Result<bool> {
self.add_to_graph(statement, GraphName::DefaultGraph)
}
pub fn add_to_graph(
&self,
statement: impl Into<Triple>,
graph_name: impl Into<GraphName>,
) -> Result<bool> {
let triple = statement.into();
let quad = Quad::new(triple.subject, triple.predicate, triple.object, graph_name);
self.insert_quad(quad)
}
pub fn insert_quad(&self, quad: Quad) -> Result<bool> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let inserted = !self
.store
.contains(quad.as_ref())
.map_err(|error| Error::Storage(error.to_string()))?;
if !inserted {
return Ok(false);
}
self.store
.insert(&quad)
.map_err(|error| Error::Storage(error.to_string()))?;
if let Some(disk) = &self.disk {
let canonical = persist::stored_matching_quad(&self.store, &quad)?;
if let Err(error) = disk.insert(&canonical) {
if let Err(reload_error) = self.reload_store_from_disk_unlocked(disk) {
let _ = self.store.remove(canonical.as_ref());
return Err(Error::Storage(format!(
"durable insert failed ({error}); rollback from disk also failed ({reload_error})"
)));
}
return Err(error);
}
}
Ok(true)
}
pub fn remove_quad(&self, quad: &Quad) -> Result<bool> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let removed = self
.store
.contains(quad.as_ref())
.map_err(|error| Error::Storage(error.to_string()))?;
if !removed {
return Ok(false);
}
let canonical = persist::stored_matching_quad(&self.store, quad)?;
self.store
.remove(quad.as_ref())
.map_err(|error| Error::Storage(error.to_string()))?;
if let Some(disk) = &self.disk {
if let Err(error) = disk.remove_rdf_equal(&canonical) {
if let Err(reload_error) = self.reload_store_from_disk_unlocked(disk) {
let _ = self.store.insert(&canonical);
return Err(Error::Storage(format!(
"durable remove failed ({error}); rollback from disk also failed ({reload_error})"
)));
}
return Err(error);
}
}
Ok(true)
}
pub fn remove(&self, statement: impl Into<Triple>) -> Result<bool> {
self.remove_from_graph(statement, GraphName::DefaultGraph)
}
pub fn remove_from_graph(
&self,
statement: impl Into<Triple>,
graph_name: impl Into<GraphName>,
) -> Result<bool> {
let triple = statement.into();
let quad = Quad::new(triple.subject, triple.predicate, triple.object, graph_name);
self.remove_quad(&quad)
}
pub fn contains(&self, statement: TripleRef<'_>) -> Result<bool> {
self.contains_in_graph(statement, GraphNameRef::DefaultGraph)
}
pub fn contains_in_graph(
&self,
statement: TripleRef<'_>,
graph_name: GraphNameRef<'_>,
) -> Result<bool> {
self.store
.contains(oxigraph::model::QuadRef::new(
statement.subject,
statement.predicate,
statement.object,
graph_name,
))
.map_err(|error| Error::Storage(error.to_string()))
}
pub fn len(&self) -> Result<usize> {
self.store
.len()
.map_err(|error| Error::Storage(error.to_string()))
}
pub fn is_empty(&self) -> Result<bool> {
self.store
.is_empty()
.map_err(|error| Error::Storage(error.to_string()))
}
pub fn find(&self, pattern: StatementPattern<'_>) -> StatementMatches {
StatementMatches {
inner: self.store.quads_for_pattern(
pattern.subject,
pattern.predicate,
pattern.object,
pattern.graph_name,
),
}
}
pub(crate) fn run_sparql_update(
&self,
update: impl FnOnce(&Store) -> Result<()>,
) -> Result<()> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
update(&self.store)?;
let Some(disk) = &self.disk else {
return Ok(());
};
if let Err(error) = disk.replace_all_from_store(&self.store) {
if let Err(reload_error) = self.reload_store_from_disk_unlocked(disk) {
return Err(Error::Storage(format!(
"durable sync failed after SPARQL Update ({error}); rollback from disk also failed ({reload_error})"
)));
}
return Err(Error::Storage(format!(
"durable sync failed after SPARQL Update; in-memory store rolled back to disk: {error}"
)));
}
Ok(())
}
fn reload_store_from_disk_unlocked(&self, disk: &DiskStore) -> Result<()> {
self.store
.clear()
.map_err(|error| Error::Storage(error.to_string()))?;
disk.load_into(&self.store)
}
}