1// Utilities for spawning rsync and consuming its output as a `progress.Node`
2// A headless parser is available with `Parse`
3
4export type Line =
5 | { kind: "ignore" }
6 | { kind: "log"; level: "info" | "warn" | "error"; message: string }
7 | { kind: "count"; files: number }
8 | {
9 kind: "progress";
10 currentFile: string;
11 bytesTransferred: number;
12 percentage: number;
13 timeElapsed: string | null;
14 transferNumber: number;
15 filesToCheck: number;
16 totalFiles: number;
17 speed: string | null;
18 };
19
20export const defaultExtraOptions = [
21 "--progress",
22];
23
24export interface SpawnOptions {
25 args: string[];
26 rsync?: string;
27 progress: progress.Node;
28 cwd?: Path;
29}
30
31export async function spawn(options: SpawnOptions) {
32 const { rsync = "rsync", args, cwd } = options;
33 using node = options.progress;
34 const proc = child_process.spawn(rsync, [...defaultExtraOptions, ...args], {
35 stdio: ["ignore", "pipe", "pipe"],
36 cwd: cwd?.toString(),
37 });
38 const parser = new Parse();
39 let running = true;
40 let fileNode: progress.Node | null = null;
41 let fileNodeTransferNumber = 0;
42
43 const stdoutSplitter = readline.createInterface({ input: proc.stdout });
44 const stderrSplitter = readline.createInterface({ input: proc.stderr });
45
46 const title = node.text;
47
48 const handleLine = (line: string) => {
49 const result = parser.onLine(line);
50 if (result.kind === "ignore") {
51 return;
52 } else if (result.kind === "log") {
53 node.log[result.level](result.message);
54 } else if (result.kind === "count") {
55 if (!running) return;
56 node.total = result.files;
57 } else if (result.kind === "progress") {
58 if (!running) return;
59 const {
60 transferNumber,
61 bytesTransferred,
62 totalFiles,
63 filesToCheck,
64 currentFile,
65 speed,
66 } = result;
67 node.value = transferNumber;
68 node.total = totalFiles;
69 node.text = filesToCheck > 0
70 ? `${title} (${filesToCheck} left)`
71 : title;
72 const fileName = currentFile.length > 20
73 ? `${currentFile.slice(0, 3)}..${currentFile.slice(-15)}`
74 : currentFile;
75 if (!fileNode || fileNodeTransferNumber !== transferNumber) {
76 fileNode?.end();
77 fileNode = node.start(fileName, { units: "bytes" });
78 fileNodeTransferNumber = transferNumber;
79 }
80 const estimatedTotal = result.percentage > 0
81 ? Math.ceil(bytesTransferred / (result.percentage / 100))
82 : 0;
83 fileNode.text = speed ? `${fileName} (${speed})` : fileName;
84 fileNode.value = bytesTransferred;
85 fileNode.total = estimatedTotal;
86 } else result satisfies never;
87 };
88
89 stdoutSplitter.on("line", handleLine);
90 stderrSplitter.on("line", handleLine);
91
92 const [code, signal] = await events.once(proc, "close");
93 running = false;
94 fileNode?.end();
95 if (code !== 0) {
96 const fmt = code ? `code ${code}` : `signal ${signal}`;
97 const e: any = new Error(`rsync failed with ${fmt}`);
98 e.args = [rsync, ...args].join(" ");
99 e.code = code;
100 e.signal = signal;
101 throw e;
102 }
103}
104
105export class Parse {
106 totalFiles = 0;
107 currentTransfer = 0;
108 toCheck = 0;
109
110 onLine(line: string): Line {
111 line = line.trimEnd();
112
113 // Parse progress lines like:
114 // 20c83c16735608fc3de4aac61e36770d7774e0c6/au26.m4s
115 // 238,377 100% 460.06kB/s 0:00:00 (xfr#557, to-chk=194111/194690)
116 const progressMatch = line.match(
117 /^\s+([\d,]+)\s+(\d+)%\s+(\S+)\s+(?:(\S+)\s+)?(?:\(xfr#(\d+), to-chk=(\d+)\/(\d+)\))?/,
118 );
119 if (progressMatch) {
120 const [
121 ,
122 bytesStr,
123 percentageStr,
124 speed,
125 timeElapsed,
126 transferStr,
127 toCheckStr,
128 totalStr,
129 ] = progressMatch;
130
131 if (transferStr) this.currentTransfer = Number(transferStr);
132
133 return {
134 kind: "progress",
135 currentFile: this.lastSeenFile || "",
136 bytesTransferred: Number(UNWRAP(bytesStr).replaceAll(",", "")),
137 percentage: Number(percentageStr),
138 timeElapsed: timeElapsed ?? null,
139 transferNumber: this.currentTransfer,
140 filesToCheck: toCheckStr
141 ? this.toCheck = Number(toCheckStr)
142 : this.toCheck,
143 totalFiles: totalStr
144 ? this.totalFiles = Number(totalStr)
145 : this.totalFiles,
146 speed: speed || null,
147 };
148 }
149
150 // Skip common rsync info lines
151 if (!line.startsWith(" ") && !line.startsWith("rsync")) {
152 if (
153 line.startsWith("sending incremental file list")
154 || line.startsWith("sent ")
155 || line.startsWith("total size is ")
156 || line.includes("speedup is ")
157 || line.startsWith("building file list")
158 ) {
159 return { kind: "ignore" };
160 }
161 if (line.trim().length > 0) {
162 this.lastSeenFile = line;
163 }
164 return { kind: "ignore" };
165 }
166 if (line.startsWith(" ")) {
167 const match = line.match(/ (\d+) files.../);
168 if (match) {
169 return { kind: "count", files: Number(match[1]) };
170 }
171 }
172 if (
173 line.toLowerCase().includes("error")
174 || line.toLowerCase().includes("failed")
175 ) {
176 return { kind: "log", level: "error", message: line };
177 }
178 if (
179 line.toLowerCase().includes("warning")
180 || line.toLowerCase().includes("skipping")
181 ) {
182 return { kind: "log", level: "warn", message: line };
183 }
184 return { kind: "log", level: "info", message: line };
185 }
186
187 private lastSeenFile: string | null = null;
188}
189
190import { Path } from "#sitegen/path";
191import { UNWRAP } from "@clo/lib/assert";
192import * as progress from "@clo/lib/progress";
193import * as child_process from "node:child_process";
194import events from "node:events";
195import * as readline from "node:readline";