use oxigraph::model::{GraphNameRef, Literal, NamedNode, NamedOrBlankNode, Quad, QuadRef, Term};
use oxigraph::store::Store;
use crate::error::CacheError;
use crate::key::KeyRef;
use crate::tier::TierId;
use crate::tier::backends::ByteValue;
use crate::tier::tier_trait::{BackendKind, CacheTier, TierHealth};
const ENTRY_NS: &str = "urn:thesix:entry:";
const VALUE_PRED: &str = "urn:thesix:value";
fn hex_encode(bytes: &[u8]) -> String {
use std::fmt::Write;
let mut out = String::with_capacity(bytes.len() * 2);
for b in bytes {
let _ = write!(out, "{b:02x}");
}
out
}
fn hex_decode(text: &str) -> Result<Vec<u8>, CacheError> {
if !text.len().is_multiple_of(2) {
return Err(CacheError::SerializationFailed);
}
let bytes = text.as_bytes();
let mut out = Vec::with_capacity(bytes.len() / 2);
for pair in bytes.chunks(2) {
let hi = (pair[0] as char)
.to_digit(16)
.ok_or(CacheError::SerializationFailed)?;
let lo = (pair[1] as char)
.to_digit(16)
.ok_or(CacheError::SerializationFailed)?;
out.push(((hi << 4) | lo) as u8);
}
Ok(out)
}
pub struct L5OxigraphBackend<V> {
store: Store,
_marker: std::marker::PhantomData<V>,
}
impl<V> std::fmt::Debug for L5OxigraphBackend<V> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("L5OxigraphBackend").finish_non_exhaustive()
}
}
impl<V: ByteValue> L5OxigraphBackend<V> {
pub fn new() -> Result<Self, CacheError> {
Ok(Self {
store: Store::new().map_err(map_store_error)?,
_marker: std::marker::PhantomData,
})
}
fn subject_for(key: &KeyRef<'_>) -> Result<NamedNode, CacheError> {
let mut iri = String::with_capacity(key.0.len() * 2 + ENTRY_NS.len());
iri.push_str(ENTRY_NS);
iri.push_str(&hex_encode(key.0));
NamedNode::new(iri).map_err(map_store_error)
}
fn value_predicate() -> Result<NamedNode, CacheError> {
NamedNode::new(VALUE_PRED).map_err(map_store_error)
}
fn find_value(&self, key: &KeyRef<'_>) -> Result<Option<Vec<u8>>, CacheError> {
let subject = Self::subject_for(key)?;
let pattern_subject: NamedOrBlankNode = subject.clone().into();
for found in self
.store
.quads_for_pattern(Some(pattern_subject.as_ref()), None, None, None)
{
let q = found.map_err(map_store_error)?;
if let NamedOrBlankNode::NamedNode(node) = &q.subject
&& node.as_str() == subject.as_str()
{
return Ok(Some(match &q.object {
Term::Literal(l) => l.value().as_bytes().to_vec(),
_ => Vec::new(),
}));
}
}
Ok(None)
}
}
fn map_store_error(e: impl std::fmt::Display) -> CacheError {
let _ = e;
CacheError::PopulationFailed
}
#[async_trait::async_trait]
impl<V: ByteValue> CacheTier<V> for L5OxigraphBackend<V> {
fn name(&self) -> String {
"L5-oxigraph".into()
}
fn backend(&self) -> BackendKind {
BackendKind::Oxigraph
}
fn capability(&self) -> crate::capability::TierCapability {
crate::capability::TierCapability::new(
BackendKind::Oxigraph,
crate::capability::CapabilityFlags::BLOCKING_IO,
crate::capability::OperationalState::Healthy,
crate::capability::DurabilityClass::Volatile,
)
}
async fn get(&self, key: &KeyRef<'_>) -> Result<Option<V>, CacheError> {
match self.find_value(key)? {
Some(encoded) => {
let text =
std::str::from_utf8(&encoded).map_err(|_| CacheError::SerializationFailed)?;
V::decode_bytes(&hex_decode(text)?).map(Some)
}
None => Ok(None),
}
}
async fn set(
&self,
key: &KeyRef<'_>,
value: V,
_ttl: Option<std::time::Duration>,
) -> Result<(), CacheError> {
let bytes = value.encode_bytes()?;
let subject = Self::subject_for(key)?;
let predicate = Self::value_predicate()?;
let object = Literal::new_simple_literal(hex_encode(&bytes));
let quad = QuadRef::new(
subject.as_ref(),
predicate.as_ref(),
object.as_ref(),
GraphNameRef::DefaultGraph,
);
self.store.remove(quad).map_err(map_store_error)?;
self.store.insert(quad).map_err(map_store_error)?;
Ok(())
}
async fn remove(&self, key: &KeyRef<'_>) -> Result<(), CacheError> {
let subject = Self::subject_for(key)?;
let pattern_subject: NamedOrBlankNode = subject.into();
let mut to_remove: Vec<Quad> = Vec::new();
for found in self
.store
.quads_for_pattern(Some(pattern_subject.as_ref()), None, None, None)
{
to_remove.push(found.map_err(map_store_error)?);
}
for q in &to_remove {
self.store.remove(q).map_err(map_store_error)?;
}
Ok(())
}
async fn contains(&self, key: &KeyRef<'_>) -> Result<bool, CacheError> {
Ok(self.find_value(key)?.is_some())
}
fn health(&self) -> TierHealth {
TierHealth::default()
}
fn tier_id(&self) -> TierId {
TierId::L5
}
}