Follow a feed that has moved for good to its new address (#115)

A feed whose address answered with a permanent redirect was read through it on every check,
and the catalogue kept the old address: 28 of 147 feeds in production, most http to https, some
to a new path or domain.

Feeds are now fetched with a client of their own that follows no redirects (Ctx::feed_client),
and feed::fetch follows them itself, up to 10 hops, so it sees each one. When every hop was
permanent (301 or 308) it says where the feed ended up, and the scan moves the feed there in the
catalogue (follow_move). A temporary hop (302, 307) anywhere moves nothing. A feed an OPML lists
is left alone, as the OPML would put the old address back, and so is a move onto an address
another feed has. A password goes only to the feed's own host, never to a redirect elsewhere;
reqwest's own following dropped it the same way. Ten hops is a loop, worded as reqwest worded
it so it still reads as redirect_loop.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-10-02 20:05:43 +00:00
parent ec8fd5dd86
commit f7f4b466dc
3 changed files with 137 additions and 36 deletions

View File

@@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed ### Changed
- A feed that has moved for good (a permanent redirect) is followed to its new address, which iPX then reads from. A temporary redirect changes nothing.
- The daemon sleeps until the next feed is due, at most ten minutes, instead of looking every minute; refreshing or adding a feed still wakes it at once. - The daemon sleeps until the next feed is due, at most ten minutes, instead of looking every minute; refreshing or adding a feed still wakes it at once.
- A pinned feed's pin sits on the corner of its artwork, as a failing feed's mark does, instead of before its name. - A pinned feed's pin sits on the corner of its artwork, as a failing feed's mark does, instead of before its name.
- The Directory lists feeds nobody subscribes to yet; Popular still lists what people subscribe to. - The Directory lists feeds nobody subscribes to yet; Popular still lists what people subscribe to.

View File

@@ -59,43 +59,66 @@ pub enum Fetched {
/// downloads, which are not, have no such limit. /// downloads, which are not, have no such limit.
pub const FEED_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); pub const FEED_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
/// Conditional GET. reqwest handles gzip and redirects; the original's hand-rolled /// Conditional GET, following redirects itself: `client` must follow none (`feed_client`), so
/// CONNECT/socket.ssl proxy path is gone -- `system-proxy` reads http_proxy/https_proxy. /// each hop is seen. The second value is where the feed now is when every hop said so for good
/// (301 or 308): a publisher that moved its feed, which the catalogue should follow rather
/// than be redirected on every read. A temporary redirect (302, 307) moves nothing.
/// `system-proxy` reads http_proxy/https_proxy.
#[tracing::instrument(skip_all, fields(url = %cfg.url))] #[tracing::instrument(skip_all, fields(url = %cfg.url))]
pub async fn fetch( pub async fn fetch(
client: &reqwest::Client, client: &reqwest::Client,
cfg: &FeedCfg, cfg: &FeedCfg,
etag: Option<&str>, etag: Option<&str>,
last_modified: Option<&str>, last_modified: Option<&str>,
) -> Result<Fetched> { ) -> Result<(Fetched, Option<String>)> {
let mut req = client.get(&cfg.url).timeout(FEED_TIMEOUT); let start = reqwest::Url::parse(&cfg.url).context("the feed's address")?;
if let Some(tag) = etag { let mut url = start.clone();
req = req.header(IF_NONE_MATCH, tag); let mut permanent = true;
} for _ in 0..10 {
if let Some(lm) = last_modified { let mut req = client.get(url.clone()).timeout(FEED_TIMEOUT);
req = req.header(IF_MODIFIED_SINCE, lm); if let Some(tag) = etag {
} req = req.header(IF_NONE_MATCH, tag);
if let Some(user) = &cfg.username { }
req = req.basic_auth(user, cfg.password()); if let Some(lm) = last_modified {
} req = req.header(IF_MODIFIED_SINCE, lm);
}
// The feed's own host only: a redirect elsewhere must not be handed the password.
if let Some(user) = &cfg.username
&& url.host_str() == start.host_str()
{
req = req.basic_auth(user, cfg.password());
}
let resp = req.send().await.context("connecting")?; let resp = req.send().await.context("connecting")?;
if resp.status() == StatusCode::NOT_MODIFIED { let status = resp.status();
return Ok(Fetched::NotModified); if status.is_redirection() && status != StatusCode::NOT_MODIFIED {
let to = resp
.headers()
.get(reqwest::header::LOCATION)
.and_then(|v| v.to_str().ok())
.ok_or_else(|| anyhow!("HTTP {status} without a Location to go to"))?;
url = url.join(to).with_context(|| format!("redirected to {to:?}, which is not an address"))?;
permanent &= matches!(status, StatusCode::MOVED_PERMANENTLY | StatusCode::PERMANENT_REDIRECT);
continue;
}
let moved = (permanent && url != start).then(|| url.to_string());
if status == StatusCode::NOT_MODIFIED {
return Ok((Fetched::NotModified, moved));
}
if !status.is_success() {
// The original surfaced 401/407 specially; the code is enough for a UI to switch on.
return Err(anyhow!("HTTP {status}"));
}
let header = |h: reqwest::header::HeaderName| {
resp.headers().get(&h).and_then(|v| v.to_str().ok()).map(str::to_owned)
};
let etag = header(ETAG);
let last_modified = header(LAST_MODIFIED);
let bytes = resp.bytes().await.context("reading body")?.to_vec();
return Ok((Fetched::Body { bytes, etag, last_modified }, moved));
} }
let status = resp.status(); // Worded as reqwest worded it, which failure_kind reads as a redirect loop.
if !status.is_success() { Err(anyhow!("error following redirect for url ({url}): too many redirects"))
// The original surfaced 401/407 specially; the code is enough for a UI to switch on.
return Err(anyhow!("HTTP {status}"));
}
let header = |h: reqwest::header::HeaderName| {
resp.headers().get(&h).and_then(|v| v.to_str().ok()).map(str::to_owned)
};
let etag = header(ETAG);
let last_modified = header(LAST_MODIFIED);
let bytes = resp.bytes().await.context("reading body")?.to_vec();
Ok(Fetched::Body { bytes, etag, last_modified })
} }
/// A stored `last_error`, translated into plain words for whoever subscribes: whose problem /// A stored `last_error`, translated into plain words for whoever subscribes: whose problem
@@ -987,6 +1010,54 @@ fn parse_date(s: &str) -> Option<i64> {
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
/// A server answering by path: /old moves for good to /new, /tmp for now, /chain for good
/// to /tmp, /new is the feed.
async fn redirecting_server() -> String {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
loop {
let Ok((mut sock, _)) = listener.accept().await else { return };
tokio::spawn(async move {
let mut buf = [0u8; 2048];
let n = sock.read(&mut buf).await.unwrap_or(0);
let req = String::from_utf8_lossy(&buf[..n]);
let path = req.split_whitespace().nth(1).unwrap_or("/").to_owned();
let body = "<?xml version=\"1.0\"?><rss version=\"2.0\"><channel><title>T</title></channel></rss>";
let resp = match path.as_str() {
"/old" => "HTTP/1.1 301 Moved Permanently\r\nLocation: /new\r\nContent-Length: 0\r\n\r\n".to_owned(),
"/tmp" => "HTTP/1.1 302 Found\r\nLocation: /new\r\nContent-Length: 0\r\n\r\n".to_owned(),
"/chain" => "HTTP/1.1 308 Permanent Redirect\r\nLocation: /tmp\r\nContent-Length: 0\r\n\r\n".to_owned(),
"/loop" => "HTTP/1.1 301 Moved Permanently\r\nLocation: /loop\r\nContent-Length: 0\r\n\r\n".to_owned(),
_ => format!("HTTP/1.1 200 OK\r\nContent-Length: {}\r\n\r\n{body}", body.len()),
};
let _ = sock.write_all(resp.as_bytes()).await;
});
}
});
format!("http://{addr}")
}
#[tokio::test]
async fn a_feed_that_moved_for_good_says_where_and_one_moved_for_now_does_not() {
let base = redirecting_server().await;
let client = reqwest::Client::builder().redirect(reqwest::redirect::Policy::none()).build().unwrap();
let get = |path: &str| {
let cfg: crate::config::Feed = serde_json::from_value(serde_json::json!({ "url": format!("{base}{path}") })).unwrap();
let client = client.clone();
async move { super::fetch(&client, &cfg, None, None).await }
};
let (got, moved) = get("/old").await.unwrap();
assert!(matches!(got, super::Fetched::Body { .. }));
assert_eq!(moved, Some(format!("{base}/new")));
assert_eq!(get("/tmp").await.unwrap().1, None); // 302: for now
assert_eq!(get("/chain").await.unwrap().1, None); // 308 then 302: not for good
assert_eq!(get("/new").await.unwrap().1, None); // never moved
let looped = get("/loop").await.err().unwrap().to_string();
assert_eq!(super::failure_kind(&looped).0, "redirect_loop");
}
use super::*; use super::*;
#[test] #[test]

View File

@@ -125,6 +125,9 @@ pub struct Ctx {
pub cfg: std::sync::RwLock<std::sync::Arc<config::Config>>, pub cfg: std::sync::RwLock<std::sync::Arc<config::Config>>,
pub db: db::Db, pub db: db::Db,
pub client: reqwest::Client, pub client: reqwest::Client,
/// For feeds only: follows no redirects, so `feed::fetch` sees each hop and can tell a feed
/// that moved for good from one sent elsewhere for now.
pub feed_client: reqwest::Client,
pub out: Emitter, pub out: Emitter,
/// Started on first use: a BitTorrent session binds ports and starts a DHT, which is /// Started on first use: a BitTorrent session binds ports and starts a DHT, which is
/// rude to do for a config that has never seen a torrent. /// rude to do for a config that has never seen a torrent.
@@ -258,6 +261,11 @@ async fn main() -> Result<()> {
// Connecting only, so it bounds a download's start, not a long download (#108). // Connecting only, so it bounds a download's start, not a long download (#108).
.connect_timeout(std::time::Duration::from_secs(10)) .connect_timeout(std::time::Duration::from_secs(10))
.build()?, .build()?,
feed_client: reqwest::Client::builder()
.user_agent(concat!("ipx/", env!("CARGO_PKG_VERSION")))
.connect_timeout(std::time::Duration::from_secs(10))
.redirect(reqwest::redirect::Policy::none())
.build()?,
out: if is_daemon { Emitter::socket(events.clone(), false) } else { Emitter::terminal() }, out: if is_daemon { Emitter::socket(events.clone(), false) } else { Emitter::terminal() },
torrents: tokio::sync::OnceCell::new(), torrents: tokio::sync::OnceCell::new(),
torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)), torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)),
@@ -741,7 +749,7 @@ pub async fn add_one(
listed: false, listed: false,
}; };
let title = match feed::fetch(&ctx.client, &probe, None, None).await { let title = match feed::fetch(&ctx.feed_client, &probe, None, None).await.map(|(got, _)| got) {
// An OPML subscription is named from its own <head><title>, not by trying to // An OPML subscription is named from its own <head><title>, not by trying to
// parse it as a feed and falling back to the hostname. // parse it as a feed and falling back to the hostname.
Ok(feed::Fetched::Body { bytes, .. }) if feed::is_opml(&bytes) => { Ok(feed::Fetched::Body { bytes, .. }) if feed::is_opml(&bytes) => {
@@ -1077,7 +1085,7 @@ async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]
const AHEAD: usize = 6; const AHEAD: usize = 6;
let mut bodies = futures_util::StreamExt::buffered( let mut bodies = futures_util::StreamExt::buffered(
futures_util::stream::iter(due.iter().map(|(sub, state)| { futures_util::stream::iter(due.iter().map(|(sub, state)| {
let (client, feed_cfg) = (ctx.client.clone(), sub.cfg.clone()); let (client, feed_cfg) = (ctx.feed_client.clone(), sub.cfg.clone());
let (etag, modified) = (state.etag.clone(), state.last_modified.clone()); let (etag, modified) = (state.etag.clone(), state.last_modified.clone());
let early = !feed::is_patreon_creator(&feed_cfg.url); let early = !feed::is_patreon_creator(&feed_cfg.url);
tokio::spawn(tracing::Instrument::instrument( tokio::spawn(tracing::Instrument::instrument(
@@ -1361,6 +1369,23 @@ enum Outcome {
Opml { added: Vec<String>, removed: usize, kept: usize, total: usize }, Opml { added: Vec<String>, removed: usize, kept: usize, total: usize },
} }
/// A feed that answered from a new address after permanent redirects moves there in the
/// catalogue, so it is read from there and no longer redirected every time. One an OPML lists
/// is left alone, since the OPML would only put the old address back; so is a move onto an
/// address another feed already has.
async fn follow_move(ctx: &Ctx, id: &str, to: &str) -> Result<()> {
if subscriptions(ctx).await?.iter().any(|s| s.id != id && feed::same_feed(&s.cfg.url, to)) {
tracing::info!(feed = id, to, "the feed moved to an address another feed already has; leaving it");
return Ok(());
}
let mut cfg = (*ctx.cfg()).clone();
let Some(f) = cfg.feeds.get_mut(id) else { return Ok(()) };
let from = std::mem::replace(&mut f.url, to.to_owned());
ctx.store_cfg(cfg).await?;
tracing::info!(feed = id, from, to, "the feed moved for good; following it to its new address");
Ok(())
}
#[tracing::instrument(name = "feed", skip_all, fields(feed = id))] #[tracing::instrument(name = "feed", skip_all, fields(feed = id))]
async fn scan_one( async fn scan_one(
ctx: &Arc<Ctx>, ctx: &Arc<Ctx>,
@@ -1368,8 +1393,8 @@ async fn scan_one(
feed_cfg: &config::Feed, feed_cfg: &config::Feed,
state: &db::HttpState, state: &db::HttpState,
force: bool, force: bool,
// The body the scan already fetched for it, if it did (#104). // The body the scan already fetched for it, if it did (#104), and where it moved, if it did.
prefetched: Option<Result<feed::Fetched>>, prefetched: Option<Result<(feed::Fetched, Option<String>)>>,
) -> Result<Outcome> { ) -> Result<Outcome> {
// The queue follows the feed's settings each time it is due, changed or not, so files a // The queue follows the feed's settings each time it is due, changed or not, so files a
// limit or auto-download no longer reaches stop counting as waiting (#113). // limit or auto-download no longer reaches stop counting as waiting (#113).
@@ -1400,10 +1425,13 @@ async fn scan_one(
} }
} }
let mut fetched = match prefetched { let (mut fetched, moved) = match prefetched {
Some(got) => got?, Some(got) => got?,
None => feed::fetch(&ctx.client, feed_cfg, state.etag.as_deref(), state.last_modified.as_deref()).await?, None => feed::fetch(&ctx.feed_client, feed_cfg, state.etag.as_deref(), state.last_modified.as_deref()).await?,
}; };
if let Some(to) = moved {
follow_move(ctx, id, &to).await?;
}
// A 304 while nothing is stored means the validator has outlived the data -- a restore // A 304 while nothing is stored means the validator has outlived the data -- a restore
// from backup, a manual edit, a cleanup that removed entries. Believe the database over // from backup, a manual edit, a cleanup that removed entries. Believe the database over
@@ -1414,7 +1442,7 @@ async fn scan_one(
if matches!(fetched, feed::Fetched::NotModified) && (stored.entries == 0 || stored.image.is_none()) { if matches!(fetched, feed::Fetched::NotModified) && (stored.entries == 0 || stored.image.is_none()) {
tracing::info!(feed = id, "not modified, but nothing stored; refetching without the validator"); tracing::info!(feed = id, "not modified, but nothing stored; refetching without the validator");
ctx.db.clear_validators(id).await?; ctx.db.clear_validators(id).await?;
fetched = feed::fetch(&ctx.client, feed_cfg, None, None).await?; fetched = feed::fetch(&ctx.feed_client, feed_cfg, None, None).await?.0;
} }
let (bytes, etag, last_modified) = match fetched { let (bytes, etag, last_modified) = match fetched {
@@ -2187,6 +2215,7 @@ mod tests {
cfg: std::sync::RwLock::new(Arc::new(cfg)), cfg: std::sync::RwLock::new(Arc::new(cfg)),
db: db::Db::memory().await.unwrap(), db: db::Db::memory().await.unwrap(),
client: reqwest::Client::new(), client: reqwest::Client::new(),
feed_client: reqwest::Client::builder().redirect(reqwest::redirect::Policy::none()).build().unwrap(),
out: Emitter::terminal(), out: Emitter::terminal(),
torrents: tokio::sync::OnceCell::new(), torrents: tokio::sync::OnceCell::new(),
torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)), torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)),