`docker exec iPX ipx rm make` printed "removed make" and took it out of the catalogue table, but the web UI went on listing it. add, rm and import run in a process of their own and write the whole catalogue there; the daemon kept its own copy in memory, and its next change of its own -- an admin's edit, a feed categorised (#117), a moved feed followed (#115) -- wrote the removed feed back and dropped any the command line had added. A new socket command, reload, has the daemon read the catalogue and settings again from the database. The socket answers it at once, as it answers status, with the counts after it: queued behind a scan, the web UI could write the old copy back in the meantime. add, rm and import send it when a daemon is live, and say so if it does not answer. User commands need nothing: the daemon reads accounts from the database on every request. A browser test runs the real `ipx add --list` and `ipx rm` against the test daemon, with an admin's edit after each, and fails without this. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
456 lines
19 KiB
Rust
456 lines
19 KiB
Rust
//! Unix-socket control and event stream. Replaces printMSG's `;;1;;1;;100.00;;42.31`.
|
|
|
|
use anyhow::{Context, Result};
|
|
use serde::{Deserialize, Serialize};
|
|
use std::path::{Path, PathBuf};
|
|
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
|
use tokio::net::{UnixListener, UnixStream};
|
|
use tokio::sync::{broadcast, mpsc};
|
|
|
|
/// One JSON object per line, `ev` naming the variant.
|
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
#[serde(tag = "ev", rename_all = "snake_case")]
|
|
pub enum Event {
|
|
FeedStart { feed: String },
|
|
FeedSkip { feed: String, reason: String },
|
|
FeedDone { feed: String, new: usize, downloaded: usize, failed: usize, torrents: usize },
|
|
FeedError { feed: String, msg: String },
|
|
/// The feed answered from a new address after redirects that all said it moved for good,
|
|
/// and the catalogue now has that address (#115).
|
|
FeedMoved { feed: String, from: String, to: String },
|
|
Progress {
|
|
feed: String,
|
|
/// Which enclosure this is about. Without it a UI cannot tell one download's
|
|
/// progress from another's and ends up animating every pending row.
|
|
enclosure: i64,
|
|
url: String,
|
|
file: String,
|
|
done: u64,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
total: Option<u64>,
|
|
},
|
|
DownloadDone { feed: String, enclosure: i64, url: String, path: String, bytes: u64 },
|
|
DownloadError { feed: String, enclosure: i64, url: String, msg: String },
|
|
TorrentDeferred { feed: String, url: String },
|
|
Reaped { path: String, bytes: u64 },
|
|
/// Terminal: a client that asked for work stops reading here.
|
|
ScanDone { feeds: usize },
|
|
ReapDone { files: usize, bytes: u64 },
|
|
Status { feeds: usize, pending: i64, downloaded: i64 },
|
|
Error { msg: String },
|
|
}
|
|
|
|
impl Event {
|
|
pub fn is_terminal(&self) -> bool {
|
|
matches!(self, Event::ScanDone { .. } | Event::ReapDone { .. } | Event::Status { .. })
|
|
}
|
|
|
|
/// The human rendering, for a terminal rather than a UI.
|
|
pub fn human(&self) -> Option<String> {
|
|
Some(match self {
|
|
Event::FeedSkip { feed, reason } => format!("{feed}: {reason}"),
|
|
Event::FeedDone { feed, new, downloaded, failed, torrents } => format!(
|
|
"{feed}: {new} new entries, {downloaded} downloaded, {failed} failed, {torrents} torrents deferred"
|
|
),
|
|
Event::FeedError { feed, msg } => format!("{feed}: error: {msg}"),
|
|
Event::FeedMoved { feed, from, to } => format!("{feed}: moved for good from {from} to {to}; following it"),
|
|
Event::Progress { file, done, total, .. } => match total {
|
|
Some(t) if *t > 0 => format!(
|
|
" {file}: {:.1}% ({:.1}/{:.1} MB)",
|
|
*done as f64 / *t as f64 * 100.0,
|
|
*done as f64 / 1_048_576.0,
|
|
*t as f64 / 1_048_576.0
|
|
),
|
|
_ => format!(" {file}: {:.1} MB", *done as f64 / 1_048_576.0),
|
|
},
|
|
Event::DownloadDone { path, .. } => format!(" saved {path}"),
|
|
Event::DownloadError { url, msg, .. } => format!(" failed {url}: {msg}"),
|
|
Event::Reaped { path, bytes } => {
|
|
format!("deleted {path} ({:.1} MB)", *bytes as f64 / 1_048_576.0)
|
|
}
|
|
Event::ReapDone { files, bytes } => format!(
|
|
"deleted {files} old file(s), {:.1} MB",
|
|
*bytes as f64 / 1_048_576.0
|
|
),
|
|
Event::Status { feeds, pending, downloaded } => {
|
|
format!("{feeds} feeds, {pending} queued to download, {downloaded} downloaded")
|
|
}
|
|
Event::Error { msg } => format!("error: {msg}"),
|
|
// Noise in a terminal; a UI still gets them on the socket.
|
|
Event::FeedStart { feed } => format!("{feed}: checking"),
|
|
Event::TorrentDeferred { feed, .. } => format!("{feed}: torrent deferred"),
|
|
Event::ScanDone { feeds } => format!("scan complete, {feeds} feed(s)"),
|
|
})
|
|
}
|
|
}
|
|
|
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
|
#[serde(tag = "cmd", rename_all = "snake_case")]
|
|
pub enum Command {
|
|
Fetch {
|
|
#[serde(default)]
|
|
feed: Option<String>,
|
|
#[serde(default)]
|
|
force: bool,
|
|
/// Only these feeds, and the feeds inside any of them that is an OPML: "check every feed"
|
|
/// from the web UI is every feed of the person asking, not of everyone (issue #37).
|
|
/// Empty is every feed, as the schedule and the CLI mean it.
|
|
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
|
feeds: Vec<String>,
|
|
},
|
|
Reap {
|
|
#[serde(default)]
|
|
dry_run: bool,
|
|
},
|
|
/// Fetch one specific enclosure now, ignoring max_new_per_check and the queue order.
|
|
/// A scan cannot express "this one, now": it takes the lowest-id pending rows up to
|
|
/// the per-scan cap, so an explicit request has to bypass both.
|
|
Download {
|
|
enclosure: i64,
|
|
},
|
|
Status,
|
|
/// Read the catalogue and server settings again from the database, which `ipx add`, `rm`
|
|
/// and `import` change in a process of their own. A daemon keeps them in memory, and with
|
|
/// its old copy wrote a removed feed back, and dropped an added one, at its next change of
|
|
/// its own (#120). Answered, as status is, with the counts after it.
|
|
Reload,
|
|
}
|
|
|
|
/// Where events go: the socket, the terminal, or both.
|
|
#[derive(Clone)]
|
|
pub struct Emitter {
|
|
tx: Option<broadcast::Sender<Event>>,
|
|
print: bool,
|
|
}
|
|
|
|
impl Emitter {
|
|
pub fn terminal() -> Self {
|
|
Self { tx: None, print: true }
|
|
}
|
|
|
|
pub fn socket(tx: broadcast::Sender<Event>, print: bool) -> Self {
|
|
Self { tx: Some(tx), print }
|
|
}
|
|
|
|
pub fn emit(&self, e: Event) {
|
|
// Also log it. Scans and downloads travel as events, not tracing calls, so
|
|
// without this the log view shows only startup and HTTP lines and none of the
|
|
// work the daemon is actually doing. Progress goes to debug: it fires on every
|
|
// whole percent and would otherwise crowd everything else out of the buffer.
|
|
// Level by how much it matters. With 80-odd feeds in an OPML subscription, one
|
|
// line per feed per tick for "not due yet" would push everything worth reading
|
|
// out of the buffer within a few minutes.
|
|
log_event(&e, self.tx.is_some());
|
|
if let Some(tx) = &self.tx {
|
|
// An error here only means nobody is listening yet.
|
|
let _ = tx.send(e.clone());
|
|
}
|
|
if self.print && let Some(line) = e.human() {
|
|
println!("{line}");
|
|
}
|
|
}
|
|
}
|
|
|
|
/// An event, logged once (#91): its words as the message, and `ev`, `feed`, `new` and the rest as
|
|
/// fields, so Loki reads them without parsing the message. A failure also gets `error.type` and,
|
|
/// from an HTTP error, `http.response.status_code`, so failures group by kind without a regex.
|
|
/// Level by how much it matters: with 80-odd feeds in an OPML subscription, a line per feed per
|
|
/// tick for "not due yet" would push everything worth reading out of the log view in minutes,
|
|
/// and Progress fires on every whole percent. `wire` also logs the event as it goes on the
|
|
/// socket, at debug: the admin page's Daemon I/O tab shows it, production's log leaves it out.
|
|
fn log_event(e: &Event, wire: bool) {
|
|
let json = serde_json::to_string(e).unwrap_or_default();
|
|
let v: serde_json::Value = serde_json::from_str(&json).unwrap_or_default();
|
|
let s = |k: &str| v.get(k).and_then(|x| x.as_str());
|
|
let n = |k: &str| v.get(k).and_then(|x| x.as_u64());
|
|
let (kind, code) = match e {
|
|
Event::FeedError { msg, .. } | Event::DownloadError { msg, .. } | Event::Error { msg } => {
|
|
let (k, c) = crate::feed::failure_kind(msg);
|
|
(Some(k), c)
|
|
}
|
|
_ => (None, None),
|
|
};
|
|
let routine = match e {
|
|
Event::Progress { .. } | Event::FeedSkip { .. } | Event::FeedStart { .. } => true,
|
|
Event::FeedDone { new, downloaded, failed, torrents, .. } => {
|
|
*new == 0 && *downloaded == 0 && *failed == 0 && *torrents == 0
|
|
}
|
|
_ => false,
|
|
};
|
|
macro_rules! line {
|
|
($level:ident, $target:literal, $text:expr) => {
|
|
tracing::$level!(
|
|
target: $target,
|
|
ev = s("ev"), feed = s("feed"), msg = s("msg"), url = s("url"), reason = s("reason"), from = s("from"), to = s("to"),
|
|
new = n("new"), downloaded = n("downloaded"), failed = n("failed"),
|
|
torrents = n("torrents"), bytes = n("bytes"), feeds = n("feeds"),
|
|
pending = n("pending"), enclosure = n("enclosure"), files = n("files"),
|
|
"error.type" = kind, "http.response.status_code" = code,
|
|
"{}", $text
|
|
)
|
|
};
|
|
}
|
|
// The healthcheck's answer, every 30s: a reply on the socket rather than work done, so it
|
|
// stays with the rest of the conversation, and the Scans tab stays about scans.
|
|
if matches!(e, Event::Status { .. }) {
|
|
return line!(info, "ipx::io", format!("<- {json}"));
|
|
}
|
|
if wire {
|
|
tracing::debug!(target: "ipx::io", "<- {json}");
|
|
}
|
|
let text = e.human().map(|l| l.trim().to_owned()).unwrap_or_else(|| json.clone());
|
|
if kind.is_some() {
|
|
line!(warn, "ipx::scan", text)
|
|
} else if routine {
|
|
line!(debug, "ipx::scan", text)
|
|
} else {
|
|
line!(info, "ipx::scan", text)
|
|
}
|
|
}
|
|
|
|
/// True when something is already listening -- i.e. a daemon owns this socket.
|
|
pub async fn daemon_is_live(path: &Path) -> bool {
|
|
UnixStream::connect(path).await.is_ok()
|
|
}
|
|
|
|
/// Answers `status` and `reload` for the socket, without the worker. The worker runs one job at a
|
|
/// time, and a healthcheck left waiting behind a scan or a long download timed out and called a
|
|
/// busy daemon dead; a reload left waiting would let the web UI write the old catalogue back in
|
|
/// the meantime. The answer goes to the client that asked and no one else: broadcast, it ended
|
|
/// any `ipx fetch` that was watching a scan, since `status` is a terminal event.
|
|
/// A future, since reading the counts is a database query.
|
|
pub type StatusFn = std::sync::Arc<
|
|
dyn Fn(Command) -> std::pin::Pin<Box<dyn std::future::Future<Output = Event> + Send>> + Send + Sync,
|
|
>;
|
|
|
|
/// Accepts connections, feeding commands to `cmds` and events from `events` back out.
|
|
pub async fn serve(
|
|
path: PathBuf,
|
|
events: broadcast::Sender<Event>,
|
|
cmds: mpsc::Sender<Command>,
|
|
status: StatusFn,
|
|
) -> Result<()> {
|
|
// A socket file left by a crashed daemon would block the bind; a live one was already
|
|
// rejected by the caller's daemon_is_live() check.
|
|
if path.exists() {
|
|
std::fs::remove_file(&path)
|
|
.with_context(|| format!("removing stale socket {}", path.display()))?;
|
|
}
|
|
if let Some(dir) = path.parent() {
|
|
std::fs::create_dir_all(dir)?;
|
|
}
|
|
let listener = UnixListener::bind(&path)
|
|
.with_context(|| format!("binding {}", path.display()))?;
|
|
tracing::info!(socket = %path.display(), "listening");
|
|
|
|
loop {
|
|
let (stream, _) = listener.accept().await?;
|
|
let rx = events.subscribe();
|
|
let cmds = cmds.clone();
|
|
let status = status.clone();
|
|
tokio::spawn(async move {
|
|
if let Err(e) = handle(stream, rx, cmds, status).await {
|
|
tracing::debug!(error = %e, "client gone");
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
async fn handle(
|
|
stream: UnixStream,
|
|
mut rx: broadcast::Receiver<Event>,
|
|
cmds: mpsc::Sender<Command>,
|
|
status: StatusFn,
|
|
) -> Result<()> {
|
|
let (read, mut write) = stream.into_split();
|
|
|
|
// Events out: everything broadcast, and the answers meant for this client alone.
|
|
let (reply, mut replies) = mpsc::channel::<Event>(4);
|
|
let writer = tokio::spawn(async move {
|
|
loop {
|
|
let ev = tokio::select! {
|
|
Some(ev) = replies.recv() => ev,
|
|
got = rx.recv() => match got {
|
|
Ok(ev) => ev,
|
|
Err(_) => break,
|
|
},
|
|
};
|
|
let mut line = serde_json::to_string(&ev).unwrap_or_default();
|
|
line.push('\n');
|
|
if write.write_all(line.as_bytes()).await.is_err() {
|
|
break;
|
|
}
|
|
}
|
|
});
|
|
|
|
// Commands in.
|
|
let mut lines = BufReader::new(read).lines();
|
|
while let Some(line) = lines.next_line().await? {
|
|
let line = line.trim();
|
|
if line.is_empty() {
|
|
continue;
|
|
}
|
|
match serde_json::from_str::<Command>(line) {
|
|
// Answered here, not queued behind whatever the worker is on: see StatusFn.
|
|
Ok(cmd @ (Command::Status | Command::Reload)) => {
|
|
tracing::info!(target: "ipx::io", "-> {line}");
|
|
let ev = status(cmd).await;
|
|
log_event(&ev, true);
|
|
let _ = reply.send(ev).await;
|
|
}
|
|
Ok(cmd) => {
|
|
if cmds.send(cmd).await.is_err() {
|
|
break; // Worker is gone; so are we.
|
|
}
|
|
}
|
|
Err(e) => tracing::warn!(error = %e, line, "bad command"),
|
|
}
|
|
}
|
|
writer.abort();
|
|
Ok(())
|
|
}
|
|
|
|
/// Tells a running daemon to read the catalogue again (see Command::Reload), and waits for it.
|
|
pub async fn reload(path: &Path) -> Result<()> {
|
|
let stream = UnixStream::connect(path).await?;
|
|
let (read, mut write) = stream.into_split();
|
|
write.write_all(b"{\"cmd\":\"reload\"}\n").await?;
|
|
let mut lines = BufReader::new(read).lines();
|
|
let answer = async {
|
|
// Other clients' events come here too; the answer is the status, or an error.
|
|
while let Some(line) = lines.next_line().await? {
|
|
match serde_json::from_str::<Event>(&line) {
|
|
Ok(Event::Status { .. }) => return Ok(()),
|
|
Ok(Event::Error { msg }) => anyhow::bail!("{msg}"),
|
|
_ => continue,
|
|
}
|
|
}
|
|
anyhow::bail!("the daemon closed the connection without answering")
|
|
};
|
|
tokio::time::timeout(std::time::Duration::from_secs(10), answer).await.context("the daemon did not answer in 10s")?
|
|
}
|
|
|
|
/// Sends one command to a running daemon and prints the events it produces.
|
|
pub async fn proxy(path: &Path, cmd: &Command) -> Result<()> {
|
|
let stream = UnixStream::connect(path).await?;
|
|
let (read, mut write) = stream.into_split();
|
|
let mut line = serde_json::to_string(cmd)?;
|
|
line.push('\n');
|
|
write.write_all(line.as_bytes()).await?;
|
|
|
|
// ponytail: the event stream is a broadcast, so a busy daemon's other work shows up
|
|
// here too. Fine for a CLI; a UI that cares would want per-request ids.
|
|
let mut lines = BufReader::new(read).lines();
|
|
while let Some(line) = lines.next_line().await? {
|
|
let Ok(ev) = serde_json::from_str::<Event>(&line) else {
|
|
continue;
|
|
};
|
|
if let Some(text) = ev.human() {
|
|
println!("{text}");
|
|
}
|
|
if ev.is_terminal() {
|
|
break;
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn a_feed_that_moved_says_so_on_the_wire_and_in_words() {
|
|
let e = Event::FeedMoved { feed: "x".into(), from: "http://a/f".into(), to: "https://a/f".into() };
|
|
let wire = serde_json::to_value(&e).unwrap();
|
|
assert_eq!((wire["ev"].as_str(), wire["from"].as_str(), wire["to"].as_str()), (Some("feed_moved"), Some("http://a/f"), Some("https://a/f")));
|
|
assert_eq!(e.human().unwrap(), "x: moved for good from http://a/f to https://a/f; following it");
|
|
}
|
|
|
|
#[test]
|
|
fn commands_parse_from_the_wire_form() {
|
|
let got: Command = serde_json::from_str(r#"{"cmd":"fetch"}"#).unwrap();
|
|
assert!(matches!(got, Command::Fetch { feed: None, force: false, .. }));
|
|
|
|
let got: Command = serde_json::from_str(r#"{"cmd":"fetch","feed":"atp","force":true}"#).unwrap();
|
|
assert!(matches!(got, Command::Fetch { feed: Some(f), force: true, .. } if f == "atp"));
|
|
|
|
let got: Command = serde_json::from_str(r#"{"cmd":"reap","dry_run":true}"#).unwrap();
|
|
assert!(matches!(got, Command::Reap { dry_run: true }));
|
|
|
|
// "Download this one now" is its own command precisely because a scan cannot
|
|
// express it: a scan takes the lowest-id pending rows up to max_new_per_check.
|
|
let got: Command = serde_json::from_str(r#"{"cmd":"download","enclosure":11}"#).unwrap();
|
|
assert!(matches!(got, Command::Download { enclosure: 11 }));
|
|
|
|
assert!(serde_json::from_str::<Command>(r#"{"cmd":"nope"}"#).is_err());
|
|
}
|
|
|
|
#[test]
|
|
fn events_serialise_to_the_documented_shape() {
|
|
let ev = Event::Progress {
|
|
feed: "atp".into(),
|
|
enclosure: 42,
|
|
url: "https://x/ep.mp3".into(),
|
|
file: "ep.mp3".into(),
|
|
done: 10_485_760,
|
|
total: Some(52_428_800),
|
|
};
|
|
let json = serde_json::to_string(&ev).unwrap();
|
|
assert!(json.starts_with(r#"{"ev":"progress""#), "got {json}");
|
|
assert!(json.contains(r#""done":10485760"#));
|
|
assert!(json.contains(r#""enclosure":42"#), "a UI needs this to target one row");
|
|
|
|
// total is omitted rather than null when the server sent no length.
|
|
let ev = Event::Progress {
|
|
feed: "a".into(),
|
|
enclosure: 1,
|
|
url: "u".into(),
|
|
file: "f".into(),
|
|
done: 1,
|
|
total: None,
|
|
};
|
|
assert!(!serde_json::to_string(&ev).unwrap().contains("total"));
|
|
}
|
|
|
|
#[test]
|
|
fn only_completion_events_end_a_client_session() {
|
|
assert!(Event::ScanDone { feeds: 1 }.is_terminal());
|
|
assert!(Event::ReapDone { files: 0, bytes: 0 }.is_terminal());
|
|
assert!(!Event::FeedDone {
|
|
feed: "a".into(),
|
|
new: 0,
|
|
downloaded: 0,
|
|
failed: 0,
|
|
torrents: 0
|
|
}
|
|
.is_terminal());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn status_is_answered_while_the_worker_is_busy() {
|
|
// The queue is full and nobody drains it, as when the worker is deep in a long download:
|
|
// anything sent to it would wait for ever.
|
|
let (cmds, _worker) = mpsc::channel::<Command>(1);
|
|
cmds.send(Command::Reap { dry_run: true }).await.unwrap();
|
|
let (events, _) = broadcast::channel::<Event>(8);
|
|
// Another client, watching a scan: it must not be handed someone else's answer, which
|
|
// would end its session.
|
|
let mut watcher = events.subscribe();
|
|
let status: StatusFn =
|
|
std::sync::Arc::new(|_| Box::pin(async { Event::Status { feeds: 1, pending: 2, downloaded: 3 } }));
|
|
let (client, server) = UnixStream::pair().unwrap();
|
|
tokio::spawn(handle(server, events.subscribe(), cmds, status));
|
|
|
|
let (read, mut write) = client.into_split();
|
|
write.write_all(b"{\"cmd\":\"status\"}\n").await.unwrap();
|
|
let line = tokio::time::timeout(std::time::Duration::from_secs(2), BufReader::new(read).lines().next_line())
|
|
.await
|
|
.expect("status waited behind the worker")
|
|
.unwrap()
|
|
.unwrap();
|
|
assert!(line.contains(r#""ev":"status""#), "{line}");
|
|
assert!(watcher.try_recv().is_err(), "the answer went to every client, not just the one asking");
|
|
}
|
|
}
|