diff --git a/CHANGELOG.md b/CHANGELOG.md index 12a950e..cc8d504 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - A feed whose address answers with nothing, as a lapsed domain does behind a DNS filter's block page, shows as failing, with what to do about it, instead of as a show that has not posted yet. - The search box finds episodes in Currently Listening. It did nothing there. +- `ipx add`, `ipx rm` and `ipx import`, run while the daemon runs, reach it at once. It kept its old list of feeds and wrote it back at its next change, undoing them. - On Postgres, a write no longer waits for the database server's disk. A scan read about one feed a second and now reads 1,500 in 20 seconds. - Checking every feed no longer makes every icon in the feed list flash while it runs: the list keeps the icons it has already drawn. diff --git a/src/ipc.rs b/src/ipc.rs index 9884d35..6ebd6eb 100644 --- a/src/ipc.rs +++ b/src/ipc.rs @@ -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 std::pin::Pin + Send>> + Send + Sync>; +pub type StatusFn = std::sync::Arc< + dyn Fn(Command) -> std::pin::Pin + 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::(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::(&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)); diff --git a/src/main.rs b/src/main.rs index 3701f00..37ddafd 100644 --- a/src/main.rs +++ b/src/main.rs @@ -249,9 +249,13 @@ async fn main() -> Result<()> { return ipc::proxy(&cfg.general.socket, cmd).await; } + let socket = cfg.general.socket.clone(); let cfg = assemble_config(&db, cfg, &config_path).await?; let is_daemon = matches!(cli.command, Command::Daemon { .. }); + // These write the catalogue from a process of their own, which a running daemon, holding its + // own copy, would write over at its next change (#120). + let changes_catalogue = matches!(cli.command, Command::Add { .. } | Command::Rm { .. } | Command::Import { .. }); let (events, _) = broadcast::channel(1024); let ctx = Arc::new(Ctx { cfg: std::sync::RwLock::new(std::sync::Arc::new(cfg)), @@ -285,6 +289,11 @@ async fn main() -> Result<()> { Command::Export { file } => export(&ctx, &file).await, _ => run(&ctx, wire_cmd.expect("only List and Daemon have no wire form")).await, }; + if result.is_ok() && changes_catalogue && ipc::daemon_is_live(&socket).await { + if let Err(e) = ipc::reload(&socket).await { + eprintln!("the running daemon did not take the change ({e:#}); restart it, or it may undo it"); + } + } // The batch exporter holds the last few seconds of spans; without this they are lost. if let Some(p) = otel { let _ = p.shutdown(); @@ -452,9 +461,19 @@ async fn run(ctx: &Arc, cmd: Cmd) -> Result<()> { ctx.out.emit(status(ctx).await); Ok(()) } + Cmd::Reload => reload(ctx).await, } } +/// The catalogue and server settings as the database now has them, in place of the copy in +/// memory, for a daemon told another ipx changed them (#120); config.toml, read again, for where +/// things are and who gets in. +async fn reload(ctx: &Ctx) -> Result<()> { + let cfg = assemble_config(&ctx.db, config::Config::load(&ctx.config_path)?, &ctx.config_path).await?; + ctx.set_cfg(cfg); + 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 { @@ -511,12 +530,20 @@ async fn daemon( let (tx_cmd, mut rx_cmd) = mpsc::channel::(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. + // status and reload are answered by the socket itself; everything else waits its turn in the + // queue. let answer: ipc::StatusFn = { let ctx = ctx.clone(); - Arc::new(move || { + Arc::new(move |cmd| { let ctx = ctx.clone(); - Box::pin(async move { status(&ctx).await }) + Box::pin(async move { + if matches!(cmd, Cmd::Reload) + && let Err(e) = reload(&ctx).await + { + return Event::Error { msg: format!("reading the catalogue again: {e:#}") }; + } + status(&ctx).await + }) }) }; let server = tokio::spawn(ipc::serve(socket.clone(), events.clone(), tx_cmd, answer)); @@ -2255,6 +2282,23 @@ mod tests { } } + #[tokio::test] + async fn a_reload_takes_the_catalogue_another_ipx_wrote() { + let with = |id: &str| { + let mut cfg = config::Config::default(); + cfg.feeds.insert(id.into(), config::Feed { url: format!("http://x/{id}.xml"), ..feed() }); + cfg + }; + let held = with("kept-in-memory"); + let ctx = test_ctx(held.clone()).await; + ctx.db.store_config(&held).await.unwrap(); + // `ipx add` and `ipx rm`, from a process of their own (#120). + ctx.db.store_config(&with("added-by-the-cli")).await.unwrap(); + reload(&ctx).await.unwrap(); + let urls: Vec = ctx.cfg().feeds.values().map(|f| f.url.clone()).collect(); + assert_eq!(urls, vec!["http://x/added-by-the-cli.xml".to_string()]); + } + /// Answers /empty with a 200 and nothing, anything else with a feed of one item. async fn empty_or_feed_server() -> String { use tokio::io::{AsyncReadExt, AsyncWriteExt}; diff --git a/tests/ui/app.spec.js b/tests/ui/app.spec.js index db03b93..7257b47 100644 --- a/tests/ui/app.spec.js +++ b/tests/ui/app.spec.js @@ -993,6 +993,28 @@ test('the feed list keeps the icons it has drawn when a feed\'s new row comes in expect(kept).toEqual({ drawn: true, redrawn: true, same: true }); }); +test('ipx add and rm while the daemon runs reach it, and its next change keeps them (#120)', async ({ page }) => { + const { execFileSync } = require('child_process'); + const setup = require('./global-setup'); + const env = { ...process.env, IPX_CONFIG: `${setup.root}/config/config.toml`, IPX_DATA_DIR: `${setup.root}/data` }; + const ipx = (...args) => execFileSync('./target/debug/ipx', args, { env, encoding: 'utf8' }); + const listed = async () => (await page.evaluate(() => api('/api/directory'))).map(p => p.id); + // A change of the daemon's own: it writes the whole catalogue from the copy it holds, which + // wrote a removed feed back and dropped an added one when that copy was old. + const daemonWrites = category => page.evaluate(c => + api('/api/feeds/picture-blog', { method: 'PATCH', body: JSON.stringify({ category: c }) }), category); + + const id = ipx('add', '--list', 'http://127.0.0.1:8792/cli.xml').match(/^added (\S+?),? /m)[1]; + await expect.poll(listed).toContain(id); + await daemonWrites('Comedy'); + expect(await listed()).toContain(id); + + ipx('rm', id); + await expect.poll(listed).not.toContain(id); + await daemonWrites(null); + expect(await listed()).not.toContain(id); +}); + test('adding a feed scans it straight away', async ({ page }) => { await page.locator('#addFeed').click(); await page.locator('#nurl').fill('http://127.0.0.1:8792/fresh.xml'); diff --git a/tests/ui/fixtures/cli.xml b/tests/ui/fixtures/cli.xml new file mode 100644 index 0000000..10426c6 --- /dev/null +++ b/tests/ui/fixtures/cli.xml @@ -0,0 +1,5 @@ + +Command Line Showhttp://127.0.0.1:8792/ +Added with ipx add while the daemon runs. +From the Command Linecli-1x +