.
This commit is contained in:
@@ -0,0 +1,20 @@
|
||||
[package]
|
||||
name = "gl-inference"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
color-eyre.workspace = true
|
||||
num-traits.workspace = true
|
||||
oxigraph.workspace = true
|
||||
prost.workspace = true
|
||||
thiserror.workspace = true
|
||||
tonic.workspace = true
|
||||
tonic-prost.workspace = true
|
||||
tokio.workspace = true
|
||||
tracing.workspace = true
|
||||
tracing-appender.workspace = true
|
||||
tracing-subscriber.workspace = true
|
||||
|
||||
[build-dependencies]
|
||||
tonic-prost-build.workspace = true
|
||||
@@ -0,0 +1,7 @@
|
||||
fn main() {
|
||||
tonic_prost_build::compile_protos("../proto/inference.proto")
|
||||
.unwrap_or_else(|e| panic!("Failed to compile protos {:?}", e));
|
||||
|
||||
println!("cargo::rerun-if-changed=build.rs");
|
||||
println!("cargo::rerun-if-changed=proto/inference.proto");
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
use thiserror::Error;
|
||||
|
||||
pub type Result<R> = std::result::Result<R, Error>;
|
||||
|
||||
#[derive(Error, Debug)]
|
||||
pub enum Error {
|
||||
#[error(transparent)]
|
||||
IriParse(#[from] oxigraph::model::IriParseError),
|
||||
|
||||
#[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),
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
pub mod proto {
|
||||
tonic::include_proto!("org.graphofliberty.inference");
|
||||
}
|
||||
|
||||
pub use tonic::transport::Channel;
|
||||
@@ -0,0 +1,58 @@
|
||||
use oxigraph::model::GraphNameRef;
|
||||
use oxigraph::sparql::{SparqlEvaluator};
|
||||
use oxigraph::store::Store;
|
||||
use tracing::debug_span;
|
||||
|
||||
const RDF_SCHEMA_PREFIX: &str = "http://www.w3.org/2000/01/rdf-schema#";
|
||||
const GL_PREFIX: &str = "https://graphofliberty.org/";
|
||||
|
||||
/// `prp-spo1`
|
||||
const SUB_PROPERTY_OF_UPDATE: &str = r#"INSERT {
|
||||
?x ?p2 ?y .
|
||||
} WHERE {
|
||||
GRAPH gl:ontology { ?p1 rdfs:subPropertyOf ?p2 . }
|
||||
?x ?p1 ?y .
|
||||
}"#;
|
||||
|
||||
/// `cax-sco`
|
||||
const SUB_CLASS_OF_UPDATE: &str = r#"INSERT {
|
||||
?x a ?c2 .
|
||||
} WHERE {
|
||||
GRAPH gl:ontology { ?c1 rdfs:subClassOf ?c2 . }
|
||||
?x a ?c1 .
|
||||
}"#;
|
||||
|
||||
/// `prp-dom`
|
||||
const DOMAIN_UPDATE: &str = r#"INSERT {
|
||||
?x a ?c .
|
||||
} WHERE {
|
||||
GRAPH gl:ontology { ?p rdfs:domain ?c . }
|
||||
?x ?p ?y .
|
||||
}"#;
|
||||
|
||||
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()?;
|
||||
|
||||
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)?;
|
||||
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_inferences_total = inferences + (new_count - old_count);
|
||||
if new_count > old_count {
|
||||
infer(new_inferences_total, store, graph_name)
|
||||
} else {
|
||||
Ok(new_inferences_total)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
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;
|
||||
|
||||
mod service;
|
||||
mod logic;
|
||||
mod error;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> color_eyre::Result<()> {
|
||||
let appender = tracing_appender::rolling::never("/tmp", "inference");
|
||||
tracing_subscriber::registry()
|
||||
.with(
|
||||
fmt::layer()
|
||||
.with_span_events(FmtSpan::CLOSE)
|
||||
.with_writer(appender),
|
||||
)
|
||||
.with(EnvFilter::from_default_env())
|
||||
.init();
|
||||
color_eyre::install()?;
|
||||
|
||||
let ontology_service = OntologyService::new();
|
||||
let server = OntologyServer::new(ontology_service);
|
||||
|
||||
let addr = String::from("[::1]:3000");
|
||||
println!("Inference engine listening on: {addr}");
|
||||
|
||||
Server::builder()
|
||||
.add_service(server)
|
||||
.serve(addr.parse().unwrap())
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,164 @@
|
||||
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::store::Store;
|
||||
use tonic::{Request, Response, Status};
|
||||
use tracing::{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",
|
||||
));
|
||||
|
||||
pub struct OntologyService {
|
||||
ontology: Store,
|
||||
}
|
||||
|
||||
impl OntologyService {
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
ontology: Store::new().unwrap(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tonic::async_trait]
|
||||
impl Ontology for OntologyService {
|
||||
async fn load(&self, request: Request<OntologyLoadRequest>) -> Result<Response<OntologyLoadResponse>, Status> {
|
||||
let span = debug_span!("Load Ontology", input_size = field::Empty, old_size = field::Empty, new_size = field::Empty).entered();
|
||||
|
||||
let request = request.get_ref();
|
||||
let mut response = OntologyLoadResponse::default();
|
||||
|
||||
response.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();
|
||||
response.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
|
||||
});
|
||||
self.ontology.extend(quads).unwrap();
|
||||
|
||||
response.new_size = self.ontology.len()
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX);
|
||||
|
||||
span.record("input_size", response.input_size);
|
||||
span.record("old_size", response.old_size);
|
||||
span.record("new_size", response.new_size);
|
||||
|
||||
Ok(Response::new(response))
|
||||
}
|
||||
|
||||
async fn clear(&self, _request: Request<()>) -> Result<Response<OntologyClearResponse>, Status> {
|
||||
let span = debug_span!("Clear Ontology", size = field::Empty).entered();
|
||||
|
||||
let mut response = OntologyClearResponse::default();
|
||||
response.size = self.ontology.len()
|
||||
.map_err(|err| Status::internal(err.to_string()))?
|
||||
.to_u64()
|
||||
.unwrap_or(u64::MAX);
|
||||
|
||||
self.ontology.clear()
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
span.record("size", response.size);
|
||||
|
||||
Ok(Response::new(response))
|
||||
}
|
||||
|
||||
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();
|
||||
let mut response = OntologyQueryResponse::default();
|
||||
|
||||
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 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>();
|
||||
|
||||
self.ontology.extend(&dataset)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
let inferences = crate::logic::infer(0, &self.ontology, graph_name.as_ref())
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
span.record("inferences", inferences);
|
||||
graph_name
|
||||
} else {
|
||||
ONTOLOGY_GRAPH.into_owned()
|
||||
};
|
||||
|
||||
let mut query = SparqlEvaluator::new()
|
||||
.parse_query(&request.sparql_query)
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
query.dataset_mut().set_default_graph(vec![graph_name.clone()]);
|
||||
|
||||
let results = query.on_store(&self.ontology)
|
||||
.execute()
|
||||
.map_err(|err| Status::internal(err.to_string()))?;
|
||||
|
||||
let mut output_buffer = Vec::new();
|
||||
match results {
|
||||
QueryResults::Boolean(result) => {
|
||||
let serializer = QueryResultsSerializer::from_format(QueryResultsFormat::Json);
|
||||
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)?;
|
||||
for solution in solutions.filter_map(Result::ok) {
|
||||
writer.serialize(&solution)?;
|
||||
}
|
||||
output_buffer = writer.finish()?;
|
||||
}
|
||||
QueryResults::Graph(graph) => {
|
||||
let mut serializer = RdfSerializer::from_format(RdfFormat::Turtle)
|
||||
.with_prefix("rdfs", "http://www.w3.org/2000/01/rdf-schema#").map_err(|err| Status::internal(err.to_string()))?
|
||||
.with_prefix("rdac", "http://rdaregistry.info/Elements/c/").map_err(|err| Status::internal(err.to_string()))?
|
||||
.with_prefix("ldp", "http://www.w3.org/ns/ldp#").map_err(|err| Status::internal(err.to_string()))?
|
||||
.with_prefix("fedora", "http://fedora.info/definitions/v4/repository#").map_err(|err| Status::internal(err.to_string()))?
|
||||
.with_prefix("lrmer", "http://iflastandards.info/ns/lrm/lrmer/").map_err(|err| Status::internal(err.to_string()))?
|
||||
.with_prefix("owl", "http://www.w3.org/2002/07/owl#").map_err(|err| Status::internal(err.to_string()))?
|
||||
.with_base_iri("http://fedora.quill.lan/rest/").map_err(|err| Status::internal(err.to_string()))?
|
||||
.for_writer(output_buffer);
|
||||
for triple in graph.filter_map(Result::ok) {
|
||||
serializer.serialize_triple(triple.as_ref())?;
|
||||
}
|
||||
output_buffer = serializer.finish()?;
|
||||
}
|
||||
}
|
||||
response.results = String::from_utf8_lossy(&output_buffer).to_string();
|
||||
|
||||
if request.turtle.is_some() {
|
||||
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