use std::collections::HashMap;
use std::sync::Mutex;
use std::time::{Duration, Instant};
use arrow::record_batch::RecordBatch;
use super::error::{MAX_ERROR_CHARS, MAX_FEED_TEXT_CHARS, truncate};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum FeedStatus {
#[default]
Never,
Fresh,
Revalidated,
StaleError,
Error,
}
impl FeedStatus {
pub fn as_str(&self) -> &'static str {
match self {
FeedStatus::Never => "never",
FeedStatus::Fresh => "fresh",
FeedStatus::Revalidated => "revalidated",
FeedStatus::StaleError => "stale-error",
FeedStatus::Error => "error",
}
}
pub fn window_status_str(&self) -> Option<&'static str> {
match self {
FeedStatus::Fresh => Some("fresh"),
FeedStatus::Revalidated => Some("revalidated"),
FeedStatus::StaleError => Some("stale-error"),
FeedStatus::Never | FeedStatus::Error => None,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct FeedObservation {
pub last_fetch_ms: Option<i64>,
pub last_status: FeedStatus,
pub http_status: Option<u16>,
pub last_error: Option<String>,
pub dialect: Option<String>,
pub dialect_declared: Option<String>,
pub conformance_notes: Option<String>,
pub title: Option<String>,
pub site_url: Option<String>,
pub description: Option<String>,
pub item_count: Option<u64>,
}
impl FeedObservation {
pub fn capped(mut self) -> Self {
fn cap(text: Option<String>, max_chars: usize) -> Option<String> {
text.map(|t| truncate(&t, max_chars))
}
self.last_error = cap(self.last_error, MAX_ERROR_CHARS);
self.dialect_declared = cap(self.dialect_declared, MAX_ERROR_CHARS);
self.conformance_notes = cap(self.conformance_notes, MAX_ERROR_CHARS);
self.site_url = cap(self.site_url, MAX_ERROR_CHARS);
self.title = cap(self.title, MAX_FEED_TEXT_CHARS);
self.description = cap(self.description, MAX_FEED_TEXT_CHARS);
self
}
}
#[derive(Debug, Clone)]
pub struct CachedWindow {
pub batch: RecordBatch,
pub etag: Option<String>,
pub last_modified: Option<String>,
pub fetched_from: String,
pub probed_at: Instant,
}
#[derive(Debug, Clone)]
pub struct FeedSnapshot {
pub observation: FeedObservation,
pub window: Option<CachedWindow>,
pub within_ttl: bool,
pub generation: u64,
}
pub trait FeedCache: Send + Sync {
fn snapshot(&self, feed: &str, now: Instant) -> FeedSnapshot;
fn record_success(
&self,
feed: &str,
expected_generation: u64,
window: CachedWindow,
observation: FeedObservation,
armed_until: Instant,
);
fn record_not_modified(
&self,
feed: &str,
expected_generation: u64,
http_status: u16,
last_fetch_ms: i64,
armed_until: Instant,
);
fn record_failure(
&self,
feed: &str,
expected_generation: u64,
http_status: Option<u16>,
error: String,
dialect_declared: Option<String>,
last_fetch_ms: i64,
armed_until: Instant,
);
fn record_egress_denial(
&self,
feed: &str,
expected_generation: u64,
error: String,
last_fetch_ms: i64,
armed_until: Instant,
);
}
pub const WINDOW_EVICTED_ON_REVALIDATION: &str = "revalidated (304) but the cached window had already been evicted from the feed \
cache; the next attempt refetches it unconditionally";
pub fn failure_fuse(ttl: Duration) -> Duration {
(ttl / 4).clamp(Duration::from_secs(30), Duration::from_secs(300))
}
const MAX_OBSERVATIONS_MULTIPLIER: usize = 8;
const MIN_OBSERVATIONS: usize = 8;
struct WindowEntry {
window: CachedWindow,
bytes: usize,
}
struct Entry {
observation: FeedObservation,
window: Option<WindowEntry>,
armed_until: Instant,
generation: u64,
last_used: u64,
window_last_used: u64,
}
impl Entry {
fn new(armed_until: Instant) -> Self {
Self {
observation: FeedObservation::default(),
window: None,
armed_until,
generation: 0,
last_used: 0,
window_last_used: 0,
}
}
}
struct Inner {
map: HashMap<String, Entry>,
window_bytes: usize,
windowed: usize,
clock: u64,
commit_counter: u64,
}
impl Inner {
fn tick(&mut self) -> u64 {
self.clock += 1;
self.clock
}
fn next_generation(&mut self) -> u64 {
self.commit_counter += 1;
self.commit_counter
}
fn store_window(&mut self, feed: &str, window: CachedWindow, bytes: usize) {
let tick = self.tick();
let Some(entry) = self.map.get_mut(feed) else {
return;
};
entry.window_last_used = tick;
let previous = entry.window.replace(WindowEntry { window, bytes });
match previous {
Some(old) => self.window_bytes = self.window_bytes.saturating_sub(old.bytes),
None => self.windowed += 1,
}
self.window_bytes += bytes;
}
fn drop_window(&mut self, feed: &str) {
let Some(entry) = self.map.get_mut(feed) else {
return;
};
let Some(w) = entry.window.take() else {
return;
};
self.window_bytes = self.window_bytes.saturating_sub(w.bytes);
self.windowed -= 1;
}
fn touch_window(&mut self, feed: &str) {
let tick = self.tick();
if let Some(entry) = self.map.get_mut(feed) {
entry.window_last_used = tick;
}
}
fn touch_entry(&mut self, feed: &str) {
let tick = self.tick();
if let Some(entry) = self.map.get_mut(feed) {
entry.last_used = tick;
}
}
fn oldest_windowed(&self) -> Option<String> {
self.map
.iter()
.filter(|(_, entry)| entry.window.is_some())
.min_by_key(|(_, entry)| entry.window_last_used)
.map(|(feed, _)| feed.clone())
}
fn oldest_entry(&self) -> Option<String> {
self.map
.iter()
.min_by_key(|(_, entry)| entry.last_used)
.map(|(feed, _)| feed.clone())
}
fn evict(&mut self, max_bytes: usize, max_entries: usize) {
while self.windowed > max_entries || self.window_bytes > max_bytes {
let Some(oldest) = self.oldest_windowed() else {
break;
};
self.drop_window(&oldest);
}
}
fn evict_observations(&mut self, max_observations: usize) {
while self.map.len() > max_observations {
let Some(oldest) = self.oldest_entry() else {
break;
};
self.drop_window(&oldest);
self.map.remove(&oldest);
}
}
}
pub struct MemoryFeedCache {
inner: Mutex<Inner>,
max_bytes: usize,
max_entries: usize,
max_observations: usize,
}
impl MemoryFeedCache {
pub fn new(max_bytes: usize, max_entries: usize) -> Self {
Self {
inner: Mutex::new(Inner {
map: HashMap::new(),
window_bytes: 0,
windowed: 0,
clock: 0,
commit_counter: 0,
}),
max_bytes,
max_entries,
max_observations: max_entries
.saturating_mul(MAX_OBSERVATIONS_MULTIPLIER)
.max(MIN_OBSERVATIONS),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, Inner> {
self.inner.lock().unwrap_or_else(|p| p.into_inner())
}
}
impl FeedCache for MemoryFeedCache {
fn snapshot(&self, feed: &str, now: Instant) -> FeedSnapshot {
let mut inner = self.lock();
let Some(entry) = inner.map.get(feed) else {
return FeedSnapshot {
observation: FeedObservation::default(),
window: None,
within_ttl: false,
generation: 0,
};
};
let within_ttl = now < entry.armed_until;
let observation = entry.observation.clone();
let window = entry.window.as_ref().map(|w| w.window.clone());
let generation = entry.generation;
let has_window = window.is_some();
if has_window {
inner.touch_window(feed);
}
inner.touch_entry(feed);
FeedSnapshot {
observation,
window,
within_ttl,
generation,
}
}
fn record_success(
&self,
feed: &str,
expected_generation: u64,
window: CachedWindow,
observation: FeedObservation,
armed_until: Instant,
) {
let bytes = window.batch.get_array_memory_size();
let mut inner = self.lock();
let (current, status) = inner.map.get(feed).map_or((0, None), |e| {
(e.generation, Some(e.observation.last_status))
});
if current != expected_generation
&& !matches!(
status,
None | Some(FeedStatus::Error | FeedStatus::StaleError)
)
{
tracing::debug!(
feed,
commit = "success",
expected_generation,
current_generation = current,
"rss stale cache commit dropped"
);
return;
}
let generation = inner.next_generation();
inner.drop_window(feed);
let entry = inner
.map
.entry(feed.to_string())
.or_insert_with(|| Entry::new(armed_until));
entry.generation = generation;
entry.observation = observation.capped();
entry.armed_until = armed_until;
if bytes <= self.max_bytes {
inner.store_window(feed, window, bytes);
inner.evict(self.max_bytes, self.max_entries);
}
inner.touch_entry(feed);
inner.evict_observations(self.max_observations);
}
fn record_not_modified(
&self,
feed: &str,
expected_generation: u64,
http_status: u16,
last_fetch_ms: i64,
armed_until: Instant,
) {
let mut inner = self.lock();
let current = inner.map.get(feed).map_or(0, |e| e.generation);
if current != expected_generation {
tracing::debug!(
feed,
commit = "not-modified",
expected_generation,
current_generation = current,
"rss stale cache commit dropped"
);
return;
}
let generation = inner.next_generation();
let entry = inner
.map
.entry(feed.to_string())
.or_insert_with(|| Entry::new(armed_until));
entry.generation = generation;
entry.observation.http_status = Some(http_status);
entry.observation.last_fetch_ms = Some(last_fetch_ms);
entry.armed_until = armed_until;
if entry.window.is_some() {
entry.observation.last_status = FeedStatus::Revalidated;
entry.observation.last_error = None;
} else {
entry.observation.last_status = FeedStatus::Error;
entry.observation.last_error =
Some(truncate(WINDOW_EVICTED_ON_REVALIDATION, MAX_ERROR_CHARS));
tracing::warn!(
feed,
http_status,
"rss feed revalidated but its cached window had already been evicted"
);
}
inner.touch_entry(feed);
inner.evict_observations(self.max_observations);
}
fn record_failure(
&self,
feed: &str,
expected_generation: u64,
http_status: Option<u16>,
error: String,
dialect_declared: Option<String>,
last_fetch_ms: i64,
armed_until: Instant,
) {
let mut inner = self.lock();
let current = inner.map.get(feed).map_or(0, |e| e.generation);
if current != expected_generation {
tracing::debug!(
feed,
commit = "failure",
expected_generation,
current_generation = current,
"rss stale cache commit dropped"
);
return;
}
let generation = inner.next_generation();
let entry = inner
.map
.entry(feed.to_string())
.or_insert_with(|| Entry::new(armed_until));
entry.generation = generation;
entry.observation.last_status = if entry.window.is_some() {
FeedStatus::StaleError
} else {
FeedStatus::Error
};
entry.observation.http_status = http_status;
entry.observation.last_error = Some(error);
entry.observation.last_fetch_ms = Some(last_fetch_ms);
if dialect_declared.is_some() {
entry.observation.dialect_declared = dialect_declared;
}
entry.armed_until = armed_until;
let observation = std::mem::take(&mut entry.observation);
entry.observation = observation.capped();
inner.touch_entry(feed);
inner.evict_observations(self.max_observations);
}
fn record_egress_denial(
&self,
feed: &str,
expected_generation: u64,
error: String,
last_fetch_ms: i64,
armed_until: Instant,
) {
let mut inner = self.lock();
let current = inner.map.get(feed).map_or(0, |e| e.generation);
if current != expected_generation {
tracing::debug!(
feed,
commit = "egress-denial",
expected_generation,
current_generation = current,
"rss stale cache commit dropped"
);
return;
}
let generation = inner.next_generation();
inner.drop_window(feed);
let entry = inner
.map
.entry(feed.to_string())
.or_insert_with(|| Entry::new(armed_until));
entry.generation = generation;
entry.observation.last_status = FeedStatus::Error;
entry.observation.http_status = None;
entry.observation.last_error = Some(error);
entry.observation.last_fetch_ms = Some(last_fetch_ms);
entry.armed_until = armed_until;
let observation = std::mem::take(&mut entry.observation);
entry.observation = observation.capped();
inner.touch_entry(feed);
inner.evict_observations(self.max_observations);
}
}
#[cfg(test)]
mod tests {
use super::*;
use arrow::array::UInt64Array;
use arrow::datatypes::{DataType, Field, Schema};
use std::sync::Arc;
fn window_with_rows(rows: usize) -> CachedWindow {
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::UInt64, false)]));
let ids: UInt64Array = (0..rows as u64).collect();
let batch = RecordBatch::try_new(schema, vec![Arc::new(ids)]).unwrap();
CachedWindow {
batch,
etag: Some("\"etag\"".into()),
last_modified: Some("Mon, 20 Jul 2026 10:00:00 GMT".into()),
fetched_from: "https://feed.example/f.xml".into(),
probed_at: Instant::now(),
}
}
fn obs_fresh(item_count: u64) -> FeedObservation {
FeedObservation {
last_fetch_ms: Some(1_700_000_000_000),
last_status: FeedStatus::Fresh,
http_status: Some(200),
last_error: None,
dialect: Some("rss2".into()),
dialect_declared: Some("rss2".into()),
conformance_notes: None,
title: Some("Feed".into()),
site_url: Some("https://example.com".into()),
description: Some("desc".into()),
item_count: Some(item_count),
}
}
fn obs_all_strings_long(chars: usize) -> FeedObservation {
let long = "x".repeat(chars);
FeedObservation {
last_fetch_ms: Some(1_700_000_000_000),
last_status: FeedStatus::Fresh,
http_status: Some(200),
last_error: Some(long.clone()),
dialect: Some("rss-2.0".into()),
dialect_declared: Some(long.clone()),
conformance_notes: Some(long.clone()),
title: Some(long.clone()),
site_url: Some(long.clone()),
description: Some(long),
item_count: Some(7),
}
}
#[test]
fn capped_bounds_every_feed_controlled_string_at_its_own_cap() {
let capped = obs_all_strings_long(10_000).capped();
assert_eq!(
capped.title.as_ref().unwrap().chars().count(),
MAX_FEED_TEXT_CHARS
);
assert_eq!(
capped.description.as_ref().unwrap().chars().count(),
MAX_FEED_TEXT_CHARS
);
assert_eq!(
capped.site_url.as_ref().unwrap().chars().count(),
MAX_ERROR_CHARS
);
assert_eq!(
capped.dialect_declared.as_ref().unwrap().chars().count(),
MAX_ERROR_CHARS
);
assert_eq!(
capped.conformance_notes.as_ref().unwrap().chars().count(),
MAX_ERROR_CHARS
);
assert_eq!(
capped.last_error.as_ref().unwrap().chars().count(),
MAX_ERROR_CHARS,
"capped again here even though its writers cap it, so the bound at \
this boundary does not depend on them"
);
}
#[test]
fn capped_leaves_values_within_bounds_byte_identical() {
let before = FeedObservation {
last_fetch_ms: Some(1_700_000_000_000),
last_status: FeedStatus::Fresh,
http_status: Some(200),
last_error: None,
dialect: Some("rss-2.0".into()),
dialect_declared: Some("rss-2.0".into()),
conformance_notes: Some("duplicate-identity: 1".into()),
title: Some("Daily Notes on Distributed Systems".into()),
site_url: Some("https://example.com/blog/index.html".into()),
description: Some("Occasional writing about storage engines.".into()),
item_count: Some(42),
};
let after = before.clone().capped();
assert_eq!(after.title, before.title);
assert_eq!(after.description, before.description);
assert_eq!(after.site_url, before.site_url);
assert_eq!(after.dialect_declared, before.dialect_declared);
assert_eq!(after.conformance_notes, before.conformance_notes);
assert_eq!(after.last_error, before.last_error);
assert_eq!(
after.dialect, before.dialect,
"never feed-controlled, never cut"
);
}
#[test]
fn capped_counts_characters_on_multi_byte_text() {
let title = "标题".repeat(10_000);
let description = "\u{1F680}".repeat(10_000);
let capped = FeedObservation {
title: Some(title),
description: Some(description),
..FeedObservation::default()
}
.capped();
let title = capped.title.expect("title survives");
assert_eq!(title.chars().count(), MAX_FEED_TEXT_CHARS);
assert_eq!(title.len(), MAX_FEED_TEXT_CHARS * 3);
let description = capped.description.expect("description survives");
assert_eq!(description.chars().count(), MAX_FEED_TEXT_CHARS);
assert_eq!(description.len(), MAX_FEED_TEXT_CHARS * 4);
}
#[test]
fn capped_leaves_non_string_fields_alone() {
let before = obs_all_strings_long(10_000);
let after = before.clone().capped();
assert_eq!(after.item_count, before.item_count);
assert_eq!(after.http_status, before.http_status);
assert_eq!(after.last_fetch_ms, before.last_fetch_ms);
assert!(matches!(after.last_status, FeedStatus::Fresh));
}
#[test]
fn record_success_caps_the_observation_it_stores_and_leaves_the_window_whole() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success(
"a",
0,
window_with_rows(3),
obs_all_strings_long(10_000),
t0 + Duration::from_secs(900),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert_eq!(
snap.observation.title.as_ref().unwrap().chars().count(),
MAX_FEED_TEXT_CHARS
);
assert_eq!(
snap.observation.site_url.as_ref().unwrap().chars().count(),
MAX_ERROR_CHARS
);
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 3);
}
#[test]
fn record_failure_caps_the_strings_it_writes() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_failure(
"a",
0,
Some(500),
"x".repeat(10_000),
Some(format!("unknown:{}", "y".repeat(10_000))),
1,
t0 + Duration::from_secs(30),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert_eq!(
snap.observation
.last_error
.as_ref()
.unwrap()
.chars()
.count(),
MAX_ERROR_CHARS
);
assert_eq!(
snap.observation
.dialect_declared
.as_ref()
.unwrap()
.chars()
.count(),
MAX_ERROR_CHARS
);
}
#[test]
fn success_arms_ttl_and_snapshot_reports_within_ttl() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success(
"a",
0,
window_with_rows(2),
obs_fresh(2),
t0 + Duration::from_secs(900),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(snap.within_ttl);
assert!(matches!(snap.observation.last_status, FeedStatus::Fresh));
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 2);
assert!(
!cache
.snapshot("a", t0 + Duration::from_secs(901))
.within_ttl
);
}
#[test]
fn failure_is_negative_cached_with_window_kept() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(2), obs_fresh(2), t0); cache.record_failure(
"a",
1,
Some(503),
"http status 503".into(),
None,
1,
t0 + Duration::from_secs(30),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap.within_ttl,
"failure re-armed the timer (negative cache)"
);
assert!(matches!(
snap.observation.last_status,
FeedStatus::StaleError
));
assert!(
snap.window.is_some(),
"stale window retained for serve-stale"
);
assert_eq!(
snap.observation.last_error.as_deref(),
Some("http status 503")
);
}
#[test]
fn failure_without_window_is_error_status() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_failure(
"a",
0,
Some(500),
"http status 500".into(),
None,
1,
t0 + Duration::from_secs(30),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(matches!(snap.observation.last_status, FeedStatus::Error));
assert!(snap.window.is_none());
assert_eq!(
snap.observation.last_error.as_deref(),
Some("http status 500")
);
assert!(
snap.within_ttl,
"failure still arms the negative-cache timer"
);
}
#[test]
fn eviction_drops_window_and_validators_but_keeps_observation() {
let one = window_with_rows(2);
let bytes = one.batch.get_array_memory_size();
let cache = MemoryFeedCache::new(bytes + 8, 64);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(900);
cache.record_success("a", 0, window_with_rows(2), obs_fresh(2), armed);
cache.record_success("b", 0, window_with_rows(2), obs_fresh(3), armed);
let snap_a = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap_a.window.is_none(),
"a's window (and its validators) must be evicted to make room for b"
);
assert!(
matches!(snap_a.observation.last_status, FeedStatus::Fresh),
"the observation survives eviction of its window"
);
assert_eq!(snap_a.observation.item_count, Some(2));
let snap_b = cache.snapshot("b", t0 + Duration::from_secs(1));
assert_eq!(snap_b.window.as_ref().unwrap().batch.num_rows(), 2);
}
#[test]
fn unknown_feed_snapshot_is_never() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let snap = cache.snapshot("ghost", Instant::now());
assert!(matches!(snap.observation.last_status, FeedStatus::Never));
assert!(snap.window.is_none());
assert!(!snap.within_ttl);
}
#[test]
fn failure_fuse_is_clamped() {
assert_eq!(
failure_fuse(Duration::from_secs(0)),
Duration::from_secs(30)
);
assert_eq!(
failure_fuse(Duration::from_secs(900)),
Duration::from_secs(225)
);
assert_eq!(
failure_fuse(Duration::from_secs(10_000)),
Duration::from_secs(300)
);
}
#[test]
fn not_modified_rearms_and_flips_to_revalidated() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(2), obs_fresh(2), t0); cache.record_not_modified("a", 1, 304, 1, t0 + Duration::from_secs(900));
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap.within_ttl,
"a 304 re-arms the timer exactly like success/failure do"
);
assert!(matches!(
snap.observation.last_status,
FeedStatus::Revalidated
));
assert_eq!(snap.observation.http_status, Some(304));
assert_eq!(snap.observation.last_fetch_ms, Some(1));
assert!(
snap.window.is_some(),
"the existing window is kept across revalidation"
);
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 2);
}
#[test]
fn record_not_modified_after_window_eviction_rearms_and_records_it() {
let one = window_with_rows(1);
let bytes = one.batch.get_array_memory_size();
let cache = MemoryFeedCache::new(bytes + 8, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), t0);
assert!(cache.snapshot("a", t0).window.is_some());
assert!(!cache.snapshot("a", t0 + Duration::from_secs(1)).within_ttl);
cache.record_success(
"b",
0,
window_with_rows(1),
obs_fresh(1),
t0 + Duration::from_secs(900),
);
assert!(
cache.snapshot("a", t0).window.is_none(),
"b's window evicted a's under the byte budget"
);
cache.record_not_modified("a", 1, 304, 42, t0 + Duration::from_secs(600));
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap.within_ttl,
"the 304 re-armed the timer even though the window was gone — \
otherwise the feed refetches on every scan"
);
assert!(
matches!(snap.observation.last_status, FeedStatus::Error),
"no window means zero rows, and Error is the status that says so: {:?}",
snap.observation.last_status
);
assert_eq!(
snap.observation.last_error.as_deref(),
Some(WINDOW_EVICTED_ON_REVALIDATION),
"the zero rows an operator will see have a stated reason"
);
assert_eq!(snap.observation.http_status, Some(304));
assert_eq!(snap.observation.last_fetch_ms, Some(42));
assert!(snap.window.is_none());
assert_eq!(snap.observation.item_count, Some(1));
}
#[test]
fn lru_touch_order_respected() {
let one = window_with_rows(1);
let bytes = one.batch.get_array_memory_size();
let cache = MemoryFeedCache::new(bytes * 2 + 8, 64);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(900);
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), armed);
cache.record_success("b", 0, window_with_rows(1), obs_fresh(1), armed);
assert!(
cache
.snapshot("a", t0 + Duration::from_secs(1))
.window
.is_some()
);
cache.record_success("c", 0, window_with_rows(1), obs_fresh(1), armed);
assert!(
cache
.snapshot("a", t0 + Duration::from_secs(1))
.window
.is_some(),
"recently touched entry stays"
);
assert!(
cache
.snapshot("b", t0 + Duration::from_secs(1))
.window
.is_none(),
"least-recently-used entry is evicted, not the touched one"
);
assert!(
cache
.snapshot("c", t0 + Duration::from_secs(1))
.window
.is_some()
);
}
#[test]
fn max_entries_bound_evicts_lru_window() {
let cache = MemoryFeedCache::new(1 << 20, 2);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(900);
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), armed);
cache.record_success("b", 0, window_with_rows(1), obs_fresh(1), armed);
cache.record_success("c", 0, window_with_rows(1), obs_fresh(1), armed);
let snap_a = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap_a.window.is_none(),
"the entry-count bound must evict the LRU window, not just the byte budget"
);
assert!(matches!(snap_a.observation.last_status, FeedStatus::Fresh));
assert!(
cache
.snapshot("b", t0 + Duration::from_secs(1))
.window
.is_some()
);
assert!(
cache
.snapshot("c", t0 + Duration::from_secs(1))
.window
.is_some()
);
}
#[test]
fn max_observations_backstop_evicts_lru_whole_entry() {
let cache = MemoryFeedCache::new(1 << 20, 1);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(30);
for i in 0..8 {
cache.record_failure(
&format!("feed-{i}"),
0,
Some(500),
"http status 500".into(),
None,
1,
armed,
);
}
let touched = cache.snapshot("feed-0", t0 + Duration::from_secs(1));
assert!(matches!(touched.observation.last_status, FeedStatus::Error));
cache.record_failure(
"feed-8",
0,
Some(500),
"http status 500".into(),
None,
1,
armed,
);
let evicted = cache.snapshot("feed-1", t0 + Duration::from_secs(1));
assert!(
matches!(evicted.observation.last_status, FeedStatus::Never),
"the least-recently-used whole entry, observation included, must be gone"
);
let survivor = cache.snapshot("feed-0", t0 + Duration::from_secs(1));
assert!(
matches!(survivor.observation.last_status, FeedStatus::Error),
"the recently touched entry survives with its observation intact"
);
}
#[test]
fn window_accounting_agrees_with_the_map_after_every_path() {
let bytes = window_with_rows(1).batch.get_array_memory_size();
let cache = MemoryFeedCache::new(bytes * 2 + 8, 8);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(900);
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), armed);
cache.record_success("b", 0, window_with_rows(1), obs_fresh(1), armed);
cache.record_success("a", 1, window_with_rows(1), obs_fresh(1), armed);
cache.record_failure("a", 3, Some(500), "boom".into(), None, 1, armed);
cache.record_success("c", 0, window_with_rows(1), obs_fresh(1), armed);
for i in 0..70 {
cache.record_failure(
&format!("f{i}"),
0,
Some(500),
"boom".into(),
None,
1,
armed,
);
}
let inner = cache.lock();
assert_eq!(
inner.windowed,
inner.map.values().filter(|e| e.window.is_some()).count(),
"the windowed count must equal the windows the map actually holds"
);
assert_eq!(
inner.window_bytes,
inner
.map
.values()
.filter_map(|e| e.window.as_ref())
.map(|w| w.bytes)
.sum::<usize>(),
"the byte total must equal the bytes the map actually holds"
);
}
#[test]
fn max_entries_zero_disables_window_cache_but_keeps_observations() {
let cache = MemoryFeedCache::new(1 << 20, 0);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(900);
cache.record_success("a", 0, window_with_rows(2), obs_fresh(2), armed);
let snap_a = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap_a.window.is_none(),
"max_entries = 0 must disable window caching, not just shrink it"
);
assert!(matches!(snap_a.observation.last_status, FeedStatus::Fresh));
cache.record_success("b", 0, window_with_rows(1), obs_fresh(1), armed);
let snap_a_after = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
matches!(snap_a_after.observation.last_status, FeedStatus::Fresh),
"the floor must keep a's observation alive across a later, unrelated insert"
);
let snap_b = cache.snapshot("b", t0 + Duration::from_secs(1));
assert!(matches!(snap_b.observation.last_status, FeedStatus::Fresh));
}
#[test]
fn a_second_success_from_the_same_snapshot_is_dropped() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
let armed = t0 + Duration::from_secs(900);
cache.record_success("a", 0, window_with_rows(2), obs_fresh(2), armed);
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), armed);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert_eq!(
snap.window.as_ref().unwrap().batch.num_rows(),
2,
"the first-committed window stays; the stale commit is dropped"
);
assert_eq!(snap.observation.item_count, Some(2));
}
#[test]
fn a_stale_304_does_not_label_a_window_it_never_validated() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), t0);
cache.record_success(
"a",
1,
window_with_rows(2),
obs_fresh(2),
t0 + Duration::from_secs(900),
);
cache.record_not_modified("a", 1, 304, 99, t0 + Duration::from_secs(600));
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
matches!(snap.observation.last_status, FeedStatus::Fresh),
"not `revalidated`: the 304 never saw the two-row window"
);
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 2);
assert_ne!(
snap.observation.last_fetch_ms,
Some(99),
"the dropped 304 must not update fetch metadata either"
);
}
#[test]
fn a_stale_failure_does_not_degrade_a_fresher_success() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), t0);
cache.record_success(
"a",
1,
window_with_rows(2),
obs_fresh(2),
t0 + Duration::from_secs(900),
);
cache.record_failure(
"a",
1,
Some(500),
"http status 500".into(),
None,
7,
t0 + Duration::from_secs(30),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(matches!(snap.observation.last_status, FeedStatus::Fresh));
assert_eq!(snap.observation.last_error, None);
assert!(
cache
.snapshot("a", t0 + Duration::from_secs(600))
.within_ttl,
"the success's 900s arm stands; the dropped failure's 30s fuse does not"
);
}
#[test]
fn a_stale_success_supersedes_a_failure_commit() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_failure(
"a",
0,
None,
"transport error: connection refused".into(),
None,
1,
t0 + Duration::from_secs(30),
);
cache.record_success(
"a",
0,
window_with_rows(2),
obs_fresh(2),
t0 + Duration::from_secs(900),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(matches!(snap.observation.last_status, FeedStatus::Fresh));
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 2);
assert_eq!(snap.observation.last_error, None);
}
#[test]
fn an_egress_denial_drops_the_window_and_negative_caches() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(2), obs_fresh(2), t0); cache.record_egress_denial(
"a",
1,
"egress blocked: host 'feed.example' resolves to private address 10.0.0.1".into(),
7,
t0 + Duration::from_secs(30),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(
snap.window.is_none(),
"the refused feed's window (and validators) must be purged"
);
assert!(
matches!(snap.observation.last_status, FeedStatus::Error),
"`error`, not `stale-error`: there is deliberately nothing to serve stale"
);
assert!(
snap.observation
.last_error
.as_deref()
.unwrap()
.contains("egress blocked")
);
assert!(
snap.within_ttl,
"a refusal negative-caches exactly like any other failure"
);
}
#[test]
fn a_stale_egress_denial_does_not_purge_a_fresher_success() {
let cache = MemoryFeedCache::new(1 << 20, 64);
let t0 = Instant::now();
cache.record_success("a", 0, window_with_rows(1), obs_fresh(1), t0);
cache.record_success(
"a",
1,
window_with_rows(2),
obs_fresh(2),
t0 + Duration::from_secs(900),
);
cache.record_egress_denial(
"a",
1,
"egress blocked: host 'feed.example' resolves to private address 10.0.0.1".into(),
7,
t0 + Duration::from_secs(30),
);
let snap = cache.snapshot("a", t0 + Duration::from_secs(1));
assert!(matches!(snap.observation.last_status, FeedStatus::Fresh));
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 2);
}
}