Send a changed feed's row to the page instead of it reloading the list (#105)
After a feed was checked, failed or downloaded a file, and after every item read, the page fetched /api/feeds whole, about 60 ms for 160 rows, though one row had changed. The live event stream now knows who is connected and, after an event that changes a feed, sends that person its row (feed_row), built by the same code as the list (feed_rows, with Db::feed_list asked for one feed). Marking an item read answers with the feed's row. The page puts the row in place and redraws once a frame. A routine skip of a feed not due sends nothing. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
38
src/db.rs
38
src/db.rs
@@ -338,8 +338,14 @@ impl Db {
|
||||
/// What the feed list shows of every feed, for one person, in five queries whatever the
|
||||
/// number of feeds. Asked feed by feed (feed_summary, http_state, blocklist, unread_count)
|
||||
/// it was six round trips a feed, about 950 for 160 feeds and 320 ms a page load (#94).
|
||||
pub async fn feed_list(&self, user_id: i64) -> Result<std::collections::HashMap<String, FeedListing>> {
|
||||
pub async fn feed_list(
|
||||
&self,
|
||||
user_id: i64,
|
||||
// One feed only, for the row a live update sends (`web::feed_rows`).
|
||||
only: Option<&str>,
|
||||
) -> Result<std::collections::HashMap<String, FeedListing>> {
|
||||
use std::collections::HashMap;
|
||||
let one = || sea_orm::Value::from(only.map(str::to_owned));
|
||||
let counts = |sql: &'static str, args: Vec<sea_orm::Value>| async move {
|
||||
self.rows(sql, args)
|
||||
.await?
|
||||
@@ -347,20 +353,29 @@ impl Db {
|
||||
.map(|r| Ok((r.try_get::<String>("", "feed_id")?, r.try_get::<i64>("", "n")?)))
|
||||
.collect::<Result<HashMap<_, _>>>()
|
||||
};
|
||||
let entries = counts("SELECT feed_id, count(*) AS n FROM entries GROUP BY feed_id", vec![]).await?;
|
||||
let downloaded =
|
||||
counts("SELECT feed_id, count(*) AS n FROM enclosures WHERE path IS NOT NULL GROUP BY feed_id", vec![])
|
||||
.await?;
|
||||
let entries = counts(
|
||||
"SELECT feed_id, count(*) AS n FROM entries
|
||||
WHERE (CAST($1 AS TEXT) IS NULL OR feed_id = $1) GROUP BY feed_id",
|
||||
vec![one()],
|
||||
)
|
||||
.await?;
|
||||
let downloaded = counts(
|
||||
"SELECT feed_id, count(*) AS n FROM enclosures
|
||||
WHERE path IS NOT NULL AND (CAST($1 AS TEXT) IS NULL OR feed_id = $1) GROUP BY feed_id",
|
||||
vec![one()],
|
||||
)
|
||||
.await?;
|
||||
let unread = counts(
|
||||
"SELECT e.feed_id, count(*) AS n FROM entries e
|
||||
JOIN subscriptions sub ON sub.user_id = $1 AND sub.feed_id = e.feed_id
|
||||
LEFT JOIN entry_state s
|
||||
ON s.user_id = $1 AND s.feed_id = e.feed_id AND s.guid = e.guid
|
||||
WHERE NOT coalesce(s.read, false)
|
||||
AND (CAST($2 AS TEXT) IS NULL OR e.feed_id = $2)
|
||||
AND NOT EXISTS (SELECT 1 FROM hidden h
|
||||
WHERE h.user_id = $1 AND h.feed_id = e.feed_id AND h.guid = e.guid)
|
||||
GROUP BY e.feed_id",
|
||||
vec![user_id.into()],
|
||||
vec![user_id.into(), one()],
|
||||
)
|
||||
.await?;
|
||||
let mut blocked: HashMap<String, Vec<String>> = blocklists::Entity::find()
|
||||
@@ -370,7 +385,11 @@ impl Db {
|
||||
.into_iter()
|
||||
.map(|b| (b.feed_id, keywords(Some(b.words)).unwrap_or_default()))
|
||||
.collect();
|
||||
Ok(feeds::Entity::find()
|
||||
let mut rows = feeds::Entity::find();
|
||||
if let Some(id) = only {
|
||||
rows = rows.filter(feeds::Column::Id.eq(id));
|
||||
}
|
||||
Ok(rows
|
||||
.all(&self.orm)
|
||||
.await?
|
||||
.into_iter()
|
||||
@@ -1964,7 +1983,10 @@ mod tests {
|
||||
INSERT INTO blocklists (user_id, feed_id, words) VALUES (1,'g','[\"spoiler\"]');",
|
||||
).await
|
||||
.unwrap();
|
||||
let list = db.feed_list(1).await.unwrap();
|
||||
let list = db.feed_list(1, None).await.unwrap();
|
||||
let g = db.feed_list(1, Some("g")).await.unwrap();
|
||||
assert_eq!(g.keys().collect::<Vec<_>>(), ["g"]);
|
||||
assert_eq!((g["g"].unread, g["g"].summary.entries, g["g"].blocked.clone()), (1, 1, vec!["spoiler".to_string()]));
|
||||
for id in ["f", "g"] {
|
||||
let one = &list[id];
|
||||
let s = db.feed_summary(id).await.unwrap();
|
||||
|
||||
52
src/web.rs
52
src/web.rs
@@ -731,6 +731,12 @@ async fn feeds(
|
||||
State(state): State<WebState>,
|
||||
user: crate::db::User,
|
||||
) -> Result<Json<Vec<FeedRow>>, ApiError> {
|
||||
Ok(Json(feed_rows(&state, &user, None).await?))
|
||||
}
|
||||
|
||||
/// The feed list as this person sees it, or with `only` the one row for that feed: what a live
|
||||
/// update sends after the feed changes, so the page redraws a row instead of reloading the list.
|
||||
async fn feed_rows(state: &WebState, user: &crate::db::User, only: Option<&str>) -> anyhow::Result<Vec<FeedRow>> {
|
||||
let cfg = state.ctx.cfg();
|
||||
// Config entries plus the feeds derived from OPML subscriptions -- the catalogue.
|
||||
// What comes back is only the part of it this person subscribes to.
|
||||
@@ -744,9 +750,9 @@ async fn feeds(
|
||||
.collect();
|
||||
let counts = state.ctx.db.subscriber_counts().await?;
|
||||
let pinned = state.ctx.db.pinned_feeds(user.id).await?;
|
||||
let mut listed = state.ctx.db.feed_list(user.id).await?;
|
||||
let mut listed = state.ctx.db.feed_list(user.id, only).await?;
|
||||
let mut out = Vec::with_capacity(mine.len());
|
||||
for sub in &subs {
|
||||
for sub in subs.iter().filter(|s| only.is_none_or(|o| o == s.id)) {
|
||||
let (id, feed) = (&sub.id, &sub.cfg);
|
||||
// In a group, what you have not set on the feed comes from your settings on the group,
|
||||
// the same fallback the scanner uses (`Db::subscribers`).
|
||||
@@ -809,7 +815,7 @@ async fn feeds(
|
||||
pinned: pinned.contains(id),
|
||||
});
|
||||
}
|
||||
Ok(Json(out))
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
// ---- popular on this server ----
|
||||
@@ -1523,7 +1529,7 @@ async fn set_flags(
|
||||
Path((feed_id, guid)): Path<(String, String)>,
|
||||
user: crate::db::User,
|
||||
Json(body): Json<Flags>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
) -> Result<Json<Option<FeedRow>>, ApiError> {
|
||||
use crate::db::EntryFlag;
|
||||
if let Some(v) = body.read {
|
||||
state.ctx.db.set_entry_flag(user.id, &feed_id, &guid, EntryFlag::Read, v).await?;
|
||||
@@ -1531,7 +1537,9 @@ async fn set_flags(
|
||||
if let Some(v) = body.flagged {
|
||||
state.ctx.db.set_entry_flag(user.id, &feed_id, &guid, EntryFlag::Flagged, v).await?;
|
||||
}
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
// The feed's row with its new unread count, for the page to put in place of the old one
|
||||
// rather than reloading the whole list after every item read.
|
||||
Ok(Json(feed_rows(&state, &user, Some(&feed_id)).await?.into_iter().next()))
|
||||
}
|
||||
|
||||
/// Downloads one enclosure now. This cannot be "requeue and scan": a scan takes the
|
||||
@@ -1644,23 +1652,45 @@ async fn fetch_now(
|
||||
}
|
||||
|
||||
/// The same broadcast the socket clients read, as server-sent events.
|
||||
async fn events(State(state): State<WebState>) -> Sse<impl futures_util::Stream<Item = Result<SseEvent, std::convert::Infallible>>> {
|
||||
async fn events(
|
||||
State(state): State<WebState>,
|
||||
user: crate::db::User,
|
||||
) -> Sse<impl futures_util::Stream<Item = Result<SseEvent, std::convert::Infallible>>> {
|
||||
let rx = state.events.subscribe();
|
||||
// A client that falls behind skips what it missed rather than being cut off.
|
||||
let stream = futures_util::stream::unfold(state.events.subscribe(), |mut rx| async move {
|
||||
let stream = futures_util::stream::unfold((rx, state, user), |(mut rx, state, user)| async move {
|
||||
loop {
|
||||
match rx.recv().await {
|
||||
Ok(ev) => {
|
||||
if let Ok(data) = serde_json::to_string(&ev) {
|
||||
let ev = Ok::<_, std::convert::Infallible>(SseEvent::default().data(data));
|
||||
return Some((ev, rx));
|
||||
let Ok(data) = serde_json::to_string(&ev) else { continue };
|
||||
let mut out = vec![Ok::<_, std::convert::Infallible>(SseEvent::default().data(data))];
|
||||
// A feed that changed goes out as this person's row for it, which the page
|
||||
// puts in place of the old one: it reloaded the whole list after each.
|
||||
if let Some(feed) = changed_feed(&ev)
|
||||
&& let Ok(rows) = feed_rows(&state, &user, Some(feed)).await
|
||||
&& let Some(row) = rows.first()
|
||||
&& let Ok(data) = serde_json::to_string(&serde_json::json!({ "ev": "feed_row", "row": row }))
|
||||
{
|
||||
out.push(Ok(SseEvent::default().data(data)));
|
||||
}
|
||||
return Some((futures_util::stream::iter(out), (rx, state, user)));
|
||||
}
|
||||
Err(broadcast::error::RecvError::Lagged(_)) => {}
|
||||
Err(broadcast::error::RecvError::Closed) => return None,
|
||||
}
|
||||
}
|
||||
});
|
||||
Sse::new(stream).keep_alive(axum::response::sse::KeepAlive::default())
|
||||
Sse::new(futures_util::StreamExt::flatten(stream)).keep_alive(axum::response::sse::KeepAlive::default())
|
||||
}
|
||||
|
||||
/// The feed an event changed what the list shows of: its counts, error or last check. Not the
|
||||
/// routine skip of a feed not due, dozens a minute that change nothing.
|
||||
fn changed_feed(ev: &Event) -> Option<&str> {
|
||||
match ev {
|
||||
Event::FeedDone { feed, .. } | Event::FeedError { feed, .. } | Event::DownloadDone { feed, .. } => Some(feed),
|
||||
Event::FeedSkip { feed, reason } if !reason.starts_with("not due") => Some(feed),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Audio, served by ServeFile so Range requests work and the player can seek.
|
||||
|
||||
Reference in New Issue
Block a user