.
This commit is contained in:
+1
-1
@@ -4,4 +4,4 @@ fn main() {
|
||||
|
||||
println!("cargo::rerun-if-changed=build.rs");
|
||||
println!("cargo::rerun-if-changed=proto/inference.proto");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,13 +9,13 @@ pub enum Error {
|
||||
|
||||
#[error(transparent)]
|
||||
SparqlSyntax(#[from] oxigraph::sparql::SparqlSyntaxError),
|
||||
|
||||
|
||||
#[error(transparent)]
|
||||
QueryEvaluation(#[from] oxigraph::sparql::QueryEvaluationError),
|
||||
|
||||
|
||||
#[error(transparent)]
|
||||
UpdateEvaluation(#[from] oxigraph::sparql::UpdateEvaluationError),
|
||||
|
||||
#[error(transparent)]
|
||||
Storage(#[from] oxigraph::store::StorageError),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,4 +2,4 @@ pub mod proto {
|
||||
tonic::include_proto!("org.graphofliberty.inference");
|
||||
}
|
||||
|
||||
pub use tonic;
|
||||
pub use tonic;
|
||||
|
||||
+27
-8
@@ -1,5 +1,5 @@
|
||||
use oxigraph::model::GraphNameRef;
|
||||
use oxigraph::sparql::{SparqlEvaluator};
|
||||
use oxigraph::sparql::SparqlEvaluator;
|
||||
use oxigraph::store::Store;
|
||||
use tracing::debug_span;
|
||||
|
||||
@@ -30,24 +30,43 @@ const DOMAIN_UPDATE: &str = r#"INSERT {
|
||||
?x ?p ?y .
|
||||
}"#;
|
||||
|
||||
fn run_update(query_name: &str, query: &str, store: &Store, graph_name: GraphNameRef<'_>) -> crate::error::Result<()> {
|
||||
fn run_update(
|
||||
query_name: &str,
|
||||
query: &str,
|
||||
store: &Store,
|
||||
graph_name: GraphNameRef<'_>,
|
||||
) -> crate::error::Result<()> {
|
||||
let _span = debug_span!("Update", name = query_name).entered();
|
||||
|
||||
SparqlEvaluator::new()
|
||||
.with_prefix("rdfs", RDF_SCHEMA_PREFIX)?
|
||||
.with_prefix("gl", GL_PREFIX)?
|
||||
.parse_update(&format!("WITH {graph_name} {query}"))?
|
||||
.on_store(store).execute()?;
|
||||
.on_store(store)
|
||||
.execute()?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) fn infer(inferences: usize, store: &Store, graph_name: GraphNameRef<'_>) -> crate::error::Result<usize> {
|
||||
let old_count = store.quads_for_pattern(None, None, None, Some(graph_name)).count();
|
||||
run_update("rdfs:subPropertyOf", SUB_PROPERTY_OF_UPDATE, store, graph_name)?;
|
||||
pub(crate) fn infer(
|
||||
inferences: usize,
|
||||
store: &Store,
|
||||
graph_name: GraphNameRef<'_>,
|
||||
) -> crate::error::Result<usize> {
|
||||
let old_count = store
|
||||
.quads_for_pattern(None, None, None, Some(graph_name))
|
||||
.count();
|
||||
run_update(
|
||||
"rdfs:subPropertyOf",
|
||||
SUB_PROPERTY_OF_UPDATE,
|
||||
store,
|
||||
graph_name,
|
||||
)?;
|
||||
run_update("rdfs:subClassOf", SUB_CLASS_OF_UPDATE, store, graph_name)?;
|
||||
run_update("rdfs:domain", DOMAIN_UPDATE, store, graph_name)?;
|
||||
let new_count = store.quads_for_pattern(None, None, None, Some(graph_name)).count();
|
||||
let new_count = store
|
||||
.quads_for_pattern(None, None, None, Some(graph_name))
|
||||
.count();
|
||||
|
||||
let new_inferences_total = inferences + (new_count - old_count);
|
||||
if new_count > old_count {
|
||||
@@ -55,4 +74,4 @@ pub(crate) fn infer(inferences: usize, store: &Store, graph_name: GraphNameRef<'
|
||||
} else {
|
||||
Ok(new_inferences_total)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,14 +1,14 @@
|
||||
use crate::service::OntologyService;
|
||||
use gl_inference::proto::ontology_server::OntologyServer;
|
||||
use tonic::transport::Server;
|
||||
use tracing_subscriber::{fmt, EnvFilter};
|
||||
use tracing_subscriber::fmt::format::FmtSpan;
|
||||
use tracing_subscriber::layer::SubscriberExt;
|
||||
use tracing_subscriber::util::SubscriberInitExt;
|
||||
use gl_inference::proto::ontology_server::OntologyServer;
|
||||
use crate::service::OntologyService;
|
||||
use tracing_subscriber::{EnvFilter, fmt};
|
||||
|
||||
mod service;
|
||||
mod logic;
|
||||
mod error;
|
||||
mod logic;
|
||||
mod service;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> color_eyre::Result<()> {
|
||||
@@ -35,4 +35,4 @@ async fn main() -> color_eyre::Result<()> {
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
+83
-39
@@ -1,13 +1,16 @@
|
||||
use gl_inference::proto::ontology_server::Ontology;
|
||||
use gl_inference::proto::{
|
||||
OntologyClearResponse, OntologyLoadRequest, OntologyLoadResponse, OntologyQueryRequest,
|
||||
OntologyQueryResponse,
|
||||
};
|
||||
use num_traits::ToPrimitive;
|
||||
use oxigraph::io::{RdfFormat, RdfParser, RdfSerializer};
|
||||
use oxigraph::model::{BlankNode, Dataset, GraphName, GraphNameRef, NamedNode, NamedNodeRef};
|
||||
use oxigraph::sparql::{QueryResults, SparqlEvaluator};
|
||||
use oxigraph::sparql::results::{QueryResultsFormat, QueryResultsSerializer};
|
||||
use oxigraph::sparql::{QueryResults, SparqlEvaluator};
|
||||
use oxigraph::store::Store;
|
||||
use tonic::{Request, Response, Status};
|
||||
use tracing::{debug, debug_span, field};
|
||||
use gl_inference::proto::ontology_server::Ontology;
|
||||
use gl_inference::proto::{OntologyClearResponse, OntologyLoadRequest, OntologyLoadResponse, OntologyQueryRequest, OntologyQueryResponse};
|
||||
|
||||
const ONTOLOGY_GRAPH: GraphNameRef = GraphNameRef::NamedNode(NamedNodeRef::new_unchecked(
|
||||
"https://graphofliberty.org/ontology",
|
||||
@@ -27,29 +30,40 @@ impl OntologyService {
|
||||
|
||||
#[tonic::async_trait]
|
||||
impl Ontology for OntologyService {
|
||||
async fn load(&self, request: Request<OntologyLoadRequest>) -> Result<Response<OntologyLoadResponse>, Status> {
|
||||
let span = debug_span!("Load Ontology", old_size = field::Empty, input_size = field::Empty, inferences = field::Empty, new_size = field::Empty).entered();
|
||||
async fn load(
|
||||
&self,
|
||||
request: Request<OntologyLoadRequest>,
|
||||
) -> Result<Response<OntologyLoadResponse>, Status> {
|
||||
let span = debug_span!(
|
||||
"Load Ontology",
|
||||
old_size = field::Empty,
|
||||
input_size = field::Empty,
|
||||
inferences = field::Empty,
|
||||
new_size = field::Empty
|
||||
)
|
||||
.entered();
|
||||
|
||||
let request = request.get_ref();
|
||||
|
||||
let old_size = self.ontology.len()
|
||||
let old_size = self
|
||||
.ontology
|
||||
.len()
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX);
|
||||
|
||||
let path = std::path::Path::new(&request.path);
|
||||
let store = Store::open_read_only(path).unwrap();
|
||||
let input_size = store.len()
|
||||
let input_size = store
|
||||
.len()
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX);
|
||||
|
||||
let quads = store.iter()
|
||||
.filter_map(Result::ok)
|
||||
.map(|mut quad| {
|
||||
quad.graph_name = ONTOLOGY_GRAPH.into_owned();
|
||||
quad
|
||||
});
|
||||
let quads = store.iter().filter_map(Result::ok).map(|mut quad| {
|
||||
quad.graph_name = ONTOLOGY_GRAPH.into_owned();
|
||||
quad
|
||||
});
|
||||
self.ontology.extend(quads).unwrap();
|
||||
|
||||
let inferences = if request.infer {
|
||||
@@ -57,9 +71,13 @@ impl Ontology for OntologyService {
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX)
|
||||
} else { 0 };
|
||||
} else {
|
||||
0
|
||||
};
|
||||
|
||||
let new_size = self.ontology.len()
|
||||
let new_size = self
|
||||
.ontology
|
||||
.len()
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX);
|
||||
@@ -79,27 +97,34 @@ impl Ontology for OntologyService {
|
||||
Ok(Response::new(response))
|
||||
}
|
||||
|
||||
async fn clear(&self, _request: Request<()>) -> Result<Response<OntologyClearResponse>, Status> {
|
||||
async fn clear(
|
||||
&self,
|
||||
_request: Request<()>,
|
||||
) -> Result<Response<OntologyClearResponse>, Status> {
|
||||
let span = debug_span!("Clear Ontology", size = field::Empty).entered();
|
||||
|
||||
let size = self.ontology.len()
|
||||
let size = self
|
||||
.ontology
|
||||
.len()
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX);
|
||||
|
||||
self.ontology.clear()
|
||||
self.ontology
|
||||
.clear()
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
let response = OntologyClearResponse {
|
||||
size,
|
||||
};
|
||||
let response = OntologyClearResponse { size };
|
||||
|
||||
span.record("size", response.size);
|
||||
|
||||
Ok(Response::new(response))
|
||||
}
|
||||
|
||||
async fn query(&self, request: Request<OntologyQueryRequest>) -> Result<Response<OntologyQueryResponse>, Status> {
|
||||
async fn query(
|
||||
&self,
|
||||
request: Request<OntologyQueryRequest>,
|
||||
) -> Result<Response<OntologyQueryResponse>, Status> {
|
||||
let span = debug_span!("Query", inferences = field::Empty).entered();
|
||||
|
||||
let request = request.get_ref();
|
||||
@@ -108,16 +133,21 @@ impl Ontology for OntologyService {
|
||||
let provided_dataset;
|
||||
let graph_name = if let Some(turtle) = &request.turtle {
|
||||
let random_graph_identifier = BlankNode::default();
|
||||
let graph_name = GraphName::NamedNode(NamedNode::new_unchecked(format!("https://graphofliberty.org/inference/{}", random_graph_identifier.as_str())));
|
||||
let graph_name = GraphName::NamedNode(NamedNode::new_unchecked(format!(
|
||||
"https://graphofliberty.org/inference/{}",
|
||||
random_graph_identifier.as_str()
|
||||
)));
|
||||
provided_dataset = RdfParser::from_format(RdfFormat::Turtle)
|
||||
.for_slice(&turtle)
|
||||
.filter_map(Result::ok)
|
||||
.map(|mut quad| {
|
||||
quad.graph_name = graph_name.clone();
|
||||
quad
|
||||
}).collect::<Dataset>();
|
||||
})
|
||||
.collect::<Dataset>();
|
||||
|
||||
self.ontology.extend(&provided_dataset)
|
||||
self.ontology
|
||||
.extend(&provided_dataset)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
let inferences = crate::logic::infer(0, &self.ontology, graph_name.as_ref())
|
||||
@@ -135,34 +165,42 @@ impl Ontology for OntologyService {
|
||||
|
||||
let mut evaluator = SparqlEvaluator::new();
|
||||
for (name, iri) in &request.prefixes {
|
||||
evaluator = evaluator.with_prefix(name, iri)
|
||||
evaluator = evaluator
|
||||
.with_prefix(name, iri)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
|
||||
if let Some(base) = &request.base {
|
||||
evaluator = evaluator.with_base_iri(base)
|
||||
evaluator = evaluator
|
||||
.with_base_iri(base)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
|
||||
let mut query = evaluator.parse_query(sparql_query)
|
||||
let mut query = evaluator
|
||||
.parse_query(sparql_query)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
query.dataset_mut().set_default_graph(vec![graph_name.clone()]);
|
||||
query
|
||||
.dataset_mut()
|
||||
.set_default_graph(vec![graph_name.clone()]);
|
||||
|
||||
let results = query.on_store(&self.ontology)
|
||||
let results = query
|
||||
.on_store(&self.ontology)
|
||||
.execute()
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
match results {
|
||||
QueryResults::Boolean(result) => {
|
||||
let serializer = QueryResultsSerializer::from_format(QueryResultsFormat::Json);
|
||||
output_buffer = serializer.serialize_boolean_to_writer(output_buffer, result)
|
||||
output_buffer = serializer
|
||||
.serialize_boolean_to_writer(output_buffer, result)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
QueryResults::Solutions(solutions) => {
|
||||
let serializer = QueryResultsSerializer::from_format(QueryResultsFormat::Json);
|
||||
let variables = Vec::from_iter(solutions.variables().iter().cloned());
|
||||
let mut writer = serializer.serialize_solutions_to_writer(output_buffer, variables)?;
|
||||
let mut writer =
|
||||
serializer.serialize_solutions_to_writer(output_buffer, variables)?;
|
||||
for solution in solutions.filter_map(Result::ok) {
|
||||
writer.serialize(&solution)?;
|
||||
}
|
||||
@@ -171,11 +209,13 @@ impl Ontology for OntologyService {
|
||||
QueryResults::Graph(graph) => {
|
||||
let mut serializer = RdfSerializer::from_format(RdfFormat::Turtle);
|
||||
for (name, iri) in &request.prefixes {
|
||||
serializer = serializer.with_prefix(name, iri)
|
||||
serializer = serializer
|
||||
.with_prefix(name, iri)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
if let Some(base) = &request.base {
|
||||
serializer = serializer.with_base_iri(base)
|
||||
serializer = serializer
|
||||
.with_base_iri(base)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
|
||||
@@ -189,16 +229,19 @@ impl Ontology for OntologyService {
|
||||
} else {
|
||||
let mut serializer = RdfSerializer::from_format(RdfFormat::Turtle);
|
||||
for (name, iri) in &request.prefixes {
|
||||
serializer = serializer.with_prefix(name, iri)
|
||||
serializer = serializer
|
||||
.with_prefix(name, iri)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
if let Some(base) = &request.base {
|
||||
serializer = serializer.with_base_iri(base)
|
||||
serializer = serializer
|
||||
.with_base_iri(base)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
|
||||
let mut serializer = serializer.for_writer(output_buffer);
|
||||
let resulting_dataset = self.ontology
|
||||
let resulting_dataset = self
|
||||
.ontology
|
||||
.quads_for_pattern(None, None, None, Some(graph_name.as_ref()))
|
||||
.filter_map(Result::ok);
|
||||
for quad in resulting_dataset {
|
||||
@@ -212,10 +255,11 @@ impl Ontology for OntologyService {
|
||||
response.results = String::from_utf8_lossy(&output_buffer).to_string();
|
||||
|
||||
if request.turtle.is_some() {
|
||||
self.ontology.clear_graph(graph_name.as_ref())
|
||||
self.ontology
|
||||
.clear_graph(graph_name.as_ref())
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
}
|
||||
|
||||
Ok(Response::new(response))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user