use std::sync::Arc;
use std::sync::atomic::{AtomicI64, Ordering};
use dashmap::DashMap;
use crate::options::RemoveByTagBehavior;
use crate::time::Timestamp;
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct Tag(Arc<str>);
impl Tag {
#[must_use]
pub fn new(tag: impl AsRef<str>) -> Option<Self> {
let s = tag.as_ref();
if s.trim().is_empty() {
None
} else {
Some(Self(Arc::from(s)))
}
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
impl std::fmt::Display for Tag {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
#[must_use]
pub fn collect_tags<I, S>(tags: I) -> Box<[Tag]>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
tags.into_iter().filter_map(Tag::new).collect()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TagVerdict {
Valid,
Expire,
Remove,
}
const NEVER: i64 = i64::MIN;
#[derive(Debug)]
pub struct TagRegistry {
markers: DashMap<Tag, i64>,
clear_expire: AtomicI64,
clear_remove: AtomicI64,
}
impl Default for TagRegistry {
fn default() -> Self {
Self {
markers: DashMap::new(),
clear_expire: AtomicI64::new(NEVER),
clear_remove: AtomicI64::new(NEVER),
}
}
}
impl TagRegistry {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn mark_tag(&self, tag: Tag, at: Timestamp) {
self.markers
.entry(tag)
.and_modify(|existing| *existing = (*existing).max(at.ticks()))
.or_insert_with(|| at.ticks());
}
#[must_use]
pub fn tag_marker(&self, tag: &Tag) -> Option<Timestamp> {
self.markers.get(tag).map(|v| Timestamp::from_ticks(*v))
}
pub fn mark_clear_expire(&self, at: Timestamp) {
bump_max(&self.clear_expire, at.ticks());
}
pub fn mark_clear_remove(&self, at: Timestamp) {
bump_max(&self.clear_remove, at.ticks());
}
#[must_use]
pub fn evaluate(
&self,
entry_created: Timestamp,
tags: &[Tag],
behavior: RemoveByTagBehavior,
) -> TagVerdict {
let created = entry_created.ticks();
if created < self.clear_remove.load(Ordering::Relaxed) {
return TagVerdict::Remove;
}
for tag in tags {
if let Some(marker) = self.markers.get(tag).map(|m| *m)
&& created < marker
{
return match behavior {
RemoveByTagBehavior::Expire => TagVerdict::Expire,
RemoveByTagBehavior::Remove => TagVerdict::Remove,
};
}
}
if created < self.clear_expire.load(Ordering::Relaxed) {
return TagVerdict::Expire;
}
TagVerdict::Valid
}
}
fn bump_max(slot: &AtomicI64, value: i64) {
let mut current = slot.load(Ordering::Relaxed);
while value > current {
match slot.compare_exchange_weak(current, value, Ordering::Relaxed, Ordering::Relaxed) {
Ok(_) => break,
Err(observed) => current = observed,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn tag(s: &str) -> Tag {
Tag::new(s).unwrap()
}
#[test]
fn blank_tags_rejected() {
assert!(Tag::new("").is_none());
assert!(Tag::new(" ").is_none());
assert!(Tag::new("x").is_some());
}
#[test]
fn entry_before_marker_is_invalid() {
let reg = TagRegistry::new();
let created = Timestamp::from_ticks(100);
reg.mark_tag(tag("a"), Timestamp::from_ticks(200));
assert_eq!(
reg.evaluate(created, &[tag("a")], RemoveByTagBehavior::Expire),
TagVerdict::Expire
);
let newer = Timestamp::from_ticks(300);
assert_eq!(
reg.evaluate(newer, &[tag("a")], RemoveByTagBehavior::Expire),
TagVerdict::Valid
);
}
#[test]
fn entry_at_marker_tick_survives() {
let reg = TagRegistry::new();
let at = Timestamp::from_ticks(200);
reg.mark_tag(tag("a"), at);
assert_eq!(
reg.evaluate(at, &[tag("a")], RemoveByTagBehavior::Expire),
TagVerdict::Valid,
"an entry created in the marker tick is kept"
);
assert_eq!(
reg.evaluate(
Timestamp::from_ticks(199),
&[tag("a")],
RemoveByTagBehavior::Expire
),
TagVerdict::Expire,
);
}
#[test]
fn remove_behavior_dominates() {
let reg = TagRegistry::new();
reg.mark_tag(tag("a"), Timestamp::from_ticks(200));
assert_eq!(
reg.evaluate(
Timestamp::from_ticks(100),
&[tag("a")],
RemoveByTagBehavior::Remove
),
TagVerdict::Remove
);
}
#[test]
fn clear_remove_outranks_clear_expire() {
let reg = TagRegistry::new();
reg.mark_clear_expire(Timestamp::from_ticks(200));
reg.mark_clear_remove(Timestamp::from_ticks(200));
assert_eq!(
reg.evaluate(Timestamp::from_ticks(100), &[], RemoveByTagBehavior::Expire),
TagVerdict::Remove
);
}
#[test]
fn clear_markers_are_strict_at_the_boundary() {
let reg = TagRegistry::new();
let at = Timestamp::from_ticks(200);
reg.mark_clear_expire(at);
reg.mark_clear_remove(at);
assert_eq!(
reg.evaluate(at, &[], RemoveByTagBehavior::Expire),
TagVerdict::Valid,
);
assert_eq!(
reg.evaluate(Timestamp::from_ticks(199), &[], RemoveByTagBehavior::Expire),
TagVerdict::Remove,
);
}
}