use crate::container::ContainerType;
use crate::rdf_source::RdfSource;
use crate::{Preference, prefer, vocab};
use bytes::Bytes;
use futures::Stream;
use http::Method;
use oxigraph::io::{RdfFormat, RdfParser};
use oxigraph::model::{GraphNameRef, NamedNodeRef, Quad};
use reqwest_middleware::reqwest::{Client, Response, StatusCode, Url, header};
use reqwest_middleware::{ClientBuilder, ClientWithMiddleware, RequestBuilder};
use sfv::{Item, TokenRef};
use std::collections::BTreeSet;
use std::str::FromStr;
use tracing::error;
#[derive(Clone, Debug)]
pub struct ResourceRequestBuilder {
client: ClientWithMiddleware,
url: Url,
follow_described_by: bool,
validate_support: bool,
formats: Vec<RdfFormat>,
include_preferences: Vec<Preference>,
omit_preferences: Vec<Preference>,
}
impl ResourceRequestBuilder {
pub fn new(url: Url) -> Self {
Self::with_client_and_url(ClientBuilder::new(Client::new()).build(), url)
}
#[must_use]
pub fn with_client_and_url(client: ClientWithMiddleware, url: Url) -> Self {
Self {
client,
url,
follow_described_by: true,
validate_support: true,
formats: Vec::new(),
include_preferences: Vec::new(),
omit_preferences: Vec::new(),
}
}
#[must_use]
pub fn follow_described_by(mut self, value: bool) -> Self {
self.follow_described_by = value;
self
}
#[must_use]
pub fn validate_support(mut self, value: bool) -> Self {
self.validate_support = value;
self
}
#[must_use]
pub fn accept_rdf_format(mut self, format: RdfFormat) -> Self {
self.formats.push(format);
self
}
#[must_use]
pub fn accept_all_rdf_formats(mut self) -> Self {
self.formats = vec![
RdfFormat::N3,
RdfFormat::NQuads,
RdfFormat::NTriples,
RdfFormat::RdfXml,
RdfFormat::TriG,
RdfFormat::Turtle,
];
self
}
pub fn include_preference(mut self, preference: Preference) -> Self {
self.include_preferences.push(preference);
self
}
pub fn omit_preference(mut self, preference: Preference) -> Self {
self.omit_preferences.push(preference);
self
}
pub fn build(self) -> ResourceRequest {
ResourceRequest { builder: self }
}
}
#[derive(Clone, Debug)]
pub struct ResourceRequest {
builder: ResourceRequestBuilder,
}
impl ResourceRequest {
fn extract_described_by(response: &Response) -> Option<Url> {
if response.status() == StatusCode::OK {
let headers = response.headers();
for link in headers.get_all(header::LINK) {
let link = link.to_str().unwrap_or("");
match parse_link_header::parse(link) {
Ok(link_map) => {
if let Some(metadata_url) = link_map.get(&Some("describedby".to_string())) {
let raw_url = metadata_url.raw_uri.as_str();
return Url::parse(raw_url).ok();
}
}
Err(err) => error!(err = ?err, "Failed to parse Link header"),
}
}
}
None
}
fn ensure_ldp_support(response: &Response) -> crate::Result<()> {
if response.status() == StatusCode::OK {
let headers = response.headers();
for link in headers.get_all(header::LINK) {
let link = link.to_str().unwrap_or_default();
match parse_link_header::parse(link) {
Ok(link_map) => {
if let Some(metadata_url) = link_map.get(&Some("type".to_string())) {
let raw_url = metadata_url.raw_uri.as_str();
if raw_url == vocab::ldp::RESOURCE {
return Ok(());
}
}
}
Err(err) => error!(err = ?err, "Failed to parse Link header"),
}
}
}
Err(crate::Error::LDPUnsupported)
}
fn resource_type(response: &Response) -> Option<ResourceType> {
let headers = response.headers();
let links = headers
.get_all(header::LINK)
.iter()
.map(|hv| hv.to_str().unwrap_or_default())
.filter_map(|link| parse_link_header::parse_with_rel(link).ok())
.filter_map(|link_map| link_map.get("type").map(|item| item.raw_uri.clone()))
.collect::<BTreeSet<_>>();
if links.contains(vocab::ldp::BASIC_CONTAINER.as_str()) {
Some(ResourceType::RdfSource(Some(ContainerType::Basic)))
} else if links.contains(vocab::ldp::DIRECT_CONTAINER.as_str()) {
Some(ResourceType::RdfSource(Some(ContainerType::Direct)))
} else if links.contains(vocab::ldp::INDIRECT_CONTAINER.as_str()) {
Some(ResourceType::RdfSource(Some(ContainerType::Indirect)))
} else if links.contains(vocab::ldp::RDF_SOURCE.as_str()) {
Some(ResourceType::RdfSource(None))
} else if links.contains(vocab::ldp::NON_RDF_SOURCE.as_str()) {
Some(ResourceType::NonRdfSource)
} else {
None
}
}
fn add_media_types(&self, mut request_builder: RequestBuilder) -> RequestBuilder {
let media_types = self.builder.formats.iter().map(|f| f.media_type());
for media_type in media_types {
request_builder = request_builder.header(header::ACCEPT, media_type);
}
request_builder
}
fn is_method_allowed(response: &Response, method: Method) -> crate::Result<bool> {
for header in response.headers().get_all(header::ALLOW) {
for value in header.to_str()?.replace(" ", "").split(',') {
if value == method {
return Ok(true);
}
}
}
Ok(false)
}
pub async fn send(&self) -> crate::Result<Resource> {
let request_builder = self.builder.client.head(self.builder.url.clone());
let mut response = request_builder.send().await?;
match response.error_for_status_ref() {
Ok(response) => {
if self.builder.validate_support {
Self::ensure_ldp_support(response)?;
}
}
Err(err) if err.status() == Some(StatusCode::METHOD_NOT_ALLOWED) => {
if !Self::is_method_allowed(&response, Method::GET)? {
return Err(err.into());
}
}
err => {
err?;
}
}
let size = response
.headers()
.get(header::CONTENT_LENGTH)
.and_then(|hv| hv.to_str().ok())
.and_then(|et| usize::from_str(et).ok());
let content_disposition = response
.headers()
.get(header::CONTENT_DISPOSITION)
.and_then(|hv| hv.to_str().ok());
let file_name = if let Some(content_disposition) = content_disposition {
sfv::Parser::new(content_disposition)
.parse::<Item>()
.ok()
.and_then(|item| {
if item.bare_item.as_token() == Some(TokenRef::constant("attachment")) {
item.params
.get("filename")
.and_then(|item| item.as_string().map(|s| s.to_string()))
} else {
None
}
})
} else {
None
};
let resource_type = Self::resource_type(&response);
let url_to_get;
let described_by = Self::extract_described_by(&response);
if let Some(new_url) = &described_by
&& self.builder.follow_described_by
{
url_to_get = new_url.clone();
} else {
url_to_get = self.builder.url.clone();
}
let mut request_builder = self.builder.client.get(url_to_get);
request_builder = self.add_media_types(request_builder);
if !self.builder.include_preferences.is_empty() || !self.builder.omit_preferences.is_empty()
{
let header = prefer::header_for_preferences(
&self.builder.include_preferences,
&self.builder.omit_preferences,
)?;
request_builder = request_builder.headers(header);
}
response = request_builder.send().await?.error_for_status()?;
if self.builder.validate_support {
Self::ensure_ldp_support(&response)?;
}
let state_token = response
.headers()
.get(crate::header::X_STATE_TOKEN)
.and_then(|hv| hv.to_str().ok().map(|et| et.to_string()));
let format = response
.headers()
.get(header::CONTENT_TYPE)
.and_then(|hv| hv.to_str().ok())
.map(|value| {
if let Some(format) = RdfFormat::from_media_type(value) {
ResponseFormat::RdfFormat(format)
} else {
ResponseFormat::Other(value.to_string())
}
})
.unwrap_or(ResponseFormat::Unspecified);
Ok(Resource {
origin: self.builder.url.clone(),
described_by,
state_token,
resource_type,
format,
file_name,
size,
response,
})
}
}
#[derive(Clone, Debug)]
pub enum ResourceType {
NonRdfSource,
RdfSource(Option<ContainerType>),
}
impl ResourceType {
pub fn as_named_node_ref(&self) -> NamedNodeRef<'_> {
match self {
ResourceType::NonRdfSource => vocab::ldp::NON_RDF_SOURCE,
ResourceType::RdfSource(None) => vocab::ldp::RDF_SOURCE,
ResourceType::RdfSource(Some(container_type)) => container_type.as_named_node_ref(),
}
}
pub fn from_named_node_ref(node: NamedNodeRef<'_>) -> Option<Self> {
match node {
vocab::ldp::NON_RDF_SOURCE => Some(ResourceType::NonRdfSource),
vocab::ldp::RDF_SOURCE => Some(ResourceType::RdfSource(None)),
_ => {
let container_type = ContainerType::from_named_node_ref(node);
if container_type.is_some() {
Some(ResourceType::RdfSource(container_type))
} else {
None
}
}
}
}
}
pub enum ResponseFormat {
RdfFormat(RdfFormat),
Other(String),
Unspecified,
}
pub struct Resource {
origin: Url,
described_by: Option<Url>,
state_token: Option<String>,
resource_type: Option<ResourceType>,
format: ResponseFormat,
file_name: Option<String>,
size: Option<usize>,
response: Response,
}
impl Resource {
pub fn origin(&self) -> &Url {
&self.origin
}
pub fn described_by(&self) -> Option<&Url> {
self.described_by.as_ref()
}
pub fn format(&self) -> &ResponseFormat {
&self.format
}
pub fn size(&self) -> Option<usize> {
self.size
}
pub fn file_name(&self) -> Option<&str> {
self.file_name.as_deref()
}
pub fn resource_type(&self) -> Option<&ResourceType> {
self.resource_type.as_ref()
}
pub fn state_token(&self) -> Option<&String> {
self.state_token.as_ref()
}
pub fn into_stream(self) -> impl Stream<Item = reqwest_middleware::reqwest::Result<Bytes>> {
self.response.bytes_stream()
}
pub async fn into_rdf_source<D>(self) -> crate::Result<RdfSource<D>>
where
D: FromIterator<Quad>,
{
if let ResponseFormat::RdfFormat(format) = self.format {
let graph_url = self.described_by.as_ref().unwrap_or(&self.origin);
let graph = GraphNameRef::NamedNode(NamedNodeRef::new_unchecked(graph_url.as_str()));
let parser = RdfParser::from_format(format).with_default_graph(graph);
let body = self.response.bytes().await?;
let quads = parser.for_slice(&body);
let dataset = quads.collect::<Result<_, _>>()?;
Ok(RdfSource {
origin: self.origin,
origin_type: self.resource_type,
state_token: self.state_token,
described_by: self.described_by,
file_name: self.file_name,
size: self.size,
dataset,
})
} else {
Err(crate::Error::UnsupportedFormat)
}
}
}