use crate::{
auth::{AuthUser, Permission},
error::{FusekiError, FusekiResult},
federated_query_optimizer::FederatedQueryOptimizer,
server::AppState,
};
use axum::{
extract::{Query, State},
http::{header::ACCEPT, HeaderMap, StatusCode},
response::{IntoResponse, Json, Response},
Form,
};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tracing::{error, instrument, warn};
use oxirs_arq::query::parse_query;
fn deserialize_string_or_seq<'de, D>(de: D) -> Result<Option<Vec<String>>, D::Error>
where
D: serde::Deserializer<'de>,
{
use serde::de::{self, SeqAccess, Visitor};
use std::fmt;
struct V;
impl<'de> Visitor<'de> for V {
type Value = Option<Vec<String>>;
fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("a string or a sequence of strings")
}
fn visit_str<E: de::Error>(self, v: &str) -> Result<Self::Value, E> {
Ok(Some(vec![v.to_string()]))
}
fn visit_string<E: de::Error>(self, v: String) -> Result<Self::Value, E> {
Ok(Some(vec![v]))
}
fn visit_none<E: de::Error>(self) -> Result<Self::Value, E> {
Ok(None)
}
fn visit_unit<E: de::Error>(self) -> Result<Self::Value, E> {
Ok(None)
}
fn visit_some<D: serde::Deserializer<'de>>(self, d: D) -> Result<Self::Value, D::Error> {
d.deserialize_any(V)
}
fn visit_seq<S: SeqAccess<'de>>(self, mut seq: S) -> Result<Self::Value, S::Error> {
let mut out: Vec<String> = Vec::new();
while let Some(item) = seq.next_element::<String>()? {
out.push(item);
}
Ok(if out.is_empty() { None } else { Some(out) })
}
}
de.deserialize_any(V)
}
#[derive(Debug, Deserialize)]
pub struct SparqlQueryParams {
pub query: Option<String>,
#[serde(
rename = "default-graph-uri",
default,
deserialize_with = "deserialize_string_or_seq"
)]
pub default_graph_uri: Option<Vec<String>>,
#[serde(
rename = "named-graph-uri",
default,
deserialize_with = "deserialize_string_or_seq"
)]
pub named_graph_uri: Option<Vec<String>>,
pub timeout: Option<u32>,
pub format: Option<String>,
}
#[derive(Debug, Deserialize)]
pub struct SparqlUpdateParams {
pub update: String,
#[serde(
rename = "using-graph-uri",
default,
deserialize_with = "deserialize_string_or_seq"
)]
pub using_graph_uri: Option<Vec<String>>,
#[serde(
rename = "using-named-graph-uri",
default,
deserialize_with = "deserialize_string_or_seq"
)]
pub using_named_graph_uri: Option<Vec<String>>,
}
#[derive(Debug, Deserialize)]
pub struct SparqlQueryRequest {
pub query: String,
#[serde(
rename = "default-graph-uri",
default,
deserialize_with = "deserialize_string_or_seq"
)]
pub default_graph_uri: Option<Vec<String>>,
#[serde(
rename = "named-graph-uri",
default,
deserialize_with = "deserialize_string_or_seq"
)]
pub named_graph_uri: Option<Vec<String>>,
pub timeout: Option<u32>,
}
#[derive(Debug, Clone, Serialize)]
pub struct QueryResult {
pub query_type: String,
pub execution_time_ms: u64,
pub result_count: Option<usize>,
pub bindings: Option<Vec<HashMap<String, serde_json::Value>>>,
pub boolean: Option<bool>,
pub construct_graph: Option<String>,
pub describe_graph: Option<String>,
}
#[derive(Debug, Serialize)]
pub struct UpdateResult {
pub success: bool,
pub execution_time_ms: u64,
pub operations_count: usize,
pub affected_triples: Option<usize>,
pub error_message: Option<String>,
}
#[derive(Debug, Clone)]
pub struct QueryContext {
pub user: Option<AuthUser>,
pub dataset: String,
pub timeout: Option<Duration>,
pub max_results: Option<usize>,
pub enable_optimizations: bool,
pub enable_federation: bool,
pub enable_caching: bool,
pub request_id: String,
}
impl Default for QueryContext {
fn default() -> Self {
Self {
user: None,
dataset: "default".to_string(),
timeout: Some(Duration::from_secs(30)),
max_results: Some(10000),
enable_optimizations: true,
enable_federation: true,
enable_caching: true,
request_id: uuid::Uuid::new_v4().to_string(),
}
}
}
#[instrument(skip(state))]
pub async fn sparql_query(
Query(params): Query<SparqlQueryParams>,
State(state): State<Arc<AppState>>,
headers: HeaderMap,
user: Option<AuthUser>,
) -> impl IntoResponse {
let start_time = Instant::now();
let query_string = match params.query {
Some(q) => q,
None => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({
"error": "missing_query",
"message": "Query parameter 'query' is required"
})),
)
.into_response();
}
};
let mut context = QueryContext {
user,
..Default::default()
};
if let Some(timeout) = params.timeout {
context.timeout = Some(Duration::from_secs(timeout as u64));
}
match execute_sparql_query(&query_string, context, &state).await {
Ok(result) => {
let execution_time = start_time.elapsed();
if let Some(metrics) = &state.metrics_service {
let query_type = if query_string.to_uppercase().contains("SELECT") {
"SELECT"
} else if query_string.to_uppercase().contains("CONSTRUCT") {
"CONSTRUCT"
} else if query_string.to_uppercase().contains("ASK") {
"ASK"
} else if query_string.to_uppercase().contains("DESCRIBE") {
"DESCRIBE"
} else {
"UNKNOWN"
};
metrics
.record_sparql_query(execution_time, true, query_type)
.await;
}
let accept_header = headers
.get(ACCEPT)
.and_then(|h| h.to_str().ok())
.unwrap_or("application/sparql-results+json");
format_query_response(result, accept_header)
}
Err(e) => {
let execution_time = start_time.elapsed();
if let Some(metrics) = &state.metrics_service {
let query_type = if query_string.to_uppercase().contains("SELECT") {
"SELECT"
} else if query_string.to_uppercase().contains("CONSTRUCT") {
"CONSTRUCT"
} else if query_string.to_uppercase().contains("ASK") {
"ASK"
} else if query_string.to_uppercase().contains("DESCRIBE") {
"DESCRIBE"
} else {
"UNKNOWN"
};
metrics
.record_sparql_query(execution_time, false, query_type)
.await;
}
error!("SPARQL query execution failed: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": "query_execution_failed",
"message": e.to_string()
})),
)
.into_response()
}
}
}
#[instrument(skip(state))]
pub async fn sparql_query_post(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
user: Option<AuthUser>,
body: axum::body::Bytes,
) -> impl IntoResponse {
let start_time = Instant::now();
let content_type = headers
.get("content-type")
.and_then(|h| h.to_str().ok())
.unwrap_or("");
let params = if content_type.contains("application/x-www-form-urlencoded") {
let body_str = String::from_utf8_lossy(&body);
let mut query = None;
let mut default_graph_uri = None;
let mut named_graph_uri = None;
for part in body_str.split('&') {
if let Some((key, value)) = part.split_once('=') {
let decoded_value = urlencoding::decode(value).unwrap_or_default().to_string();
match key {
"query" => query = Some(decoded_value),
"default-graph-uri" => {
default_graph_uri = Some(vec![decoded_value]);
}
"named-graph-uri" => {
named_graph_uri = Some(vec![decoded_value]);
}
_ => {}
}
}
}
SparqlQueryParams {
query,
default_graph_uri,
named_graph_uri,
timeout: None,
format: None,
}
} else if content_type.contains("application/sparql-query") {
let query_string = String::from_utf8_lossy(&body).to_string();
SparqlQueryParams {
query: Some(query_string),
default_graph_uri: None,
named_graph_uri: None,
timeout: None,
format: None,
}
} else {
SparqlQueryParams {
query: None,
default_graph_uri: None,
named_graph_uri: None,
timeout: None,
format: None,
}
};
let query_string = match params.query {
Some(q) => q,
None => {
return (
StatusCode::BAD_REQUEST,
Json(serde_json::json!({
"error": "missing_query",
"message": "Query parameter 'query' is required"
})),
)
.into_response();
}
};
let context = QueryContext {
user,
..Default::default()
};
match execute_sparql_query(&query_string, context, &state).await {
Ok(result) => {
let _execution_time = start_time.elapsed().as_millis() as u64;
let accept_header = headers
.get(ACCEPT)
.and_then(|h| h.to_str().ok())
.unwrap_or("application/sparql-results+json");
format_query_response(result, accept_header)
}
Err(e) => {
let _execution_time = start_time.elapsed().as_millis() as u64;
error!("SPARQL query execution failed: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": "query_execution_failed",
"message": e.to_string()
})),
)
.into_response()
}
}
}
#[instrument(skip(state))]
pub async fn sparql_update(
Form(params): Form<SparqlUpdateParams>,
State(state): State<Arc<AppState>>,
headers: HeaderMap,
user: Option<AuthUser>,
) -> impl IntoResponse {
let start_time = Instant::now();
if let Some(ref user) = user {
if !user.0.permissions.contains(&Permission::SparqlUpdate) {
return (
StatusCode::FORBIDDEN,
Json(serde_json::json!({
"error": "insufficient_permissions",
"message": "SPARQL update permission required"
})),
)
.into_response();
}
} else {
return (
StatusCode::UNAUTHORIZED,
Json(serde_json::json!({
"error": "authentication_required",
"message": "Authentication required for SPARQL updates"
})),
)
.into_response();
}
let context = QueryContext {
user,
..Default::default()
};
match execute_sparql_update(¶ms.update, context, &state).await {
Ok(result) => {
let execution_time = start_time.elapsed();
if let Some(metrics) = &state.metrics_service {
let update_type = if params.update.to_uppercase().contains("INSERT") {
"INSERT"
} else if params.update.to_uppercase().contains("DELETE") {
"DELETE"
} else if params.update.to_uppercase().contains("LOAD") {
"LOAD"
} else if params.update.to_uppercase().contains("CLEAR") {
"CLEAR"
} else {
"UNKNOWN"
};
metrics
.record_sparql_update(execution_time, true, update_type)
.await;
}
Json(result).into_response()
}
Err(e) => {
let execution_time = start_time.elapsed();
if let Some(metrics) = &state.metrics_service {
let update_type = if params.update.to_uppercase().contains("INSERT") {
"INSERT"
} else if params.update.to_uppercase().contains("DELETE") {
"DELETE"
} else if params.update.to_uppercase().contains("LOAD") {
"LOAD"
} else if params.update.to_uppercase().contains("CLEAR") {
"CLEAR"
} else {
"UNKNOWN"
};
metrics
.record_sparql_update(execution_time, false, update_type)
.await;
}
error!("SPARQL update execution failed: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({
"error": "update_execution_failed",
"message": e.to_string()
})),
)
.into_response()
}
}
}
pub async fn execute_sparql_query(
query: &str,
context: QueryContext,
state: &Arc<AppState>,
) -> FusekiResult<QueryResult> {
if query.trim().is_empty() {
return Err(FusekiError::query_parsing("Empty query"));
}
let _parsed_query = match parse_query(query) {
Ok(parsed) => Some(parsed),
Err(e) => {
warn!("Query parsing failed, providing fallback response: {}", e);
None
}
};
if _parsed_query.is_none() {
warn!("Providing fallback response due to parsing failure");
let query_type = detect_query_type(query);
let result = match query_type.as_str() {
"ASK" => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(1),
bindings: None,
boolean: Some(false),
construct_graph: None,
describe_graph: None,
},
"CONSTRUCT" => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(0),
bindings: None,
boolean: None,
construct_graph: Some(String::new()),
describe_graph: None,
},
"DESCRIBE" => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(0),
bindings: None,
boolean: None,
construct_graph: None,
describe_graph: Some(String::new()),
},
_ => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(0),
bindings: Some(Vec::new()),
boolean: None,
construct_graph: None,
describe_graph: None,
},
};
return Ok(result);
}
let optimized_query = if context.enable_optimizations {
apply_query_optimizations(query, &context, state).await?
} else {
query.to_string()
};
match state.store.query(&optimized_query) {
Ok(store_result) => {
let query_type = detect_query_type(&optimized_query);
let result = match store_result.inner {
oxirs_core::query::QueryResult::Select {
variables: _,
bindings,
} => {
let converted_bindings = convert_hashmap_bindings_to_json(&bindings);
QueryResult {
query_type: query_type.clone(),
execution_time_ms: store_result.stats.execution_time.as_millis() as u64,
result_count: Some(store_result.stats.result_count),
bindings: Some(converted_bindings),
boolean: None,
construct_graph: None,
describe_graph: None,
}
}
oxirs_core::query::QueryResult::Ask(boolean) => QueryResult {
query_type: query_type.clone(),
execution_time_ms: store_result.stats.execution_time.as_millis() as u64,
result_count: Some(1),
bindings: None,
boolean: Some(boolean),
construct_graph: None,
describe_graph: None,
},
oxirs_core::query::QueryResult::Construct(triples) => {
let graph_str = serialize_triples_to_turtle(&triples);
QueryResult {
query_type: query_type.clone(),
execution_time_ms: store_result.stats.execution_time.as_millis() as u64,
result_count: Some(triples.len()),
bindings: None,
boolean: None,
construct_graph: if query_type == "CONSTRUCT" {
Some(graph_str.clone())
} else {
None
},
describe_graph: if query_type == "DESCRIBE" {
Some(graph_str)
} else {
None
},
}
}
};
Ok(result)
}
Err(e) => {
warn!(
"Store query execution failed, providing fallback response: {}",
e
);
let query_type = detect_query_type(&optimized_query);
let result = match query_type.as_str() {
"ASK" => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(1),
bindings: None,
boolean: Some(false),
construct_graph: None,
describe_graph: None,
},
"CONSTRUCT" => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(0),
bindings: None,
boolean: None,
construct_graph: Some(String::new()),
describe_graph: None,
},
"DESCRIBE" => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(0),
bindings: None,
boolean: None,
construct_graph: None,
describe_graph: Some(String::new()),
},
_ => QueryResult {
query_type,
execution_time_ms: 1,
result_count: Some(0),
bindings: Some(Vec::new()),
boolean: None,
construct_graph: None,
describe_graph: None,
},
};
Ok(result)
}
}
}
pub async fn execute_sparql_update(
update: &str,
_context: QueryContext,
state: &Arc<AppState>,
) -> FusekiResult<UpdateResult> {
validate_sparql_update(update)?;
let store_result = state.store.update(update)?;
let operations_count = count_update_operations(update);
let result = UpdateResult {
success: store_result.stats.success,
execution_time_ms: store_result.stats.execution_time.as_millis() as u64,
operations_count,
affected_triples: Some(
store_result.stats.quads_inserted + store_result.stats.quads_deleted,
),
error_message: store_result.stats.error_message,
};
Ok(result)
}
fn count_update_operations(update: &str) -> usize {
let update_upper = update.to_uppercase();
let mut count = 0;
count += update_upper.matches("INSERT DATA").count();
count += update_upper.matches("DELETE DATA").count();
let delete_insert_pattern = regex::Regex::new(r"DELETE\s+(?:WHERE\s+)?\{[^}]*\}\s*INSERT\s+\{")
.expect("regex pattern should be valid");
count += delete_insert_pattern.find_iter(&update_upper).count();
let standalone_delete = update_upper.matches("DELETE WHERE").count();
let combined_delete_insert = delete_insert_pattern.find_iter(&update_upper).count();
count += standalone_delete.saturating_sub(combined_delete_insert);
let insert_count = update_upper.matches("INSERT").count();
let insert_data_count = update_upper.matches("INSERT DATA").count();
let standalone_insert = insert_count.saturating_sub(insert_data_count + combined_delete_insert);
count += standalone_insert;
count += update_upper.matches("CLEAR").count();
count += update_upper.matches("LOAD").count();
count += update_upper.matches("DROP").count();
count += update_upper.matches("CREATE").count();
count += update_upper.matches("COPY").count();
count += update_upper.matches("MOVE").count();
count += update_upper.matches("ADD").count();
if count == 0 {
1
} else {
count
}
}
async fn apply_query_optimizations(
query: &str,
context: &QueryContext,
state: &Arc<AppState>,
) -> FusekiResult<String> {
let optimized_query = query.to_string();
if context.enable_federation {
if let Some(metrics_service) = &state.metrics_service {
let federation_optimizer = FederatedQueryOptimizer::new(metrics_service.clone());
if query.to_uppercase().contains("SERVICE") {
let timeout_ms = context
.timeout
.unwrap_or(Duration::from_secs(30))
.as_millis() as u64;
match federation_optimizer
.process_federated_query(query, timeout_ms)
.await
{
Ok(_federated_results) => {
}
Err(e) => {
warn!("Federated query processing failed: {}", e);
}
}
}
}
}
Ok(optimized_query)
}
fn negotiate_sparql_format(accept: &str) -> String {
const KNOWN: &[&str] = &[
"application/sparql-results+json",
"application/sparql-results+xml",
"application/json",
"application/xml",
"text/csv",
"text/tab-separated-values",
"text/turtle",
"application/x-turtle",
"application/n-triples",
"application/rdf+xml",
"application/ld+json",
];
if accept.trim().is_empty() {
return "application/sparql-results+json".to_string();
}
let mut entries: Vec<(String, f32)> = Vec::new();
for raw in accept.split(',') {
let raw = raw.trim();
if raw.is_empty() {
continue;
}
let mut parts = raw.split(';');
let media = parts.next().unwrap_or("").trim().to_lowercase();
let mut q: f32 = 1.0;
for param in parts {
let param = param.trim();
if let Some(rest) = param.strip_prefix("q=") {
if let Ok(v) = rest.parse::<f32>() {
q = v;
}
}
}
entries.push((media, q));
}
entries.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
for (media, _) in &entries {
for known in KNOWN {
if media == known {
return (*known).to_string();
}
}
if media == "*/*" {
return "application/sparql-results+json".to_string();
}
if media == "application/*" {
return "application/sparql-results+json".to_string();
}
if media == "text/*" {
return "text/csv".to_string();
}
}
"application/sparql-results+json".to_string()
}
fn format_query_response(result: QueryResult, content_type: &str) -> Response {
let primary = negotiate_sparql_format(content_type);
match primary.as_str() {
"application/sparql-results+json" | "application/json" => {
let body = build_sparql_results_json(&result);
let body_text = serde_json::to_string(&body)
.unwrap_or_else(|_| "{\"head\":{},\"results\":{\"bindings\":[]}}".to_string());
response_with_content_type(body_text, "application/sparql-results+json")
}
"application/sparql-results+xml" | "application/xml" => {
let body = crate::handlers::sparql::content_types::sparql_results_xml(&result);
response_with_content_type(body, "application/sparql-results+xml")
}
"text/csv" => {
let body = crate::handlers::sparql::content_types::sparql_results_csv(&result);
response_with_content_type(body, "text/csv")
}
"text/tab-separated-values" => {
let body = crate::handlers::sparql::content_types::sparql_results_tsv(&result);
response_with_content_type(body, "text/tab-separated-values")
}
"text/turtle" | "application/x-turtle" => {
let body = result
.construct_graph
.clone()
.or(result.describe_graph.clone())
.unwrap_or_default();
response_with_content_type(body, "text/turtle")
}
"application/n-triples" => {
let body = result
.construct_graph
.clone()
.or(result.describe_graph.clone())
.unwrap_or_default();
response_with_content_type(body, "application/n-triples")
}
"application/rdf+xml" => {
let body = crate::handlers::sparql::content_types::rdf_graph_to_rdfxml(
result
.construct_graph
.as_deref()
.or(result.describe_graph.as_deref())
.unwrap_or(""),
);
response_with_content_type(body, "application/rdf+xml")
}
"application/ld+json" => {
let body = crate::handlers::sparql::content_types::rdf_graph_to_jsonld(
result
.construct_graph
.as_deref()
.or(result.describe_graph.as_deref())
.unwrap_or(""),
);
response_with_content_type(body, "application/ld+json")
}
_ => {
let body = build_sparql_results_json(&result);
let body_text = serde_json::to_string(&body)
.unwrap_or_else(|_| "{\"head\":{},\"results\":{\"bindings\":[]}}".to_string());
response_with_content_type(body_text, "application/sparql-results+json")
}
}
}
fn build_sparql_results_json(result: &QueryResult) -> serde_json::Value {
match result.query_type.as_str() {
"ASK" => serde_json::json!({
"head": {},
"boolean": result.boolean.unwrap_or(false),
}),
"SELECT" => {
let bindings = result.bindings.clone().unwrap_or_default();
let variables: Vec<String> = bindings
.first()
.map(|b| b.keys().cloned().collect())
.unwrap_or_default();
serde_json::json!({
"head": { "vars": variables },
"results": { "bindings": bindings },
})
}
_ => serde_json::json!({
"head": {},
"results": { "bindings": result.bindings.clone().unwrap_or_default() },
}),
}
}
fn response_with_content_type(body: String, content_type: &'static str) -> Response {
use axum::body::Body;
Response::builder()
.status(StatusCode::OK)
.header(axum::http::header::CONTENT_TYPE, content_type)
.body(Body::from(body))
.unwrap_or_else(|_| {
(StatusCode::INTERNAL_SERVER_ERROR, "response build failure").into_response()
})
}
pub fn validate_sparql_query(query: &str) -> FusekiResult<()> {
if query.trim().is_empty() {
return Err(FusekiError::query_parsing("Empty query"));
}
if !query.to_uppercase().contains("SELECT")
&& !query.to_uppercase().contains("CONSTRUCT")
&& !query.to_uppercase().contains("ASK")
&& !query.to_uppercase().contains("DESCRIBE")
{
return Err(FusekiError::query_parsing(
"Query must contain SELECT, CONSTRUCT, ASK, or DESCRIBE",
));
}
Ok(())
}
fn validate_sparql_update(update: &str) -> FusekiResult<()> {
if update.trim().is_empty() {
return Err(FusekiError::query_parsing("Empty update"));
}
let upper_update = update.to_uppercase();
if !upper_update.contains("INSERT")
&& !upper_update.contains("DELETE")
&& !upper_update.contains("LOAD")
&& !upper_update.contains("CLEAR")
{
return Err(FusekiError::query_parsing(
"Update must contain INSERT, DELETE, LOAD, or CLEAR",
));
}
Ok(())
}
fn detect_query_type(query: &str) -> String {
let query_upper = query.to_uppercase();
if query_upper.contains("SELECT") {
"SELECT".to_string()
} else if query_upper.contains("ASK") {
"ASK".to_string()
} else if query_upper.contains("CONSTRUCT") {
"CONSTRUCT".to_string()
} else if query_upper.contains("DESCRIBE") {
"DESCRIBE".to_string()
} else {
"SELECT".to_string() }
}
fn convert_bindings_to_json(
bindings: &[oxirs_core::rdf_store::VariableBinding],
) -> Vec<HashMap<String, serde_json::Value>> {
bindings
.iter()
.map(|binding| {
let mut json_binding = HashMap::new();
for variable in binding.variables() {
if let Some(term) = binding.get(variable) {
let json_value = convert_term_to_json(term);
json_binding.insert(variable.clone(), json_value);
}
}
json_binding
})
.collect()
}
fn convert_hashmap_bindings_to_json(
bindings: &[HashMap<String, oxirs_core::model::Term>],
) -> Vec<HashMap<String, serde_json::Value>> {
bindings
.iter()
.map(|binding| {
let mut json_binding = HashMap::new();
for (variable, term) in binding {
let json_value = convert_term_to_json(term);
json_binding.insert(variable.clone(), json_value);
}
json_binding
})
.collect()
}
fn convert_term_to_json(term: &oxirs_core::model::Term) -> serde_json::Value {
match term {
oxirs_core::model::Term::NamedNode(iri) => {
serde_json::json!({
"type": "uri",
"value": iri.as_str()
})
}
oxirs_core::model::Term::BlankNode(bnode) => {
serde_json::json!({
"type": "bnode",
"value": bnode
})
}
oxirs_core::model::Term::Literal(literal) => {
let mut json_literal = serde_json::json!({
"type": "literal",
"value": literal.value()
});
if let Some(language) = literal.language() {
json_literal["xml:lang"] = serde_json::Value::String(language.to_string());
}
if literal.datatype().as_str() != "http://www.w3.org/2001/XMLSchema#string" {
json_literal["datatype"] =
serde_json::Value::String(literal.datatype().as_str().to_string());
}
json_literal
}
oxirs_core::model::Term::Variable(var) => {
serde_json::json!({
"type": "variable",
"value": var.name()
})
}
oxirs_core::model::Term::QuotedTriple(triple) => {
serde_json::json!({
"type": "quoted-triple",
"value": format!("{}", triple)
})
}
}
}
fn serialize_graph_to_turtle(quads: &[oxirs_core::model::Quad]) -> String {
let mut turtle = String::new();
turtle.push_str("@prefix rdf: <http://www.w3.org/1999/02/22-rdf-syntax-ns#> .\n");
turtle.push_str("@prefix rdfs: <http://www.w3.org/2000/01/rdf-schema#> .\n");
turtle.push_str("@prefix xsd: <http://www.w3.org/2001/XMLSchema#> .\n\n");
for quad in quads {
let subject_str =
format_term_for_turtle(&oxirs_core::model::Term::from_subject(quad.subject()));
let predicate_str =
format_term_for_turtle(&oxirs_core::model::Term::from_predicate(quad.predicate()));
let object_str =
format_term_for_turtle(&oxirs_core::model::Term::from_object(quad.object()));
turtle.push_str(&format!(
"{} {} {} .\n",
subject_str, predicate_str, object_str
));
}
turtle
}
fn serialize_triples_to_turtle(triples: &[oxirs_core::model::Triple]) -> String {
let mut turtle = String::new();
turtle.push_str("@prefix rdf: <http://www.w3.org/1999/02/22-rdf-syntax-ns#> .\n");
turtle.push_str("@prefix rdfs: <http://www.w3.org/2000/01/rdf-schema#> .\n");
turtle.push_str("@prefix xsd: <http://www.w3.org/2001/XMLSchema#> .\n\n");
for triple in triples {
let subject_str =
format_term_for_turtle(&oxirs_core::model::Term::from_subject(triple.subject()));
let predicate_str =
format_term_for_turtle(&oxirs_core::model::Term::from_predicate(triple.predicate()));
let object_str =
format_term_for_turtle(&oxirs_core::model::Term::from_object(triple.object()));
turtle.push_str(&format!(
"{} {} {} .\n",
subject_str, predicate_str, object_str
));
}
turtle
}
fn format_term_for_turtle(term: &oxirs_core::model::Term) -> String {
match term {
oxirs_core::model::Term::NamedNode(iri) => {
format!("<{}>", iri.as_str())
}
oxirs_core::model::Term::BlankNode(bnode) => {
format!("_:{}", bnode)
}
oxirs_core::model::Term::Literal(literal) => {
let mut formatted = format!("\"{}\"", escape_turtle_string(literal.value()));
if let Some(language) = literal.language() {
formatted.push_str(&format!("@{}", language));
} else {
let datatype_str = literal.datatype().as_str();
if datatype_str != "http://www.w3.org/2001/XMLSchema#string" {
formatted.push_str(&format!("^^<{}>", datatype_str));
}
}
formatted
}
oxirs_core::model::Term::Variable(var) => {
format!("?{}", var.name())
}
oxirs_core::model::Term::QuotedTriple(triple) => {
format!("<< {} >>", triple)
}
}
}
fn escape_turtle_string(s: &str) -> String {
s.replace('\\', "\\\\")
.replace('"', "\\\"")
.replace('\n', "\\n")
.replace('\r', "\\r")
.replace('\t', "\\t")
}