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.
This commit is contained in:
+99
-72
@@ -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,15 +393,15 @@ 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",
|
||||
]
|
||||
["--create-empty-src-dirs", "--progress", "--stats", "1s"]
|
||||
.iter()
|
||||
.map(|s| s.to_string()),
|
||||
);
|
||||
@@ -709,9 +708,48 @@ fn resolve_fatsort_target(mount: &Path) -> Result<String, String> {
|
||||
|
||||
// ── subprocess plumbing ───────────────────────────────────────────────────────
|
||||
|
||||
/// Read lines from a subprocess pipe and forward them (tagged by origin) to `tx`.
|
||||
fn pump_lines<R: std::io::Read + Send + 'static>(
|
||||
mut reader: R,
|
||||
is_stderr: bool,
|
||||
tx: mpsc::Sender<(bool, String)>,
|
||||
) {
|
||||
let mut buf = [0u8; 8192];
|
||||
let mut partial: Vec<u8> = 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<F>(
|
||||
ctx: &WorkerCtx,
|
||||
cmd: &str,
|
||||
@@ -729,77 +767,24 @@ 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::<String>();
|
||||
let err_handle = thread::spawn(move || {
|
||||
let mut buf = [0u8; 8192];
|
||||
let mut partial: Vec<u8> = 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<u8> = 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;
|
||||
}
|
||||
}
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
let _ = err_handle.join();
|
||||
for line in err_rx.try_iter() {
|
||||
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
|
||||
@@ -807,8 +792,18 @@ where
|
||||
LogLevel::Info
|
||||
};
|
||||
ctx.log(level, line);
|
||||
} else {
|
||||
ctx.log(LogLevel::Info, line);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(mpsc::RecvTimeoutError::Timeout) => continue,
|
||||
Err(mpsc::RecvTimeoutError::Disconnected) => break,
|
||||
}
|
||||
}
|
||||
|
||||
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::<EngineEvent>();
|
||||
let ctx = Arc::new(WorkerCtx {
|
||||
state: Arc::new(Mutex::new(SyncState::default())),
|
||||
tx,
|
||||
abort: Arc::new(AtomicBool::new(false)),
|
||||
});
|
||||
let args: Vec<String> = vec![
|
||||
"-c".to_string(),
|
||||
"echo 'NOTICE: first stats line'; sleep 0.5; echo 'NOTICE: second stats line'".to_string(),
|
||||
];
|
||||
let seen: Arc<Mutex<Vec<String>>> = 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"));
|
||||
}
|
||||
}
|
||||
|
||||
+55
-7
@@ -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<RcloneStats> {
|
||||
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::<u16>().ok()
|
||||
})
|
||||
.and_then(|p| p.trim().trim_end_matches('%').parse::<u16>().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());
|
||||
|
||||
Reference in New Issue
Block a user