use std::sync::Arc;
use crate::connector::cache_backend::CacheBackend;
use crate::errors::OrionError;
pub const MAX_NAMESPACES: usize = 8;
pub const MAX_NAMESPACE_LEN: usize = 64;
pub fn version_key(namespace: &str) -> String {
format!("orion:rc:ns:{namespace}")
}
pub fn entry_key(plain_key: &str) -> String {
match plain_key.strip_prefix("cache:") {
Some(rest) => format!("cache-ns:{rest}"),
None => format!("cache-ns:{plain_key}"),
}
}
pub fn check_name(name: &str) -> Result<(), String> {
if name.is_empty() || name.len() > MAX_NAMESPACE_LEN {
return Err(format!(
"namespace '{name}' must be 1 to {MAX_NAMESPACE_LEN} characters"
));
}
if !name
.bytes()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b"_-.:".contains(&b))
{
return Err(format!(
"namespace '{name}' may contain only lowercase letters, digits, '_', '-', '.' and ':'"
));
}
Ok(())
}
pub fn encode_entry(versions: &[i64], body: &str) -> String {
let mut out = String::with_capacity(body.len() + versions.len() * 4 + 1);
for (i, v) in versions.iter().enumerate() {
if i > 0 {
out.push(',');
}
out.push_str(&v.to_string());
}
out.push('\n');
out.push_str(body);
out
}
pub fn decode_entry(mut stored: String, versions: &[i64]) -> Option<String> {
let newline = stored.find('\n')?;
let header = &stored[..newline];
let mut stored_versions = header.split(',').map(|v| v.parse::<i64>());
let current = versions
.iter()
.all(|v| stored_versions.next().and_then(Result::ok) == Some(*v))
&& stored_versions.next().is_none();
if !current {
return None;
}
stored.drain(..=newline);
Some(stored)
}
pub fn parse_version(raw: Option<&str>) -> Option<i64> {
match raw {
None => Some(0),
Some(s) => s.trim().parse().ok(),
}
}
pub async fn invalidate(
targets: &[Arc<dyn CacheBackend>],
namespaces: &[String],
source: &'static str,
) -> Result<(), OrionError> {
let keys: Vec<String> = namespaces.iter().map(|ns| version_key(ns)).collect();
let bumps = targets.iter().map(|backend| {
let keys = &keys;
async move {
for key in keys {
backend.incr_by(key, 1, None).await?;
}
Ok::<(), OrionError>(())
}
});
let mut first_error = None;
for result in futures::future::join_all(bumps).await {
if let Err(e) = result {
tracing::warn!(
namespaces = ?namespaces,
error = %e,
"Failed to bump a response-cache namespace version"
);
first_error.get_or_insert(e);
}
}
crate::metrics::record_cache_invalidations(source, namespaces.len() as u64);
first_error.map_or(Ok(()), Err)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn names_are_checked() {
assert!(check_name("ladder").is_ok());
assert!(check_name("game:chess.v2_x-y").is_ok());
assert!(check_name("").is_err());
assert!(check_name("Ladder").is_err());
assert!(check_name("a b").is_err());
assert!(check_name(&"x".repeat(65)).is_err());
}
#[test]
fn a_namespaced_entry_never_shares_a_key_with_a_plain_one() {
assert_eq!(entry_key("cache:orders:abc"), "cache-ns:orders:abc");
assert_ne!(entry_key("cache:orders:abc"), "cache:orders:abc");
}
#[test]
fn entries_round_trip_only_under_their_own_versions() {
let body = r#"{"data":{"q":"a\nb"}}"#;
let stored = encode_entry(&[3, 0], body);
assert_eq!(decode_entry(stored.clone(), &[3, 0]).as_deref(), Some(body));
assert_eq!(decode_entry(stored.clone(), &[4, 0]), None);
assert_eq!(decode_entry(stored.clone(), &[3]), None);
assert_eq!(decode_entry(stored, &[3, 0, 1]), None);
assert_eq!(decode_entry("no header".to_string(), &[0]), None);
}
#[test]
fn versions_parse() {
assert_eq!(parse_version(None), Some(0));
assert_eq!(parse_version(Some("7")), Some(7));
assert_eq!(parse_version(Some("\"x\"")), None);
}
#[tokio::test]
async fn invalidate_bumps_each_namespace_in_each_store() {
use crate::connector::cache_backend::MemoryCacheBackend;
let a: Arc<dyn CacheBackend> = MemoryCacheBackend::new(60, 0);
let b: Arc<dyn CacheBackend> = MemoryCacheBackend::new(60, 0);
let ns = vec!["ladder".to_string(), "season".to_string()];
invalidate(&[a.clone(), b.clone()], &ns, "workflow")
.await
.expect("test");
assert_eq!(
a.get(&version_key("ladder")).await.expect("test"),
Some("1".to_string())
);
assert_eq!(
b.get(&version_key("season")).await.expect("test"),
Some("1".to_string())
);
}
}