use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use arrow::record_batch::RecordBatch;
use tokio::sync::Semaphore;
use super::ResolvedSubscription;
use super::cache::{
CachedWindow, FeedCache, FeedObservation, FeedSnapshot, FeedStatus, MemoryFeedCache,
failure_fuse,
};
use super::config::{DEFAULT_MAX_RESPONSE_BYTES, RssConfig};
use super::egress::EgressPolicy;
use super::error::{MAX_ERROR_CHARS, RssError, truncate};
use super::fetch::{FeedFetcher, FetchError, FetchOutcome, Validators};
use super::parse::{ParseFailure, ParsedDocument, parse_feed_document};
use super::schema::{FeedsRow, build_feeds_batch, build_items_batch, with_window_status};
pub const CACHE_MAX_BYTES: usize = 64 * 1024 * 1024;
const WINDOW_ENTRY_HEADROOM: usize = 8;
const MAX_TTL: Duration = Duration::from_secs(365 * 24 * 60 * 60);
const PARSE_TIMEOUT: Duration = Duration::from_secs(10);
const REPROBE_EVERY_REFRESHES: u32 = 24;
const MIN_REPROBE_INTERVAL: Duration = Duration::from_secs(6 * 60 * 60);
const MAX_REPROBE_INTERVAL: Duration = Duration::from_secs(7 * 24 * 60 * 60);
fn reprobe_interval(ttl: Duration) -> Duration {
ttl.saturating_mul(REPROBE_EVERY_REFRESHES)
.clamp(MIN_REPROBE_INTERVAL, MAX_REPROBE_INTERVAL)
}
const MAX_PARSE_FUSE: Duration = Duration::from_secs(60 * 60);
fn parse_fuse(max_response_bytes: u64) -> Duration {
let units = max_response_bytes
.div_ceil(DEFAULT_MAX_RESPONSE_BYTES)
.max(1);
Duration::from_secs(PARSE_TIMEOUT.as_secs().saturating_mul(units)).min(MAX_PARSE_FUSE)
}
pub struct RssEngine {
source_name: String,
subscriptions: Vec<ResolvedSubscription>,
by_name: HashMap<String, usize>,
fetcher: FeedFetcher,
cache: Arc<dyn FeedCache>,
semaphore: Arc<Semaphore>,
ttl: Duration,
scan_timeout: Duration,
parse_fuse: Duration,
reprobe_interval: Duration,
}
impl RssEngine {
pub fn new(
source_name: String,
subscriptions: Vec<ResolvedSubscription>,
config: &RssConfig,
policy: Option<Arc<dyn EgressPolicy>>,
) -> Result<Self, RssError> {
let fetcher = FeedFetcher::new(
policy,
Duration::from_secs(config.request_timeout_seconds),
config.max_response_bytes,
config.user_agent.clone(),
)?;
let cache = Arc::new(MemoryFeedCache::new(
CACHE_MAX_BYTES,
subscriptions.len().saturating_add(WINDOW_ENTRY_HEADROOM),
));
Ok(Self::with_parts(
source_name,
subscriptions,
config,
fetcher,
cache,
))
}
pub(crate) fn with_parts(
source_name: String,
subscriptions: Vec<ResolvedSubscription>,
config: &RssConfig,
fetcher: FeedFetcher,
cache: Arc<dyn FeedCache>,
) -> Self {
let by_name = subscriptions
.iter()
.enumerate()
.map(|(index, sub)| (sub.name.clone(), index))
.collect();
let configured_ttl = Duration::from_secs(config.ttl_seconds);
let ttl = configured_ttl.min(MAX_TTL);
if ttl != configured_ttl {
tracing::warn!(
source = %source_name,
configured_ttl_seconds = config.ttl_seconds,
effective_ttl_seconds = ttl.as_secs(),
"rss ttl_seconds clamped to the engine's ceiling"
);
}
let max_concurrent = config.max_concurrent.clamp(1, Semaphore::MAX_PERMITS);
if config.max_concurrent > Semaphore::MAX_PERMITS {
tracing::warn!(
source = %source_name,
configured_max_concurrent = config.max_concurrent,
effective_max_concurrent = max_concurrent,
"rss max_concurrent clamped to the semaphore's ceiling"
);
}
let scan_timeout = Duration::from_secs(config.scan_timeout_seconds);
let parse_fuse = parse_fuse(config.max_response_bytes);
if parse_fuse >= scan_timeout {
tracing::warn!(
source = %source_name,
parse_fuse_seconds = parse_fuse.as_secs(),
scan_timeout_seconds = config.scan_timeout_seconds,
"rss parse fuse meets or exceeds the scan timeout; raise scan_timeout_seconds"
);
}
Self {
source_name,
subscriptions,
by_name,
fetcher,
cache,
semaphore: Arc::new(Semaphore::new(max_concurrent)),
ttl,
scan_timeout,
parse_fuse,
reprobe_interval: reprobe_interval(ttl),
}
}
pub fn subscriptions(&self) -> &[ResolvedSubscription] {
&self.subscriptions
}
pub fn scan_timeout(&self) -> Duration {
self.scan_timeout
}
pub async fn serve_feed(
&self,
feed: &str,
launch_gate: impl Fn() -> bool + Send,
) -> Option<RecordBatch> {
let Some(sub) = self.subscription(feed) else {
tracing::debug!(
source = %self.source_name,
feed,
outcome = "unknown-feed",
"rss feed served"
);
return None;
};
let started = Instant::now();
let (batch, log) = self.serve_subscription(sub, launch_gate).await;
tracing::debug!(
source = %self.source_name,
feed = %sub.name,
url = %sub.url,
outcome = log.outcome,
http_status = ?log.http_status,
bytes = log.bytes,
rows = log.rows,
notes = log.notes,
elapsed_ms = started.elapsed().as_millis() as u64,
"rss feed served"
);
batch
}
async fn serve_subscription(
&self,
sub: &ResolvedSubscription,
launch_gate: impl Fn() -> bool + Send,
) -> (Option<RecordBatch>, ServeLog) {
let snapshot = self.cache.snapshot(&sub.name, Instant::now());
if snapshot.within_ttl && !window_lost(&snapshot) {
let batch = snapshot.window.as_ref().and_then(|window| {
with_window_status(&window.batch, snapshot.observation.last_status)
});
let rows = batch.as_ref().map_or(0, RecordBatch::num_rows);
return (
batch,
ServeLog {
outcome: "cache-hit",
http_status: snapshot.observation.http_status,
bytes: 0,
rows,
notes: 0,
},
);
}
let _permit = match self.semaphore.acquire().await {
Ok(permit) => permit,
Err(_) => return (None, ServeLog::bare("error")),
};
if !launch_gate() {
return (None, ServeLog::bare("gate-closed"));
}
let window = snapshot.window.as_ref();
let target = match window {
Some(cached) if !self.reprobe_due(cached, &sub.url) => cached.fetched_from.clone(),
_ => sub.url.clone(),
};
let validators = window
.filter(|cached| cached.fetched_from == target)
.map(|cached| Validators {
etag: cached.etag.clone(),
last_modified: cached.last_modified.clone(),
});
let probed_at = match window {
Some(cached) if target != sub.url => cached.probed_at,
_ => Instant::now(),
};
match self.fetcher.fetch(&target, validators.as_ref()).await {
Ok(FetchOutcome::NotModified { http_status }) => {
self.cache.record_not_modified(
&sub.name,
snapshot.generation,
http_status,
now_ms(),
arm(Instant::now(), self.ttl),
);
let batch = snapshot
.window
.as_ref()
.and_then(|window| with_window_status(&window.batch, FeedStatus::Revalidated));
let rows = batch.as_ref().map_or(0, RecordBatch::num_rows);
(
batch,
ServeLog {
outcome: "revalidated",
http_status: Some(http_status),
bytes: 0,
rows,
notes: 0,
},
)
}
Ok(FetchOutcome::Fetched {
body,
http_status,
etag,
last_modified,
content_type,
final_url,
}) => {
let bytes = body.len();
match parse_off_worker(body, content_type, self.parse_fuse).await {
Ok(document) => {
let notes = document.conformance_notes.len();
let batch = self.record_fresh_window(
sub,
snapshot.generation,
document,
http_status,
etag,
last_modified,
final_url,
probed_at,
);
let rows = batch.num_rows();
(
Some(batch),
ServeLog {
outcome: "fetched",
http_status: Some(http_status),
bytes,
rows,
notes,
},
)
}
Err(failure) => self.degrade(
sub,
snapshot.generation,
Some(http_status),
parse_error_message(failure.stage, &failure.reason),
failure.dialect_declared,
bytes,
),
}
}
Err(FetchError::Egress(denied)) => self.deny(
sub,
snapshot.generation,
truncate(&denied.to_string(), MAX_ERROR_CHARS),
),
Err(error) => {
let http_status = match &error {
FetchError::Status { status } => Some(*status),
_ => None,
};
self.degrade(
sub,
snapshot.generation,
http_status,
truncate(&error.to_string(), MAX_ERROR_CHARS),
None,
0,
)
}
}
}
fn reprobe_due(&self, window: &CachedWindow, configured_url: &str) -> bool {
window.fetched_from != configured_url
&& Instant::now().saturating_duration_since(window.probed_at) >= self.reprobe_interval
}
fn record_fresh_window(
&self,
sub: &ResolvedSubscription,
expected_generation: u64,
document: ParsedDocument,
http_status: u16,
etag: Option<String>,
last_modified: Option<String>,
fetched_from: String,
probed_at: Instant,
) -> RecordBatch {
let batch = build_items_batch(&sub.name, &sub.url, &document.items);
let observation = FeedObservation {
last_fetch_ms: Some(now_ms()),
last_status: FeedStatus::Fresh,
http_status: Some(http_status),
last_error: None,
dialect: Some(document.dialect.to_string()),
dialect_declared: document.dialect_declared,
conformance_notes: serde_json::to_string(&document.conformance_notes).ok(),
title: document.meta.title,
site_url: document.meta.site_url,
description: document.meta.description,
item_count: Some(batch.num_rows() as u64),
};
self.cache.record_success(
&sub.name,
expected_generation,
CachedWindow {
batch: batch.clone(),
etag,
last_modified,
fetched_from,
probed_at,
},
observation,
arm(Instant::now(), self.ttl),
);
batch
}
fn degrade(
&self,
sub: &ResolvedSubscription,
expected_generation: u64,
http_status: Option<u16>,
error: String,
dialect_declared: Option<String>,
bytes: usize,
) -> (Option<RecordBatch>, ServeLog) {
self.cache.record_failure(
&sub.name,
expected_generation,
http_status,
error.clone(),
dialect_declared,
now_ms(),
arm(Instant::now(), failure_fuse(self.ttl)),
);
self.degraded_serve(sub, http_status, error, bytes)
}
fn deny(
&self,
sub: &ResolvedSubscription,
expected_generation: u64,
error: String,
) -> (Option<RecordBatch>, ServeLog) {
self.cache.record_egress_denial(
&sub.name,
expected_generation,
error.clone(),
now_ms(),
arm(Instant::now(), failure_fuse(self.ttl)),
);
self.degraded_serve(sub, None, error, 0)
}
fn degraded_serve(
&self,
sub: &ResolvedSubscription,
http_status: Option<u16>,
error: String,
bytes: usize,
) -> (Option<RecordBatch>, ServeLog) {
tracing::warn!(
source = %self.source_name,
feed = %sub.name,
%error,
"rss feed degraded"
);
let snapshot = self.cache.snapshot(&sub.name, Instant::now());
let batch = snapshot
.window
.as_ref()
.and_then(|window| with_window_status(&window.batch, snapshot.observation.last_status));
let rows = batch.as_ref().map_or(0, RecordBatch::num_rows);
(
batch,
ServeLog {
outcome: snapshot.observation.last_status.as_str(),
http_status,
bytes,
rows,
notes: 0,
},
)
}
pub fn feeds_row(&self, feed: &str) -> RecordBatch {
let Some(sub) = self.subscription(feed) else {
return build_feeds_batch(&[]);
};
let snapshot = self.cache.snapshot(&sub.name, Instant::now());
let observation = snapshot.observation;
let (etag, last_modified) = snapshot
.window
.map_or((None, None), |window| (window.etag, window.last_modified));
build_feeds_batch(&[FeedsRow {
name: sub.name.clone(),
url: sub.url.clone(),
title: observation.title,
site_url: observation.site_url,
description: observation.description,
last_fetch_ms: observation.last_fetch_ms,
last_status: observation.last_status,
http_status: observation.http_status,
last_error: observation.last_error,
etag,
last_modified,
dialect: observation.dialect,
dialect_declared: observation.dialect_declared,
conformance_notes: observation.conformance_notes,
item_count: observation.item_count,
}])
}
fn subscription(&self, feed: &str) -> Option<&ResolvedSubscription> {
let index = *self.by_name.get(feed)?;
self.subscriptions.get(index)
}
}
struct ServeLog {
outcome: &'static str,
http_status: Option<u16>,
bytes: usize,
rows: usize,
notes: usize,
}
impl ServeLog {
fn bare(outcome: &'static str) -> Self {
Self {
outcome,
http_status: None,
bytes: 0,
rows: 0,
notes: 0,
}
}
}
fn window_lost(snapshot: &FeedSnapshot) -> bool {
snapshot.window.is_none()
&& snapshot
.observation
.last_status
.window_status_str()
.is_some()
}
async fn parse_off_worker(
body: Vec<u8>,
content_type: Option<String>,
fuse: Duration,
) -> Result<ParsedDocument, ParseFailure> {
let parse =
tokio::task::spawn_blocking(move || parse_feed_document(&body, content_type.as_deref()));
match tokio::time::timeout(fuse, parse).await {
Ok(Ok(result)) => result,
Ok(Err(join_error)) => Err(join_failure(&join_error)),
Err(_elapsed) => Err(ParseFailure {
stage: "timeout",
reason: format!(
"feed parse did not finish within {}s; abandoned",
fuse.as_secs()
),
dialect_declared: None,
}),
}
}
fn join_failure(join_error: &tokio::task::JoinError) -> ParseFailure {
ParseFailure {
stage: "panic",
reason: if join_error.is_panic() {
"feed parse panicked; the payload and backtrace are in the server log".to_string()
} else {
"feed parse did not complete: its task was cancelled".to_string()
},
dialect_declared: None,
}
}
fn parse_error_message(stage: &str, reason: &str) -> String {
truncate(
&format!("parse failed at {stage}: {reason}"),
MAX_ERROR_CHARS,
)
}
fn now_ms() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.ok()
.and_then(|since| i64::try_from(since.as_millis()).ok())
.unwrap_or(0)
}
fn arm(now: Instant, after: Duration) -> Instant {
now.checked_add(after).unwrap_or(now)
}
#[cfg(test)]
mod tests {
use std::net::IpAddr;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use arrow::array::{Array, UInt64Array};
use super::*;
use crate::sources::providers::open_connector::testutil::{CapturedEvent, capture_events};
use crate::sources::providers::rss::config::{FeedSubscription, inline_config};
use crate::sources::providers::rss::egress::{AllowAll, EgressReason};
use crate::sources::providers::rss::testutil::{
MockFeedServer, MockResponse, MockResponseExt, RSS2_MINIMAL, str_col, str_opt_col,
};
#[derive(Debug)]
struct DenyList(Vec<IpAddr>);
impl EgressPolicy for DenyList {
fn check_ip(&self, ip: IpAddr) -> Result<(), EgressReason> {
if self.0.contains(&ip) {
Err("test-denied".into())
} else {
Ok(())
}
}
}
#[derive(Debug)]
struct TogglePolicy(AtomicBool);
impl EgressPolicy for TogglePolicy {
fn check_ip(&self, _ip: IpAddr) -> Result<(), EgressReason> {
if self.0.load(Ordering::SeqCst) {
Err("test-denied".into())
} else {
Ok(())
}
}
}
fn test_engine(server: &MockFeedServer, feeds: &[(&str, &str)], ttl_seconds: u64) -> RssEngine {
let urls: Vec<(String, String)> = feeds
.iter()
.map(|(name, path)| ((*name).to_string(), format!("{}{path}", server.url())))
.collect();
engine_over(&urls, ttl_seconds, 4)
}
fn engine_over(
feeds: &[(String, String)],
ttl_seconds: u64,
max_concurrent: usize,
) -> RssEngine {
let cache = Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, feeds.len() + 8));
engine_with_cache(feeds, ttl_seconds, max_concurrent, cache)
}
struct AlwaysExpired(MemoryFeedCache);
impl FeedCache for AlwaysExpired {
fn snapshot(&self, feed: &str, now: Instant) -> FeedSnapshot {
FeedSnapshot {
within_ttl: false,
..self.0.snapshot(feed, now)
}
}
fn record_success(
&self,
feed: &str,
expected_generation: u64,
window: CachedWindow,
observation: FeedObservation,
armed_until: Instant,
) {
self.0
.record_success(feed, expected_generation, window, observation, armed_until);
}
fn record_not_modified(
&self,
feed: &str,
expected_generation: u64,
http_status: u16,
last_fetch_ms: i64,
armed_until: Instant,
) {
self.0.record_not_modified(
feed,
expected_generation,
http_status,
last_fetch_ms,
armed_until,
);
}
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,
) {
self.0.record_failure(
feed,
expected_generation,
http_status,
error,
dialect_declared,
last_fetch_ms,
armed_until,
);
}
fn record_egress_denial(
&self,
feed: &str,
expected_generation: u64,
error: String,
last_fetch_ms: i64,
armed_until: Instant,
) {
self.0.record_egress_denial(
feed,
expected_generation,
error,
last_fetch_ms,
armed_until,
);
}
}
struct RecordsFailureArm {
inner: MemoryFeedCache,
armed: std::sync::Mutex<Vec<Instant>>,
}
impl RecordsFailureArm {
fn new(inner: MemoryFeedCache) -> Self {
Self {
inner,
armed: std::sync::Mutex::new(Vec::new()),
}
}
fn armed_instants(&self) -> Vec<Instant> {
self.armed.lock().unwrap_or_else(|p| p.into_inner()).clone()
}
}
impl FeedCache for RecordsFailureArm {
fn snapshot(&self, feed: &str, now: Instant) -> FeedSnapshot {
self.inner.snapshot(feed, now)
}
fn record_success(
&self,
feed: &str,
expected_generation: u64,
window: CachedWindow,
observation: FeedObservation,
armed_until: Instant,
) {
self.inner
.record_success(feed, expected_generation, window, observation, armed_until);
}
fn record_not_modified(
&self,
feed: &str,
expected_generation: u64,
http_status: u16,
last_fetch_ms: i64,
armed_until: Instant,
) {
self.inner.record_not_modified(
feed,
expected_generation,
http_status,
last_fetch_ms,
armed_until,
);
}
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,
) {
self.armed
.lock()
.unwrap_or_else(|p| p.into_inner())
.push(armed_until);
self.inner.record_failure(
feed,
expected_generation,
http_status,
error,
dialect_declared,
last_fetch_ms,
armed_until,
);
}
fn record_egress_denial(
&self,
feed: &str,
expected_generation: u64,
error: String,
last_fetch_ms: i64,
armed_until: Instant,
) {
self.inner.record_egress_denial(
feed,
expected_generation,
error,
last_fetch_ms,
armed_until,
);
}
}
fn engine_with_cache(
feeds: &[(String, String)],
ttl_seconds: u64,
max_concurrent: usize,
cache: Arc<dyn FeedCache>,
) -> RssEngine {
engine_with_cache_and_policy(
feeds,
ttl_seconds,
max_concurrent,
cache,
Arc::new(AllowAll),
)
}
fn engine_with_cache_and_policy(
feeds: &[(String, String)],
ttl_seconds: u64,
max_concurrent: usize,
cache: Arc<dyn FeedCache>,
policy: Arc<dyn EgressPolicy>,
) -> RssEngine {
let subscriptions: Vec<ResolvedSubscription> = feeds
.iter()
.map(|(name, url)| ResolvedSubscription {
name: name.clone(),
url: url.clone(),
})
.collect();
let mut config = inline_config(
subscriptions
.iter()
.map(|sub| FeedSubscription {
url: sub.url.clone(),
name: Some(sub.name.clone()),
})
.collect(),
);
config.ttl_seconds = ttl_seconds;
config.max_concurrent = max_concurrent;
config.request_timeout_seconds = 5;
let fetcher = FeedFetcher::new(
Some(policy),
Duration::from_secs(config.request_timeout_seconds),
config.max_response_bytes,
config.user_agent.clone(),
)
.expect("build the test fetcher");
RssEngine::with_parts(
"rss_test".to_string(),
subscriptions,
&config,
fetcher,
cache,
)
}
fn events_with_message(
events: &Arc<std::sync::Mutex<Vec<CapturedEvent>>>,
message: &str,
) -> Vec<CapturedEvent> {
events
.lock()
.unwrap_or_else(|p| p.into_inner())
.iter()
.filter(|event| event.message == message)
.cloned()
.collect()
}
fn u64_col(batch: &RecordBatch, name: &str) -> Vec<Option<u64>> {
let index = batch.schema().index_of(name).expect("column exists");
let column = batch
.column(index)
.as_any()
.downcast_ref::<UInt64Array>()
.expect("column is UInt64");
(0..column.len())
.map(|row| column.is_valid(row).then(|| column.value(row)))
.collect()
}
#[tokio::test]
async fn fresh_fetch_parses_and_stamps_fresh() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
let batch = engine.serve_feed("a", || true).await.expect("rows served");
assert_eq!(batch.num_rows(), 1);
assert_eq!(str_col(&batch, "window_status"), vec!["fresh"]);
assert_eq!(str_col(&batch, "feed"), vec!["a"]);
assert_eq!(server.requests().len(), 1);
let again = engine.serve_feed("a", || true).await.expect("rows served");
assert_eq!(str_col(&again, "window_status"), vec!["fresh"]);
assert_eq!(server.requests().len(), 1);
}
#[tokio::test]
async fn expired_with_etag_takes_304_and_stamps_revalidated() {
let server = MockFeedServer::start(|req| {
if req.header("if-none-match").is_some() {
MockResponse::status(304)
} else {
MockResponse::xml(RSS2_MINIMAL).with_header("etag", "\"v1\"")
}
})
.await;
let engine = test_engine(&server, &[("a", "/f.xml")], 0);
engine.serve_feed("a", || true).await.expect("first serve");
let batch = engine.serve_feed("a", || true).await.expect("second serve");
assert_eq!(str_col(&batch, "window_status"), vec!["revalidated"]);
assert_eq!(batch.num_rows(), 1, "the cached window is still served");
assert_eq!(server.requests().len(), 2);
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["revalidated"]);
assert_eq!(str_opt_col(&row, "etag"), vec![Some("\"v1\"".to_string())]);
}
#[tokio::test]
async fn failed_refetch_serves_stale_rows_and_records_error() {
let (_interest_guard, _) = capture_events();
let hits = Arc::new(AtomicUsize::new(0));
let h = Arc::clone(&hits);
let server = MockFeedServer::start(move |_| {
if h.fetch_add(1, Ordering::SeqCst) == 0 {
MockResponse::xml(RSS2_MINIMAL)
} else {
MockResponse::status(500)
}
})
.await;
let engine = test_engine(&server, &[("a", "/f.xml")], 0);
engine.serve_feed("a", || true).await.expect("first serve");
let batch = engine.serve_feed("a", || true).await.expect("stale rows");
assert_eq!(str_col(&batch, "window_status"), vec!["stale-error"]);
assert_eq!(batch.num_rows(), 1);
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["stale-error"]);
let error = str_opt_col(&row, "last_error")[0]
.clone()
.expect("last_error recorded");
assert!(
error.contains("500"),
"last_error names the status: {error}"
);
let n = server.requests().len();
let cached = engine
.serve_feed("a", || true)
.await
.expect("stale rows again, from the cache");
assert_eq!(
str_col(&cached, "window_status"),
vec!["stale-error"],
"the cache-hit path stamps the stored status, not the batch's build-time label"
);
engine.feeds_row("a");
assert_eq!(server.requests().len(), n);
}
#[tokio::test]
async fn stale_error_then_304_returns_to_revalidated_and_clears_the_error() {
let (_interest_guard, _) = capture_events();
let hits = Arc::new(AtomicUsize::new(0));
let h = Arc::clone(&hits);
let server = MockFeedServer::start(move |req| {
match h.fetch_add(1, Ordering::SeqCst) {
0 => MockResponse::xml(RSS2_MINIMAL).with_header("etag", "\"v1\""),
1..=3 => MockResponse::status(500),
_ if req.header("if-none-match").is_some() => MockResponse::status(304),
_ => MockResponse::status(400),
}
})
.await;
let urls = vec![("a".to_string(), format!("{}/f.xml", server.url()))];
let cache = Arc::new(AlwaysExpired(MemoryFeedCache::new(CACHE_MAX_BYTES, 8)));
let engine = engine_with_cache(&urls, 900, 4, cache);
engine.serve_feed("a", || true).await.expect("first serve");
let stale = engine.serve_feed("a", || true).await.expect("stale rows");
assert_eq!(str_col(&stale, "window_status"), vec!["stale-error"]);
assert!(str_opt_col(&engine.feeds_row("a"), "last_error")[0].is_some());
let revalidated = engine.serve_feed("a", || true).await.expect("304 serve");
assert_eq!(str_col(&revalidated, "window_status"), vec!["revalidated"]);
assert_eq!(revalidated.num_rows(), 1);
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["revalidated"]);
assert_eq!(
str_opt_col(&row, "last_error"),
vec![None],
"a successful revalidation clears the previous failure's error"
);
assert_eq!(
str_opt_col(&row, "etag"),
vec![Some("\"v1\"".to_string())],
"the validators that earned the 304 survive it"
);
}
#[tokio::test]
async fn a_cancelled_serve_releases_the_politeness_permit() {
let server = MockFeedServer::start(|_| {
MockResponse::xml(RSS2_MINIMAL).with_delay(Duration::from_millis(300))
})
.await;
let urls = vec![
("a".to_string(), format!("{}/a.xml", server.url())),
("b".to_string(), format!("{}/b.xml", server.url())),
];
let engine = engine_over(&urls, 900, 1);
let cancelled =
tokio::time::timeout(Duration::from_millis(30), engine.serve_feed("a", || true)).await;
assert!(
cancelled.is_err(),
"the serve must still have been in flight when the timeout dropped it"
);
let served = tokio::time::timeout(Duration::from_secs(10), engine.serve_feed("b", || true))
.await
.expect("b acquired the permit a released on cancellation");
assert!(served.is_some());
assert_eq!(
str_col(&engine.feeds_row("b"), "last_status"),
vec!["fresh"]
);
assert_eq!(
str_col(&engine.feeds_row("a"), "last_status"),
vec!["never"]
);
}
#[tokio::test]
async fn never_fetched_failure_yields_zero_rows_and_error_status() {
let (_interest_guard, _) = capture_events();
let server = MockFeedServer::start(|_| MockResponse::status(500)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["error"]);
assert_eq!(u64_col(&row, "item_count"), vec![None]);
let n = server.requests().len();
assert!(
engine.serve_feed("a", || true).await.is_none(),
"still no window to serve"
);
assert_eq!(
server.requests().len(),
n,
"the second serve went to the negative cache, not the network"
);
}
#[tokio::test]
async fn feeds_row_before_any_scan_is_never_and_issues_no_requests() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
let row = engine.feeds_row("a");
assert_eq!(row.num_rows(), 1);
assert_eq!(str_col(&row, "last_status"), vec!["never"]);
assert_eq!(str_col(&row, "name"), vec!["a"]);
assert_eq!(str_opt_col(&row, "last_error"), vec![None]);
assert_eq!(server.requests().len(), 0);
}
#[tokio::test]
async fn parse_failure_records_stage_and_declared_dialect() {
let (_interest_guard, _) = capture_events();
let server = MockFeedServer::start(|_| {
MockResponse::xml("<rss version=\"2.0\"><channel><title>truncat")
})
.await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["error"]);
let error = str_opt_col(&row, "last_error")[0]
.clone()
.expect("last_error recorded");
assert!(
error.contains("parse failed at strict-parse"),
"last_error names the stage: {error}"
);
assert_eq!(
str_opt_col(&row, "dialect_declared"),
vec![Some("rss-2.0".to_string())],
"the declared-dialect sniff survives a failed parse"
);
}
#[tokio::test]
async fn parse_failure_last_error_quotes_structure_not_prose() {
let (_interest_guard, _) = capture_events();
const SENTINEL: &str = "SHOULD-NOT-LEAK";
#[derive(Debug, PartialEq, Eq)]
enum Fate {
Absent,
Kept,
NoError,
}
let shapes: &[(&str, &str, bool, Fate)] = &[
(
"character data, document truncated",
"<rss version=\"2.0\"><channel><title>SHOULD-NOT-LEAK secret prose",
false,
Fate::Absent,
),
(
"undefined entity in character data, document truncated",
"<rss version=\"2.0\"><channel><title>&SHOULD-NOT-LEAK; truncat",
false,
Fate::Absent,
),
(
"character data beside an out-of-range character reference",
concat!(
r#"<rss version="2.0"><channel>SHOULD-NOT-LEAK �"#,
r#"<title>t</title><link>https://e.example/</link>"#,
r#"<description>d</description></channel></rss>"#,
),
false,
Fate::Absent,
),
(
"entry prose, document failing structurally elsewhere",
concat!(
r#"<feed xmlns="http://www.w3.org/2005/Atom"><title>t</title>"#,
r#"<id>i</id><entry><id>e</id>"#,
r#"<summary>SHOULD-NOT-LEAK the whole article body</summary>"#,
r#"</mismatched></feed>"#,
),
false,
Fate::Absent,
),
(
"feed title, JSON failing on a type elsewhere",
r#"{"version":"https://jsonfeed.org/version/1.1","title":"SHOULD-NOT-LEAK prose","items":"x"}"#,
true,
Fate::Absent,
),
(
"undefined entity in a well-formed title",
concat!(
r#"<rss version="2.0"><channel><title>&SHOULD-NOT-LEAK;</title>"#,
r#"<link>https://e.example/</link><description>d</description>"#,
r#"</channel></rss>"#,
),
false,
Fate::NoError,
),
(
"Atom content type attribute",
concat!(
r#"<feed xmlns="http://www.w3.org/2005/Atom"><title>t</title>"#,
r#"<id>i</id><entry><id>e</id>"#,
r#"<content type="SHOULD-NOT-LEAK">x</content></entry></feed>"#,
),
false,
Fate::Kept,
),
(
"mismatched end tag names the element",
concat!(
r#"<feed xmlns="http://www.w3.org/2005/Atom"><title>t</title>"#,
r#"<id>i</id><entry><id>e</id></SHOULD-NOT-LEAK></feed>"#,
),
false,
Fate::Kept,
),
(
"JSON string where a u64 was declared",
concat!(
r#"{"version":"https://jsonfeed.org/version/1.1","title":"t","#,
r#""items":[{"id":"1","attachments":[{"url":"u","#,
r#""size_in_bytes":"SHOULD-NOT-LEAK"}]}]}"#,
),
true,
Fate::Kept,
),
(
"JSON string where a sequence was declared",
r#"{"version":"https://jsonfeed.org/version/1.1","title":"t","items":"SHOULD-NOT-LEAK"}"#,
true,
Fate::Kept,
),
];
let mut errors_seen = 0;
for (label, body, is_json, fate) in shapes {
let body = (*body).to_string();
let is_json = *is_json;
let server = MockFeedServer::start(move |_| {
if is_json {
MockResponse::new(200, body.clone().into_bytes())
.with_header("content-type", "application/json")
} else {
MockResponse::xml(&body)
}
})
.await;
let engine = test_engine(&server, &[("a", "/f")], 900);
engine.serve_feed("a", || true).await;
let recorded = str_opt_col(&engine.feeds_row("a"), "last_error")[0].clone();
match (fate, &recorded) {
(Fate::NoError, Some(error)) => panic!(
"shape {label:?} was expected to parse cleanly but recorded an error — \
a dependency changed and this row must be re-derived by hand: {error}"
),
(Fate::NoError, None) => {}
(_, None) => panic!(
"shape {label:?} was expected to fail to parse and did not; the assertion \
it exists for never ran"
),
(Fate::Absent, Some(error)) => {
errors_seen += 1;
assert!(
!error.contains(SENTINEL),
"feed content the provider reads as prose reached last_error for \
{label:?}: {error}"
);
}
(Fate::Kept, Some(error)) => {
errors_seen += 1;
assert!(
error.contains(SENTINEL),
"shape {label:?} is documented as quoting its structural fragment and \
no longer does; the module doc's measured list is now wrong: {error}"
);
}
}
if let Some(error) = recorded {
assert!(
error.chars().count() <= MAX_ERROR_CHARS,
"stored error for {label:?} exceeds the cap: {}",
error.chars().count()
);
}
}
assert_eq!(
errors_seen,
shapes
.iter()
.filter(|(_, _, _, f)| *f != Fate::NoError)
.count(),
"every shape but the well-formed one must record an error"
);
assert_eq!(errors_seen, 9);
}
#[tokio::test]
async fn a_json_type_mismatch_quotes_arbitrary_text_up_to_the_cap() {
let (_interest_guard, _) = capture_events();
let filler = "arbitrary feed-chosen text ".repeat(40);
let body = format!(
concat!(
r#"{{"version":"https://jsonfeed.org/version/1.1","title":"t","#,
r#""items":[{{"id":"1","tags":"{}"}}]}}"#,
),
filler
);
assert!(
filler.len() > MAX_ERROR_CHARS,
"the filler must overrun the cap"
);
let server = MockFeedServer::start(move |_| {
MockResponse::new(200, body.clone().into_bytes())
.with_header("content-type", "application/json")
})
.await;
let engine = test_engine(&server, &[("a", "/f.json")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
let error = str_opt_col(&engine.feeds_row("a"), "last_error")[0]
.clone()
.expect("a type mismatch is a parse failure");
assert!(
error.contains("arbitrary feed-chosen text"),
"the offending value is quoted verbatim: {error}"
);
assert_eq!(
error.chars().count(),
MAX_ERROR_CHARS,
"and the cap is what stops it"
);
}
#[tokio::test]
async fn json_unsupported_version_is_body_text_kept_in_last_error() {
let (_interest_guard, _) = capture_events();
let body = r#"{"version":"SHOULD-NOT-LEAK-1.9","title":"t","items":[]}"#;
let server = MockFeedServer::start(move |_| {
MockResponse::new(200, body.as_bytes().to_vec())
.with_header("content-type", "application/json")
})
.await;
let engine = test_engine(&server, &[("a", "/f.json")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
let error = str_opt_col(&engine.feeds_row("a"), "last_error")[0]
.clone()
.expect("an unsupported version is a parse failure");
assert!(
error.contains("unsupported version: SHOULD-NOT-LEAK-1.9"),
"the declared version is kept, because the error is undiagnosable \
without it: {error}"
);
assert!(error.chars().count() <= MAX_ERROR_CHARS);
}
#[tokio::test]
async fn a_huge_json_version_is_still_capped() {
let (_interest_guard, _) = capture_events();
let body = format!(
r#"{{"version":"{}","title":"t","items":[]}}"#,
"v".repeat(20_000)
);
let server = MockFeedServer::start(move |_| {
MockResponse::new(200, body.clone().into_bytes())
.with_header("content-type", "application/json")
})
.await;
let engine = test_engine(&server, &[("a", "/f.json")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
let error = str_opt_col(&engine.feeds_row("a"), "last_error")[0]
.clone()
.expect("last_error recorded");
assert_eq!(
error.chars().count(),
512,
"the documented cap on feeds.last_error, spelled out"
);
}
#[tokio::test]
async fn a_degraded_feed_emits_a_warning_naming_the_feed_and_reason() {
let (_guard, events) = capture_events();
let server = MockFeedServer::start(|_| MockResponse::status(503)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
let warnings = events_with_message(&events, "rss feed degraded");
assert_eq!(
warnings.len(),
1,
"exactly one warning per degraded serve, not zero and not one per attempt"
);
let warning = &warnings[0];
assert_eq!(warning.level, tracing::Level::WARN);
assert_eq!(warning.fields.get("feed").map(String::as_str), Some("a"));
assert_eq!(
warning.fields.get("source").map(String::as_str),
Some("rss_test")
);
assert!(
warning.fields.get("url").is_none(),
"no URL at warn — a subscription URL can carry a private query token"
);
let error = warning
.fields
.get("error")
.expect("the warning carries the reason, not just the feed name");
assert!(
error.contains("503"),
"and the reason is the one recorded in last_error: {error}"
);
assert_eq!(
warning.fields.get("error").cloned(),
str_opt_col(&engine.feeds_row("a"), "last_error")[0].clone(),
"the log and the column carry the same string, so neither can drift"
);
let ok = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let healthy = test_engine(&ok, &[("b", "/f.xml")], 900);
healthy.serve_feed("b", || true).await.expect("b served");
assert_eq!(
events_with_message(&events, "rss feed degraded").len(),
1,
"a successful serve adds no degradation warning"
);
}
#[tokio::test]
async fn egress_blocked_feed_degrades_like_unreachable() {
let (_interest_guard, _) = capture_events();
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let feeds = vec![("a".to_string(), "http://10.1.2.3/f".to_string())];
let cache = Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, feeds.len() + 8));
let policy: Arc<dyn EgressPolicy> = Arc::new(DenyList(vec!["10.1.2.3".parse().unwrap()]));
let engine = engine_with_cache_and_policy(&feeds, 900, 4, cache, policy);
assert!(engine.serve_feed("a", || true).await.is_none());
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["error"]);
let error = str_opt_col(&row, "last_error")[0]
.clone()
.expect("last_error recorded");
assert!(
error.contains("egress blocked"),
"last_error names the refusal: {error}"
);
assert_eq!(
server.requests().len(),
0,
"nothing was connected to at all"
);
}
#[tokio::test]
async fn a_denial_after_a_warm_cache_serves_zero_rows_not_stale() {
let (_interest_guard, _) = capture_events();
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let feeds = vec![("a".to_string(), format!("{}/f.xml", server.url()))];
let cache = Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, 64));
let policy = Arc::new(TogglePolicy(AtomicBool::new(false)));
let engine = engine_with_cache_and_policy(
&feeds,
0,
4,
cache,
Arc::clone(&policy) as Arc<dyn EgressPolicy>,
);
assert_eq!(
engine
.serve_feed("a", || true)
.await
.expect("the allowed fetch warms the cache")
.num_rows(),
1
);
policy.0.store(true, Ordering::SeqCst);
assert!(
engine.serve_feed("a", || true).await.is_none(),
"no stale rows from a refused destination"
);
let row = engine.feeds_row("a");
assert_eq!(str_col(&row, "last_status"), vec!["error"]);
let error = str_opt_col(&row, "last_error")[0]
.clone()
.expect("last_error records the refusal");
assert!(error.contains("egress blocked"), "{error}");
let requests_after_denial = server.requests().len();
assert!(engine.serve_feed("a", || true).await.is_none());
assert_eq!(
server.requests().len(),
requests_after_denial,
"the fuse holds: a denied feed is not re-poked within it"
);
}
#[tokio::test]
async fn conformance_notes_land_in_feeds_row() {
let naked_amp = concat!(
r#"<rss version="2.0"><channel><title>Fish & Chips</title>"#,
r#"<link>https://e.example/</link><description>d</description>"#,
r#"<item><guid>g1</guid><title>t</title></item></channel></rss>"#,
);
let server = MockFeedServer::start(|_| MockResponse::xml(naked_amp)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
engine.serve_feed("a", || true).await.expect("rows served");
let notes = str_opt_col(&engine.feeds_row("a"), "conformance_notes")[0]
.clone()
.expect("notes recorded");
assert!(
notes.contains("sanitation: escaped-naked-ampersands"),
"notes carry the repair: {notes}"
);
}
#[tokio::test]
async fn item_count_and_meta_populate() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
engine.serve_feed("a", || true).await.expect("rows served");
let row = engine.feeds_row("a");
assert_eq!(u64_col(&row, "item_count"), vec![Some(1)]);
assert_eq!(
str_opt_col(&row, "title"),
vec![Some("Minimal Feed".to_string())]
);
assert_eq!(
str_opt_col(&row, "site_url"),
vec![Some("https://feed.example/".to_string())]
);
assert_eq!(
str_opt_col(&row, "description"),
vec![Some("A minimal feed.".to_string())]
);
assert_eq!(
str_opt_col(&row, "dialect"),
vec![Some("rss-2.0".to_string())]
);
assert_eq!(
str_opt_col(&row, "conformance_notes"),
vec![Some("[]".to_string())],
"a clean feed records an empty note list, not NULL"
);
}
#[tokio::test]
async fn false_launch_gate_skips_fetch_and_health_write() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 0);
engine.serve_feed("a", || true).await.expect("first serve");
let before = engine.feeds_row("a");
let requests = server.requests().len();
assert!(
engine.serve_feed("a", || false).await.is_none(),
"a closed gate serves nothing"
);
assert_eq!(server.requests().len(), requests, "and fetches nothing");
let after = engine.feeds_row("a");
assert_eq!(
str_col(&after, "last_status"),
str_col(&before, "last_status"),
"health is untouched: neither fetched nor health-refreshed"
);
assert_eq!(
str_opt_col(&after, "last_error"),
str_opt_col(&before, "last_error")
);
assert_eq!(
u64_col(&after, "item_count"),
u64_col(&before, "item_count")
);
}
#[tokio::test]
async fn within_ttl_serve_refetches_after_its_window_was_evicted() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let urls = vec![
("a".to_string(), format!("{}/a.xml", server.url())),
("b".to_string(), format!("{}/b.xml", server.url())),
];
let probe = engine_over(&urls, 900, 4);
let window_bytes = probe
.serve_feed("a", || true)
.await
.expect("probe serve")
.get_array_memory_size();
let cache = Arc::new(MemoryFeedCache::new(window_bytes + 8, 64));
let engine = engine_with_cache(&urls, 900, 4, cache);
let first = engine.serve_feed("a", || true).await.expect("a served");
assert_eq!(first.num_rows(), 1);
engine.serve_feed("b", || true).await.expect("b served");
let health = engine.feeds_row("a");
assert_eq!(
str_col(&health, "last_status"),
vec!["fresh"],
"eviction keeps the observation, so health still reads fresh"
);
assert_eq!(u64_col(&health, "item_count"), vec![Some(1)]);
assert_eq!(str_opt_col(&health, "last_error"), vec![None]);
assert_eq!(
str_opt_col(&health, "etag"),
vec![None],
"the validators went with the window, so the refetch below has none"
);
let before = server.requests().len();
let again = engine
.serve_feed("a", || true)
.await
.expect("an evicted window inside its TTL refetches instead of serving nothing");
assert_eq!(again.num_rows(), 1, "the rows are actually back");
assert_eq!(str_col(&again, "window_status"), vec!["fresh"]);
let refetches = &server.requests()[before..];
assert_eq!(
refetches.len(),
1,
"exactly one request, not zero and not a retry storm"
);
assert_eq!(refetches[0].path, "/a.xml");
assert_eq!(
refetches[0].header("if-none-match"),
None,
"no window means no validators: an unconditional GET, which is the \
only request that can refill the window"
);
assert_eq!(refetches[0].header("if-modified-since"), None);
}
#[tokio::test]
async fn within_ttl_error_state_still_short_circuits_the_network() {
let (_interest_guard, _) = capture_events();
let server = MockFeedServer::start(|_| MockResponse::status(500)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
assert!(engine.serve_feed("a", || true).await.is_none());
assert_eq!(
str_col(&engine.feeds_row("a"), "last_status"),
vec!["error"]
);
let after_failure = server.requests().len();
assert!(engine.serve_feed("a", || true).await.is_none());
assert_eq!(
server.requests().len(),
after_failure,
"a window-less `error` inside its fuse must not refetch"
);
}
#[tokio::test]
async fn false_gate_still_serves_within_ttl_cache() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
engine.serve_feed("a", || true).await.expect("first serve");
let batch = engine
.serve_feed("a", || false)
.await
.expect("a cache hit has no side effect to gate");
assert_eq!(str_col(&batch, "window_status"), vec!["fresh"]);
assert_eq!(server.requests().len(), 1);
}
#[tokio::test]
async fn launch_gate_is_rechecked_after_acquiring_the_permit() {
let gate_open = Arc::new(AtomicBool::new(true));
let closer = Arc::clone(&gate_open);
let server = MockFeedServer::start(move |_| {
closer.store(false, Ordering::SeqCst);
MockResponse::xml(RSS2_MINIMAL).with_delay(Duration::from_millis(50))
})
.await;
let urls = vec![
("a".to_string(), format!("{}/a.xml", server.url())),
("b".to_string(), format!("{}/b.xml", server.url())),
];
let engine = engine_over(&urls, 900, 1);
let gate = Arc::clone(&gate_open);
let (first, second) = tokio::join!(
biased;
engine.serve_feed("a", || true),
engine.serve_feed("b", move || gate.load(Ordering::SeqCst)),
);
assert!(first.is_some(), "a held the permit and fetched");
assert!(
second.is_none(),
"b's gate closed while it queued for the permit"
);
let paths: Vec<String> = server
.requests()
.iter()
.map(|req| req.path.clone())
.collect();
assert_eq!(
paths,
vec!["/a.xml".to_string()],
"b must never have been fetched"
);
assert_eq!(
str_col(&engine.feeds_row("b"), "last_status"),
vec!["never"]
);
}
#[tokio::test]
async fn unknown_feed_name_serves_nothing_and_reaches_no_state() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let engine = test_engine(&server, &[("a", "/f.xml")], 900);
assert!(engine.serve_feed("' OR 1=1 --", || true).await.is_none());
assert_eq!(engine.feeds_row("' OR 1=1 --").num_rows(), 0);
assert_eq!(server.requests().len(), 0);
}
#[test]
fn parse_error_message_caps_the_composed_string_not_just_the_reason() {
let message = parse_error_message("strict-parse", &"x".repeat(4_000));
assert_eq!(message.chars().count(), MAX_ERROR_CHARS);
assert!(message.starts_with("parse failed at strict-parse: "));
assert_eq!(
parse_error_message("refused-internal-dtd", "internal DTD subset refused"),
"parse failed at refused-internal-dtd: internal DTD subset refused"
);
}
#[tokio::test]
async fn a_parse_that_outlives_the_fuse_degrades_to_a_timeout_failure() {
let mut body = String::from(r#"<rss version="2.0"><channel><title>t</title>"#);
for i in 0..2_000 {
body.push_str(&format!("<item><guid>g{i}</guid><title>x</title></item>"));
}
body.push_str("</channel></rss>");
let failure = parse_off_worker(body.into_bytes(), None, Duration::ZERO)
.await
.unwrap_err();
assert_eq!(failure.stage, "timeout");
assert!(
failure.reason.contains("did not finish"),
"{}",
failure.reason
);
}
#[tokio::test]
async fn a_parse_within_the_fuse_returns_the_document() {
let document = parse_off_worker(RSS2_MINIMAL.as_bytes().to_vec(), None, PARSE_TIMEOUT)
.await
.expect("a well-formed document parses within the fuse");
assert_eq!(document.items.len(), 1);
}
#[tokio::test]
async fn a_redirected_feed_revalidates_at_its_landing_url() {
let server = MockFeedServer::start(|req| {
if req.path.starts_with("/a") {
MockResponse::status(301).with_header("location", "/b.xml")
} else if req.header("if-none-match").as_deref() == Some("\"v1\"") {
MockResponse::status(304)
} else {
MockResponse::xml(RSS2_MINIMAL).with_header("etag", "\"v1\"")
}
})
.await;
let engine = test_engine(&server, &[("a", "/a.xml")], 0);
assert_eq!(
engine
.serve_feed("a", || true)
.await
.expect("the redirect is followed and the body served")
.num_rows(),
1
);
let first = server.requests();
assert_eq!(
first.iter().map(|r| r.path.as_str()).collect::<Vec<_>>(),
vec!["/a.xml", "/b.xml"],
"scan 1 follows the redirect"
);
let batch = engine
.serve_feed("a", || true)
.await
.expect("the 304 serves the cached window");
assert_eq!(str_col(&batch, "window_status"), vec!["revalidated"]);
let second: Vec<String> = server.requests()[first.len()..]
.iter()
.map(|r| r.path.clone())
.collect();
assert_eq!(
second,
vec!["/b.xml"],
"no hop through the redirector: the conditional request goes to \
the URL that issued the validators"
);
assert_eq!(
server.requests().last().unwrap().header("if-none-match"),
Some("\"v1\"".to_string()),
"and it is conditional — otherwise the 304 is unreachable"
);
assert_eq!(
str_col(&engine.feeds_row("a"), "last_status"),
vec!["revalidated"]
);
assert!(str_col(&engine.feeds_row("a"), "url")[0].ends_with("/a.xml"));
let before_third = server.requests().len();
engine.serve_feed("a", || true).await.expect("served again");
let third: Vec<String> = server.requests()[before_third..]
.iter()
.map(|r| r.path.clone())
.collect();
assert_eq!(third, vec!["/b.xml"]);
}
#[tokio::test]
async fn a_reprobe_starts_from_the_configured_url_and_follows_drift() {
let calls = Arc::new(AtomicUsize::new(0));
let calls2 = Arc::clone(&calls);
let server = MockFeedServer::start(move |req| {
if req.path.starts_with("/a") {
let nth = calls2.fetch_add(1, Ordering::SeqCst);
let to = if nth == 0 { "/b.xml" } else { "/c.xml" };
MockResponse::status(301).with_header("location", to)
} else {
MockResponse::xml(RSS2_MINIMAL).with_header("etag", "\"v1\"")
}
})
.await;
let mut engine = test_engine(&server, &[("a", "/a.xml")], 0);
engine.reprobe_interval = Duration::ZERO;
engine.serve_feed("a", || true).await.expect("first serve");
let before = server.requests().len();
engine.serve_feed("a", || true).await.expect("second serve");
let probe: Vec<String> = server.requests()[before..]
.iter()
.map(|r| r.path.clone())
.collect();
assert_eq!(
probe,
vec!["/a.xml", "/c.xml"],
"the probe re-derives the landing URL and follows it to its new target"
);
assert_eq!(
server.requests()[before].header("if-none-match"),
None,
"a probe is unconditional: /b's validators mean nothing to /a"
);
engine.reprobe_interval = Duration::from_secs(3600);
let before_third = server.requests().len();
engine.serve_feed("a", || true).await.expect("third serve");
let third: Vec<String> = server.requests()[before_third..]
.iter()
.map(|r| r.path.clone())
.collect();
assert_eq!(third, vec!["/c.xml"], "drift adopted");
}
#[tokio::test]
async fn a_content_refresh_does_not_reset_the_drift_clock() {
let server = MockFeedServer::start(|req| {
if req.path.starts_with("/a") {
MockResponse::status(301).with_header("location", "/b.xml")
} else {
MockResponse::xml(RSS2_MINIMAL)
}
})
.await;
let mut engine = test_engine(&server, &[("a", "/a.xml")], 0);
engine.reprobe_interval = Duration::from_millis(400);
engine.serve_feed("a", || true).await.expect("first serve");
tokio::time::sleep(Duration::from_millis(300)).await;
engine.serve_feed("a", || true).await.expect("refresh");
let before = server.requests().len();
assert_eq!(
server.requests()[before - 1].path,
"/b.xml",
"the refresh went straight to the landing URL"
);
tokio::time::sleep(Duration::from_millis(200)).await;
engine.serve_feed("a", || true).await.expect("probe due");
let probe: Vec<String> = server.requests()[before..]
.iter()
.map(|r| r.path.clone())
.collect();
assert_eq!(
probe,
vec!["/a.xml", "/b.xml"],
"the interval elapsed since the *probe*, not since the last refresh"
);
}
#[test]
fn the_reprobe_interval_is_a_clamped_multiple_of_the_ttl() {
assert_eq!(
reprobe_interval(Duration::from_secs(900)),
Duration::from_secs(900 * 24),
"the default TTL gives 24 refreshes between probes"
);
assert_eq!(
reprobe_interval(Duration::ZERO),
MIN_REPROBE_INTERVAL,
"an always-live feed still does not probe on every scan"
);
assert_eq!(reprobe_interval(MAX_TTL), MAX_REPROBE_INTERVAL);
}
#[tokio::test]
async fn a_panicked_parse_does_not_quote_the_payload() {
let join_error = tokio::task::spawn_blocking(|| panic!("feed-authored payload"))
.await
.expect_err("the task panicked");
assert!(join_error.is_panic());
assert!(
join_error.to_string().contains("feed-authored payload"),
"guard: if JoinError stops quoting payloads, this test's premise is gone"
);
let failure = join_failure(&join_error);
assert_eq!(failure.stage, "panic");
assert!(
!failure.reason.contains("feed-authored payload"),
"the payload must not reach last_error: {}",
failure.reason
);
assert!(
failure.reason.contains("server log"),
"and the reason must say where the payload actually is: {}",
failure.reason
);
}
#[test]
fn the_parse_fuse_tracks_max_response_bytes() {
assert_eq!(
parse_fuse(DEFAULT_MAX_RESPONSE_BYTES),
Duration::from_secs(10),
"the default cap keeps the original ten-second fuse"
);
assert_eq!(
parse_fuse(1),
Duration::from_secs(10),
"floored at one unit"
);
assert_eq!(
parse_fuse(DEFAULT_MAX_RESPONSE_BYTES + 1),
Duration::from_secs(20),
"partial units round up"
);
assert_eq!(
parse_fuse(DEFAULT_MAX_RESPONSE_BYTES * 10),
Duration::from_secs(100)
);
assert_eq!(
parse_fuse(u64::MAX),
MAX_PARSE_FUSE,
"capped, not overflowing"
);
}
#[tokio::test]
async fn a_parse_fuse_past_the_scan_timeout_warns_at_construction() {
let (_guard, events) = capture_events();
let subscriptions = vec![ResolvedSubscription {
name: "a".to_string(),
url: "https://feed.example/f.xml".to_string(),
}];
let mut config = inline_config(vec![FeedSubscription {
url: "https://feed.example/f.xml".to_string(),
name: Some("a".to_string()),
}]);
config.max_response_bytes = DEFAULT_MAX_RESPONSE_BYTES * 7;
let fetcher = FeedFetcher::new(
None,
Duration::from_secs(5),
config.max_response_bytes,
config.user_agent.clone(),
)
.expect("build the test fetcher");
let _engine = RssEngine::with_parts(
"rss_test".to_string(),
subscriptions,
&config,
fetcher,
Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, 8)),
);
let warned = events_with_message(
&events,
"rss parse fuse meets or exceeds the scan timeout; raise scan_timeout_seconds",
);
assert_eq!(warned.len(), 1, "one warning per engine built");
assert_eq!(warned[0].level, tracing::Level::WARN);
assert_eq!(warned[0].field("parse_fuse_seconds"), Some("70"));
assert_eq!(warned[0].field("scan_timeout_seconds"), Some("60"));
}
const RSS2_TWO_ITEMS: &str = concat!(
r#"<rss version="2.0"><channel>"#,
r#"<title>Minimal Feed</title>"#,
r#"<link>https://feed.example/</link>"#,
r#"<description>d</description>"#,
r#"<item><guid>g1</guid><title>one</title></item>"#,
r#"<item><guid>g2</guid><title>two</title></item>"#,
r#"</channel></rss>"#
);
#[tokio::test]
async fn a_slower_concurrent_fetch_cannot_regress_the_window() {
let calls = Arc::new(AtomicUsize::new(0));
let calls2 = Arc::clone(&calls);
let server = MockFeedServer::start(move |_req| {
if calls2.fetch_add(1, Ordering::SeqCst) == 0 {
MockResponse::xml(RSS2_MINIMAL).with_delay(Duration::from_millis(400))
} else {
MockResponse::xml(RSS2_TWO_ITEMS)
}
})
.await;
let cache = Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, 64));
let feeds = vec![("a".to_string(), format!("{}/f.xml", server.url()))];
let engine = engine_with_cache(&feeds, 0, 6, Arc::clone(&cache) as Arc<dyn FeedCache>);
let slow = engine.serve_feed("a", || true);
let fast = async {
tokio::time::sleep(Duration::from_millis(100)).await;
engine.serve_feed("a", || true).await
};
let (slow_batch, fast_batch) = tokio::join!(slow, fast);
assert_eq!(slow_batch.expect("slow serve emits its read").num_rows(), 1);
assert_eq!(fast_batch.expect("fast serve emits its read").num_rows(), 2);
let snap = cache.snapshot("a", Instant::now());
assert_eq!(
snap.window.as_ref().unwrap().batch.num_rows(),
2,
"the slower fetch's commit must be dropped, not applied last"
);
assert_eq!(snap.observation.item_count, Some(2));
}
#[tokio::test]
async fn a_stale_304_end_to_end_does_not_relabel_the_new_window() {
let calls = Arc::new(AtomicUsize::new(0));
let calls2 = Arc::clone(&calls);
let server = MockFeedServer::start(move |_req| {
match calls2.fetch_add(1, Ordering::SeqCst) {
0 => MockResponse::xml(RSS2_MINIMAL).with_header("etag", "\"v1\""),
1 => MockResponse::status(304).with_delay(Duration::from_millis(400)),
_ => MockResponse::xml(RSS2_TWO_ITEMS),
}
})
.await;
let cache = Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, 64));
let feeds = vec![("a".to_string(), format!("{}/f.xml", server.url()))];
let engine = engine_with_cache(&feeds, 0, 6, Arc::clone(&cache) as Arc<dyn FeedCache>);
engine.serve_feed("a", || true).await.expect("primed");
let slow_304 = engine.serve_feed("a", || true);
let fast_200 = async {
tokio::time::sleep(Duration::from_millis(100)).await;
engine.serve_feed("a", || true).await
};
let (revalidated_read, fresh_read) = tokio::join!(slow_304, fast_200);
assert_eq!(
revalidated_read
.expect("the 304 serves its snapshot window")
.num_rows(),
1
);
assert_eq!(fresh_read.expect("the 200 serves its fetch").num_rows(), 2);
let snap = cache.snapshot("a", Instant::now());
assert!(
matches!(snap.observation.last_status, FeedStatus::Fresh),
"not `revalidated` — the 304 never saw this window: {:?}",
snap.observation.last_status
);
assert_eq!(snap.window.as_ref().unwrap().batch.num_rows(), 2);
}
#[tokio::test]
async fn an_absurd_ttl_is_clamped_rather_than_overflowing() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let (_guard, events) = capture_events();
let engine = test_engine(&server, &[("a", "/f.xml")], u64::MAX);
let batch = engine.serve_feed("a", || true).await.expect("rows served");
assert_eq!(str_col(&batch, "window_status"), vec!["fresh"]);
let clamped =
events_with_message(&events, "rss ttl_seconds clamped to the engine's ceiling");
assert_eq!(clamped.len(), 1, "one warning per engine built");
assert_eq!(clamped[0].level, tracing::Level::WARN);
assert_eq!(
clamped[0].field("configured_ttl_seconds"),
Some(u64::MAX.to_string().as_str())
);
assert_eq!(
clamped[0].field("effective_ttl_seconds"),
Some(MAX_TTL.as_secs().to_string().as_str())
);
}
#[tokio::test]
async fn an_absurd_max_concurrent_is_clamped_rather_than_panicking() {
let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
let (_guard, events) = capture_events();
let feeds = vec![("a".to_string(), format!("{}/f.xml", server.url()))];
let cache = Arc::new(MemoryFeedCache::new(CACHE_MAX_BYTES, 64));
let engine = engine_with_cache(&feeds, 900, usize::MAX, cache);
let batch = engine.serve_feed("a", || true).await.expect("rows served");
assert_eq!(str_col(&batch, "window_status"), vec!["fresh"]);
let clamped = events_with_message(
&events,
"rss max_concurrent clamped to the semaphore's ceiling",
);
assert_eq!(clamped.len(), 1, "one warning per engine built");
assert_eq!(clamped[0].level, tracing::Level::WARN);
assert_eq!(
clamped[0].field("configured_max_concurrent"),
Some(usize::MAX.to_string().as_str())
);
assert_eq!(
clamped[0].field("effective_max_concurrent"),
Some(Semaphore::MAX_PERMITS.to_string().as_str())
);
}
#[test]
fn arm_saturates_to_now_rather_than_panicking_on_an_unrepresentable_add() {
let now = Instant::now();
assert_eq!(
arm(now, Duration::MAX),
now,
"an add no platform Instant can represent leaves the feed due immediately"
);
assert!(arm(now, Duration::from_secs(30)) > now);
}
#[tokio::test]
async fn a_failure_arms_the_ttls_quarter_not_the_floor() {
let (_interest_guard, _) = capture_events();
let server = MockFeedServer::start(|_| MockResponse::status(500)).await;
let urls = vec![("a".to_string(), format!("{}/f.xml", server.url()))];
let cache = Arc::new(RecordsFailureArm::new(MemoryFeedCache::new(
CACHE_MAX_BYTES,
8,
)));
let engine = engine_with_cache(&urls, 900, 4, Arc::clone(&cache) as Arc<dyn FeedCache>);
let before = Instant::now();
assert!(engine.serve_feed("a", || true).await.is_none());
let after = Instant::now();
let armed = cache.armed_instants();
assert_eq!(armed.len(), 1, "one failed serve records one arm");
let lower = armed[0].duration_since(after);
let upper = armed[0].duration_since(before);
let expected = Duration::from_secs(225);
assert!(
lower <= expected && expected <= upper,
"the engine armed a fuse in {lower:?}..={upper:?}; ttl_seconds 900 must arm \
failure_fuse(900s) = 225s — neither the 30s floor nor the 300s ceiling"
);
}
}