From 5b1e8b742b072480932414fef05e3731ffa8b31c Mon Sep 17 00:00:00 2001 From: Sebastian Date: Sat, 5 Sep 2026 23:09:33 +0200 Subject: [PATCH] rclone mirror: live progress via --progress instead of --stats-one-line --stats-one-line is carriage-return-rewritten when stdout/stderr are piped (non-TTY), so the cloud->local mirror step showed no progress and jumped straight to 100% at the end. --progress emits newline-terminated Transferred: stats that stream live; also strip the 'Transferred:' label from the done column in parse_rclone_stats. --- src/engine.rs | 183 +++++++++++++++++++++++++++--------------------- src/progress.rs | 62 ++++++++++++++-- 2 files changed, 160 insertions(+), 85 deletions(-) diff --git a/src/engine.rs b/src/engine.rs index 8583ae5..7861eac 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -3,7 +3,6 @@ use crate::progress::{parse_rclone_stats, parse_rsync_line, RsyncLine}; use chrono::Local; use std::collections::HashMap; use std::fs; -use std::io::Read; use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; use std::sync::atomic::{AtomicBool, Ordering}; @@ -394,17 +393,17 @@ fn run_mirror(config: &Config, device: &Device, ctx: &WorkerCtx) -> Result<(), S args.push("--exclude".to_string()); args.push(pat.clone()); } + // NOTE: do NOT use `--stats-one-line` here. When stdout/stderr are not a + // TTY (which is always the case for us — the child is piped), rclone emits + // the one-line stats as a single carriage-return-rewritten line that only + // surfaces once at the very end (100%). `--progress` instead emits the + // multi-line `Transferred: … / …, N%, speed, ETA …` block newline-terminated, + // so it streams live through the pipe and we can update the gauge in real + // time while a cloud→local mirror is running. args.extend( - [ - "--create-empty-src-dirs", - "--stats", - "1s", - "--stats-one-line", - "--stats-log-level", - "NOTICE", - ] - .iter() - .map(|s| s.to_string()), + ["--create-empty-src-dirs", "--progress", "--stats", "1s"] + .iter() + .map(|s| s.to_string()), ); stream(ctx, "rclone", &args, |line| { @@ -709,9 +708,48 @@ fn resolve_fatsort_target(mount: &Path) -> Result { // ── subprocess plumbing ─────────────────────────────────────────────────────── +/// Read lines from a subprocess pipe and forward them (tagged by origin) to `tx`. +fn pump_lines( + mut reader: R, + is_stderr: bool, + tx: mpsc::Sender<(bool, String)>, +) { + let mut buf = [0u8; 8192]; + let mut partial: Vec = Vec::new(); + loop { + match reader.read(&mut buf) { + Ok(0) | Err(_) => break, + Ok(n) => { + partial.extend_from_slice(&buf[..n]); + let mut pos = 0; + for (i, b) in partial.iter().enumerate() { + if *b == b'\n' || *b == b'\r' { + let line = String::from_utf8_lossy(&partial[pos..i]).trim().to_string(); + if !line.is_empty() { + let _ = tx.send((is_stderr, line)); + } + pos = i + 1; + } + } + partial.drain(0..pos); + } + } + } + if !partial.is_empty() { + let line = String::from_utf8_lossy(&partial).trim().to_string(); + if !line.is_empty() { + let _ = tx.send((is_stderr, line)); + } + } +} + /// Run a command, streaming its output lines to `on_line`. Returns Err on /// failure or abort. Lines returned `true` from `on_line` are treated as /// consumed (progress); unconsumed lines are logged. +/// +/// Both stdout and stderr are pumped into a shared channel so lines from +/// stderr (e.g. rclone stats) are consumed live instead of being buffered +/// until the process exits. fn stream( ctx: &WorkerCtx, cmd: &str, @@ -729,86 +767,43 @@ where .spawn() .map_err(|e| format!("could not start {cmd}: {e}"))?; - let mut stdout = child.stdout.take().unwrap(); - let mut stderr = child.stderr.take().unwrap(); + let stdout = child.stdout.take().unwrap(); + let stderr = child.stderr.take().unwrap(); - let (err_tx, err_rx) = mpsc::channel::(); - let err_handle = thread::spawn(move || { - let mut buf = [0u8; 8192]; - let mut partial: Vec = Vec::new(); - loop { - match stderr.read(&mut buf) { - Ok(0) | Err(_) => break, - Ok(n) => { - partial.extend_from_slice(&buf[..n]); - let mut pos = 0; - for (i, b) in partial.iter().enumerate() { - if *b == b'\n' || *b == b'\r' { - let line = String::from_utf8_lossy(&partial[pos..i]).trim().to_string(); - if !line.is_empty() { - let _ = err_tx.send(line); - } - pos = i + 1; - } - } - partial.drain(0..pos); - } - } - } - if !partial.is_empty() { - let line = String::from_utf8_lossy(&partial).trim().to_string(); - if !line.is_empty() { - let _ = err_tx.send(line); - } - } - }); + let (tx, rx) = mpsc::channel::<(bool, String)>(); // (is_stderr, line) + let out_tx = tx.clone(); + let out_thread = thread::spawn(move || pump_lines(stdout, false, out_tx)); + let err_thread = thread::spawn(move || pump_lines(stderr, true, tx)); - let mut buf = [0u8; 8192]; - let mut partial: Vec = Vec::new(); loop { if ctx.abort.load(Ordering::Relaxed) { let _ = child.kill(); let _ = child.wait(); return Err("aborted".to_string()); } - match stdout.read(&mut buf) { - Ok(0) => break, - Err(_) => break, - Ok(n) => { - partial.extend_from_slice(&buf[..n]); - let mut pos = 0; - for (i, b) in partial.iter().enumerate() { - if *b == b'\n' || *b == b'\r' { - let line = String::from_utf8_lossy(&partial[pos..i]).trim().to_string(); - if !line.is_empty() && !on_line(&line) { - ctx.log(LogLevel::Info, line); - } - pos = i + 1; + match rx.recv_timeout(std::time::Duration::from_millis(200)) { + Ok((is_stderr, line)) => { + if !line.is_empty() && !on_line(&line) { + if is_stderr { + let lower = line.to_lowercase(); + let level = if lower.contains("error") || lower.contains("failed") { + LogLevel::Error + } else { + LogLevel::Info + }; + ctx.log(level, line); + } else { + ctx.log(LogLevel::Info, line); } } - partial.drain(0..pos); } - } - } - if !partial.is_empty() { - let line = String::from_utf8_lossy(&partial).trim().to_string(); - if !line.is_empty() && !on_line(&line) { - ctx.log(LogLevel::Info, line); + Err(mpsc::RecvTimeoutError::Timeout) => continue, + Err(mpsc::RecvTimeoutError::Disconnected) => break, } } - let _ = err_handle.join(); - for line in err_rx.try_iter() { - if !line.is_empty() && !on_line(&line) { - let lower = line.to_lowercase(); - let level = if lower.contains("error") || lower.contains("failed") { - LogLevel::Error - } else { - LogLevel::Info - }; - ctx.log(level, line); - } - } + let _ = out_thread.join(); + let _ = err_thread.join(); let status = child.wait().map_err(|e| format!("{cmd} failed: {e}"))?; if ctx.abort.load(Ordering::Relaxed) { @@ -825,5 +820,37 @@ where #[cfg(test)] mod tests { - // engine logic is exercised through parser + integration tests + use super::*; + use std::sync::atomic::AtomicBool; + + #[test] + fn stream_consumes_stderr_lines_live() { + // rclone prints stats to stderr; they must reach on_line while the + // process is still running (not only after exit). + let (tx, _rx) = mpsc::channel::(); + let ctx = Arc::new(WorkerCtx { + state: Arc::new(Mutex::new(SyncState::default())), + tx, + abort: Arc::new(AtomicBool::new(false)), + }); + let args: Vec = vec![ + "-c".to_string(), + "echo 'NOTICE: first stats line'; sleep 0.5; echo 'NOTICE: second stats line'".to_string(), + ]; + let seen: Arc>> = Arc::new(Mutex::new(Vec::new())); + let seen2 = seen.clone(); + let result = stream(&ctx, "sh", &args, move |line| { + if line.starts_with("NOTICE:") { + seen2.lock().unwrap().push(line.to_string()); + true + } else { + false + } + }); + assert!(result.is_ok()); + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 2, "stderr lines should arrive live"); + assert!(seen[0].contains("first stats line")); + assert!(seen[1].contains("second stats line")); + } } diff --git a/src/progress.rs b/src/progress.rs index 523bffe..239c060 100644 --- a/src/progress.rs +++ b/src/progress.rs @@ -155,31 +155,50 @@ pub struct RcloneStats { pub eta: String, } -/// Parse a rclone `--stats-one-line` "Transferred:" line. +/// Parse a rclone stats line. Two shapes are supported: +/// +/// 1. `--progress` (used by the engine): a `Transferred:` line such as +/// `Transferred:\t 20.141 MiB / 500 MiB, 4%, 20.139 MiB/s, ETA 23s`. +/// These are newline-terminated and stream live even when piped. +/// +/// 2. `--stats-one-line` (legacy): rclone prefixes each line with +/// "timestamp LEVEL:" (e.g. "2026/08/22 13:27:19 NOTICE:"), which we strip +/// before parsing the comma-separated body. pub fn parse_rclone_stats(line: &str) -> Option { let line = line.trim(); + + // Fast reject: the stats body always carries "x / y" and a percentage. if !line.contains(" / ") || !line.contains('%') { return None; } + let body = line - .split_once(':') + .split_once("NOTICE:") + .or_else(|| line.split_once("INFO:")) + .or_else(|| line.split_once("DEBUG:")) + .or_else(|| line.split_once("ERROR:")) .map(|(_, rest)| rest.trim()) .unwrap_or(line); if !body.contains(" / ") { return None; } + let mut parts = body.split(','); let sizes = parts.next()?.trim(); let (done, total) = sizes.split_once('/')?; - let done = done.trim().to_string(); + // `--progress` lines carry a leading `Transferred:` label; strip it so the + // gauge shows a clean byte amount instead of "Transferred: 512.000 MiB". + let done = done + .trim() + .strip_prefix("Transferred:") + .map(|s| s.trim()) + .unwrap_or_else(|| done.trim()) + .to_string(); let total = total.trim().to_string(); let pct = parts .next() - .and_then(|p| { - let p = p.trim().trim_end_matches('%'); - p.parse::().ok() - }) + .and_then(|p| p.trim().trim_end_matches('%').parse::().ok()) .unwrap_or(0) .min(100); @@ -210,6 +229,35 @@ mod tests { assert_eq!(s.eta, "20s"); } + #[test] + fn parses_rclone_progress_line_with_tabs() { + // real `rclone --progress` output when piped (non-TTY): the done column + // is padded with spaces and tabs after the `Transferred:` label. + let s = parse_rclone_stats( + "Transferred: \t 20.141 MiB / 500 MiB, 4%, 20.139 MiB/s, ETA 23s", + ) + .unwrap(); + assert_eq!(s.done, "20.141 MiB"); + assert_eq!(s.total, "500 MiB"); + assert_eq!(s.pct, 4); + assert_eq!(s.speed, "20.139 MiB/s"); + assert_eq!(s.eta, "23s"); + } + + #[test] + fn parses_rclone_stats_with_log_prefix() { + // real `rclone copy --stats-one-line --stats-log-level NOTICE` output + let s = parse_rclone_stats( + "2026/08/22 13:27:21 NOTICE: 13.996 MiB / 20 MiB, 70%, 5.497 MiB/s, ETA 1s", + ) + .unwrap(); + assert_eq!(s.done, "13.996 MiB"); + assert_eq!(s.total, "20 MiB"); + assert_eq!(s.pct, 70); + assert_eq!(s.speed, "5.497 MiB/s"); + assert_eq!(s.eta, "1s"); + } + #[test] fn rejects_non_stats_line() { assert!(parse_rclone_stats("Elapsed time: 5.0s").is_none());