Files
ipx/src/main.rs
rays e7c59489ee Announce a followed move as a feed_moved event (#115)
A feed moved to its new address was only logged by follow_move. It is now an event, feed_moved
with the feed and its old and new addresses, so it goes where every other event goes: the log,
with from and to as fields, the admin page's Scans view, `ipx fetch`, and the page, which gets
the feed's new row.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-02 20:22:28 +00:00

2277 lines
93 KiB
Rust

mod access;
mod art;
mod auth;
mod config;
mod db;
mod entity;
mod download;
mod feed;
mod ipc;
mod logbuf;
mod retention;
mod torrent;
mod web;
use anyhow::{Context, Result};
use clap::{Parser, Subcommand};
use ipc::{Command as Cmd, Emitter, Event};
use std::path::PathBuf;
use std::sync::Arc;
use std::future::Future;
use tokio::sync::{broadcast, mpsc, watch};
#[derive(Parser)]
#[command(name = "ipx", version, about = "A headless podcatcher")]
struct Cli {
/// Config file (default: $XDG_CONFIG_HOME/ipx/config.toml)
#[arg(long, global = true)]
config: Option<PathBuf>,
/// Do the work here even if a daemon is running
#[arg(long, global = true)]
local: bool,
#[command(subcommand)]
command: Command,
}
#[derive(Subcommand)]
enum Command {
/// Show configured feeds and their state
List,
/// Scan feeds for new entries
Fetch {
/// Only this feed id
feed: Option<String>,
/// Read the feed in full now, even when it is not due and has not changed
#[arg(long)]
force: bool,
},
/// Delete old or over-quota downloads
Reap {
/// Show what would go, delete nothing
#[arg(long)]
dry_run: bool,
},
/// Counts of feeds, pending and downloaded enclosures
Status,
/// Subscribe to a feed
Add {
url: String,
/// Download folder name (default: the feed title)
#[arg(long)]
folder: Option<String>,
/// Only take enclosures matching these keywords
#[arg(long, value_delimiter = ',')]
keywords: Vec<String>,
/// The Directory's category, for a feed that names none of its own (News, Technology)
#[arg(long)]
category: Option<String>,
/// Put it in the Directory for anyone to subscribe to, and keep it there when nobody does
#[arg(long)]
list: bool,
},
/// Unsubscribe. Downloads and history are left alone.
Rm { feed: String },
/// Add every feed in an OPML file
Import { file: PathBuf },
/// Write subscriptions out as OPML
Export { file: PathBuf },
/// Add, list or remove the accounts that can sign in to the web UI
User {
#[command(subcommand)]
cmd: UserCmd,
},
/// Run the scheduler and serve the control socket
Daemon {
/// Serve the web UI on this address, overriding [web] in the config
#[arg(long, value_name = "ADDR")]
web: Option<String>,
},
}
#[derive(Subcommand)]
enum UserCmd {
/// Create an account. The password is read from stdin: `echo -n hunter2 | ipx user add ray`
Add {
name: String,
/// May manage other accounts, and is who the shared web token signs in as
#[arg(long)]
admin: bool,
/// Sign-in comes from the proxy instead, so there is no password to set
#[arg(long)]
no_password: bool,
},
/// Show the accounts and how each one signs in
List,
/// Replace a password, read from stdin
Passwd { name: String },
/// Delete an account and everything it knows: its subscriptions and read state
Rm { name: String },
/// Rename an account, keeping its feeds, read state and admin rights. This is how an
/// account made before the proxy takes the name the proxy signs it in as
Rename { name: String, new_name: String },
}
/// What a brand new database starts with, so there is always a way in. Announced loudly
/// in the log, and the first thing the settings page nags about.
const DEFAULT_PASSWORD: &str = "ipodderx";
/// Everything a command needs. One per process.
pub struct Ctx {
/// Swapped wholesale when the web UI rewrites config.toml, so a running daemon picks
/// up feed changes without a restart. Callers take a snapshot; no guard is ever held
/// across an await.
pub cfg: std::sync::RwLock<std::sync::Arc<config::Config>>,
pub db: db::Db,
pub client: reqwest::Client,
/// For feeds only: follows no redirects, so `feed::fetch` sees each hop and can tell a feed
/// that moved for good from one sent elsewhere for now.
pub feed_client: reqwest::Client,
pub out: Emitter,
/// Started on first use: a BitTorrent session binds ports and starts a DHT, which is
/// rude to do for a config that has never seen a torrent.
pub torrents: tokio::sync::OnceCell<torrent::Torrents>,
/// Caps how many torrents run at once when they are detached.
pub torrent_slots: Arc<tokio::sync::Semaphore>,
/// Needed because a subscribed OPML rewrites the feed list as it syncs.
pub config_path: PathBuf,
/// Run torrents off the command worker. A torrent takes minutes to fetch metadata and
/// then seeds for up to an hour, and the worker is sequential -- inline, one torrent
/// stops every feed scan, every HTTP download and every status command behind it.
/// A one-shot CLI run keeps them inline, or the process would exit mid-download.
pub detach_torrents: bool,
}
impl Ctx {
pub fn cfg(&self) -> std::sync::Arc<config::Config> {
self.cfg.read().unwrap().clone()
}
/// Re-reads config.toml into the live snapshot.
/// Keeps a changed catalogue or server settings: in the database, and for everything running
/// here from now on. What config.toml's save and reload did, before the database held them.
pub async fn store_cfg(&self, cfg: config::Config) -> Result<()> {
self.db.store_config(&cfg).await?;
self.set_cfg(cfg);
Ok(())
}
fn set_cfg(&self, cfg: config::Config) {
*self.cfg.write().unwrap() = std::sync::Arc::new(cfg);
tracing::info!("config reloaded");
}
async fn torrents(&self) -> Result<&torrent::Torrents> {
let cfg = self.cfg();
self.torrents
.get_or_try_init(|| torrent::Torrents::new(&cfg))
.await
}
}
#[tokio::main]
async fn main() -> Result<()> {
let cli = Cli::parse();
// The daemon only: `ipx status` runs every half minute as the healthcheck, and a trace
// for each would bury the ones worth looking at.
let otel = match cli.command {
Command::Daemon { .. } => otel_provider()?,
_ => None,
};
// Everything goes to stderr as before, and is mirrored into a ring the UI can read.
{
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use tracing_subscriber::Layer;
// Two filters, deliberately different. stderr follows IPX_LOG; the in-app buffer
// keeps debug as well, so the log view can show protocol traffic and routine
// skips that would be noise on a terminal. IPX_UI_LOG overrides it.
let stderr_filter = || {
tracing_subscriber::EnvFilter::try_from_env("IPX_LOG").unwrap_or_else(|_| "ipx=info".into())
};
// One JSON object a line for Loki (#91), with each event's fields as its own; text
// otherwise, for someone reading a terminal.
let json = std::env::var("IPX_LOG_FORMAT").is_ok_and(|f| f.eq_ignore_ascii_case("json"));
let ui_filter = tracing_subscriber::EnvFilter::try_from_env("IPX_UI_LOG")
.unwrap_or_else(|_| "ipx=debug".into());
tracing_subscriber::registry()
.with((!json).then(|| {
tracing_subscriber::fmt::layer()
.with_writer(std::io::stderr)
// Colour for a terminal only: in docker logs and Loki the escapes are noise
// every query has to strip (#88).
.with_ansi(std::io::IsTerminal::is_terminal(&std::io::stderr()))
.with_filter(stderr_filter())
}))
.with(json.then(|| {
tracing_subscriber::fmt::layer()
.fmt_fields(tracing_subscriber::fmt::format::JsonFields::new())
.event_format(WithTrace(
tracing_subscriber::fmt::format().json().flatten_event(true).with_current_span(true).with_span_list(false),
))
.with_writer(std::io::stderr)
.with_filter(stderr_filter())
}))
.with(logbuf::RingLayer.with_filter(ui_filter))
.with(otel.as_ref().map(|p| {
use opentelemetry::trace::TracerProvider;
tracing_opentelemetry::layer()
.with_tracer(p.tracer("ipx"))
.with_filter(tracing_subscriber::EnvFilter::new("ipx=info"))
}))
.init();
}
let config_path = cli.config.clone().unwrap_or_else(config::config_path);
let cfg = config::Config::load(&config_path)?;
let db = db::Db::open(&db::location()).await?;
// A daemon owns the state; don't have two processes downloading the same thing.
let wire_cmd = match &cli.command {
Command::Fetch { feed, force } => {
Some(Cmd::Fetch { feed: feed.clone(), force: *force, feeds: vec![] })
}
Command::Reap { dry_run } => Some(Cmd::Reap { dry_run: *dry_run }),
Command::Status => Some(Cmd::Status),
Command::List
| Command::Daemon { .. }
| Command::User { .. }
| Command::Add { .. }
| Command::Rm { .. }
| Command::Import { .. }
| Command::Export { .. } => None,
};
if let Some(cmd) = &wire_cmd
&& !cli.local
&& ipc::daemon_is_live(&cfg.general.socket).await
{
return ipc::proxy(&cfg.general.socket, cmd).await;
}
let cfg = assemble_config(&db, cfg, &config_path).await?;
let is_daemon = matches!(cli.command, Command::Daemon { .. });
let (events, _) = broadcast::channel(1024);
let ctx = Arc::new(Ctx {
cfg: std::sync::RwLock::new(std::sync::Arc::new(cfg)),
db,
client: reqwest::Client::builder()
.user_agent(concat!("ipx/", env!("CARGO_PKG_VERSION")))
// Connecting only, so it bounds a download's start, not a long download (#108).
.connect_timeout(std::time::Duration::from_secs(10))
.build()?,
feed_client: reqwest::Client::builder()
.user_agent(concat!("ipx/", env!("CARGO_PKG_VERSION")))
.connect_timeout(std::time::Duration::from_secs(10))
.redirect(reqwest::redirect::Policy::none())
.build()?,
out: if is_daemon { Emitter::socket(events.clone(), false) } else { Emitter::terminal() },
torrents: tokio::sync::OnceCell::new(),
torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)),
config_path: config_path.clone(),
detach_torrents: is_daemon,
});
let result = match cli.command {
Command::List => list(&ctx).await,
Command::Daemon { web } => daemon(ctx, config_path, web, events).await,
Command::Add { url, folder, keywords, category, list } => {
add(&ctx, &url, folder, keywords, category, list).await
}
Command::Rm { feed } => rm(&ctx, &feed).await,
Command::User { cmd } => user_cmd(&ctx, cmd).await,
Command::Import { file } => import(&ctx, &file).await,
Command::Export { file } => export(&ctx, &file).await,
_ => run(&ctx, wire_cmd.expect("only List and Daemon have no wire form")).await,
};
// The batch exporter holds the last few seconds of spans; without this they are lost.
if let Some(p) = otel {
let _ = p.shutdown();
}
result
}
/// The JSON log line with the trace and span it belongs to (#91), so a line in Loki leads to its
/// trace in Tempo. The JSON formatter cannot take a field of its own, so the ids go on the end of
/// the object it writes. A line outside any traced span, or with no trace exporter, is unchanged.
struct WithTrace<F>(F);
impl<S, N, F> tracing_subscriber::fmt::FormatEvent<S, N> for WithTrace<F>
where
F: tracing_subscriber::fmt::FormatEvent<S, N>,
S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
N: for<'w> tracing_subscriber::fmt::FormatFields<'w> + 'static,
{
fn format_event(
&self,
ctx: &tracing_subscriber::fmt::FmtContext<'_, S, N>,
mut w: tracing_subscriber::fmt::format::Writer<'_>,
ev: &tracing::Event<'_>,
) -> std::fmt::Result {
use opentelemetry::trace::TraceContextExt;
use tracing_opentelemetry::OpenTelemetrySpanExt;
let current = tracing::Span::current();
let sc = if current.is_none() { None } else { Some(current.context().span().span_context().clone()) };
let Some(sc) = sc.filter(|c| c.is_valid()) else {
return self.0.format_event(ctx, w, ev);
};
let mut line = String::new();
self.0.format_event(ctx, tracing_subscriber::fmt::format::Writer::new(&mut line), ev)?;
match line.trim_end().strip_suffix('}') {
Some(body) => writeln!(w, r#"{body},"trace_id":"{}","span_id":"{}"}}"#, sc.trace_id(), sc.span_id()),
None => w.write_str(&line),
}
}
}
/// Traces over OTLP, to Tempo for one, when OTEL_EXPORTER_OTLP_ENDPOINT names a collector
/// (`http://host:4318`: the exporter speaks OTLP over HTTP and adds `/v1/traces`). The exporter
/// reads that and the other `OTEL_` variables itself.
fn otel_provider() -> Result<Option<opentelemetry_sdk::trace::SdkTracerProvider>> {
if std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").unwrap_or_default().is_empty() {
return Ok(None);
}
let exporter = opentelemetry_otlp::SpanExporter::builder().with_http().build()?;
let resource = opentelemetry_sdk::Resource::builder().with_service_name("ipx").build();
Ok(Some(
opentelemetry_sdk::trace::SdkTracerProvider::builder()
.with_batch_exporter(exporter)
.with_resource(resource)
.build(),
))
}
/// Accounts. Passwords come in on stdin so they never reach a shell history or a `ps`
/// listing.
async fn user_cmd(ctx: &Arc<Ctx>, cmd: UserCmd) -> Result<()> {
let read_password = || -> Result<String> {
use std::io::Read;
let mut buf = String::new();
std::io::stdin().read_to_string(&mut buf)?;
let pw = buf.trim_end_matches(['\n', '\r']).to_string();
if pw.is_empty() {
anyhow::bail!("no password on stdin: try `echo -n secret | ipx user ...`");
}
Ok(pw)
};
match cmd {
UserCmd::Add { name, admin, no_password } => {
let name = name.trim().to_ascii_lowercase();
if name.is_empty() {
anyhow::bail!("a name is required");
}
if ctx.db.user_by_name(&name).await?.is_some() {
anyhow::bail!("{name} already exists");
}
let hash = if no_password {
None
} else {
Some(crate::auth::hash_password(&read_password()?)?)
};
// The first account runs the place; there is nobody else to grant it.
let first = ctx.db.users().await?.is_empty();
ctx.db.create_user(&name, hash.as_deref(), admin || first).await?;
println!(
"added {name}{}{}",
if admin || first { " (admin)" } else { "" },
if no_password { ", signs in through the proxy" } else { "" }
);
Ok(())
}
UserCmd::List => {
let users = ctx.db.users().await?;
if users.is_empty() {
println!("no accounts yet: ipx user add <name>");
}
for u in users {
let added = u
.created
.and_then(|t| chrono::DateTime::from_timestamp(t, 0))
.map_or("?".into(), |d| d.format("%Y-%m-%d").to_string());
let seen = u.last_login.map_or("never signed in".into(), |t| format!("signed in {}", ago(Some(t))));
println!(
"{:<20} {:<6} {:<11} added {added} {seen}",
u.name,
if u.is_admin { "admin" } else { "" },
if u.pass_hash.is_some() { "password" } else { "proxy only" },
);
}
Ok(())
}
UserCmd::Passwd { name } => {
let name = name.trim().to_ascii_lowercase();
let user = ctx
.db
.user_by_name(&name).await?
.ok_or_else(|| anyhow::anyhow!("no such account: {name}"))?;
ctx.db.set_password(user.id, &crate::auth::hash_password(&read_password()?)?).await?;
println!("password changed for {name}");
Ok(())
}
UserCmd::Rename { name, new_name } => {
let name = name.trim().to_ascii_lowercase();
// The same rules as a name the proxy vouches for, or the proxy would never find it.
let new_name = crate::auth::name_from_header(&new_name)
.ok_or_else(|| anyhow::anyhow!("not a usable name: no commas, semicolons or line breaks"))?;
let user = ctx
.db
.user_by_name(&name).await?
.ok_or_else(|| anyhow::anyhow!("no such account: {name}"))?;
if ctx.db.user_by_name(&new_name).await?.is_some() {
anyhow::bail!("{new_name} already exists");
}
ctx.db.rename_user(user.id, &new_name).await?;
println!("renamed {name} to {new_name}");
Ok(())
}
UserCmd::Rm { name } => {
let name = name.trim().to_ascii_lowercase();
let user = ctx
.db
.user_by_name(&name).await?
.ok_or_else(|| anyhow::anyhow!("no such account: {name}"))?;
ctx.db.delete_user(user.id).await?;
println!("removed {name}");
Ok(())
}
}
}
async fn run(ctx: &Arc<Ctx>, cmd: Cmd) -> Result<()> {
match cmd {
Cmd::Fetch { feed, force, feeds } => {
// Make room before pulling more down, as the original did per download.
reap(ctx, false, false).await?;
fetch(ctx, feed.as_deref(), force, &feeds).await
}
Cmd::Reap { dry_run } => reap(ctx, dry_run, true).await,
Cmd::Download { enclosure } => download_one(ctx, enclosure).await,
Cmd::Status => {
ctx.out.emit(status(ctx).await);
Ok(())
}
}
}
/// The counts `ipx status` prints, and /api/status serves. A running daemon's socket answers with
/// this directly rather than through the job queue.
pub(crate) async fn status(ctx: &Ctx) -> Event {
match ctx.db.counts().await {
Ok((pending, downloaded)) => {
let feeds = subscriptions(ctx).await.map(|s| s.len()).unwrap_or(0);
Event::Status { feeds, pending, downloaded }
}
Err(e) => Event::Error { msg: format!("{e:#}") },
}
}
async fn daemon(
ctx: Arc<Ctx>,
config_path: PathBuf,
web_addr: Option<String>,
events: broadcast::Sender<Event>,
) -> Result<()> {
let socket = ctx.cfg().general.socket.clone();
if ipc::daemon_is_live(&socket).await {
anyhow::bail!("a daemon is already listening on {}", socket.display());
}
// A database with nobody in it cannot be signed into.
if ctx.db.users().await?.is_empty() {
ctx.db.create_user("admin", Some(&crate::auth::hash_password(DEFAULT_PASSWORD)?), true).await?;
tracing::warn!(
"no accounts yet: created 'admin' with the default password '{DEFAULT_PASSWORD}'. \
Change it with `echo -n <password> | ipx user passwd admin`"
);
}
if let Some(admin) = ctx.db.users().await?.into_iter().find(|u| u.is_admin) {
let catalogue: Vec<String> = ctx.cfg().feeds.keys().cloned().collect();
match ctx.db.adopt_catalogue(admin.id, &catalogue).await {
Ok(0) => {}
Ok(n) => tracing::info!(user = %admin.name, feeds = n, "subscribed the first admin to the catalogue"),
Err(e) => tracing::error!(error = %e, "could not subscribe the first admin to the catalogue"),
}
}
match ctx.db.requeue_interrupted().await {
Ok(n) if n > 0 => tracing::info!(count = n, "requeued downloads interrupted by a restart"),
Ok(_) => {}
Err(e) => tracing::warn!(error = %format!("{e:#}"), "could not requeue interrupted downloads"),
}
match retire_stranded(&ctx).await {
Ok(0) => {}
Ok(n) => tracing::info!(feeds = n, "retired feeds whose OPML is no longer in config"),
Err(e) => tracing::warn!(error = %format!("{e:#}"), "could not retire feeds whose OPML is no longer in config"),
}
let (tx_cmd, mut rx_cmd) = mpsc::channel::<Cmd>(64);
let web = start_web(&ctx, &config_path, web_addr, &tx_cmd, &events).await?;
// status is answered by the socket itself; everything else waits its turn in the queue.
let answer: ipc::StatusFn = {
let ctx = ctx.clone();
Arc::new(move || {
let ctx = ctx.clone();
Box::pin(async move { status(&ctx).await })
})
};
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.
tracing::info!(
feeds = subscriptions(&ctx).await.map(|s| s.len()).unwrap_or(0),
"daemon started"
);
// A signal has to be able to interrupt work in progress, not just the wait between
// jobs. Racing `shutdown()` in the outer select only cancels branch selection: once
// inside a long download the daemon stopped listening and had to be SIGKILLed.
let (tx_stop, rx_stop) = tokio::sync::watch::channel(false);
tokio::spawn(async move {
shutdown().await;
let _ = tx_stop.send(true);
});
// Runs one job, abandoning it if a signal arrives. Returns false to end the loop.
async fn until_stopped(
ctx: &Ctx,
rx: &watch::Receiver<bool>,
job: impl Future<Output = Result<()>>,
) -> bool {
let mut stop = rx.clone();
tokio::select! {
_ = stop.changed() => {
tracing::info!("signal received; abandoning the job in progress");
false
}
result = job => {
if let Err(e) = result {
ctx.out.emit(Event::Error { msg: format!("{e:#}") });
}
true
}
}
}
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,
Some(cmd) = rx_cmd.recv() => {
// Both halves of the protocol are logged under one target so the UI can
// show the conversation on its own: this is everything arriving, whatever
// the source -- a socket client, the CLI proxying, or the web UI.
tracing::info!(
target: "ipx::io",
"-> {}",
serde_json::to_string(&cmd).unwrap_or_else(|_| format!("{cmd:?}"))
);
if !until_stopped(&ctx, &rx_stop, run(&ctx, cmd)).await {
break;
}
}
_ = 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 {
break;
}
}
}
}
server.abort();
if let Some(w) = web {
w.abort();
}
let _ = std::fs::remove_file(&socket);
tracing::info!("daemon stopped");
Ok(())
}
/// Starts the web UI when it is switched on, minting and saving a token if there is none.
async fn start_web(
ctx: &Arc<Ctx>,
config_path: &std::path::Path,
web_addr: Option<String>,
cmds: &mpsc::Sender<Cmd>,
events: &broadcast::Sender<Event>,
) -> Result<Option<tokio::task::JoinHandle<()>>> {
let cfg = ctx.cfg();
let enabled = cfg.web.enabled || web_addr.is_some();
if !enabled {
return Ok(None);
}
let bind = web_addr.unwrap_or_else(|| cfg.web.bind.clone());
if cfg.web.token.is_empty() {
let mut fresh = (*cfg).clone();
fresh.web.enabled = true;
fresh.web.bind = bind.clone();
fresh.web.token = crate::auth::new_session_token();
// The token is config.toml's, not the database's: it decides who gets in.
fresh.save_bootstrap(config_path)?;
ctx.set_cfg(fresh.clone());
// The token signs in as the admin, and whatever reads this process's output (docker logs,
// for one) is wider than who reads config.toml. So say where it is, never what it is.
// Logged, not printed, so a JSON log stays one object a line (#91).
tracing::info!(
"web ui token generated and saved to {} as [web] token. Open http://{bind}/?token=<that token>",
config_path.display()
);
} else {
tracing::info!(
"web ui at http://{bind}/ (the sign-in token is [web] token in {})",
config_path.display()
);
}
// Bound beyond localhost, as a container has to be for its port to be published. Worth a
// warning only when nothing but a password stands in front of it: it said "the token is all
// that guards it" on every start of production, behind Cloudflare Access, and was the only
// warning in a healthy log (#106).
if ctx.cfg().web.binds_publicly() {
if ctx.cfg().web.access().is_some() {
tracing::info!(
bind,
"web ui is reachable off this machine; signing in takes an account's password or the admin token, or Cloudflare Access through a trusted proxy"
);
} else {
tracing::warn!(
bind,
"web ui is reachable off this machine; an account's password or the admin token is all that guards it"
);
}
}
let access = Arc::new(access::Keys::default());
if let Some((team, _)) = ctx.cfg().web.access() {
let (access, ctx, team) = (access.clone(), ctx.clone(), team.to_owned());
tokio::spawn(async move { access.prefetch(&ctx.client, &team).await });
}
let state = web::WebState {
ctx: ctx.clone(),
cmds: cmds.clone(),
events: events.clone(),
access,
};
Ok(Some(tokio::spawn(async move {
if let Err(e) = web::serve(state, &bind).await {
tracing::error!(error = %format!("{e:#}"), "web ui stopped");
}
})))
}
async fn shutdown() {
use tokio::signal::unix::{SignalKind, signal};
let mut term = match signal(SignalKind::terminate()) {
Ok(s) => s,
Err(_) => return std::future::pending().await,
};
tokio::select! {
_ = tokio::signal::ctrl_c() => {}
_ = term.recv() => {}
}
}
/// Subscribes to one feed, naming it from its own title.
async fn add(
ctx: &Ctx,
url: &str,
folder: Option<String>,
keywords: Vec<String>,
category: Option<String>,
list: bool,
) -> Result<()> {
let mut cfg = (*ctx.cfg()).clone();
let url = &feed::find_feed(&ctx.client, &feed::expand_input(url)).await?;
// Includes feeds derived from an OPML, or the same show could be added twice.
if let Some(existing) = subscriptions(ctx).await?.iter().find(|s| feed::same_feed(&s.cfg.url, url)) {
// Listing a feed the catalogue already has, or giving it a category, is the point of
// asking again: it changes those, and nothing else.
if let Some(f) = cfg.feeds.get_mut(&existing.id)
&& (list || category.is_some())
{
f.listed |= list;
if category.is_some() {
f.category = category;
}
ctx.store_cfg(cfg).await?;
println!("{} is already in the catalogue; updated its Directory listing", existing.id);
return Ok(());
}
anyhow::bail!("already subscribed as {:?}", existing.id);
}
let id = add_one(ctx, &mut cfg, url, folder, keywords).await?;
if let Some(f) = cfg.feeds.get_mut(&id) {
f.category = category;
f.listed = list;
}
ctx.store_cfg(cfg).await?;
println!("added {id}{}", if list { ", listed in the Directory" } else { "" });
Ok(())
}
/// Returns the new feed id. Adding checks the address is a feed first (`feed::find_feed`); this
/// still names one it cannot read from its URL rather than failing.
pub async fn add_one(
ctx: &Ctx,
cfg: &mut config::Config,
url: &str,
folder: Option<String>,
keywords: Vec<String>,
) -> Result<String> {
let probe = config::Feed {
url: url.to_owned(),
folder: folder.clone(),
group: None,
media_types: None,
schedule: None,
keywords: keywords.clone(),
allow_explicit: false,
auto_download: true,
max_new_per_check: None,
username: None,
password: None,
password_env: None,
category: None,
listed: false,
};
let title = match feed::fetch(&ctx.feed_client, &probe, None, None).await.map(|(got, _)| got) {
// An OPML subscription is named from its own <head><title>, not by trying to
// parse it as a feed and falling back to the hostname.
Ok(feed::Fetched::Body { bytes, .. }) if feed::is_opml(&bytes) => {
feed::opml_title(&bytes).unwrap_or_else(|| url_stem(url))
}
Ok(feed::Fetched::Body { bytes, .. }) => feed::parse(&bytes)
.ok()
.and_then(|f| f.title)
.unwrap_or_else(|| url_stem(url)),
_ => {
tracing::warn!(url, "could not read the feed; naming it from its URL");
url_stem(url)
}
};
// Slugs must be unique across derived feeds too, or a new feed can collide with one
// an OPML already introduced.
let mut taken: std::collections::BTreeMap<String, config::Feed> = subscriptions(ctx).await?
.into_iter()
.map(|s| (s.id, s.cfg))
.collect();
// A removed feed keeps its rows, so its id is only free again for the same feed: re-adding
// it gets its history back, and a different feed does not inherit someone else's.
for (id, other) in ctx.db.feed_urls().await? {
if !feed::same_feed(&other, url) {
taken.entry(id).or_insert_with(|| probe.clone());
}
}
let id = config::unique_slug(&title, &taken);
cfg.feeds.insert(id.clone(), probe);
Ok(id)
}
/// Host plus last path segment, for naming a feed we could not read.
fn url_stem(url: &str) -> String {
url::Url::parse(url)
.ok()
.and_then(|u| u.host_str().map(str::to_owned))
.unwrap_or_else(|| url.to_owned())
}
async fn rm(ctx: &Ctx, feed: &str) -> Result<()> {
let mut cfg = (*ctx.cfg()).clone();
if cfg.feeds.remove(feed).is_none() {
// Derived from an OPML: drop it here, though the subscription will list it again
// on the next read unless the OPML itself goes.
ctx.db.drop_managed(feed).await?;
println!("removed {feed}; it came from an OPML subscription and may return on the next read");
return Ok(());
}
ctx.store_cfg(cfg).await?;
// State and files stay: re-adding the feed should not re-download its back catalogue.
println!("removed {feed}; downloads and history kept");
retire_group(ctx, feed).await?;
Ok(())
}
async fn import(ctx: &Ctx, file: &std::path::Path) -> Result<()> {
let text = std::fs::read_to_string(file)
.with_context(|| format!("reading {}", file.display()))?;
// The CLI speaks for the operator, as the shared web token does.
let admin = ctx
.db
.users().await?
.into_iter()
.find(|u| u.is_admin)
.ok_or_else(|| anyhow::anyhow!("no admin account to subscribe: ipx user add <name> --admin"))?;
let doc = opml::OPML::from_str(&text)
.map_err(|e| anyhow::anyhow!("{} is not OPML: {e}", file.display()))?;
let (added, had) = subscribe_opml(ctx, &doc, admin.id).await?;
println!("subscribed {} to {added} feed(s); {had} already there", admin.name);
Ok(())
}
/// Subscribes one person to every feed in an OPML document, for the CLI and the web alike.
/// A feed already in the catalogue costs nothing; an unknown one is added under the OPML's
/// title rather than refetching each. Returns (newly subscribed, already subscribed).
///
/// Before accounts, importing only added unknown URLs to config.toml. Once subscriptions
/// decided what each person sees, that imported nothing at all for a feed someone else
/// already had, and a new one had no subscriber, so it was never scanned.
///
/// The caller parses the document, so each refuses a file that is not OPML in its own terms,
/// before anything is touched: a 400 from the web, a message from the CLI.
pub async fn subscribe_opml(
ctx: &Ctx,
doc: &opml::OPML,
user_id: i64,
) -> Result<(usize, usize)> {
let mut found = vec![];
collect_outlines(&doc.body.outlines, &mut found);
let known = subscriptions(ctx).await?;
let mut cfg = (*ctx.cfg()).clone();
let mut ids = vec![];
let mut grew = false;
for (title, url) in found {
let existing = known
.iter()
.find(|s| s.cfg.url == url)
.map(|s| s.id.clone())
// The same URL listed twice in one file.
.or_else(|| cfg.feeds.iter().find(|(_, f)| f.url == url).map(|(id, _)| id.clone()));
let id = match existing {
Some(id) => id,
None => {
let id = config::unique_slug(&title, &cfg.feeds);
cfg.feeds.insert(
id.clone(),
config::Feed {
url,
folder: None,
group: None,
media_types: None,
schedule: None,
keywords: vec![],
allow_explicit: false,
auto_download: true,
max_new_per_check: None,
username: None,
password: None,
password_env: None,
category: None,
listed: false,
},
);
grew = true;
id
}
};
ids.push(id);
}
if grew {
ctx.store_cfg(cfg).await?;
}
let (mut added, mut had) = (0, 0);
for id in ids {
if ctx.db.subscription(user_id, &id).await?.is_some() {
had += 1;
} else {
ctx.db.subscribe(user_id, &id).await?;
added += 1;
}
}
Ok((added, had))
}
/// OPML nests feeds inside folder outlines, so this walks the whole tree.
pub fn collect_outlines(outlines: &[opml::Outline], out: &mut Vec<(String, String)>) {
for o in outlines {
if let Some(url) = &o.xml_url {
let title = o.title.clone().unwrap_or_else(|| o.text.clone());
out.push((title, url.clone()));
}
collect_outlines(&o.outlines, out);
}
}
/// The configuration ipx runs with: config.toml for where things are and who may sign in, the
/// database for the feeds and the server settings (issue #18). The first time a database holds
/// neither, it takes them from config.toml, which is then cut down to the rest, the original kept
/// beside it as config.toml.pre-database.
async fn assemble_config(db: &db::Db, mut cfg: config::Config, path: &std::path::Path) -> Result<config::Config> {
for _ in 0..2 {
if let Some((stored, feeds)) = db.stored_config().await? {
if config::Config::file_holds_stored(path) {
tracing::warn!(
"config.toml still lists feeds or server settings; they are ignored, since the \
database holds them now. Change them in the web UI, or with ipx add and rm."
);
}
stored.apply(&mut cfg);
cfg.feeds = feeds;
return Ok(cfg);
}
if db.import_config(&cfg).await? {
if path.exists() {
let original = path.with_extension("toml.pre-database");
if !original.exists() {
std::fs::copy(path, &original)
.with_context(|| format!("keeping the original as {}", original.display()))?;
}
cfg.save_bootstrap(path)?;
}
tracing::info!(feeds = cfg.feeds.len(), "moved the feeds and server settings from config.toml into the database");
return Ok(cfg);
}
// Another ipx imported between our look and our insert; the next pass reads its copy.
}
anyhow::bail!("the database says it holds the configuration and then that it does not")
}
async fn export(ctx: &Ctx, file: &std::path::Path) -> Result<()> {
let mut doc = opml::OPML::default();
doc.head = Some(opml::Head {
title: Some("ipx subscriptions".into()),
..Default::default()
});
for (id, feed) in &ctx.cfg().feeds {
let title = ctx
.db
.feed_summary(id).await
.ok()
.and_then(|s| s.title)
.unwrap_or_else(|| id.clone());
doc.add_feed(&title, &feed.url);
}
let xml = doc.to_string().map_err(|e| anyhow::anyhow!("writing OPML: {e}"))?;
std::fs::write(file, xml).with_context(|| format!("writing {}", file.display()))?;
println!("exported {} feed(s) to {}", ctx.cfg().feeds.len(), file.display());
Ok(())
}
async fn list(ctx: &Ctx) -> Result<()> {
let cfg = ctx.cfg();
if cfg.feeds.is_empty() {
println!("No feeds configured. Add one with `ipx add <url>`, or in the web UI.");
return Ok(());
}
for (id, feed) in &cfg.feeds {
let s = ctx.db.feed_summary(id).await?;
// The id on a line of its own, labelled: first on the title's line, antirez.com's id
// "feed" read as a heading and 'ipx fetch antirez' was tried instead (issue #82).
println!("{}", s.title.as_deref().filter(|t| !t.is_empty()).unwrap_or(id));
println!(" id {id}");
println!(" url {}", feed.url);
println!(" last checked {}", ago(s.last_checked));
println!(" entries {} ({} downloaded)", s.entries, s.downloaded);
if let Some(err) = &s.last_error {
println!(" last error {err}");
}
}
Ok(())
}
/// Removes feeds nobody subscribes to that are dead (failing for 30 days) or quiet (nothing new
/// in a year), from the catalogue and the database, so the Directory lists feeds worth taking.
/// A feed listed in it with no subscribers stays for as long as it works and publishes.
async fn clean_directory(ctx: &Ctx) -> Result<()> {
const DEAD: i64 = 30 * 86_400;
const QUIET: i64 = 365 * 86_400;
let now = db::now();
let stale = ctx.db.stale_unsubscribed(now - DEAD, now - QUIET).await?;
if stale.is_empty() {
return Ok(());
}
let mut cfg = (*ctx.cfg()).clone();
if stale.iter().fold(false, |any, id| cfg.feeds.remove(id).is_some() || any) {
ctx.store_cfg(cfg).await?;
}
for id in &stale {
retire_group(ctx, id).await?;
ctx.db.forget_feed(id).await?;
tracing::info!(feed = id, "removed a feed nobody subscribes to that is dead or has published nothing in a year");
}
Ok(())
}
/// `standalone` false means this is the sweep that runs before a scan: it reports what it
/// deleted, but must not emit the terminal ReapDone, or a client waiting on its `fetch`
/// would stop reading before the scan had even started.
#[tracing::instrument(name = "reap", skip_all, fields(dry_run = dry_run))]
async fn reap(ctx: &Ctx, dry_run: bool, standalone: bool) -> Result<()> {
let r = retention::run(&ctx.cfg(), &ctx.db, dry_run).await?;
if !dry_run {
clean_directory(ctx).await?;
let gone = art::trim(&art::dir(), ctx.cfg().general.art_cache_mb * 1_048_576);
if gone > 0 {
tracing::debug!(gone, "trimmed the artwork kept on disk");
}
}
for c in r.aged_out.iter().chain(r.over_quota.iter()) {
ctx.out.emit(Event::Reaped {
path: c.path.clone(),
bytes: c.bytes.max(0) as u64,
});
}
if standalone {
ctx.out.emit(Event::ReapDone {
files: r.aged_out.len() + r.over_quota.len(),
bytes: r.bytes_freed,
});
}
Ok(())
}
/// `scope`, when not empty, narrows the scan to those feeds and the feeds inside any of them.
#[tracing::instrument(name = "scan", skip_all, fields(only = ?only, force = force))]
async fn fetch(ctx: &Arc<Ctx>, only: Option<&str>, force: bool, scope: &[String]) -> Result<()> {
let cfg = ctx.cfg();
let subs = subscriptions(ctx).await?;
if let Some(id) = only
&& !subs.iter().any(|s| s.id == id)
{
anyhow::bail!("no feed with id {id:?}");
}
let subscribed = ctx.db.subscriber_counts().await?;
let mut scanned = 0;
let mut fresh: Vec<String> = vec![];
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 = 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 {
state.etag = None;
state.last_modified = None;
}
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(),
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.feed_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, fetched).await {
Ok(Outcome::Feed(s)) => ctx.out.emit(Event::FeedDone {
feed: id.clone(),
new: s.new_entries,
downloaded: s.downloaded,
failed: s.failed,
torrents: s.torrents,
}),
Ok(Outcome::NotModified) => ctx.out.emit(Event::FeedSkip {
feed: id.clone(),
reason: "not modified".into(),
}),
Ok(Outcome::Empty) => ctx.out.emit(Event::FeedSkip {
feed: id.clone(),
reason: "nothing yet".into(),
}),
Ok(Outcome::Opml { added, removed, kept, total }) => {
ctx.out.emit(Event::FeedSkip {
feed: id.clone(),
reason: format!(
"{total} feed(s) listed, {} added, {removed} unsubscribed, {kept} kept without a listing",
added.len()
),
});
// Read them in the same pass, as the original did, rather than making
// the user wait a whole interval for a newly listed show.
fresh.extend(added);
}
Err(e) => {
// One bad feed must not end the scan.
let msg = format!("{e:#}");
ctx.out.emit(Event::FeedError { feed: id.clone(), msg: msg.clone() });
ctx.db.set_feed_error(id, &feed_cfg.url, &msg).await?;
}
}
}
// Feeds a subscribed OPML just introduced: scan them now, in this run.
if !fresh.is_empty() {
let subs = subscriptions(ctx).await?;
for id in &fresh {
let Some(feed_cfg) = subs.iter().find(|s| &s.id == id).map(|s| &s.cfg) else {
continue;
};
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, None).await {
Ok(Outcome::Feed(s)) => ctx.out.emit(Event::FeedDone {
feed: id.clone(),
new: s.new_entries,
downloaded: s.downloaded,
failed: s.failed,
torrents: s.torrents,
}),
Ok(_) => {}
Err(e) => {
let msg = format!("{e:#}");
ctx.out.emit(Event::FeedError { feed: id.clone(), msg: msg.clone() });
ctx.db.set_feed_error(id, &feed_cfg.url, &msg).await?;
}
}
}
}
ctx.out.emit(Event::ScanDone { feeds: scanned });
Ok(())
}
/// One feed to scan: either an entry you wrote in config.toml, or one derived from an
/// OPML subscription and held only in the database.
pub struct Sub {
pub id: String,
pub cfg: config::Feed,
/// True when it came from an OPML and has no config entry of its own.
pub managed: bool,
}
/// Everything to scan: your config entries, plus whatever the OPML subscriptions listed.
///
/// A derived feed borrows its parent's settings wholesale. That is why it needs no config
/// entry -- there is nothing to store but its URL and where it came from.
pub async fn subscriptions(ctx: &Ctx) -> Result<Vec<Sub>> {
let cfg = ctx.cfg();
let mut out: Vec<Sub> = cfg
.feeds
.iter()
.map(|(id, f)| Sub { id: id.clone(), cfg: f.clone(), managed: false })
.collect();
for m in ctx.db.managed_feeds().await? {
if cfg.feeds.contains_key(&m.id) {
continue; // promoted to config at some point; that entry wins
}
let parent = cfg.feeds.get(&m.group_id);
if parent.is_none() {
// The OPML or Patreon feed this was derived from is no longer in config --
// removing it should have retired these rows too (see `retire_group`), but
// skip them here regardless so a row that slips through is never scanned.
continue;
}
let base = match parent.and_then(|p| p.folder.clone()) {
Some(folder) => folder,
None => ctx.db.feed_summary(&m.group_id).await.ok().and_then(|s| s.title).unwrap_or_else(|| m.group_id.clone()),
};
let title = m.title.clone().unwrap_or_else(|| m.id.clone());
out.push(Sub {
id: m.id.clone(),
cfg: config::Feed {
url: m.url.clone(),
folder: Some(format!("{base}/{title}")),
group: Some(m.group_id.clone()),
media_types: parent.and_then(|p| p.media_types.clone()),
schedule: parent.and_then(|p| p.schedule.clone()),
keywords: parent.map(|p| p.keywords.clone()).unwrap_or_default(),
allow_explicit: parent.is_some_and(|p| p.allow_explicit),
auto_download: parent.is_none_or(|p| p.auto_download),
max_new_per_check: parent.and_then(|p| p.max_new_per_check),
username: parent.and_then(|p| p.username.clone()),
password: parent.and_then(|p| p.password.clone()),
password_env: parent.and_then(|p| p.password_env.clone()),
category: None,
listed: false,
},
managed: true,
});
}
out.sort_by(|a, b| a.id.cmp(&b.id));
Ok(out)
}
/// Retires every feed derived from `parent_id`, now that nothing subscribes to the OPML or
/// Patreon feed that listed them: the same rule `sync_group` applies to one the list drops --
/// removed if nothing was downloaded, orphaned and kept otherwise. Called once the parent
/// itself is removed, since `subscriptions()` would otherwise keep scanning them under a
/// fallback policy meant for a feed with no parent at all. A feed promoted to config is not
/// derived any more, so it is only unmanaged.
pub async fn retire_group(ctx: &Ctx, parent_id: &str) -> Result<()> {
let cfg = ctx.cfg();
for m in ctx.db.managed_feeds().await?.into_iter().filter(|m| m.group_id == parent_id) {
if cfg.feeds.contains_key(&m.id) {
// Scanned from its config entry and still read. Dropped as derived, its stored
// entries would go with it: davewiner's 11 were promoted without being unmanaged.
ctx.db.unmanage(&m.id).await?;
} else if ctx.db.downloaded_count(&m.id).await.unwrap_or(1) > 0 {
ctx.db.set_orphaned(&m.id, true).await?;
} else {
ctx.db.drop_managed(&m.id).await?;
}
}
Ok(())
}
/// Retires every group whose parent is gone from config, and returns how many derived rows that
/// dropped or unmanaged. An OPML removed before `retire_group` existed left its feeds behind:
/// davewiner's 922 were skipped by every scan and never cleared, and their stale errors were
/// most of the ones stored.
async fn retire_stranded(ctx: &Ctx) -> Result<usize> {
let cfg = ctx.cfg();
let before = ctx.db.managed_feeds().await?;
let stranded: std::collections::BTreeSet<&str> = before
.iter()
.map(|m| m.group_id.as_str())
.filter(|g| !cfg.feeds.contains_key(*g))
.collect();
for group in stranded {
retire_group(ctx, group).await?;
}
Ok(before.len() - ctx.db.managed_feeds().await?.len())
}
/// 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
/// 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
/// its run of failures began, to `last_checked`, when it last failed.
pub fn due_after(
cfg: &config::Config,
feed: &config::Feed,
ttl_mins: Option<u64>,
error_since: Option<i64>,
last_checked: i64,
) -> u64 {
let mins = match feed.schedule.as_deref().and_then(config::parse_interval) {
Some(explicit) => explicit,
None => ttl_mins.unwrap_or(0).max(cfg.general.interval()),
};
let failing_for = error_since.map(|since| (last_checked - since).max(0) as u64);
backoff(mins * 60, failing_for)
}
/// A failing feed waits as long as it has been failing, so the wait doubles with each failure
/// (1h, 1h, 2h, 4h, ... on an hourly schedule), never less than its usual interval and, past that,
/// never more than a day. Retried hourly, a feed dead for good cost a request, a warning and scan
/// time every hour (#99). The first success clears error_since, and with it the backoff; a
/// refresh someone asks for is not held back by it.
fn backoff(usual: u64, failing_for: Option<u64>) -> u64 {
const CEILING: u64 = 86_400;
match failing_for {
Some(f) => usual.max(f.min(CEILING)),
None => usual,
}
}
#[derive(Default)]
struct Scan {
new_entries: usize,
downloaded: usize,
failed: usize,
torrents: usize,
}
/// What a scan of one feed turned out to be.
enum Outcome {
NotModified,
/// A response with nothing in it -- the British Antarctic Survey answers a 202 with an
/// empty body when it has nothing new to publish. Not a parse failure; try again later.
Empty,
Feed(Scan),
/// The URL is a list of feeds rather than a feed: an OPML, or a Patreon creator's shows.
Opml { added: Vec<String>, removed: usize, kept: usize, total: usize },
}
/// A feed that answered from a new address after permanent redirects moves there in the
/// catalogue, so it is read from there and no longer redirected every time. One an OPML lists
/// is left alone, since the OPML would only put the old address back; so is a move onto an
/// address another feed already has.
async fn follow_move(ctx: &Ctx, id: &str, to: &str) -> Result<()> {
// Said in the event below, not logged here: the event is the log line.
if subscriptions(ctx).await?.iter().any(|s| s.id != id && feed::same_feed(&s.cfg.url, to)) {
tracing::info!(feed = id, to, "the feed moved to an address another feed already has; leaving it");
return Ok(());
}
let mut cfg = (*ctx.cfg()).clone();
let Some(f) = cfg.feeds.get_mut(id) else { return Ok(()) };
let from = std::mem::replace(&mut f.url, to.to_owned());
ctx.store_cfg(cfg).await?;
ctx.out.emit(Event::FeedMoved { feed: id.to_owned(), from, to: to.to_owned() });
Ok(())
}
#[tracing::instrument(name = "feed", skip_all, fields(feed = id))]
async fn scan_one(
ctx: &Arc<Ctx>,
id: &str,
feed_cfg: &config::Feed,
state: &db::HttpState,
force: bool,
// The body the scan already fetched for it, if it did (#104), and where it moved, if it did.
prefetched: Option<Result<(feed::Fetched, Option<String>)>>,
) -> Result<Outcome> {
// The queue follows the feed's settings each time it is due, changed or not, so files a
// limit or auto-download no longer reaches stop counting as waiting (#113).
let policy = policy_for(ctx, id, feed_cfg).await?;
ctx.db.hold_back(id, if policy.auto_download { policy.budget } else { 0 }).await?;
// A Patreon creator with more than one show is a list of feeds, like an OPML.
if feed::is_patreon_creator(&feed_cfg.url) {
match feed::patreon_shows(&ctx.client, &feed_cfg.url).await {
Ok((name, shows)) if shows.len() > 1 => {
ctx.db.touch_feed(id, &feed_cfg.url).await?;
if let Some(name) = name {
ctx.db.set_title(id, &name).await?;
}
// Read as one feed before it was split, it listed every show's items in one
// heap. The items go; its files and read state move to each show as the show
// lists them (`Db::adopt`), so no show comes up empty for want of a URL.
ctx.db.clear_entries(id).await?;
return sync_group(ctx, id, feed_cfg, &shows).await;
}
Ok(_) => {} // One show: the creator's feed is that show.
// Already split: keep the shows it has rather than read the creator as one heap.
Err(e) if ctx.db.managed_feeds().await?.iter().any(|m| m.group_id == id) => return Err(e),
Err(e) => tracing::warn!(
feed = id,
error = %format!("{e:#}"),
"could not list the Patreon shows; reading it as one feed"
),
}
}
let (mut fetched, moved) = match prefetched {
Some(got) => got?,
None => feed::fetch(&ctx.feed_client, feed_cfg, state.etag.as_deref(), state.last_modified.as_deref()).await?,
};
if let Some(to) = moved {
follow_move(ctx, id, &to).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
// the validator: drop it and ask again, or the feed stays empty until the publisher
// happens to change something. The same for artwork never looked for: a feed from before
// site icons (#73) would otherwise wait for its next post to get one.
let stored = ctx.db.feed_summary(id).await?;
if matches!(fetched, feed::Fetched::NotModified) && (stored.entries == 0 || stored.image.is_none()) {
tracing::info!(feed = id, "not modified, but nothing stored; refetching without the validator");
ctx.db.clear_validators(id).await?;
fetched = feed::fetch(&ctx.feed_client, feed_cfg, None, None).await?.0;
}
let (bytes, etag, last_modified) = match fetched {
feed::Fetched::NotModified => {
ctx.db.touch_feed(id, &feed_cfg.url).await?;
return Ok(Outcome::NotModified);
}
feed::Fetched::Body { bytes, etag, last_modified } => (bytes, etag, last_modified),
};
if bytes.iter().all(u8::is_ascii_whitespace) {
ctx.db.touch_feed(id, &feed_cfg.url).await?;
return Ok(Outcome::Empty);
}
// A subscribed OPML is a list of feeds, not a feed. The original matched on a ".opml"
// URL; sniffing the body also catches one served from a URL without that extension.
if feed::is_opml(&bytes) {
ctx.db.touch_feed(id, &feed_cfg.url).await?;
return sync_opml(ctx, id, feed_cfg, &bytes).await;
}
let mut parsed = feed::parse(&bytes)?;
// Artwork on http moves to https where its host serves it (#110), before it is compared
// with what is stored, which is then the https address too.
let mut secure = std::collections::HashMap::new();
if let Some(art) = &parsed.image {
parsed.image = Some(feed::prefer_https(&ctx.client, art, &mut secure).await);
}
for entry in parsed.entries.iter_mut() {
if let Some(art) = &entry.image {
entry.image = Some(feed::prefer_https(&ctx.client, art, &mut secure).await);
}
}
// Artwork is looked at when it may have changed: the feed names different artwork from what
// is stored, or someone asked for a refresh, so an icon the site changes or fixes still
// follows it (#80). Looked at on every full read, a feed without validators asked its site
// on every scan: 2 s a scan for lfg.co (#95). A miss is stored as "", which the page draws
// as no art, and which stops the refetch above.
// A feed's own artwork has to be there too: Ken and Robin's names a 404, and stored unasked
// it stood in the way of the site's icon, which works (#89).
let art_span = tracing::info_span!("artwork");
tracing::Instrument::instrument(async {
if let Some(art) = &parsed.image
&& stored.image.as_deref() != Some(art.as_str())
&& !feed::is_image(&ctx.client, art).await
{
tracing::info!(feed = id, art = %art, "the feed's artwork is not an image; trying its site's icon");
parsed.image = None;
}
if parsed.image.is_none() {
parsed.image = Some(match (&stored.image, &parsed.site) {
(Some(known), _) if !force => known.clone(),
(_, Some(site)) => feed::site_icon(&ctx.client, site).await.unwrap_or_default(),
(_, None) => String::new(),
});
}
}, art_span).await;
// A site's icon comes back on whatever the site is on.
if let Some(art) = &parsed.image {
parsed.image = Some(feed::prefer_https(&ctx.client, art, &mut secure).await);
}
// Items stored before, which a scan does not write again, move with their host, and a host
// only they still name is asked too.
for art in ctx.db.http_images(id).await? {
feed::prefer_https(&ctx.client, &art, &mut secure).await;
}
for (host, ok) in &secure {
if *ok {
ctx.db.secure_images(id, host).await?;
}
}
ctx.db.record_feed(
id,
&feed_cfg.url,
parsed.title.as_deref(),
etag.as_deref(),
last_modified.as_deref(),
parsed.ttl_mins,
parsed.image.as_deref(),
parsed.category.as_deref(),
).await?;
if let Some(parent) = &feed_cfg.group {
let listed: Vec<(&str, &str)> = parsed
.entries
.iter()
.flat_map(|e| e.enclosures.iter().map(move |x| (e.guid.as_str(), x.url.as_str())))
.collect();
ctx.db.adopt(parent, id, &listed).await?;
}
// Verdicts are recorded in `state`, so the download queue below is just "everything still
// pending". A filter's verdict is looked at again on every scan, though: made once, at
// discovery, it outlived the setting behind it, and allowing explicit items afterwards
// changed nothing however often the feed was scanned.
let skipped = ctx.db.skipped_by_filter(id).await?;
let (known_items, known_files) = ctx.db.stored_items(id).await?;
let mut scan = Scan::default();
// Its own span: the time a feed spends after its fetch was untraced (#96).
let store = tracing::info_span!("store", items = parsed.entries.len());
tracing::Instrument::instrument(async {
for entry in &parsed.entries {
// Only what is not stored yet is inserted; the insert would find the rest and do nothing.
if !known_items.contains(&entry.guid) && ctx.db.record_entry(id, entry).await? {
scan.new_entries += 1;
}
for enc in &entry.enclosures {
// A URL not among this feed's files may still be another feed's: the insert says.
let was = if !known_files.contains(&enc.url) && ctx.db.record_enclosure(id, &entry.guid, enc).await? {
None
} else if let Some(reason) = skipped.get(&enc.url) {
Some(reason.as_str())
} else {
continue; // Settled: queued, downloaded, reaped, or another feed's file.
};
let now = reject(&ctx.cfg(), feed_cfg, &policy, entry, enc);
if now != was {
match now {
Some(reason) => ctx.db.mark_enclosure(&enc.url, "skipped", Some(reason)).await?,
None => ctx.db.mark_enclosure(&enc.url, "pending", None).await?,
}
}
}
}
anyhow::Ok(())
}, store).await?;
if scan.new_entries > 0 {
ctx.db.rehide(id).await?; // what is new may hold someone's blocked words
}
// A new item takes a place among the newest and the oldest of them leaves the queue; a file a
// filter lets through again may be outside them.
ctx.db.hold_back(id, if policy.auto_download { policy.budget } else { 0 }).await?;
let budget = policy.budget;
if policy.auto_download && budget > 0 {
let cfg = ctx.cfg();
let folder = download::folder_for(&cfg, id, feed_cfg, parsed.title.as_deref());
let dest_dir = cfg.general.download_dir.join(&folder);
for item in ctx.db.pending(id, budget).await? {
if download::looks_like_torrent(&item.url, item.mime.as_deref()) {
if !ctx.cfg().torrent.enabled {
ctx.db.mark_enclosure(&item.url, "skipped", Some("torrents disabled")).await?;
ctx.out.emit(Event::TorrentDeferred {
feed: id.to_string(),
url: item.url.clone(),
});
scan.torrents += 1;
continue;
}
if ctx.detach_torrents {
// 'downloading' keeps the next scan from queueing it a second time.
ctx.db.mark_enclosure(&item.url, "downloading", None).await?;
spawn_torrent(ctx, id.to_string(), item.id, item.url.clone(), dest_dir.clone());
scan.torrents += 1;
continue;
}
match torrent_one(ctx, id, item.id, &item.url, &dest_dir).await {
Ok((path, bytes)) => {
ctx.db.mark_downloaded(&item.url, &path, bytes).await?;
ctx.out.emit(Event::DownloadDone {
feed: id.to_string(),
enclosure: item.id,
url: item.url.clone(),
path: path.display().to_string(),
bytes,
});
scan.downloaded += 1;
}
Err(e) => {
let msg = format!("{e:#}");
ctx.out.emit(Event::DownloadError {
feed: id.to_string(),
enclosure: item.id,
url: item.url.clone(),
msg: msg.clone(),
});
ctx.db.mark_enclosure(&item.url, "error", Some(&msg)).await?;
scan.failed += 1;
}
}
continue;
}
match fetch_one(ctx, id, item.id, feed_cfg, &item.url, &dest_dir).await {
Ok((path, bytes)) => {
ctx.out.emit(Event::DownloadDone {
feed: id.to_string(),
enclosure: item.id,
url: item.url.clone(),
path: path.display().to_string(),
bytes,
});
scan.downloaded += 1;
}
Err(e) => {
let msg = format!("{e:#}");
ctx.out.emit(Event::DownloadError {
feed: id.to_string(),
enclosure: item.id,
url: item.url.clone(),
msg: msg.clone(),
});
ctx.db.mark_enclosure(&item.url, "error", Some(&msg)).await?;
scan.failed += 1;
}
}
}
}
Ok(Outcome::Feed(scan))
}
async fn sync_opml(
ctx: &Arc<Ctx>,
parent_id: &str,
parent: &config::Feed,
bytes: &[u8],
) -> Result<Outcome> {
let listed = feed::parse_opml(bytes)?;
if let Some(title) = feed::opml_title(bytes) {
ctx.db.set_title(parent_id, &title).await?;
}
sync_group(ctx, parent_id, parent, &listed).await
}
/// Brings the feed list in step with a list of feeds: a subscribed OPML, or a Patreon
/// creator's shows.
///
/// New entries are added under the list's group and folder. An entry that has gone from
/// the list is unsubscribed *only if nothing was ever downloaded for it* -- otherwise it
/// is kept and flagged, because dropping it would orphan files on disk with nothing in
/// the UI to explain them.
async fn sync_group(
ctx: &Arc<Ctx>,
parent_id: &str,
parent: &config::Feed,
listed: &[(String, String)],
) -> Result<Outcome> {
let cfg = ctx.cfg();
let existing = ctx.db.managed_feeds().await?;
let mut added = vec![];
for (title, url) in listed {
// Already known, whether derived or promoted into the config.
if let Some(m) = existing.iter().find(|m| &m.url == url) {
ctx.db.upsert_managed(&m.id, url, title, parent_id).await?;
continue;
}
// A Patreon show you added by hand may be spelled differently from the one listed.
if cfg.feeds.values().any(|f| feed::same_feed(&f.url, url)) {
continue;
}
// A removed feed keeps its rows, so its id is only free again for the same feed.
let known = ctx.db.feed_urls().await?;
let taken: std::collections::BTreeMap<String, config::Feed> = cfg
.feeds
.keys()
.chain(existing.iter().map(|m| &m.id))
.chain(added.iter())
.chain(known.iter().filter(|(_, u)| !feed::same_feed(u, url)).map(|(id, _)| id))
.map(|id| (id.clone(), parent.clone()))
.collect();
let id = config::unique_slug(title, &taken);
ctx.db.upsert_managed(&id, url, title, parent_id).await?;
added.push(id);
}
// Whoever subscribes to the OPML subscribes to what it lists: that is what taking a
// subscription means. Their own feeds are untouched.
for id in ctx
.db
.managed_feeds().await?
.iter()
.filter(|m| m.group_id == parent_id)
.map(|m| m.id.clone())
.chain(std::iter::once(parent_id.to_string()))
{
for user in ctx.db.users().await? {
if ctx.db.subscription(user.id, parent_id).await?.is_some() {
ctx.db.subscribe(user.id, &id).await?;
}
}
}
// Anything in this group the OPML no longer lists.
let mut removed = 0;
let mut kept = 0;
for m in existing.iter().filter(|m| m.group_id == parent_id) {
if listed.iter().any(|(_, u)| u == &m.url) {
continue;
}
if ctx.db.downloaded_count(&m.id).await.unwrap_or(1) > 0 {
// Never orphan a downloaded file: keep the feed and say why in the UI.
ctx.db.set_orphaned(&m.id, true).await?;
kept += 1;
tracing::info!(feed = %m.id, "dropped from the OPML but has downloads; keeping it");
} else {
ctx.db.drop_managed(&m.id).await?;
removed += 1;
tracing::info!(feed = %m.id, "dropped from the OPML with nothing downloaded; removed");
}
}
Ok(Outcome::Opml { added, removed, kept, total: listed.len() })
}
/// Why this enclosure should not be downloaded, if it should not be.
fn reject(
cfg: &config::Config,
feed_cfg: &config::Feed,
policy: &Policy,
entry: &feed::Entry,
enc: &feed::Enclosure,
) -> Option<&'static str> {
let url = enc.url.as_str();
if !policy.auto_download {
return Some("auto_download is off");
}
// Blog feeds put the article's header image in an <enclosure>; without this a text
// feed reads as a podcast full of episodes and fills the disk with artwork.
let wanted = feed_cfg
.media_types
.as_deref()
.unwrap_or(&cfg.general.media_types);
if !config::wanted_media(enc.mime.as_deref(), wanted) {
return Some("not audio or video");
}
if entry.explicit && !policy.allow_explicit {
return Some("explicit");
}
let categories = entry.categories.join(" ");
let haystacks = [
url,
entry.title.as_deref().unwrap_or(""),
entry.description.as_deref().unwrap_or(""),
categories.as_str(),
];
let text = [entry.title.as_deref().unwrap_or(""), entry.description.as_deref().unwrap_or("")];
if !policy.wanted(&haystacks, &text) {
// Worded for whichever filter could have let it through, the keywords if none blocked it.
return Some(if policy.wanted(&haystacks, &[]) { "blocked word" } else { "no keyword match" });
}
None
}
/// What the scanner should do for a feed, merged across everyone subscribed to it. The
/// feed is fetched once and its files are downloaded once, so the merge is a union: if
/// one person wants a thing, it is fetched, and everyone else simply sees it listed.
///
/// With no subscribers at all -- a hand-written config entry nobody has claimed yet --
/// the feed's own settings stand, which is how a single-user install behaves.
pub struct Policy {
pub auto_download: bool,
pub allow_explicit: bool,
/// What each subscriber fetching the feed wants.
pub wants: Vec<Want>,
pub budget: usize,
}
/// One person's filters: an item is theirs if it matches their keywords (all of it, with none)
/// and none of their blocked words.
pub struct Want {
pub keywords: Vec<String>,
pub blocked: Vec<String>,
}
impl Policy {
/// One file serves everyone subscribed, so an item is wanted if anyone wants it. Keywords
/// look at `haystacks`, blocked words at `text`, the item's title and body only.
fn wanted(&self, haystacks: &[&str], text: &[&str]) -> bool {
self.wants.iter().any(|w| {
download::matches_keywords(&w.keywords, haystacks) && !download::blocked(&w.blocked, text)
})
}
}
async fn policy_for(ctx: &Ctx, id: &str, feed_cfg: &config::Feed) -> Result<Policy> {
let global = ctx.cfg().general.max_new_per_check;
Ok(merge_policy(&ctx.db.subscribers(id, feed_cfg.group.as_deref()).await?, feed_cfg, global))
}
fn merge_policy(subs: &[db::Sub], feed_cfg: &config::Feed, global: usize) -> Policy {
let cap = |n: Option<usize>| n.unwrap_or(if global == 0 { usize::MAX } else { global });
if subs.is_empty() {
return Policy {
// A feed listed in the Directory with nobody subscribed is there to be found, not
// downloaded: its files would be for no one. Once someone subscribes, the files
// skipped for it are judged again on the next scan, by their settings (#107).
auto_download: feed_cfg.auto_download && !feed_cfg.listed,
allow_explicit: feed_cfg.allow_explicit,
wants: vec![Want { keywords: feed_cfg.keywords.clone(), blocked: vec![] }],
budget: cap(feed_cfg.max_new_per_check),
};
}
let mut policy = Policy {
auto_download: false,
allow_explicit: false,
wants: vec![],
budget: 0,
};
for sub in subs {
if !sub.auto_download.unwrap_or(feed_cfg.auto_download) {
continue; // Not fetching for this person, so their wants add nothing.
}
policy.auto_download = true;
policy.allow_explicit |= sub.allow_explicit.unwrap_or(feed_cfg.allow_explicit);
policy.budget = policy
.budget
.max(cap(sub.max_new_per_check.map(|n| n as usize).or(feed_cfg.max_new_per_check)));
policy.wants.push(Want {
keywords: sub.keywords.clone().unwrap_or_else(|| feed_cfg.keywords.clone()),
blocked: sub.blocked.clone(),
});
}
policy
}
#[tracing::instrument(name = "download", skip_all, fields(feed = feed_id, enclosure = enclosure, url = url))]
async fn fetch_one(
ctx: &Arc<Ctx>,
feed_id: &str,
enclosure: i64,
feed_cfg: &config::Feed,
url: &str,
dest_dir: &std::path::Path,
) -> Result<(PathBuf, u64)> {
// Throttled to whole percents, as the original's lastDLStepSize guard did.
let mut last_pct = -1i64;
let name = download::filename_for(url, None);
let cfg = ctx.cfg();
let got = download::download(&ctx.client, &cfg, feed_cfg, url, |done, total| {
if let Some(t) = total.filter(|t| *t > 0) {
let pct = (done * 100 / t) as i64;
if pct > last_pct {
last_pct = pct;
ctx.out.emit(Event::Progress {
feed: feed_id.to_string(),
enclosure,
url: url.to_string(),
file: name.clone(),
done,
total,
});
}
}
})
.await?;
if matches!(got.kind, download::Kind::Torrent) {
// The MIME lied. Hand the URL to the torrent session instead of filing a .torrent
// as if it were an episode.
let _ = tokio::fs::remove_file(&got.tmp).await;
if !ctx.cfg().torrent.enabled {
ctx.db.mark_enclosure(url, "skipped", Some("torrents disabled")).await?;
anyhow::bail!("body is a torrent and torrents are disabled");
}
return torrent_one(ctx, feed_id, enclosure, url, dest_dir).await;
}
let path = download::place(&got, dest_dir).await?;
ctx.db.mark_downloaded(url, &path, got.bytes).await?;
Ok((path, got.bytes))
}
/// Runs a torrent off the command worker, so feed scans and HTTP downloads keep moving
/// while it fetches metadata, transfers and then seeds.
fn spawn_torrent(ctx: &Arc<Ctx>, feed_id: String, enclosure: i64, url: String, dest_dir: PathBuf) {
let ctx = ctx.clone();
tokio::spawn(async move {
// Held for the whole job, so a feed full of torrents cannot open hundreds at once.
let _permit = match ctx.torrent_slots.clone().acquire_owned().await {
Ok(p) => p,
Err(_) => return,
};
let outcome = torrent_one(&ctx, &feed_id, enclosure, &url, &dest_dir).await;
let db = &ctx.db;
match outcome {
Ok((path, bytes)) => {
if let Err(e) = db.mark_downloaded(&url, &path, bytes).await {
tracing::warn!(error = %format!("{e:#}"), "could not record the finished torrent");
}
ctx.out.emit(Event::DownloadDone {
feed: feed_id,
enclosure,
url,
path: path.display().to_string(),
bytes,
});
}
Err(e) => {
let msg = format!("{e:#}");
let _ = db.mark_enclosure(&url, "error", Some(&msg)).await;
ctx.out.emit(Event::DownloadError { feed: feed_id, enclosure, url, msg });
}
}
});
}
/// Downloads one specific enclosure immediately, whatever the per-scan cap says and
/// wherever it sits in the queue.
#[tracing::instrument(skip(ctx))]
async fn download_one(ctx: &Arc<Ctx>, id: i64) -> Result<()> {
let cfg = ctx.cfg();
let enc = ctx
.db
.enclosure(id).await?
.ok_or_else(|| anyhow::anyhow!("no enclosure {id}"))?;
if enc.path.is_some() {
return Ok(()); // Already here.
}
// Must look through the derived feeds too: anything inside an OPML subscription has
// no config entry, so a config-only lookup called every one of them "unsubscribed".
let subs = subscriptions(ctx).await?;
let feed_cfg = subs
.iter()
.find(|s| s.id == enc.feed_id)
.map(|s| s.cfg.clone())
.ok_or_else(|| anyhow::anyhow!("enclosure {id} belongs to unsubscribed feed {:?}", enc.feed_id))?;
let feed_cfg = &feed_cfg;
let title = ctx.db.feed_summary(&enc.feed_id).await?.title;
let folder = download::folder_for(&cfg, &enc.feed_id, feed_cfg, title.as_deref());
let dest_dir = cfg.general.download_dir.join(&folder);
ctx.out.emit(Event::FeedStart { feed: enc.feed_id.clone() });
let is_torrent = download::looks_like_torrent(&enc.url, enc.mime.as_deref());
if is_torrent && cfg.torrent.enabled && ctx.detach_torrents {
ctx.db.mark_enclosure(&enc.url, "downloading", None).await?;
spawn_torrent(ctx, enc.feed_id.clone(), enc.id, enc.url.clone(), dest_dir);
return Ok(());
}
let result = if is_torrent {
if !cfg.torrent.enabled {
Err(anyhow::anyhow!("torrents are disabled"))
} else {
torrent_one(ctx, &enc.feed_id, enc.id, &enc.url, &dest_dir).await
}
} else {
fetch_one(ctx, &enc.feed_id, enc.id, feed_cfg, &enc.url, &dest_dir).await
};
match result {
Ok((path, bytes)) => {
ctx.db.mark_downloaded(&enc.url, &path, bytes).await?;
ctx.out.emit(Event::DownloadDone {
feed: enc.feed_id.clone(),
enclosure: enc.id,
url: enc.url.clone(),
path: path.display().to_string(),
bytes,
});
}
Err(e) => {
let msg = format!("{e:#}");
ctx.db.mark_enclosure(&enc.url, "error", Some(&msg)).await?;
ctx.out.emit(Event::DownloadError {
feed: enc.feed_id.clone(),
enclosure: enc.id,
url: enc.url.clone(),
msg,
});
}
}
// Terminal, so a UI waiting on this request stops here.
ctx.out.emit(Event::ScanDone { feeds: 1 });
Ok(())
}
/// Torrent progress is reported the same way an HTTP download's is, throttled to whole
/// percents so a UI is not flooded.
#[tracing::instrument(name = "torrent", skip_all, fields(feed = feed_id, enclosure = enclosure))]
async fn torrent_one(
ctx: &Arc<Ctx>,
feed_id: &str,
enclosure: i64,
url: &str,
dest_dir: &std::path::Path,
) -> Result<(PathBuf, u64)> {
let name = download::filename_for(url, None);
let mut last_pct = -1i64;
let cfg = ctx.cfg();
ctx.torrents()
.await?
.fetch(&cfg, url, dest_dir, |done, total| {
if total > 0 {
let pct = (done * 100 / total) as i64;
if pct > last_pct {
last_pct = pct;
ctx.out.emit(Event::Progress {
feed: feed_id.to_string(),
enclosure,
url: url.to_string(),
file: name.clone(),
done,
total: Some(total),
});
}
}
})
.await
}
fn ago(t: Option<i64>) -> String {
let Some(t) = t else { return "never".into() };
format!("{} ago", duration((db::now() - t).max(0) as u64))
}
fn duration(secs: u64) -> String {
match secs {
s if s < 90 => format!("{s}s"),
s if s < 5400 => format!("{}m", s / 60),
s if s < 172_800 => format!("{}h", s / 3600),
s => format!("{}d", s / 86_400),
}
}
#[cfg(test)]
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<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]
fn a_listed_feed_nobody_subscribes_to_downloads_nothing() {
let mut f = feed();
assert!(merge_policy(&[], &f, 3).auto_download);
f.listed = true;
assert!(!merge_policy(&[], &f, 3).auto_download);
}
#[test]
fn a_failing_feed_backs_off_doubling_up_to_a_day() {
let hour = 3600;
assert_eq!(backoff(hour, None), hour);
// Failing since 0: each check waits as long as the failure has lasted so far.
let (mut at, mut waits) = (0, vec![]);
for _ in 0..8 {
let w = backoff(hour, Some(at));
waits.push(w / hour);
at += w;
}
assert_eq!(waits, [1, 1, 2, 4, 8, 16, 24, 24]);
// A weekly schedule is longer than the ceiling and stays as it is.
assert_eq!(backoff(7 * 86_400, Some(30 * 86_400)), 7 * 86_400);
}
#[tokio::test]
async fn the_first_start_moves_the_configuration_in_and_trims_the_file() {
let dir = std::env::temp_dir().join(format!("ipx-assemble-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("config.toml");
std::fs::write(
&path,
"[general]\nschedule = \"every 2h\"\n[web]\ntoken = \"t\"\n[feeds.show]\nurl = \"http://x/show.xml\"\n",
)
.unwrap();
let db = db::Db::memory().await.unwrap();
let first = assemble_config(&db, config::Config::load(&path).unwrap(), &path).await.unwrap();
assert_eq!(first.feeds.keys().collect::<Vec<_>>(), ["show"]);
assert_eq!(first.general.schedule, "every 2h");
assert!(dir.join("config.toml.pre-database").exists(), "the original is kept");
assert!(!config::Config::file_holds_stored(&path), "and the file no longer lists them");
// The next start reads them from the database, the trimmed file notwithstanding.
let next = assemble_config(&db, config::Config::load(&path).unwrap(), &path).await.unwrap();
assert_eq!(next.feeds.keys().collect::<Vec<_>>(), ["show"]);
assert_eq!(next.general.schedule, "every 2h");
assert_eq!(next.web.token, "t");
std::fs::remove_dir_all(&dir).unwrap();
}
fn feed() -> config::Feed {
// Whatever `ipx add` would write, which is the shape every code path sees.
let mut cfg = config::Config::default();
let f = add_one_cfg(&mut cfg, "http://x/f.xml", None, vec![]);
f
}
/// The feed entry `add` builds, without the network round trip it does for a title.
fn add_one_cfg(
_cfg: &mut config::Config,
url: &str,
folder: Option<String>,
keywords: Vec<String>,
) -> config::Feed {
config::Feed {
url: url.into(),
folder,
keywords,
allow_explicit: false,
auto_download: true,
group: None,
media_types: None,
schedule: None,
max_new_per_check: None,
username: None,
password: None,
password_env: None,
category: None,
listed: false,
}
}
fn sub(kw: Option<&[&str]>, auto: Option<bool>, max: Option<i64>) -> db::Sub {
db::Sub {
feed_id: "f".into(),
keywords: kw.map(|k| k.iter().map(|s| s.to_string()).collect()),
auto_download: auto,
allow_explicit: None,
max_new_per_check: max,
blocked: vec![],
}
}
#[tokio::test]
async fn a_shared_feed_is_fetched_for_whoever_wants_the_most() {
// Nobody subscribed: the feed's own settings stand, as in a single-user install.
let p = merge_policy(&[], &feed(), 3);
assert!(p.auto_download);
assert_eq!(p.budget, 3);
assert!(p.wanted(&["anything"], &[]));
// Two filters: an item wanted by either of them is fetched, since one file serves
// both. The larger per-scan cap wins for the same reason.
let p = merge_policy(
&[sub(Some(&["rust"]), None, Some(2)), sub(Some(&["sqlite"]), None, Some(9))],
&feed(),
3,
);
assert!(p.wanted(&["rust"], &[]) && p.wanted(&["sqlite"], &[]));
assert!(!p.wanted(&["python"], &[]));
assert_eq!(p.budget, 9);
// One person taking everything removes the filter for the shared copy.
let p = merge_policy(&[sub(Some(&["rust"]), None, None), sub(Some(&[]), None, None)], &feed(), 3);
assert!(p.wanted(&["python"], &[]));
// A blocked word keeps an item from being fetched for that person, not for the rest.
let blocking = |w: &str| db::Sub { blocked: vec![w.into()], ..sub(None, None, None) };
let p = merge_policy(&[blocking("politics")], &feed(), 3);
assert!(!p.wanted(&[], &["Politics today"]));
assert!(p.wanted(&[], &["Cooking today"]));
let p = merge_policy(&[blocking("politics"), sub(None, None, None)], &feed(), 3);
assert!(p.wanted(&[], &["Politics today"]), "someone else still wants it");
// Everyone has auto-download off: nothing is fetched automatically.
let p = merge_policy(&[sub(None, Some(false), None), sub(None, Some(false), None)], &feed(), 3);
assert!(!p.auto_download);
// One of them wants it, so it is fetched.
let p = merge_policy(&[sub(None, Some(false), None), sub(None, Some(true), None)], &feed(), 3);
assert!(p.auto_download);
}
async fn test_ctx(cfg: config::Config) -> Ctx {
Ctx {
cfg: std::sync::RwLock::new(Arc::new(cfg)),
db: db::Db::memory().await.unwrap(),
client: reqwest::Client::new(),
feed_client: reqwest::Client::builder().redirect(reqwest::redirect::Policy::none()).build().unwrap(),
out: Emitter::terminal(),
torrents: tokio::sync::OnceCell::new(),
torrent_slots: Arc::new(tokio::sync::Semaphore::new(2)),
config_path: PathBuf::new(),
detach_torrents: false,
}
}
#[tokio::test]
async fn a_derived_feed_is_not_scanned_once_its_opml_leaves_config() {
// davewiner: the OPML subscription left config.toml, but its 922 derived rows
// stayed in the database and kept being scanned under the no-parent fallback.
let ctx = test_ctx(config::Config::default()).await;
ctx.db.upsert_managed("child", "http://x/child.xml", "Child", "gone-opml").await.unwrap();
assert!(
subscriptions(&ctx).await.unwrap().iter().all(|s| s.id != "child"),
"a derived feed whose parent is gone from config must not be scanned"
);
}
#[tokio::test]
async fn retiring_a_group_drops_what_was_never_downloaded_and_orphans_the_rest() {
let ctx = test_ctx(config::Config::default()).await;
ctx.db.upsert_managed("empty", "http://x/empty.xml", "Empty", "parent").await.unwrap();
ctx.db.upsert_managed("has-file", "http://x/has-file.xml", "Has File", "parent").await.unwrap();
let enc = feed::Enclosure { url: "http://x/ep.mp3".into(), mime: None, length: None };
ctx.db.record_enclosure("has-file", "g1", &enc).await.unwrap();
ctx.db.mark_downloaded(&enc.url, std::path::Path::new("/downloads/ep.mp3"), 1).await.unwrap();
retire_group(&ctx, "parent").await.unwrap();
let managed = ctx.db.managed_feeds().await.unwrap();
assert!(!managed.iter().any(|m| m.id == "empty"), "nothing downloaded, so it is forgotten");
assert!(managed.iter().any(|m| m.id == "has-file"), "has a file on disk, so it is kept");
assert!(ctx.db.feed_summary("has-file").await.unwrap().orphaned, "and flagged as orphaned");
}
#[tokio::test]
async fn a_stranded_group_is_retired_but_a_promoted_feed_keeps_its_entries() {
// davewiner: the OPML left config before retire_group existed, and 11 of its feeds
// promoted to config since still said managed = 1.
let mut cfg = config::Config::default();
cfg.feeds.insert("promoted".into(), feed());
cfg.feeds.insert("live-opml".into(), feed());
let ctx = test_ctx(cfg).await;
ctx.db.upsert_managed("promoted", "http://x/p.xml", "Promoted", "gone-opml").await.unwrap();
ctx.db.upsert_managed("empty", "http://x/e.xml", "Empty", "gone-opml").await.unwrap();
ctx.db.upsert_managed("listed", "http://x/l.xml", "Listed", "live-opml").await.unwrap();
ctx.db.record_entry("promoted", &feed::Entry { guid: "g1".into(), ..Default::default() }).await.unwrap();
assert_eq!(retire_stranded(&ctx).await.unwrap(), 2, "empty dropped, promoted unmanaged");
let managed: Vec<String> = ctx.db.managed_feeds().await.unwrap().into_iter().map(|m| m.id).collect();
assert_eq!(managed, ["listed"], "a group still in config is left alone");
assert_eq!(ctx.db.feed_summary("promoted").await.unwrap().entries, 1, "its entries survive");
}
}