mod app; mod args; mod error; mod navigator; mod rdf; mod theme; mod widget; mod tasks; use std::collections::HashMap; use crate::app::Publisher; use crate::args::{AppArgs, Command}; use clap::Parser; use gl_search::{Schema, SearchIndex}; use iced::futures::StreamExt; use ldp::middleware::BasicAuthMiddleware; use ldp::reqwest::Client; use ldp::reqwest_middleware::ClientBuilder; use ldp::traverse::Traverse; use oxigraph::io::{RdfFormat, RdfParser, RdfSerializer}; use oxigraph::model::{Dataset, Graph, Triple, TripleRef}; use tracing::{debug_span, error, field, Instrument}; use tracing_subscriber::fmt::format::FmtSpan; use tracing_subscriber::layer::SubscriberExt; use tracing_subscriber::util::SubscriberInitExt; use tracing_subscriber::{EnvFilter, fmt}; use url::Url; use gl_graph::indexer::Indexer; use gl_graph::{language, CurieHelper, vocab}; use gl_graph::ontology::{OntologyBuilder, ResourceSelector}; use gl_inference::proto::ontology_client::OntologyClient; use gl_inference::proto::OntologyQueryRequest; fn main() -> color_eyre::Result<()> { let appender = tracing_appender::rolling::never("/tmp", "publisher-log"); tracing_subscriber::registry() .with( fmt::layer() .with_span_events(FmtSpan::CLOSE) .with_writer(appender), ) .with(EnvFilter::from_default_env()) .init(); color_eyre::install()?; let args = AppArgs::parse(); match args.command { Some(Command::Query(args)) => { let raw_query = if let Some(query) = &args.query_path { Some(String::from_utf8(std::fs::read(query)?)?) } else { None }; let graph = if let Some(dataset_path) = &args.dataset_path { Some(String::from_utf8(std::fs::read(dataset_path)?)?) } else { None }; let mut request = OntologyQueryRequest::default(); request.sparql_query = raw_query; request.turtle = graph; request.prefixes = HashMap::from_iter(gl_graph::PREFIXES.iter().map(|(name, iri)| (name.clone(), iri.clone()))); request.base = args.base.clone(); request.inferences_only = args.inferences_only; let runtime = tokio::runtime::Builder::new_multi_thread() .enable_all() .name("inference-client") .build()?; runtime.block_on(async { let mut client = OntologyClient::connect("http://[::1]:3000") .await .unwrap() .max_decoding_message_size(1024 * 1024 * 1024); let response = client.query(request).await.unwrap(); println!("{}", response.get_ref().results); }); } Some(Command::Search(args)) => { let mut index = SearchIndex::builder() .with_path("/home/alex/.local/share/org.graphofliberty.desktop/index") .build() .expect("Failed to build search index"); let limit = args.limit.unwrap_or(5); for document in index.query(args.discriminant, &args.query, Schema::all_fields(), limit)? { println!("{}", gl_search::to_json(document)); } } Some(Command::Reindex) => { let mut index = SearchIndex::builder() .with_path("/home/alex/.local/share/org.graphofliberty.desktop/index") .build() .expect("Failed to build search index"); let mut writer = index.writer()?; debug_span!("Clear Index").in_scope(|| { writer.delete_all_documents()?; writer.commit() })?; let runtime = tokio::runtime::Builder::new_multi_thread() .enable_all() .name("reindex") .build()?; runtime.block_on(async { let curie_helper = CurieHelper::new(gl_graph::PREFIXES.clone()); let indexer = Indexer::new(language::ENGLISH_OR_UNTAGGED.clone(), &curie_helper); let mut client = OntologyBuilder::from_string("http://[::1]:3000", language::ENGLISH_OR_UNTAGGED.clone())? .connect() .await?; let resource_descriptions = client.resource_descriptions(ResourceSelector::Classes, [ vocab::rdf::PROPERTY.into_owned(), vocab::rdfs::CLASS.into_owned(), vocab::skos::CONCEPT.into_owned(), ]).await?; let documents = indexer.ontology(resource_descriptions); debug_span!("Index Ontology", documents = field::Empty).in_scope(|| { for document in documents { writer.add_document(document).unwrap(); } }); Ok::<_, crate::error::Error>(()) })?; let mut writer = runtime.block_on(async move { let http_client = ClientBuilder::new(Client::new()) .with(BasicAuthMiddleware::new( "fedoraAdmin".to_string(), Some("fedoraAdmin".to_string()), )) .build(); let starting_url = Url::parse("http://fedora.quill.lan/rest/")?; let mut graph = Graph::new(); let mut traversal = Traverse::new(http_client, starting_url, None); let mut rdf_source_count = 0usize; while let Some(result) = traversal.next().await { match result { Ok(rdf_source) => { let triples = rdf_source.dataset() .iter() .map(|quad| TripleRef::from(quad)); graph.extend(triples); rdf_source_count += 1; }, Err(err) => error!(?err), } } let span = debug_span!("Index Repository", rdf_sources = rdf_source_count, triples = field::Empty, documents = field::Empty); let triple_count = async { let curie_helper = CurieHelper::new(gl_graph::PREFIXES.clone()); let indexer = Indexer::new(language::ENGLISH_OR_UNTAGGED.clone(), &curie_helper); let mut client = OntologyBuilder::from_string("http://[::1]:3000", language::ENGLISH_OR_UNTAGGED.clone())? .connect() .await?; let graph_with_inferences = client.run_inference(&graph).await?; for document in indexer.graph(&graph_with_inferences) { writer.add_document(document)?; } Ok::(graph_with_inferences.len()) }.instrument(span.clone()).await?; span.record("triples", triple_count); Ok::<_, crate::error::Error>(writer) })?; debug_span!("Commit").in_scope(|| writer.commit())?; } None => { iced::daemon(Publisher::new, Publisher::update, Publisher::view) .title(Publisher::title) .subscription(Publisher::subscription) .run()?; } } Ok(()) }