use crate::server::Server;
use crate::server::configuration::BackendConfiguration;
use crate::server::configuration::RequestMethod;
use crate::server::lsp::ExecuteUpdateResponseResult;
use crate::server::lsp::SparqlEngine;
use crate::server::sparql_operations::ConnectionError;
use crate::server::sparql_operations::HttpError;
use crate::server::sparql_operations::SparqlRequestError;
use crate::server::sparql_operations::utils::add_limit_offset_to_query;
use crate::server::sparql_operations::utils::health_check_url;
use crate::sparql::results::SparqlResult;
use futures::lock::Mutex;
use reqwest::Client;
use std::rc::Rc;
use std::time::Duration;
use tokio::time::timeout;
use urlencoding::encode;
#[allow(clippy::too_many_arguments)]
pub(crate) async fn execute_query(
_server_rc: Rc<Mutex<Server>>,
url: String,
mut query: String,
_query_id: Option<&str>,
_engine: Option<SparqlEngine>,
timeout_ms: Option<u32>,
method: RequestMethod,
limit: Option<usize>,
offset: usize,
lazy: bool,
) -> Result<Option<SparqlResult>, SparqlRequestError> {
if lazy {
tracing::warn!("Lazy Query execution is not implemented for non wasm targets");
}
if let Some(new_query) = add_limit_offset_to_query(&query, limit, offset) {
query = new_query;
}
let request = match method {
RequestMethod::GET => Client::new()
.get(format!("{}?query={}", url, encode(&query)))
.header(
"Content-Type",
"application/x-www-form-urlencoded;charset=UTF-8",
)
.header("Accept", "application/sparql-results+json")
.header("User-Agent", "qlue-ls/1.0")
.send(),
RequestMethod::POST => Client::new()
.post(url)
.header(
"Content-Type",
"application/x-www-form-urlencoded;charset=UTF-8",
)
.header("Accept", "application/sparql-results+json")
.header("User-Agent", "qlue-ls/1.0")
.form(&[("query", &query)])
.send(),
};
let duration = Duration::from_millis(timeout_ms.unwrap_or(5000) as u64);
let request = timeout(duration, request);
let response = request
.await
.map_err(|_| SparqlRequestError::Timeout)?
.map_err(|err| {
SparqlRequestError::Connection(ConnectionError {
message: err.to_string(),
query,
})
})?;
let status = response.status();
if !status.is_success() {
let body = response.text().await.unwrap_or_default();
return Err(match serde_json::from_str(&body) {
Ok(exception) => SparqlRequestError::QLeverException(exception),
Err(_) => SparqlRequestError::Http(HttpError {
status: status.as_u16(),
status_text: status
.canonical_reason()
.unwrap_or("Unknown Status")
.to_string(),
body,
}),
});
}
let result = response
.json::<SparqlResult>()
.await
.map_err(|err| SparqlRequestError::Deserialization(err.to_string()))?;
Ok(Some(result))
}
pub(crate) async fn check_server_availability(backend: &BackendConfiguration) -> bool {
let url = health_check_url(backend);
let request = Client::new()
.get(&url)
.header("Accept", "application/sparql-results+json")
.header("User-Agent", "qlue-ls/1.0")
.send();
match timeout(Duration::from_secs(5), request).await {
Ok(Ok(response)) => response.status().is_success(),
Ok(Err(err)) => {
tracing::info!("health check for \"{}\" failed: {}", url, err);
false
}
Err(_) => {
tracing::info!("health check for \"{}\" timed out", url);
false
}
}
}
pub(crate) async fn execute_construct_query(
_server_rc: Rc<Mutex<Server>>,
_url: &str,
_query: &str,
_query_id: Option<&str>,
_engine: Option<SparqlEngine>,
_lazy: bool,
) -> Result<Option<SparqlResult>, SparqlRequestError> {
todo!()
}
pub(crate) async fn execute_update(
_server_rc: Rc<Mutex<Server>>,
_url: &str,
_query: &str,
_query_id: Option<&str>,
_access_token: Option<&str>,
) -> Result<ExecuteUpdateResponseResult, SparqlRequestError> {
todo!()
}