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 <noreply@anthropic.com>
This commit is contained in:
2026-10-02 19:56:13 +00:00
parent a387012c69
commit ec8fd5dd86
4 changed files with 86 additions and 21 deletions

View File

@@ -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}}'` 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 and the lines just before it. A pending count that only grows means downloads are not
keeping up or not running. 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.** 5. **Slow and failed traces.**
``` ```
$Q traces '{resource.deployment.environment.name="production" && duration > 5s}' $Q traces '{resource.deployment.environment.name="production" && duration > 5s}'

View File

@@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed ### 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. - 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.
A feed nobody subscribes to is checked once a day, and at once when someone subscribes. A A feed nobody subscribes to is checked once a day, and at once when someone subscribes. A

View File

@@ -438,19 +438,27 @@ pub struct HttpState {
pub error_since: Option<i64>, pub error_since: Option<i64>,
} }
impl Db { impl From<feeds::Model> for HttpState {
pub async fn http_state(&self, feed_id: &str) -> Result<HttpState> { fn from(f: feeds::Model) -> Self {
Ok(feeds::Entity::find_by_id(feed_id.to_owned()) HttpState {
.one(&self.orm)
.await?
.map(|f| HttpState {
etag: f.etag, etag: f.etag,
last_modified: f.last_modified, last_modified: f.last_modified,
last_checked: f.last_checked, last_checked: f.last_checked,
ttl_mins: f.ttl_mins.map(|t| t.max(0) as u64), ttl_mins: f.ttl_mins.map(|t| t.max(0) as u64),
error_since: f.error_since, error_since: f.error_since,
}) }
.unwrap_or_default()) }
}
impl Db {
pub async fn http_state(&self, feed_id: &str) -> Result<HttpState> {
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<std::collections::HashMap<String, HttpState>> {
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. /// Upsert after a successful poll. Clears any previous error.

View File

@@ -514,8 +514,6 @@ async fn daemon(
let server = tokio::spawn(ipc::serve(socket.clone(), events.clone(), tx_cmd, answer)); 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. // 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!( tracing::info!(
feeds = subscriptions(&ctx).await.map(|s| s.len()).unwrap_or(0), feeds = subscriptions(&ctx).await.map(|s| s.len()).unwrap_or(0),
"daemon started" "daemon started"
@@ -552,10 +550,13 @@ async fn daemon(
} }
let mut stop = rx_stop.clone(); 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 { loop {
if *stop.borrow() { if *stop.borrow() {
break; break;
} }
let wait = if std::mem::take(&mut first) { std::time::Duration::ZERO } else { until_next_scan(&ctx).await };
tokio::select! { tokio::select! {
biased; biased;
_ = stop.changed() => break, _ = stop.changed() => break,
@@ -572,7 +573,7 @@ async fn daemon(
break; break;
} }
} }
_ = ticker.tick() => { _ = tokio::time::sleep(wait) => {
// Per-feed schedule and TTL decide what actually gets polled. // Per-feed schedule and TTL decide what actually gets polled.
let job = run(&ctx, Cmd::Fetch { feed: None, force: false, feeds: vec![] }); let job = run(&ctx, Cmd::Fetch { feed: None, force: false, feeds: vec![] });
if !until_stopped(&ctx, &rx_stop, job).await { if !until_stopped(&ctx, &rx_stop, job).await {
@@ -1045,10 +1046,11 @@ async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]
let in_scope = |s: &Sub| { let in_scope = |s: &Sub| {
scope.is_empty() || scope.contains(&s.id) || s.cfg.group.as_ref().is_some_and(|g| scope.contains(g)) 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![]; let mut due = vec![];
for sub in subs.iter().filter(|s| only.is_none_or(|o| o == s.id) && in_scope(s)) { for sub in subs.iter().filter(|s| only.is_none_or(|o| o == s.id) && in_scope(s)) {
let id = &sub.id; 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 // 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). // wait: a feed that had not changed answered 304 and nothing was read (#77).
if force { if force {
@@ -1056,12 +1058,7 @@ async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]
state.last_modified = None; state.last_modified = None;
} }
if !force && let Some(last) = state.last_checked { if !force && let Some(at) = due_at(&cfg, &sub.cfg, &state, subscribed.contains_key(id)) {
// 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 at > db::now() { if at > db::now() {
ctx.out.emit(Event::FeedSkip { ctx.out.emit(Event::FeedSkip {
feed: id.clone(), feed: id.clone(),
@@ -1272,6 +1269,47 @@ async fn retire_stranded(ctx: &Ctx) -> Result<usize> {
/// Seconds to wait before re-checking a feed. /// 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<i64> {
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 /// A per-feed schedule is an explicit instruction and wins outright. Without one, the
/// global schedule applies, but the feed's own <ttl> raises it when the publisher asks to /// global schedule applies, but the feed's own <ttl> raises it when the publisher asks to
/// be polled less often. A feed that is failing backs off (`backoff`), from `error_since`, when /// 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 { mod tests {
use super::*; 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<i64>, 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] #[test]
fn a_listed_feed_nobody_subscribes_to_downloads_nothing() { fn a_listed_feed_nobody_subscribes_to_downloads_nothing() {
let mut f = feed(); let mut f = feed();