use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use arrow::array::{Array, RecordBatch, UInt16Array};
use datafusion::prelude::SessionContext;
use super::config::{FeedSubscription, RssConfig, inline_config};
use super::egress::AllowAll;
use super::fetch::MAX_ATTEMPTS;
use super::register_rss_tables_with_policy;
use super::testutil::{
MockFeedServer, MockResponse, RecordedRequest, str_col, str_opt_col, total_rows,
};
use crate::sources::hierarchy::HierarchyLevel;
const SOURCE: &str = "news";
const QUERY_CEILING: Duration = Duration::from_secs(60);
fn rss_with(guids: &[&str]) -> String {
let mut doc = String::from(r#"<rss version="2.0"><channel><title>Mock Feed</title>"#);
doc.push_str(r#"<link>https://feed.example/</link>"#);
doc.push_str(r#"<description>A mock feed.</description>"#);
for guid in guids {
doc.push_str(&format!(
"<item><guid>{guid}</guid><title>{guid} title</title>\
<link>https://feed.example/{guid}</link></item>"
));
}
doc.push_str("</channel></rss>");
doc
}
fn default_body(path: &str) -> String {
rss_with(&[&format!("{path}#1")])
}
fn default_guid(path: &str) -> String {
format!("{path}#1")
}
#[derive(Clone)]
struct Canned {
status: u16,
headers: Vec<(String, String)>,
body: Vec<u8>,
delay: Option<Duration>,
}
impl Canned {
fn xml(body: &str) -> Self {
Self {
status: 200,
headers: vec![("content-type".to_string(), "application/xml".to_string())],
body: body.as_bytes().to_vec(),
delay: None,
}
}
fn status(status: u16) -> Self {
Self {
status,
headers: Vec::new(),
body: Vec::new(),
delay: None,
}
}
fn bytes(status: u16, body: Vec<u8>) -> Self {
Self {
status,
headers: Vec::new(),
body,
delay: None,
}
}
fn with_header(mut self, name: &str, value: &str) -> Self {
self.headers.push((name.to_string(), value.to_string()));
self
}
fn with_delay(mut self, delay: Duration) -> Self {
self.delay = Some(delay);
self
}
fn into_response(self) -> MockResponse {
let mut response = MockResponse::new(self.status, self.body);
for (name, value) in self.headers {
response = response.with_header(&name, &value);
}
match self.delay {
Some(delay) => response.with_delay(delay),
None => response,
}
}
}
struct FeedScript {
steps: Vec<Canned>,
served: usize,
etag: Option<String>,
}
impl FeedScript {
fn always(canned: Canned) -> Self {
Self::steps(vec![canned])
}
fn steps(steps: Vec<Canned>) -> Self {
assert!(!steps.is_empty(), "a feed script needs at least one step");
Self {
steps,
served: 0,
etag: None,
}
}
fn with_etag(mut self, etag: &str) -> Self {
self.etag = Some(etag.to_string());
self
}
fn answer(&mut self, request: &RecordedRequest) -> MockResponse {
self.served += 1;
if let Some(etag) = &self.etag
&& request.header("if-none-match").as_deref() == Some(etag.as_str())
{
return MockResponse::status(304);
}
let step = self.steps[(self.served - 1).min(self.steps.len() - 1)].clone();
match (&self.etag, step.status) {
(Some(etag), 200) => step.with_header("etag", etag).into_response(),
_ => step.into_response(),
}
}
}
type Scripts = Arc<Mutex<HashMap<String, FeedScript>>>;
struct TestNews {
server: MockFeedServer,
ctx: SessionContext,
scripts: Scripts,
}
impl TestNews {
async fn start(feeds: &[(&str, &str)], tune: impl FnOnce(&mut RssConfig)) -> Self {
let scripts: Scripts = Arc::new(Mutex::new(HashMap::new()));
let handler_scripts = Arc::clone(&scripts);
let server = MockFeedServer::start(move |request| {
let mut scripts = handler_scripts
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match scripts.get_mut(&request.path) {
Some(script) => script.answer(request),
None => MockResponse::status(404),
}
})
.await;
let mut subscriptions = Vec::with_capacity(feeds.len());
for (name, target) in feeds {
let url = if target.starts_with("http://") || target.starts_with("https://") {
(*target).to_string()
} else {
scripts
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.insert(
(*target).to_string(),
FeedScript::always(Canned::xml(&default_body(target))),
);
format!("{}{target}", server.url())
};
subscriptions.push(FeedSubscription {
url,
name: Some((*name).to_string()),
});
}
let mut config = inline_config(subscriptions);
config.request_timeout_seconds = 5;
config.scan_timeout_seconds = 20;
tune(&mut config);
let mut ctx = SessionContext::new();
register_rss_tables_with_policy(
&mut ctx,
SOURCE,
Some(&config),
false,
HierarchyLevel::Catalog,
Arc::new(AllowAll),
)
.await
.expect("registration succeeds");
Self {
server,
ctx,
scripts,
}
}
fn script(&self, path: &str, script: FeedScript) {
self.scripts
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.insert(path.to_string(), script);
}
async fn sql(&self, sql: &str) -> Vec<RecordBatch> {
tokio::time::timeout(QUERY_CEILING, async {
self.ctx
.sql(sql)
.await
.unwrap_or_else(|e| panic!("plan {sql:?}: {e}"))
.collect()
.await
.unwrap_or_else(|e| panic!("execute {sql:?}: {e}"))
})
.await
.unwrap_or_else(|_| panic!("{sql:?} did not finish within {QUERY_CEILING:?}"))
}
fn requests(&self) -> Vec<RecordedRequest> {
self.server.requests()
}
fn paths(&self) -> Vec<String> {
self.requests()
.into_iter()
.map(|request| request.path)
.collect()
}
fn request_count(&self) -> usize {
self.server.requests().len()
}
fn context(&self) -> SessionContext {
self.ctx.clone()
}
async fn await_requests(&self, count: usize, within: Duration) {
let deadline = Instant::now() + within;
while self.request_count() < count {
assert!(
Instant::now() < deadline,
"only {} of {count} requests arrived within {within:?}",
self.request_count()
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
}
fn col(batches: &[RecordBatch], name: &str) -> Vec<String> {
batches.iter().flat_map(|b| str_col(b, name)).collect()
}
fn opt_col(batches: &[RecordBatch], name: &str) -> Vec<Option<String>> {
batches.iter().flat_map(|b| str_opt_col(b, name)).collect()
}
fn u16_col(batches: &[RecordBatch], name: &str) -> Vec<Option<u16>> {
batches
.iter()
.flat_map(|batch| {
let index = batch
.schema()
.index_of(name)
.unwrap_or_else(|e| panic!("batch has no column {name:?}: {e}"));
let column = batch
.column(index)
.as_any()
.downcast_ref::<UInt16Array>()
.unwrap_or_else(|| panic!("column {name:?} is not UInt16"));
(0..column.len())
.map(|row| column.is_valid(row).then(|| column.value(row)))
.collect::<Vec<_>>()
})
.collect()
}
fn sorted(mut paths: Vec<String>) -> Vec<String> {
paths.sort();
paths
}
#[track_caller]
fn assert_contains(actual: &str, needle: &str) {
assert!(actual.contains(needle), "expected {needle:?} in {actual:?}");
}
#[tokio::test]
async fn ac1_registration_is_zero_network() {
let news = TestNews::start(&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")], |_| {}).await;
assert_eq!(
news.request_count(),
0,
"registration performed network I/O"
);
let feeds = news
.sql("SELECT name, last_status, last_error FROM news.main.feeds ORDER BY name")
.await;
assert_eq!(col(&feeds, "name"), vec!["a", "b", "c"]);
assert_eq!(
col(&feeds, "last_status"),
vec!["never", "never", "never"],
"no subscription has been attempted"
);
assert_eq!(opt_col(&feeds, "last_error"), vec![None, None, None]);
assert_eq!(
news.request_count(),
0,
"a feeds scan performed network I/O"
);
}
#[tokio::test]
async fn ac2_full_scan_fetches_all_where_prunes_to_one() {
let news = TestNews::start(
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
|config| config.ttl_seconds = 0,
)
.await;
let all = news
.sql("SELECT feed, guid FROM news.main.items ORDER BY feed")
.await;
assert_eq!(col(&all, "feed"), vec!["a", "b", "c"]);
assert_eq!(
col(&all, "guid"),
vec![
default_guid("/a.xml"),
default_guid("/b.xml"),
default_guid("/c.xml")
],
"each row carries the guid of the feed it was fetched from"
);
assert_eq!(
sorted(news.paths()),
vec!["/a.xml", "/b.xml", "/c.xml"],
"a full scan visits every subscription"
);
let one = news
.sql("SELECT feed, guid FROM news.main.items WHERE feed = 'b'")
.await;
assert_eq!(col(&one, "feed"), vec!["b"]);
assert_eq!(col(&one, "guid"), vec![default_guid("/b.xml")]);
assert_eq!(
news.paths()[3..],
["/b.xml".to_string()],
"the pruned scan fetched exactly the one feed the predicate names"
);
}
#[tokio::test]
async fn ac3_ttl_and_304_paths() {
let cached = TestNews::start(
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
|config| config.ttl_seconds = 900,
)
.await;
let first = cached
.sql("SELECT feed, window_status FROM news.main.items")
.await;
assert_eq!(total_rows(&first), 3);
assert_eq!(cached.request_count(), 3, "one fetch per feed");
let second = cached
.sql("SELECT feed, window_status FROM news.main.items")
.await;
assert_eq!(total_rows(&second), 3, "the cached windows still serve");
assert_eq!(
cached.request_count(),
3,
"a second scan within the TTL issues no request"
);
let live = TestNews::start(&[("a", "/a.xml")], |config| config.ttl_seconds = 0).await;
live.script(
"/a.xml",
FeedScript::always(Canned::xml(&default_body("/a.xml"))).with_etag("\"v1\""),
);
let fresh = live
.sql("SELECT guid, window_status FROM news.main.items")
.await;
assert_eq!(col(&fresh, "window_status"), vec!["fresh"]);
assert_eq!(live.request_count(), 1);
assert_eq!(
live.requests()[0].header("if-none-match"),
None,
"the first request has no validator to send"
);
let revalidated = live
.sql("SELECT guid, window_status FROM news.main.items")
.await;
assert_eq!(live.request_count(), 2, "the expired feed was revalidated");
assert_eq!(
live.requests()[1].header("if-none-match").as_deref(),
Some("\"v1\""),
"the second request carries the etag the first response set"
);
assert_eq!(
col(&revalidated, "window_status"),
vec!["revalidated"],
"a 304 serves the cached window, relabelled"
);
assert_eq!(
col(&revalidated, "guid"),
vec![default_guid("/a.xml")],
"and the rows are the ones the 304 vouches for"
);
let feeds = live
.sql("SELECT last_status, http_status, etag FROM news.main.feeds")
.await;
assert_eq!(col(&feeds, "last_status"), vec!["revalidated"]);
assert_eq!(u16_col(&feeds, "http_status"), vec![Some(304)]);
assert_eq!(
opt_col(&feeds, "etag"),
vec![Some("\"v1\"".to_string())],
"the validator survives the revalidation"
);
}
#[tokio::test]
async fn ac4_dead_feed_isolation_with_stale_stamp() {
let news = TestNews::start(
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
|config| config.ttl_seconds = 0,
)
.await;
let healthy = news
.sql("SELECT feed, guid, title, window_status FROM news.main.items ORDER BY feed")
.await;
assert_eq!(col(&healthy, "feed"), vec!["a", "b", "c"]);
assert_eq!(
col(&healthy, "window_status"),
vec!["fresh", "fresh", "fresh"]
);
news.script("/b.xml", FeedScript::always(Canned::status(500)));
let degraded = news
.sql("SELECT feed, guid, title, window_status FROM news.main.items ORDER BY feed")
.await;
assert_eq!(
col(°raded, "feed"),
vec!["a", "b", "c"],
"the dead feed still contributes its cached window"
);
assert_eq!(
col(°raded, "window_status"),
vec!["fresh", "stale-error", "fresh"],
"and only that feed's rows are relabelled"
);
assert_eq!(
col(°raded, "guid"),
col(&healthy, "guid"),
"the served rows are the same rows, dead feed included"
);
assert_eq!(
opt_col(°raded, "title"),
opt_col(&healthy, "title"),
"a neighbour's values are unaffected by another feed's failure"
);
let feeds = news
.sql(
"SELECT name, last_status, http_status, last_error, item_count \
FROM news.main.feeds ORDER BY name",
)
.await;
assert_eq!(
col(&feeds, "last_status"),
vec!["fresh", "stale-error", "fresh"]
);
assert_eq!(
u16_col(&feeds, "http_status"),
vec![Some(200), Some(500), Some(200)]
);
let errors = opt_col(&feeds, "last_error");
assert_eq!(errors[0], None, "a healthy feed carries no error");
assert_eq!(errors[2], None, "a healthy feed carries no error");
assert_contains(
errors[1]
.as_deref()
.expect("the dead feed records an error"),
"http status 500",
);
assert_eq!(
news.paths().iter().filter(|path| *path == "/b.xml").count(),
1 + MAX_ATTEMPTS as usize,
"the first scan's one success, then one exhausted attempt budget"
);
}
#[tokio::test]
async fn ac13_feeds_scan_zero_requests_even_after_failure() {
let news = TestNews::start(&[("a", "/a.xml"), ("b", "/b.xml")], |config| {
config.ttl_seconds = 0
})
.await;
news.sql("SELECT guid FROM news.main.items").await;
news.script("/b.xml", FeedScript::always(Canned::status(500)));
news.sql("SELECT guid FROM news.main.items").await;
let after_failure = news.request_count();
let feeds = news
.sql("SELECT name, last_status FROM news.main.feeds ORDER BY name")
.await;
assert_eq!(
col(&feeds, "last_status"),
vec!["fresh", "stale-error"],
"the failure is visible in the health table"
);
assert_eq!(
news.request_count(),
after_failure,
"a feeds scan right after a failure issued a request"
);
news.sql("SELECT * FROM news.main.feeds").await;
assert_eq!(
news.request_count(),
after_failure,
"a second feeds scan issued a request"
);
let items = news
.sql("SELECT feed, window_status FROM news.main.items ORDER BY feed")
.await;
assert_eq!(
col(&items, "window_status"),
vec!["fresh", "stale-error"],
"the dead feed still serves its stale window from cache"
);
assert_eq!(
news.paths().iter().filter(|path| *path == "/b.xml").count(),
1 + MAX_ATTEMPTS as usize,
"the dead feed was re-poked inside its failure fuse"
);
}
#[tokio::test]
async fn decompressed_cap_rejects_gzip_bomb() {
const BOMB: &[u8] = include_bytes!("fixtures/bomb.xml.gz");
let news = TestNews::start(
&[("bomb", "/bomb.xml"), ("healthy", "/healthy.xml")],
|_| {},
)
.await;
assert!(
(BOMB.len() as u64) < 64 * 1024,
"the body served is {} bytes, which is not a decompression test",
BOMB.len()
);
news.script(
"/bomb.xml",
FeedScript::always(
Canned::bytes(200, BOMB.to_vec())
.with_header("content-type", "application/xml")
.with_header("content-encoding", "gzip"),
),
);
let items = news
.sql("SELECT feed, guid FROM news.main.items ORDER BY feed")
.await;
assert_eq!(
col(&items, "feed"),
vec!["healthy"],
"the bomb contributes no rows; its healthy sibling is unaffected"
);
assert_eq!(col(&items, "guid"), vec![default_guid("/healthy.xml")]);
let feeds = news
.sql("SELECT name, last_status, last_error FROM news.main.feeds ORDER BY name")
.await;
assert_eq!(col(&feeds, "last_status"), vec!["error", "fresh"]);
assert_eq!(
opt_col(&feeds, "last_error")[0].as_deref(),
Some("response exceeded 5242880 bytes"),
"the refusal names the decoded-byte cap it broke"
);
assert_eq!(
news.paths()
.iter()
.filter(|path| *path == "/bomb.xml")
.count(),
1,
"an over-cap body is terminal, not retried"
);
}
#[tokio::test]
async fn request_timeout_isolates_slow_feed() {
let news = TestNews::start(&[("fast", "/fast.xml"), ("slow", "/slow.xml")], |config| {
config.request_timeout_seconds = 1;
})
.await;
news.script(
"/slow.xml",
FeedScript::always(
Canned::xml(&default_body("/slow.xml")).with_delay(Duration::from_secs(3)),
),
);
let started = Instant::now();
let items = tokio::time::timeout(
Duration::from_secs(30),
news.sql("SELECT feed, guid FROM news.main.items ORDER BY feed"),
)
.await
.expect("the slow feed must time out rather than hang the scan");
let elapsed = started.elapsed();
assert_eq!(
col(&items, "feed"),
vec!["fast"],
"the fast feed serves while the slow one is still being retried"
);
assert!(
elapsed < Duration::from_secs(20),
"the query took {elapsed:?}, so it was the scan deadline that cut it, not the \
per-request timeout"
);
let feeds = news
.sql("SELECT name, last_status, last_error FROM news.main.feeds ORDER BY name")
.await;
assert_eq!(col(&feeds, "last_status"), vec!["fresh", "error"]);
assert_eq!(
opt_col(&feeds, "last_error")[1].as_deref(),
Some("request timed out after 1s"),
"the recorded reason is the per-request timeout, not the scan deadline"
);
assert_eq!(
news.paths()
.iter()
.filter(|path| *path == "/slow.xml")
.count(),
MAX_ATTEMPTS as usize,
"a timeout is retryable, so the slow feed spends its whole attempt budget"
);
}
#[tokio::test]
async fn retry_after_is_honored_within_scan() {
let news = TestNews::start(&[("a", "/a.xml")], |_| {}).await;
news.script(
"/a.xml",
FeedScript::steps(vec![
Canned::status(429).with_header("retry-after", "1"),
Canned::xml(&default_body("/a.xml")),
]),
);
let started = Instant::now();
let items = news.sql("SELECT feed, guid FROM news.main.items").await;
let elapsed = started.elapsed();
assert_eq!(col(&items, "guid"), vec![default_guid("/a.xml")]);
assert_eq!(
news.paths(),
vec!["/a.xml", "/a.xml"],
"the retry happened, and only one of them"
);
assert!(
elapsed >= Duration::from_secs(1),
"the retry waited {elapsed:?}, which is less than the Retry-After it was given"
);
let feeds = news
.sql("SELECT last_status, http_status FROM news.main.feeds")
.await;
assert_eq!(col(&feeds, "last_status"), vec!["fresh"]);
assert_eq!(u16_col(&feeds, "http_status"), vec![Some(200)]);
}
#[tokio::test]
async fn limit_stops_launching_fetches() {
let news = TestNews::start(
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
|config| {
config.ttl_seconds = 0;
config.max_concurrent = 1;
},
)
.await;
let limited = news.sql("SELECT guid FROM news.main.items LIMIT 1").await;
assert_eq!(total_rows(&limited), 1);
assert_eq!(
news.request_count(),
1,
"one row of LIMIT is one fetch: the partitions past it must not launch one"
);
let top_k = news
.sql("SELECT guid FROM news.main.items ORDER BY guid LIMIT 1")
.await;
assert_eq!(total_rows(&top_k), 1);
assert_eq!(
col(&top_k, "guid"),
vec![default_guid("/a.xml")],
"the Top-K's winner is the smallest guid across every feed"
);
assert_eq!(
sorted(news.paths()[1..].to_vec()),
vec!["/a.xml", "/b.xml", "/c.xml"],
"a Top-K consumes every partition, so the limit cannot gate any fetch"
);
}
#[tokio::test]
async fn cancellation_stops_further_fetches() {
let news = TestNews::start(
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
|config| config.max_concurrent = 1,
)
.await;
for path in ["/a.xml", "/b.xml", "/c.xml"] {
news.script(
path,
FeedScript::always(Canned::xml(&default_body(path)).with_delay(Duration::from_secs(5))),
);
}
let ctx = news.context();
let query = tokio::spawn(async move {
ctx.sql("SELECT guid FROM news.main.items")
.await
.expect("plan")
.collect()
.await
});
news.await_requests(1, Duration::from_secs(10)).await;
query.abort();
assert!(
query
.await
.expect_err("the query was aborted")
.is_cancelled(),
"the task ended for some reason other than the abort"
);
tokio::time::sleep(Duration::from_secs(8)).await;
assert_eq!(
news.paths().len(),
1,
"the cancelled scan launched further fetches: {:?}",
news.paths()
);
}
const FIVE_FEEDS: [(&str, &str); 5] = [
("a", "/a.xml"),
("b", "/b.xml"),
("c", "/c.xml"),
("d", "/d.xml"),
("e", "/e.xml"),
];
#[tokio::test]
async fn a_long_in_list_prunes_to_its_members() {
let news = TestNews::start(&FIVE_FEEDS, |_| {}).await;
let items = news
.sql("SELECT feed FROM news.main.items WHERE feed IN ('a', 'b', 'c', 'd') ORDER BY feed")
.await;
assert_eq!(col(&items, "feed"), vec!["a", "b", "c", "d"]);
assert_eq!(
sorted(news.paths()),
vec!["/a.xml", "/b.xml", "/c.xml", "/d.xml"],
"the feed the IN list omits is never fetched"
);
}
#[tokio::test]
async fn a_short_in_list_is_rewritten_to_a_disjunction_and_still_prunes() {
let news = TestNews::start(&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")], |_| {}).await;
let items = news
.sql("SELECT feed FROM news.main.items WHERE feed IN ('a', 'c') ORDER BY feed")
.await;
assert_eq!(col(&items, "feed"), vec!["a", "c"]);
assert_eq!(
sorted(news.paths()),
vec!["/a.xml", "/c.xml"],
"the rewritten short IN list prunes to its members, so `b` is never fetched"
);
}
#[tokio::test]
async fn a_hand_written_disjunction_of_feed_equalities_prunes() {
let news = TestNews::start(&FIVE_FEEDS, |_| {}).await;
let items = news
.sql(
"SELECT feed FROM news.main.items \
WHERE feed = 'a' OR feed = 'c' ORDER BY feed",
)
.await;
assert_eq!(col(&items, "feed"), vec!["a", "c"]);
assert_eq!(
sorted(news.paths()),
vec!["/a.xml", "/c.xml"],
"a hand-written OR over one feed column prunes to the union of its leaves"
);
}
#[tokio::test]
async fn a_disjunction_mixing_feed_and_feed_url_does_not_prune() {
let news = TestNews::start(&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")], |_| {}).await;
let c_url = format!("{}/c.xml", news.server.url());
let items = news
.sql(&format!(
"SELECT feed FROM news.main.items \
WHERE feed = 'a' OR feed_url = '{c_url}' ORDER BY feed"
))
.await;
assert_eq!(
col(&items, "feed"),
vec!["a", "c"],
"the predicate is still applied — above the scan, by DataFusion"
);
assert_eq!(
sorted(news.paths()),
vec!["/a.xml", "/b.xml", "/c.xml"],
"a disjunction over two feed columns prunes nothing, so every subscription is fetched"
);
}
#[tokio::test]
async fn conjunction_of_feed_predicates_prunes_to_the_intersection() {
let news = TestNews::start(&FIVE_FEEDS, |config| config.ttl_seconds = 0).await;
let in_then_eq = news
.sql("SELECT feed FROM news.main.items WHERE feed IN ('a','b','c','d') AND feed = 'b'")
.await;
assert_eq!(col(&in_then_eq, "feed"), vec!["b"]);
assert_eq!(
news.paths(),
vec!["/b.xml"],
"the intersection of the two predicates is one feed, so one fetch"
);
let eq_then_in = news
.sql("SELECT feed FROM news.main.items WHERE feed = 'b' AND feed IN ('a','b','c','d')")
.await;
assert_eq!(col(&eq_then_in, "feed"), vec!["b"]);
assert_eq!(
news.paths()[1..],
["/b.xml".to_string()],
"the intersection does not depend on which predicate came first"
);
}
#[tokio::test]
async fn pruning_by_a_shared_feed_url_visits_every_subscription_using_it() {
let news = TestNews::start(
&[
("primary", "/shared.xml"),
("mirror", "/shared.xml"),
("other", "/other.xml"),
],
|_| {},
)
.await;
let shared = format!("{}/shared.xml", news.server.url());
let items = news
.sql(&format!(
"SELECT feed, feed_url FROM news.main.items WHERE feed_url = '{shared}' ORDER BY feed"
))
.await;
assert_eq!(
col(&items, "feed"),
vec!["mirror", "primary"],
"both subscriptions on the shared URL must contribute rows"
);
assert_eq!(col(&items, "feed_url"), vec![shared.clone(), shared]);
assert_eq!(
news.paths(),
vec!["/shared.xml", "/shared.xml"],
"one fetch per surviving subscription, and the unrelated feed is pruned away"
);
}
#[tokio::test]
async fn user_agent_is_sent() {
let news = TestNews::start(&[("a", "/a.xml"), ("b", "/b.xml")], |_| {}).await;
news.sql("SELECT guid FROM news.main.items").await;
let expected = format!(
"skardi-rss/{} (+https://github.com/SkardiLabs/skardi)",
env!("CARGO_PKG_VERSION")
);
let requests = news.requests();
assert_eq!(requests.len(), 2, "both feeds were fetched");
for request in &requests {
assert_eq!(
request.header("user-agent").as_deref(),
Some(expected.as_str()),
"request to {} sent the wrong User-Agent",
request.path
);
}
}
#[tokio::test]
async fn count_star_over_items_is_row_accurate() {
let news = TestNews::start(&[("a", "/a.xml"), ("b", "/b.xml")], |_| {}).await;
news.script(
"/a.xml",
FeedScript::always(Canned::xml(&rss_with(&["a1", "a2", "a3"]))),
);
let counted = news.sql("SELECT count(*) AS n FROM news.main.items").await;
let n = counted[0]
.column_by_name("n")
.expect("count(*) column")
.as_any()
.downcast_ref::<arrow::array::Int64Array>()
.expect("count(*) is Int64")
.value(0);
assert_eq!(n, 4, "three items from feed a plus one from feed b");
let listed = news.sql("SELECT guid FROM news.main.items").await;
assert_eq!(
total_rows(&listed) as i64,
n,
"the empty projection and the ordinary one disagree about the row count"
);
assert_eq!(
sorted(col(&listed, "guid")),
vec!["/b.xml#1", "a1", "a2", "a3"]
);
}
#[tokio::test]
async fn absence_check_pattern_works() {
let news = TestNews::start(&[("alive", "/alive.xml"), ("dead", "/dead.xml")], |_| {}).await;
news.script("/dead.xml", FeedScript::always(Canned::status(404)));
let warm = news.sql("SELECT feed FROM news.main.items").await;
assert_eq!(col(&warm, "feed"), vec!["alive"]);
let warmed_requests = news.request_count();
let missing = news
.sql(
"SELECT f.name, f.last_status FROM news.main.feeds f \
LEFT JOIN news.main.items i ON i.feed = f.name \
WHERE i.feed IS NULL ORDER BY f.name",
)
.await;
assert_eq!(
col(&missing, "name"),
vec!["dead"],
"the subscription with no items is the one the anti-join returns"
);
assert_eq!(
col(&missing, "last_status"),
vec!["error"],
"and the health table says why it has none"
);
assert_eq!(
news.request_count(),
warmed_requests,
"the anti-join's scans were served from cache, so the state it read is the warmed one"
);
}