use bytes::BufMut;
use http::{StatusCode, header};
use oxigraph::io::{RdfFormat, RdfSerializer};
use oxigraph::model::Dataset;
use reqwest_middleware::ClientWithMiddleware;
use reqwest_middleware::reqwest::Url;
#[derive(Clone, Debug)]
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_update(&self, format: RdfFormat) -> crate::Result<RdfSourceUpdateRequest> {
let url = self.described_by.clone().unwrap_or(self.origin.clone());
let media_type = format.media_type().to_string();
let body = self.serialize(format)?;
Ok(RdfSourceUpdateRequest {
url,
state_token: self.state_token.clone(),
media_type,
body,
})
}
}
pub struct RdfSourceUpdateRequest {
url: Url,
state_token: Option<String>,
media_type: String,
body: bytes::Bytes,
}
pub enum RdfSourceUpdateResponse {
Success,
DocumentModified(RdfSourceUpdateRequest),
}
impl RdfSourceUpdateRequest {
pub async fn send(
self,
client: ClientWithMiddleware,
overwrite: bool,
) -> crate::Result<RdfSourceUpdateResponse> {
if overwrite {
client
.put(self.url)
.header(header::CONTENT_TYPE, self.media_type)
.body(self.body)
.send()
.await?
.error_for_status()?;
Ok(RdfSourceUpdateResponse::Success)
} else {
let mut builder = client
.put(self.url.clone())
.header(header::CONTENT_TYPE, self.media_type.clone());
if let Some(state_token) = &self.state_token {
builder = builder.header(crate::header::X_IF_STATE_TOKEN, state_token.as_str());
}
let response = builder.body(self.body.clone()).send().await?;
match response.status() {
StatusCode::PRECONDITION_FAILED => {
Ok(RdfSourceUpdateResponse::DocumentModified(self))
}
_ => {
response.error_for_status()?;
Ok(RdfSourceUpdateResponse::Success)
}
}
}
}
}