1use std::sync::Arc;
15use std::sync::atomic::{AtomicI64, Ordering};
16
17use dashmap::DashMap;
18
19use crate::options::RemoveByTagBehavior;
20use crate::time::Timestamp;
21
22#[derive(Debug, Clone, PartialEq, Eq, Hash)]
27pub struct Tag(Arc<str>);
28
29impl Tag {
30 #[must_use]
32 pub fn new(tag: impl AsRef<str>) -> Option<Self> {
33 let s = tag.as_ref();
34 if s.trim().is_empty() {
35 None
36 } else {
37 Some(Self(Arc::from(s)))
38 }
39 }
40
41 #[must_use]
43 pub fn as_str(&self) -> &str {
44 &self.0
45 }
46}
47
48impl std::fmt::Display for Tag {
49 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
50 f.write_str(&self.0)
51 }
52}
53
54#[must_use]
56pub fn collect_tags<I, S>(tags: I) -> Box<[Tag]>
57where
58 I: IntoIterator<Item = S>,
59 S: AsRef<str>,
60{
61 tags.into_iter().filter_map(Tag::new).collect()
62}
63
64#[derive(Debug, Clone, Copy, PartialEq, Eq)]
66pub enum TagVerdict {
67 Valid,
69 Expire,
71 Remove,
73}
74
75const NEVER: i64 = i64::MIN;
76
77#[derive(Debug)]
79pub struct TagRegistry {
80 markers: DashMap<Tag, i64>,
82 clear_expire: AtomicI64,
84 clear_remove: AtomicI64,
86}
87
88impl Default for TagRegistry {
89 fn default() -> Self {
90 Self {
91 markers: DashMap::new(),
92 clear_expire: AtomicI64::new(NEVER),
93 clear_remove: AtomicI64::new(NEVER),
94 }
95 }
96}
97
98impl TagRegistry {
99 #[must_use]
101 pub fn new() -> Self {
102 Self::default()
103 }
104
105 pub fn mark_tag(&self, tag: Tag, at: Timestamp) {
108 self.markers
109 .entry(tag)
110 .and_modify(|existing| *existing = (*existing).max(at.ticks()))
111 .or_insert_with(|| at.ticks());
112 }
113
114 #[must_use]
116 pub fn tag_marker(&self, tag: &Tag) -> Option<Timestamp> {
117 self.markers.get(tag).map(|v| Timestamp::from_ticks(*v))
118 }
119
120 pub fn mark_clear_expire(&self, at: Timestamp) {
123 bump_max(&self.clear_expire, at.ticks());
124 }
125
126 pub fn mark_clear_remove(&self, at: Timestamp) {
129 bump_max(&self.clear_remove, at.ticks());
130 }
131
132 #[must_use]
141 pub fn evaluate(
142 &self,
143 entry_created: Timestamp,
144 tags: &[Tag],
145 behavior: RemoveByTagBehavior,
146 ) -> TagVerdict {
147 let created = entry_created.ticks();
148
149 if created < self.clear_remove.load(Ordering::Relaxed) {
150 return TagVerdict::Remove;
151 }
152
153 for tag in tags {
154 if let Some(marker) = self.markers.get(tag).map(|m| *m)
155 && created < marker
156 {
157 return match behavior {
158 RemoveByTagBehavior::Expire => TagVerdict::Expire,
159 RemoveByTagBehavior::Remove => TagVerdict::Remove,
160 };
161 }
162 }
163
164 if created < self.clear_expire.load(Ordering::Relaxed) {
165 return TagVerdict::Expire;
166 }
167
168 TagVerdict::Valid
169 }
170}
171
172fn bump_max(slot: &AtomicI64, value: i64) {
173 let mut current = slot.load(Ordering::Relaxed);
174 while value > current {
175 match slot.compare_exchange_weak(current, value, Ordering::Relaxed, Ordering::Relaxed) {
176 Ok(_) => break,
177 Err(observed) => current = observed,
178 }
179 }
180}
181
182#[cfg(test)]
183mod tests {
184 use super::*;
185
186 fn tag(s: &str) -> Tag {
187 Tag::new(s).unwrap()
188 }
189
190 #[test]
191 fn blank_tags_rejected() {
192 assert!(Tag::new("").is_none());
193 assert!(Tag::new(" ").is_none());
194 assert!(Tag::new("x").is_some());
195 }
196
197 #[test]
198 fn entry_before_marker_is_invalid() {
199 let reg = TagRegistry::new();
200 let created = Timestamp::from_ticks(100);
201 reg.mark_tag(tag("a"), Timestamp::from_ticks(200));
202 assert_eq!(
203 reg.evaluate(created, &[tag("a")], RemoveByTagBehavior::Expire),
204 TagVerdict::Expire
205 );
206 let newer = Timestamp::from_ticks(300);
208 assert_eq!(
209 reg.evaluate(newer, &[tag("a")], RemoveByTagBehavior::Expire),
210 TagVerdict::Valid
211 );
212 }
213
214 #[test]
215 fn entry_at_marker_tick_survives() {
216 let reg = TagRegistry::new();
220 let at = Timestamp::from_ticks(200);
221 reg.mark_tag(tag("a"), at);
222 assert_eq!(
223 reg.evaluate(at, &[tag("a")], RemoveByTagBehavior::Expire),
224 TagVerdict::Valid,
225 "an entry created in the marker tick is kept"
226 );
227 assert_eq!(
229 reg.evaluate(
230 Timestamp::from_ticks(199),
231 &[tag("a")],
232 RemoveByTagBehavior::Expire
233 ),
234 TagVerdict::Expire,
235 );
236 }
237
238 #[test]
239 fn remove_behavior_dominates() {
240 let reg = TagRegistry::new();
241 reg.mark_tag(tag("a"), Timestamp::from_ticks(200));
242 assert_eq!(
243 reg.evaluate(
244 Timestamp::from_ticks(100),
245 &[tag("a")],
246 RemoveByTagBehavior::Remove
247 ),
248 TagVerdict::Remove
249 );
250 }
251
252 #[test]
253 fn clear_remove_outranks_clear_expire() {
254 let reg = TagRegistry::new();
255 reg.mark_clear_expire(Timestamp::from_ticks(200));
256 reg.mark_clear_remove(Timestamp::from_ticks(200));
257 assert_eq!(
258 reg.evaluate(Timestamp::from_ticks(100), &[], RemoveByTagBehavior::Expire),
259 TagVerdict::Remove
260 );
261 }
262
263 #[test]
264 fn clear_markers_are_strict_at_the_boundary() {
265 let reg = TagRegistry::new();
266 let at = Timestamp::from_ticks(200);
267 reg.mark_clear_expire(at);
268 reg.mark_clear_remove(at);
269 assert_eq!(
271 reg.evaluate(at, &[], RemoveByTagBehavior::Expire),
272 TagVerdict::Valid,
273 );
274 assert_eq!(
276 reg.evaluate(Timestamp::from_ticks(199), &[], RemoveByTagBehavior::Expire),
277 TagVerdict::Remove,
278 );
279 }
280}