diff --git a/.claude/skills/ipx-prod-check/SKILL.md b/.claude/skills/ipx-prod-check/SKILL.md index 8aa9897..ab7cb8b 100644 --- a/.claude/skills/ipx-prod-check/SKILL.md +++ b/.claude/skills/ipx-prod-check/SKILL.md @@ -85,6 +85,12 @@ a count says something happened, the lines and traces say why. crash or an OOM kill: check `docker inspect iPX -f '{{.State.OOMKilled}} {{.RestartCount}}'` and the lines just before it. A pending count that only grows means downloads are not keeping up or not running. + + Since 2026-10-02 the daemon sleeps until the next feed is due (at most 10 minutes) instead of + scanning every minute (#114), so expect tens of scans in 6 hours, not 360, nearly all with + `feeds` above 0; none at all for over 10 minutes means the worker is stuck. And `pending` is + the real queue (#113): files a scan will download on its own. Back-catalogue files are + `held`, listed but not counted, so it is usually 0 or a handful. 5. **Slow and failed traces.** ``` $Q traces '{resource.deployment.environment.name="production" && duration > 5s}' diff --git a/CHANGELOG.md b/CHANGELOG.md index 6200acf..b855e49 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 +- 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. A feed nobody subscribes to is checked once a day, and at once when someone subscribes. A diff --git a/src/db.rs b/src/db.rs index 1f9ebce..20a24c4 100644 --- a/src/db.rs +++ b/src/db.rs @@ -438,19 +438,27 @@ pub struct HttpState { pub error_since: Option, } +impl From for HttpState { + fn from(f: feeds::Model) -> Self { + HttpState { + etag: f.etag, + last_modified: f.last_modified, + last_checked: f.last_checked, + ttl_mins: f.ttl_mins.map(|t| t.max(0) as u64), + error_since: f.error_since, + } + } +} + impl Db { pub async fn http_state(&self, feed_id: &str) -> Result { - Ok(feeds::Entity::find_by_id(feed_id.to_owned()) - .one(&self.orm) - .await? - .map(|f| HttpState { - etag: f.etag, - last_modified: f.last_modified, - last_checked: f.last_checked, - ttl_mins: f.ttl_mins.map(|t| t.max(0) as u64), - error_since: f.error_since, - }) - .unwrap_or_default()) + Ok(feeds::Entity::find_by_id(feed_id.to_owned()).one(&self.orm).await?.map(HttpState::from).unwrap_or_default()) + } + + /// Every feed's, in one query: a scan, and the daemon working out when the next is due, asked + /// feed by feed, 180 queries a minute (#114). + pub async fn http_states(&self) -> Result> { + Ok(feeds::Entity::find().all(&self.orm).await?.into_iter().map(|f| (f.id.clone(), HttpState::from(f))).collect()) } /// Upsert after a successful poll. Clears any previous error. diff --git a/src/main.rs b/src/main.rs index a5de343..6134f18 100644 --- a/src/main.rs +++ b/src/main.rs @@ -514,8 +514,6 @@ async fn daemon( let server = tokio::spawn(ipc::serve(socket.clone(), events.clone(), tx_cmd, answer)); // One command at a time: the queue is what keeps two scans from overlapping. - let mut ticker = tokio::time::interval(std::time::Duration::from_secs(60)); - ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); tracing::info!( feeds = subscriptions(&ctx).await.map(|s| s.len()).unwrap_or(0), "daemon started" @@ -552,10 +550,13 @@ async fn daemon( } let mut stop = rx_stop.clone(); + // The first pass at once, as the minute's tick did: what came due while it was down. + let mut first = true; loop { if *stop.borrow() { break; } + let wait = if std::mem::take(&mut first) { std::time::Duration::ZERO } else { until_next_scan(&ctx).await }; tokio::select! { biased; _ = stop.changed() => break, @@ -572,7 +573,7 @@ async fn daemon( break; } } - _ = ticker.tick() => { + _ = tokio::time::sleep(wait) => { // Per-feed schedule and TTL decide what actually gets polled. let job = run(&ctx, Cmd::Fetch { feed: None, force: false, feeds: vec![] }); if !until_stopped(&ctx, &rx_stop, job).await { @@ -1045,10 +1046,11 @@ async fn fetch(ctx: &Arc, only: Option<&str>, force: bool, scope: &[String] let in_scope = |s: &Sub| { scope.is_empty() || scope.contains(&s.id) || s.cfg.group.as_ref().is_some_and(|g| scope.contains(g)) }; + let mut states = ctx.db.http_states().await?; let mut due = vec![]; for sub in subs.iter().filter(|s| only.is_none_or(|o| o == s.id) && in_scope(s)) { let id = &sub.id; - let mut state = ctx.db.http_state(id).await?; + let mut state = states.remove(id).unwrap_or_default(); // A scan someone asked for reads the feed in full. With the validators it only skipped the // wait: a feed that had not changed answered 304 and nothing was read (#77). if force { @@ -1056,12 +1058,7 @@ async fn fetch(ctx: &Arc, only: Option<&str>, force: bool, scope: &[String] state.last_modified = None; } - if !force && let Some(last) = state.last_checked { - // A feed nobody subscribes to, one listed in the Directory, is read once a day: enough - // to keep its entry current, without fetching it hourly for no one. - let floor = if subscribed.contains_key(id) { 0 } else { 86_400 }; - let wait = due_after(&cfg, &sub.cfg, state.ttl_mins, state.error_since, last).max(floor); - let at = last + wait as i64; + if !force && let Some(at) = due_at(&cfg, &sub.cfg, &state, subscribed.contains_key(id)) { if at > db::now() { ctx.out.emit(Event::FeedSkip { feed: id.clone(), @@ -1272,6 +1269,47 @@ async fn retire_stranded(ctx: &Ctx) -> Result { /// Seconds to wait before re-checking a feed. /// +/// When a feed is next due, or None for one never checked, which is due now. A feed nobody +/// subscribes to, one listed in the Directory, is read once a day: enough to keep its entry +/// current, without fetching it hourly for no one. +pub fn due_at(cfg: &config::Config, feed: &config::Feed, state: &db::HttpState, subscribed: bool) -> Option { + let last = state.last_checked?; + let floor = if subscribed { 0 } else { 86_400 }; + Some(last + due_after(cfg, feed, state.ttl_mins, state.error_since, last).max(floor) as i64) +} + +/// How long the daemon may sleep before a feed is due: until the earliest one, at least 30 s and +/// at most 10 minutes. It ticked every minute and ran a scan pass each time, 80% of them finding +/// nothing due (#114). The floor keeps a feed that never gets a check time from spinning it; the +/// ceiling picks up within ten minutes what no command announces, such as `ipx add` or a +/// shorter schedule. A command, a refresh or a feed added on the page, wakes it at once anyway. +async fn until_next_scan(ctx: &Ctx) -> std::time::Duration { + let next = async { + let cfg = ctx.cfg(); + let subscribed = ctx.db.subscriber_counts().await?; + let states = ctx.db.http_states().await?; + anyhow::Ok( + subscriptions(ctx) + .await? + .iter() + // A feed with no row yet has never been checked: due now. + .map(|s| { + states.get(&s.id).and_then(|st| due_at(&cfg, &s.cfg, st, subscribed.contains_key(&s.id))).unwrap_or(0) + }) + .min(), + ) + }; + let wait = match next.await { + Ok(Some(at)) => (at - db::now()).max(0) as u64, + Ok(None) => u64::MAX, + Err(e) => { + tracing::warn!(error = %format!("{e:#}"), "could not work out when the next feed is due"); + 0 + } + }; + std::time::Duration::from_secs(wait.clamp(30, 600)) +} + /// A per-feed schedule is an explicit instruction and wins outright. Without one, the /// global schedule applies, but the feed's own raises it when the publisher asks to /// be polled less often. A feed that is failing backs off (`backoff`), from `error_since`, when @@ -1999,6 +2037,18 @@ fn duration(secs: u64) -> String { mod tests { use super::*; + #[test] + fn a_feed_is_due_after_its_wait_and_daily_with_nobody_subscribed() { + let cfg = config::Config::default(); + let f = feed(); + let at = |last: Option, subscribed| { + due_at(&cfg, &f, &db::HttpState { last_checked: last, ..Default::default() }, subscribed) + }; + assert_eq!(at(None, true), None); // never checked: due now + assert_eq!(at(Some(1000), true), Some(1000 + 3600)); // the hourly default + assert_eq!(at(Some(1000), false), Some(1000 + 86_400)); // listed, nobody subscribed + } + #[test] fn a_listed_feed_nobody_subscribes_to_downloads_nothing() { let mut f = feed();