diff --git a/CHANGELOG.md b/CHANGELOG.md index b855e49..9684103 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/src/feed.rs b/src/feed.rs index 9771d6e..fd9363d 100644 --- a/src/feed.rs +++ b/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 { - 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)> { + 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 { #[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 = "T"; + 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] diff --git a/src/main.rs b/src/main.rs index 6134f18..fdfc9c7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -125,6 +125,9 @@ pub struct Ctx { pub cfg: std::sync::RwLock>, 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 , 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)),