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::storage::{
self, DurableStore, DurableStoreOps, OpenOptions, StorageBackend, StorageCapabilities,
StorageFacade,
};
use crate::world::{BridgeToken, FeatureMap, FeatureValue, World};
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: Option<QuadIter<'static>>,
prelude_error: Option<Error>,
}
impl StatementMatches {
fn failed(error: Error) -> Self {
Self {
inner: None,
prelude_error: Some(error),
}
}
}
impl Iterator for StatementMatches {
type Item = Result<Quad>;
fn next(&mut self) -> Option<Self::Item> {
if let Some(error) = self.prelude_error.take() {
return Some(Err(error));
}
self.inner
.as_mut()?
.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<DurableStore>,
lock: Arc<RwLock<()>>,
pending_inserts: Arc<Mutex<Vec<Quad>>>,
read_only: bool,
in_transaction: Arc<AtomicBool>,
txn_owner: Arc<Mutex<Option<ThreadId>>>,
world: World,
feature_map: FeatureMap,
storage_features: FeatureMap,
storage_instance: Arc<RwLock<Option<BridgeToken>>>,
}
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 Drop for Model {
fn drop(&mut self) {
let _ = self.flush_pending_inserts();
}
}
impl Model {
pub fn new() -> Result<Self> {
Self::with_world(World::new())
}
pub fn with_world(world: World) -> Result<Self> {
Store::new()
.map(|store| Self::from_parts(store, None, false, world))
.map_err(|error| Error::Storage(error.to_string()))
}
fn from_parts(store: Store, disk: Option<DurableStore>, read_only: bool, world: World) -> Self {
Self {
store,
disk,
lock: Arc::new(RwLock::new(())),
pending_inserts: Arc::new(Mutex::new(Vec::new())),
read_only,
in_transaction: Arc::new(AtomicBool::new(false)),
txn_owner: Arc::new(Mutex::new(None)),
world,
feature_map: FeatureMap::new(),
storage_features: FeatureMap::new(),
storage_instance: Arc::new(RwLock::new(None)),
}
}
#[must_use]
pub fn world(&self) -> &World {
&self.world
}
pub fn set_feature(&self, iri: impl Into<String>, value: FeatureValue) {
self.feature_map.set(iri, value);
}
#[must_use]
pub fn feature(&self, iri: &str) -> Option<FeatureValue> {
self.feature_map.get(iri)
}
pub fn set_storage_feature(&self, iri: impl Into<String>, value: FeatureValue) {
self.storage_features.set(iri, value);
}
#[must_use]
pub fn storage_feature(&self, iri: &str) -> Option<FeatureValue> {
self.storage_features.get(iri)
}
#[must_use]
pub fn storage_instance(&self) -> Option<BridgeToken> {
*self
.storage_instance
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
pub fn set_storage_instance(&self, token: Option<BridgeToken>) {
*self
.storage_instance
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = token;
}
#[must_use]
pub fn as_storage(&self) -> StorageFacade<'_> {
StorageFacade::new(self)
}
pub fn open(path: impl AsRef<Path>) -> Result<Self> {
Self::open_with(OpenOptions::fjall(path))
}
pub fn open_with(options: OpenOptions) -> Result<Self> {
if options.backend() == StorageBackend::Memory {
return Store::new()
.map(|store| Self::from_parts(store, None, options.is_read_only(), World::new()))
.map_err(|error| Error::Storage(error.to_string()));
}
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 !DurableStore::looks_like_store(options.backend(), path) {
return Err(Error::OpenStore {
path: path.to_owned(),
message: "read-only open cannot initialize a new Oxiland store".into(),
});
}
}
let disk = DurableStore::open(options.backend(), 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::from_parts(
store,
Some(disk),
options.is_read_only(),
World::new(),
))
}
pub fn copy_to(&self, options: OpenOptions) -> Result<Self> {
if options.backend() == StorageBackend::Memory {
let dest = Self::new()?;
for quad in self.find(StatementPattern::default()) {
dest.insert_quad(quad?)?;
}
return Ok(dest);
}
if options.is_read_only() {
return Err(Error::Unsupported(
"Model::copy_to requires a writable destination".into(),
));
}
if !options.can_create() {
return Err(Error::Unsupported(
"Model::copy_to requires OpenOptions::create(true)".into(),
));
}
let path = options.path();
if DurableStore::looks_like_store(options.backend(), path) {
return Err(Error::OpenStore {
path: path.to_owned(),
message:
"copy_to destination already looks like an Oxiland store; refuse to overwrite"
.into(),
});
}
let dest = Self::open_with(options)?;
for quad in self.find(StatementPattern::default()) {
dest.insert_quad(quad?)?;
}
dest.sync()?;
Ok(dest)
}
pub fn migrate_legacy_store(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref();
{
let disk = DurableStore::open(StorageBackend::Fjall, 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::for_backend(StorageBackend::Memory, self.read_only),
Some(disk) => disk.capabilities(self.read_only),
}
}
#[must_use]
pub fn backend(&self) -> StorageBackend {
match &self.disk {
None => StorageBackend::Memory,
Some(disk) => disk.backend_id(),
}
}
#[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(crate) fn flush_pending_inserts(&self) -> Result<()> {
let batch = {
let mut pending = self
.pending_inserts
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if pending.is_empty() {
return Ok(());
}
std::mem::take(&mut *pending)
};
let _guard = self
.lock
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let result = (|| {
let mut transaction = self
.store
.start_transaction()
.map_err(|error| Error::Storage(error.to_string()))?;
for quad in &batch {
transaction.insert(quad);
}
transaction
.commit()
.map_err(|error| Error::Storage(error.to_string()))
})();
if result.is_err() {
self.pending_inserts
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.splice(0..0, batch);
}
result
}
pub fn storage_backend_available(name: &str) -> Result<bool> {
let backend = StorageBackend::from_name(name)?;
Ok(storage::compiled_backends().contains(&backend))
}
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(),
));
}
self.flush_pending_inserts()?;
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(),
));
}
self.flush_pending_inserts()?;
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()?;
self.flush_pending_inserts()?;
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()?;
self.flush_pending_inserts()?;
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()?;
self.flush_pending_inserts()?;
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 = storage::stored_matching_quad(&self.store, &quad)?;
if let Err(error) = disk.insert_quad(&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 insert_quad_unchecked(&self, quad: Quad) -> Result<()> {
self.ensure_writable()?;
if self.disk.is_none() {
let mut pending = self
.pending_inserts
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if pending.capacity() == 0 {
pending.reserve(1_024);
}
pending.push(quad);
drop(pending);
return Ok(());
}
self.flush_pending_inserts()?;
let _guard = self
.lock
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.store
.insert(&quad)
.map_err(|error| Error::Storage(error.to_string()))?;
if let Some(disk) = &self.disk {
let canonical = storage::stored_matching_quad(&self.store, &quad)?;
if let Err(error) = disk.insert_quad(&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(())
}
pub fn remove_quad(&self, quad: &Quad) -> Result<bool> {
self.ensure_writable()?;
self.flush_pending_inserts()?;
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 = storage::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.flush_pending_inserts()?;
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.flush_pending_inserts()?;
self.with_read_lock(|| {
self.store
.len()
.map_err(|error| Error::Storage(error.to_string()))
})
}
pub fn is_empty(&self) -> Result<bool> {
self.flush_pending_inserts()?;
self.with_read_lock(|| {
self.store
.is_empty()
.map_err(|error| Error::Storage(error.to_string()))
})
}
pub fn find(&self, pattern: StatementPattern<'_>) -> StatementMatches {
if let Err(error) = self.flush_pending_inserts() {
return StatementMatches::failed(error);
}
self.with_read_lock(|| StatementMatches {
inner: Some(self.store.quads_for_pattern(
pattern.subject,
pattern.predicate,
pattern.object,
pattern.graph_name,
)),
prelude_error: None,
})
}
pub(crate) fn run_sparql_update(
&self,
update: impl FnOnce(&Store) -> Result<()>,
) -> Result<()> {
self.ensure_writable()?;
self.flush_pending_inserts()?;
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: &DurableStore) -> Result<()> {
self.store
.clear()
.map_err(|error| Error::Storage(error.to_string()))?;
disk.load_into(&self.store)
}
}