Add basic logic to fetch and update RDF Sources
This commit is contained in:
@@ -1,2 +1,3 @@
|
||||
/target
|
||||
/.idea
|
||||
/.cargo
|
||||
|
||||
Generated
+1844
File diff suppressed because it is too large
Load Diff
+16
@@ -4,3 +4,19 @@ members = [
|
||||
"ldp",
|
||||
"ldctl",
|
||||
]
|
||||
|
||||
[workspace.dependencies]
|
||||
ldp = { path = "ldp" }
|
||||
|
||||
async-trait = "0.1"
|
||||
base64 = "0.22"
|
||||
bytes = "1.11"
|
||||
futures = "0.3"
|
||||
http = "1.4"
|
||||
color-eyre = "0.6"
|
||||
oxigraph = "0.5"
|
||||
parse_link_header = "0.4"
|
||||
reqwest-middleware = { git = "https://github.com/TrueLayer/reqwest-middleware.git", features = ["stream"] }
|
||||
thiserror = "2"
|
||||
tokio = { version = "1", features = ["full"] }
|
||||
tracing = "0.1"
|
||||
@@ -4,3 +4,13 @@ version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
async-trait.workspace = true
|
||||
base64.workspace = true
|
||||
bytes.workspace = true
|
||||
futures.workspace = true
|
||||
http.workspace = true
|
||||
oxigraph.workspace = true
|
||||
parse_link_header.workspace = true
|
||||
reqwest-middleware.workspace = true
|
||||
thiserror.workspace = true
|
||||
tracing.workspace = true
|
||||
@@ -0,0 +1,27 @@
|
||||
use thiserror::Error;
|
||||
|
||||
pub type Result<T> = std::result::Result<T, Error>;
|
||||
|
||||
#[derive(Error, Debug)]
|
||||
pub enum Error {
|
||||
#[error(transparent)]
|
||||
Io(#[from] std::io::Error),
|
||||
|
||||
#[error(transparent)]
|
||||
Reqwest(#[from] reqwest_middleware::reqwest::Error),
|
||||
|
||||
#[error(transparent)]
|
||||
ReqwestMiddleware(#[from] reqwest_middleware::Error),
|
||||
|
||||
#[error(transparent)]
|
||||
InvalidHeaderValue(#[from] http::header::InvalidHeaderValue),
|
||||
|
||||
#[error("Server did not advertise LDP support")]
|
||||
LDPUnsupported,
|
||||
|
||||
#[error("Response was not in a supported RDF format")]
|
||||
UnsupportedFormat,
|
||||
|
||||
#[error("Document has been modified since last fetch, and overwrite was not enabled")]
|
||||
DocumentModified,
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
use reqwest_middleware::reqwest::header::HeaderName;
|
||||
|
||||
pub const X_STATE_TOKEN: HeaderName = HeaderName::from_static("x-state-token");
|
||||
pub const X_IF_STATE_TOKEN: HeaderName = HeaderName::from_static("x-if-state-token");
|
||||
pub const PREFER: HeaderName = HeaderName::from_static("prefer");
|
||||
pub const PREFERENCE_APPLIED: HeaderName = HeaderName::from_static("preference-applied");
|
||||
@@ -0,0 +1,16 @@
|
||||
mod error;
|
||||
pub mod header;
|
||||
pub mod middleware;
|
||||
mod rdf_source;
|
||||
mod resource;
|
||||
pub mod vocab;
|
||||
|
||||
pub use http::{HeaderName, HeaderValue};
|
||||
pub use oxigraph;
|
||||
pub use reqwest_middleware::ClientBuilder;
|
||||
pub use reqwest_middleware::reqwest::Client;
|
||||
pub use reqwest_middleware::reqwest::Url;
|
||||
|
||||
pub use error::Result;
|
||||
pub use rdf_source::RdfSource;
|
||||
pub use resource::{Resource, ResourceRequestBuilder};
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
use bytes::Bytes;
|
||||
use http::Extensions;
|
||||
use reqwest_middleware::reqwest::header::HeaderValue;
|
||||
use reqwest_middleware::reqwest::{Request, Response, header};
|
||||
use reqwest_middleware::{Middleware, Next};
|
||||
|
||||
pub struct BasicAuthMiddleware {
|
||||
username: String,
|
||||
password: Option<String>,
|
||||
}
|
||||
|
||||
impl BasicAuthMiddleware {
|
||||
pub fn new(username: String, password: Option<String>) -> Self {
|
||||
Self { username, password }
|
||||
}
|
||||
fn basic_auth(&self) -> HeaderValue {
|
||||
use base64::prelude::BASE64_STANDARD;
|
||||
use base64::write::EncoderWriter;
|
||||
use std::io::Write;
|
||||
|
||||
let mut buf = b"Basic ".to_vec();
|
||||
{
|
||||
let mut encoder = EncoderWriter::new(&mut buf, &BASE64_STANDARD);
|
||||
let _ = write!(encoder, "{}:", self.username);
|
||||
if let Some(password) = &self.password {
|
||||
let _ = write!(encoder, "{password}");
|
||||
}
|
||||
}
|
||||
let mut header = HeaderValue::from_maybe_shared(Bytes::from(buf))
|
||||
.expect("base64 is always valid HeaderValue");
|
||||
header.set_sensitive(true);
|
||||
header
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Middleware for BasicAuthMiddleware {
|
||||
async fn handle(
|
||||
&self,
|
||||
mut req: Request,
|
||||
extensions: &mut Extensions,
|
||||
next: Next<'_>,
|
||||
) -> reqwest_middleware::Result<Response> {
|
||||
req.headers_mut()
|
||||
.insert(header::AUTHORIZATION, self.basic_auth());
|
||||
next.run(req, extensions).await
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
use crate::{Url, error};
|
||||
use bytes::BufMut;
|
||||
use http::{HeaderValue, StatusCode, header};
|
||||
use oxigraph::io::{RdfFormat, RdfSerializer};
|
||||
use oxigraph::model::Dataset;
|
||||
use reqwest_middleware::ClientWithMiddleware;
|
||||
use reqwest_middleware::reqwest::Request;
|
||||
|
||||
pub struct RdfSource {
|
||||
pub(crate) origin: Url,
|
||||
pub(crate) described_by: Option<Url>,
|
||||
pub(crate) state_token: Option<String>,
|
||||
pub(crate) dataset: Dataset,
|
||||
}
|
||||
|
||||
impl RdfSource {
|
||||
pub fn origin(&self) -> &Url {
|
||||
&self.origin
|
||||
}
|
||||
|
||||
pub fn described_by(&self) -> Option<&Url> {
|
||||
self.described_by.as_ref()
|
||||
}
|
||||
|
||||
pub fn state_token(&self) -> Option<&str> {
|
||||
self.state_token.as_deref()
|
||||
}
|
||||
|
||||
pub fn dataset(&self) -> &Dataset {
|
||||
&self.dataset
|
||||
}
|
||||
|
||||
pub fn serialize(&self, format: RdfFormat) -> crate::Result<bytes::Bytes> {
|
||||
let writer = bytes::BytesMut::new().writer();
|
||||
let mut serializer = RdfSerializer::from_format(format).for_writer(writer);
|
||||
|
||||
if format.supports_datasets() {
|
||||
for quad in &self.dataset {
|
||||
serializer.serialize_quad(quad)?;
|
||||
}
|
||||
} else {
|
||||
for quad in &self.dataset {
|
||||
serializer.serialize_triple(quad)?;
|
||||
}
|
||||
}
|
||||
|
||||
let finished_writer = serializer.finish()?;
|
||||
Ok(finished_writer.into_inner().freeze())
|
||||
}
|
||||
|
||||
pub fn to_request(
|
||||
&self,
|
||||
client: ClientWithMiddleware,
|
||||
format: RdfFormat,
|
||||
) -> crate::Result<Request> {
|
||||
let url = self.described_by.clone().unwrap_or(self.origin.clone());
|
||||
let body = self.serialize(format)?;
|
||||
Ok(client
|
||||
.put(url)
|
||||
.header(header::CONTENT_TYPE, format.media_type())
|
||||
.body(body)
|
||||
.build()?)
|
||||
}
|
||||
|
||||
pub async fn send(
|
||||
&self,
|
||||
client: ClientWithMiddleware,
|
||||
mut request: Request,
|
||||
overwrite: bool,
|
||||
) -> crate::Result<()> {
|
||||
if !overwrite && let Some(state_token) = &self.state_token {
|
||||
let value = HeaderValue::from_str(state_token.as_str())?;
|
||||
request
|
||||
.headers_mut()
|
||||
.insert(crate::header::X_IF_STATE_TOKEN, value);
|
||||
}
|
||||
|
||||
let response = client.execute(request).await?;
|
||||
|
||||
match response.status() {
|
||||
StatusCode::PRECONDITION_FAILED => Err(error::Error::DocumentModified),
|
||||
_ => {
|
||||
response.error_for_status()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,185 @@
|
||||
use crate::rdf_source::RdfSource;
|
||||
use crate::{error, vocab};
|
||||
use bytes::Bytes;
|
||||
use futures::Stream;
|
||||
use oxigraph::io::{RdfFormat, RdfParser};
|
||||
use oxigraph::model::{Dataset, GraphNameRef, NamedNodeRef};
|
||||
use reqwest_middleware::reqwest::{Client, Response, StatusCode, Url, header};
|
||||
use reqwest_middleware::{ClientBuilder, ClientWithMiddleware, RequestBuilder};
|
||||
use tracing::error;
|
||||
|
||||
pub struct ResourceRequestBuilder {
|
||||
client: ClientWithMiddleware,
|
||||
url: Url,
|
||||
follow_described_by: bool,
|
||||
formats: Vec<RdfFormat>,
|
||||
}
|
||||
|
||||
impl ResourceRequestBuilder {
|
||||
pub fn new(url: Url) -> Self {
|
||||
Self::with_client_and_url(ClientBuilder::new(Client::new()).build(), url)
|
||||
}
|
||||
|
||||
pub fn with_client_and_url(client: ClientWithMiddleware, url: Url) -> Self {
|
||||
Self {
|
||||
client,
|
||||
url,
|
||||
follow_described_by: true,
|
||||
formats: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn follow_described_by(mut self, value: bool) -> Self {
|
||||
self.follow_described_by = value;
|
||||
self
|
||||
}
|
||||
|
||||
pub fn allow_format(mut self, format: RdfFormat) -> Self {
|
||||
self.formats.push(format);
|
||||
self
|
||||
}
|
||||
|
||||
pub fn allow_all_supported_formats(mut self) -> Self {
|
||||
self.formats = vec![
|
||||
RdfFormat::N3,
|
||||
RdfFormat::NQuads,
|
||||
RdfFormat::NTriples,
|
||||
RdfFormat::RdfXml,
|
||||
RdfFormat::TriG,
|
||||
RdfFormat::Turtle,
|
||||
];
|
||||
self
|
||||
}
|
||||
|
||||
pub async fn send(self) -> crate::Result<Resource> {
|
||||
Resource::from_builder(self).await
|
||||
}
|
||||
}
|
||||
|
||||
pub struct Resource {
|
||||
origin: Url,
|
||||
described_by: Option<Url>,
|
||||
state_token: Option<String>,
|
||||
format: Option<RdfFormat>,
|
||||
response: Response,
|
||||
}
|
||||
|
||||
impl Resource {
|
||||
/// Described by example:
|
||||
/// Link: <http://fedora.quill.lan/rest/E2/fcr:metadata>; rel="describedby"
|
||||
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("");
|
||||
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(error::Error::LDPUnsupported)
|
||||
}
|
||||
|
||||
fn add_media_types(
|
||||
formats: Vec<RdfFormat>,
|
||||
mut request_builder: RequestBuilder,
|
||||
) -> RequestBuilder {
|
||||
let media_types = formats.iter().map(|f| f.media_type());
|
||||
for media_type in media_types {
|
||||
request_builder = request_builder.header(header::ACCEPT, media_type);
|
||||
}
|
||||
request_builder
|
||||
}
|
||||
|
||||
pub async fn from_builder(builder: ResourceRequestBuilder) -> crate::Result<Self> {
|
||||
let request_builder = builder.client.head(builder.url.clone());
|
||||
let mut response = request_builder.send().await?.error_for_status()?;
|
||||
Resource::ensure_ldp_support(&response)?;
|
||||
|
||||
let url_to_get;
|
||||
let described_by = Self::extract_described_by(&response);
|
||||
if let Some(new_url) = &described_by
|
||||
&& builder.follow_described_by
|
||||
{
|
||||
url_to_get = new_url.clone();
|
||||
} else {
|
||||
url_to_get = builder.url.clone();
|
||||
}
|
||||
|
||||
let mut request_builder = builder.client.get(url_to_get);
|
||||
if !builder.formats.is_empty() {
|
||||
request_builder = Self::add_media_types(builder.formats, request_builder);
|
||||
}
|
||||
response = request_builder.send().await?.error_for_status()?;
|
||||
Resource::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)
|
||||
.map(|hv| hv.to_str().unwrap_or_default())
|
||||
.and_then(RdfFormat::from_media_type);
|
||||
|
||||
Ok(Self {
|
||||
origin: builder.url,
|
||||
described_by,
|
||||
state_token,
|
||||
format,
|
||||
response,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn into_stream(self) -> impl Stream<Item = reqwest_middleware::reqwest::Result<Bytes>> {
|
||||
self.response.bytes_stream()
|
||||
}
|
||||
|
||||
pub async fn into_rdf_source(self) -> crate::Result<RdfSource> {
|
||||
if let Some(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.filter_map(Result::ok).collect::<Dataset>();
|
||||
Ok(RdfSource {
|
||||
origin: self.origin,
|
||||
described_by: self.described_by,
|
||||
state_token: self.state_token,
|
||||
dataset,
|
||||
})
|
||||
} else {
|
||||
Err(error::Error::UnsupportedFormat)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
pub mod ldp {
|
||||
use oxigraph::model::NamedNodeRef;
|
||||
|
||||
pub const CONTAINS: NamedNodeRef<'_> =
|
||||
NamedNodeRef::new_unchecked("http://www.w3.org/ns/ldp#contains");
|
||||
pub const RESOURCE: NamedNodeRef<'_> =
|
||||
NamedNodeRef::new_unchecked("http://www.w3.org/ns/ldp#Resource");
|
||||
}
|
||||
Reference in New Issue
Block a user