pub mod concrete;
pub(crate) mod dictionary;
pub mod persistence;
pub mod storage;
pub mod types;
#[cfg(feature = "async-tokio")]
pub mod async_store;
pub use concrete::*;
pub use persistence::{PersistentState, SyncPolicy};
pub use storage::*;
pub use types::*;
#[cfg(feature = "async-tokio")]
pub use async_store::AsyncRdfStore;
use crate::indexing::{IndexStats, MemoryUsage};
use crate::model::*;
use crate::parser::RdfFormat;
use crate::serializer::Serializer;
use crate::sparql::extract_and_expand_prefixes; use crate::{OxirsError, Result};
use async_trait::async_trait;
use std::collections::HashSet;
use std::path::Path;
use std::sync::{Arc, RwLock};
#[async_trait]
pub trait Store: Send + Sync {
fn insert_quad(&self, quad: Quad) -> Result<bool>;
fn remove_quad(&self, quad: &Quad) -> Result<bool>;
fn find_quads(
&self,
subject: Option<&Subject>,
predicate: Option<&Predicate>,
object: Option<&Object>,
graph_name: Option<&GraphName>,
) -> Result<Vec<Quad>>;
fn is_ready(&self) -> bool;
fn len(&self) -> Result<usize>;
fn is_empty(&self) -> Result<bool>;
fn query(&self, sparql: &str) -> Result<OxirsQueryResults>;
fn prepare_query(&self, sparql: &str) -> Result<PreparedQuery>;
fn insert_triple(&self, triple: Triple) -> Result<bool> {
let quad = Quad::from_triple(triple);
self.insert_quad(quad)
}
fn insert(&self, quad: &Quad) -> Result<()> {
self.insert_quad(quad.clone())?;
Ok(())
}
fn remove(&self, quad: &Quad) -> Result<bool> {
self.remove_quad(quad)
}
fn quads(&self) -> Result<Vec<Quad>> {
self.find_quads(None, None, None, None)
}
fn named_graphs(&self) -> Result<Vec<NamedNode>> {
Ok(Vec::new())
}
fn graphs(&self) -> Result<Vec<NamedNode>> {
self.named_graphs()
}
fn named_graph_quads(&self) -> Result<Vec<Quad>> {
let all_quads = self.quads()?;
Ok(all_quads
.into_iter()
.filter(|quad| matches!(quad.graph_name(), GraphName::NamedNode(_)))
.collect())
}
fn default_graph_quads(&self) -> Result<Vec<Quad>> {
let default_graph = GraphName::DefaultGraph;
self.find_quads(None, None, None, Some(&default_graph))
}
fn graph_quads(&self, graph: Option<&NamedNode>) -> Result<Vec<Quad>> {
let graph_name = graph
.map(|g| GraphName::NamedNode(g.clone()))
.unwrap_or(GraphName::DefaultGraph);
self.find_quads(None, None, None, Some(&graph_name))
}
fn clear_all(&self) -> Result<usize> {
Err(OxirsError::NotSupported(
"clear_all requires mutable access".to_string(),
))
}
fn clear_named_graphs(&self) -> Result<usize> {
Err(OxirsError::NotSupported(
"clear_named_graphs requires mutable access".to_string(),
))
}
fn clear_default_graph(&self) -> Result<usize> {
Err(OxirsError::NotSupported(
"clear_default_graph requires mutable access".to_string(),
))
}
fn clear_graph(&self, _graph: Option<&GraphName>) -> Result<usize> {
Err(OxirsError::NotSupported(
"clear_graph requires mutable access".to_string(),
))
}
fn create_graph(&self, _graph: Option<&NamedNode>) -> Result<()> {
Err(OxirsError::NotSupported(
"create_graph requires mutable access".to_string(),
))
}
fn drop_graph(&self, _graph: Option<&GraphName>) -> Result<()> {
Err(OxirsError::NotSupported(
"drop_graph requires mutable access".to_string(),
))
}
fn load_from_url(&self, _url: &str, _graph: Option<&NamedNode>) -> Result<usize> {
Err(OxirsError::NotSupported(
"load_from_url requires mutable access".to_string(),
))
}
fn triples(&self) -> Result<Vec<Triple>> {
let quads = self.find_quads(None, None, None, None)?;
Ok(quads
.into_iter()
.map(|quad| {
Triple::new(
quad.subject().clone(),
quad.predicate().clone(),
quad.object().clone(),
)
})
.collect())
}
fn bulk_insert_quads(&self, quads: Vec<Quad>) -> Result<usize> {
let mut inserted = 0usize;
for quad in quads {
if self.insert_quad(quad)? {
inserted += 1;
}
}
Ok(inserted)
}
fn for_each_quad(
&self,
subject: Option<&Subject>,
predicate: Option<&Predicate>,
object: Option<&Object>,
graph_name: Option<&GraphName>,
f: &mut dyn FnMut(Quad),
) -> Result<()> {
for quad in self.find_quads(subject, predicate, object, graph_name)? {
f(quad);
}
Ok(())
}
}
pub struct PreparedQuery {
sparql: String,
backend: Option<StorageBackend>,
}
impl PreparedQuery {
pub fn new(sparql: String) -> Self {
Self {
sparql,
backend: None,
}
}
pub fn with_backend(sparql: String, backend: StorageBackend) -> Self {
Self {
sparql,
backend: Some(backend),
}
}
pub fn exec(&self) -> Result<QueryResultsIterator> {
let backend = self.backend.as_ref().ok_or_else(|| {
OxirsError::Query(
"PreparedQuery has no bound store backend; construct it via Store::prepare_query"
.to_string(),
)
})?;
let executor = crate::sparql::QueryExecutor::new(backend);
let results = executor.execute(&self.sparql)?;
Ok(QueryResultsIterator::from_results(&results))
}
}
pub struct QueryResultsIterator {
results: Vec<SolutionMapping>,
index: usize,
}
impl QueryResultsIterator {
pub fn empty() -> Self {
Self {
results: Vec::new(),
index: 0,
}
}
pub fn from_results(results: &OxirsQueryResults) -> Self {
let rows = match results.results() {
QueryResults::Bindings(bindings) => bindings
.iter()
.map(|binding| SolutionMapping::from_bindings(binding.bindings.clone()))
.collect(),
QueryResults::Boolean(_) | QueryResults::Graph(_) => Vec::new(),
};
Self {
results: rows,
index: 0,
}
}
}
impl Iterator for QueryResultsIterator {
type Item = SolutionMapping;
fn next(&mut self) -> Option<Self::Item> {
if self.index < self.results.len() {
let result = self.results[self.index].clone();
self.index += 1;
Some(result)
} else {
None
}
}
}
#[derive(Debug, Clone, Default)]
pub struct SolutionMapping {
bindings: std::collections::HashMap<String, Term>,
}
impl SolutionMapping {
pub fn new() -> Self {
Self {
bindings: std::collections::HashMap::new(),
}
}
pub fn from_bindings(bindings: std::collections::HashMap<String, Term>) -> Self {
Self { bindings }
}
pub fn iter(&self) -> impl Iterator<Item = (&String, &Term)> {
self.bindings.iter()
}
}
#[derive(Debug)]
pub struct RdfStore {
backend: StorageBackend,
}
impl RdfStore {
pub fn new() -> Result<Self> {
Ok(RdfStore {
backend: StorageBackend::Memory(Arc::new(RwLock::new(MemoryStorage::new()))),
})
}
pub fn new_legacy() -> Result<Self> {
Ok(RdfStore {
backend: StorageBackend::Memory(Arc::new(RwLock::new(MemoryStorage::new()))),
})
}
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
Self::open_with_sync_policy(path, SyncPolicy::default())
}
pub fn open_with_sync_policy<P: AsRef<Path>>(path: P, sync_policy: SyncPolicy) -> Result<Self> {
let path_buf = path.as_ref().to_path_buf();
let data_file = path_buf.join("data.nq");
let (storage, load_had_errors) = if data_file.exists() {
persistence::load_from_disk(&data_file)?
} else {
(MemoryStorage::new(), false)
};
let state = PersistentState::open(path_buf, sync_policy, load_had_errors)?;
Ok(RdfStore {
backend: StorageBackend::Persistent(Arc::new(RwLock::new(storage)), Arc::new(state)),
})
}
fn quad_to_nquads_line(quad: &Quad) -> Result<String> {
Serializer::new(RdfFormat::NQuads).serialize_quad_to_nquads(quad)
}
pub fn flush(&self) -> Result<()> {
if let StorageBackend::Persistent(storage, state) = &self.backend {
let guard = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
state.flush(&guard)?;
}
Ok(())
}
pub fn insert_quad(&mut self, quad: Quad) -> Result<bool> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let existing = index.find_quads(
Some(quad.subject()),
Some(quad.predicate()),
Some(quad.object()),
Some(quad.graph_name()),
);
if !existing.is_empty() {
return Ok(false); }
let _id = index.insert_quad(&quad);
Ok(true) }
StorageBackend::Memory(storage) => {
let mut storage = storage
.write()
.map_err(|e| OxirsError::Store(format!("Failed to acquire write lock: {e}")))?;
Ok(storage.insert_quad(quad))
}
StorageBackend::Persistent(storage, state) => {
let is_new = {
let mut storage = storage.write().map_err(|e| {
OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.insert_quad(quad.clone())
};
if is_new {
let line = Self::quad_to_nquads_line(&quad)?;
state.append_line(&line)?;
}
Ok(is_new)
}
}
}
pub fn bulk_insert_quads(&mut self, quads: Vec<Quad>) -> Result<Vec<u64>> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let ids = index.bulk_insert_quads(&quads);
Ok(ids)
}
StorageBackend::Memory(storage) => {
let mut guard = storage
.write()
.map_err(|e| OxirsError::Store(format!("Failed to acquire write lock: {e}")))?;
let mut ids = Vec::with_capacity(quads.len());
for quad in quads {
ids.push(u64::from(guard.insert_quad(quad)));
}
Ok(ids)
}
StorageBackend::Persistent(storage, state) => {
let mut ids = Vec::with_capacity(quads.len());
let mut lines = Vec::new();
{
let mut guard = storage.write().map_err(|e| {
OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
for quad in quads {
let is_new = guard.insert_quad(quad.clone());
if is_new {
lines.push(Self::quad_to_nquads_line(&quad)?);
}
ids.push(u64::from(is_new));
}
}
state.append_lines(&lines)?;
Ok(ids)
}
}
}
pub fn insert_triple(&mut self, triple: Triple) -> Result<bool> {
self.insert_quad(Quad::from_triple(triple))
}
pub fn insert_string_triple(
&mut self,
subject: &str,
predicate: &str,
object: &str,
) -> Result<bool> {
let subject_node = NamedNode::new(subject)?;
let predicate_node = NamedNode::new(predicate)?;
let object_literal = Literal::new(object);
let triple = Triple::new(subject_node, predicate_node, object_literal);
self.insert_triple(triple)
}
pub fn remove_quad(&mut self, quad: &Quad) -> Result<bool> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => Ok(index.remove_quad(quad)),
StorageBackend::Memory(storage) => {
let mut storage = storage
.write()
.map_err(|e| OxirsError::Store(format!("Failed to acquire write lock: {e}")))?;
Ok(storage.remove_quad(quad))
}
StorageBackend::Persistent(storage, state) => {
let removed = {
let mut storage = storage.write().map_err(|e| {
OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.remove_quad(quad)
};
if removed {
state.mark_dirty();
}
Ok(removed)
}
}
}
pub fn contains_quad(&self, quad: &Quad) -> Result<bool> {
match &self.backend {
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let storage = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
Ok(storage.contains_quad(quad))
}
StorageBackend::UltraMemory(index, _) => {
let results = index.find_quads(
Some(quad.subject()),
Some(quad.predicate()),
Some(quad.object()),
Some(quad.graph_name()),
);
Ok(!results.is_empty())
}
}
}
pub fn query_quads(
&self,
subject: Option<&Subject>,
predicate: Option<&Predicate>,
object: Option<&Object>,
graph_name: Option<&GraphName>,
) -> Result<Vec<Quad>> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let results = index.find_quads(subject, predicate, object, graph_name);
Ok(results)
}
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let storage = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
Ok(storage.query_quads(subject, predicate, object, graph_name))
}
}
}
pub fn query_triples(
&self,
subject: Option<&Subject>,
predicate: Option<&Predicate>,
object: Option<&Object>,
) -> Result<Vec<Triple>> {
let default_graph = GraphName::DefaultGraph;
let quads = self.query_quads(subject, predicate, object, Some(&default_graph))?;
Ok(quads.into_iter().map(|quad| quad.to_triple()).collect())
}
pub fn iter_quads(&self) -> Result<Vec<Quad>> {
self.query_quads(None, None, None, None)
}
pub fn triples(&self) -> Result<Vec<Triple>> {
self.query_triples(None, None, None)
}
pub fn len(&self) -> Result<usize> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => Ok(index.len()),
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let storage = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
Ok(storage.len())
}
}
}
pub fn is_empty(&self) -> Result<bool> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => Ok(index.is_empty()),
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let storage = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
Ok(storage.is_empty())
}
}
}
pub fn stats(&self) -> Option<Arc<IndexStats>> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => Some(index.stats()),
_ => None,
}
}
pub fn memory_usage(&self) -> Option<MemoryUsage> {
match &self.backend {
StorageBackend::UltraMemory(index, arena) => {
let mut usage = index.memory_usage();
usage.arena_bytes = arena.allocated_bytes();
Some(usage)
}
_ => None,
}
}
pub fn clear_arena(&self) {
if let StorageBackend::UltraMemory(index, _arena) = &self.backend {
index.clear_arena();
}
}
pub fn clear(&mut self) -> Result<()> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
index.clear();
Ok(())
}
StorageBackend::Memory(storage) => {
let mut storage = storage
.write()
.map_err(|e| OxirsError::Store(format!("Failed to acquire write lock: {e}")))?;
*storage = MemoryStorage::new();
Ok(())
}
StorageBackend::Persistent(storage, state) => {
{
let mut storage = storage.write().map_err(|e| {
OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
*storage = MemoryStorage::new();
}
state.mark_dirty();
Ok(())
}
}
}
#[allow(dead_code)]
fn extract_and_expand_prefixes(
&self,
sparql: &str,
) -> Result<(std::collections::HashMap<String, String>, String)> {
extract_and_expand_prefixes(sparql)
}
pub fn query(&self, sparql: &str) -> Result<OxirsQueryResults> {
let executor = crate::sparql::QueryExecutor::new(&self.backend);
executor.execute(sparql)
}
pub fn insert(&mut self, quad: &Quad) -> Result<()> {
self.insert_quad(quad.clone())?;
Ok(())
}
pub fn remove(&mut self, quad: &Quad) -> Result<bool> {
self.remove_quad(quad)
}
pub fn quads(&self) -> Result<Vec<Quad>> {
self.iter_quads()
}
pub fn named_graph_quads(&self) -> Result<Vec<Quad>> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let mut result = Vec::new();
let graphs = self.named_graphs()?;
for graph in graphs {
let graph_name = GraphName::NamedNode(graph);
let quads = index.find_quads(None, None, None, Some(&graph_name));
result.extend(quads);
}
Ok(result)
}
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let storage = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
let mut result = Vec::new();
for graph in &storage.named_graphs {
let graph_name = GraphName::NamedNode(graph.clone());
let quads = storage.query_quads(None, None, None, Some(&graph_name));
result.extend(quads);
}
Ok(result)
}
}
}
pub fn default_graph_quads(&self) -> Result<Vec<Quad>> {
let default_graph = GraphName::DefaultGraph;
self.query_quads(None, None, None, Some(&default_graph))
}
pub fn graph_quads(&self, graph: Option<&NamedNode>) -> Result<Vec<Quad>> {
let graph_name = graph
.map(|g| GraphName::NamedNode(g.clone()))
.unwrap_or(GraphName::DefaultGraph);
self.query_quads(None, None, None, Some(&graph_name))
}
pub fn clear_all(&mut self) -> Result<usize> {
let count = self.len()?;
self.clear()?;
Ok(count)
}
pub fn clear_named_graphs(&mut self) -> Result<usize> {
let mut deleted = 0;
let graphs = self.named_graphs()?;
for graph in graphs {
let graph_name = GraphName::NamedNode(graph);
deleted += self.clear_graph(Some(&graph_name))?;
}
Ok(deleted)
}
pub fn clear_default_graph(&mut self) -> Result<usize> {
self.clear_graph(None)
}
pub fn clear_graph(&mut self, graph: Option<&GraphName>) -> Result<usize> {
let graph_name = graph.cloned().unwrap_or(GraphName::DefaultGraph);
let quads = self.query_quads(None, None, None, Some(&graph_name))?;
let count = quads.len();
for quad in quads {
self.remove_quad(&quad)?;
}
Ok(count)
}
pub fn graphs(&self) -> Result<Vec<NamedNode>> {
let mut graphs = self.named_graphs()?;
let default_graph = GraphName::DefaultGraph;
let default_quads = self.query_quads(None, None, None, Some(&default_graph))?;
if !default_quads.is_empty() {
if let Ok(default_marker) = NamedNode::new("urn:x-oxirs:default-graph") {
graphs.push(default_marker);
}
}
Ok(graphs)
}
pub fn named_graphs(&self) -> Result<Vec<NamedNode>> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let mut graphs = HashSet::new();
let all_quads = index.find_quads(None, None, None, None);
for quad in all_quads {
if let GraphName::NamedNode(graph) = quad.graph_name() {
graphs.insert(graph.clone());
}
}
Ok(graphs.into_iter().collect())
}
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let storage = storage
.read()
.map_err(|e| OxirsError::Store(format!("Failed to acquire read lock: {e}")))?;
Ok(storage.named_graphs.iter().cloned().collect())
}
}
}
pub fn create_graph(&mut self, graph: Option<&NamedNode>) -> Result<()> {
if let Some(graph_name) = graph {
match &self.backend {
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let mut storage = storage.write().map_err(|e| {
OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.named_graphs.insert(graph_name.clone());
}
StorageBackend::UltraMemory(_, _) => {
}
}
}
Ok(())
}
pub fn drop_graph(&mut self, graph: Option<&GraphName>) -> Result<()> {
self.clear_graph(graph)?;
if let Some(GraphName::NamedNode(graph_name)) = graph {
match &self.backend {
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let mut storage = storage.write().map_err(|e| {
OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.named_graphs.remove(graph_name);
}
StorageBackend::UltraMemory(_, _) => {
}
}
}
Ok(())
}
pub fn load_from_url(&mut self, url: &str, graph: Option<&NamedNode>) -> Result<usize> {
use crate::parser::Parser;
let url_path = url.split('?').next().unwrap_or(url);
let extension = url_path
.split('/')
.next_back()
.and_then(|filename| filename.rsplit('.').next());
let url_owned = url.to_string();
let fetch = async move {
let response = reqwest::get(&url_owned)
.await
.map_err(|e| OxirsError::Store(format!("Failed to fetch URL {url_owned}: {e}")))?;
if !response.status().is_success() {
return Err(OxirsError::Store(format!(
"HTTP error {} when fetching {url_owned}",
response.status()
)));
}
let content_type = response
.headers()
.get("content-type")
.and_then(|v| v.to_str().ok())
.map(|s| s.split(';').next().unwrap_or(s).trim().to_string());
let text = response
.text()
.await
.map_err(|e| OxirsError::Store(format!("Failed to read response body: {e}")))?;
Ok::<_, OxirsError>((text, content_type))
};
let (content, content_type) = match tokio::runtime::Handle::try_current() {
Ok(handle) => tokio::task::block_in_place(|| handle.block_on(fetch)),
Err(_) => {
let runtime = tokio::runtime::Runtime::new()
.map_err(|e| OxirsError::Store(format!("Failed to create runtime: {e}")))?;
runtime.block_on(fetch)
}
}?;
let format = Self::detect_format_from_url(&content_type, extension, &content)?;
let parser = Parser::new(format);
let quads = parser
.parse_str_to_quads(&content)
.map_err(|e| OxirsError::Store(format!("Failed to parse RDF data from {url}: {e}")))?;
let target_graph = graph.cloned().map(GraphName::NamedNode);
let mut inserted_count = 0;
for quad in quads {
let final_quad = if let Some(ref target) = target_graph {
Quad::new(
quad.subject().clone(),
quad.predicate().clone(),
quad.object().clone(),
target.clone(),
)
} else {
quad
};
if self.insert_quad(final_quad)? {
inserted_count += 1;
}
}
Ok(inserted_count)
}
fn detect_format_from_url(
content_type: &Option<String>,
extension: Option<&str>,
content: &str,
) -> Result<RdfFormat> {
if let Some(ct) = content_type {
let ct_lower = ct.to_lowercase();
if let Some(format) = Self::format_from_media_type(&ct_lower) {
return Ok(format);
}
}
if let Some(ext) = extension {
if let Some(format) = RdfFormat::from_extension(ext) {
return Ok(format);
}
}
if let Some(format) = crate::parser::detect_format_from_content(content) {
return Ok(format);
}
Err(OxirsError::Store(
"Could not detect RDF format from URL, content type, or content".to_string(),
))
}
fn format_from_media_type(media_type: &str) -> Option<RdfFormat> {
match media_type {
"text/turtle" | "application/x-turtle" => Some(RdfFormat::Turtle),
"application/n-triples" | "text/plain" => Some(RdfFormat::NTriples),
"application/trig" | "application/x-trig" => Some(RdfFormat::TriG),
"application/n-quads" | "text/x-nquads" => Some(RdfFormat::NQuads),
"application/rdf+xml" | "application/xml" | "text/xml" => Some(RdfFormat::RdfXml),
"application/ld+json" | "application/json" => Some(RdfFormat::JsonLd),
_ => None,
}
}
#[allow(dead_code)]
fn string_to_subject(s: &str) -> Option<Subject> {
if s.starts_with('<') && s.ends_with('>') {
let iri = &s[1..s.len() - 1];
NamedNode::new(iri).ok().map(Subject::NamedNode)
} else if let Some(blank_id) = s.strip_prefix("_:") {
BlankNode::new(blank_id).ok().map(Subject::BlankNode)
} else {
None
}
}
#[allow(dead_code)]
fn string_to_predicate(p: &str) -> Option<Predicate> {
if p.starts_with('<') && p.ends_with('>') {
let iri = &p[1..p.len() - 1];
NamedNode::new(iri).ok().map(Predicate::NamedNode)
} else {
None
}
}
#[allow(dead_code)]
fn string_to_object(o: &str) -> Option<Object> {
if o.starts_with('<') && o.ends_with('>') {
let iri = &o[1..o.len() - 1];
NamedNode::new(iri).ok().map(Object::NamedNode)
} else if let Some(blank_id) = o.strip_prefix("_:") {
BlankNode::new(blank_id).ok().map(Object::BlankNode)
} else if o.starts_with('"') {
let literal_content = &o[1..o.len() - 1];
Some(Object::Literal(Literal::new(literal_content)))
} else {
None
}
}
}
impl Default for RdfStore {
fn default() -> Self {
RdfStore {
backend: StorageBackend::Memory(Arc::new(RwLock::new(MemoryStorage::new()))),
}
}
}
impl Drop for RdfStore {
fn drop(&mut self) {
if let StorageBackend::Persistent(..) = &self.backend {
if let Err(e) = self.flush() {
tracing::error!("RdfStore flush on drop failed: {e}");
}
}
}
}
#[async_trait]
impl Store for RdfStore {
fn insert_quad(&self, quad: Quad) -> Result<bool> {
let inserted = match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let existing = index.find_quads(
Some(quad.subject()),
Some(quad.predicate()),
Some(quad.object()),
Some(quad.graph_name()),
);
if !existing.is_empty() {
return Ok(false); }
let _id = index.insert_quad(&quad);
true }
StorageBackend::Memory(storage) => {
let mut storage = storage.write().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.insert_quad(quad)
}
StorageBackend::Persistent(storage, state) => {
let is_new = {
let mut storage = storage.write().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.insert_quad(quad.clone())
};
if is_new {
let line = Self::quad_to_nquads_line(&quad)?;
state.append_line(&line)?;
}
is_new
}
};
Ok(inserted)
}
fn remove_quad(&self, quad: &Quad) -> Result<bool> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => Ok(index.remove_quad(quad)),
StorageBackend::Memory(storage) => {
let mut storage = storage.write().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
Ok(storage.remove_quad(quad))
}
StorageBackend::Persistent(storage, state) => {
let removed = {
let mut storage = storage.write().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
storage.remove_quad(quad)
};
if removed {
state.mark_dirty();
}
Ok(removed)
}
}
}
fn find_quads(
&self,
subject: Option<&Subject>,
predicate: Option<&Predicate>,
object: Option<&Object>,
graph_name: Option<&GraphName>,
) -> Result<Vec<Quad>> {
self.query_quads(subject, predicate, object, graph_name)
}
fn is_ready(&self) -> bool {
true }
fn len(&self) -> Result<usize> {
self.len()
}
fn is_empty(&self) -> Result<bool> {
self.is_empty()
}
fn query(&self, sparql: &str) -> Result<OxirsQueryResults> {
self.query(sparql)
}
fn prepare_query(&self, sparql: &str) -> Result<PreparedQuery> {
Ok(PreparedQuery::with_backend(
sparql.to_string(),
self.backend.clone(),
))
}
fn named_graphs(&self) -> Result<Vec<NamedNode>> {
self.named_graphs()
}
fn bulk_insert_quads(&self, quads: Vec<Quad>) -> Result<usize> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
let mut inserted = 0usize;
for quad in quads {
let existing = index.find_quads(
Some(quad.subject()),
Some(quad.predicate()),
Some(quad.object()),
Some(quad.graph_name()),
);
if existing.is_empty() {
let _id = index.insert_quad(&quad);
inserted += 1;
}
}
Ok(inserted)
}
StorageBackend::Memory(storage) => {
let mut guard = storage.write().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
let mut inserted = 0usize;
for quad in quads {
if guard.insert_quad(quad) {
inserted += 1;
}
}
Ok(inserted)
}
StorageBackend::Persistent(storage, state) => {
let mut lines = Vec::new();
{
let mut guard = storage.write().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire write lock: {e}"))
})?;
for quad in quads {
if guard.insert_quad(quad.clone()) {
lines.push(Self::quad_to_nquads_line(&quad)?);
}
}
}
let inserted = lines.len();
state.append_lines(&lines)?;
Ok(inserted)
}
}
}
fn for_each_quad(
&self,
subject: Option<&Subject>,
predicate: Option<&Predicate>,
object: Option<&Object>,
graph_name: Option<&GraphName>,
f: &mut dyn FnMut(Quad),
) -> Result<()> {
match &self.backend {
StorageBackend::UltraMemory(index, _arena) => {
for quad in index.find_quads(subject, predicate, object, graph_name) {
f(quad);
}
Ok(())
}
StorageBackend::Memory(storage) | StorageBackend::Persistent(storage, _) => {
let guard = storage.read().map_err(|e| {
crate::OxirsError::Store(format!("Failed to acquire read lock: {e}"))
})?;
guard.for_each_quad(subject, predicate, object, graph_name, f);
Ok(())
}
}
}
}
#[cfg(test)]
mod bulk_and_scan_tests {
use super::*;
use crate::model::{GraphName, Literal, NamedNode, Object, Quad};
fn quad(s: &str, o: &str, g: GraphName) -> Quad {
Quad::new(
NamedNode::new(s).expect("subject IRI"),
NamedNode::new("http://example.org/p").expect("predicate IRI"),
Literal::new_simple_literal(o),
g,
)
}
fn unique_dir(tag: &str) -> std::path::PathBuf {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let dir = std::env::temp_dir().join(format!("oxirs_bulk_scan_{tag}_{nanos}"));
std::fs::create_dir_all(&dir).expect("create temp dir");
dir
}
#[test]
fn rdfstore_trait_bulk_insert_persists_all_in_one_file() {
let dir = unique_dir("persist");
{
let store = RdfStore::open(&dir).expect("open persistent store");
let quads = vec![
quad("http://example.org/a", "1", GraphName::DefaultGraph),
quad("http://example.org/b", "2", GraphName::DefaultGraph),
quad(
"http://example.org/c",
"3",
GraphName::NamedNode(
NamedNode::new("http://example.org/g").expect("graph IRI"),
),
),
quad("http://example.org/a", "1", GraphName::DefaultGraph),
];
let inserted =
Store::bulk_insert_quads(&store, quads).expect("bulk insert through trait");
assert_eq!(inserted, 3, "only three distinct quads are new");
store.flush().expect("flush");
let data = std::fs::read_to_string(dir.join("data.nq")).expect("read data.nq");
let lines = data.lines().filter(|l| !l.trim().is_empty()).count();
assert_eq!(lines, 3, "one durable line per newly-inserted quad");
}
let reopened = RdfStore::open(&dir).expect("reopen persistent store");
assert_eq!(reopened.len().expect("len"), 3);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn persistent_delete_compacts_from_interned_state_and_reopens() {
let dir = unique_dir("delete_compact");
let g = GraphName::NamedNode(NamedNode::new("http://example.org/g").expect("graph IRI"));
let keep_a = quad("http://example.org/a", "1", GraphName::DefaultGraph);
let drop_b = quad("http://example.org/b", "2", GraphName::DefaultGraph);
let keep_c = quad("http://example.org/c", "3", g.clone());
{
let store = RdfStore::open(&dir).expect("open persistent store");
store.insert_quad(keep_a.clone()).expect("insert a");
store.insert_quad(drop_b.clone()).expect("insert b");
store.insert_quad(keep_c.clone()).expect("insert c");
store.flush().expect("flush appends");
assert_eq!(store.len().expect("len"), 3);
assert!(store.remove_quad(&drop_b).expect("remove b"));
assert_eq!(store.len().expect("len"), 2);
assert!(!store.contains_quad(&drop_b).expect("contains b"));
store.flush().expect("flush compaction");
let data = std::fs::read_to_string(dir.join("data.nq")).expect("read data.nq");
let lines = data.lines().filter(|l| !l.trim().is_empty()).count();
assert_eq!(lines, 2, "compacted file holds only the surviving quads");
}
let reopened = RdfStore::open(&dir).expect("reopen persistent store");
assert_eq!(reopened.len().expect("len"), 2);
assert!(reopened.contains_quad(&keep_a).expect("contains a"));
assert!(reopened.contains_quad(&keep_c).expect("contains c"));
assert!(!reopened.contains_quad(&drop_b).expect("not contains b"));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn for_each_quad_visits_exactly_matching_quads() {
use std::collections::BTreeSet;
let store = RdfStore::new().expect("in-memory store");
let g = GraphName::NamedNode(NamedNode::new("http://example.org/g").expect("graph IRI"));
Store::bulk_insert_quads(
&store,
vec![
quad("http://example.org/a", "1", GraphName::DefaultGraph),
quad("http://example.org/b", "2", GraphName::DefaultGraph),
quad("http://example.org/c", "3", g.clone()),
],
)
.expect("seed");
let mut visited: Vec<Quad> = Vec::new();
Store::for_each_quad(&store, None, None, None, None, &mut |q| visited.push(q))
.expect("scan all");
let via_find = store.find_quads(None, None, None, None).expect("find all");
let visited_set: BTreeSet<Quad> = visited.iter().cloned().collect();
let find_set: BTreeSet<Quad> = via_find.into_iter().collect();
assert_eq!(
visited_set, find_set,
"scan must match find_quads content (set-equal)"
);
assert_eq!(visited.len(), 3);
let mut named: Vec<Quad> = Vec::new();
Store::for_each_quad(&store, None, None, None, Some(&g), &mut |q| named.push(q))
.expect("scan named");
assert_eq!(named.len(), 1);
assert_eq!(named[0].graph_name(), &g);
let obj = Object::Literal(Literal::new_simple_literal("2"));
let mut by_obj: Vec<Quad> = Vec::new();
Store::for_each_quad(&store, None, None, Some(&obj), None, &mut |q| {
by_obj.push(q)
})
.expect("scan by object");
assert_eq!(by_obj.len(), 1);
assert_eq!(by_obj[0].object(), &obj);
}
}
#[cfg(test)]
mod named_graphs_trait_tests {
use super::*;
use crate::model::{GraphName, Literal, NamedNode, Quad};
fn quad(s: &str, o: &str, g: GraphName) -> Quad {
Quad::new(
NamedNode::new(s).expect("subject IRI"),
NamedNode::new("http://example.org/p").expect("predicate IRI"),
Literal::new_simple_literal(o),
g,
)
}
fn seed(store: &dyn Store) {
let g1 = GraphName::NamedNode(NamedNode::new("http://example.org/g1").expect("g1 IRI"));
let g2 = GraphName::NamedNode(NamedNode::new("http://example.org/g2").expect("g2 IRI"));
let g3 = GraphName::NamedNode(NamedNode::new("http://example.org/g3").expect("g3 IRI"));
store
.insert_quad(quad("http://example.org/a", "1", GraphName::DefaultGraph))
.expect("insert default graph quad 1");
store
.insert_quad(quad("http://example.org/b", "2", GraphName::DefaultGraph))
.expect("insert default graph quad 2");
store
.insert_quad(quad("http://example.org/c", "3", g1.clone()))
.expect("insert g1 quad 1");
store
.insert_quad(quad("http://example.org/d", "4", g1.clone()))
.expect("insert g1 quad 2");
store
.insert_quad(quad("http://example.org/e", "5", g2.clone()))
.expect("insert g2 quad");
store
.insert_quad(quad("http://example.org/f", "6", g3.clone()))
.expect("insert g3 quad");
}
fn assert_exactly_three_named_graphs(store: &dyn Store) {
let mut graphs: Vec<String> = store
.named_graphs()
.expect("named_graphs via Store trait")
.into_iter()
.map(|n| n.as_str().to_string())
.collect();
graphs.sort();
graphs.dedup();
assert_eq!(
graphs,
vec![
"http://example.org/g1".to_string(),
"http://example.org/g2".to_string(),
"http://example.org/g3".to_string(),
],
"named_graphs() through the Store trait must return exactly the \
named graphs (default graph excluded, no duplicates)"
);
}
#[test]
fn dyn_store_named_graphs_in_memory_excludes_default_and_dedups() {
let store = RdfStore::new().expect("in-memory store");
seed(&store);
let dyn_store: &dyn Store = &store;
assert_exactly_three_named_graphs(dyn_store);
}
#[test]
fn dyn_store_named_graphs_persistent_excludes_default_and_dedups() {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let dir = std::env::temp_dir().join(format!("oxirs_named_graphs_trait_{nanos}"));
std::fs::create_dir_all(&dir).expect("create temp dir");
{
let store = RdfStore::open(&dir).expect("open persistent store");
seed(&store);
let dyn_store: &dyn Store = &store;
assert_exactly_three_named_graphs(dyn_store);
}
let reopened = RdfStore::open(&dir).expect("reopen persistent store");
let dyn_reopened: &dyn Store = &reopened;
assert_exactly_three_named_graphs(dyn_reopened);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn dyn_store_named_graphs_concrete_store_delegates() {
let store = ConcreteStore::new().expect("concrete store");
seed(&store);
let dyn_store: &dyn Store = &store;
assert_exactly_three_named_graphs(dyn_store);
}
}