use std::fs::File;
use std::io::{BufReader, BufWriter};
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, RwLock};
use std::thread::ThreadId;
use oxigraph::io::{RdfFormat, RdfParser, RdfSerializer};
use oxigraph::model::{
GraphName, GraphNameRef, NamedOrBlankNodeRef, Quad, QuadRef, TermRef, Triple, TripleRef,
};
use oxigraph::store::{QuadIter, Store, Transaction as OxigraphTransaction};
use crate::io::{BomStrippingReader, map_rdf_parse_error};
use crate::persist::{self, DiskStore};
use crate::storage::{OpenOptions, StorageBackend, StorageCapabilities};
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())))
}
}
pub struct ModelTransaction<'a> {
inner: OxigraphTransaction<'a>,
}
impl ModelTransaction<'_> {
pub fn add(&mut self, statement: impl Into<Triple>) -> Result<bool> {
self.add_to_graph(statement, GraphName::DefaultGraph)
}
pub fn add_to_graph(
&mut 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(&mut self, quad: Quad) -> Result<bool> {
let inserted = !self
.inner
.contains(quad.as_ref())
.map_err(|error| Error::Storage(error.to_string()))?;
self.inner.insert(quad.as_ref());
Ok(inserted)
}
pub fn remove(&mut self, statement: impl Into<Triple>) -> Result<bool> {
self.remove_from_graph(statement, GraphName::DefaultGraph)
}
pub fn remove_from_graph(
&mut 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 remove_quad(&mut self, quad: &Quad) -> Result<bool> {
let removed = self
.inner
.contains(quad.as_ref())
.map_err(|error| Error::Storage(error.to_string()))?;
if removed {
self.inner.remove(quad.as_ref());
}
Ok(removed)
}
pub fn clear(&mut self) -> Result<()> {
self.inner
.clear()
.map_err(|error| Error::Storage(error.to_string()))
}
pub fn clear_graph(&mut self, graph_name: impl Into<GraphName>) -> Result<()> {
let graph_name = graph_name.into();
self.inner
.clear_graph(graph_name.as_ref())
.map_err(|error| Error::Storage(error.to_string()))
}
}
#[derive(Clone)]
pub struct Model {
store: Store,
disk: Option<DiskStore>,
lock: Arc<RwLock<()>>,
read_only: bool,
in_transaction: Arc<AtomicBool>,
txn_owner: Arc<Mutex<Option<ThreadId>>>,
}
struct InTransactionGuard<'a> {
model: &'a Model,
}
impl Drop for InTransactionGuard<'_> {
fn drop(&mut self) {
self.model.in_transaction.store(false, Ordering::Release);
*self
.model
.txn_owner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
}
}
impl Model {
pub fn new() -> Result<Self> {
Store::new()
.map(|store| Self {
store,
disk: None,
lock: Arc::new(RwLock::new(())),
read_only: false,
in_transaction: Arc::new(AtomicBool::new(false)),
txn_owner: Arc::new(Mutex::new(None)),
})
.map_err(|error| Error::Storage(error.to_string()))
}
pub fn open(path: impl AsRef<Path>) -> Result<Self> {
Self::open_with(OpenOptions::fjall(path))
}
pub fn open_with(options: OpenOptions) -> Result<Self> {
let path = options.path();
if options.is_read_only() {
if !path.exists() {
return Err(Error::OpenStore {
path: path.to_owned(),
message: "read-only open requires an existing store path".into(),
});
}
if !persist::looks_like_fjall_store(path) {
return Err(Error::OpenStore {
path: path.to_owned(),
message: "read-only open cannot initialize a new Oxiland store".into(),
});
}
}
let disk = DiskStore::open_with_create(path, options.can_create())?;
let allow_init = !options.is_read_only() && options.can_create();
disk.ensure_format_v1(path, allow_init)?;
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),
lock: Arc::new(RwLock::new(())),
read_only: options.is_read_only(),
in_transaction: Arc::new(AtomicBool::new(false)),
txn_owner: Arc::new(Mutex::new(None)),
})
}
pub fn migrate_legacy_store(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref();
{
let disk = DiskStore::open_with_create(path, false)?;
disk.migrate_legacy_to_v1()?;
}
Self::open_with(OpenOptions::fjall(path))
}
#[must_use]
pub fn capabilities(&self) -> StorageCapabilities {
match &self.disk {
None => StorageCapabilities::memory(),
Some(_) => StorageCapabilities::fjall(self.read_only),
}
}
#[must_use]
pub fn backend(&self) -> StorageBackend {
self.capabilities().backend
}
#[must_use]
pub fn store(&self) -> &Store {
&self.store
}
pub(crate) fn with_read_lock<R>(&self, f: impl FnOnce() -> R) -> R {
if self.same_thread_in_transaction() {
return f();
}
let _guard = self
.lock
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
f()
}
pub fn storage_backend_available(name: &str) -> Result<bool> {
StorageBackend::from_name(name).map(|_| true)
}
fn same_thread_in_transaction(&self) -> bool {
if !self.in_transaction.load(Ordering::Acquire) {
return false;
}
let owner = self
.txn_owner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
*owner == Some(std::thread::current().id())
}
fn ensure_writable(&self) -> Result<()> {
if self.read_only {
return Err(Error::Unsupported(
"model was opened read-only; mutating APIs are unavailable".into(),
));
}
if self.in_transaction.load(Ordering::Acquire) {
return Err(Error::Unsupported(
"auto-commit mutation is unavailable while a Model::transaction is open; use the transaction handle"
.into(),
));
}
Ok(())
}
pub fn transaction<R>(
&self,
f: impl FnOnce(&mut ModelTransaction<'_>) -> Result<R>,
) -> Result<R> {
if self.read_only {
return Err(Error::Unsupported(
"model was opened read-only; mutating APIs are unavailable".into(),
));
}
if self.same_thread_in_transaction() {
return Err(Error::Unsupported(
"nested Model::transaction is unsupported".into(),
));
}
let _guard = self
.lock
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self
.in_transaction
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return Err(Error::Unsupported(
"nested Model::transaction is unsupported".into(),
));
}
*self
.txn_owner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(std::thread::current().id());
let _txn_flag = InTransactionGuard { model: self };
let oxi = self
.store
.start_transaction()
.map_err(|error| Error::Storage(error.to_string()))?;
let mut tx = ModelTransaction { inner: oxi };
let value = f(&mut tx)?;
let ModelTransaction { inner } = tx;
inner
.commit()
.map_err(|error| Error::Storage(error.to_string()))?;
if let Some(disk) = &self.disk {
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 transaction ({error}); rollback from disk also failed ({reload_error})"
)));
}
return Err(Error::Storage(format!(
"durable sync failed after transaction; in-memory store rolled back to disk: {error}"
)));
}
}
Ok(value)
}
pub fn sync(&self) -> Result<()> {
if self.same_thread_in_transaction() {
return Err(Error::Unsupported(
"sync is unavailable while a Model::transaction is open on this thread".into(),
));
}
let _guard = self
.lock
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match &self.disk {
Some(disk) => disk.sync(),
None => Ok(()),
}
}
pub fn clear(&self) -> Result<()> {
self.ensure_writable()?;
let _guard = self
.lock
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.store
.clear()
.map_err(|error| Error::Storage(error.to_string()))?;
if let Some(disk) = &self.disk {
if let Err(error) = disk.clear_quads() {
if let Err(reload_error) = self.reload_store_from_disk_unlocked(disk) {
return Err(Error::Storage(format!(
"durable clear failed ({error}); rollback from disk also failed ({reload_error})"
)));
}
return Err(error);
}
}
Ok(())
}
pub fn clear_graph(&self, graph_name: impl Into<GraphName>) -> Result<()> {
self.ensure_writable()?;
let graph_name = graph_name.into();
let _guard = self
.lock
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.store
.clear_graph(graph_name.as_ref())
.map_err(|error| Error::Storage(error.to_string()))?;
if let Some(disk) = &self.disk {
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 clear_graph failed ({error}); rollback from disk also failed ({reload_error})"
)));
}
return Err(error);
}
}
Ok(())
}
pub fn bulk_insert_quads(&self, quads: impl IntoIterator<Item = Quad>) -> Result<usize> {
let quads: Vec<_> = quads.into_iter().collect();
let total = quads.len();
self.transaction(|tx| {
for quad in quads {
tx.insert_quad(quad)?;
}
Ok(total)
})
}
pub fn export_nquads_to_path(&self, path: impl AsRef<Path>) -> Result<()> {
let path = path.as_ref();
let file = File::create(path).map_err(|error| {
Error::Io(std::io::Error::new(
error.kind(),
format!("{}: {}", path.display(), error),
))
})?;
let mut serializer =
RdfSerializer::from_format(RdfFormat::NQuads).for_writer(BufWriter::new(file));
for item in self.find(StatementPattern::default()) {
let quad = item?;
serializer
.serialize_quad(QuadRef::from(&quad))
.map_err(Error::Io)?;
}
let writer = serializer.finish().map_err(Error::Io)?;
writer
.into_inner()
.map_err(|error| Error::Io(error.into_error()))?;
Ok(())
}
pub fn import_nquads_from_path(&self, path: impl AsRef<Path>) -> Result<usize> {
let path = path.as_ref();
let file = File::open(path).map_err(|error| {
Error::Io(std::io::Error::new(
error.kind(),
format!("{}: {}", path.display(), error),
))
})?;
let reader = BomStrippingReader::new(BufReader::new(file));
let quads = RdfParser::from_format(RdfFormat::NQuads)
.rename_blank_nodes()
.for_reader(reader)
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(map_rdf_parse_error)?;
let total = quads.len();
self.bulk_insert_quads(quads)?;
Ok(total)
}
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> {
self.ensure_writable()?;
let _guard = self
.lock
.write()
.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> {
self.ensure_writable()?;
let _guard = self
.lock
.write()
.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.with_read_lock(|| {
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.with_read_lock(|| {
self.store
.len()
.map_err(|error| Error::Storage(error.to_string()))
})
}
pub fn is_empty(&self) -> Result<bool> {
self.with_read_lock(|| {
self.store
.is_empty()
.map_err(|error| Error::Storage(error.to_string()))
})
}
pub fn find(&self, pattern: StatementPattern<'_>) -> StatementMatches {
self.with_read_lock(|| 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<()> {
self.ensure_writable()?;
let _guard = self
.lock
.write()
.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)
}
}