1/**
2 * a progress reporting API that allows creating CLI spinners and progress bars
3 * that can be easily nested to visualize complex, parallel work. in addition
4 * to providing rich feedback in the terminal, a headless {@linkcode Root} node
5 * can be constructed, which can be used to bind to {@link formatHtml|the DOM}
6 * or transmit progress information {@link encodeStream|across a network}.
7 *
8 * the progress API takes advantage of the `using` keyword to make nodes
9 * scoped. the convention is to make functions take a {@linkcode Ref|progress.Ref},
10 * to act as the destination for the sub-node.
11 *
12 * ```ts
13 * async function doSomething(p: progress.Ref) {
14 * using node = p.start("my task");
15 * // ...
16 * await doSomethingElse(node); // node satisfies progress.Ref
17 * // ...
18 * }
19 * async function doSomethingElse(p: progress.Ref) {
20 * using node = p.start("my sub task");
21 * // ...
22 * }
23 *
24 * doSomething(progress); // the progress module satisfies progress.Ref
25 *
26 * import * as progress from "@clo/lib/progress";
27 * ```
28 *
29 * in some cases, such as in `subprocess/ffmpeg.ts`, it may make more sense
30 * to pass in a pre-started node. the `ffmpeg` module, for example, uses the
31 * starting text as a base title and decorates it with encoding stats.
32 *
33 * when using the top-level {@linkcode start|progress.start}, proper
34 * integration is made with `lib/log.ts` so that the progress tree does not
35 * intersect with regular log messages, as well as forwarding each
36 * {@linkcode Node} custom logging scope to the global log, including any
37 * custom redirections.
38 *
39 * this module is under construction. while i am happy with the overall API, it
40 * needs more work and feature development. the API of `Node` is stable, though.
41 *
42 * inspired by the [Zig Progress API](https://andrewkelley.me/post/zig-new-cli-progress-bar-explained.html).
43 *
44 * @module
45 */
46
47/**
48 * creates a new trackable unit of work, visible in the terminal as a spinner
49 * that rests at the end of the log. if `estimate` is given, a progress bar is
50 * shown.
51 */
52/* node:coverage ignore next 3 */
53export function start(text: string, opts?: StartOptions): Node {
54 return global.start(text, opts);
55}
56
57/**
58 * a reference point to create sub-items. this interface is satisfied by
59 * - the `progress` module itself
60 * - a node recieved by `p.start("label")`
61 * - a headless progress runner from `progress.headless()`
62 *
63 * prefer taking `Ref` as a function parameter, so that the function can be
64 * started at the top level, or as a sub-task of an existing node.
65 *
66 * ```ts
67 * function build(p: progress.Ref = progress) {
68 * using _ = p.start("build process");
69 * // ...
70 * }
71 * ```
72 *
73 * sometimes it makes sense to have an API recieve `Node` instead of `Ref`,
74 * specifically if the caller should have some control over the node name.
75 *
76 * ```ts
77 * await ffmpeg.spawn({
78 * cmd: ["-i", "hello.mov", "-c:v", "svtav1", "hello.mp4"],
79 * progress: progress.start("encode hello.mov"),
80 * });
81 * ```
82 */
83export interface Ref {
84 /**
85 * creates a new trackable unit of work as a child of this one.
86 * when given an estimate, a progress bar is rendered.
87 */
88 start(text: string, opts?: StartOptions): Node;
89}
90
91/**
92 * control for a progress nooe. all fields are primed with setters to trigger
93 * ui dispatch automatically. to make tools that use `progress.ts` more
94 * modular to a custom progress implementation, many fields are optional.
95 */
96export interface Node extends Ref, Disposable {
97 /** end this Node or RootNode. ending collapses children. */
98 end(this: Node): void;
99 /** one line of text display. */
100 text: string;
101 /**
102 * number of items completed. by default this is a count of items, but the
103 * `units` property can be passed to `start` to change how the `value` is
104 * interpretted when displaying this progress node.
105 */
106 value: number;
107 /**
108 * an estimate of how many items this task contains.
109 * when set above `0`, it is shown with a progress bar.
110 */
111 total: number;
112 /**
113 * defaults to increasing by 1. if incrementing causes `value >= estimate`,
114 * it also calls `end`. use `node.value += 1` if that behavior is not wanted.
115 */
116 inc(this: Node, value?: number): void;
117 /** a scoped logger set to output inline on this item. */
118 readonly log: log.RootScope;
119
120 /**
121 * a time for when this node is estimated to be completed, in milliseconds
122 * relative to the unix epoch. by default, this is automatically managed and
123 * made read-only. if `StartOptions` specifies `estimateCompletions: false`,
124 * this field may be mutated.
125 */
126 estimatedTime?: number | null;
127 /**
128 * default `true`. `false` will hides a bar + printing total. this is useful
129 * for preseving auto-end behavior of `inc` without showing the total.
130 */
131 showTotal?: boolean;
132 /**
133 * default `false`. when true, this progress item is hidden if
134 * there are no children.
135 */
136 passive?: boolean;
137 /** default `false`. when true, this progress item is hidden */
138 hidden?: boolean;
139 /** override sorting for children */
140 sortChildren?: ((a: ReadOnlyNode, b: ReadOnlyNode) => number) | null;
141 /**
142 * specify how `value` and `total` is formed into text.
143 * this has some discoverable presets on `StartOptions.units`.
144 */
145 valueFormatter?: (value: number, total: number | null) => string;
146}
147
148/** override the default field values of {@linkcode Node} on creation. */
149export interface StartOptions {
150 /** number of items already completed. */
151 value?: number | undefined | null;
152 /** when set above `0`, it is shown with a progress bar. */
153 total?: number | undefined | null;
154 /**
155 * presets for common formattings of `value` and `total`.
156 * @default "count"
157 */
158 units?: "count" | "bytes" | "percent";
159 /**
160 * customize how an estimate is given for progress bars. to make
161 * `Node.estimatedTime` mutable, pass `false` here.
162 * @default `null`, which is equivilent to `new progress.Ema()`
163 */
164 estimateCompletion?: EstimationAlgorithm | boolean | null;
165 /**
166 * when true, this progress item is hidden if there are no children.
167 * this makes sense when scheduling items with a priority queue, where a
168 * passive node is only used for grouping, but children are used for
169 * indicating status.
170 * @default false
171 */
172 passive?: boolean | undefined | null;
173 /** default `false`. when true, this progress item is hidden */
174 hidden?: boolean | undefined | null;
175 /**
176 * default `true`. `false` will hides a bar + printing estimate,
177 * but preseving auto-end behavior of `inc`
178 */
179 showTotal?: boolean | undefined | null;
180 /** specify sorting for children */
181 sortChildren?: Node["sortChildren"] | undefined | null;
182 /**
183 * specify how `value` and `total` is formed into text.
184 * this has some discoverable presets on `StartOptions.units`.
185 */
186 valueFormatter?: ValueFormatter;
187 /** pre-fill the list of log messages */
188 messages?: log.Message[];
189}
190
191/**
192 * value` will be at least zero, total will be `null` or greater than zero. a
193 * custom value formatter can be passed to {@linkcode start} or mutated on an
194 * existing node.
195 */
196export type ValueFormatter = (value: number, total: number | null) => string;
197
198/**
199 * creates a function to use as a `Node.valueFormatter` by using a single
200 * number formatter, and then joining the node's value and total with a "/",
201 * but leaving it out if there is no total.
202 */
203export function slashValueFormatter(
204 fmtNumber: (value: number) => string,
205): ValueFormatter {
206 return (value: number, total: number | null) =>
207 value > 0
208 ? fmtNumber(value)
209 + (total != null && total > 0 ? "/" + fmtNumber(total) : "")
210 : "";
211}
212
213/** the formatter used when passing `units: "count"` to {@linkcode start} */
214export const defaultValueFormatter: ValueFormatter = /* @__PURE__ */
215 slashValueFormatter(String);
216
217/** the formatter used when passing `units: "bytes"` to {@linkcode start} */
218export const bytesValueFormatter: ValueFormatter = /* @__PURE__ */
219 slashValueFormatter(string.formatByteSize);
220
221/** the formatter used when passing `units: "percent"` to {@linkcode start} */
222export function percentValueFormatter(
223 value: number,
224 total: number | null,
225): string {
226 if (total != null && total > 0) {
227 value /= total;
228 }
229 return `${Math.round(value * 1000) / 10}%`;
230}
231
232/**
233 * {@linkcode Root} allows overriding it's timing functions, which are used for
234 * mostly for rendering and streaming adapters.
235 */
236export interface RootOptions {
237 delay?: typeof async.delay;
238 /** used for debouncing and UI. time estimation does NOT use this method */
239 now?: typeof performance.now;
240}
241
242/**
243 * a headless progress node. when the state of the tree changes, the `change`
244 * event is emitted (batched with many changes). the primary use case of using
245 * a custom `Root`s is to stream a task's progress over a network or process
246 * IPC using {@linkcode encodeByteStream}/{@linkcode encodeEventStream}.
247 * another use case is to render a progress task differently.
248 *
249 * custom events may be dispatched on the `Root`, but this is only useful when
250 * streaming, since the streams will carry any custom events.
251 */
252export class Root<
253 Result = void,
254 Map extends Events.Map = ts.EmptyObject,
255> extends Events<MergeRootEvents<Result, Map>> {
256 delay: typeof async.delay;
257 now: typeof performance.now;
258 #status: "active" | "end" | "error" = "active";
259 #active: Internal[] = [];
260 #debounce: async.Cancelable<void> | null = null;
261
262 constructor({ delay, now }: RootOptions = {}) {
263 super();
264 this.delay = delay ?? async.delay;
265 this.now = now ?? performance.now.bind(performance);
266 this.on("node-end", (state) => {
267 if (state.parent) return;
268 const i = this.#active.indexOf(state as Internal);
269 ASSERT(i !== -1);
270 this.#active.splice(i, 1);
271 });
272 }
273
274 /** the caller may not mutate `Root`'s state */
275 get active(): readonly ReadOnlyNode[] {
276 return this.#active;
277 }
278
279 /**
280 * creates a new trackable unit of work at the top level of this root.
281 * when given an estimate, a progress bar is rendered.
282 */
283 start(text: string, opts?: StartOptions): Node {
284 ASSERT(this.#status === "active");
285 const [state, node] = newNode(this, text, opts);
286 this.#active.push(state);
287 this.emit("node-start", state, null);
288 return node;
289 }
290
291 /**
292 * end this root, as well as all children nodes. this emits the `end` event
293 * with the provided value, JSON-serializing it if connected via
294 * {@linkcode encodeEventStream} or {@linkcode encodeByteStream}.
295 *
296 * this function is pre-bound so that it can be passed to `then`:
297 * ```ts
298 * const root = new Root();
299 * void asyncFunction().then(root.end, root.error);
300 * ```
301 */
302 end = (result: Result): void => {
303 this.emit("end", result);
304 };
305
306 /**
307 * end this root, as well as all children nodes. this emits the `error` event.
308 */
309 error = (error: unknown): void => {
310 this.emit("error", error);
311 };
312
313 /** convert this root into a promise. */
314 asPromise(): Promise<Result> {
315 return this.once("end").then((x) => x[0]);
316 }
317
318 /**
319 * emit an event on the specified channel. behavior is not defined when manually
320 * emitting one of the built in events.
321 */
322 override emit<C extends keyof MergeRootEvents<Result, Map>>(
323 channel: C,
324 ...args: MergeRootEvents<Result, Map>[C]
325 ): void {
326 ASSERT(this.#status === "active");
327 if (channel !== "change" && rootEvents.has(channel)) this.#emitChangeSoon();
328 if (channel === "error" || channel === "end") {
329 this.#end();
330 this.#status = channel;
331 }
332 super.emit(channel, ...args);
333 }
334
335 #end() {
336 const hasActive = this.#active.length > 0 || !this.#debounce;
337 for (const active of this.#active) endNode(this, active);
338 this.#active.splice(0, this.#active.length);
339 this.#debounce?.cancel();
340 if (hasActive) this.emit("change", this.#active);
341 }
342
343 #emitChangeSoon = () => {
344 if (this.#debounce) return;
345 (this.#debounce = this.delay(1))
346 .then(() => {
347 this.#debounce = null;
348 this.emit("change", this.#active);
349 });
350 };
351}
352
353/** the built in events provided by {@linkcode Root} */
354export type RootEventMap<Result = void> = {
355 // this event is the debounced "ready to re-render" event. do not emit.
356 "change": [rootNodes: readonly ReadOnlyNode[]];
357 // these events are emitted without a debounce, and are used for general
358 // communication within `Root`'s implementation. do not emit.
359 "node-start": [newNode: ReadOnlyNode, parentNode: ReadOnlyNode | null];
360 "node-change": [node: ReadOnlyNode, key: keyof ReadOnlyNode];
361 "node-end": [node: ReadOnlyNode];
362 "node-detached-log": [msg: log.Message];
363 // these events are used for result communication. they can be sent via
364 // `emit` or the shorthand methods `end` and `error`. sending either will put
365 // the `Root` into a "done" state.
366 "end": [result: Result];
367 "error": [error: unknown];
368};
369
370const rootEvents: ReadonlySet<unknown> = new Set([
371 "change",
372 "node-start",
373 "node-change",
374 "node-end",
375 "node-detached-log",
376 "end",
377 "error",
378]);
379
380type MergeRootEvents<Result, Map extends Events.Map> = {
381 [K in keyof Map | keyof RootEventMap]: K extends keyof RootEventMap ? RootEventMap<Result>[K]
382 : Map[K];
383};
384
385/** a read-only version of {@linkcode Node}, provided to renderers. */
386export type ReadOnlyNode = Readonly<
387 & Pick<
388 Required<Node>,
389 | "text"
390 | "value"
391 | "total"
392 | "estimatedTime"
393 | "passive"
394 | "hidden"
395 | "showTotal"
396 | "sortChildren"
397 | "valueFormatter"
398 >
399 & {
400 /** unique identifier. once {@linkcode Node} is disposed, the key is reused. */
401 key: number;
402 logs: readonly log.Message[];
403 /** traverse down the tree. */
404 children: readonly ReadOnlyNode[];
405 /** traverse up the tree. */
406 parent: ReadOnlyNode | null;
407 }
408>;
409
410/** internal state is read-write. */
411interface Internal extends ts.Writeable<ReadOnlyNode> {
412 logs: log.Message[];
413 children: Internal[];
414 parent: Internal | null;
415 detached: boolean;
416}
417
418// TypeScript Test
419void function(unknown: unknown) {
420 (unknown as Internal) satisfies ReadOnlyNode;
421};
422
423/**
424 * constructs a linked `Node` and `Internal` that dispatches events to `owner`.
425 * the `Internal` is the non-reactive source of truth that can be blindly cast
426 * to `ReadOnlyNode`, and `Node` is the Node.
427 */
428function newNode<Result, EventMap extends Events.Map>(
429 owner: Root<Result, EventMap>,
430 text: string,
431 opts: StartOptions = {},
432): [Internal, Node] {
433 let estimator: EstimationAlgorithm | null = null;
434
435 const state: Internal = {
436 key: globalKeyPool.get(),
437 text,
438 value: opts.value ?? 0,
439 total: opts.total ?? 0,
440 estimatedTime: null,
441 showTotal: opts.showTotal ?? true,
442 passive: opts.passive ?? false,
443 hidden: opts.hidden ?? false,
444 logs: [],
445 children: [],
446 parent: null,
447 sortChildren: opts.sortChildren ?? null,
448 valueFormatter: opts.valueFormatter ?? {
449 count: defaultValueFormatter,
450 bytes: bytesValueFormatter,
451 percent: percentValueFormatter,
452 }[opts.units ?? "count"],
453 detached: false,
454 };
455
456 if (opts.estimateCompletion !== false) {
457 const { estimateCompletion } = opts;
458 if (estimateCompletion && typeof estimateCompletion === "object") {
459 estimator = estimateCompletion;
460 } else estimator = new Ema();
461 }
462 estimator?.sample(Date.now(), 0, 0);
463
464 function mutate(key: keyof ReadOnlyNode) {
465 if (state.detached) return;
466 if (
467 key === "value" && estimator && state.total > 0 && state.value > 0
468 && state.value < state.total
469 ) {
470 // originally, auto-estimation would use the Root's timing primitives,
471 // but this makes it so that `state.estimatedTime` is in terms of the
472 // Root, which is unintuitive for everyone but the renderer.
473 const time = estimator.sample(Date.now(), state.value, state.total);
474 if (time && (!state.estimatedTime || state.estimatedTime !== time)) {
475 state.estimatedTime = time;
476 owner.emit("node-change", state, "estimatedTime");
477 }
478 }
479 const { parent } = state;
480 parent?.sortChildren && parent.children.sort(parent.sortChildren);
481 owner.emit("node-change", state, key);
482 }
483
484 const scope = log.headlessScope((m) => {
485 if (state.detached) return owner.emit("node-detached-log", m);
486 state.logs.push(m);
487 owner.emit("node-change", state, "logs");
488 });
489
490 const binding: Node = {
491 start(text, opts) {
492 if (state.detached) return nullNode;
493 const [child, node] = newNode(owner, text, opts);
494 state.children.push(child);
495 child.parent = state;
496 state.sortChildren && state.children.sort(state.sortChildren);
497 owner.emit("node-start", child, state);
498 owner.emit("node-change", state, "children");
499 return node;
500 },
501 get text() {
502 return state.text;
503 },
504 set text(value) {
505 state.text = value;
506 mutate("text");
507 },
508 get value() {
509 return state.value;
510 },
511 set value(value) {
512 state.value = value;
513 mutate("value");
514 },
515 get total() {
516 return state.total;
517 },
518 set total(value) {
519 state.total = value;
520 mutate("total");
521 },
522 get showTotal() {
523 return state.showTotal;
524 },
525 set showTotal(value) {
526 state.showTotal = value;
527 mutate("showTotal");
528 },
529 get passive() {
530 return state.passive;
531 },
532 set passive(value) {
533 state.passive = value;
534 mutate("passive");
535 },
536 get hidden() {
537 return state.hidden;
538 },
539 set hidden(value) {
540 state.hidden = value;
541 mutate("hidden");
542 },
543 get sortChildren() {
544 return state.sortChildren;
545 },
546 set sortChildren(value) {
547 state.sortChildren = value;
548 mutate("sortChildren");
549 },
550 get estimatedTime() {
551 return state.estimatedTime;
552 },
553 set estimatedTime(date) {
554 if (!estimator) {
555 state.estimatedTime = date;
556 mutate("estimatedTime");
557 }
558 },
559 get valueFormatter() {
560 return state.valueFormatter;
561 },
562 set valueFormatter(value) {
563 state.valueFormatter = value;
564 mutate("valueFormatter");
565 },
566 inc(delta = 1) {
567 state.value += delta;
568 if (state.total > 0 && state.value >= state.total) {
569 endNode(owner, state);
570 } else {
571 mutate("value");
572 }
573 },
574 log: scope,
575 end: () => void endNode(owner, state),
576 [Symbol.dispose]: () => void endNode(owner, state),
577 };
578 return [state, binding];
579}
580
581const noop = () => {};
582noop[Symbol.dispose] = noop;
583
584/** minimal implementation of {@linkcode Node} that has no I/O */
585export const nullNode: Node = {
586 start: () => nullNode,
587 get text() {
588 return "[detached]";
589 },
590 set text(_) {
591 },
592 get value() {
593 return 0;
594 },
595 set value(_) {
596 },
597 get total() {
598 return 0;
599 },
600 set total(_) {
601 },
602 inc: noop,
603 log: {
604 info: noop,
605 warn: noop,
606 error: noop,
607 log: noop,
608 debug: noop,
609 writeMessage: noop,
610 write: noop,
611 scoped() {
612 return this;
613 },
614 tee: () => noop,
615 },
616 end: () => {},
617 [Symbol.dispose]: () => {},
618};
619
620function endNode<R, M extends Events.Map>(owner: Root<R, M>, state: Internal) {
621 if (state.detached) return;
622 state.children.forEach((child) => endNode(owner, child));
623 state.detached = true;
624 globalKeyPool.recycle(state.key);
625 if (state.parent) {
626 const i = state.parent.children.indexOf(state);
627 ASSERT(i !== -1);
628 state.parent.children.splice(i, 1);
629 }
630 owner.emit("node-end", state);
631 state.logs = [];
632 state.parent = null;
633}
634
635/** Convert the top level {@linkcode ReadOnlyNode|ReadOnlyNode[]} into ANSI text. */
636export function formatAnsi(now: number, list: readonly ReadOnlyNode[]): string {
637 let out = "";
638 for (const top of list) {
639 if (top.passive && !hasChildren(top)) continue;
640 out += renderAnsiMainLine(top, now, 0) + "\n";
641 out += renderChildren(top, now, []);
642 }
643 return out.trimEnd();
644}
645
646const spinnerFps = 12.5;
647const box = { tee: "├─ ", line: "│ ", langle: "└─ " };
648const barChars = [" ", "▏", "▎", "▍", "▌", "▋", "▊", "▉"];
649const fullBar = "█";
650const spinner = ["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"]
651 .map((frame) => ansi.style(ansi.fgBlue, frame));
652
653function renderAnsiMainLine(state: ReadOnlyNode, now: number, depth: number) {
654 const { text, total, showTotal, value, estimatedTime } = state;
655 const dateNow = Date.now();
656 const showEstimate = estimatedTime && !findLongerEstimate(state, estimatedTime);
657 const estimate = showEstimate && (estimatedTime > dateNow + 1000)
658 ? ", "
659 + string.formatDurationLetters(Math.round((estimatedTime - dateNow) / 1000))
660 : "";
661 const valueFormatted = state.valueFormatter(value, total);
662 if (total && showTotal) {
663 return ansi.style(
664 ansi.bgBrightBlack + ansi.fgBlue,
665 formatUnicodeBar(
666 value / total,
667 Math.max(4, depth > 0 ? 12 : 25),
668 ),
669 )
670 + (valueFormatted || estimate ? ` [${valueFormatted}${estimate}] ` : " ")
671 + text;
672 }
673 const frame = Math.floor(now / (1000 / spinnerFps)) % spinner.length;
674 return ((depth === 0 ? spinner[frame] + " " : "")
675 + (valueFormatted ? `[${valueFormatted}] ` : "") + text);
676}
677
678function hasChildren(states: ReadOnlyNode): boolean {
679 return states.children.some((x) => !x.hidden && (!x.passive || hasChildren(x)));
680}
681
682function findLongerEstimate(
683 state: ReadOnlyNode,
684 estimatedTime: number,
685): boolean {
686 for (const child of state.children) {
687 if (child.hidden || (child.passive && !hasChildren(child))) continue;
688 if (child.estimatedTime != null && child.estimatedTime > estimatedTime) {
689 return true;
690 }
691 if (findLongerEstimate(child, estimatedTime)) return true;
692 }
693 return false;
694}
695
696function renderChildren(state: ReadOnlyNode, now: number, depth: boolean[]) {
697 let maxHeight = 50; // TODO: flexible layout
698 let truncated = 0;
699 let out = "";
700 let { children } = state;
701 children = children.filter((x) => !x.hidden && (!x.passive || hasChildren(x)));
702 if (children.length === 0) return "";
703 for (let i = 0, { length } = children; i < length; i += 1) {
704 const child = UNWRAP(children[i]);
705 let item = "";
706 const left = depth.map((x) => x ? box.line : " ").join("");
707 item += left + (i === length - 1 && !truncated ? box.langle : box.tee);
708 item += renderAnsiMainLine(child, now, 1 + depth.length) + "\n";
709 item += renderChildren(child, now, depth.concat(i < length - 1));
710 const h = string.countNewlines(item);
711 if (h > maxHeight) {
712 truncated += 1;
713 continue;
714 }
715 const logLines = child.logs
716 .map((msg) => log.formatAnsiMessage(msg, true))
717 .join("")
718 .trim();
719 if (logLines) {
720 for (const line of logLines.split("\n").slice(-Math.min(3, maxHeight))) {
721 item += left + (i === length - 1 && !truncated ? " " : box.line)
722 + " "
723 + ansi.style(ansi.fgBrightBlack, ">")
724 + " " + line + "\n";
725 }
726 }
727 maxHeight -= h + string.countNewlines(logLines);
728 out += item;
729 }
730 if (truncated) {
731 out += depth.map((x) => x ? box.line : " ").join("") + box.langle;
732 out += ansi.style(ansi.fgBrightBlack, `[${truncated} more]`) + "\n";
733 }
734 return out;
735}
736
737/**
738 * This function is derived from an old program I wrote back in 2020 called `f`
739 * which did ffmpeg handling. It is probably one of the coolest progress bars
740 * ever imagined.
741 */
742export function formatUnicodeBar(progress: number, width: number): string {
743 if (progress >= 1) return fullBar.repeat(width);
744 if (progress <= 0 || Number.isNaN(progress)) return " ".repeat(width);
745
746 const wholeWidth = Math.floor(progress * width);
747 const remainderWidth = (progress * width) % 1;
748 const partWidth = Math.floor(remainderWidth * 8);
749 let partChar = barChars[partWidth];
750 if (width - wholeWidth - 1 < 0) partChar = "";
751
752 const fill = fullBar.repeat(wholeWidth);
753 const empty = " ".repeat(width - wholeWidth - 1);
754
755 return fill + partChar + empty;
756}
757
758/**
759 * configure a {@linkcode Root} to display its contents to a
760 * {@linkcode log.HeadlessWidgetHost}. this is used by the global progress
761 * instance to output to the terminal.
762 * ```ts
763 * const globalProgress = new progress.Root();
764 * progress.attachToScreen(log);
765 * ```
766 * unit tests use a mock screen instead of the terminal.
767 * ```ts
768 * const screen = new TestWidgetHost();
769 * const root = new progress.Root(screen.delay);
770 * root.attachToScreen(screen);
771 * // can use the progress API
772 * root.start("hello");
773 * // and then examine the contents
774 * screen.expectFrame(...);
775 * ```
776 */
777export function attachToScreen(
778 root: Root,
779 { writeOutput, startWidget }: Pick<
780 log.WidgetHost,
781 "writeOutput" | "startWidget"
782 >,
783): ts.Dispose {
784 const stack = new DisposableStack();
785 let widget: log.WidgetInstance | null = null;
786
787 stack.use(root.on("change", (items) => {
788 if (items.length > 0) {
789 widget ??= startWidget({
790 format: ({ now }) => formatAnsi(now, root.active),
791 }) ?? null;
792 if (!widget) return;
793 // scan for spinners and visible estimations
794 widget.fps = items.some((x) => !x.hidden && !(x.showTotal !== false && x.total > 0))
795 ? spinnerFps
796 : hasEstimate(items)
797 ? 1
798 : null;
799 widget.redraw();
800 } else {
801 widget?.stop();
802 widget = null;
803 }
804 }));
805 stack.use(root.on("node-detached-log", (msg) => {
806 writeOutput(log.formatAnsiMessage(msg, true));
807 }));
808 stack.use(root.on("node-end", (node) => {
809 let title = node.text;
810 let p: ReadOnlyNode | null = node;
811 while (p = p.parent) title = p.text + " / " + title;
812 const { logs } = node;
813 if (logs.length > 0) {
814 const header = `[logs from ${title}]`;
815 writeOutput(ansi.style(ansi.fgBrightBlack, header) + "\n");
816 writeOutput(logs.map((msg) => log.formatAnsiMessage(msg, true)).join(""));
817 }
818 }));
819
820 stack.defer(() => {
821 widget?.stop();
822 widget = null;
823 });
824
825 return ts.defer(() => stack.dispose());
826}
827
828function hasEstimate(nodes: readonly ReadOnlyNode[]): boolean {
829 return nodes.some((n) => n.estimatedTime != null || hasEstimate(n.children));
830}
831
832const header = /* @__PURE__ */ string.encodeUtf8("clover's progress <3\n");
833
834/** the 'Content-Type' of {@linkcode encodeByteStream}'s return value */
835export const contentType = "application/x-clover-progress";
836
837/**
838 * opaque but JSON-compatible serialized state, either containing the entire
839 * progress tree or just a delta update. this data structure always be future
840 * and backwards compatbile with changes to this library. if a breaking change
841 * is made, the APIs in this file will be renamed.
842 *
843 * for details on the format, read the source of {@linkcode encodeEventStream} or
844 * {@linkcode encodeByteStream}
845 */
846export type StreamEvent = Array<
847 number | StreamNode | StreamCustomEvent | StreamRootChildren
848>;
849const kNode = Symbol("underlyingNode");
850/** @internal */
851export interface StreamNode {
852 k: EncodedKey;
853 /** text, overwrite */
854 t?: string | undefined;
855 /** value, overwrite */
856 v?: number | undefined;
857 /** total, overwrite */
858 e?: number | undefined;
859 /** showTotal, overwrite */
860 s?: boolean | undefined;
861 /** passive, overwrite */
862 p?: boolean | undefined;
863 /** logs, append */
864 l?: log.Message[] | undefined;
865 /** hidden, overwrite */
866 h?: boolean | undefined;
867 /** children, overwrite */
868 c?: EncodedKey[] | undefined;
869 /** estimated time, in unix milliseconds, overwrite */
870 E?: number | null | undefined;
871 /** formatted value if not defaultValueFormatter, overwrite */
872 V?: string | undefined;
873
874 /** access to the underlying node is used by {@linkcode encodeByteStream} */
875 [kNode]?: ReadOnlyNode;
876}
877type StreamCustomEvent = [
878 /** custom event name or "end" */
879 e: "end" | string,
880 /** custom event payload */
881 ...args: ts.Json[],
882];
883type StreamRootChildren = number[];
884
885/** options for both {@linkcode encodeEventStream} and {@linkcode encodeByteStream} */
886export interface EncodeStreamOptions {
887 /**
888 * set the minimum time between packets.
889 * @default 1000 / 30 (30fps)
890 */
891 throttleMs?: number;
892 /**
893 * `log.ts` is capable of capturing stack traces from all log calls. by
894 * seeing this to true, those traces will be serialized. disabled by default
895 * for privacy reasons.
896 */
897 serializeStackTraces?: boolean;
898}
899
900/**
901 * converts a {@linkcode Root|progress.Root} into an object stream
902 * for communicating progress over the process or network boundary.
903 *
904 * if the given root adds custom event handlers, they must all have
905 * json-serializable payloads.
906 */
907export function encodeEventStream<
908 Result extends ts.Json,
909 Map extends { [key: string]: ts.Json[] },
910>(
911 // fun fact, this type is impossible to write without generics
912 root: Root<Result, Map>,
913 {
914 throttleMs = 1000 / 30, /* 30fps */
915 serializeStackTraces = false,
916 }: EncodeStreamOptions = {},
917): ReadableStream<StreamEvent> {
918 const { delay, now } = root;
919 const stack = new DisposableStack();
920 const s = new Encoder(root, serializeStackTraces);
921 const ready = new async.Watch(false);
922
923 let lastEvent = now();
924 let timer: async.Cancelable<void> | null = null;
925 function emitSoon() {
926 if (timer || ready.value) return;
927 const remaining = lastEvent + throttleMs - now();
928 if (remaining > 0) {
929 (timer = delay(remaining)).then(() => {
930 timer = null;
931 ready.value = true;
932 lastEvent = now();
933 });
934 } else {
935 ready.value = true;
936 lastEvent = now();
937 }
938 }
939 /** terminal events skip the throttle so the stream ends promptly */
940 function emitNow() {
941 timer?.cancel();
942 timer = null;
943 ready.value = true;
944 lastEvent = now();
945 }
946 stack.defer(() => {
947 if (timer) timer.cancel();
948 s.changed.clear();
949 s.deleted.clear();
950 });
951
952 let isFirst = true;
953 return new ReadableStream({
954 start(controller) {
955 if (root.active.length > 0) {
956 controller.enqueue(s.getDelta(isFirst));
957 isFirst = false;
958 }
959
960 // TODO: why does ts require two args
961 stack.use(root.on("node-start", (newNode, _) => {
962 if (!newNode.parent) s.rootChildrenUpdated = true;
963 s.changed.set(newNode, new Set());
964 emitSoon();
965 }));
966 stack.use(root.on("node-change", (node, key) => {
967 let set = s.changed.get(node);
968 if (!set) s.changed.set(node, set = new Set());
969 set.add(key);
970 emitSoon();
971 }));
972 stack.use(root.on("node-end", (node) => {
973 if (!node.parent) s.rootChildrenUpdated = true;
974 s.changed.delete(node);
975 s.deleted.add(node); // handle add and remove in same frame
976 emitSoon();
977 }));
978 // custom events pass through the stream verbatim
979 stack.use(root.onAny((channel, args) => {
980 if (rootEvents.has(channel)) return;
981 s.pendingEvents.push([channel as string, ...args as ts.Json[]]);
982 emitSoon();
983 }));
984 stack.use(root.on("end", (result) => {
985 s.pendingEvents.push(["end", result]);
986 s.ended = true;
987 emitNow();
988 }));
989 stack.use(root.on("error", (error) => {
990 s.pendingEvents.push(["error", exceptions.message(error)]);
991 s.ended = true;
992 emitNow();
993 }));
994 },
995 async pull(controller) {
996 await ready.until((x) => x === true);
997 ready.value = false;
998 if (timer) {
999 timer.cancel();
1000 timer = null;
1001 }
1002 controller.enqueue(s.getDelta(isFirst));
1003 isFirst = false;
1004 if (s.ended) {
1005 controller.close();
1006 stack.dispose();
1007 }
1008 },
1009 cancel() {
1010 stack.dispose();
1011 },
1012 }, { highWaterMark: 0 });
1013}
1014
1015/**
1016 * decodes a stream created by {@linkcode encodeEventStream} into managed calls
1017 * to `target.start`. when cancelling, all the managed nodes are destroyed.
1018 */
1019export function decodeEventStream<
1020 Result = void,
1021 EventMap extends Events.Map = ts.EmptyObject,
1022>(
1023 encoded: ReadableStream<StreamEvent>,
1024 target: Ref,
1025): async.Cancelable<Result> & Events<EventMap> {
1026 const { resolve, reject, promise } = Promise.withResolvers<Result>();
1027 const events = new Events<EventMap>();
1028 const ctrl = new AbortController();
1029 const decoder = new Decoder(target, events, resolve);
1030 const { signal } = ctrl;
1031 (async () => {
1032 if (signal.aborted) return;
1033 const reader = encoded.getReader();
1034 try {
1035 let hasEmittedEnd = false;
1036 while (true) {
1037 signal.throwIfAborted();
1038 const { value, done } = await reader.read();
1039 signal.throwIfAborted();
1040 if (done) break;
1041 if (decoder.processEvent(value)) hasEmittedEnd = true;
1042 }
1043 ASSERT(hasEmittedEnd, "Stream terminated early.");
1044 } catch (err) {
1045 reader.cancel(err);
1046 if (!signal.aborted) reject(err);
1047 } finally {
1048 reader.releaseLock();
1049 }
1050 })();
1051 return Object.assign(ts.mixin(promise, events), {
1052 cancel: (e?: unknown) => ctrl.abort(e),
1053 [Symbol.dispose]: () => ctrl.abort(),
1054 });
1055}
1056
1057/**
1058 * by listening to the event interface in {@linkcode Root}, compute delta
1059 * events for a progress stream. the caller is expected to fill `changed` and `deleted`
1060 */
1061class Encoder<
1062 Result extends ts.Json,
1063 EventMap extends { [key: string]: ts.Json[] },
1064> {
1065 root: Root<Result, EventMap>;
1066 pool: KeyPool<EncodedKey> = new KeyPool();
1067 keyMap: Map<ReadOnlyNode, EncodedKey> = new Map();
1068 logsLength: WeakMap<ReadOnlyNode, number> = new WeakMap();
1069 serializeStackTraces: boolean;
1070
1071 rootChildrenUpdated = true;
1072 changed = new Map<ReadOnlyNode, Set<keyof ReadOnlyNode>>();
1073 deleted = new Set<ReadOnlyNode>();
1074 /** custom and terminal ("end"/"error") events awaiting the next delta */
1075 pendingEvents: StreamCustomEvent[] = [];
1076 /** set when "end" or "error" was recorded; the stream closes after flushing */
1077 ended = false;
1078
1079 constructor(
1080 root: Root<Result, EventMap>,
1081 serializeStackTraces: boolean = false,
1082 ) {
1083 this.root = root;
1084 this.serializeStackTraces = serializeStackTraces;
1085 for (const node of root.active) {
1086 this.changed.set(node, new Set());
1087 }
1088 }
1089
1090 getDelta(isFirst = false) {
1091 const deleted = Array.from(
1092 this.deleted,
1093 (node) => [node, this.keyMap.get(node)] as const,
1094 )
1095 .filter(([node, dst]) => {
1096 if (dst == null) return false; // started and stopped in same tick
1097 this.pool.recycle(dst);
1098 this.keyMap.delete(node);
1099 return true;
1100 })
1101 .map((x) => UNWRAP(x[1]));
1102 ASSERT(!isFirst || deleted.length === 0);
1103
1104 const { updated = [], started = [] } = Object.groupBy(
1105 this.changed,
1106 ([node]) => this.keyMap.has(node) ? "updated" : "started",
1107 );
1108
1109 const payload: StreamEvent = [
1110 isFirst ? -1 : deleted.length,
1111 ...deleted,
1112 ];
1113
1114 if (started.length > 0) {
1115 this.beginNodesDepthFirst(
1116 started.reverse().map(([node]) => node),
1117 payload,
1118 );
1119 }
1120
1121 for (const [node, props] of updated) {
1122 let logs = undefined;
1123 if (
1124 props.has("logs")
1125 && node.logs.length > (this.logsLength.get(node) ?? 0)
1126 ) {
1127 logs = this.redactLogsIfNeeded(
1128 node.logs.slice(this.logsLength.get(node) ?? 0),
1129 );
1130 this.logsLength.set(node, node.logs.length);
1131 }
1132 payload.push({
1133 k: UNWRAP(this.keyMap.get(node)),
1134 t: props.has("text") ? node.text : undefined,
1135 v: props.has("value") ? node.value : undefined,
1136 e: props.has("total") ? node.total : undefined,
1137 s: props.has("showTotal") ? node.showTotal : undefined,
1138 p: props.has("passive") ? node.passive : undefined,
1139 l: logs,
1140 h: props.has("hidden") ? node.hidden : undefined,
1141 c: props.has("children")
1142 ? node.children.map((node) => UNWRAP(this.keyMap.get(node)))
1143 : undefined,
1144 E: props.has("estimatedTime") ? node.estimatedTime : undefined,
1145 V: node.valueFormatter !== defaultValueFormatter
1146 && (props.has("value") || props.has("total")
1147 || props.has("valueFormatter"))
1148 ? node.valueFormatter(node.value, node.total)
1149 : undefined,
1150 [kNode]: node,
1151 });
1152 }
1153
1154 if (this.rootChildrenUpdated) {
1155 payload.push(
1156 this.root.active.map((node) => UNWRAP(this.keyMap.get(node))),
1157 );
1158 }
1159
1160 payload.push(...this.pendingEvents);
1161
1162 this.deleted.clear();
1163 this.changed.clear();
1164 this.pendingEvents = [];
1165 this.rootChildrenUpdated = false;
1166
1167 return payload;
1168 }
1169
1170 beginNodesDepthFirst(
1171 nodes: ReadOnlyNode[],
1172 out: StreamEvent,
1173 ) {
1174 for (let i = 0; i < nodes.length; i += 1) {
1175 const node = UNWRAP(nodes[i]);
1176 nodes.push(...node.children);
1177 }
1178 let next;
1179 while (next = nodes.pop()) {
1180 if (!this.keyMap.has(next)) {
1181 out.push(this.beginNode(next));
1182 }
1183 }
1184 }
1185
1186 beginNode(node: ReadOnlyNode): StreamNode {
1187 const k = this.pool.get();
1188 this.keyMap.set(node, k);
1189
1190 return {
1191 k,
1192 v: node.value !== 0 ? node.value : undefined,
1193 e: node.total !== 0 ? node.total : undefined,
1194 t: node.text.length > 0 ? node.text : undefined,
1195 l: node.logs.length > 0
1196 ? (this.logsLength.set(node, node.logs.length), this.redactLogsIfNeeded(node.logs.slice()))
1197 : undefined,
1198 c: node.children.length > 0
1199 ? node.children.map((node) => UNWRAP(this.keyMap.get(node), () => [node.children, this.keyMap]))
1200 : undefined,
1201 p: node.passive === true ? true : undefined,
1202 s: node.showTotal === false ? false : undefined,
1203 h: node.hidden === true ? true : undefined,
1204 E: node.estimatedTime != null ? node.estimatedTime : undefined,
1205 V: node.valueFormatter !== defaultValueFormatter
1206 ? node.valueFormatter(node.value, node.total)
1207 : undefined,
1208 [kNode]: node,
1209 };
1210 }
1211
1212 redactLogsIfNeeded(messages: log.Message[]) {
1213 if (this.serializeStackTraces) return messages;
1214 return messages.map((x) => ({ ...x, stack: undefined }));
1215 }
1216}
1217
1218class Decoder<
1219 Result = void,
1220 EventMap extends Events.Map = ts.EmptyObject,
1221> {
1222 target: Ref;
1223 events: Events<EventMap>;
1224 resolve: (result: Result) => void;
1225 active = new Map<EncodedKey, Node>();
1226 pendingStart = new Map<EncodedKey, PendingStart>();
1227 pendingParents = new Map<EncodedKey, EncodedKey>();
1228
1229 constructor(
1230 target: Ref,
1231 events: Events<EventMap>,
1232 resolve: (result: Result) => void,
1233 ) {
1234 this.target = target;
1235 this.events = events;
1236 this.resolve = resolve;
1237 }
1238
1239 reset() {
1240 for (const v of this.active.values()) v.end();
1241 this.active.clear();
1242 }
1243
1244 /** returns whether the stream signalled "end"; throws on "error" */
1245 processEvent(event: StreamEvent): boolean {
1246 event = event.slice();
1247
1248 // delete unreferenced nodes first, since they may get re-allocated
1249 // in the next phase.
1250 const deletedCount = event.shift();
1251 ASSERT(typeof deletedCount === "number");
1252 if (deletedCount === -1) {
1253 // connection reset, delete all nodes
1254 this.reset();
1255 } else {
1256 for (const id of event.splice(0, deletedCount)) {
1257 ASSERT(typeof id === "number");
1258 UNWRAP(
1259 this.active.get(id as EncodedKey),
1260 () => `dstkey ${id} not found, ${[...this.active.keys()]}`,
1261 ).end();
1262 ASSERT(this.active.delete(id as EncodedKey));
1263 }
1264 }
1265
1266 // parse all of the events, updating all nodes and extracting creations
1267 let resolving = false;
1268 let resolvingValue: unknown = null;
1269 let updateChildren: Array<{
1270 key: EncodedKey;
1271 children: EncodedKey[];
1272 }> = [];
1273 for (const chunk of event) {
1274 ASSERT(typeof chunk === "object");
1275 if (Array.isArray(chunk)) {
1276 if (typeof chunk[0] === "string") {
1277 if (chunk[0] === "end") {
1278 resolving = true;
1279 resolvingValue = chunk[1];
1280 } else if (chunk[0] === "error") {
1281 throw new Error(`Progress stream error: ${chunk[1]}`);
1282 } else {
1283 // custom event
1284 // @ts-expect-error TODO: typescript soundness
1285 this.events.emit(...chunk);
1286 }
1287 } else {
1288 const ids = chunk as EncodedKey[];
1289 for (const id of ids) this.pendingParents.set(id, 0 as EncodedKey);
1290 }
1291 } else {
1292 // node update or creation
1293 const {
1294 k: key,
1295 t: text,
1296 v: value,
1297 e: total,
1298 s: showTotal,
1299 p: passive,
1300 h: hidden,
1301 l: messages,
1302 c: children,
1303 E: estimatedTime,
1304 V: formattedValue,
1305 } = chunk;
1306 const node = this.active.get(key);
1307 if (children) {
1308 for (const id of children) this.pendingParents.set(id, key);
1309 updateChildren.push({ key, children });
1310 }
1311 if (node) {
1312 if (text) node.text = text;
1313 if (value) node.value = value;
1314 if (total) node.total = total;
1315 if (showTotal != null) node.showTotal = showTotal;
1316 if (passive != null) node.passive = passive;
1317 if (hidden != null) node.hidden = hidden;
1318 for (const message of messages ?? []) node.log.writeMessage(message);
1319 if (estimatedTime) node.estimatedTime = estimatedTime;
1320 if (formattedValue) node.valueFormatter = () => formattedValue;
1321 for (const id of children ?? []) this.pendingParents.set(id, key);
1322 } else {
1323 ASSERT(text != null, `DstKey ${key} not active but not created`);
1324 this.pendingStart.set(key, {
1325 text,
1326 value,
1327 total,
1328 showTotal,
1329 passive,
1330 hidden,
1331 messages,
1332 estimatedTime,
1333 formattedValue,
1334 });
1335 }
1336 }
1337 }
1338
1339 let entry;
1340 while (entry = this.pendingStart.entries().next().value) {
1341 const [key, value] = entry;
1342 this.startRecursive(key, value);
1343 }
1344 this.pendingParents.clear();
1345
1346 // validate that every referenced child became active; a failure here
1347 // means the stream is desynced
1348 for (const { key, children } of updateChildren) {
1349 UNWRAP(this.active.get(key), () => `parent ${key} not active`);
1350 for (const id of children) {
1351 UNWRAP(this.active.get(id), () => `child ${id} not active`);
1352 }
1353 }
1354
1355 if (resolving) {
1356 this.reset();
1357 this.resolve(resolvingValue as Result);
1358 return true;
1359 }
1360 return false;
1361 }
1362
1363 startRecursive(key: EncodedKey, opts: PendingStart): Ref {
1364 ASSERT(this.pendingStart.delete(key));
1365 // a node whose parent linkage is missing (malformed stream) attaches at
1366 // the root rather than killing the whole decoder
1367 const parentId = this.pendingParents.get(key) ?? (0 as EncodedKey);
1368 let ref: Ref | null = parentId === 0
1369 ? this.target
1370 : this.active.get(parentId) ?? null;
1371 if (!ref) {
1372 ref = this.startRecursive(
1373 parentId,
1374 UNWRAP(this.pendingStart.get(parentId)),
1375 );
1376 }
1377 const { text, estimatedTime, formattedValue, ...options } = opts;
1378 const node = ref.start(text, {
1379 ...options,
1380 estimateCompletion: false,
1381 });
1382 this.active.set(key, node);
1383 node.estimatedTime = estimatedTime;
1384 if (formattedValue) node.valueFormatter = () => formattedValue;
1385 this.pendingParents.delete(key);
1386 return node;
1387 }
1388}
1389
1390interface PendingStart {
1391 text: string;
1392 value?: number;
1393 total?: number;
1394 showTotal?: boolean;
1395 passive?: boolean;
1396 hidden?: boolean;
1397 messages?: log.Message[];
1398 estimatedTime?: number | null;
1399 formattedValue?: string;
1400}
1401
1402/** this interface must not have more than 16 properties */
1403interface DeltaFlags {
1404 // the new status is always present since it's one bit
1405 hidden: boolean;
1406 passive: boolean;
1407 showTotal: boolean;
1408
1409 // for each packet that changed it stores out of line
1410 changedValue: boolean;
1411 changedTotal: boolean;
1412 changedLogs: boolean;
1413 changedText: boolean;
1414 changedChildren: boolean;
1415 changedEstimatedTime: boolean;
1416 changedFormattedValue: boolean;
1417
1418 estimatedTimeNonNull: boolean;
1419
1420 futureProofingPacket: boolean;
1421 deleted: boolean;
1422}
1423
1424/**
1425 * NOTE: this function contains bugs and its format is not yet stabilized.
1426 *
1427 * converts a {@linkcode Root|progress.Root} into an byte stream for
1428 * communicating progress over the process or network boundary. the output is a
1429 * raw binary payload that uses an extremely compact representation for the
1430 * nodes. if trying to transmit UTF-8, consider the JSON-encoding
1431 * {@linkcode encodeEventStream}
1432 *
1433 * if the given root adds custom event handlers, they must all have
1434 * json-serializable payloads.
1435 */
1436export function encodeByteStream<
1437 Result extends ts.Json,
1438 EventMap extends { [key: string]: ts.Json[] },
1439>(
1440 root: Root<Result, EventMap>,
1441 options?: EncodeStreamOptions,
1442): ReadableStream<Uint8Array> {
1443 const events = encodeEventStream(root, options);
1444 const transform = new TransformStream({ transform: writeStreamEvent });
1445 return events.pipeThrough(transform);
1446}
1447
1448/** convert {@linkcode StreamEvent} into a byte payload. Mutates the input event. */
1449function writeStreamEvent(
1450 event: StreamEvent,
1451 writer: stream.WriteTarget,
1452) {
1453 using w = new stream.BufferedWriter(writer);
1454
1455 let deleteCount = event.shift();
1456 ASSERT(typeof deleteCount === "number");
1457
1458 if (deleteCount === -1) {
1459 deleteCount = 0;
1460 w.write(header);
1461 } else {
1462 w.u8("\n".charCodeAt(0));
1463
1464 const deletedNodes = event.splice(0, deleteCount);
1465 w.varUint(deletedNodes.length);
1466 for (const dstKey of deletedNodes) {
1467 ASSERT(typeof dstKey === "number");
1468 w.varUint(dstKey);
1469 }
1470 }
1471
1472 const rootChildren = event.find((x): x is StreamRootChildren => Array.isArray(x) && typeof x[0] === "number");
1473 const changedNodes = event.filter((x): x is StreamNode => typeof x === "object" && !Array.isArray(x));
1474 w.varUint(changedNodes.length + +!!rootChildren);
1475 for (const obj of changedNodes) {
1476 ASSERT(typeof obj === "object" && "k" in obj);
1477 const {
1478 k: key,
1479 t: text,
1480 v: value,
1481 e: total,
1482 l: messages,
1483 c: childrenSet,
1484 E: estimatedTime,
1485 V: formattedValue,
1486 } = obj;
1487 const { hidden, showTotal, passive } = UNWRAP(obj[kNode]);
1488 ASSERT(key > 0);
1489 ASSERT(messages == null || messages.length > 0);
1490
1491 w.varUint(key);
1492 w.u16(encodeDeltaFlags({
1493 hidden,
1494 passive,
1495 showTotal,
1496 changedValue: value !== undefined,
1497 changedTotal: total !== undefined,
1498 changedLogs: messages !== undefined,
1499 changedText: text !== undefined,
1500 changedChildren: childrenSet !== undefined,
1501 changedEstimatedTime: estimatedTime !== undefined,
1502 changedFormattedValue: formattedValue !== undefined,
1503 estimatedTimeNonNull: estimatedTime !== null,
1504 }));
1505 if (value !== undefined) w.f64(value);
1506 if (total !== undefined) w.f64(total);
1507 if (text !== undefined) w.stringWithLength(text);
1508 if (messages !== undefined) {
1509 w.varUint(messages.length);
1510 for (const msg of messages) {
1511 let level = logLevelSerialize.indexOf(msg.level ?? "info");
1512 if (level === -1) level = 0;
1513 w.u8(
1514 level
1515 + (msg.scope ? 1 << 5 : 0)
1516 + (msg.stack && msg.stack.length > 0 ? 1 << 6 : 0)
1517 + (msg.custom !== undefined ? 1 << 7 : 0),
1518 );
1519 w.stringWithLength(msg.text);
1520 w.varUint(Math.floor(msg.time));
1521 if (msg.scope) w.stringWithLength(msg.scope);
1522 if (msg.stack && msg.stack.length > 0) {
1523 w.varUint(msg.stack.length);
1524 for (const { fn, file, line, col } of msg.stack) {
1525 w.stringWithLength(fn ?? "");
1526 w.stringWithLength(file ?? "");
1527 w.varUint(line ?? 0);
1528 w.varUint(col ?? 0);
1529 }
1530 }
1531 if (msg.custom !== undefined) {
1532 w.stringWithLength(JSON.stringify(msg.custom));
1533 }
1534 // TODO: serialize newline attribute
1535 }
1536 }
1537 if (childrenSet !== undefined) {
1538 w.varUint(childrenSet.length);
1539 for (const key of childrenSet) w.varUint(key);
1540 }
1541 if (estimatedTime != null) w.varUint(Math.floor(estimatedTime));
1542 if (formattedValue !== undefined) {
1543 w.stringWithLength(formattedValue);
1544 }
1545 }
1546 if (rootChildren) {
1547 w.varUint(0);
1548 w.varUint(rootChildren.length);
1549 for (const key of rootChildren) {
1550 w.varUint(key);
1551 }
1552 }
1553
1554 const customEvents = event.filter((x): x is StreamCustomEvent => Array.isArray(x) && typeof x[0] === "string");
1555 w.varUint(customEvents.length);
1556 for (const [event, ...args] of customEvents) {
1557 w.stringWithLength(event);
1558 w.stringWithLength(JSON.stringify(args));
1559 }
1560}
1561
1562/**
1563 * decodes a stream created by {@linkcode encodeByteStream} into managed calls to
1564 * `target.start`. when cancelling, all the managed nodes are destroyed.
1565 */
1566export function decodeByteStream<
1567 Result = void,
1568 EventMap extends Events.Map = ts.EmptyObject,
1569>(
1570 encoded: ReadableStream<Uint8Array>,
1571 target: Ref,
1572): async.Cancelable<Result> & Events<EventMap> {
1573 let cancelled = false;
1574 let reader: stream.BufferedReader | null = null;
1575 return decodeEventStream<Result, EventMap>(
1576 new ReadableStream<StreamEvent>({
1577 async start(controller) {
1578 reader = new stream.BufferedReader(encoded.getReader());
1579 try {
1580 // an event boundary with no further bytes is a clean end of stream
1581 while (reader && await reader.ensureAvailableOrFalse(1)) {
1582 controller.enqueue(await readStreamEvent(reader));
1583 }
1584 if (!cancelled) controller.close();
1585 } catch (e) {
1586 if (!cancelled) {
1587 reader?.cancel(e);
1588 throw e;
1589 }
1590 } finally {
1591 reader?.releaseLock();
1592 reader = null;
1593 }
1594 },
1595 cancel(reason) {
1596 cancelled = true;
1597 reader?.releaseLock();
1598 reader = null;
1599 encoded.cancel(reason);
1600 },
1601 }),
1602 target,
1603 );
1604}
1605
1606async function readStreamEvent(r: stream.BufferedReader): Promise<StreamEvent> {
1607 const event: StreamEvent = [];
1608 const start = await r.peekU8();
1609 if (start === "\n".charCodeAt(0)) {
1610 await r.u8();
1611 let deletedCount = await r.varUint();
1612 event.push(deletedCount);
1613 while (deletedCount-- > 0) event.push(await r.varUint());
1614 } else {
1615 const readHeader = await r.readExactly(header.byteLength);
1616 if (!bytes.eql(readHeader, header)) {
1617 const m = JSON.stringify(string.decodeUtf8(readHeader));
1618 throw new Error(`Header mismatch, got ${m}`);
1619 }
1620 event.push(-1);
1621 }
1622 const changedNodes = await r.varUint();
1623 for (let i = 0; i < changedNodes; i += 1) {
1624 const key = await r.varUint();
1625 if (key === 0) {
1626 let count = await r.varUint();
1627 const array: StreamRootChildren = [];
1628 while (count-- > 0) array.push(await r.varUint());
1629 event.push(array);
1630 continue;
1631 }
1632 const flags = decodeDeltaFlags(await r.u16());
1633 const v = flags.changedValue ? await r.f64() : undefined;
1634 const e = flags.changedTotal ? await r.f64() : undefined;
1635 const t = flags.changedText ? await r.stringWithLength() : undefined;
1636 let l: undefined | log.Message[];
1637 if (flags.changedLogs) {
1638 l = [];
1639 let len = await r.varUint();
1640 while (len-- > 0) {
1641 const msgFlags = await r.u8();
1642 const level = logLevelSerialize[msgFlags & 0b1111] ?? "info";
1643 // bitwise tests: `msgFlags && (1 << 5) > 0` parsed as
1644 // `msgFlags && true`, which fabricated scope/stack/custom reads for
1645 // any message with a non-zero level and desynced the entire stream
1646 const hasScope = (msgFlags & (1 << 5)) !== 0;
1647 const hasStack = (msgFlags & (1 << 6)) !== 0;
1648 const hasCustom = (msgFlags & (1 << 7)) !== 0;
1649 const text = await r.stringWithLength();
1650 const time = await r.varUint();
1651 const scope = hasScope ? await r.stringWithLength() : undefined;
1652 let stack: stack.Frame[] | undefined = undefined;
1653 if (hasStack) {
1654 stack = [];
1655 let frames = await r.varUint();
1656 while (frames-- > 0) {
1657 const fn = await r.stringWithLength();
1658 const file = await r.stringWithLength();
1659 const line = await r.varUint();
1660 const col = await r.varUint();
1661 stack.push({
1662 fn: fn || null,
1663 file: file || null,
1664 line: line || null,
1665 col: col || null,
1666 });
1667 }
1668 }
1669 const custom = hasCustom
1670 ? JSON.parse(await r.stringWithLength())
1671 : undefined;
1672 l.push({ level, text, time, scope, stack, custom });
1673 }
1674 }
1675 let c: undefined | EncodedKey[] = undefined;
1676 if (flags.changedChildren) {
1677 c = [];
1678 let len = await r.varUint();
1679 while (len-- > 0) c.push(await r.varUint() as EncodedKey);
1680 }
1681 const E = flags.changedEstimatedTime
1682 ? flags.estimatedTimeNonNull ? await r.varUint() : null
1683 : undefined;
1684 const V = flags.changedFormattedValue
1685 ? await r.stringWithLength()
1686 : undefined;
1687
1688 // a "v2"'s packet will go here, contents ignored by this version.
1689 if (flags.futureProofingPacket) {
1690 const len = await r.varUint();
1691 void await r.readExactly(len);
1692 }
1693
1694 event.push({
1695 k: key as EncodedKey,
1696 t,
1697 v,
1698 e,
1699 s: flags.showTotal,
1700 p: flags.passive,
1701 l,
1702 h: flags.hidden,
1703 c,
1704 E,
1705 V,
1706 });
1707 }
1708
1709 // events
1710 let eventCount = await r.varUint();
1711 while (eventCount-- > 0) {
1712 const name = await r.stringWithLength();
1713 const args = JSON.parse(await r.stringWithLength());
1714 event.push([name, ...args]);
1715 }
1716
1717 return event;
1718}
1719
1720const logLevelSerialize = ["info", "error", "warn", "debug"] as const;
1721
1722function encodeDeltaFlags(node: Partial<DeltaFlags>) {
1723 return (
1724 (node.hidden ? 1 << 0 : 0)
1725 + (node.passive ? 1 << 1 : 0)
1726 + (node.showTotal ? 1 << 2 : 0)
1727 + (node.changedValue ? 1 << 3 : 0)
1728 + (node.changedTotal ? 1 << 4 : 0)
1729 + (node.changedLogs ? 1 << 6 : 0)
1730 + (node.changedText ? 1 << 7 : 0)
1731 + (node.changedChildren ? 1 << 8 : 0)
1732 + (node.changedEstimatedTime ? 1 << 9 : 0)
1733 + (node.changedFormattedValue ? 1 << 10 : 0)
1734 + (node.estimatedTimeNonNull ? 1 << 11 : 0)
1735 + (node.futureProofingPacket ? 1 << 14 : 0)
1736 + (node.deleted ? 1 << 15 : 0)
1737 );
1738}
1739
1740function decodeDeltaFlags(flags: number): DeltaFlags {
1741 return {
1742 hidden: (flags & 1 << 0) !== 0,
1743 passive: (flags & 1 << 1) !== 0,
1744 showTotal: (flags & 1 << 2) !== 0,
1745 changedValue: (flags & 1 << 3) !== 0,
1746 changedTotal: (flags & 1 << 4) !== 0,
1747 changedLogs: (flags & 1 << 6) !== 0,
1748 changedText: (flags & 1 << 7) !== 0,
1749 changedChildren: (flags & 1 << 8) !== 0,
1750 changedEstimatedTime: (flags & 1 << 9) !== 0,
1751 changedFormattedValue: (flags & 1 << 10) !== 0,
1752 estimatedTimeNonNull: (flags & 1 << 11) !== 0,
1753 futureProofingPacket: (flags & 1 << 14) !== 0,
1754 deleted: (flags & 1 << 15) !== 0,
1755 };
1756}
1757
1758class KeyPool<T extends number> {
1759 next = 0;
1760 old: T[] = [];
1761
1762 get() {
1763 const existing = this.old.shift();
1764 if (existing) return existing as T;
1765 return (this.next += 1) as T;
1766 }
1767
1768 recycle(k: T) {
1769 this.old.push(k);
1770 }
1771}
1772
1773/** stateful estimation by providing samples */
1774export interface EstimationAlgorithm {
1775 /**
1776 * given a timestamp `time` and a progress value `progress`, return the
1777 * timestamp that the action will be finished on. called with (t, 0, 0)
1778 * to initialize the algorithm.
1779 */
1780 sample(time: number, current: number, total: number): number | null;
1781}
1782
1783/**
1784 * Exponential Moving Average.
1785 * this is the algorithm used for the default progress estimator.
1786 * https://en.wikipedia.org/wiki/Exponential_smoothing
1787 */
1788export class Ema implements EstimationAlgorithm {
1789 /** milliseconds */
1790 start: number | null = null;
1791 estimate: number | null = null;
1792 /** smoothing factor (0-1) */
1793 alpha: number;
1794 samples = 0;
1795
1796 constructor(alpha: number = 0.25) {
1797 this.alpha = alpha;
1798 }
1799
1800 sample(time: number, current: number, total: number): number | null {
1801 this.samples += 1;
1802
1803 ASSERT(current >= 0 && current <= total);
1804 if (this.start == null) {
1805 this.start = time;
1806 this.estimate = null;
1807 return null;
1808 }
1809
1810 const elapsed = time - this.start;
1811 if (elapsed === 0) return null; // too quick
1812
1813 const rate = current / elapsed;
1814 this.estimate = this.estimate != null
1815 ? (1 - this.alpha) * this.estimate + this.alpha * rate
1816 : rate;
1817 if (elapsed < 10_000 || this.estimate <= 0 || this.samples < 3) return null;
1818 const remaining = total - current;
1819 return time + remaining / this.estimate;
1820 }
1821
1822 reset(): void {
1823 this.start = null;
1824 this.estimate = null;
1825 }
1826}
1827
1828/** according to a stream */
1829type EncodedKey = number & { brand: typeof kNode };
1830
1831const globalKeyPool = new KeyPool<number>();
1832
1833const global: Ref = /** @__PURE__ */ ((root = new Root()) => (attachToScreen(root, log), root))();
1834
1835/**
1836 * a {@linkcode Ref} to the global progress root. unlike referencing the
1837 * namespace import, this value is tree-shakable.
1838 */
1839export const globalRoot: Ref = { start: (text, opts) => global.start(text, opts) };
1840
1841/**
1842 * for testing. not covered by semver
1843 * @internal
1844 */
1845export const internals: {
1846 readStreamEvent: typeof readStreamEvent;
1847 writeStreamEvent: typeof writeStreamEvent;
1848 kNode: symbol;
1849 EncodedKey: EncodedKey;
1850} = /** @__PURE__ */ (() => ({
1851 readStreamEvent,
1852 writeStreamEvent,
1853 kNode: kNode as symbol,
1854 EncodedKey: 0 as EncodedKey,
1855}))();
1856
1857import { ASSERT, UNWRAP } from "./assert.ts";
1858import * as async from "./async.ts";
1859import * as bytes from "./bytes.ts";
1860import * as exceptions from "./exception.ts";
1861import { Events } from "./Events.ts";
1862import * as log from "./log.ts";
1863import * as stack from "./log/stack.ts";
1864import * as stream from "./stream.ts";
1865import * as string from "./string.ts";
1866import * as ansi from "./string/ansi.ts";
1867import * as ts from "./ts.ts";