use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use crate::kv::{ChangePublisher, KvStore};
pub const INVAL_PREFIX: &str = "_inval/";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct Entry {
writer: String,
keys: Vec<String>,
}
pub struct Changelog {
store: Arc<dyn KvStore>,
writer: String,
counter: AtomicU64,
retention_secs: u64,
}
impl Changelog {
pub fn new(store: Arc<dyn KvStore>, retention_secs: u64) -> Self {
Self {
store,
writer: random_writer_id(),
counter: AtomicU64::new(0),
retention_secs: retention_secs.max(1),
}
}
pub fn writer_id(&self) -> &str {
&self.writer
}
pub async fn current_cursor(&self) -> String {
self.list_entry_keys()
.await
.into_iter()
.max()
.unwrap_or_default()
}
pub async fn publish(&self, keys: &[String]) {
let keys: Vec<String> = keys
.iter()
.filter(|k| !k.starts_with(INVAL_PREFIX))
.cloned()
.collect();
if keys.is_empty() {
return;
}
let entry = Entry {
writer: self.writer.clone(),
keys,
};
if let Ok(value) = serde_json::to_vec(&entry) {
let key = self.entry_key();
let _ = self.store.put(&key, value).await;
}
}
pub async fn poll(&self, cursor: &mut String) -> Vec<String> {
let mut entry_keys: Vec<String> = self
.list_entry_keys()
.await
.into_iter()
.filter(|k| *k > *cursor)
.collect();
entry_keys.sort();
let mut changed = Vec::new();
for entry_key in &entry_keys {
if let Ok(Some(bytes)) = self.store.get(entry_key).await {
if let Ok(entry) = serde_json::from_slice::<Entry>(&bytes) {
if entry.writer != self.writer {
changed.extend(entry.keys);
}
}
}
}
if let Some(max) = entry_keys.into_iter().max() {
*cursor = max;
}
changed
}
pub async fn trim(&self) {
let cutoff = now_millis().saturating_sub(self.retention_secs * 1000);
for key in self.list_entry_keys().await {
if entry_millis(&key).is_some_and(|ms| ms < cutoff) {
let _ = self.store.delete(&key).await;
}
}
}
async fn list_entry_keys(&self) -> Vec<String> {
self.store
.list_prefix(INVAL_PREFIX)
.await
.unwrap_or_default()
}
fn entry_key(&self) -> String {
let n = self.counter.fetch_add(1, Ordering::Relaxed);
format!(
"{INVAL_PREFIX}{:013}-{}-{:020}",
now_millis(),
self.writer,
n
)
}
}
#[async_trait]
impl ChangePublisher for Changelog {
async fn publish(&self, keys: &[String]) {
Self::publish(self, keys).await;
}
}
fn now_millis() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
fn entry_millis(key: &str) -> Option<u64> {
key.strip_prefix(INVAL_PREFIX)?
.split('-')
.next()?
.parse()
.ok()
}
fn random_writer_id() -> String {
let mut bytes = [0u8; 8];
getrandom::getrandom(&mut bytes).expect("system RNG");
hex::encode(bytes)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::kv::{CachedKv, MemoryKv};
#[tokio::test]
async fn peer_write_invalidates_only_the_changed_key() {
let shared: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
shared.put("site/a/config", b"v1".to_vec()).await.unwrap();
shared.put("site/b/config", b"b1".to_vec()).await.unwrap();
let log_a = Arc::new(Changelog::new(shared.clone(), 60));
let cache_a: Arc<dyn KvStore> =
Arc::new(CachedKv::new(shared.clone(), 64).with_publisher(log_a.clone()));
let log_b = Arc::new(Changelog::new(shared.clone(), 60));
let cache_b: Arc<dyn KvStore> =
Arc::new(CachedKv::new(shared.clone(), 64).with_publisher(log_b.clone()));
let mut cursor_b = log_b.current_cursor().await;
assert_eq!(
cache_b.get("site/a/config").await.unwrap(),
Some(b"v1".to_vec())
);
assert_eq!(
cache_b.get("site/b/config").await.unwrap(),
Some(b"b1".to_vec())
);
cache_a.put("site/a/config", b"v2".to_vec()).await.unwrap();
assert_eq!(
cache_b.get("site/a/config").await.unwrap(),
Some(b"v1".to_vec())
);
let changed = log_b.poll(&mut cursor_b).await;
assert_eq!(changed, vec!["site/a/config".to_string()]);
cache_b.invalidate_keys(&changed);
assert_eq!(
cache_b.get("site/a/config").await.unwrap(),
Some(b"v2".to_vec())
);
assert_eq!(
cache_b.get("site/b/config").await.unwrap(),
Some(b"b1".to_vec())
);
}
#[tokio::test]
async fn poll_skips_own_writes() {
let shared: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let log = Arc::new(Changelog::new(shared.clone(), 60));
let cache: Arc<dyn KvStore> =
Arc::new(CachedKv::new(shared.clone(), 64).with_publisher(log.clone()));
let mut cursor = log.current_cursor().await;
cache.put("k", b"v".to_vec()).await.unwrap();
assert!(log.poll(&mut cursor).await.is_empty());
}
#[tokio::test]
async fn batch_publishes_one_entry_with_all_keys() {
let shared: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let writer = Arc::new(Changelog::new(shared.clone(), 60));
let cache: Arc<dyn KvStore> =
Arc::new(CachedKv::new(shared.clone(), 64).with_publisher(writer.clone()));
let reader = Arc::new(Changelog::new(shared.clone(), 60));
let mut cursor = reader.current_cursor().await;
cache
.write_batch(vec![
crate::kv::WriteOp::Put("current/x".into(), b"id".to_vec()),
crate::kv::WriteOp::Put("site/x/config".into(), b"c".to_vec()),
])
.await
.unwrap();
let mut changed = reader.poll(&mut cursor).await;
changed.sort();
assert_eq!(
changed,
vec!["current/x".to_string(), "site/x/config".to_string()]
);
assert_eq!(shared.list_prefix(INVAL_PREFIX).await.unwrap().len(), 1);
}
#[tokio::test]
async fn trim_drops_entries_outside_retention() {
let shared: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let log = Changelog::new(shared.clone(), 1);
log.publish(&["k".to_string()]).await;
assert_eq!(shared.list_prefix(INVAL_PREFIX).await.unwrap().len(), 1);
shared
.put(
&format!("{INVAL_PREFIX}0000000000001-old-0"),
b"{\"writer\":\"x\",\"keys\":[]}".to_vec(),
)
.await
.unwrap();
log.trim().await;
let remaining = shared.list_prefix(INVAL_PREFIX).await.unwrap();
assert!(
remaining.iter().all(|k| !k.contains("-old-")),
"ancient entry trimmed"
);
}
}