use std::fmt;
use std::io::Write;
use oxigraph::model::{GraphName, NamedOrBlankNode};
use oxigraph::sparql::results::{QueryResultsFormat, QueryResultsSerializer};
use oxigraph::sparql::{
CancellationToken, PreparedSparqlQuery, QuerySolutionIter, QueryTripleIter, SparqlEvaluator,
};
use spargebra::algebra::GraphPattern;
use spargebra::{Query as SparqlAlgebraQuery, SparqlParser};
use crate::io::Serializer;
use crate::{Error, Model, Result};
pub enum QueryResults<'a> {
Boolean(bool),
Solutions(QuerySolutionIter<'a>),
Graph(QueryTripleIter<'a>),
}
impl fmt::Debug for QueryResults<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Boolean(value) => f.debug_tuple("Boolean").field(value).finish(),
Self::Solutions(_) => f.write_str("Solutions(..)"),
Self::Graph(_) => f.write_str("Graph(..)"),
}
}
}
fn map_query_results(results: oxigraph::sparql::QueryResults<'_>) -> QueryResults<'_> {
match results {
oxigraph::sparql::QueryResults::Boolean(value) => QueryResults::Boolean(value),
oxigraph::sparql::QueryResults::Solutions(solutions) => QueryResults::Solutions(solutions),
oxigraph::sparql::QueryResults::Graph(graph) => QueryResults::Graph(graph),
}
}
#[derive(Clone, Debug, Default)]
struct DatasetConfig {
default_graphs: Option<Vec<GraphName>>,
default_as_union: bool,
named_graphs: Option<Vec<NamedOrBlankNode>>,
}
impl DatasetConfig {
fn is_configured(&self) -> bool {
self.default_as_union || self.default_graphs.is_some() || self.named_graphs.is_some()
}
}
#[derive(Clone)]
pub struct Query {
text: String,
base_iri: Option<String>,
prefixes: Vec<(String, String)>,
limit: Option<usize>,
offset: usize,
dataset: DatasetConfig,
cancellation: Option<CancellationToken>,
}
impl Query {
#[must_use]
pub fn new(text: impl Into<String>) -> Self {
Self {
text: text.into(),
base_iri: None,
prefixes: Vec::new(),
limit: None,
offset: 0,
dataset: DatasetConfig::default(),
cancellation: None,
}
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.text
}
pub fn base_iri(mut self, base_iri: impl Into<String>) -> Result<Self> {
let base_iri = base_iri.into();
let _ = SparqlParser::new()
.with_base_iri(&base_iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
self.base_iri = Some(base_iri);
Ok(self)
}
pub fn prefix(
mut self,
prefix_name: impl Into<String>,
prefix_iri: impl Into<String>,
) -> Result<Self> {
let prefix_name = prefix_name.into();
let prefix_iri = prefix_iri.into();
let _ = SparqlParser::new()
.with_prefix(&prefix_name, &prefix_iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
self.prefixes.push((prefix_name, prefix_iri));
Ok(self)
}
pub fn limit(mut self, limit: usize) -> Result<Self> {
self.ensure_not_ask_for_slice()?;
self.limit = Some(limit);
Ok(self)
}
pub fn offset(mut self, offset: usize) -> Result<Self> {
self.ensure_not_ask_for_slice()?;
self.offset = offset;
Ok(self)
}
#[must_use]
pub fn default_graph(mut self, graphs: impl IntoIterator<Item = GraphName>) -> Self {
self.dataset.default_graphs = Some(graphs.into_iter().collect());
self.dataset.default_as_union = false;
self
}
#[must_use]
pub fn default_graph_as_union(mut self) -> Self {
self.dataset.default_as_union = true;
self.dataset.default_graphs = None;
self
}
#[must_use]
pub fn available_named_graphs(
mut self,
graphs: impl IntoIterator<Item = NamedOrBlankNode>,
) -> Self {
self.dataset.named_graphs = Some(graphs.into_iter().collect());
self
}
#[must_use]
pub fn cancellation_token(mut self, token: CancellationToken) -> Self {
self.cancellation = Some(token);
self
}
pub fn execute<'a>(&self, model: &'a Model) -> Result<QueryResults<'a>> {
let prepared = self.prepare()?;
model.with_read_lock(|| {
let results = prepared
.on_store(model.store())
.execute()
.map_err(|error| Error::SparqlEvaluation(error.to_string()))?;
Ok(map_query_results(results))
})
}
fn ensure_not_ask_for_slice(&self) -> Result<()> {
if looks_like_ask_query(&self.text) {
return Err(Error::Unsupported(
"API-level limit/offset cannot be applied to ASK queries; put LIMIT in the query text if needed"
.into(),
));
}
Ok(())
}
fn prepare(&self) -> Result<PreparedSparqlQuery> {
let mut parser = SparqlParser::new();
if let Some(base_iri) = &self.base_iri {
parser = parser
.with_base_iri(base_iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
}
for (name, iri) in &self.prefixes {
parser = parser
.with_prefix(name, iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
}
let mut algebra = parser
.parse_query(&self.text)
.map_err(|error| Error::SparqlParse(error.to_string()))?;
algebra = apply_slice(algebra, self.offset, self.limit)?;
let mut evaluator = SparqlEvaluator::new();
if let Some(token) = &self.cancellation {
evaluator = evaluator.with_cancellation_token(token.clone());
}
let mut prepared = evaluator.for_query(algebra);
apply_dataset(prepared.dataset_mut(), &self.dataset);
Ok(prepared)
}
}
#[derive(Clone)]
pub struct Update {
text: String,
base_iri: Option<String>,
prefixes: Vec<(String, String)>,
dataset: DatasetConfig,
cancellation: Option<CancellationToken>,
}
impl Update {
#[must_use]
pub fn new(text: impl Into<String>) -> Self {
Self {
text: text.into(),
base_iri: None,
prefixes: Vec::new(),
dataset: DatasetConfig::default(),
cancellation: None,
}
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.text
}
pub fn base_iri(mut self, base_iri: impl Into<String>) -> Result<Self> {
let base_iri = base_iri.into();
let _ = SparqlParser::new()
.with_base_iri(&base_iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
self.base_iri = Some(base_iri);
Ok(self)
}
pub fn prefix(
mut self,
prefix_name: impl Into<String>,
prefix_iri: impl Into<String>,
) -> Result<Self> {
let prefix_name = prefix_name.into();
let prefix_iri = prefix_iri.into();
let _ = SparqlParser::new()
.with_prefix(&prefix_name, &prefix_iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
self.prefixes.push((prefix_name, prefix_iri));
Ok(self)
}
#[must_use]
pub fn default_graph(mut self, graphs: impl IntoIterator<Item = GraphName>) -> Self {
self.dataset.default_graphs = Some(graphs.into_iter().collect());
self.dataset.default_as_union = false;
self
}
#[must_use]
pub fn default_graph_as_union(mut self) -> Self {
self.dataset.default_as_union = true;
self.dataset.default_graphs = None;
self
}
#[must_use]
pub fn cancellation_token(mut self, token: CancellationToken) -> Self {
self.cancellation = Some(token);
self
}
pub fn execute(self, model: &Model) -> Result<()> {
let mut evaluator = SparqlEvaluator::new();
if let Some(base_iri) = &self.base_iri {
evaluator = evaluator
.with_base_iri(base_iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
}
for (name, iri) in &self.prefixes {
evaluator = evaluator
.with_prefix(name, iri)
.map_err(|error| Error::InvalidRdf(error.to_string()))?;
}
if let Some(token) = &self.cancellation {
evaluator = evaluator.with_cancellation_token(token.clone());
}
let mut prepared = evaluator
.parse_update(&self.text)
.map_err(|error| Error::SparqlParse(error.to_string()))?;
let mut applied_dataset = false;
for dataset in prepared.using_datasets_mut() {
apply_dataset(dataset, &self.dataset);
applied_dataset = true;
}
if self.dataset.is_configured() && !applied_dataset {
return Err(Error::Unsupported(
"Update dataset configuration requires DELETE/INSERT operations that expose USING datasets; INSERT DATA and DELETE DATA do not accept dataset overrides"
.into(),
));
}
model.run_sparql_update(|store| {
prepared
.on_store(store)
.execute()
.map_err(|error| Error::SparqlEvaluation(error.to_string()))
})
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
pub enum ResultsFormat {
Xml,
Json,
Csv,
Tsv,
}
impl ResultsFormat {
#[must_use]
pub const fn name(self) -> &'static str {
match self {
Self::Xml => "xml",
Self::Json => "json",
Self::Csv => "csv",
Self::Tsv => "tsv",
}
}
#[must_use]
pub const fn media_type(self) -> &'static str {
match self {
Self::Xml => "application/sparql-results+xml",
Self::Json => "application/sparql-results+json",
Self::Csv => "text/csv",
Self::Tsv => "text/tab-separated-values",
}
}
pub fn from_name(name: &str) -> Result<Self> {
match name.trim().to_ascii_lowercase().as_str() {
"xml" | "sparql-results+xml" => Ok(Self::Xml),
"json" | "sparql-results+json" => Ok(Self::Json),
"csv" => Ok(Self::Csv),
"tsv" | "tab-separated-values" => Ok(Self::Tsv),
other => Err(Error::Unsupported(format!(
"SPARQL results format '{other}' is not supported"
))),
}
}
pub fn from_media_type(media_type: &str) -> Result<Self> {
let base = media_type
.split(';')
.next()
.unwrap_or(media_type)
.trim()
.to_ascii_lowercase();
match base.as_str() {
"application/sparql-results+xml" => Ok(Self::Xml),
"application/sparql-results+json" => Ok(Self::Json),
"text/csv" => Ok(Self::Csv),
"text/tab-separated-values" | "text/tsv" => Ok(Self::Tsv),
other => Err(Error::Unsupported(format!(
"SPARQL results media type '{other}' is not supported"
))),
}
}
fn to_oxigraph(self) -> QueryResultsFormat {
match self {
Self::Xml => QueryResultsFormat::Xml,
Self::Json => QueryResultsFormat::Json,
Self::Csv => QueryResultsFormat::Csv,
Self::Tsv => QueryResultsFormat::Tsv,
}
}
}
pub fn serialize_query_results_to_writer<W: Write>(
results: QueryResults<'_>,
format: ResultsFormat,
writer: W,
) -> Result<W> {
let serializer = QueryResultsSerializer::from_format(format.to_oxigraph());
match results {
QueryResults::Boolean(value) => serializer
.serialize_boolean_to_writer(writer, value)
.map_err(|error| Error::Serialize(error.to_string())),
QueryResults::Solutions(mut solutions) => {
let variables = solutions.variables().to_vec();
let mut solutions_writer = serializer
.serialize_solutions_to_writer(writer, variables)
.map_err(|error| Error::Serialize(error.to_string()))?;
for solution in solutions.by_ref() {
let solution =
solution.map_err(|error| Error::SparqlEvaluation(error.to_string()))?;
solutions_writer
.serialize(&solution)
.map_err(|error| Error::Serialize(error.to_string()))?;
}
solutions_writer
.finish()
.map_err(|error| Error::Serialize(error.to_string()))
}
QueryResults::Graph(_) => Err(Error::Unsupported(
"graph query results must be serialized with serialize_graph_results_to_writer / oxiland::io::Serializer, not ResultsFormat"
.into(),
)),
}
}
pub fn serialize_query_results_to_string(
results: QueryResults<'_>,
format: ResultsFormat,
) -> Result<String> {
let buffer = serialize_query_results_to_writer(results, format, Vec::new())?;
String::from_utf8(buffer).map_err(|error| Error::Serialize(error.to_string()))
}
pub fn serialize_graph_results_to_writer<W: Write>(
results: QueryResults<'_>,
serializer: &Serializer,
writer: W,
) -> Result<W> {
let QueryResults::Graph(graph) = results else {
return Err(Error::Unsupported(
"serialize_graph_results_to_writer requires QueryResults::Graph".into(),
));
};
serializer.serialize_triples_fallible_to_writer(
writer,
graph.map(|triple| triple.map_err(|error| Error::SparqlEvaluation(error.to_string()))),
)
}
fn looks_like_ask_query(text: &str) -> bool {
let trimmed = strip_sparql_prologue(text).trim_start();
let Some(rest) = trimmed.get(..3) else {
return false;
};
if !rest.eq_ignore_ascii_case("ask") {
return false;
}
match trimmed.as_bytes().get(3) {
None => true,
Some(b) => b.is_ascii_whitespace() || *b == b'{',
}
}
fn strip_sparql_prologue(text: &str) -> &str {
let text = text.strip_prefix('\u{feff}').unwrap_or(text);
let bytes = text.as_bytes();
let mut i = 0;
while i < bytes.len() {
while i < bytes.len() && bytes[i].is_ascii_whitespace() {
i += 1;
}
if i >= bytes.len() {
break;
}
if bytes[i] == b'#' {
while i < bytes.len() && bytes[i] != b'\n' {
i += 1;
}
continue;
}
let rest = &text[i..];
if rest.len() >= 4
&& rest.is_char_boundary(4)
&& rest[..4].eq_ignore_ascii_case("base")
&& rest
.as_bytes()
.get(4)
.is_some_and(|b| b.is_ascii_whitespace() || *b == b'<')
{
if let Some(rel) = rest.find('>') {
i += rel + 1;
continue;
}
break;
}
if rest.len() >= 6
&& rest.is_char_boundary(6)
&& rest[..6].eq_ignore_ascii_case("prefix")
&& rest
.as_bytes()
.get(6)
.is_some_and(|b| b.is_ascii_whitespace())
{
if let Some(rel) = rest.find('>') {
i += rel + 1;
continue;
}
break;
}
break;
}
&text[i..]
}
fn apply_slice(
mut query: SparqlAlgebraQuery,
offset: usize,
limit: Option<usize>,
) -> Result<SparqlAlgebraQuery> {
if offset == 0 && limit.is_none() {
return Ok(query);
}
match &mut query {
SparqlAlgebraQuery::Select { pattern, .. }
| SparqlAlgebraQuery::Construct { pattern, .. }
| SparqlAlgebraQuery::Describe { pattern, .. } => {
let inner = strip_slices(std::mem::replace(
pattern,
GraphPattern::Bgp {
patterns: Vec::new(),
},
));
*pattern = GraphPattern::Slice {
inner: Box::new(inner),
start: offset,
length: limit,
};
Ok(query)
}
SparqlAlgebraQuery::Ask { .. } => Err(Error::Unsupported(
"API-level limit/offset cannot be applied to ASK queries; put LIMIT in the query text if needed"
.into(),
)),
}
}
fn strip_slices(mut pattern: GraphPattern) -> GraphPattern {
while let GraphPattern::Slice { inner, .. } = pattern {
pattern = *inner;
}
pattern
}
fn apply_dataset(dataset: &mut oxigraph::sparql::QueryDataset, config: &DatasetConfig) {
if config.default_as_union {
dataset.set_default_graph_as_union();
} else if let Some(graphs) = &config.default_graphs {
dataset.set_default_graph(graphs.clone());
}
if let Some(named) = &config.named_graphs {
dataset.set_available_named_graphs(named.clone());
}
}