commit 718c6fd115aa54938e051bf6a15548b4578422bf
parent fff366c4bfff0ea9bdaa84337985ba9e5eb6ff9d
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sun, 20 Sep 2026 17:44:01 -0400
relocate: rsync already knew how far it had got — say so
`--info=progress2` has been on the copy since the move shipped, and every
frame of it went to the job log as a carriage-return redraw of one line. So a
131 GB move and a 3 MB one looked identical from /jobs: a spinner.
parseRsyncProgress reads one frame (bytes, percent, rate, h:mm:ss ETA);
rsyncProgressFraction divides by the tree the controller ALREADY measured for
its space check rather than by rsync's own percentage, which under incremental
recursion is a percentage of what has been enumerated so far and walks
backwards. formatRsyncProgressDetail is the one wording both readouts use.
relocateChannelMedia gains `onProgress`, keeps the frames out of the log, and
writes one decile line instead — eleven lines where there were thousands. The
job turns the frames into a `relocate` JobTask, added on the FIRST frame so the
preflight and a resume's verify pass (which transfer nothing) do not draw a bar
sitting at zero.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
7 files changed, 342 insertions(+), 3 deletions(-)
diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts
@@ -16,6 +16,11 @@ import {
writeFile,
} from "node:fs/promises";
import { execa } from "execa";
+import {
+ formatRsyncProgressDetail,
+ parseRsyncProgress,
+ rsyncProgressFraction,
+} from "../jobs/progressParsers";
import type { Paths } from "../lib/paths";
import { isSocialChannel } from "../lib/channelConfig";
import { getFreeBytes } from "../lib/diskSpace";
@@ -67,6 +72,29 @@ export type RelocateChannelMediaResult = {
retried: boolean;
};
+// WHAT A COPY IN FLIGHT LOOKS LIKE FROM OUTSIDE, once per rsync redraw.
+//
+// rsync has always printed this (the copy runs with `--info=progress2`) and it
+// has always gone nowhere but the job log, as several thousand carriage-return
+// redraws of one line. This is the same fact as a value: the job turns it into
+// the task bar on /jobs and the panel's own read-out, and the controller logs
+// one line per decile so the log afterwards says how it went without holding
+// every frame of it.
+//
+// `fraction` is against the MEASURED tree, never rsync's own percentage — see
+// parseRsyncProgress for why that one walks backwards under incremental
+// recursion.
+export type RelocationProgress = {
+ bytes: number;
+ totalBytes: number;
+ fraction: number;
+ rate: string;
+ etaSeconds: number;
+ // "12.3 GB of 45.6 GB · 27 % · 110.50MB/s · ETA 5:32". One wording, shared by
+ // the task detail and the decile log line.
+ detail: string;
+};
+
type RelocateOpts = {
paths: Paths;
slug: string;
@@ -74,9 +102,44 @@ type RelocateOpts = {
// Required for "out"; ignored for "back", which reads the target from config.
root?: string;
onLog?: (line: string) => void;
+ // Called on every rsync redraw of the COPY phase (never the verify, which
+ // transfers nothing). Optional: a bin script that only wants the log passes
+ // nothing and pays only the parse.
+ onProgress?: (p: RelocationProgress) => void;
signal?: AbortSignal;
};
+// Turn rsync's redraws into `RelocationProgress`, and into one log line per
+// decile. The decile latch is why this is a closure and not a function: "every
+// 10 %" is a fact about the run, and a copy that resumes at 80 % must not
+// re-announce the eight deciles it did not do.
+function makeProgressSink(opts: {
+ totalBytes: number;
+ log: (m: string) => void;
+ onProgress?: (p: RelocationProgress) => void;
+}): (line: string) => void {
+ let lastDecile = -1;
+ return (line: string) => {
+ const raw = parseRsyncProgress(line);
+ if (!raw) return;
+ const fraction = rsyncProgressFraction(raw, opts.totalBytes);
+ const detail = formatRsyncProgressDetail(raw, opts.totalBytes, formatBytes);
+ opts.onProgress?.({
+ bytes: raw.bytes,
+ totalBytes: opts.totalBytes,
+ fraction,
+ rate: raw.rate,
+ etaSeconds: raw.etaSeconds,
+ detail,
+ });
+ const decile = Math.min(10, Math.floor(fraction * 10));
+ if (decile > lastDecile) {
+ lastDecile = decile;
+ opts.log(`Copying… ${detail}`);
+ }
+ };
+}
+
export type RelocationPreview = {
slug: string;
target: string;
@@ -313,6 +376,11 @@ async function rsyncTree(opts: {
dest: string;
args: string[];
log: (m: string) => void;
+ // Fed every `--info=progress2` redraw. When set, those redraws are kept OUT
+ // of the log: one 131 GB copy is several thousand frames of one line, and the
+ // sink writes a decile line instead. Everything rsync says that is not a
+ // progress frame still goes to the log verbatim.
+ progress?: (line: string) => void;
signal?: AbortSignal;
}): Promise<{ exitCode: number; output: string }> {
// Trailing slash on src: copy the CONTENTS, so <src>/ -> <dest>/ and not
@@ -328,8 +396,26 @@ async function rsyncTree(opts: {
let output = "";
child.all?.on("data", (c: Buffer) => {
const text = c.toString("utf8");
+ // The verify reads this whole buffer back (driftLines), so it is
+ // accumulated verbatim whatever the log ends up holding.
output += text;
- opts.log(text);
+ if (!opts.progress) {
+ opts.log(text);
+ return;
+ }
+ // rsync rewrites the progress line in place with carriage returns, so one
+ // chunk carries many frames. Split on both, exactly as taskHooks does for
+ // yt-dlp.
+ let passedThrough = "";
+ for (const part of text.split(/[\r\n]+/)) {
+ if (part.trim() === "") continue;
+ if (parseRsyncProgress(part) === null) {
+ passedThrough += `${part}\n`;
+ continue;
+ }
+ opts.progress(part);
+ }
+ if (passedThrough) opts.log(passedThrough);
});
const result = await child;
return { exitCode: result.exitCode ?? 1, output };
@@ -465,6 +551,7 @@ export async function relocateChannelMedia(
): Promise<RelocateChannelMediaResult> {
const { paths, slug, direction, signal } = opts;
const log = opts.onLog ?? ((m: string) => console.log(m));
+ const onProgress = opts.onProgress;
const config = await readChannelConfig(paths, slug);
if (!config) throw new Error(`Channel "${slug}" not found`);
if (isSocialChannel(config)) {
@@ -514,6 +601,7 @@ export async function relocateChannelMedia(
dataDir,
target,
log,
+ onProgress,
signal,
resumed: Boolean(resume),
phase: resume?.phase ?? "copy",
@@ -558,6 +646,7 @@ export async function relocateChannelMedia(
root,
target,
log,
+ onProgress,
signal,
// `resumed` relaxes the "must be in-place" precondition, so it must mean
// "this run continues an interrupted move OUT" and nothing looser.
@@ -575,6 +664,7 @@ async function moveOut(args: {
root: string;
target: string;
log: (m: string) => void;
+ onProgress?: (p: RelocationProgress) => void;
signal?: AbortSignal;
resumed: boolean;
phase: RelocationPhase;
@@ -669,6 +759,11 @@ async function moveOut(args: {
dest: target,
args: ["-a", "--partial", "--info=progress2"],
log,
+ progress: makeProgressSink({
+ totalBytes: measured.bytes,
+ log,
+ onProgress: args.onProgress,
+ }),
signal,
});
if (signal?.aborted) {
@@ -790,6 +885,7 @@ async function moveBack(args: {
dataDir: string;
target: string;
log: (m: string) => void;
+ onProgress?: (p: RelocationProgress) => void;
signal?: AbortSignal;
resumed: boolean;
phase: RelocationPhase;
@@ -883,6 +979,11 @@ async function moveBack(args: {
dest: incoming,
args: ["-a", "--partial", "--info=progress2"],
log,
+ progress: makeProgressSink({
+ totalBytes: measured.bytes,
+ log,
+ onProgress: args.onProgress,
+ }),
signal,
});
if (signal?.aborted) {
diff --git a/common/jobs/progressParsers.test.ts b/common/jobs/progressParsers.test.ts
@@ -3,7 +3,12 @@ import assert from "node:assert/strict";
import {
createDiarizeProgressParser,
createDownloadProgressParser,
+ formatRsyncProgressDetail,
+ isRsyncProgressLine,
+ parseRsyncProgress,
+ rsyncProgressFraction,
} from "./progressParsers";
+import { formatBytes } from "../lib/format";
// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/progressParsers.test.ts
@@ -138,3 +143,90 @@ test("the diarize parser ignores everything that is not its own output", () => {
);
assert.equal(p.feed("[download] 50.0% of 1.00GiB at 2.31MiB/s ETA 00:30"), null);
});
+
+// --- rsync --info=progress2 -------------------------------------------------
+//
+// The exact bytes rsync 3.x writes, captured off a real copy (two files, 35 MB)
+// and pasted verbatim: one \r-separated buffer, the early frame with no xfr#
+// suffix, a mid frame, and the two identical final frames rsync always prints.
+const RSYNC_CHUNK =
+ "\r 32,768 0% 0.00kB/s 0:00:00 " +
+ "\r 30,000,000 85% 1.03GB/s 0:00:03 (xfr#1, to-chk=1/3)" +
+ "\r 35,000,000 100% 1.05GB/s 0:00:00 (xfr#2, to-chk=0/3)";
+
+test("rsync progress: bytes, percent, rate and ETA off one frame", () => {
+ const p = parseRsyncProgress(
+ " 30,000,000 85% 1.03GB/s 0:01:03 (xfr#1, to-chk=1/3)",
+ );
+ assert.deepEqual(p, {
+ bytes: 30_000_000,
+ percent: 85,
+ rate: "1.03GB/s",
+ etaSeconds: 63,
+ });
+});
+
+test("rsync progress: a whole \\r buffer splits into frames", () => {
+ const frames = RSYNC_CHUNK.split(/[\r\n]+/)
+ .filter((l) => l.trim() !== "")
+ .map(parseRsyncProgress);
+ assert.equal(frames.length, 3);
+ assert.deepEqual(
+ frames.map((f) => f?.bytes),
+ [32_768, 30_000_000, 35_000_000],
+ );
+ assert.deepEqual(
+ frames.map((f) => f?.percent),
+ [0, 85, 100],
+ );
+});
+
+test("rsync progress: an hours-long ETA is seconds, not a string", () => {
+ assert.equal(
+ parseRsyncProgress(" 1,000 10% 1.00MB/s 2:03:04")?.etaSeconds,
+ 2 * 3600 + 3 * 60 + 4,
+ );
+});
+
+test("rsync progress: everything that is not a frame is not a frame", () => {
+ for (const line of [
+ "sending incremental file list",
+ "20240101_test1234567/",
+ "sent 35,012,345 bytes received 4,321 bytes 70,033,332.00 bytes/sec",
+ "total size is 35,000,000 speedup is 1.00",
+ ".d..t...... 20240101_test1234567/",
+ "",
+ "rsync: [sender] change_dir failed: No such file or directory (2)",
+ ]) {
+ assert.equal(parseRsyncProgress(line), null, line);
+ assert.equal(isRsyncProgressLine(line), false, line);
+ }
+});
+
+// THE FRACTION IS AGAINST THE MEASURED TREE, not rsync's percentage — which
+// under incremental recursion is a percentage of what it has enumerated so far
+// and walks backwards. The controller has already measured the whole tree for
+// its space check, so it has a denominator that only moves one way.
+test("rsync progress: the fraction divides by the measured tree", () => {
+ const p = parseRsyncProgress(" 30,000,000 85% 1.03GB/s 0:00:03")!;
+ assert.equal(rsyncProgressFraction(p, 60_000_000), 0.5);
+ // No measurement (a caller that never walked the tree) falls back to rsync's
+ // own percentage rather than reporting zero.
+ assert.equal(rsyncProgressFraction(p, 0), 0.85);
+ // Clamped: a tree that grew under the copy must not report 140 %.
+ assert.equal(rsyncProgressFraction(p, 10_000_000), 1);
+});
+
+test("rsync progress: the detail line is the one wording", () => {
+ const p = parseRsyncProgress(" 30,000,000 85% 1.03GB/s 0:05:32")!;
+ assert.equal(
+ formatRsyncProgressDetail(p, 60_000_000, formatBytes),
+ "28.6 MB of 57.2 MB · 50 % · 1.03GB/s · ETA 5:32",
+ );
+ // A finished frame prints 0:00:00, and an ETA of zero is noise.
+ const done = parseRsyncProgress(" 60,000,000 100% 1.05GB/s 0:00:00")!;
+ assert.equal(
+ formatRsyncProgressDetail(done, 60_000_000, formatBytes),
+ "57.2 MB of 57.2 MB · 100 % · 1.05GB/s",
+ );
+});
diff --git a/common/jobs/progressParsers.ts b/common/jobs/progressParsers.ts
@@ -455,3 +455,98 @@ export function createParakeetProgressParser(): {
},
};
}
+
+// ---------------------------------------------------------------------------
+// rsync --info=progress2, for the relocate job
+// ---------------------------------------------------------------------------
+
+// ONE LINE OF `rsync --info=progress2`, parsed.
+//
+// rsync rewrites this line in place with a carriage return, so a chunk read off
+// the child's stdout holds several of them; the caller splits on /[\r\n]/ and
+// feeds each part. The shape (measured, not assumed — rsync 3.x, no
+// --human-readable):
+//
+// 30,000,000 85% 1.03GB/s 0:00:00 (xfr#1, to-chk=1/3)
+//
+// Four fields: bytes transferred so far (grouped with commas), rsync's own
+// percentage, a rate, and an ETA as h:mm:ss. The trailing `(xfr#…)` appears only
+// once a file completes and is deliberately not parsed — `to-chk` counts FILES,
+// and a relocate is priced in bytes.
+//
+// WHY THE PERCENTAGE IS NOT THE FRACTION WE REPORT. rsync's own percentage is
+// against the total it has scanned SO FAR: with incremental recursion (the
+// default) an early line reads "85%" of a tree it has only half enumerated, and
+// the bar then walks backwards. The relocate controller has already measured the
+// whole tree for its space check, so it divides by that instead and the fraction
+// only ever moves forward. This parser returns both and lets the caller choose.
+export type RsyncProgress = {
+ // Bytes rsync says it has transferred so far.
+ bytes: number;
+ // rsync's own percentage, 0..100. See above for why it is not the fraction.
+ percent: number;
+ // Verbatim, e.g. "1.03GB/s". Not re-formatted: it is already the unit an
+ // operator watching a copy reads in.
+ rate: string;
+ // Seconds, from rsync's h:mm:ss ETA.
+ etaSeconds: number;
+};
+
+const RSYNC_PROGRESS_RE =
+ /^\s*([\d,._ ]*\d)\s+(\d{1,3})%\s+(\S+)\s+(\d+):([0-5]?\d):([0-5]?\d)/;
+
+export function parseRsyncProgress(line: string): RsyncProgress | null {
+ const m = line.match(RSYNC_PROGRESS_RE);
+ if (!m) return null;
+ // Grouping separators vary with the locale rsync was built against; strip
+ // everything that is not a digit rather than assuming a comma.
+ const bytes = Number.parseInt(m[1].replace(/\D/g, ""), 10);
+ const percent = Number.parseInt(m[2], 10);
+ if (!Number.isFinite(bytes) || !Number.isFinite(percent)) return null;
+ const etaSeconds =
+ Number.parseInt(m[4], 10) * 3600 +
+ Number.parseInt(m[5], 10) * 60 +
+ Number.parseInt(m[6], 10);
+ return { bytes, percent, rate: m[3], etaSeconds };
+}
+
+// Is this line ONLY a progress redraw? Used by the relocate controller to keep
+// several thousand carriage-return redraws out of a job log that an operator
+// reads afterwards — the parsed decile lines say the same thing in eleven lines.
+export function isRsyncProgressLine(line: string): boolean {
+ return parseRsyncProgress(line) !== null;
+}
+
+// The one wording of a copy's progress, shared by the job task's detail string
+// and the decile line in the log — so the /jobs row and the log cannot word the
+// same instant differently.
+//
+// "12.3 GB of 45.6 GB · 27 % · 110.50MB/s · ETA 5:32"
+//
+// `totalBytes` is the measured source tree, not rsync's running total.
+export function formatRsyncProgressDetail(
+ p: RsyncProgress,
+ totalBytes: number,
+ formatBytes: (n: number) => string,
+): string {
+ const pct =
+ totalBytes > 0
+ ? Math.round(clamp01(p.bytes / totalBytes) * 100)
+ : p.percent;
+ const bits = [
+ `${formatBytes(p.bytes)} of ${totalBytes > 0 ? formatBytes(totalBytes) : "?"}`,
+ `${pct} %`,
+ p.rate,
+ ];
+ if (p.etaSeconds > 0) bits.push(`ETA ${formatClock(p.etaSeconds)}`);
+ return bits.join(" · ");
+}
+
+// 0..1 against the MEASURED tree, with rsync's own percentage as the fallback
+// for a caller that never measured one.
+export function rsyncProgressFraction(
+ p: RsyncProgress,
+ totalBytes: number,
+): number {
+ return totalBytes > 0 ? clamp01(p.bytes / totalBytes) : clamp01(p.percent / 100);
+}
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -49,7 +49,16 @@ export type JobProgress = {
remainingAudioSeconds?: number;
};
-export type JobTaskKind = "download" | "transcribe" | "digest" | "backfill";
+// "relocate" is one CHANNEL'S media move, not one video: the relocate job has
+// exactly one sub-operation and its progress is bytes copied by rsync. It is a
+// task rather than a JobProgress metric because JobProgress counts artifacts
+// re-countable from disk, and a copy in flight is neither.
+export type JobTaskKind =
+ | "download"
+ | "transcribe"
+ | "digest"
+ | "backfill"
+ | "relocate";
// A single in-flight sub-operation within a job (one video download or one
// transcription). Only currently-running tasks are kept on the record — they
diff --git a/editor/app/channels/lib/relocationJob.ts b/editor/app/channels/lib/relocationJob.ts
@@ -44,14 +44,51 @@ export async function enqueueRelocation(opts: {
queueKey: relocationQueueKey(),
paths,
channelSlug: slug,
- fn: async (onLog, signal) => {
+ // THE COPY IS ONE TASK, and it is the only one this job has. rsync already
+ // prints its progress (`--info=progress2`); until now that went nowhere but
+ // the log, as thousands of carriage-return redraws of one line, so a move
+ // of 131 GB looked from /jobs exactly like a move of 3 MB — a spinner. The
+ // controller parses those frames (see makeProgressSink) and hands them back
+ // as `RelocationProgress`; this turns them into the task bar every other
+ // long job on that page already draws.
+ //
+ // ADDED ON THE FIRST FRAME, NOT AT START. A move spends its first seconds in
+ // preflight (space, writability, movable state) and, on a resume, in a
+ // verify pass that transfers nothing — a task bar sitting at 0 % through
+ // that would be claiming a copy had begun.
+ fn: async (onLog, signal, _setProgress, ctx) => {
+ const taskId = `relocate:${slug}`;
+ let started = false;
const result = await relocateChannelMedia({
paths,
slug,
direction,
root,
onLog,
+ onProgress: (p) => {
+ if (!started) {
+ started = true;
+ ctx?.addTask({
+ id: taskId,
+ label:
+ direction === "out"
+ ? `${slug} → ${root ?? "destination"}`
+ : `${slug} → in place`,
+ kind: "relocate",
+ startedAt: Date.now(),
+ });
+ }
+ ctx?.updateTask(taskId, {
+ fraction: p.fraction,
+ detail: p.detail,
+ });
+ },
signal,
+ }).finally(() => {
+ // In a `finally` so a failed or cancelled move does not leave a task on
+ // the record for the life of the process — the same contract
+ // taskHooks.ts's `end()` has.
+ if (started) ctx?.removeTask(taskId);
});
onLog(
`${direction === "out" ? "Moved" : "Moved back"} ${result.files} file(s) / ` +
diff --git a/editor/app/jobs/components/JobProgressBars.tsx b/editor/app/jobs/components/JobProgressBars.tsx
@@ -32,6 +32,7 @@ const TASK_KIND_VERB: Record<JobTaskKind, string> = {
transcribe: "Transcribing",
digest: "Digesting",
backfill: "Backfilling",
+ relocate: "Moving media",
};
const METRIC_FILL_BY_TASK: Record<JobTaskKind, string> = {
@@ -39,6 +40,7 @@ const METRIC_FILL_BY_TASK: Record<JobTaskKind, string> = {
transcribe: "bg-success",
digest: "bg-info",
backfill: "bg-warning",
+ relocate: "bg-info/70",
};
export function TaskProgressBar({ task }: { task: JobRowTask }) {
diff --git a/editor/app/widget/components/MonitorWidget.tsx b/editor/app/widget/components/MonitorWidget.tsx
@@ -1095,6 +1095,8 @@ const TASK_KIND_VERB: Record<JobTaskKind, string> = {
transcribe: "\u270e",
digest: "\u00b6",
backfill: "\u21ba",
+ // A right arrow: a relocate task is a channel's media moving to another disk.
+ relocate: "\u21e2",
};
const TASK_KIND_FILL: Record<JobTaskKind, string> = {
@@ -1102,6 +1104,7 @@ const TASK_KIND_FILL: Record<JobTaskKind, string> = {
transcribe: "bg-success",
digest: "bg-info",
backfill: "bg-warning",
+ relocate: "bg-info/70",
};
// Per-metric glyph for the compact widget line. A Record over JobProgressMetric