rsync step: smooth byte-level progress + scan feedback

- Parse rsync live per-file byte progress (\r segments without xfr#) so the
  progress bar advances continuously during large file transfers instead of
  only per completed file; also stops logging those lines as noise
- Update the active step detail during scan/remove/copy phases
- smooth_pct helper combines completed files + current file's byte fraction
- unit tests for byte-progress parsing and smooth_pct
This commit is contained in:
Sebastian
2026-08-22 09:42:54 +02:00
parent 16a5a4327e
commit cf2285fd56
2 changed files with 92 additions and 6 deletions
+57 -5
View File
@@ -129,7 +129,7 @@ impl WorkerCtx {
self.changed(); self.changed();
} }
fn set_device_progress(&self, done: u64, total: u64, speed: &str, eta: &str, current: &str) { fn set_device_progress(&self, done: u64, total: u64, speed: &str, eta: &str, current: &str, pct: Option<u16>) {
if let Ok(mut s) = self.state.lock() { if let Ok(mut s) = self.state.lock() {
let p = &mut s.device_prog; let p = &mut s.device_prog;
p.done = done.to_string(); p.done = done.to_string();
@@ -139,7 +139,9 @@ impl WorkerCtx {
if !current.is_empty() { if !current.is_empty() {
p.current = current.to_string(); p.current = current.to_string();
} }
if total > 0 { if let Some(pct) = pct {
p.pct = pct.min(100);
} else if total > 0 {
p.pct = ((done as f64 / total as f64) * 100.0) as u16; p.pct = ((done as f64 / total as f64) * 100.0) as u16;
} }
} }
@@ -353,6 +355,14 @@ fn run_mirror(config: &Config, device: &Device, ctx: &WorkerCtx) -> Result<(), S
} }
true true
} }
RsyncLine::ByteProgress(p) => {
if let Ok(mut s) = ctx.state.lock() {
s.mirror.pct = p.pct;
s.mirror.speed = p.speed.clone();
s.mirror.eta = p.eta.clone();
}
true
}
RsyncLine::Total(_) | RsyncLine::File(_) => true, RsyncLine::Total(_) | RsyncLine::File(_) => true,
} }
} else { } else {
@@ -475,8 +485,10 @@ fn run_rsync_device(config: &Config, device: &Device, ctx: &WorkerCtx) -> Result
fs::create_dir_all(&dest).map_err(|e| format!("cannot create {}: {e}", dest.display()))?; fs::create_dir_all(&dest).map_err(|e| format!("cannot create {}: {e}", dest.display()))?;
ctx.log(LogLevel::Info, format!("Scanning local {label} library…")); ctx.log(LogLevel::Info, format!("Scanning local {label} library…"));
ctx.set_step(1, StepStatus::Active, format!("scanning local {label} library…"));
let src_map = scan_music_dir(&src_dir); let src_map = scan_music_dir(&src_dir);
ctx.log(LogLevel::Info, format!("Scanning device {}…", dest.display())); ctx.log(LogLevel::Info, format!("Scanning device {}…", dest.display()));
ctx.set_step(1, StepStatus::Active, format!("scanning device {}…", dest.display()));
let dst_map = scan_music_dir(&dest); let dst_map = scan_music_dir(&dest);
let mut to_copy: Vec<String> = Vec::new(); let mut to_copy: Vec<String> = Vec::new();
@@ -510,11 +522,13 @@ fn run_rsync_device(config: &Config, device: &Device, ctx: &WorkerCtx) -> Result
} }
if !to_delete.is_empty() { if !to_delete.is_empty() {
ctx.set_step(1, StepStatus::Active, format!("removing {} stale file(s)…", to_delete.len()));
remove_stale(&dest, &to_delete); remove_stale(&dest, &to_delete);
ctx.log(LogLevel::Info, format!("Removed {} stale file(s).", to_delete.len())); ctx.log(LogLevel::Info, format!("Removed {} stale file(s).", to_delete.len()));
} }
if !to_copy.is_empty() { if !to_copy.is_empty() {
ctx.set_step(1, StepStatus::Active, format!("copying {} file(s)…", to_copy.len()));
let tmp = write_file_list(&to_copy)?; let tmp = write_file_list(&to_copy)?;
let res = { let res = {
let args = vec![ let args = vec![
@@ -528,16 +542,26 @@ fn run_rsync_device(config: &Config, device: &Device, ctx: &WorkerCtx) -> Result
dest.display().to_string(), dest.display().to_string(),
]; ];
let total = to_copy.len() as u64; let total = to_copy.len() as u64;
let mut xfr_done: u64 = 0;
let mut cur_pct: u16 = 0;
stream(ctx, "rsync", &args, |line| { stream(ctx, "rsync", &args, |line| {
match parse_rsync_line(line) { match parse_rsync_line(line) {
Some(RsyncLine::Progress(p)) => { Some(RsyncLine::Progress(p)) => {
let done = if p.done > 0 { p.done } else { total }; xfr_done = p.done.max(1);
ctx.set_device_progress(done, total, &p.speed, &p.eta, ""); cur_pct = 100;
let pct = smooth_pct(xfr_done, cur_pct, total);
ctx.set_device_progress(xfr_done, total, &p.speed, &p.eta, "", Some(pct));
true
}
Some(RsyncLine::ByteProgress(p)) => {
cur_pct = p.pct;
let pct = smooth_pct(xfr_done, cur_pct, total);
ctx.set_device_progress(xfr_done, total, &p.speed, &p.eta, "", Some(pct));
true true
} }
Some(RsyncLine::Total(_)) => true, Some(RsyncLine::Total(_)) => true,
Some(RsyncLine::File(path)) => { Some(RsyncLine::File(path)) => {
ctx.set_device_progress(0, total, "", "", &path); ctx.set_device_progress(0, total, "", "", &path, None);
ctx.log(LogLevel::Transfer, path); ctx.log(LogLevel::Transfer, path);
true true
} }
@@ -602,6 +626,16 @@ fn write_file_list(items: &[String]) -> Result<String, String> {
// ── mlove finalize ──────────────────────────────────────────────────────────── // ── mlove finalize ────────────────────────────────────────────────────────────
/// Combine completed files with the current file's byte fraction for a smooth
/// overall percentage, e.g. 3 files done + file at 50% over 5 total → 70%.
fn smooth_pct(xfr_done: u64, cur_pct: u16, total: u64) -> u16 {
if total == 0 {
return 0;
}
let ratio = xfr_done as f64 + (cur_pct as f64 / 100.0);
((ratio / total as f64) * 100.0).clamp(0.0, 100.0) as u16
}
fn capture(cmd: &str, args: &[&str]) -> Result<(i32, String), String> { fn capture(cmd: &str, args: &[&str]) -> Result<(i32, String), String> {
let out = Command::new(cmd) let out = Command::new(cmd)
.args(args) .args(args)
@@ -774,3 +808,21 @@ where
} }
Ok(()) Ok(())
} }
#[cfg(test)]
mod tests {
use super::smooth_pct;
#[test]
fn smooth_progress_advances_within_a_file() {
// 3 files done, current at 50%, 5 total → (3 + 0.5)/5 = 70%
assert_eq!(smooth_pct(3, 50, 5), 70);
// 3 files done, current at 0% → 60%
assert_eq!(smooth_pct(3, 0, 5), 60);
// 1 file done, current at 100% → 40% (xfr#1 line arrives before next file)
assert_eq!(smooth_pct(1, 100, 5), 40);
// everything done
assert_eq!(smooth_pct(5, 100, 5), 100);
assert_eq!(smooth_pct(0, 100, 0), 0);
}
}
+35 -1
View File
@@ -6,8 +6,10 @@ const AUDIO_EXTS: &[&str] = &[".m4a", ".mp3", ".aac", ".ogg", ".wav", ".wv", ".a
pub enum RsyncLine { pub enum RsyncLine {
/// Total file count announced by rsync ("Transfer starting: N files") /// Total file count announced by rsync ("Transfer starting: N files")
Total(u64), Total(u64),
/// A per-file progress update /// A per-file completion update ("… (xfr#N, to-check=M/T)")
Progress(RsyncProgress), Progress(RsyncProgress),
/// Live byte progress of the current file ("bytes pct% speed ETA")
ByteProgress(RsyncProgress),
/// A transferred audio file path (current file being handled) /// A transferred audio file path (current file being handled)
File(String), File(String),
} }
@@ -113,6 +115,26 @@ pub fn parse_rsync_line(line: &str) -> Option<RsyncLine> {
})); }));
} }
// Live byte progress of the current file: " 32.768 0% 999,51MB/s 0:00:00 "
// (no xfr# yet — rsync emits these with \r as the file transfers).
if line.contains('%') {
let starts_with_digit = line
.chars()
.next()
.map(|c| c.is_ascii_digit())
.unwrap_or(false);
let speed = find_speed(line);
if starts_with_digit || !speed.is_empty() {
return Some(RsyncLine::ByteProgress(RsyncProgress {
done: 0,
total: 0,
pct: find_pct(line).unwrap_or(0),
speed,
eta: find_eta(line),
}));
}
}
// Plain audio file path line (transferred file being processed). // Plain audio file path line (transferred file being processed).
if ends_with_audio_ext(line) { if ends_with_audio_ext(line) {
// Skip stats-looking lines that just happen to end in an audio ext. // Skip stats-looking lines that just happen to end in an audio ext.
@@ -222,4 +244,16 @@ mod tests {
other => panic!("expected file, got {other:?}"), other => panic!("expected file, got {other:?}"),
} }
} }
#[test]
fn parses_live_byte_progress_without_xfr() {
// rsync \r segment during a large-file transfer
match parse_rsync_line(" 32.768 0% 0,00kB/s 0:00:00") {
Some(RsyncLine::ByteProgress(p)) => {
assert_eq!(p.pct, 0);
assert_eq!(p.speed, "0,00kB/s");
}
other => panic!("expected byte progress, got {other:?}"),
}
}
} }