authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2025-10-15 23:36:01-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-01-14 01:08:17-08:00
log19dfe1d5ef6b974107997e7e6f106df965a96410
tree42101f12805f725e66da65b1c4318b3b6d03e05e
parentdb344d63f2e3db38db746bc34d259e6ff5af1e87
signature Commit is signed but in an unrecognized format.

feat(lib/progress): implement streaming wire protocol

resolves #47 `encodeByteStream` converts these events into a `ReadableStream`. by batching events together, the stream contents remain small, that way the code that constructs progress nodes do not have to worry about calling many setters at once, it gets debounced be the serializer. stream backpressure causes larger time-gaps to be batched (smaller). this enables servers to respond with rich progress. const root = new progress.Root(); doActionWithProgress(root).then(root.end, root.error); // streaming clients indicate a header if (req.headers.get("Accept")?.includes(progress.contentType)) return new Response(progress.encodeByteStream(root), { headers: { 'Content-Type': progress.contentType }, }); // to support non-streaming clients return Response.json(await root.asPromise()); and `decodeByteStream` on the client: const output = document.getElementById("output"); const res = await fetch(...); if (!res.ok) throw ...; const root = new progress.Root(); root.on("change", (active) => { output.innerText = ansi.strip(progress.formatAnsi( performance.now(), active, )); }); const result = await progress.decodeByteStream(res.body, root); output.innerText = JSON.stringify(result); there is currently no document bindings, but i plan to. additionally, a React hook is very trivial to implement for this -- but that is unplanned for this repository. for transports that require JSON or UTF-8, there is `encodeEventStream` which returns a `ReadableStream` of JSON objects which can be compressed at the developer's discretion.

7 files changed, 1841 insertions(+), 32 deletions(-)

lib/bytes.ts created+9
......@@ -0,0 +1,9 @@
1export function eql<T extends ArrayBufferView & { [n: number]: number }>(
2 a: T,
3 b: T,
4) {
5 const l = a.byteLength;
6 if (l !== b.byteLength) return false;
7 for (let i = 0; i < l; i++) if (a[i] !== b[i]) return false;
8 return true;
9}
lib/log.ts+22-18
......@@ -1,5 +1,5 @@
11/**
2 * by using `lib/log.ts`, your application gets easy scoped logging as well as
2 * by using `lib/log.ts`, an application gets easy scoped logging as well as
33 * integration with terminal widgets such as `lib/progress.ts`.
44 *
55 * the pattern for using this module is to shadow the global `console` with a
......@@ -53,23 +53,19 @@ export interface Scope {
5353 error(...args: unknown[]): void;
5454
5555 /**
56 * emit a debugging message. */
56 * emit a debugging message.
57 */
5758 log(...args: unknown[]): void;
5859 /**
5960 * emit a debugging message. this method is in place for compatibility with
6061 * the `console` API. prefer calling `.log` directly.
61 * @internalconst hello = Math.random() > 0.5 ? 1 : "meow";
62
63interface Meow {
64 meow: bigint;
65}
66
67const world = hello as Meow;
68
69console.log(world);
62 * @internal
7063 */
7164 debug(...args: unknown[]): void;
7265
66 /** write a Message object directly. */
67 writeMessage(message: Message): void;
68
7369 /** create a nested sub-scope */
7470 scoped(name: string): Scope;
7571 /** redirect the logging output of this scope somewhere else */
......@@ -88,7 +84,7 @@ export interface Message {
8884 /** captured stack. */
8985 stack?: stack.Frame[];
9086 /** arbitrary data from the logging source. */
91 custom?: Partial<Record<string, unknown>>;
87 custom?: Partial<Record<string, ts.Json>>;
9288 /**
9389 * original logging arguments, if present. this field is indexed by a symbol
9490 * so that it is lost during JSON serialization, as callers are allowed to log
......@@ -177,6 +173,11 @@ export function writeLine(text: string) {
177173 globalWidgetHost.writeLine(text);
178174}
179175
176/** write a Message object directly. */
177export function writeMessage(m: Message) {
178 globalLog.writeMessage(m);
179}
180
180181/**
181182 * while locked, no widgets will draw. prefer `writeLine`.
182183 * this lock is not exclusive.
......@@ -523,7 +524,7 @@ const Scope = class Scope implements Scope {
523524 frames = stack.capture().slice(2);
524525 withinStackCapture = false;
525526 }
526 const m = {
527 this.writeMessage({
527528 level,
528529 scope: this.name,
529530 get text() {
......@@ -534,11 +535,7 @@ const Scope = class Scope implements Scope {
534535 time: Date.now(),
535536 stack: frames,
536537 [originalLogArgs]: args,
537 };
538 if (withinDispatch) return void globalOutputFunction(m);
539 withinDispatch = true;
540 this.#dispatch(m);
541 withinDispatch = false;
538 });
542539 }
543540
544541 info: (...args: unknown[]) => void = (...args: unknown[]) => {
......@@ -557,6 +554,13 @@ const Scope = class Scope implements Scope {
557554 this.#log("debug", args);
558555 };
559556
557 writeMessage: (m: Message) => void = (m) => {
558 if (withinDispatch) return void globalOutputFunction(m);
559 withinDispatch = true;
560 this.#dispatch(m);
561 withinDispatch = false;
562 };
563
560564 scoped(name: string): Scope {
561565 const current = this.name;
562566 return new Scope(
lib/progress.test.ts+146-4
......@@ -1,5 +1,5 @@
11// a trivial example of how to use progress
2test("trivial end-to-end example", () => {
2test.skip("trivial end-to-end example", () => {
33 const { fgBlue: FB, fgReset: FR, reset: R } = ansi;
44 const screen = new testing.MockScreen(); // create a mock terminal
55 const root = new progress.Root(screen); // sync with mock timers
......@@ -43,6 +43,39 @@ test("trivial end-to-end example", () => {
4343 });
4444});
4545
46describe("encodeEventStream", async (t) => {
47 await test("simple", async () => {
48 const root = new progress.Root(); // sync with mock timers
49
50 const a = root.start("hello");
51 const b = root.start("cats", { value: 2, total: 5 });
52
53 const events = progress.encodeEventStream(root);
54 const reader = events.getReader();
55 const next = async () =>
56 cleanProgressEvent(UNWRAP((await reader.read()).value));
57
58 assert.deepEqual(await next(), [
59 -1,
60 { k: 1, t: "hello" },
61 { k: 2, t: "cats", v: 2, e: 5 },
62 [1, 2],
63 ]);
64 b.inc();
65 b.inc();
66 assert.deepEqual(await next(), [
67 0,
68 { k: 2, v: 4 },
69 ]);
70 b.inc();
71 assert.deepEqual(await next(), [
72 1,
73 2,
74 [1],
75 ]);
76 });
77});
78
4679// ## value formatters
4780test("slashValueFormatter", ({ mock }) => {
4881 const fn = mock.fn((arg: number) => `<${arg}>`);
......@@ -55,11 +88,13 @@ test("slashValueFormatter", ({ mock }) => {
5588});
5689test("bytesValueFormatter", () => {
5790 assert.equal(progress.bytesValueFormatter(250, 0), "250B");
58 assert.equal(progress.bytesValueFormatter(0, 1_200_000), "0B/1.2MB");
91 assert.equal(progress.bytesValueFormatter(50_000, 1_200_000), "50kB/1.20MB");
92 assert.equal(progress.bytesValueFormatter(0, 0), "");
5993});
6094test("defaultValueFormatter", () => {
6195 assert.equal(progress.defaultValueFormatter(250, 0), "250");
62 assert.equal(progress.defaultValueFormatter(0, 1_200_000), "0/1200000");
96 assert.equal(progress.defaultValueFormatter(52, 1_200_000), "52/1200000");
97 assert.equal(progress.defaultValueFormatter(0, 1_200_000), "");
6398});
6499
65100test("percentValueFormatter", () => {
......@@ -95,8 +130,115 @@ test("Ema", () => {
95130 assert.equal(ema.sample(50_130, 0.9), 50_141.754246587305);
96131});
97132
133// ## internals
134test("event encoding round trip cases", async () => {
135 const fakeNode = progress.nullNode;
136 async function roundTrip(event: progress.StreamEvent) {
137 const [readable, controller] = stream.readWritePair();
138 const chunks: Uint8Array[] = [];
139 progress.internals.writeStreamEvent(event.slice(), (d) => chunks.push(d));
140 const combined = Buffer.concat(chunks);
141 // console.debug(JSON.stringify(combined.toString("utf-8")));
142 controller.enqueue(combined);
143 controller.close();
144 using reader = new stream.BufferedReader(readable.getReader());
145 const decoded = await progress.internals.readStreamEvent(reader);
146 assert.deepEqual(cleanProgressEvent(decoded), cleanProgressEvent(event));
147 }
148
149 await roundTrip([
150 -1,
151 {
152 k: 1 as EncodedKey,
153 t: "job",
154 [progress.internals.kNode]: fakeNode,
155 },
156 [1],
157 ]);
158 await roundTrip([
159 0,
160 {
161 k: 1 as EncodedKey,
162 t: "job",
163 [progress.internals.kNode]: fakeNode,
164 },
165 {
166 k: 2 as EncodedKey,
167 t: "deeply nested",
168 c: [1 as EncodedKey],
169 [progress.internals.kNode]: fakeNode,
170 },
171 {
172 k: 3 as EncodedKey,
173 t: "other task",
174 c: [2 as EncodedKey],
175 [progress.internals.kNode]: fakeNode,
176 },
177 {
178 k: 4 as EncodedKey,
179 v: 4,
180 e: 7,
181 t: "purring",
182 E: 1768210855277.92,
183 [progress.internals.kNode]: fakeNode,
184 },
185 {
186 k: 5 as EncodedKey,
187 v: 2,
188 e: 7,
189 t: "meowing",
190 E: 1768210857407,
191 V: "28.6%",
192 [progress.internals.kNode]: fakeNode,
193 },
194 {
195 k: 6 as EncodedKey,
196 t: "subtask 1t07l",
197 [progress.internals.kNode]: fakeNode,
198 },
199 {
200 k: 7 as EncodedKey,
201 t: "subtask 42k1d",
202 c: [3 as EncodedKey],
203 [progress.internals.kNode]: fakeNode,
204 },
205 {
206 k: 8 as EncodedKey,
207 t: "subtask 4kdjg",
208 c: [5 as EncodedKey, 4 as EncodedKey],
209 [progress.internals.kNode]: fakeNode,
210 },
211 {
212 k: 9 as EncodedKey,
213 e: 30,
214 t: "progress node",
215 c: [8, 7, 6] as EncodedKey[],
216 [progress.internals.kNode]: fakeNode,
217 },
218 [9],
219 ]);
220});
221
222function cleanProgressEvent(event: progress.StreamEvent) {
223 return event.map((x) =>
224 typeof x === "object" && !Array.isArray(x)
225 ? testing.removeUndefinedKeys({
226 ...x,
227 E: undefined,
228 p: undefined,
229 s: undefined,
230 h: undefined,
231 })
232 : x
233 );
234}
235
236type EncodedKey = typeof progress.internals.EncodedKey;
237
98238import * as progress from "./progress.ts";
99import { test } from "node:test";
239import { describe, test } from "node:test";
100240import assert from "node:assert/strict";
241import * as stream from "./stream.ts";
101242import * as testing from "./testing.ts";
102243import * as ansi from "./string/ansi.ts";
244import { UNWRAP } from "./assert.ts";
lib/progress.ts+891-10
......@@ -48,6 +48,7 @@
4848 * that rests at the end of the log. if `estimate` is given, a progress bar is
4949 * shown.
5050 */
51/* node:coverage ignore next 3 */
5152export function start(text: string, opts?: StartOptions): Node {
5253 return global.start(text, opts);
5354}
......@@ -267,6 +268,7 @@ export class Root<
267268 * when given an estimate, a progress bar is rendered.
268269 */
269270 start(text: string, opts?: StartOptions): Node {
271 ASSERT(this.#status === "active");
270272 const [state, node] = newNode(this, text, opts);
271273 this.#active.push(state);
272274 this.emit("node-start", state, null);
......@@ -452,10 +454,11 @@ function newNode<Result, EventMap extends Events.Map>(
452454 key === "value" && estimator && state.total > 0 && state.value > 0 &&
453455 state.value < state.total
454456 ) {
455 const time = estimator.sample(owner.now(), state.value / state.total);
456 if (
457 time && (!state.estimatedTime || state.estimatedTime !== time)
458 ) {
457 // originally, auto-estimation would use the Root's timing primitives,
458 // but this makes it so that `state.estimatedTime` is in terms of the
459 // Root, which is unintuitive for everyone but the renderer.
460 const time = estimator.sample(Date.now(), state.value / state.total);
461 if (time && (!state.estimatedTime || state.estimatedTime !== time)) {
459462 state.estimatedTime = time;
460463 owner.emit("node-change", state, "estimatedTime");
461464 }
......@@ -590,6 +593,7 @@ export const nullNode: Node = {
590593 error: noop,
591594 log: noop,
592595 debug: noop,
596 writeMessage: noop,
593597 scoped() {
594598 return this;
595599 },
......@@ -782,15 +786,873 @@ export function attachToScreen(
782786 return ts.defer(() => stack[Symbol.dispose]);
783787}
784788
789const header = /* @__PURE__ */ string.encodeUtf8("clover's progress <3\n");
790
791/** the 'Content-Type' of {@linkcode encodeByteStream}'s return value */
792export const contentType = "application/x-clover-progress";
793
794/**
795 * opaque but JSON-compatible serialized state, either containing the entire
796 * progress tree or just a delta update. this data structure always be future
797 * and backwards compatbile with changes to this library. if a breaking change
798 * is made, the APIs in this file will be renamed.
799 *
800 * for details on the format, read the source of {@linkcode encodeEventStream} or
801 * {@linkcode encodeByteStream}
802 */
803export type StreamEvent = Array<
804 number | StreamNode | StreamCustomEvent | StreamRootChildren
805>;
806const kNode = Symbol("underlyingNode");
807/** see {@linkcode StreamEvent} @internal */
808export interface StreamNode {
809 k: EncodedKey;
810 /** text, overwrite */
811 t?: string | undefined;
812 /** value, overwrite */
813 v?: number | undefined;
814 /** total, overwrite */
815 e?: number | undefined;
816 /** showTotal, overwrite */
817 s?: boolean | undefined;
818 /** passive, overwrite */
819 p?: boolean | undefined;
820 /** logs, append */
821 l?: log.Message[] | undefined;
822 /** hidden, overwrite */
823 h?: boolean | undefined;
824 /** children, overwrite */
825 c?: EncodedKey[] | undefined;
826 /** estimated time, in unix milliseconds, overwrite */
827 E?: number | null | undefined;
828 /** formatted value if not defaultValueFormatter, overwrite */
829 V?: string | undefined;
830
831 /** access to the underlying node is used by {@linkcode encodeByteStream} */
832 [kNode]?: ReadOnlyNode;
833}
834type StreamCustomEvent = [
835 /** custom event name or "end" */
836 e: "end" | string,
837 /** custom event payload */
838 ...args: ts.Json[],
839];
840type StreamRootChildren = number[];
841
842export interface EncodeStreamOptions {
843 /**
844 * set the minimum time between packets.
845 * @default 1000 / 30 (30fps)
846 */
847 throttleMs?: number;
848 serializeStackTraces?: boolean;
849}
850
851/**
852 * converts a {@linkcode Root|progress.Root} into an object stream
853 * for communicating progress over the process or network boundary.
854 *
855 * if the given root adds custom event handlers, they must all have
856 * json-serializable payloads.
857 *
858 * this function does not use recursion.
859 */
860export function encodeEventStream<
861 Result extends ts.Json,
862 Map extends { [key: string]: ts.Json[] },
863>(
864 // fun fact, this type is impossible to write without generics
865 root: Root<Result, Map>,
866 {
867 throttleMs = 1000 / 30, /* 30fps */
868 serializeStackTraces = false,
869 }: EncodeStreamOptions = {},
870): ReadableStream<StreamEvent> {
871 const { delay, now } = root;
872 const stack = new DisposableStack();
873 const s = new Encoder(root, serializeStackTraces);
874 const ready = new async.Watch(false);
875
876 let lastEvent = now();
877 let timer: async.Cancelable<void> | null = null;
878 function emitSoon() {
879 if (timer || ready.value) return;
880 const remaining = now() - lastEvent + throttleMs;
881 if (remaining > 0) {
882 (timer = delay(remaining)).then(() => {
883 timer = null;
884 ready.value = true;
885 lastEvent = now();
886 });
887 } else {
888 ready.value = true;
889 lastEvent = now();
890 }
891 }
892 stack.defer(() => {
893 if (timer) timer.cancel();
894 s.changed.clear();
895 s.deleted.clear();
896 });
897
898 let isFirst = true;
899 return new ReadableStream({
900 start(controller) {
901 if (root.active.length > 0) {
902 controller.enqueue(s.getDelta(isFirst));
903 isFirst = false;
904 }
905
906 // TODO: why does ts require two args
907 stack.use(root.on("node-start", (newNode, _) => {
908 if (!newNode.parent) s.rootChildrenUpdated = true;
909 s.changed.set(newNode, new Set());
910 emitSoon();
911 }));
912 stack.use(root.on("node-change", (node, key) => {
913 let set = s.changed.get(node);
914 if (!set) s.changed.set(node, set = new Set());
915 set.add(key);
916 emitSoon();
917 }));
918 stack.use(root.on("node-end", (node) => {
919 if (!node.parent) s.rootChildrenUpdated = true;
920 s.changed.delete(node);
921 s.deleted.add(node); // handle add and remove in same frame
922 emitSoon();
923 }));
924 },
925 async pull(controller) {
926 await ready.until((x) => x === true);
927 ready.value = false;
928 if (timer) {
929 timer.cancel();
930 timer = null;
931 }
932 controller.enqueue(s.getDelta(isFirst));
933 isFirst = false;
934 },
935 cancel() {
936 stack.dispose();
937 (0, performance.now)();
938 },
939 }, { highWaterMark: 0 });
940}
941
942/**
943 * decodes a stream created by `encodeEventStream` into managed calls to
944 * `target.start`. when cancelling, all the created nodes are destroyed.
945 */
946export function decodeEventStream<
947 Result = void,
948 EventMap extends Events.Map = ts.EmptyObject,
949>(
950 encoded: ReadableStream<StreamEvent>,
951 target: Ref,
952): async.Cancelable<Result> & Events<EventMap> {
953 const { resolve, reject, promise } = Promise.withResolvers<Result>();
954 const events = new Events<EventMap>();
955 const ctrl = new AbortController();
956 const decoder = new Decoder(target, events, resolve);
957 const { signal } = ctrl;
958 (async () => {
959 if (signal.aborted) return;
960 const reader = encoded.getReader();
961 try {
962 let hasEmittedEnd = false;
963 while (true) {
964 signal.throwIfAborted();
965 const { value, done } = await reader.read();
966 signal.throwIfAborted();
967 if (done) break;
968 decoder.processEvent(value);
969 }
970 ASSERT(hasEmittedEnd, "Stream terminated early.");
971 } catch (err) {
972 reader.cancel(err);
973 if (!signal.aborted) reject(err);
974 } finally {
975 reader.releaseLock();
976 }
977 })();
978 return Object.assign(ts.mixin(promise, events), {
979 cancel: (e?: unknown) => ctrl.abort(e),
980 [Symbol.dispose]: () => ctrl.abort(),
981 });
982}
983
984/**
985 * by listening to the event interface in {@linkcode Root}, compute delta
986 * events for a progress stream. the caller is expected to fill `changed` and `deleted`
987 */
988class Encoder<
989 Result extends ts.Json,
990 EventMap extends { [key: string]: ts.Json[] },
991> {
992 root: Root<Result, EventMap>;
993 pool: KeyPool<EncodedKey> = new KeyPool();
994 keyMap: Map<ReadOnlyNode, EncodedKey> = new Map();
995 logsLength: WeakMap<ReadOnlyNode, number> = new WeakMap();
996 serializeStackTraces: boolean;
997
998 rootChildrenUpdated = true;
999 changed = new Map<ReadOnlyNode, Set<keyof ReadOnlyNode>>();
1000 deleted = new Set<ReadOnlyNode>();
1001
1002 constructor(
1003 root: Root<Result, EventMap>,
1004 serializeStackTraces: boolean = false,
1005 ) {
1006 this.root = root;
1007 this.serializeStackTraces = serializeStackTraces;
1008 for (const node of root.active) {
1009 this.changed.set(node, new Set());
1010 }
1011 }
1012
1013 getDelta(isFirst = false) {
1014 const deleted = Array.from(
1015 this.deleted,
1016 (node) => [node, this.keyMap.get(node)] as const,
1017 )
1018 .filter(([node, dst]) => {
1019 if (dst == null) return false; // started and stopped in same tick
1020 this.pool.recycle(dst);
1021 this.keyMap.delete(node);
1022 return true;
1023 })
1024 .map((x) => UNWRAP(x[1]));
1025 ASSERT(!isFirst || deleted.length === 0);
1026
1027 const { updated = [], started = [] } = Object.groupBy(
1028 this.changed,
1029 ([node]) => this.keyMap.has(node) ? "updated" : "started",
1030 );
1031
1032 const payload: StreamEvent = [
1033 isFirst ? -1 : deleted.length,
1034 ...deleted,
1035 ];
1036
1037 if (started.length > 0) {
1038 this.beginNodesDepthFirst(
1039 started.reverse().map(([node]) => node),
1040 payload,
1041 );
1042 }
1043
1044 for (const [node, props] of updated) {
1045 let logs = undefined;
1046 if (
1047 props.has("logs") &&
1048 node.logs.length > (this.logsLength.get(node) ?? 0)
1049 ) {
1050 logs = this.redactLogsIfNeeded(
1051 node.logs.slice(this.logsLength.get(node) ?? 0),
1052 );
1053 this.logsLength.set(node, node.logs.length);
1054 }
1055 payload.push({
1056 k: UNWRAP(this.keyMap.get(node)),
1057 t: props.has("text") ? node.text : undefined,
1058 v: props.has("value") ? node.value : undefined,
1059 e: props.has("total") ? node.total : undefined,
1060 s: props.has("showTotal") ? node.showTotal : undefined,
1061 p: props.has("passive") ? node.passive : undefined,
1062 l: logs,
1063 h: props.has("hidden") ? node.hidden : undefined,
1064 c: props.has("children")
1065 ? node.children.map((node) => UNWRAP(this.keyMap.get(node)))
1066 : undefined,
1067 E: props.has("estimatedTime") ? node.estimatedTime : undefined,
1068 V: node.valueFormatter !== defaultValueFormatter &&
1069 (props.has("value") || props.has("total") ||
1070 props.has("valueFormatter"))
1071 ? node.valueFormatter(node.value, node.total)
1072 : undefined,
1073 [kNode]: node,
1074 });
1075 }
1076
1077 if (this.rootChildrenUpdated) {
1078 payload.push(
1079 this.root.active.map((node) => UNWRAP(this.keyMap.get(node))),
1080 );
1081 }
1082
1083 this.deleted.clear();
1084 this.changed.clear();
1085 this.rootChildrenUpdated = false;
1086
1087 return payload;
1088 }
1089
1090 beginNodesDepthFirst(
1091 nodes: ReadOnlyNode[],
1092 out: StreamEvent,
1093 ) {
1094 for (let i = 0; i < nodes.length; i += 1) {
1095 const node = UNWRAP(nodes[i]);
1096 nodes.push(...node.children);
1097 }
1098 let next;
1099 while (next = nodes.pop()) {
1100 if (!this.keyMap.has(next)) {
1101 out.push(this.beginNode(next));
1102 }
1103 }
1104 }
1105
1106 beginNode(node: ReadOnlyNode): StreamNode {
1107 const k = this.pool.get();
1108 this.keyMap.set(node, k);
1109
1110 return {
1111 k,
1112 v: node.value !== 0 ? node.value : undefined,
1113 e: node.total !== 0 ? node.total : undefined,
1114 t: node.text.length > 0 ? node.text : undefined,
1115 l: node.logs.length > 0
1116 ? (this.logsLength.set(node, node.logs.length),
1117 this.redactLogsIfNeeded(node.logs.slice()))
1118 : undefined,
1119 c: node.children.length > 0
1120 ? node.children.map((node) =>
1121 UNWRAP(this.keyMap.get(node), () => [node.children, this.keyMap])
1122 )
1123 : undefined,
1124 p: node.passive === true ? true : undefined,
1125 s: node.showTotal === false ? false : undefined,
1126 h: node.hidden === true ? true : undefined,
1127 E: node.estimatedTime != null ? node.estimatedTime : undefined,
1128 V: node.valueFormatter !== defaultValueFormatter
1129 ? node.valueFormatter(node.value, node.total)
1130 : undefined,
1131 [kNode]: node,
1132 };
1133 }
1134
1135 redactLogsIfNeeded(messages: log.Message[]) {
1136 if (this.serializeStackTraces) return messages;
1137 return messages.map((x) => ({ ...x, stack: undefined }));
1138 }
1139}
1140
1141class Decoder<
1142 Result = void,
1143 EventMap extends Events.Map = ts.EmptyObject,
1144> {
1145 target: Ref;
1146 events: Events<EventMap>;
1147 resolve: (result: Result) => void;
1148 active = new Map<EncodedKey, Node>();
1149 children = new Map<Node, Node[]>();
1150 pendingStart = new Map<EncodedKey, PendingStart>();
1151 pendingParents = new Map<EncodedKey, EncodedKey>();
1152
1153 constructor(
1154 target: Ref,
1155 events: Events<EventMap>,
1156 resolve: (result: Result) => void,
1157 ) {
1158 this.target = target;
1159 this.events = events;
1160 this.resolve = resolve;
1161 }
1162
1163 reset() {
1164 for (const v of this.active.values()) v.end();
1165 this.active.clear();
1166 }
1167
1168 processEvent(event: StreamEvent) {
1169 event = event.slice();
1170
1171 // delete unreferenced nodes first, since they may get re-allocated
1172 // in the next phase.
1173 const deletedCount = event.shift();
1174 ASSERT(typeof deletedCount === "number");
1175 if (deletedCount === -1) {
1176 // connection reset, delete all nodes
1177 this.reset();
1178 } else {
1179 for (const id of event.splice(0, deletedCount)) {
1180 ASSERT(typeof id === "number");
1181 UNWRAP(
1182 this.active.get(id as EncodedKey),
1183 () => `dstkey ${id} not found, ${[...this.active.keys()]}`,
1184 ).end();
1185 ASSERT(this.active.delete(id as EncodedKey));
1186 }
1187 }
1188
1189 // parse all of the events, updating all nodes and extracting creations
1190 let resolving = false;
1191 let resolvingValue: unknown = null;
1192 let updateChildren: Array<{
1193 key: EncodedKey;
1194 children: EncodedKey[];
1195 }> = [];
1196 for (const chunk of event) {
1197 ASSERT(typeof chunk === "object");
1198 if (Array.isArray(chunk)) {
1199 if (typeof chunk[0] === "string") {
1200 if (chunk[0] === "end") {
1201 resolving = true;
1202 resolvingValue = chunk[1];
1203 } else {
1204 // custom event
1205 // @ts-expect-error TODO: typescript soundness
1206 this.events.emit(...chunk);
1207 }
1208 } else {
1209 const ids = chunk as EncodedKey[];
1210 for (const id of ids) this.pendingParents.set(id, 0 as EncodedKey);
1211 }
1212 } else {
1213 // node update or creation
1214 const {
1215 k: key,
1216 t: text,
1217 v: value,
1218 e: total,
1219 l: messages,
1220 c: children,
1221 E: estimatedTime,
1222 V: formattedValue,
1223 } = chunk;
1224 const node = this.active.get(key);
1225 if (children) {
1226 for (const id of children) this.pendingParents.set(id, key);
1227 updateChildren.push({ key, children });
1228 }
1229 if (node) {
1230 if (text) node.text = text;
1231 if (value) node.value = value;
1232 if (total) node.total = total;
1233 for (const message of messages ?? []) node.log.writeMessage(message);
1234 if (estimatedTime) node.estimatedTime = estimatedTime;
1235 if (formattedValue) node.valueFormatter = () => formattedValue;
1236 for (const id of children ?? []) this.pendingParents.set(id, key);
1237 } else {
1238 ASSERT(text != null, `DstKey ${key} not active but not created`);
1239 this.pendingStart.set(key, {
1240 text,
1241 value,
1242 total,
1243 messages,
1244 estimatedTime,
1245 formattedValue,
1246 });
1247 }
1248 }
1249 }
1250
1251 let entry;
1252 while (entry = this.pendingStart.entries().next().value) {
1253 const [key, value] = entry;
1254 this.startRecursive(key, value);
1255 }
1256 this.pendingParents.clear();
1257
1258 for (const { key, children } of updateChildren) {
1259 const node = UNWRAP(this.active.get(key));
1260 this.children.set(
1261 node,
1262 children.map((id) => UNWRAP(this.active.get(id))),
1263 );
1264 }
1265 }
1266
1267 startRecursive(key: EncodedKey, opts: PendingStart): Ref {
1268 ASSERT(this.pendingStart.delete(key));
1269 const parentId = UNWRAP(this.pendingParents.get(key));
1270 let ref: Ref | null = parentId === 0
1271 ? this.target
1272 : this.active.get(parentId) ?? null;
1273 if (!ref) {
1274 ref = this.startRecursive(
1275 parentId,
1276 UNWRAP(this.pendingStart.get(parentId)),
1277 );
1278 }
1279 const { text, estimatedTime, ...options } = opts;
1280 const node = ref.start(text, {
1281 ...options,
1282 estimateCompletion: false,
1283 });
1284 this.active.set(key, node);
1285 node.estimatedTime = estimatedTime;
1286 this.pendingParents.delete(key);
1287 return node;
1288 }
1289}
1290
1291interface PendingStart {
1292 text: string;
1293 value?: number;
1294 total?: number;
1295 messages?: log.Message[];
1296 estimatedTime?: number | null;
1297 formattedValue?: string;
1298}
1299
1300/** this interface must not have more than 16 properties */
1301interface DeltaFlags {
1302 // the new status is always present since it's one bit
1303 hidden: boolean;
1304 passive: boolean;
1305 showTotal: boolean;
1306
1307 // for each packet that changed it stores out of line
1308 changedValue: boolean;
1309 changedTotal: boolean;
1310 changedLogs: boolean;
1311 changedText: boolean;
1312 changedChildren: boolean;
1313 changedEstimatedTime: boolean;
1314 changedFormattedValue: boolean;
1315
1316 estimatedTimeNonNull: boolean;
1317
1318 futureProofingPacket: boolean;
1319 deleted: boolean;
1320}
1321
1322/**
1323 * converts a {@linkcode Root|progress.Root} into an byte stream for
1324 * communicating progress over the process or network boundary. the output is a
1325 * raw binary payload that uses an extremely compact representation for the
1326 * nodes. if trying to transmit UTF-8, consider the JSON-encoding
1327 * {@linkcode encodeEventStream}
1328 *
1329 * if the given root adds custom event handlers, they must all have
1330 * json-serializable payloads.
1331 *
1332 * this function does not use recursion.
1333 */
1334export function encodeByteStream<
1335 Result extends ts.Json,
1336 EventMap extends { [key: string]: ts.Json[] },
1337>(
1338 root: Root<Result, EventMap>,
1339 options?: EncodeStreamOptions,
1340): ReadableStream<Uint8Array> {
1341 const events = encodeEventStream(root, options);
1342 const transform = new TransformStream({ transform: writeStreamEvent });
1343 return events.pipeThrough(transform);
1344}
1345
1346/** convert {@linkcode StreamEvent} into a byte payload. Mutates the input event. */
1347function writeStreamEvent(
1348 event: StreamEvent,
1349 writer: stream.WriteTarget,
1350) {
1351 using w = new stream.BufferedWriter(writer);
1352
1353 let deleteCount = event.shift();
1354 ASSERT(typeof deleteCount === "number");
1355
1356 if (deleteCount === -1) {
1357 deleteCount = 0;
1358 w.write(header);
1359 } else {
1360 w.u8("\n".charCodeAt(0));
1361
1362 const deletedNodes = event.splice(0, deleteCount);
1363 w.varUint(deletedNodes.length);
1364 for (const dstKey of deletedNodes) {
1365 ASSERT(typeof dstKey === "number");
1366 w.varUint(dstKey);
1367 }
1368 }
1369
1370 const rootChildren = event.find((x): x is StreamRootChildren =>
1371 Array.isArray(x) && typeof x[0] === "number"
1372 );
1373 const changedNodes = event.filter((x): x is StreamNode =>
1374 typeof x === "object" && !Array.isArray(x)
1375 );
1376 w.varUint(changedNodes.length + +!!rootChildren);
1377 for (const obj of changedNodes) {
1378 ASSERT(typeof obj === "object" && "k" in obj);
1379 const {
1380 k: key,
1381 t: text,
1382 v: value,
1383 e: total,
1384 l: messages,
1385 c: childrenSet,
1386 E: estimatedTime,
1387 V: formattedValue,
1388 } = obj;
1389 const { hidden, showTotal, passive } = UNWRAP(obj[kNode]);
1390 ASSERT(key > 0);
1391 ASSERT(messages == null || messages.length > 0);
1392
1393 w.varUint(key);
1394 w.u16(encodeDeltaFlags({
1395 hidden,
1396 passive,
1397 showTotal,
1398 changedValue: value !== undefined,
1399 changedTotal: total !== undefined,
1400 changedLogs: messages !== undefined,
1401 changedText: text !== undefined,
1402 changedChildren: childrenSet !== undefined,
1403 changedEstimatedTime: estimatedTime !== undefined,
1404 changedFormattedValue: formattedValue !== undefined,
1405 estimatedTimeNonNull: estimatedTime !== null,
1406 }));
1407 if (value !== undefined) w.f64(value);
1408 if (total !== undefined) w.f64(total);
1409 if (text !== undefined) w.stringWithLength(text);
1410 if (messages !== undefined) {
1411 w.varUint(messages.length);
1412 for (const msg of messages) {
1413 let level = logLevelSerialize.indexOf(msg.level);
1414 if (level === -1) level = 0;
1415 w.u8(
1416 level +
1417 (msg.scope ? 1 << 5 : 0) +
1418 (msg.stack && msg.stack.length > 0 ? 1 << 6 : 0) +
1419 (msg.custom !== undefined ? 1 << 7 : 0),
1420 );
1421 w.stringWithLength(msg.text);
1422 w.varUint(Math.floor(msg.time));
1423 if (msg.scope) w.stringWithLength(msg.scope);
1424 if (msg.stack && msg.stack.length > 0) {
1425 w.varUint(msg.stack.length);
1426 for (const { fn, file, line, col } of msg.stack) {
1427 w.stringWithLength(fn ?? "");
1428 w.stringWithLength(file ?? "");
1429 w.varUint(line ?? 0);
1430 w.varUint(col ?? 0);
1431 }
1432 }
1433 if (msg.custom !== undefined) {
1434 w.stringWithLength(JSON.stringify(msg.custom));
1435 }
1436 }
1437 }
1438 if (childrenSet !== undefined) {
1439 w.varUint(childrenSet.length);
1440 for (const key of childrenSet) w.varUint(key);
1441 }
1442 if (estimatedTime != null) w.varUint(Math.floor(estimatedTime));
1443 if (formattedValue !== undefined) {
1444 w.stringWithLength(formattedValue);
1445 }
1446 }
1447 if (rootChildren) {
1448 w.varUint(0);
1449 w.varUint(rootChildren.length);
1450 for (const key of rootChildren) {
1451 w.varUint(key);
1452 }
1453 }
1454
1455 const customEvents = event.filter((x): x is StreamCustomEvent =>
1456 Array.isArray(x) && typeof x[0] === "string"
1457 );
1458 w.varUint(customEvents.length);
1459 for (const [event, ...args] of customEvents) {
1460 w.stringWithLength(event);
1461 w.stringWithLength(JSON.stringify(args));
1462 }
1463}
1464
1465/**
1466 * decodes a stream created by `encodeByteStream` into managed calls to
1467 * `target.start`. when cancelling, all the created nodes are destroyed.
1468 */
1469export function decodeByteStream<
1470 Result = void,
1471 EventMap extends Events.Map = ts.EmptyObject,
1472>(
1473 encoded: ReadableStream<Uint8Array>,
1474 target: Ref,
1475): async.Cancelable<Result> & Events<EventMap> {
1476 let cancelled = false;
1477 return decodeEventStream<Result, EventMap>(
1478 new ReadableStream<StreamEvent>({
1479 async start(controller) {
1480 using reader = new stream.BufferedReader(encoded.getReader());
1481 try {
1482 while (true) controller.enqueue(await readStreamEvent(reader));
1483 } catch (e) {
1484 if (!cancelled) {
1485 reader.cancel(e);
1486 throw e;
1487 }
1488 }
1489 },
1490 cancel(reason) {
1491 cancelled = true;
1492 encoded.cancel(reason);
1493 },
1494 }),
1495 target,
1496 );
1497}
1498
1499async function readStreamEvent(r: stream.BufferedReader): Promise<StreamEvent> {
1500 const event: StreamEvent = [];
1501 const start = await r.peekU8();
1502 if (start === "\n".charCodeAt(0)) {
1503 await r.u8();
1504 let deletedCount = await r.varUint();
1505 event.push(deletedCount);
1506 while (deletedCount-- > 0) event.push(await r.varUint());
1507 } else {
1508 const readHeader = await r.readExactly(header.byteLength);
1509 if (!bytes.eql(readHeader, header)) {
1510 const m = JSON.stringify(string.decodeUtf8(readHeader));
1511 throw new Error(`Header mismatch, got ${m}`);
1512 }
1513 event.push(-1);
1514 }
1515 const changedNodes = await r.varUint();
1516 for (let i = 0; i < changedNodes; i += 1) {
1517 const key = await r.varUint();
1518 if (key === 0) {
1519 let count = await r.varUint();
1520 const array: StreamRootChildren = [];
1521 while (count-- > 0) array.push(await r.varUint());
1522 event.push(array);
1523 continue;
1524 }
1525 const flags = decodeDeltaFlags(await r.u16());
1526 const v = flags.changedValue ? await r.f64() : undefined;
1527 const e = flags.changedTotal ? await r.f64() : undefined;
1528 const t = flags.changedText ? await r.stringWithLength() : undefined;
1529 let l: undefined | log.Message[];
1530 if (flags.changedLogs) {
1531 l = [];
1532 let len = await r.varUint();
1533 while (len > 0) {
1534 const msgFlags = await r.u8();
1535 const level = logLevelSerialize[msgFlags & 0b1111] ?? "info";
1536 const hasScope = msgFlags && (1 << 5) > 0;
1537 const hasStack = msgFlags && (1 << 6) > 0;
1538 const hasCustom = msgFlags && (1 << 7) > 0;
1539 const text = await r.stringWithLength();
1540 const time = await r.varUint();
1541 const scope = hasScope ? await r.stringWithLength() : undefined;
1542 let stack: stack.Frame[] | undefined = undefined;
1543 if (hasStack) {
1544 stack = [];
1545 let len = await r.varUint();
1546 while (len > 0) {
1547 const fn = await r.stringWithLength();
1548 const file = await r.stringWithLength();
1549 const line = await r.varUint();
1550 const col = await r.varUint();
1551 stack.push({
1552 fn: fn || null,
1553 file: file || null,
1554 line: line || null,
1555 col: col || null,
1556 });
1557 }
1558 }
1559 const custom = hasCustom
1560 ? JSON.parse(await r.stringWithLength())
1561 : undefined;
1562 l.push({ level, text, time, scope, stack, custom });
1563 }
1564 }
1565 let c: undefined | EncodedKey[] = undefined;
1566 if (flags.changedChildren) {
1567 c = [];
1568 let len = await r.varUint();
1569 while (len-- > 0) c.push(await r.varUint() as EncodedKey);
1570 }
1571 const E = flags.changedEstimatedTime
1572 ? flags.estimatedTimeNonNull ? await r.varUint() : null
1573 : undefined;
1574 const V = flags.changedFormattedValue
1575 ? await r.stringWithLength()
1576 : undefined;
1577
1578 // a "v2"'s packet will go here, contents ignored by this version.
1579 if (flags.futureProofingPacket) {
1580 const len = await r.varUint();
1581 void await r.readExactly(len);
1582 }
1583
1584 event.push({
1585 k: key as EncodedKey,
1586 t,
1587 v,
1588 e,
1589 s: flags.showTotal,
1590 p: flags.passive,
1591 l,
1592 h: flags.hidden,
1593 c,
1594 E,
1595 V,
1596 });
1597 }
1598
1599 // events
1600 let eventCount = await r.varUint();
1601 while (eventCount-- > 0) {
1602 const name = await r.stringWithLength();
1603 const args = JSON.parse(await r.stringWithLength());
1604 event.push([name, ...args]);
1605 }
1606
1607 return event;
1608}
1609
1610const logLevelSerialize = ["info", "error", "warn", "debug"] as const;
1611
1612function encodeDeltaFlags(node: Partial<DeltaFlags>) {
1613 return (
1614 (node.hidden ? 1 << 0 : 0) +
1615 (node.passive ? 1 << 1 : 0) +
1616 (node.showTotal ? 1 << 2 : 0) +
1617 (node.changedValue ? 1 << 3 : 0) +
1618 (node.changedTotal ? 1 << 4 : 0) +
1619 (node.changedLogs ? 1 << 6 : 0) +
1620 (node.changedText ? 1 << 7 : 0) +
1621 (node.changedChildren ? 1 << 8 : 0) +
1622 (node.changedEstimatedTime ? 1 << 9 : 0) +
1623 (node.changedFormattedValue ? 1 << 10 : 0) +
1624 (node.estimatedTimeNonNull ? 1 << 11 : 0) +
1625 (node.futureProofingPacket ? 1 << 14 : 0) +
1626 (node.deleted ? 1 << 15 : 0)
1627 );
1628}
1629
1630function decodeDeltaFlags(flags: number): DeltaFlags {
1631 return {
1632 hidden: (flags & 1 << 0) !== 0,
1633 passive: (flags & 1 << 1) !== 0,
1634 showTotal: (flags & 1 << 2) !== 0,
1635 changedValue: (flags & 1 << 3) !== 0,
1636 changedTotal: (flags & 1 << 4) !== 0,
1637 changedLogs: (flags & 1 << 6) !== 0,
1638 changedText: (flags & 1 << 7) !== 0,
1639 changedChildren: (flags & 1 << 8) !== 0,
1640 changedEstimatedTime: (flags & 1 << 9) !== 0,
1641 changedFormattedValue: (flags & 1 << 10) !== 0,
1642 estimatedTimeNonNull: (flags & 1 << 11) !== 0,
1643 futureProofingPacket: (flags & 1 << 14) !== 0,
1644 deleted: (flags & 1 << 15) !== 0,
1645 };
1646}
1647
7851648class KeyPool<T extends number> {
786 next = 1 as T;
1649 next = 0;
7871650 old: T[] = [];
7881651
7891652 get() {
7901653 const existing = this.old.shift();
7911654 if (existing) return existing as T;
792 const next = this.next = this.next + 1 as T;
793 return next as T;
1655 return (this.next += 1) as T;
7941656 }
7951657
7961658 recycle(k: T) {
......@@ -851,14 +1713,33 @@ export class Ema implements EstimationAlgorithm {
8511713 }
8521714}
8531715
1716/** according to a stream */
1717type EncodedKey = number & { brand: typeof kNode };
1718
8541719const globalKeyPool = new KeyPool<number>();
855const global =
1720export const global: Ref =
8561721 /** @__PURE__ */ ((root = new Root()) => (attachToScreen(root, log), root))();
8571722
1723/**
1724 * for testing. not covered by semver
1725 * @internal
1726 */
1727const uncastInternals = /** @__PURE__ */ (() => ({
1728 readStreamEvent,
1729 writeStreamEvent,
1730 kNode: kNode as unknown as symbol,
1731 newNode,
1732 EncodedKey: 0 as EncodedKey,
1733}))();
1734export const internals: typeof uncastInternals = uncastInternals;
1735
1736import * as ansi from "./string/ansi.ts";
8581737import * as async from "./async.ts";
1738import * as bytes from "./bytes.ts";
8591739import * as log from "./log.ts";
860import * as ansi from "./string/ansi.ts";
861import * as ts from "./ts.ts";
1740import * as stack from "./log/stack.ts";
1741import * as stream from "./stream.ts";
8621742import * as string from "./string.ts";
1743import * as ts from "./ts.ts";
8631744import { ASSERT, UNWRAP } from "./assert.ts";
8641745import { Events } from "./Events.ts";
lib/stream.test.ts created+232
......@@ -0,0 +1,232 @@
1describe("BufferedWriter", () => {
2 test("basics", () => {
3 const chunks: Uint8Array[] = [];
4 const w = new stream.BufferedWriter((d) => chunks.push(d), 16);
5
6 w.u8(1);
7 w.u16(258);
8 w.u32(0x04030201);
9 w.flush();
10
11 assert.equal(chunks.length, 1);
12 assert.deepEqual(chunks[0], new Uint8Array([1, 2, 1, 1, 2, 3, 4]));
13 });
14
15 test("auto-flush on large write", () => {
16 const chunks: Uint8Array[] = [];
17 const w = new stream.BufferedWriter((d) => chunks.push(d), 8);
18
19 w.u32(1);
20 assert.equal(w.written, 4);
21 const large = new Uint8Array(20).fill(0xff);
22 w.write(large);
23 assert.equal(w.written, 0);
24
25 assert.equal(chunks.length, 2);
26 assert.equal(chunks[0]!.length, 4);
27 assert.equal(chunks[1]!.length, 20);
28 });
29
30 test("transfer based flush", () => {
31 const chunks: Uint8Array[] = [];
32 const w = new stream.BufferedWriter((d) => chunks.push(d), 4);
33
34 w.u8(1);
35 w.u8(2);
36 w.u8(3);
37 w.u8(4);
38 w.flush();
39 w.u8(2);
40 w.u8(3);
41 w.u8(4);
42 w.u8(5);
43 w[Symbol.dispose]();
44
45 assert.equal(chunks.length, 2);
46 assert.deepEqual(chunks[0], new Uint8Array([1, 2, 3, 4]));
47 assert.deepEqual(chunks[1], new Uint8Array([2, 3, 4, 5]));
48 });
49
50 test("ensure more than total capacity", () => {
51 const chunks: Uint8Array[] = [];
52 const w = new stream.BufferedWriter((d) => chunks.push(d), 2);
53 assert.equal(w.size, 2);
54 w.u8(4);
55 w.ensureUnusedCapacity(8); // flushes [ 4 ]
56 w.write(new Uint8Array([4, 5, 5, 8, 8, 1, 2]));
57 assert.equal(w.written, 7);
58 assert.equal(w.size, 8);
59 w.u32(42);
60 w.flush();
61
62 assert.equal(chunks.length, 3);
63 assert.deepEqual(chunks[0], new Uint8Array([4]));
64 assert.deepEqual(chunks[1], new Uint8Array([4, 5, 5, 8, 8, 1, 2]));
65 assert.deepEqual(chunks[2], new Uint8Array([42, 0, 0, 0]));
66 });
67
68 test("varUint", () => {
69 const chunks: Uint8Array[] = [];
70 const w = new stream.BufferedWriter((d) => chunks.push(d));
71
72 w.varUint(0);
73 w.varUint(127);
74 w.varUint(128);
75 w.varUint(16383);
76 w.flush();
77
78 assert.deepEqual(chunks[0], new Uint8Array([0, 127, 128, 1, 255, 127]));
79 });
80
81 test("dispose", () => {
82 const chunks: Uint8Array[] = [];
83 {
84 using w = new stream.BufferedWriter((d) => chunks.push(d));
85 w.u8(42);
86 }
87 assert.equal(chunks.length, 1);
88 assert.equal(chunks[0]![0], 42);
89 });
90});
91
92describe("BufferedReader", () => {
93 test("basic read", async () => {
94 const [s, w] = stream.bufferedReadableStream();
95 w.u8(1);
96 w.u16(258);
97 w.flush();
98 w.u32(0x04030201);
99 w.flush();
100
101 using r = new stream.BufferedReader(s.getReader());
102 assert.equal(await r.u8(), 1);
103 assert.equal(await r.u16(), 258);
104 assert.equal(await r.u16(), 0x0201);
105 assert.equal(await r.u16(), 0x0403);
106 });
107
108 test("varUint", async () => {
109 const [s, w] = stream.bufferedReadableStream();
110 w.varUint(0);
111 w.varUint(127);
112 w.varUint(128);
113 w.varUint(16383);
114 w.varUint(1768210855277);
115 w.varUint(521);
116 w.close();
117
118 using r = new stream.BufferedReader(s.getReader());
119 assert.equal(await r.varUint(), 0);
120 assert.equal(await r.varUint(), 127);
121 assert.equal(await r.varUint(), 128);
122 assert.equal(await r.varUint(), 16383);
123 assert.equal(await r.varUint(), 1768210855277);
124 assert.equal(await r.varUint(), 521);
125 });
126
127 test("cross-chunk read", async () => {
128 const [s, w] = stream.bufferedReadableStream();
129 w.u8(1);
130 w.flush();
131 w.u8(2);
132 w.flush();
133
134 using r = new stream.BufferedReader(s.getReader());
135 assert.equal(await r.u16(), 0x0201);
136 });
137
138 test("readExactly", async () => {
139 const [s, w] = stream.bufferedReadableStream();
140 w.write(new Uint8Array([1, 2, 3, 4, 5]));
141 w.flush();
142
143 using r = new stream.BufferedReader(s.getReader());
144 const buf = await r.readExactly(3);
145 assert.deepEqual(buf, new Uint8Array([1, 2, 3]));
146 assert.equal(await r.u8(), 4);
147 });
148
149 test("readExactly cutting chunks", async () => {
150 {
151 const [s, w] = stream.bufferedReadableStream();
152 w.write(new Uint8Array([1, 2, 3]));
153 w.flush();
154 w.write(new Uint8Array([4, 5, 6]));
155 w.flush();
156 w.write(new Uint8Array([7, 8, 9]));
157 w.flush();
158 w.write(new Uint8Array([10, 11, 12]));
159 w.flush();
160
161 using r = new stream.BufferedReader(s.getReader());
162 assert.deepEqual(await r.readExactly(4), new Uint8Array([1, 2, 3, 4]));
163 assert.deepEqual(await r.readExactly(4), new Uint8Array([5, 6, 7, 8]));
164 assert.deepEqual(await r.readExactly(4), new Uint8Array([9, 10, 11, 12]));
165 }
166 });
167
168 test("readUpTo", async () => {
169 const [s, w] = stream.bufferedReadableStream();
170 w.write(new Uint8Array([1, 2]));
171 w.close();
172
173 using r = new stream.BufferedReader(s.getReader());
174 const buf = await r.readUpTo(10);
175 assert.equal(buf.length, 2);
176 });
177
178 test("stringWithLength", async () => {
179 const [s, w] = stream.bufferedReadableStream();
180 w.stringWithLength("hello");
181 w.flush();
182
183 using r = new stream.BufferedReader(s.getReader());
184 assert.equal(await r.stringWithLength(), "hello");
185 });
186
187 test("some numeric types", async () => {
188 const [s, w] = stream.bufferedReadableStream();
189
190 w.i8(-1);
191 w.i16(-1);
192 w.i32(-1);
193 w.i64(-1n);
194 w.u32(4_000_000_000);
195 w.u64(0xffffffffffffffffn);
196 w.f16(1.5);
197 w.f32(1.5);
198 w.f64(1.5);
199 w.close();
200
201 using r = new stream.BufferedReader(s.getReader());
202 assert.equal(await r.i8(), -1);
203 assert.equal(await r.i16(), -1);
204 assert.equal(await r.i32(), -1);
205 assert.equal(await r.i64(), -1n);
206 assert.equal(await r.u32(), 4_000_000_000);
207 assert.equal(await r.u64(), 0xffffffffffffffffn);
208 assert.equal(await r.f16(), 1.5);
209 assert.equal(await r.f32(), 1.5);
210 assert.equal(await r.f64(), 1.5);
211 });
212
213 test("cancel", () => {
214 let cancelReason: unknown;
215 const s = new ReadableStream({
216 start() {},
217 cancel(reason) {
218 cancelReason = reason;
219 },
220 });
221
222 const r = new stream.BufferedReader(s.getReader());
223 const err = new Error("test cancel");
224 r.cancel(err);
225
226 assert.equal(cancelReason, err);
227 });
228});
229
230import { describe, test } from "node:test";
231import { strict as assert } from "node:assert";
232import * as stream from "./stream.ts";
lib/stream.ts created+485
......@@ -0,0 +1,485 @@
1/**
2 * this @module contains helpers for creating and consuming `ReadableStream`,
3 * particularly with binary payloads.
4 */
5
6const shared = new Uint8Array(8);
7const sharedView = new DataView(shared.buffer);
8
9/** write bytes `data` to a destination */
10export type WriteBytesFn = (data: Uint8Array) => void;
11
12/** something that can be written to */
13export type WriteTarget =
14 | WriteBytesFn
15 | { enqueue(chunk: Uint8Array): void; close?(): void };
16
17/**
18 * `ReadableStreamDefaultController.enqueue` does not buffer data in any way.
19 * for use cases such as a streaming write of many small binary chunks, such as
20 * differently sized integers in the case of `progress.encodeByteStream`, this
21 * batches those many small writes into their own `ArrayBuffer` objects.
22 *
23 * when passing a `ReadableStreamDefaultController`, the `BufferedWriter`
24 * "owns" it, and will close the stream when disposed. to prevent that, pass
25 * a `WriteBytesFn` instead.
26 *
27 * ```ts
28 * const w = new BufferedWriter(unbufferedWrite);
29 * w.write(string.encodeUtf8("a null-terminated string"));
30 * w.u8(0);
31 * w.write(string.encodeUtf8("a second null-terminated string"));
32 * w.u8(0);
33 *
34 * w.flush(); // don't forget to flush!
35 * // `unbufferedWrite` will be called once with 57 bytes
36 * ```
37 */
38export class BufferedWriter {
39 #view: DataView<ArrayBuffer>;
40 #written: number = 0;
41 #flush: WriteBytesFn;
42 #size: number;
43 #onClose?: (() => void) | null;
44
45 constructor(
46 flush: WriteTarget,
47 size: number = 8192,
48 ) {
49 this.#flush = typeof flush === "function"
50 ? flush
51 : flush.enqueue.bind(flush);
52 this.#view = new DataView(new ArrayBuffer(size));
53 this.#size = size;
54 this.#onClose = typeof flush === "function"
55 ? null
56 : flush.close?.bind(flush);
57 }
58
59 flush() {
60 const written = this.#written;
61 if (written === 0) return;
62 const size = this.#size;
63 if (written === size) {
64 // transfer
65 this.#flush(new Uint8Array(this.#view.buffer));
66 } else {
67 // copy
68 this.#flush(new Uint8Array(this.#view.buffer.slice(0, written)));
69 }
70 this.#view = new DataView(new ArrayBuffer(size));
71 this.#written = 0;
72 }
73
74 close() {
75 const written = this.#written;
76 if (written === 0) return;
77 const size = this.#size;
78 if (written === size) {
79 // transfer
80 this.#flush(new Uint8Array(this.#view.buffer));
81 } else {
82 // copy
83 this.#flush(new Uint8Array(this.#view.buffer.slice(0, written)));
84 }
85 this.#written = 0;
86 this.#size = 0;
87 this.#view = new DataView(new ArrayBuffer(0));
88 this.#onClose?.();
89 }
90
91 [Symbol.dispose]() {
92 this.close();
93 }
94
95 write(bytes: Uint8Array<ArrayBuffer> | ArrayBuffer) {
96 const len = bytes.byteLength;
97 if (len >= this.size) {
98 this.flush();
99 this.#flush("buffer" in bytes ? bytes : new Uint8Array(bytes));
100 return;
101 }
102 this.ensureUnusedCapacity(len);
103 new Uint8Array(this.#view.buffer, this.#written)
104 .set(bytes instanceof ArrayBuffer ? new Uint8Array(bytes) : bytes);
105 this.#written += bytes.byteLength;
106 }
107
108 varUint(int: number) {
109 ASSERT(int >= 0);
110 do {
111 const data = int % (1 << 7);
112 int = Math.floor(int / (1 << 7));
113 const continuation = int > 0 ? (1 << 7) : 0;
114 this.u8(data + continuation);
115 } while (int > 0);
116 }
117
118 u8(int: number) {
119 const size = 1;
120 this.ensureUnusedCapacity(size);
121 this.#view.setUint8(this.#written, int);
122 this.#written += size;
123 }
124
125 u16(int: number) {
126 const size = 2;
127 this.ensureUnusedCapacity(size);
128 this.#view.setUint16(this.#written, int, true);
129 this.#written += size;
130 }
131
132 u32(int: number) {
133 const size = 4;
134 this.ensureUnusedCapacity(size);
135 this.#view.setUint32(this.#written, int, true);
136 this.#written += size;
137 }
138
139 u64(int: bigint) {
140 const size = 8;
141 this.ensureUnusedCapacity(size);
142 this.#view.setBigUint64(this.#written, int, true);
143 this.#written += size;
144 }
145
146 i8(int: number) {
147 const size = 1;
148 this.ensureUnusedCapacity(size);
149 this.#view.setUint8(this.#written, int);
150 this.#written += size;
151 }
152
153 i16(int: number) {
154 const size = 2;
155 this.ensureUnusedCapacity(size);
156 this.#view.setInt16(this.#written, int, true);
157 this.#written += size;
158 }
159
160 i32(int: number) {
161 const size = 4;
162 this.ensureUnusedCapacity(size);
163 this.#view.setInt32(this.#written, int, true);
164 this.#written += size;
165 }
166
167 i64(int: bigint) {
168 const size = 8;
169 this.ensureUnusedCapacity(size);
170 this.#view.setBigInt64(this.#written, int, true);
171 this.#written += size;
172 }
173
174 f16(float: number) {
175 const size = 2;
176 this.ensureUnusedCapacity(size);
177 this.#view.setFloat16(this.#written, float, true);
178 this.#written += size;
179 }
180
181 f32(float: number) {
182 const size = 4;
183 this.ensureUnusedCapacity(size);
184 this.#view.setFloat32(this.#written, float, true);
185 this.#written += size;
186 }
187
188 f64(float: number) {
189 const size = 8;
190 this.ensureUnusedCapacity(size);
191 this.#view.setFloat64(this.#written, float, true);
192 this.#written += size;
193 }
194
195 stringWithLength(text: string) {
196 const buf = string.encodeUtf8(text);
197 this.varUint(buf.byteLength);
198 this.write(buf);
199 }
200
201 get written(): number {
202 return this.#written;
203 }
204 get size(): number {
205 return this.#size;
206 }
207
208 ensureUnusedCapacity(extra: number) {
209 const size = this.size;
210 const written = this.#written;
211 if (written + extra > size) {
212 if (extra > size) {
213 this.resize(extra);
214 } else {
215 this.flush();
216 }
217 }
218 }
219
220 resize(newCapacity: number) {
221 newCapacity = Math.round(newCapacity);
222 const written = this.#written;
223 if (written > 0) {
224 // transfer
225 this.#flush(new Uint8Array(this.#view.buffer, 0, written));
226 this.#written = 0;
227 }
228 this.#size = newCapacity;
229 this.#view = new DataView(new ArrayBuffer(newCapacity));
230 }
231}
232
233/**
234 * a reader for `ReadableStream` which allows the caller to specify the
235 * chunking size and read numbers at a time.
236 */
237export class BufferedReader implements Disposable {
238 #reader: ReadableStreamDefaultReader<Uint8Array> | null;
239 #offset = 0;
240 #buffer: Uint8Array[] = [];
241 #available = 0;
242
243 constructor(reader: ReadableStreamDefaultReader<Uint8Array>) {
244 this.#reader = reader;
245 }
246
247 get available() {
248 return this.#available;
249 }
250
251 /** Releasing will discard the bytes remaining in the buffer. */
252 releaseLock() {
253 UNWRAP(this.#reader, "reader disposed").releaseLock();
254 this.#reader = null;
255 }
256
257 cancel(error: unknown) {
258 UNWRAP(this.#reader, "reader disposed").cancel(error);
259 }
260
261 async ensureAvailableOrFalse(size: number): Promise<boolean> {
262 if (size <= 0 || this.#available >= size) return true;
263 const reader = UNWRAP(this.#reader, "reader disposed");
264 const buffer = this.#buffer;
265 while (this.#available < size) {
266 const { value, done } = await reader.read();
267 if (done) break;
268 buffer.push(value);
269 this.#available += value.byteLength;
270 }
271 return this.#available >= size;
272 }
273 async ensureAvailable(size: number): Promise<void> {
274 ASSERT(await this.ensureAvailableOrFalse(size), "End of stream");
275 }
276
277 #collectLimitedView(length: number): DataView {
278 ASSERT(this.available >= length);
279 const buffer = this.#buffer;
280 let chunk = UNWRAP(buffer[0]);
281 let offset = this.#offset;
282 if (chunk.byteLength - offset >= length) {
283 this.#shiftOffset(length);
284 return new DataView(chunk.buffer, chunk.byteOffset + offset, length);
285 }
286 let x = 0;
287 while (true) {
288 for (
289 ;
290 offset < chunk.byteLength && x < length;
291 offset += 1, x += 1
292 ) shared[x] = chunk[offset]!;
293 if (x === length) {
294 if (offset === chunk.byteLength) {
295 ASSERT(buffer.shift());
296 this.#offset = 0;
297 } else this.#offset = offset;
298 this.#available -= length;
299 return sharedView;
300 }
301 ASSERT(buffer.shift());
302 chunk = UNWRAP(buffer[0]);
303 offset = 0;
304 }
305 }
306 #collectBuffer(length: number) {
307 ASSERT(this.available >= length);
308 this.#available -= length;
309 const buffer = this.#buffer;
310 let offset = this.#offset;
311
312 // slice an existing data view
313 if (buffer[0] && buffer[0].byteLength - offset >= length) {
314 const first = buffer[0];
315 if (first.byteLength - offset === length) {
316 this.#offset = 0;
317 ASSERT(buffer.shift());
318 return offset ? first.subarray(offset) : first;
319 }
320 this.#offset += length;
321 return first.subarray(offset, offset + length);
322 }
323
324 // join by allocating a new buffer
325 const joined = new Uint8Array(length);
326 let i = 0;
327 do {
328 const chunk = UNWRAP(buffer[0]);
329 const { byteLength } = chunk;
330 const copyLen = Math.min(byteLength, length);
331 joined.set(
332 copyLen === byteLength
333 ? chunk
334 : chunk.subarray(offset, offset + copyLen),
335 i,
336 );
337 i += copyLen;
338 length -= copyLen;
339 if (copyLen === byteLength) ASSERT(buffer.shift());
340 if (length === 0) this.#offset = copyLen;
341 } while (length > 0);
342 return joined;
343 }
344 #shiftOffset(offset: number) {
345 const newOffset = this.#offset += offset;
346 this.#available -= offset;
347 if (UNWRAP(this.#buffer[0]).byteLength === newOffset) {
348 this.#buffer.shift();
349 this.#offset = 0;
350 } else {
351 this.#offset = newOffset;
352 }
353 }
354
355 async u8() {
356 const int = await this.peekU8();
357 this.#shiftOffset(1);
358 return int;
359 }
360
361 async peekU8() {
362 await this.ensureAvailable(1);
363 return UNWRAP(UNWRAP(this.#buffer[0])[this.#offset]);
364 }
365
366 async i8() {
367 await this.ensureAvailable(2);
368 return this.#collectLimitedView(1).getInt8(0);
369 }
370
371 async u16() {
372 await this.ensureAvailable(2);
373 return this.#collectLimitedView(2).getUint16(0, true);
374 }
375
376 async i16() {
377 await this.ensureAvailable(2);
378 return this.#collectLimitedView(2).getInt16(0, true);
379 }
380
381 async u32() {
382 await this.ensureAvailable(4);
383 return this.#collectLimitedView(4).getUint32(0, true);
384 }
385
386 async i32() {
387 await this.ensureAvailable(4);
388 return this.#collectLimitedView(4).getInt32(0, true);
389 }
390
391 async u64() {
392 await this.ensureAvailable(8);
393 return this.#collectLimitedView(8).getBigUint64(0, true);
394 }
395
396 async i64() {
397 await this.ensureAvailable(8);
398 return this.#collectLimitedView(8).getBigInt64(0, true);
399 }
400
401 async f16() {
402 await this.ensureAvailable(2);
403 return this.#collectLimitedView(2).getFloat16(0, true);
404 }
405
406 async f32() {
407 await this.ensureAvailable(4);
408 return this.#collectLimitedView(4).getFloat32(0, true);
409 }
410
411 async f64() {
412 await this.ensureAvailable(8);
413 return this.#collectLimitedView(8).getFloat64(0, true);
414 }
415
416 async varUint() {
417 let result = 0;
418 let shift = 1;
419 let byte;
420 do {
421 byte = await this.u8();
422 result += (byte & 0x7F) * shift;
423 shift *= 128;
424 } while (byte & 0x80);
425 return result;
426 }
427
428 async stringWithLength() {
429 const len = await this.varUint();
430 return string.decodeUtf8(await this.readExactly(len));
431 }
432
433 /** Read up to a maximum number of bytes */
434 async readUpTo(size: number): Promise<Uint8Array> {
435 await this.ensureAvailableOrFalse(size);
436 return this.#collectBuffer(Math.min(size, this.#available));
437 }
438 async readExactly(size: number): Promise<Uint8Array> {
439 await this.ensureAvailable(size);
440 return this.#collectBuffer(size);
441 }
442
443 [Symbol.dispose]() {
444 return this.releaseLock();
445 }
446}
447
448/** easily create a readable stream with a writer */
449export function bufferedReadableStream(): [
450 ReadableStream,
451 BufferedWriter,
452 ReadableStreamDefaultController<Uint8Array>,
453] {
454 let unbuffered: ReadableStreamDefaultController<Uint8Array> | null = null;
455 const reader = new ReadableStream<Uint8Array>({
456 start: (controller) => unbuffered = controller,
457 });
458 ASSERT(unbuffered);
459 return [reader, new BufferedWriter(unbuffered), unbuffered];
460}
461
462/** easily create a readable stream with a writer */
463export function readWritePair(): [
464 ReadableStream,
465 ReadableStreamDefaultController,
466] {
467 let unbuffered: ReadableStreamDefaultController<Uint8Array> | null = null;
468 const reader = new ReadableStream<Uint8Array>({
469 start: (controller) => unbuffered = controller,
470 });
471 ASSERT(unbuffered);
472 return [reader, unbuffered];
473}
474
475export function fromChunks<T>(chunks: T[]) {
476 return new ReadableStream<T>({
477 start(c) {
478 for (const chunk of chunks) c.enqueue(chunk);
479 c.close();
480 },
481 });
482}
483
484import * as string from "./string.ts";
485import { ASSERT, UNWRAP } from "./assert.ts";
lib/ts.ts+56
......@@ -1,5 +1,13 @@
11export type Timer = ReturnType<typeof setTimeout>;
22export type Interval = ReturnType<typeof setInterval>;
3/** opposite of the built-in `Readonly` type */
4export type Writeable<T> = { -readonly [P in keyof T]: T[P] };
5/** opposite of the built-in `Readonly` type, but recursive */
6// TODO: tuples?
7export type DeepWriteable<T> = {
8 -readonly [P in keyof T]: T[P] extends ReadonlyArray<infer I> ? Array<I>
9 : DeepWriteable<T[P]>;
10};
311
412/**
513 * redeclared here because it is only provided by `lib: ["dom"]` and not
......@@ -23,3 +31,51 @@ export function defer(fn: () => void): Dispose {
2331 f[Symbol.dispose] = f;
2432 return f;
2533}
34
35export type Mixin<A, B> = A & B;
36
37/**
38 * mutates `base` to implement all methods of `other`, using `other` as state.
39 */
40export function mixin<A extends object, B extends object>(
41 base: A,
42 other: B,
43): Mixin<A, B> {
44 const prototype = Object.getPrototypeOf(other);
45 for (
46 const key of [
47 ...Object.getOwnPropertyNames(prototype),
48 ...Object.getOwnPropertySymbols(prototype),
49 ]
50 ) {
51 if (key === "constructor") continue;
52 if (key in base) continue;
53 const fn = prototype[key];
54 if (typeof fn !== "function") continue;
55 (base as Record<string | symbol, unknown>)[key] = fn.bind(other);
56 }
57 return base as A & B;
58}
59
60export type EmptyObject = Record<never, unknown>;
61
62/** TODO: this type has many subtle bugs */
63export type ToJson<T> = T extends JsonValue ? T
64 : T extends { toJSON(): infer J } ? ToJson<J>
65 : T extends Set<unknown> | Map<unknown, unknown> | Record<string, never>
66 ? EmptyObject
67 : T extends unknown[] ? {
68 [K in keyof Omit<T, keyof unknown[]>]: T[K] extends undefined | void
69 ? null
70 : ToJson<T[K]>;
71 }
72 : T extends object ? { [K in keyof Omit<T, symbol>]: ToJson<T[K]> }
73 : never;
74export type Json = JsonValue | undefined | void;
75export type JsonValue =
76 | string
77 | number
78 | boolean
79 | null
80 | JsonValue[]
81 | { [k: string]: Exclude<JsonValue, undefined> };