From ec8fd5dd8676d2d89d4dc01f0853d3e67f3b7322 Mon Sep 17 00:00:00 2001 From: rays Date: Fri, 2 Oct 2026 19:56:13 +0000 Subject: [PATCH] Sleep until the next feed is due instead of scanning every minute (#114) The daemon ticked every 60 s and ran a scan pass each time: the sweep, then a check-state query per feed (about 180) to find which were due. In the six hours before, 293 of 362 passes found nothing due. Now, after each pass, it works out when the earliest feed is due (due_at, shared with the scan's own check, over Db::http_states, one query) and sleeps until then: at least 30 s, so a feed that never gets a check time cannot spin it, and at most 10 minutes, so what no command announces, ipx add or a shorter schedule, is picked up. Commands still wake it at once, and the first pass after starting runs straight away, as the tick's did. The scan reads every feed's state in one query too. The prod-check skill says what to expect now: tens of scans in six hours, and pending as the real queue. Co-Authored-By: Claude Opus 5.5 --- .claude/skills/ipx-prod-check/SKILL.md | 6 +++ CHANGELOG.md | 1 + src/db.rs | 30 +++++++---- src/main.rs | 70 ++++++++++++++++++++++---- 4 files changed, 86 insertions(+), 21 deletions(-) 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();