Phase 2: web front end
axum served from inside the daemon so it reads SQLite and the event bus directly: browse feeds, read show notes, play with seeking, download and delete files, mark read/flag, and edit feed settings. Config is now hot-reloadable (Ctx.cfg behind RwLock<Arc<Config>>), so UI edits apply without a daemon restart. Access is a shared token minted from /dev/urandom, carried in a cookie because an <audio> element cannot send headers. Show notes are untrusted feed HTML and are sanitized with ammonia server-side. read/flagged finally have a writer, which retention has needed since it started ordering by them. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RPyeapneuXrCdojsaiXGbe
This commit is contained in:
@@ -5,18 +5,20 @@ use serde::{Deserialize, Serialize};
|
||||
use std::collections::BTreeMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
#[derive(Debug, Default, Deserialize, Serialize)]
|
||||
#[derive(Debug, Default, Clone, Deserialize, Serialize)]
|
||||
pub struct Config {
|
||||
#[serde(default)]
|
||||
pub general: General,
|
||||
#[serde(default)]
|
||||
pub torrent: Torrent,
|
||||
#[serde(default)]
|
||||
pub web: Web,
|
||||
/// Keyed by feed id: the TOML table name, which replaces the old genHash(feedURL).
|
||||
#[serde(default)]
|
||||
pub feeds: BTreeMap<String, Feed>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize)]
|
||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||
#[serde(default)]
|
||||
pub struct General {
|
||||
pub download_dir: PathBuf,
|
||||
@@ -39,7 +41,7 @@ pub enum Organize {
|
||||
Date,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize)]
|
||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||
#[serde(default)]
|
||||
pub struct Torrent {
|
||||
pub enabled: bool,
|
||||
@@ -51,6 +53,32 @@ pub struct Torrent {
|
||||
pub stall_mins: u64,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||
#[serde(default)]
|
||||
pub struct Web {
|
||||
pub enabled: bool,
|
||||
/// Use 0.0.0.0 to reach it from the LAN. Anything but loopback needs the token.
|
||||
pub bind: String,
|
||||
/// Shared secret. Generated and written back on first run when left empty.
|
||||
pub token: String,
|
||||
}
|
||||
|
||||
impl Default for Web {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
enabled: false,
|
||||
bind: "127.0.0.1:8080".into(),
|
||||
token: String::new(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Web {
|
||||
pub fn binds_publicly(&self) -> bool {
|
||||
!self.bind.starts_with("127.") && !self.bind.starts_with("localhost")
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||
pub struct Feed {
|
||||
pub url: String,
|
||||
|
||||
155
src/db.rs
155
src/db.rs
@@ -386,6 +386,15 @@ impl Db {
|
||||
}
|
||||
|
||||
impl Db {
|
||||
pub fn unread_count(&self, feed_id: &str) -> Result<i64> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
Ok(conn.query_row(
|
||||
"SELECT count(*) FROM entries WHERE feed_id = ?1 AND read = 0",
|
||||
[feed_id],
|
||||
|r| r.get(0),
|
||||
)?)
|
||||
}
|
||||
|
||||
/// (pending, downloaded) across all feeds, for the status command.
|
||||
pub fn counts(&self) -> Result<(i64, i64)> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
@@ -396,6 +405,152 @@ impl Db {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/// An entry plus its enclosures, for the web UI.
|
||||
#[derive(Debug, serde::Serialize)]
|
||||
pub struct EntryRow {
|
||||
pub guid: String,
|
||||
pub title: Option<String>,
|
||||
pub link: Option<String>,
|
||||
pub published: Option<i64>,
|
||||
pub description: Option<String>,
|
||||
pub read: bool,
|
||||
pub flagged: bool,
|
||||
pub enclosures: Vec<EncRow>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, serde::Serialize)]
|
||||
pub struct EncRow {
|
||||
pub id: i64,
|
||||
pub feed_id: String,
|
||||
pub guid: String,
|
||||
pub url: String,
|
||||
pub mime: Option<String>,
|
||||
pub length: Option<i64>,
|
||||
pub path: Option<String>,
|
||||
pub state: String,
|
||||
pub last_error: Option<String>,
|
||||
}
|
||||
|
||||
impl Db {
|
||||
/// One page of a feed's entries, newest first, each with its enclosures attached.
|
||||
pub fn entries(&self, feed_id: &str, offset: i64, limit: i64) -> Result<Vec<EntryRow>> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT guid, title, link, published, description, read, flagged
|
||||
FROM entries WHERE feed_id = ?1
|
||||
ORDER BY coalesce(published, first_seen) DESC, rowid DESC
|
||||
LIMIT ?3 OFFSET ?2",
|
||||
)?;
|
||||
let mut rows: Vec<EntryRow> = stmt
|
||||
.query_map(rusqlite::params![feed_id, offset, limit], |r| {
|
||||
Ok(EntryRow {
|
||||
guid: r.get(0)?,
|
||||
title: r.get(1)?,
|
||||
link: r.get(2)?,
|
||||
published: r.get(3)?,
|
||||
description: r.get(4)?,
|
||||
read: r.get::<_, i64>(5)? != 0,
|
||||
flagged: r.get::<_, i64>(6)? != 0,
|
||||
enclosures: vec![],
|
||||
})
|
||||
})?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
|
||||
if rows.is_empty() {
|
||||
return Ok(rows);
|
||||
}
|
||||
|
||||
// Only the guids on this page, so a feed with thousands of entries stays cheap.
|
||||
let placeholders = std::iter::repeat_n("?", rows.len()).collect::<Vec<_>>().join(",");
|
||||
let sql = format!(
|
||||
"SELECT id, feed_id, guid, url, mime, length, path, state, last_error
|
||||
FROM enclosures WHERE feed_id = ? AND guid IN ({placeholders}) ORDER BY id"
|
||||
);
|
||||
let mut params: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(rows.len() + 1);
|
||||
params.push(&feed_id);
|
||||
for row in &rows {
|
||||
params.push(&row.guid);
|
||||
}
|
||||
let mut stmt = conn.prepare(&sql)?;
|
||||
let encs = stmt
|
||||
.query_map(params.as_slice(), |r| {
|
||||
Ok(EncRow {
|
||||
id: r.get(0)?,
|
||||
feed_id: r.get(1)?,
|
||||
guid: r.get(2)?,
|
||||
url: r.get(3)?,
|
||||
mime: r.get(4)?,
|
||||
length: r.get(5)?,
|
||||
path: r.get(6)?,
|
||||
state: r.get(7)?,
|
||||
last_error: r.get(8)?,
|
||||
})
|
||||
})?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
|
||||
for enc in encs {
|
||||
if let Some(row) = rows.iter_mut().find(|r| r.guid == enc.guid) {
|
||||
row.enclosures.push(enc);
|
||||
}
|
||||
}
|
||||
Ok(rows)
|
||||
}
|
||||
|
||||
pub fn enclosure(&self, id: i64) -> Result<Option<EncRow>> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
Ok(conn
|
||||
.query_row(
|
||||
"SELECT id, feed_id, guid, url, mime, length, path, state, last_error
|
||||
FROM enclosures WHERE id = ?1",
|
||||
[id],
|
||||
|r| {
|
||||
Ok(EncRow {
|
||||
id: r.get(0)?,
|
||||
feed_id: r.get(1)?,
|
||||
guid: r.get(2)?,
|
||||
url: r.get(3)?,
|
||||
mime: r.get(4)?,
|
||||
length: r.get(5)?,
|
||||
path: r.get(6)?,
|
||||
state: r.get(7)?,
|
||||
last_error: r.get(8)?,
|
||||
})
|
||||
},
|
||||
)
|
||||
.optional()?)
|
||||
}
|
||||
|
||||
/// `read` and `flagged` finally get a writer: retention orders by them.
|
||||
pub fn set_entry_flag(&self, feed_id: &str, guid: &str, field: EntryFlag, on: bool) -> Result<()> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
let sql = match field {
|
||||
EntryFlag::Read => "UPDATE entries SET read = ?3 WHERE feed_id = ?1 AND guid = ?2",
|
||||
EntryFlag::Flagged => "UPDATE entries SET flagged = ?3 WHERE feed_id = ?1 AND guid = ?2",
|
||||
};
|
||||
conn.execute(sql, rusqlite::params![feed_id, guid, on as i64])?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Puts an enclosure back in the queue so the next scan picks it up. This is how a
|
||||
/// `skipped` verdict (from a filter that has since been changed) gets revisited.
|
||||
pub fn requeue(&self, id: i64) -> Result<()> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
conn.execute(
|
||||
"UPDATE enclosures SET state = 'pending', last_error = NULL
|
||||
WHERE id = ?1 AND path IS NULL",
|
||||
[id],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub enum EntryFlag {
|
||||
Read,
|
||||
Flagged,
|
||||
}
|
||||
|
||||
/// Unix seconds. Everything time-shaped in the DB is stored this way.
|
||||
pub fn now() -> i64 {
|
||||
std::time::SystemTime::now()
|
||||
|
||||
177
src/main.rs
177
src/main.rs
@@ -5,11 +5,13 @@ mod feed;
|
||||
mod ipc;
|
||||
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 tokio::sync::{broadcast, mpsc};
|
||||
|
||||
#[derive(Parser)]
|
||||
@@ -64,24 +66,44 @@ enum Command {
|
||||
/// Write subscriptions out as OPML
|
||||
Export { file: PathBuf },
|
||||
/// Run the scheduler and serve the control socket
|
||||
Daemon,
|
||||
Daemon {
|
||||
/// Serve the web UI on this address, overriding [web] in the config
|
||||
#[arg(long, value_name = "ADDR")]
|
||||
web: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
/// Everything a command needs. One per process.
|
||||
struct Ctx {
|
||||
cfg: config::Config,
|
||||
db: db::Db,
|
||||
client: reqwest::Client,
|
||||
out: Emitter,
|
||||
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,
|
||||
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.
|
||||
torrents: tokio::sync::OnceCell<torrent::Torrents>,
|
||||
pub torrents: tokio::sync::OnceCell<torrent::Torrents>,
|
||||
}
|
||||
|
||||
impl Ctx {
|
||||
pub fn cfg(&self) -> std::sync::Arc<config::Config> {
|
||||
self.cfg.read().unwrap().clone()
|
||||
}
|
||||
|
||||
/// Re-reads config.toml into the live snapshot.
|
||||
pub fn reload_cfg(&self, path: &std::path::Path) -> Result<()> {
|
||||
let fresh = config::Config::load(path)?;
|
||||
*self.cfg.write().unwrap() = std::sync::Arc::new(fresh);
|
||||
tracing::info!("config reloaded");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn torrents(&self) -> Result<&torrent::Torrents> {
|
||||
let cfg = self.cfg();
|
||||
self.torrents
|
||||
.get_or_try_init(|| torrent::Torrents::new(&self.cfg))
|
||||
.get_or_try_init(|| torrent::Torrents::new(&cfg))
|
||||
.await
|
||||
}
|
||||
}
|
||||
@@ -109,7 +131,7 @@ async fn main() -> Result<()> {
|
||||
Command::Reap { dry_run } => Some(Cmd::Reap { dry_run: *dry_run }),
|
||||
Command::Status => Some(Cmd::Status),
|
||||
Command::List
|
||||
| Command::Daemon
|
||||
| Command::Daemon { .. }
|
||||
| Command::Add { .. }
|
||||
| Command::Rm { .. }
|
||||
| Command::Import { .. }
|
||||
@@ -123,7 +145,7 @@ async fn main() -> Result<()> {
|
||||
}
|
||||
|
||||
let ctx = Ctx {
|
||||
cfg,
|
||||
cfg: std::sync::RwLock::new(std::sync::Arc::new(cfg)),
|
||||
db,
|
||||
client: reqwest::Client::builder()
|
||||
.user_agent(concat!("ipx/", env!("CARGO_PKG_VERSION")))
|
||||
@@ -134,7 +156,7 @@ async fn main() -> Result<()> {
|
||||
|
||||
match cli.command {
|
||||
Command::List => list(&ctx, &config_path),
|
||||
Command::Daemon => daemon(ctx).await,
|
||||
Command::Daemon { web } => daemon(ctx, config_path, web).await,
|
||||
Command::Add { url, folder, keywords } => {
|
||||
add(ctx, &config_path, &url, folder, keywords).await
|
||||
}
|
||||
@@ -155,28 +177,29 @@ async fn run(ctx: &Ctx, cmd: Cmd) -> Result<()> {
|
||||
Cmd::Reap { dry_run } => reap(ctx, dry_run, true),
|
||||
Cmd::Status => {
|
||||
let (pending, downloaded) = ctx.db.counts()?;
|
||||
ctx.out.emit(Event::Status { feeds: ctx.cfg.feeds.len(), pending, downloaded });
|
||||
ctx.out.emit(Event::Status { feeds: ctx.cfg().feeds.len(), pending, downloaded });
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn daemon(ctx: Ctx) -> Result<()> {
|
||||
let socket = ctx.cfg.general.socket.clone();
|
||||
async fn daemon(ctx: Ctx, config_path: PathBuf, web_addr: Option<String>) -> 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());
|
||||
}
|
||||
|
||||
let (events, _) = broadcast::channel(1024);
|
||||
let (tx_cmd, mut rx_cmd) = mpsc::channel::<Cmd>(64);
|
||||
let ctx = Ctx { out: Emitter::socket(events.clone(), false), ..ctx };
|
||||
let ctx = Arc::new(Ctx { out: Emitter::socket(events.clone(), false), ..ctx });
|
||||
|
||||
let server = tokio::spawn(ipc::serve(socket.clone(), events, tx_cmd));
|
||||
let web = start_web(&ctx, &config_path, web_addr, &tx_cmd, &events).await?;
|
||||
let server = tokio::spawn(ipc::serve(socket.clone(), events.clone(), tx_cmd));
|
||||
|
||||
// One command at a time: the queue is what keeps two scans from overlapping.
|
||||
let mut ticker = tokio::time::interval(std::time::Duration::from_secs(60));
|
||||
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
||||
tracing::info!(feeds = ctx.cfg.feeds.len(), "daemon started");
|
||||
tracing::info!(feeds = ctx.cfg().feeds.len(), "daemon started");
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
@@ -197,11 +220,61 @@ async fn daemon(ctx: Ctx) -> Result<()> {
|
||||
}
|
||||
|
||||
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 = web::generate_token();
|
||||
fresh.save(config_path)?;
|
||||
ctx.reload_cfg(config_path)?;
|
||||
println!("web ui token generated. Open:\n http://{bind}/?token={}", fresh.web.token);
|
||||
} else {
|
||||
println!(
|
||||
"web ui at http://{bind}/?token={}",
|
||||
ctx.cfg().web.token
|
||||
);
|
||||
}
|
||||
|
||||
if ctx.cfg().web.binds_publicly() {
|
||||
tracing::warn!(bind, "web ui is reachable off this machine; the token is all that guards it");
|
||||
}
|
||||
|
||||
let state = web::WebState {
|
||||
ctx: ctx.clone(),
|
||||
config_path: config_path.to_path_buf(),
|
||||
cmds: cmds.clone(),
|
||||
events: events.clone(),
|
||||
};
|
||||
Ok(Some(tokio::spawn(async move {
|
||||
if let Err(e) = web::serve(state, &bind).await {
|
||||
tracing::error!(error = ?e, "web ui stopped");
|
||||
}
|
||||
})))
|
||||
}
|
||||
|
||||
async fn shutdown() {
|
||||
use tokio::signal::unix::{SignalKind, signal};
|
||||
let mut term = match signal(SignalKind::terminate()) {
|
||||
@@ -216,25 +289,27 @@ async fn shutdown() {
|
||||
|
||||
/// Subscribes to one feed, naming it from its own title.
|
||||
async fn add(
|
||||
mut ctx: Ctx,
|
||||
ctx: Ctx,
|
||||
config_path: &std::path::Path,
|
||||
url: &str,
|
||||
folder: Option<String>,
|
||||
keywords: Vec<String>,
|
||||
) -> Result<()> {
|
||||
if let Some((id, _)) = ctx.cfg.feeds.iter().find(|(_, f)| f.url == url) {
|
||||
let mut cfg = (*ctx.cfg()).clone();
|
||||
if let Some((id, _)) = cfg.feeds.iter().find(|(_, f)| f.url == url) {
|
||||
anyhow::bail!("already subscribed as {id:?}");
|
||||
}
|
||||
let id = add_one(&mut ctx, url, folder, keywords).await?;
|
||||
ctx.cfg.save(config_path)?;
|
||||
let id = add_one(&ctx, &mut cfg, url, folder, keywords).await?;
|
||||
cfg.save(config_path)?;
|
||||
println!("added {id}");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Returns the new feed id. The title needs a fetch, so a feed that cannot be reached is
|
||||
/// still added -- under a slug derived from its URL -- rather than refused.
|
||||
async fn add_one(
|
||||
ctx: &mut Ctx,
|
||||
pub async fn add_one(
|
||||
ctx: &Ctx,
|
||||
cfg: &mut config::Config,
|
||||
url: &str,
|
||||
folder: Option<String>,
|
||||
keywords: Vec<String>,
|
||||
@@ -262,8 +337,8 @@ async fn add_one(
|
||||
}
|
||||
};
|
||||
|
||||
let id = config::unique_slug(&title, &ctx.cfg.feeds);
|
||||
ctx.cfg.feeds.insert(id.clone(), probe);
|
||||
let id = config::unique_slug(&title, &cfg.feeds);
|
||||
cfg.feeds.insert(id.clone(), probe);
|
||||
Ok(id)
|
||||
}
|
||||
|
||||
@@ -275,17 +350,19 @@ fn url_stem(url: &str) -> String {
|
||||
.unwrap_or_else(|| url.to_owned())
|
||||
}
|
||||
|
||||
fn rm(mut ctx: Ctx, config_path: &std::path::Path, feed: &str) -> Result<()> {
|
||||
if ctx.cfg.feeds.remove(feed).is_none() {
|
||||
fn rm(ctx: Ctx, config_path: &std::path::Path, feed: &str) -> Result<()> {
|
||||
let mut cfg = (*ctx.cfg()).clone();
|
||||
if cfg.feeds.remove(feed).is_none() {
|
||||
anyhow::bail!("no feed with id {feed:?}");
|
||||
}
|
||||
ctx.cfg.save(config_path)?;
|
||||
cfg.save(config_path)?;
|
||||
// State and files stay: re-adding the feed should not re-download its back catalogue.
|
||||
println!("removed {feed}; downloads and history kept");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn import(mut ctx: Ctx, config_path: &std::path::Path, file: &std::path::Path) -> Result<()> {
|
||||
async fn import(ctx: Ctx, config_path: &std::path::Path, file: &std::path::Path) -> Result<()> {
|
||||
let mut cfg = (*ctx.cfg()).clone();
|
||||
let text = std::fs::read_to_string(file)
|
||||
.with_context(|| format!("reading {}", file.display()))?;
|
||||
let doc = opml::OPML::from_str(&text).map_err(|e| anyhow::anyhow!("parsing OPML: {e}"))?;
|
||||
@@ -295,12 +372,12 @@ async fn import(mut ctx: Ctx, config_path: &std::path::Path, file: &std::path::P
|
||||
|
||||
let mut added = 0;
|
||||
for (title, url) in found {
|
||||
if ctx.cfg.feeds.values().any(|f| f.url == url) {
|
||||
if cfg.feeds.values().any(|f| f.url == url) {
|
||||
continue;
|
||||
}
|
||||
// Name it from the OPML title rather than refetching every feed.
|
||||
let id = config::unique_slug(&title, &ctx.cfg.feeds);
|
||||
ctx.cfg.feeds.insert(
|
||||
let id = config::unique_slug(&title, &cfg.feeds);
|
||||
cfg.feeds.insert(
|
||||
id.clone(),
|
||||
config::Feed {
|
||||
url,
|
||||
@@ -317,7 +394,7 @@ async fn import(mut ctx: Ctx, config_path: &std::path::Path, file: &std::path::P
|
||||
println!("added {id}");
|
||||
added += 1;
|
||||
}
|
||||
ctx.cfg.save(config_path)?;
|
||||
cfg.save(config_path)?;
|
||||
println!("{added} feed(s) imported");
|
||||
Ok(())
|
||||
}
|
||||
@@ -339,7 +416,7 @@ fn export(ctx: &Ctx, file: &std::path::Path) -> Result<()> {
|
||||
title: Some("ipx subscriptions".into()),
|
||||
..Default::default()
|
||||
});
|
||||
for (id, feed) in &ctx.cfg.feeds {
|
||||
for (id, feed) in &ctx.cfg().feeds {
|
||||
let title = ctx
|
||||
.db
|
||||
.feed_summary(id)
|
||||
@@ -350,16 +427,17 @@ fn export(ctx: &Ctx, file: &std::path::Path) -> Result<()> {
|
||||
}
|
||||
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());
|
||||
println!("exported {} feed(s) to {}", ctx.cfg().feeds.len(), file.display());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn list(ctx: &Ctx, config_path: &std::path::Path) -> Result<()> {
|
||||
if ctx.cfg.feeds.is_empty() {
|
||||
let cfg = ctx.cfg();
|
||||
if cfg.feeds.is_empty() {
|
||||
println!("No feeds configured in {}", config_path.display());
|
||||
return Ok(());
|
||||
}
|
||||
for (id, feed) in &ctx.cfg.feeds {
|
||||
for (id, feed) in &cfg.feeds {
|
||||
let s = ctx.db.feed_summary(id)?;
|
||||
println!("{id} {}", s.title.as_deref().unwrap_or("-"));
|
||||
println!(" url {}", feed.url);
|
||||
@@ -376,7 +454,7 @@ fn list(ctx: &Ctx, config_path: &std::path::Path) -> Result<()> {
|
||||
/// deleted, but must not emit the terminal ReapDone, or a client waiting on its `fetch`
|
||||
/// would stop reading before the scan had even started.
|
||||
fn reap(ctx: &Ctx, dry_run: bool, standalone: bool) -> Result<()> {
|
||||
let r = retention::run(&ctx.cfg, &ctx.db, dry_run)?;
|
||||
let r = retention::run(&ctx.cfg(), &ctx.db, dry_run)?;
|
||||
for c in r.aged_out.iter().chain(r.over_quota.iter()) {
|
||||
ctx.out.emit(Event::Reaped {
|
||||
path: c.path.clone(),
|
||||
@@ -393,15 +471,15 @@ fn reap(ctx: &Ctx, dry_run: bool, standalone: bool) -> Result<()> {
|
||||
}
|
||||
|
||||
async fn fetch(ctx: &Ctx, only: Option<&str>, force: bool) -> Result<()> {
|
||||
let cfg = ctx.cfg();
|
||||
if let Some(id) = only
|
||||
&& !ctx.cfg.feeds.contains_key(id)
|
||||
&& !cfg.feeds.contains_key(id)
|
||||
{
|
||||
anyhow::bail!("no feed with id {id:?}");
|
||||
}
|
||||
|
||||
let mut scanned = 0;
|
||||
for (id, feed_cfg) in ctx
|
||||
.cfg
|
||||
for (id, feed_cfg) in cfg
|
||||
.feeds
|
||||
.iter()
|
||||
.filter(|(id, _)| only.is_none_or(|o| o == *id))
|
||||
@@ -410,7 +488,7 @@ async fn fetch(ctx: &Ctx, only: Option<&str>, force: bool) -> Result<()> {
|
||||
|
||||
// TTL: the feed's own <ttl> wins when it is longer than our poll interval.
|
||||
if !force && let Some(last) = state.last_checked {
|
||||
let wait = state.ttl_mins.unwrap_or(0).max(ctx.cfg.general.interval_mins) * 60;
|
||||
let wait = state.ttl_mins.unwrap_or(0).max(cfg.general.interval_mins) * 60;
|
||||
let due = last + wait as i64;
|
||||
if due > db::now() {
|
||||
ctx.out.emit(Event::FeedSkip {
|
||||
@@ -507,12 +585,13 @@ async fn scan_one(
|
||||
|
||||
let budget = feed_cfg.max_new_per_check.unwrap_or(usize::MAX);
|
||||
if feed_cfg.auto_download && budget > 0 {
|
||||
let folder = download::folder_for(&ctx.cfg, id, feed_cfg, parsed.title.as_deref());
|
||||
let dest_dir = ctx.cfg.general.download_dir.join(&folder);
|
||||
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)? {
|
||||
if download::looks_like_torrent(&item.url, item.mime.as_deref()) {
|
||||
if !ctx.cfg.torrent.enabled {
|
||||
if !ctx.cfg().torrent.enabled {
|
||||
ctx.db.mark_enclosure(&item.url, "skipped", Some("torrents disabled"))?;
|
||||
ctx.out.emit(Event::TorrentDeferred {
|
||||
feed: id.to_string(),
|
||||
@@ -603,7 +682,8 @@ async fn fetch_one(
|
||||
// Throttled to whole percents, as the original's lastDLStepSize guard did.
|
||||
let mut last_pct = -1i64;
|
||||
let name = download::filename_for(url, None);
|
||||
let got = download::download(&ctx.client, &ctx.cfg, feed_cfg, url, |done, total| {
|
||||
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 {
|
||||
@@ -624,7 +704,7 @@ async fn fetch_one(
|
||||
// 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 {
|
||||
if !ctx.cfg().torrent.enabled {
|
||||
ctx.db.mark_enclosure(url, "skipped", Some("torrents disabled"))?;
|
||||
anyhow::bail!("body is a torrent and torrents are disabled");
|
||||
}
|
||||
@@ -646,9 +726,10 @@ async fn torrent_one(
|
||||
) -> Result<(PathBuf, u64)> {
|
||||
let name = download::filename_for(url, None);
|
||||
let mut last_pct = -1i64;
|
||||
let cfg = ctx.cfg();
|
||||
ctx.torrents()
|
||||
.await?
|
||||
.fetch(&ctx.cfg, url, dest_dir, |done, total| {
|
||||
.fetch(&cfg, url, dest_dir, |done, total| {
|
||||
if total > 0 {
|
||||
let pct = (done * 100 / total) as i64;
|
||||
if pct > last_pct {
|
||||
|
||||
437
src/web.rs
Normal file
437
src/web.rs
Normal file
@@ -0,0 +1,437 @@
|
||||
//! Web front end. Runs inside the daemon so it reads SQLite and the event bus directly.
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use axum::{
|
||||
Json, Router,
|
||||
extract::{Path, Query, Request, State},
|
||||
http::{StatusCode, header},
|
||||
middleware::{self, Next},
|
||||
response::{
|
||||
Html, IntoResponse, Response,
|
||||
sse::{Event as SseEvent, Sse},
|
||||
},
|
||||
routing::{delete, get, patch, post},
|
||||
};
|
||||
use futures_util::StreamExt;
|
||||
use serde::Deserialize;
|
||||
use tokio_stream::wrappers::BroadcastStream;
|
||||
use tower::ServiceExt;
|
||||
use tower_http::services::ServeFile;
|
||||
use serde::Serialize;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::{broadcast, mpsc};
|
||||
|
||||
use crate::Ctx;
|
||||
use crate::ipc::{Command, Event};
|
||||
|
||||
const COOKIE: &str = "ipx_token";
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct WebState {
|
||||
pub ctx: Arc<Ctx>,
|
||||
pub config_path: PathBuf,
|
||||
pub cmds: mpsc::Sender<Command>,
|
||||
pub events: broadcast::Sender<Event>,
|
||||
}
|
||||
|
||||
/// A 32-hex-character shared secret, generated when config.toml has none.
|
||||
///
|
||||
/// ponytail: /dev/urandom rather than a CSPRNG crate -- 16 bytes, once, on a Unix-only
|
||||
/// binary. Falls back to the clock only if urandom is somehow unreadable, which would be a
|
||||
/// weak token, so that case is logged loudly.
|
||||
pub fn generate_token() -> String {
|
||||
use std::io::Read;
|
||||
let mut bytes = [0u8; 16];
|
||||
match std::fs::File::open("/dev/urandom").and_then(|mut f| f.read_exact(&mut bytes)) {
|
||||
Ok(()) => {}
|
||||
Err(e) => {
|
||||
tracing::error!(error = %e, "could not read /dev/urandom; token is NOT secure");
|
||||
let n = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|d| d.as_nanos() as u64)
|
||||
.unwrap_or(0);
|
||||
bytes[..8].copy_from_slice(&n.to_le_bytes());
|
||||
}
|
||||
}
|
||||
bytes.iter().map(|b| format!("{b:02x}")).collect()
|
||||
}
|
||||
|
||||
pub fn router(state: WebState) -> Router {
|
||||
Router::new()
|
||||
.route("/", get(index))
|
||||
.route("/api/feeds", get(feeds).post(add_feed))
|
||||
.route("/api/feeds/{id}", patch(patch_feed).delete(remove_feed))
|
||||
.route("/api/feeds/{id}/entries", get(entries))
|
||||
.route("/api/entries/{feed_id}/{guid}/flags", post(set_flags))
|
||||
.route("/api/enclosures/{id}/download", post(download_now))
|
||||
.route("/api/enclosures/{id}", delete(delete_file))
|
||||
.route("/api/fetch", post(fetch_now))
|
||||
.route("/api/events", get(events))
|
||||
.route("/media/{id}", get(media))
|
||||
.layer(middleware::from_fn_with_state(state.clone(), auth))
|
||||
.with_state(state)
|
||||
}
|
||||
|
||||
pub async fn serve(state: WebState, bind: &str) -> Result<()> {
|
||||
let listener = tokio::net::TcpListener::bind(bind)
|
||||
.await
|
||||
.with_context(|| format!("binding {bind}"))?;
|
||||
tracing::info!(bind, "web ui listening");
|
||||
axum::serve(listener, router(state))
|
||||
.await
|
||||
.context("serving the web ui")
|
||||
}
|
||||
|
||||
/// Token in `?token=` (which then sets a cookie) or in the cookie itself.
|
||||
///
|
||||
/// It has to be a cookie rather than a header: an `<audio src>` request is issued by the
|
||||
/// browser, and there is no way to attach a header to it.
|
||||
async fn auth(State(state): State<WebState>, req: Request, next: Next) -> Response {
|
||||
let expected = state.ctx.cfg().web.token.clone();
|
||||
if expected.is_empty() {
|
||||
// Refuse to serve rather than serve unauthenticated.
|
||||
return (StatusCode::INTERNAL_SERVER_ERROR, "no web token configured").into_response();
|
||||
}
|
||||
|
||||
let from_query = req.uri().query().and_then(|q| {
|
||||
q.split('&')
|
||||
.find_map(|kv| kv.strip_prefix("token=").map(str::to_owned))
|
||||
});
|
||||
let from_cookie = req
|
||||
.headers()
|
||||
.get(header::COOKIE)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(|c| {
|
||||
c.split(';')
|
||||
.find_map(|kv| kv.trim().strip_prefix(&format!("{COOKIE}=")).map(str::to_owned))
|
||||
});
|
||||
|
||||
let supplied = from_query.clone().or(from_cookie);
|
||||
if !supplied.is_some_and(|t| constant_time_eq(&t, &expected)) {
|
||||
return (StatusCode::UNAUTHORIZED, "bad or missing token").into_response();
|
||||
}
|
||||
|
||||
let mut resp = next.run(req).await;
|
||||
if from_query.is_some() {
|
||||
// Remember it so the rest of the page (and the audio element) authenticates.
|
||||
if let Ok(v) = header::HeaderValue::from_str(&format!(
|
||||
"{COOKIE}={expected}; Path=/; SameSite=Lax; Max-Age=31536000"
|
||||
)) {
|
||||
resp.headers_mut().insert(header::SET_COOKIE, v);
|
||||
}
|
||||
}
|
||||
resp
|
||||
}
|
||||
|
||||
/// Compares without leaking length or position through timing.
|
||||
fn constant_time_eq(a: &str, b: &str) -> bool {
|
||||
let (a, b) = (a.as_bytes(), b.as_bytes());
|
||||
if a.len() != b.len() {
|
||||
return false;
|
||||
}
|
||||
a.iter().zip(b).fold(0u8, |acc, (x, y)| acc | (x ^ y)) == 0
|
||||
}
|
||||
|
||||
async fn index() -> Html<&'static str> {
|
||||
Html(include_str!("../web/index.html"))
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct FeedRow {
|
||||
id: String,
|
||||
url: String,
|
||||
title: Option<String>,
|
||||
folder: Option<String>,
|
||||
keywords: Vec<String>,
|
||||
allow_explicit: bool,
|
||||
auto_download: bool,
|
||||
max_new_per_check: Option<usize>,
|
||||
last_checked: Option<i64>,
|
||||
last_error: Option<String>,
|
||||
entries: i64,
|
||||
downloaded: i64,
|
||||
unread: i64,
|
||||
}
|
||||
|
||||
async fn feeds(State(state): State<WebState>) -> Result<Json<Vec<FeedRow>>, ApiError> {
|
||||
let cfg = state.ctx.cfg();
|
||||
let mut out = Vec::with_capacity(cfg.feeds.len());
|
||||
for (id, feed) in &cfg.feeds {
|
||||
let s = state.ctx.db.feed_summary(id)?;
|
||||
out.push(FeedRow {
|
||||
id: id.clone(),
|
||||
url: feed.url.clone(),
|
||||
title: s.title,
|
||||
folder: feed.folder.clone(),
|
||||
keywords: feed.keywords.clone(),
|
||||
allow_explicit: feed.allow_explicit,
|
||||
auto_download: feed.auto_download,
|
||||
max_new_per_check: feed.max_new_per_check,
|
||||
last_checked: s.last_checked,
|
||||
last_error: s.last_error,
|
||||
entries: s.entries,
|
||||
downloaded: s.downloaded,
|
||||
unread: state.ctx.db.unread_count(id)?,
|
||||
});
|
||||
}
|
||||
Ok(Json(out))
|
||||
}
|
||||
|
||||
/// Turns anyhow errors into a 500 with a readable body.
|
||||
pub struct ApiError(anyhow::Error);
|
||||
|
||||
impl<E: Into<anyhow::Error>> From<E> for ApiError {
|
||||
fn from(e: E) -> Self {
|
||||
Self(e.into())
|
||||
}
|
||||
}
|
||||
|
||||
impl IntoResponse for ApiError {
|
||||
fn into_response(self) -> Response {
|
||||
tracing::warn!(error = ?self.0, "api error");
|
||||
(StatusCode::INTERNAL_SERVER_ERROR, format!("{:#}", self.0)).into_response()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn token_comparison_rejects_mismatches_and_length_differences() {
|
||||
assert!(constant_time_eq("abc123", "abc123"));
|
||||
assert!(!constant_time_eq("abc123", "abc124"));
|
||||
assert!(!constant_time_eq("abc", "abc123"));
|
||||
assert!(!constant_time_eq("", "abc"));
|
||||
assert!(constant_time_eq("", ""));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn generated_tokens_are_32_hex_chars_and_not_repeated() {
|
||||
let a = generate_token();
|
||||
let b = generate_token();
|
||||
assert_eq!(a.len(), 32);
|
||||
assert!(a.chars().all(|c| c.is_ascii_hexdigit()));
|
||||
assert_ne!(a, b);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct Page {
|
||||
#[serde(default)]
|
||||
offset: i64,
|
||||
#[serde(default = "fifty")]
|
||||
limit: i64,
|
||||
}
|
||||
|
||||
fn fifty() -> i64 {
|
||||
50
|
||||
}
|
||||
|
||||
async fn entries(
|
||||
State(state): State<WebState>,
|
||||
Path(id): Path<String>,
|
||||
Query(page): Query<Page>,
|
||||
) -> Result<Json<Vec<crate::db::EntryRow>>, ApiError> {
|
||||
let mut rows = state.ctx.db.entries(&id, page.offset, page.limit.clamp(1, 200))?;
|
||||
// Feed HTML is untrusted: it reaches the page only after ammonia has been through it.
|
||||
for row in &mut rows {
|
||||
if let Some(d) = &row.description {
|
||||
row.description = Some(ammonia::clean(d));
|
||||
}
|
||||
}
|
||||
Ok(Json(rows))
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct NewFeed {
|
||||
url: String,
|
||||
#[serde(default)]
|
||||
folder: Option<String>,
|
||||
#[serde(default)]
|
||||
keywords: Vec<String>,
|
||||
}
|
||||
|
||||
async fn add_feed(
|
||||
State(state): State<WebState>,
|
||||
Json(body): Json<NewFeed>,
|
||||
) -> Result<Json<serde_json::Value>, ApiError> {
|
||||
let mut cfg = (*state.ctx.cfg()).clone();
|
||||
if let Some((id, _)) = cfg.feeds.iter().find(|(_, f)| f.url == body.url) {
|
||||
return Ok(Json(serde_json::json!({ "id": id, "existing": true })));
|
||||
}
|
||||
let id = crate::add_one(&state.ctx, &mut cfg, &body.url, body.folder, body.keywords).await?;
|
||||
cfg.save(&state.config_path)?;
|
||||
state.ctx.reload_cfg(&state.config_path)?;
|
||||
Ok(Json(serde_json::json!({ "id": id, "existing": false })))
|
||||
}
|
||||
|
||||
/// Only the fields that are present are changed.
|
||||
#[derive(Deserialize)]
|
||||
struct FeedPatch {
|
||||
folder: Option<Option<String>>,
|
||||
keywords: Option<Vec<String>>,
|
||||
allow_explicit: Option<bool>,
|
||||
auto_download: Option<bool>,
|
||||
max_new_per_check: Option<Option<usize>>,
|
||||
}
|
||||
|
||||
async fn patch_feed(
|
||||
State(state): State<WebState>,
|
||||
Path(id): Path<String>,
|
||||
Json(body): Json<FeedPatch>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
let mut cfg = (*state.ctx.cfg()).clone();
|
||||
let feed = cfg
|
||||
.feeds
|
||||
.get_mut(&id)
|
||||
.ok_or_else(|| anyhow::anyhow!("no feed with id {id:?}"))?;
|
||||
|
||||
if let Some(v) = body.folder {
|
||||
feed.folder = v.filter(|s| !s.trim().is_empty());
|
||||
}
|
||||
if let Some(v) = body.keywords {
|
||||
feed.keywords = v.into_iter().filter(|k| !k.trim().is_empty()).collect();
|
||||
}
|
||||
if let Some(v) = body.allow_explicit {
|
||||
feed.allow_explicit = v;
|
||||
}
|
||||
if let Some(v) = body.auto_download {
|
||||
feed.auto_download = v;
|
||||
}
|
||||
if let Some(v) = body.max_new_per_check {
|
||||
feed.max_new_per_check = v;
|
||||
}
|
||||
cfg.save(&state.config_path)?;
|
||||
state.ctx.reload_cfg(&state.config_path)?;
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
}
|
||||
|
||||
async fn remove_feed(
|
||||
State(state): State<WebState>,
|
||||
Path(id): Path<String>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
let mut cfg = (*state.ctx.cfg()).clone();
|
||||
if cfg.feeds.remove(&id).is_none() {
|
||||
return Err(anyhow::anyhow!("no feed with id {id:?}").into());
|
||||
}
|
||||
// Downloads and history stay, so re-adding does not re-pull the back catalogue.
|
||||
cfg.save(&state.config_path)?;
|
||||
state.ctx.reload_cfg(&state.config_path)?;
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct Flags {
|
||||
read: Option<bool>,
|
||||
flagged: Option<bool>,
|
||||
}
|
||||
|
||||
async fn set_flags(
|
||||
State(state): State<WebState>,
|
||||
Path((feed_id, guid)): Path<(String, String)>,
|
||||
Json(body): Json<Flags>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
use crate::db::EntryFlag;
|
||||
if let Some(v) = body.read {
|
||||
state.ctx.db.set_entry_flag(&feed_id, &guid, EntryFlag::Read, v)?;
|
||||
}
|
||||
if let Some(v) = body.flagged {
|
||||
state.ctx.db.set_entry_flag(&feed_id, &guid, EntryFlag::Flagged, v)?;
|
||||
}
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
}
|
||||
|
||||
/// Puts one enclosure back in the queue and kicks a scan of its feed. The queue is the
|
||||
/// table, so this is all it takes -- including for something a filter once skipped.
|
||||
async fn download_now(
|
||||
State(state): State<WebState>,
|
||||
Path(id): Path<i64>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
let enc = state
|
||||
.ctx
|
||||
.db
|
||||
.enclosure(id)?
|
||||
.ok_or_else(|| anyhow::anyhow!("no enclosure {id}"))?;
|
||||
if enc.path.is_some() {
|
||||
return Ok(StatusCode::NO_CONTENT); // Already here.
|
||||
}
|
||||
state.ctx.db.requeue(id)?;
|
||||
let _ = state
|
||||
.cmds
|
||||
.send(Command::Fetch { feed: Some(enc.feed_id), force: true })
|
||||
.await;
|
||||
Ok(StatusCode::ACCEPTED)
|
||||
}
|
||||
|
||||
async fn delete_file(
|
||||
State(state): State<WebState>,
|
||||
Path(id): Path<i64>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
let enc = state
|
||||
.ctx
|
||||
.db
|
||||
.enclosure(id)?
|
||||
.ok_or_else(|| anyhow::anyhow!("no enclosure {id}"))?;
|
||||
if let Some(path) = &enc.path
|
||||
&& let Err(e) = std::fs::remove_file(path)
|
||||
&& e.kind() != std::io::ErrorKind::NotFound
|
||||
{
|
||||
return Err(e.into());
|
||||
}
|
||||
// The row survives as 'reaped', which is what stops the next scan re-downloading it.
|
||||
state.ctx.db.mark_reaped(id)?;
|
||||
state.events.send(Event::Reaped {
|
||||
path: enc.path.unwrap_or_default(),
|
||||
bytes: enc.length.unwrap_or(0).max(0) as u64,
|
||||
}).ok();
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct FetchBody {
|
||||
#[serde(default)]
|
||||
feed: Option<String>,
|
||||
#[serde(default)]
|
||||
force: bool,
|
||||
}
|
||||
|
||||
async fn fetch_now(
|
||||
State(state): State<WebState>,
|
||||
Json(body): Json<FetchBody>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
state
|
||||
.cmds
|
||||
.send(Command::Fetch { feed: body.feed, force: body.force })
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("the daemon is not accepting commands"))?;
|
||||
Ok(StatusCode::ACCEPTED)
|
||||
}
|
||||
|
||||
/// 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>>> {
|
||||
let stream = BroadcastStream::new(state.events.subscribe()).filter_map(|ev| async move {
|
||||
let ev = ev.ok()?;
|
||||
Some(Ok(SseEvent::default().data(serde_json::to_string(&ev).ok()?)))
|
||||
});
|
||||
Sse::new(stream).keep_alive(axum::response::sse::KeepAlive::default())
|
||||
}
|
||||
|
||||
/// Audio, served by ServeFile so Range requests work and the player can seek.
|
||||
async fn media(
|
||||
State(state): State<WebState>,
|
||||
Path(id): Path<i64>,
|
||||
req: Request,
|
||||
) -> Response {
|
||||
let Ok(Some(enc)) = state.ctx.db.enclosure(id) else {
|
||||
return (StatusCode::NOT_FOUND, "no such enclosure").into_response();
|
||||
};
|
||||
let Some(path) = enc.path else {
|
||||
return (StatusCode::NOT_FOUND, "not downloaded").into_response();
|
||||
};
|
||||
match ServeFile::new(path).oneshot(req).await {
|
||||
Ok(r) => r.into_response(),
|
||||
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response(),
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user