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:
@@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
|
||||
### 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.
|
||||
- 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.
|
||||
|
||||
129
src/feed.rs
129
src/feed.rs
@@ -59,43 +59,66 @@ pub enum Fetched {
|
||||
/// downloads, which are not, have no such limit.
|
||||
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
|
||||
/// CONNECT/socket.ssl proxy path is gone -- `system-proxy` reads http_proxy/https_proxy.
|
||||
/// Conditional GET, following redirects itself: `client` must follow none (`feed_client`), so
|
||||
/// 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))]
|
||||
pub async fn fetch(
|
||||
client: &reqwest::Client,
|
||||
cfg: &FeedCfg,
|
||||
etag: Option<&str>,
|
||||
last_modified: Option<&str>,
|
||||
) -> Result<Fetched> {
|
||||
let mut req = client.get(&cfg.url).timeout(FEED_TIMEOUT);
|
||||
if let Some(tag) = etag {
|
||||
req = req.header(IF_NONE_MATCH, tag);
|
||||
}
|
||||
if let Some(lm) = last_modified {
|
||||
req = req.header(IF_MODIFIED_SINCE, lm);
|
||||
}
|
||||
if let Some(user) = &cfg.username {
|
||||
req = req.basic_auth(user, cfg.password());
|
||||
}
|
||||
) -> Result<(Fetched, Option<String>)> {
|
||||
let start = reqwest::Url::parse(&cfg.url).context("the feed's address")?;
|
||||
let mut url = start.clone();
|
||||
let mut permanent = true;
|
||||
for _ in 0..10 {
|
||||
let mut req = client.get(url.clone()).timeout(FEED_TIMEOUT);
|
||||
if let Some(tag) = etag {
|
||||
req = req.header(IF_NONE_MATCH, tag);
|
||||
}
|
||||
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")?;
|
||||
if resp.status() == StatusCode::NOT_MODIFIED {
|
||||
return Ok(Fetched::NotModified);
|
||||
let resp = req.send().await.context("connecting")?;
|
||||
let status = resp.status();
|
||||
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();
|
||||
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();
|
||||
Ok(Fetched::Body { bytes, etag, last_modified })
|
||||
// Worded as reqwest worded it, which failure_kind reads as a redirect loop.
|
||||
Err(anyhow!("error following redirect for url ({url}): too many redirects"))
|
||||
}
|
||||
|
||||
/// 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)]
|
||||
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::*;
|
||||
|
||||
#[test]
|
||||
|
||||
43
src/main.rs
43
src/main.rs
@@ -125,6 +125,9 @@ pub struct Ctx {
|
||||
pub cfg: std::sync::RwLock<std::sync::Arc<config::Config>>,
|
||||
pub db: db::Db,
|
||||
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,
|
||||
/// 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.
|
||||
@@ -258,6 +261,11 @@ async fn main() -> Result<()> {
|
||||
// Connecting only, so it bounds a download's start, not a long download (#108).
|
||||
.connect_timeout(std::time::Duration::from_secs(10))
|
||||
.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() },
|
||||
torrents: tokio::sync::OnceCell::new(),
|
||||
torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)),
|
||||
@@ -741,7 +749,7 @@ pub async fn add_one(
|
||||
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
|
||||
// parse it as a feed and falling back to the hostname.
|
||||
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;
|
||||
let mut bodies = futures_util::StreamExt::buffered(
|
||||
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 early = !feed::is_patreon_creator(&feed_cfg.url);
|
||||
tokio::spawn(tracing::Instrument::instrument(
|
||||
@@ -1361,6 +1369,23 @@ enum Outcome {
|
||||
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))]
|
||||
async fn scan_one(
|
||||
ctx: &Arc<Ctx>,
|
||||
@@ -1368,8 +1393,8 @@ async fn scan_one(
|
||||
feed_cfg: &config::Feed,
|
||||
state: &db::HttpState,
|
||||
force: bool,
|
||||
// The body the scan already fetched for it, if it did (#104).
|
||||
prefetched: Option<Result<feed::Fetched>>,
|
||||
// The body the scan already fetched for it, if it did (#104), and where it moved, if it did.
|
||||
prefetched: Option<Result<(feed::Fetched, Option<String>)>>,
|
||||
) -> Result<Outcome> {
|
||||
// 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).
|
||||
@@ -1400,10 +1425,13 @@ async fn scan_one(
|
||||
}
|
||||
}
|
||||
|
||||
let mut fetched = match prefetched {
|
||||
let (mut fetched, moved) = match prefetched {
|
||||
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
|
||||
// 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()) {
|
||||
tracing::info!(feed = id, "not modified, but nothing stored; refetching without the validator");
|
||||
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 {
|
||||
@@ -2187,6 +2215,7 @@ mod tests {
|
||||
cfg: std::sync::RwLock::new(Arc::new(cfg)),
|
||||
db: db::Db::memory().await.unwrap(),
|
||||
client: reqwest::Client::new(),
|
||||
feed_client: reqwest::Client::builder().redirect(reqwest::redirect::Policy::none()).build().unwrap(),
|
||||
out: Emitter::terminal(),
|
||||
torrents: tokio::sync::OnceCell::new(),
|
||||
torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)),
|
||||
|
||||
Reference in New Issue
Block a user