Fetch feeds six ahead in a scan; reload the list only after a scan that checked something (#103, #104)
Fetching was 65-90% of a scan, each feed waiting for the one before: 13 s of fetches in a 20 s refresh of 32 feeds. The scan now works out which feeds are due, fetches their bodies up to six ahead in tasks of their own, and handles each in order as before, so database writes, downloads and OPML syncs stay one at a time. A Patreon creator still fetches in scan_one. The page reloaded /api/feeds, and /api/settings with it, on every scan_done: the scheduler scans every minute, so each open page reloaded the list once a minute, 169 times an hour. It now reloads only when the scan checked a feed, and asks for settings once, on first load. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
54
src/main.rs
54
src/main.rs
@@ -970,8 +970,9 @@ async fn fetch(ctx: &Arc<Ctx>, 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 due = vec![];
|
||||
for sub in subs.iter().filter(|s| only.is_none_or(|o| o == s.id) && in_scope(s)) {
|
||||
let (id, feed_cfg) = (&sub.id, &sub.cfg);
|
||||
let id = &sub.id;
|
||||
let mut state = ctx.db.http_state(id).await?;
|
||||
// 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).
|
||||
@@ -981,19 +982,47 @@ async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]
|
||||
}
|
||||
|
||||
if !force && let Some(last) = state.last_checked {
|
||||
let due = last + due_after(&cfg, feed_cfg, state.ttl_mins, state.error_since, last) as i64;
|
||||
if due > db::now() {
|
||||
let at = last + due_after(&cfg, &sub.cfg, state.ttl_mins, state.error_since, last) as i64;
|
||||
if at > db::now() {
|
||||
ctx.out.emit(Event::FeedSkip {
|
||||
feed: id.clone(),
|
||||
reason: format!("not due for {}", duration((due - db::now()) as u64)),
|
||||
reason: format!("not due for {}", duration((at - db::now()) as u64)),
|
||||
});
|
||||
continue;
|
||||
}
|
||||
}
|
||||
due.push((sub, state));
|
||||
}
|
||||
|
||||
// Each feed's body is fetched a few feeds ahead of its turn, in tasks of their own, and the
|
||||
// feeds are then handled one at a time, in order, as before: database writes, downloads and
|
||||
// OPML syncs stay one at a time. Fetched one after another, a scan waited on every site in
|
||||
// turn, 65-90% of its time (#104). A Patreon creator fetches inside scan_one, its own way.
|
||||
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 (etag, modified) = (state.etag.clone(), state.last_modified.clone());
|
||||
let early = !feed::is_patreon_creator(&feed_cfg.url);
|
||||
tokio::spawn(tracing::Instrument::instrument(
|
||||
async move {
|
||||
if !early {
|
||||
return None;
|
||||
}
|
||||
Some(feed::fetch(&client, &feed_cfg, etag.as_deref(), modified.as_deref()).await)
|
||||
},
|
||||
tracing::Span::current(),
|
||||
))
|
||||
})),
|
||||
AHEAD,
|
||||
);
|
||||
for (sub, state) in &due {
|
||||
let (id, feed_cfg) = (&sub.id, &sub.cfg);
|
||||
// A task that panicked has no body; scan_one then fetches it itself.
|
||||
let fetched = futures_util::StreamExt::next(&mut bodies).await.and_then(|j| j.ok()).flatten();
|
||||
scanned += 1;
|
||||
ctx.out.emit(Event::FeedStart { feed: id.clone() });
|
||||
match scan_one(ctx, id, feed_cfg, &state, force).await {
|
||||
match scan_one(ctx, id, feed_cfg, state, force, fetched).await {
|
||||
Ok(Outcome::Feed(s)) => ctx.out.emit(Event::FeedDone {
|
||||
feed: id.clone(),
|
||||
new: s.new_entries,
|
||||
@@ -1039,7 +1068,7 @@ async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]
|
||||
scanned += 1;
|
||||
ctx.out.emit(Event::FeedStart { feed: id.clone() });
|
||||
let state = ctx.db.http_state(id).await?;
|
||||
match scan_one(ctx, id, feed_cfg, &state, force).await {
|
||||
match scan_one(ctx, id, feed_cfg, &state, force, None).await {
|
||||
Ok(Outcome::Feed(s)) => ctx.out.emit(Event::FeedDone {
|
||||
feed: id.clone(),
|
||||
new: s.new_entries,
|
||||
@@ -1221,6 +1250,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>>,
|
||||
) -> Result<Outcome> {
|
||||
// A Patreon creator with more than one show is a list of feeds, like an OPML.
|
||||
if feed::is_patreon_creator(&feed_cfg.url) {
|
||||
@@ -1247,13 +1278,10 @@ async fn scan_one(
|
||||
}
|
||||
}
|
||||
|
||||
let mut fetched = feed::fetch(
|
||||
&ctx.client,
|
||||
feed_cfg,
|
||||
state.etag.as_deref(),
|
||||
state.last_modified.as_deref(),
|
||||
)
|
||||
.await?;
|
||||
let mut fetched = match prefetched {
|
||||
Some(got) => got?,
|
||||
None => feed::fetch(&ctx.client, feed_cfg, state.etag.as_deref(), state.last_modified.as_deref()).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
|
||||
|
||||
Reference in New Issue
Block a user