Have ipx add, rm and import tell a running daemon to read the catalogue again (#120)

`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>
This commit is contained in:
2026-10-05 19:55:22 +00:00
parent 9c2f3eaf1b
commit a553d050ca
5 changed files with 111 additions and 12 deletions

View File

@@ -109,6 +109,11 @@ pub enum Command {
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.
@@ -208,13 +213,15 @@ pub async fn daemon_is_live(path: &Path) -> bool {
UnixStream::connect(path).await.is_ok()
}
/// Answers `status` 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. 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.
/// 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() -> std::pin::Pin<Box<dyn std::future::Future<Output = Event> + Send>> + Send + Sync>;
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(
@@ -285,9 +292,9 @@ async fn handle(
}
match serde_json::from_str::<Command>(line) {
// Answered here, not queued behind whatever the worker is on: see StatusFn.
Ok(Command::Status) => {
Ok(cmd @ (Command::Status | Command::Reload)) => {
tracing::info!(target: "ipx::io", "-> {line}");
let ev = status().await;
let ev = status(cmd).await;
log_event(&ev, true);
let _ = reply.send(ev).await;
}
@@ -303,6 +310,26 @@ async fn handle(
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?;
@@ -411,7 +438,7 @@ mod tests {
// 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 } }));
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));