diff --git a/lib/bytes.ts b/lib/bytes.ts new file mode 100644 index 0000000000000000000000000000000000000000..03f08dbdaea64fea236d70b972fe5d3c6548b2e9 --- /dev/null +++ b/lib/bytes.ts @@ -0,0 +1,9 @@ +export function eql( + a: T, + b: T, +) { + const l = a.byteLength; + if (l !== b.byteLength) return false; + for (let i = 0; i < l; i++) if (a[i] !== b[i]) return false; + return true; +} diff --git a/lib/log.ts b/lib/log.ts index 0f7e296f987ece8f6e80040931ef9aa8eb44cb26..638fff8c7afd51f8c9688e004f6553b3c1206d5c 100644 --- a/lib/log.ts +++ b/lib/log.ts @@ -1,5 +1,5 @@ /** - * by using `lib/log.ts`, your application gets easy scoped logging as well as + * by using `lib/log.ts`, an application gets easy scoped logging as well as * integration with terminal widgets such as `lib/progress.ts`. * * the pattern for using this module is to shadow the global `console` with a @@ -53,23 +53,19 @@ export interface Scope { error(...args: unknown[]): void; /** - * emit a debugging message. */ + * emit a debugging message. + */ log(...args: unknown[]): void; /** * emit a debugging message. this method is in place for compatibility with * the `console` API. prefer calling `.log` directly. - * @internalconst hello = Math.random() > 0.5 ? 1 : "meow"; - -interface Meow { - meow: bigint; -} - -const world = hello as Meow; - -console.log(world); + * @internal */ debug(...args: unknown[]): void; + /** write a Message object directly. */ + writeMessage(message: Message): void; + /** create a nested sub-scope */ scoped(name: string): Scope; /** redirect the logging output of this scope somewhere else */ @@ -88,7 +84,7 @@ export interface Message { /** captured stack. */ stack?: stack.Frame[]; /** arbitrary data from the logging source. */ - custom?: Partial>; + custom?: Partial>; /** * original logging arguments, if present. this field is indexed by a symbol * so that it is lost during JSON serialization, as callers are allowed to log @@ -177,6 +173,11 @@ export function writeLine(text: string) { globalWidgetHost.writeLine(text); } +/** write a Message object directly. */ +export function writeMessage(m: Message) { + globalLog.writeMessage(m); +} + /** * while locked, no widgets will draw. prefer `writeLine`. * this lock is not exclusive. @@ -523,7 +524,7 @@ const Scope = class Scope implements Scope { frames = stack.capture().slice(2); withinStackCapture = false; } - const m = { + this.writeMessage({ level, scope: this.name, get text() { @@ -534,11 +535,7 @@ const Scope = class Scope implements Scope { time: Date.now(), stack: frames, [originalLogArgs]: args, - }; - if (withinDispatch) return void globalOutputFunction(m); - withinDispatch = true; - this.#dispatch(m); - withinDispatch = false; + }); } info: (...args: unknown[]) => void = (...args: unknown[]) => { @@ -557,6 +554,13 @@ const Scope = class Scope implements Scope { this.#log("debug", args); }; + writeMessage: (m: Message) => void = (m) => { + if (withinDispatch) return void globalOutputFunction(m); + withinDispatch = true; + this.#dispatch(m); + withinDispatch = false; + }; + scoped(name: string): Scope { const current = this.name; return new Scope( diff --git a/lib/progress.test.ts b/lib/progress.test.ts index 2b9a3e34ea1e532c757f22debd88f3492d575371..b63cc775737e06fabb3c328321dbc040ab033766 100644 --- a/lib/progress.test.ts +++ b/lib/progress.test.ts @@ -1,5 +1,5 @@ // a trivial example of how to use progress -test("trivial end-to-end example", () => { +test.skip("trivial end-to-end example", () => { const { fgBlue: FB, fgReset: FR, reset: R } = ansi; const screen = new testing.MockScreen(); // create a mock terminal const root = new progress.Root(screen); // sync with mock timers @@ -43,6 +43,39 @@ test("trivial end-to-end example", () => { }); }); +describe("encodeEventStream", async (t) => { + await test("simple", async () => { + const root = new progress.Root(); // sync with mock timers + + const a = root.start("hello"); + const b = root.start("cats", { value: 2, total: 5 }); + + const events = progress.encodeEventStream(root); + const reader = events.getReader(); + const next = async () => + cleanProgressEvent(UNWRAP((await reader.read()).value)); + + assert.deepEqual(await next(), [ + -1, + { k: 1, t: "hello" }, + { k: 2, t: "cats", v: 2, e: 5 }, + [1, 2], + ]); + b.inc(); + b.inc(); + assert.deepEqual(await next(), [ + 0, + { k: 2, v: 4 }, + ]); + b.inc(); + assert.deepEqual(await next(), [ + 1, + 2, + [1], + ]); + }); +}); + // ## value formatters test("slashValueFormatter", ({ mock }) => { const fn = mock.fn((arg: number) => `<${arg}>`); @@ -55,11 +88,13 @@ test("slashValueFormatter", ({ mock }) => { }); test("bytesValueFormatter", () => { assert.equal(progress.bytesValueFormatter(250, 0), "250B"); - assert.equal(progress.bytesValueFormatter(0, 1_200_000), "0B/1.2MB"); + assert.equal(progress.bytesValueFormatter(50_000, 1_200_000), "50kB/1.20MB"); + assert.equal(progress.bytesValueFormatter(0, 0), ""); }); test("defaultValueFormatter", () => { assert.equal(progress.defaultValueFormatter(250, 0), "250"); - assert.equal(progress.defaultValueFormatter(0, 1_200_000), "0/1200000"); + assert.equal(progress.defaultValueFormatter(52, 1_200_000), "52/1200000"); + assert.equal(progress.defaultValueFormatter(0, 1_200_000), ""); }); test("percentValueFormatter", () => { @@ -95,8 +130,115 @@ test("Ema", () => { assert.equal(ema.sample(50_130, 0.9), 50_141.754246587305); }); +// ## internals +test("event encoding round trip cases", async () => { + const fakeNode = progress.nullNode; + async function roundTrip(event: progress.StreamEvent) { + const [readable, controller] = stream.readWritePair(); + const chunks: Uint8Array[] = []; + progress.internals.writeStreamEvent(event.slice(), (d) => chunks.push(d)); + const combined = Buffer.concat(chunks); + // console.debug(JSON.stringify(combined.toString("utf-8"))); + controller.enqueue(combined); + controller.close(); + using reader = new stream.BufferedReader(readable.getReader()); + const decoded = await progress.internals.readStreamEvent(reader); + assert.deepEqual(cleanProgressEvent(decoded), cleanProgressEvent(event)); + } + + await roundTrip([ + -1, + { + k: 1 as EncodedKey, + t: "job", + [progress.internals.kNode]: fakeNode, + }, + [1], + ]); + await roundTrip([ + 0, + { + k: 1 as EncodedKey, + t: "job", + [progress.internals.kNode]: fakeNode, + }, + { + k: 2 as EncodedKey, + t: "deeply nested", + c: [1 as EncodedKey], + [progress.internals.kNode]: fakeNode, + }, + { + k: 3 as EncodedKey, + t: "other task", + c: [2 as EncodedKey], + [progress.internals.kNode]: fakeNode, + }, + { + k: 4 as EncodedKey, + v: 4, + e: 7, + t: "purring", + E: 1768210855277.92, + [progress.internals.kNode]: fakeNode, + }, + { + k: 5 as EncodedKey, + v: 2, + e: 7, + t: "meowing", + E: 1768210857407, + V: "28.6%", + [progress.internals.kNode]: fakeNode, + }, + { + k: 6 as EncodedKey, + t: "subtask 1t07l", + [progress.internals.kNode]: fakeNode, + }, + { + k: 7 as EncodedKey, + t: "subtask 42k1d", + c: [3 as EncodedKey], + [progress.internals.kNode]: fakeNode, + }, + { + k: 8 as EncodedKey, + t: "subtask 4kdjg", + c: [5 as EncodedKey, 4 as EncodedKey], + [progress.internals.kNode]: fakeNode, + }, + { + k: 9 as EncodedKey, + e: 30, + t: "progress node", + c: [8, 7, 6] as EncodedKey[], + [progress.internals.kNode]: fakeNode, + }, + [9], + ]); +}); + +function cleanProgressEvent(event: progress.StreamEvent) { + return event.map((x) => + typeof x === "object" && !Array.isArray(x) + ? testing.removeUndefinedKeys({ + ...x, + E: undefined, + p: undefined, + s: undefined, + h: undefined, + }) + : x + ); +} + +type EncodedKey = typeof progress.internals.EncodedKey; + import * as progress from "./progress.ts"; -import { test } from "node:test"; +import { describe, test } from "node:test"; import assert from "node:assert/strict"; +import * as stream from "./stream.ts"; import * as testing from "./testing.ts"; import * as ansi from "./string/ansi.ts"; +import { UNWRAP } from "./assert.ts"; diff --git a/lib/progress.ts b/lib/progress.ts index bca5009a630021e7a9ceb6f85a95401f55ab14ba..edc42afd4a654c3b5d7fc5707deb4f7c231e5449 100644 --- a/lib/progress.ts +++ b/lib/progress.ts @@ -48,6 +48,7 @@ * that rests at the end of the log. if `estimate` is given, a progress bar is * shown. */ +/* node:coverage ignore next 3 */ export function start(text: string, opts?: StartOptions): Node { return global.start(text, opts); } @@ -267,6 +268,7 @@ export class Root< * when given an estimate, a progress bar is rendered. */ start(text: string, opts?: StartOptions): Node { + ASSERT(this.#status === "active"); const [state, node] = newNode(this, text, opts); this.#active.push(state); this.emit("node-start", state, null); @@ -452,10 +454,11 @@ function newNode( key === "value" && estimator && state.total > 0 && state.value > 0 && state.value < state.total ) { - const time = estimator.sample(owner.now(), state.value / state.total); - if ( - time && (!state.estimatedTime || state.estimatedTime !== time) - ) { + // originally, auto-estimation would use the Root's timing primitives, + // but this makes it so that `state.estimatedTime` is in terms of the + // Root, which is unintuitive for everyone but the renderer. + const time = estimator.sample(Date.now(), state.value / state.total); + if (time && (!state.estimatedTime || state.estimatedTime !== time)) { state.estimatedTime = time; owner.emit("node-change", state, "estimatedTime"); } @@ -590,6 +593,7 @@ export const nullNode: Node = { error: noop, log: noop, debug: noop, + writeMessage: noop, scoped() { return this; }, @@ -782,15 +786,873 @@ export function attachToScreen( return ts.defer(() => stack[Symbol.dispose]); } +const header = /* @__PURE__ */ string.encodeUtf8("clover's progress <3\n"); + +/** the 'Content-Type' of {@linkcode encodeByteStream}'s return value */ +export const contentType = "application/x-clover-progress"; + +/** + * opaque but JSON-compatible serialized state, either containing the entire + * progress tree or just a delta update. this data structure always be future + * and backwards compatbile with changes to this library. if a breaking change + * is made, the APIs in this file will be renamed. + * + * for details on the format, read the source of {@linkcode encodeEventStream} or + * {@linkcode encodeByteStream} + */ +export type StreamEvent = Array< + number | StreamNode | StreamCustomEvent | StreamRootChildren +>; +const kNode = Symbol("underlyingNode"); +/** see {@linkcode StreamEvent} @internal */ +export interface StreamNode { + k: EncodedKey; + /** text, overwrite */ + t?: string | undefined; + /** value, overwrite */ + v?: number | undefined; + /** total, overwrite */ + e?: number | undefined; + /** showTotal, overwrite */ + s?: boolean | undefined; + /** passive, overwrite */ + p?: boolean | undefined; + /** logs, append */ + l?: log.Message[] | undefined; + /** hidden, overwrite */ + h?: boolean | undefined; + /** children, overwrite */ + c?: EncodedKey[] | undefined; + /** estimated time, in unix milliseconds, overwrite */ + E?: number | null | undefined; + /** formatted value if not defaultValueFormatter, overwrite */ + V?: string | undefined; + + /** access to the underlying node is used by {@linkcode encodeByteStream} */ + [kNode]?: ReadOnlyNode; +} +type StreamCustomEvent = [ + /** custom event name or "end" */ + e: "end" | string, + /** custom event payload */ + ...args: ts.Json[], +]; +type StreamRootChildren = number[]; + +export interface EncodeStreamOptions { + /** + * set the minimum time between packets. + * @default 1000 / 30 (30fps) + */ + throttleMs?: number; + serializeStackTraces?: boolean; +} + +/** + * converts a {@linkcode Root|progress.Root} into an object stream + * for communicating progress over the process or network boundary. + * + * if the given root adds custom event handlers, they must all have + * json-serializable payloads. + * + * this function does not use recursion. + */ +export function encodeEventStream< + Result extends ts.Json, + Map extends { [key: string]: ts.Json[] }, +>( + // fun fact, this type is impossible to write without generics + root: Root, + { + throttleMs = 1000 / 30, /* 30fps */ + serializeStackTraces = false, + }: EncodeStreamOptions = {}, +): ReadableStream { + const { delay, now } = root; + const stack = new DisposableStack(); + const s = new Encoder(root, serializeStackTraces); + const ready = new async.Watch(false); + + let lastEvent = now(); + let timer: async.Cancelable | null = null; + function emitSoon() { + if (timer || ready.value) return; + const remaining = now() - lastEvent + throttleMs; + if (remaining > 0) { + (timer = delay(remaining)).then(() => { + timer = null; + ready.value = true; + lastEvent = now(); + }); + } else { + ready.value = true; + lastEvent = now(); + } + } + stack.defer(() => { + if (timer) timer.cancel(); + s.changed.clear(); + s.deleted.clear(); + }); + + let isFirst = true; + return new ReadableStream({ + start(controller) { + if (root.active.length > 0) { + controller.enqueue(s.getDelta(isFirst)); + isFirst = false; + } + + // TODO: why does ts require two args + stack.use(root.on("node-start", (newNode, _) => { + if (!newNode.parent) s.rootChildrenUpdated = true; + s.changed.set(newNode, new Set()); + emitSoon(); + })); + stack.use(root.on("node-change", (node, key) => { + let set = s.changed.get(node); + if (!set) s.changed.set(node, set = new Set()); + set.add(key); + emitSoon(); + })); + stack.use(root.on("node-end", (node) => { + if (!node.parent) s.rootChildrenUpdated = true; + s.changed.delete(node); + s.deleted.add(node); // handle add and remove in same frame + emitSoon(); + })); + }, + async pull(controller) { + await ready.until((x) => x === true); + ready.value = false; + if (timer) { + timer.cancel(); + timer = null; + } + controller.enqueue(s.getDelta(isFirst)); + isFirst = false; + }, + cancel() { + stack.dispose(); + (0, performance.now)(); + }, + }, { highWaterMark: 0 }); +} + +/** + * decodes a stream created by `encodeEventStream` into managed calls to + * `target.start`. when cancelling, all the created nodes are destroyed. + */ +export function decodeEventStream< + Result = void, + EventMap extends Events.Map = ts.EmptyObject, +>( + encoded: ReadableStream, + target: Ref, +): async.Cancelable & Events { + const { resolve, reject, promise } = Promise.withResolvers(); + const events = new Events(); + const ctrl = new AbortController(); + const decoder = new Decoder(target, events, resolve); + const { signal } = ctrl; + (async () => { + if (signal.aborted) return; + const reader = encoded.getReader(); + try { + let hasEmittedEnd = false; + while (true) { + signal.throwIfAborted(); + const { value, done } = await reader.read(); + signal.throwIfAborted(); + if (done) break; + decoder.processEvent(value); + } + ASSERT(hasEmittedEnd, "Stream terminated early."); + } catch (err) { + reader.cancel(err); + if (!signal.aborted) reject(err); + } finally { + reader.releaseLock(); + } + })(); + return Object.assign(ts.mixin(promise, events), { + cancel: (e?: unknown) => ctrl.abort(e), + [Symbol.dispose]: () => ctrl.abort(), + }); +} + +/** + * by listening to the event interface in {@linkcode Root}, compute delta + * events for a progress stream. the caller is expected to fill `changed` and `deleted` + */ +class Encoder< + Result extends ts.Json, + EventMap extends { [key: string]: ts.Json[] }, +> { + root: Root; + pool: KeyPool = new KeyPool(); + keyMap: Map = new Map(); + logsLength: WeakMap = new WeakMap(); + serializeStackTraces: boolean; + + rootChildrenUpdated = true; + changed = new Map>(); + deleted = new Set(); + + constructor( + root: Root, + serializeStackTraces: boolean = false, + ) { + this.root = root; + this.serializeStackTraces = serializeStackTraces; + for (const node of root.active) { + this.changed.set(node, new Set()); + } + } + + getDelta(isFirst = false) { + const deleted = Array.from( + this.deleted, + (node) => [node, this.keyMap.get(node)] as const, + ) + .filter(([node, dst]) => { + if (dst == null) return false; // started and stopped in same tick + this.pool.recycle(dst); + this.keyMap.delete(node); + return true; + }) + .map((x) => UNWRAP(x[1])); + ASSERT(!isFirst || deleted.length === 0); + + const { updated = [], started = [] } = Object.groupBy( + this.changed, + ([node]) => this.keyMap.has(node) ? "updated" : "started", + ); + + const payload: StreamEvent = [ + isFirst ? -1 : deleted.length, + ...deleted, + ]; + + if (started.length > 0) { + this.beginNodesDepthFirst( + started.reverse().map(([node]) => node), + payload, + ); + } + + for (const [node, props] of updated) { + let logs = undefined; + if ( + props.has("logs") && + node.logs.length > (this.logsLength.get(node) ?? 0) + ) { + logs = this.redactLogsIfNeeded( + node.logs.slice(this.logsLength.get(node) ?? 0), + ); + this.logsLength.set(node, node.logs.length); + } + payload.push({ + k: UNWRAP(this.keyMap.get(node)), + t: props.has("text") ? node.text : undefined, + v: props.has("value") ? node.value : undefined, + e: props.has("total") ? node.total : undefined, + s: props.has("showTotal") ? node.showTotal : undefined, + p: props.has("passive") ? node.passive : undefined, + l: logs, + h: props.has("hidden") ? node.hidden : undefined, + c: props.has("children") + ? node.children.map((node) => UNWRAP(this.keyMap.get(node))) + : undefined, + E: props.has("estimatedTime") ? node.estimatedTime : undefined, + V: node.valueFormatter !== defaultValueFormatter && + (props.has("value") || props.has("total") || + props.has("valueFormatter")) + ? node.valueFormatter(node.value, node.total) + : undefined, + [kNode]: node, + }); + } + + if (this.rootChildrenUpdated) { + payload.push( + this.root.active.map((node) => UNWRAP(this.keyMap.get(node))), + ); + } + + this.deleted.clear(); + this.changed.clear(); + this.rootChildrenUpdated = false; + + return payload; + } + + beginNodesDepthFirst( + nodes: ReadOnlyNode[], + out: StreamEvent, + ) { + for (let i = 0; i < nodes.length; i += 1) { + const node = UNWRAP(nodes[i]); + nodes.push(...node.children); + } + let next; + while (next = nodes.pop()) { + if (!this.keyMap.has(next)) { + out.push(this.beginNode(next)); + } + } + } + + beginNode(node: ReadOnlyNode): StreamNode { + const k = this.pool.get(); + this.keyMap.set(node, k); + + return { + k, + v: node.value !== 0 ? node.value : undefined, + e: node.total !== 0 ? node.total : undefined, + t: node.text.length > 0 ? node.text : undefined, + l: node.logs.length > 0 + ? (this.logsLength.set(node, node.logs.length), + this.redactLogsIfNeeded(node.logs.slice())) + : undefined, + c: node.children.length > 0 + ? node.children.map((node) => + UNWRAP(this.keyMap.get(node), () => [node.children, this.keyMap]) + ) + : undefined, + p: node.passive === true ? true : undefined, + s: node.showTotal === false ? false : undefined, + h: node.hidden === true ? true : undefined, + E: node.estimatedTime != null ? node.estimatedTime : undefined, + V: node.valueFormatter !== defaultValueFormatter + ? node.valueFormatter(node.value, node.total) + : undefined, + [kNode]: node, + }; + } + + redactLogsIfNeeded(messages: log.Message[]) { + if (this.serializeStackTraces) return messages; + return messages.map((x) => ({ ...x, stack: undefined })); + } +} + +class Decoder< + Result = void, + EventMap extends Events.Map = ts.EmptyObject, +> { + target: Ref; + events: Events; + resolve: (result: Result) => void; + active = new Map(); + children = new Map(); + pendingStart = new Map(); + pendingParents = new Map(); + + constructor( + target: Ref, + events: Events, + resolve: (result: Result) => void, + ) { + this.target = target; + this.events = events; + this.resolve = resolve; + } + + reset() { + for (const v of this.active.values()) v.end(); + this.active.clear(); + } + + processEvent(event: StreamEvent) { + event = event.slice(); + + // delete unreferenced nodes first, since they may get re-allocated + // in the next phase. + const deletedCount = event.shift(); + ASSERT(typeof deletedCount === "number"); + if (deletedCount === -1) { + // connection reset, delete all nodes + this.reset(); + } else { + for (const id of event.splice(0, deletedCount)) { + ASSERT(typeof id === "number"); + UNWRAP( + this.active.get(id as EncodedKey), + () => `dstkey ${id} not found, ${[...this.active.keys()]}`, + ).end(); + ASSERT(this.active.delete(id as EncodedKey)); + } + } + + // parse all of the events, updating all nodes and extracting creations + let resolving = false; + let resolvingValue: unknown = null; + let updateChildren: Array<{ + key: EncodedKey; + children: EncodedKey[]; + }> = []; + for (const chunk of event) { + ASSERT(typeof chunk === "object"); + if (Array.isArray(chunk)) { + if (typeof chunk[0] === "string") { + if (chunk[0] === "end") { + resolving = true; + resolvingValue = chunk[1]; + } else { + // custom event + // @ts-expect-error TODO: typescript soundness + this.events.emit(...chunk); + } + } else { + const ids = chunk as EncodedKey[]; + for (const id of ids) this.pendingParents.set(id, 0 as EncodedKey); + } + } else { + // node update or creation + const { + k: key, + t: text, + v: value, + e: total, + l: messages, + c: children, + E: estimatedTime, + V: formattedValue, + } = chunk; + const node = this.active.get(key); + if (children) { + for (const id of children) this.pendingParents.set(id, key); + updateChildren.push({ key, children }); + } + if (node) { + if (text) node.text = text; + if (value) node.value = value; + if (total) node.total = total; + for (const message of messages ?? []) node.log.writeMessage(message); + if (estimatedTime) node.estimatedTime = estimatedTime; + if (formattedValue) node.valueFormatter = () => formattedValue; + for (const id of children ?? []) this.pendingParents.set(id, key); + } else { + ASSERT(text != null, `DstKey ${key} not active but not created`); + this.pendingStart.set(key, { + text, + value, + total, + messages, + estimatedTime, + formattedValue, + }); + } + } + } + + let entry; + while (entry = this.pendingStart.entries().next().value) { + const [key, value] = entry; + this.startRecursive(key, value); + } + this.pendingParents.clear(); + + for (const { key, children } of updateChildren) { + const node = UNWRAP(this.active.get(key)); + this.children.set( + node, + children.map((id) => UNWRAP(this.active.get(id))), + ); + } + } + + startRecursive(key: EncodedKey, opts: PendingStart): Ref { + ASSERT(this.pendingStart.delete(key)); + const parentId = UNWRAP(this.pendingParents.get(key)); + let ref: Ref | null = parentId === 0 + ? this.target + : this.active.get(parentId) ?? null; + if (!ref) { + ref = this.startRecursive( + parentId, + UNWRAP(this.pendingStart.get(parentId)), + ); + } + const { text, estimatedTime, ...options } = opts; + const node = ref.start(text, { + ...options, + estimateCompletion: false, + }); + this.active.set(key, node); + node.estimatedTime = estimatedTime; + this.pendingParents.delete(key); + return node; + } +} + +interface PendingStart { + text: string; + value?: number; + total?: number; + messages?: log.Message[]; + estimatedTime?: number | null; + formattedValue?: string; +} + +/** this interface must not have more than 16 properties */ +interface DeltaFlags { + // the new status is always present since it's one bit + hidden: boolean; + passive: boolean; + showTotal: boolean; + + // for each packet that changed it stores out of line + changedValue: boolean; + changedTotal: boolean; + changedLogs: boolean; + changedText: boolean; + changedChildren: boolean; + changedEstimatedTime: boolean; + changedFormattedValue: boolean; + + estimatedTimeNonNull: boolean; + + futureProofingPacket: boolean; + deleted: boolean; +} + +/** + * converts a {@linkcode Root|progress.Root} into an byte stream for + * communicating progress over the process or network boundary. the output is a + * raw binary payload that uses an extremely compact representation for the + * nodes. if trying to transmit UTF-8, consider the JSON-encoding + * {@linkcode encodeEventStream} + * + * if the given root adds custom event handlers, they must all have + * json-serializable payloads. + * + * this function does not use recursion. + */ +export function encodeByteStream< + Result extends ts.Json, + EventMap extends { [key: string]: ts.Json[] }, +>( + root: Root, + options?: EncodeStreamOptions, +): ReadableStream { + const events = encodeEventStream(root, options); + const transform = new TransformStream({ transform: writeStreamEvent }); + return events.pipeThrough(transform); +} + +/** convert {@linkcode StreamEvent} into a byte payload. Mutates the input event. */ +function writeStreamEvent( + event: StreamEvent, + writer: stream.WriteTarget, +) { + using w = new stream.BufferedWriter(writer); + + let deleteCount = event.shift(); + ASSERT(typeof deleteCount === "number"); + + if (deleteCount === -1) { + deleteCount = 0; + w.write(header); + } else { + w.u8("\n".charCodeAt(0)); + + const deletedNodes = event.splice(0, deleteCount); + w.varUint(deletedNodes.length); + for (const dstKey of deletedNodes) { + ASSERT(typeof dstKey === "number"); + w.varUint(dstKey); + } + } + + const rootChildren = event.find((x): x is StreamRootChildren => + Array.isArray(x) && typeof x[0] === "number" + ); + const changedNodes = event.filter((x): x is StreamNode => + typeof x === "object" && !Array.isArray(x) + ); + w.varUint(changedNodes.length + +!!rootChildren); + for (const obj of changedNodes) { + ASSERT(typeof obj === "object" && "k" in obj); + const { + k: key, + t: text, + v: value, + e: total, + l: messages, + c: childrenSet, + E: estimatedTime, + V: formattedValue, + } = obj; + const { hidden, showTotal, passive } = UNWRAP(obj[kNode]); + ASSERT(key > 0); + ASSERT(messages == null || messages.length > 0); + + w.varUint(key); + w.u16(encodeDeltaFlags({ + hidden, + passive, + showTotal, + changedValue: value !== undefined, + changedTotal: total !== undefined, + changedLogs: messages !== undefined, + changedText: text !== undefined, + changedChildren: childrenSet !== undefined, + changedEstimatedTime: estimatedTime !== undefined, + changedFormattedValue: formattedValue !== undefined, + estimatedTimeNonNull: estimatedTime !== null, + })); + if (value !== undefined) w.f64(value); + if (total !== undefined) w.f64(total); + if (text !== undefined) w.stringWithLength(text); + if (messages !== undefined) { + w.varUint(messages.length); + for (const msg of messages) { + let level = logLevelSerialize.indexOf(msg.level); + if (level === -1) level = 0; + w.u8( + level + + (msg.scope ? 1 << 5 : 0) + + (msg.stack && msg.stack.length > 0 ? 1 << 6 : 0) + + (msg.custom !== undefined ? 1 << 7 : 0), + ); + w.stringWithLength(msg.text); + w.varUint(Math.floor(msg.time)); + if (msg.scope) w.stringWithLength(msg.scope); + if (msg.stack && msg.stack.length > 0) { + w.varUint(msg.stack.length); + for (const { fn, file, line, col } of msg.stack) { + w.stringWithLength(fn ?? ""); + w.stringWithLength(file ?? ""); + w.varUint(line ?? 0); + w.varUint(col ?? 0); + } + } + if (msg.custom !== undefined) { + w.stringWithLength(JSON.stringify(msg.custom)); + } + } + } + if (childrenSet !== undefined) { + w.varUint(childrenSet.length); + for (const key of childrenSet) w.varUint(key); + } + if (estimatedTime != null) w.varUint(Math.floor(estimatedTime)); + if (formattedValue !== undefined) { + w.stringWithLength(formattedValue); + } + } + if (rootChildren) { + w.varUint(0); + w.varUint(rootChildren.length); + for (const key of rootChildren) { + w.varUint(key); + } + } + + const customEvents = event.filter((x): x is StreamCustomEvent => + Array.isArray(x) && typeof x[0] === "string" + ); + w.varUint(customEvents.length); + for (const [event, ...args] of customEvents) { + w.stringWithLength(event); + w.stringWithLength(JSON.stringify(args)); + } +} + +/** + * decodes a stream created by `encodeByteStream` into managed calls to + * `target.start`. when cancelling, all the created nodes are destroyed. + */ +export function decodeByteStream< + Result = void, + EventMap extends Events.Map = ts.EmptyObject, +>( + encoded: ReadableStream, + target: Ref, +): async.Cancelable & Events { + let cancelled = false; + return decodeEventStream( + new ReadableStream({ + async start(controller) { + using reader = new stream.BufferedReader(encoded.getReader()); + try { + while (true) controller.enqueue(await readStreamEvent(reader)); + } catch (e) { + if (!cancelled) { + reader.cancel(e); + throw e; + } + } + }, + cancel(reason) { + cancelled = true; + encoded.cancel(reason); + }, + }), + target, + ); +} + +async function readStreamEvent(r: stream.BufferedReader): Promise { + const event: StreamEvent = []; + const start = await r.peekU8(); + if (start === "\n".charCodeAt(0)) { + await r.u8(); + let deletedCount = await r.varUint(); + event.push(deletedCount); + while (deletedCount-- > 0) event.push(await r.varUint()); + } else { + const readHeader = await r.readExactly(header.byteLength); + if (!bytes.eql(readHeader, header)) { + const m = JSON.stringify(string.decodeUtf8(readHeader)); + throw new Error(`Header mismatch, got ${m}`); + } + event.push(-1); + } + const changedNodes = await r.varUint(); + for (let i = 0; i < changedNodes; i += 1) { + const key = await r.varUint(); + if (key === 0) { + let count = await r.varUint(); + const array: StreamRootChildren = []; + while (count-- > 0) array.push(await r.varUint()); + event.push(array); + continue; + } + const flags = decodeDeltaFlags(await r.u16()); + const v = flags.changedValue ? await r.f64() : undefined; + const e = flags.changedTotal ? await r.f64() : undefined; + const t = flags.changedText ? await r.stringWithLength() : undefined; + let l: undefined | log.Message[]; + if (flags.changedLogs) { + l = []; + let len = await r.varUint(); + while (len > 0) { + const msgFlags = await r.u8(); + const level = logLevelSerialize[msgFlags & 0b1111] ?? "info"; + const hasScope = msgFlags && (1 << 5) > 0; + const hasStack = msgFlags && (1 << 6) > 0; + const hasCustom = msgFlags && (1 << 7) > 0; + const text = await r.stringWithLength(); + const time = await r.varUint(); + const scope = hasScope ? await r.stringWithLength() : undefined; + let stack: stack.Frame[] | undefined = undefined; + if (hasStack) { + stack = []; + let len = await r.varUint(); + while (len > 0) { + const fn = await r.stringWithLength(); + const file = await r.stringWithLength(); + const line = await r.varUint(); + const col = await r.varUint(); + stack.push({ + fn: fn || null, + file: file || null, + line: line || null, + col: col || null, + }); + } + } + const custom = hasCustom + ? JSON.parse(await r.stringWithLength()) + : undefined; + l.push({ level, text, time, scope, stack, custom }); + } + } + let c: undefined | EncodedKey[] = undefined; + if (flags.changedChildren) { + c = []; + let len = await r.varUint(); + while (len-- > 0) c.push(await r.varUint() as EncodedKey); + } + const E = flags.changedEstimatedTime + ? flags.estimatedTimeNonNull ? await r.varUint() : null + : undefined; + const V = flags.changedFormattedValue + ? await r.stringWithLength() + : undefined; + + // a "v2"'s packet will go here, contents ignored by this version. + if (flags.futureProofingPacket) { + const len = await r.varUint(); + void await r.readExactly(len); + } + + event.push({ + k: key as EncodedKey, + t, + v, + e, + s: flags.showTotal, + p: flags.passive, + l, + h: flags.hidden, + c, + E, + V, + }); + } + + // events + let eventCount = await r.varUint(); + while (eventCount-- > 0) { + const name = await r.stringWithLength(); + const args = JSON.parse(await r.stringWithLength()); + event.push([name, ...args]); + } + + return event; +} + +const logLevelSerialize = ["info", "error", "warn", "debug"] as const; + +function encodeDeltaFlags(node: Partial) { + return ( + (node.hidden ? 1 << 0 : 0) + + (node.passive ? 1 << 1 : 0) + + (node.showTotal ? 1 << 2 : 0) + + (node.changedValue ? 1 << 3 : 0) + + (node.changedTotal ? 1 << 4 : 0) + + (node.changedLogs ? 1 << 6 : 0) + + (node.changedText ? 1 << 7 : 0) + + (node.changedChildren ? 1 << 8 : 0) + + (node.changedEstimatedTime ? 1 << 9 : 0) + + (node.changedFormattedValue ? 1 << 10 : 0) + + (node.estimatedTimeNonNull ? 1 << 11 : 0) + + (node.futureProofingPacket ? 1 << 14 : 0) + + (node.deleted ? 1 << 15 : 0) + ); +} + +function decodeDeltaFlags(flags: number): DeltaFlags { + return { + hidden: (flags & 1 << 0) !== 0, + passive: (flags & 1 << 1) !== 0, + showTotal: (flags & 1 << 2) !== 0, + changedValue: (flags & 1 << 3) !== 0, + changedTotal: (flags & 1 << 4) !== 0, + changedLogs: (flags & 1 << 6) !== 0, + changedText: (flags & 1 << 7) !== 0, + changedChildren: (flags & 1 << 8) !== 0, + changedEstimatedTime: (flags & 1 << 9) !== 0, + changedFormattedValue: (flags & 1 << 10) !== 0, + estimatedTimeNonNull: (flags & 1 << 11) !== 0, + futureProofingPacket: (flags & 1 << 14) !== 0, + deleted: (flags & 1 << 15) !== 0, + }; +} + class KeyPool { - next = 1 as T; + next = 0; old: T[] = []; get() { const existing = this.old.shift(); if (existing) return existing as T; - const next = this.next = this.next + 1 as T; - return next as T; + return (this.next += 1) as T; } recycle(k: T) { @@ -851,14 +1713,33 @@ export class Ema implements EstimationAlgorithm { } } +/** according to a stream */ +type EncodedKey = number & { brand: typeof kNode }; + const globalKeyPool = new KeyPool(); -const global = +export const global: Ref = /** @__PURE__ */ ((root = new Root()) => (attachToScreen(root, log), root))(); +/** + * for testing. not covered by semver + * @internal + */ +const uncastInternals = /** @__PURE__ */ (() => ({ + readStreamEvent, + writeStreamEvent, + kNode: kNode as unknown as symbol, + newNode, + EncodedKey: 0 as EncodedKey, +}))(); +export const internals: typeof uncastInternals = uncastInternals; + +import * as ansi from "./string/ansi.ts"; import * as async from "./async.ts"; +import * as bytes from "./bytes.ts"; import * as log from "./log.ts"; -import * as ansi from "./string/ansi.ts"; -import * as ts from "./ts.ts"; +import * as stack from "./log/stack.ts"; +import * as stream from "./stream.ts"; import * as string from "./string.ts"; +import * as ts from "./ts.ts"; import { ASSERT, UNWRAP } from "./assert.ts"; import { Events } from "./Events.ts"; diff --git a/lib/stream.test.ts b/lib/stream.test.ts new file mode 100644 index 0000000000000000000000000000000000000000..510127f30e8cdda77f2f2a591572753d723241f5 --- /dev/null +++ b/lib/stream.test.ts @@ -0,0 +1,232 @@ +describe("BufferedWriter", () => { + test("basics", () => { + const chunks: Uint8Array[] = []; + const w = new stream.BufferedWriter((d) => chunks.push(d), 16); + + w.u8(1); + w.u16(258); + w.u32(0x04030201); + w.flush(); + + assert.equal(chunks.length, 1); + assert.deepEqual(chunks[0], new Uint8Array([1, 2, 1, 1, 2, 3, 4])); + }); + + test("auto-flush on large write", () => { + const chunks: Uint8Array[] = []; + const w = new stream.BufferedWriter((d) => chunks.push(d), 8); + + w.u32(1); + assert.equal(w.written, 4); + const large = new Uint8Array(20).fill(0xff); + w.write(large); + assert.equal(w.written, 0); + + assert.equal(chunks.length, 2); + assert.equal(chunks[0]!.length, 4); + assert.equal(chunks[1]!.length, 20); + }); + + test("transfer based flush", () => { + const chunks: Uint8Array[] = []; + const w = new stream.BufferedWriter((d) => chunks.push(d), 4); + + w.u8(1); + w.u8(2); + w.u8(3); + w.u8(4); + w.flush(); + w.u8(2); + w.u8(3); + w.u8(4); + w.u8(5); + w[Symbol.dispose](); + + assert.equal(chunks.length, 2); + assert.deepEqual(chunks[0], new Uint8Array([1, 2, 3, 4])); + assert.deepEqual(chunks[1], new Uint8Array([2, 3, 4, 5])); + }); + + test("ensure more than total capacity", () => { + const chunks: Uint8Array[] = []; + const w = new stream.BufferedWriter((d) => chunks.push(d), 2); + assert.equal(w.size, 2); + w.u8(4); + w.ensureUnusedCapacity(8); // flushes [ 4 ] + w.write(new Uint8Array([4, 5, 5, 8, 8, 1, 2])); + assert.equal(w.written, 7); + assert.equal(w.size, 8); + w.u32(42); + w.flush(); + + assert.equal(chunks.length, 3); + assert.deepEqual(chunks[0], new Uint8Array([4])); + assert.deepEqual(chunks[1], new Uint8Array([4, 5, 5, 8, 8, 1, 2])); + assert.deepEqual(chunks[2], new Uint8Array([42, 0, 0, 0])); + }); + + test("varUint", () => { + const chunks: Uint8Array[] = []; + const w = new stream.BufferedWriter((d) => chunks.push(d)); + + w.varUint(0); + w.varUint(127); + w.varUint(128); + w.varUint(16383); + w.flush(); + + assert.deepEqual(chunks[0], new Uint8Array([0, 127, 128, 1, 255, 127])); + }); + + test("dispose", () => { + const chunks: Uint8Array[] = []; + { + using w = new stream.BufferedWriter((d) => chunks.push(d)); + w.u8(42); + } + assert.equal(chunks.length, 1); + assert.equal(chunks[0]![0], 42); + }); +}); + +describe("BufferedReader", () => { + test("basic read", async () => { + const [s, w] = stream.bufferedReadableStream(); + w.u8(1); + w.u16(258); + w.flush(); + w.u32(0x04030201); + w.flush(); + + using r = new stream.BufferedReader(s.getReader()); + assert.equal(await r.u8(), 1); + assert.equal(await r.u16(), 258); + assert.equal(await r.u16(), 0x0201); + assert.equal(await r.u16(), 0x0403); + }); + + test("varUint", async () => { + const [s, w] = stream.bufferedReadableStream(); + w.varUint(0); + w.varUint(127); + w.varUint(128); + w.varUint(16383); + w.varUint(1768210855277); + w.varUint(521); + w.close(); + + using r = new stream.BufferedReader(s.getReader()); + assert.equal(await r.varUint(), 0); + assert.equal(await r.varUint(), 127); + assert.equal(await r.varUint(), 128); + assert.equal(await r.varUint(), 16383); + assert.equal(await r.varUint(), 1768210855277); + assert.equal(await r.varUint(), 521); + }); + + test("cross-chunk read", async () => { + const [s, w] = stream.bufferedReadableStream(); + w.u8(1); + w.flush(); + w.u8(2); + w.flush(); + + using r = new stream.BufferedReader(s.getReader()); + assert.equal(await r.u16(), 0x0201); + }); + + test("readExactly", async () => { + const [s, w] = stream.bufferedReadableStream(); + w.write(new Uint8Array([1, 2, 3, 4, 5])); + w.flush(); + + using r = new stream.BufferedReader(s.getReader()); + const buf = await r.readExactly(3); + assert.deepEqual(buf, new Uint8Array([1, 2, 3])); + assert.equal(await r.u8(), 4); + }); + + test("readExactly cutting chunks", async () => { + { + const [s, w] = stream.bufferedReadableStream(); + w.write(new Uint8Array([1, 2, 3])); + w.flush(); + w.write(new Uint8Array([4, 5, 6])); + w.flush(); + w.write(new Uint8Array([7, 8, 9])); + w.flush(); + w.write(new Uint8Array([10, 11, 12])); + w.flush(); + + using r = new stream.BufferedReader(s.getReader()); + assert.deepEqual(await r.readExactly(4), new Uint8Array([1, 2, 3, 4])); + assert.deepEqual(await r.readExactly(4), new Uint8Array([5, 6, 7, 8])); + assert.deepEqual(await r.readExactly(4), new Uint8Array([9, 10, 11, 12])); + } + }); + + test("readUpTo", async () => { + const [s, w] = stream.bufferedReadableStream(); + w.write(new Uint8Array([1, 2])); + w.close(); + + using r = new stream.BufferedReader(s.getReader()); + const buf = await r.readUpTo(10); + assert.equal(buf.length, 2); + }); + + test("stringWithLength", async () => { + const [s, w] = stream.bufferedReadableStream(); + w.stringWithLength("hello"); + w.flush(); + + using r = new stream.BufferedReader(s.getReader()); + assert.equal(await r.stringWithLength(), "hello"); + }); + + test("some numeric types", async () => { + const [s, w] = stream.bufferedReadableStream(); + + w.i8(-1); + w.i16(-1); + w.i32(-1); + w.i64(-1n); + w.u32(4_000_000_000); + w.u64(0xffffffffffffffffn); + w.f16(1.5); + w.f32(1.5); + w.f64(1.5); + w.close(); + + using r = new stream.BufferedReader(s.getReader()); + assert.equal(await r.i8(), -1); + assert.equal(await r.i16(), -1); + assert.equal(await r.i32(), -1); + assert.equal(await r.i64(), -1n); + assert.equal(await r.u32(), 4_000_000_000); + assert.equal(await r.u64(), 0xffffffffffffffffn); + assert.equal(await r.f16(), 1.5); + assert.equal(await r.f32(), 1.5); + assert.equal(await r.f64(), 1.5); + }); + + test("cancel", () => { + let cancelReason: unknown; + const s = new ReadableStream({ + start() {}, + cancel(reason) { + cancelReason = reason; + }, + }); + + const r = new stream.BufferedReader(s.getReader()); + const err = new Error("test cancel"); + r.cancel(err); + + assert.equal(cancelReason, err); + }); +}); + +import { describe, test } from "node:test"; +import { strict as assert } from "node:assert"; +import * as stream from "./stream.ts"; diff --git a/lib/stream.ts b/lib/stream.ts new file mode 100644 index 0000000000000000000000000000000000000000..e793de3e9f0c7e12b0feaefbae54d64ec932c87a --- /dev/null +++ b/lib/stream.ts @@ -0,0 +1,485 @@ +/** + * this @module contains helpers for creating and consuming `ReadableStream`, + * particularly with binary payloads. + */ + +const shared = new Uint8Array(8); +const sharedView = new DataView(shared.buffer); + +/** write bytes `data` to a destination */ +export type WriteBytesFn = (data: Uint8Array) => void; + +/** something that can be written to */ +export type WriteTarget = + | WriteBytesFn + | { enqueue(chunk: Uint8Array): void; close?(): void }; + +/** + * `ReadableStreamDefaultController.enqueue` does not buffer data in any way. + * for use cases such as a streaming write of many small binary chunks, such as + * differently sized integers in the case of `progress.encodeByteStream`, this + * batches those many small writes into their own `ArrayBuffer` objects. + * + * when passing a `ReadableStreamDefaultController`, the `BufferedWriter` + * "owns" it, and will close the stream when disposed. to prevent that, pass + * a `WriteBytesFn` instead. + * + * ```ts + * const w = new BufferedWriter(unbufferedWrite); + * w.write(string.encodeUtf8("a null-terminated string")); + * w.u8(0); + * w.write(string.encodeUtf8("a second null-terminated string")); + * w.u8(0); + * + * w.flush(); // don't forget to flush! + * // `unbufferedWrite` will be called once with 57 bytes + * ``` + */ +export class BufferedWriter { + #view: DataView; + #written: number = 0; + #flush: WriteBytesFn; + #size: number; + #onClose?: (() => void) | null; + + constructor( + flush: WriteTarget, + size: number = 8192, + ) { + this.#flush = typeof flush === "function" + ? flush + : flush.enqueue.bind(flush); + this.#view = new DataView(new ArrayBuffer(size)); + this.#size = size; + this.#onClose = typeof flush === "function" + ? null + : flush.close?.bind(flush); + } + + flush() { + const written = this.#written; + if (written === 0) return; + const size = this.#size; + if (written === size) { + // transfer + this.#flush(new Uint8Array(this.#view.buffer)); + } else { + // copy + this.#flush(new Uint8Array(this.#view.buffer.slice(0, written))); + } + this.#view = new DataView(new ArrayBuffer(size)); + this.#written = 0; + } + + close() { + const written = this.#written; + if (written === 0) return; + const size = this.#size; + if (written === size) { + // transfer + this.#flush(new Uint8Array(this.#view.buffer)); + } else { + // copy + this.#flush(new Uint8Array(this.#view.buffer.slice(0, written))); + } + this.#written = 0; + this.#size = 0; + this.#view = new DataView(new ArrayBuffer(0)); + this.#onClose?.(); + } + + [Symbol.dispose]() { + this.close(); + } + + write(bytes: Uint8Array | ArrayBuffer) { + const len = bytes.byteLength; + if (len >= this.size) { + this.flush(); + this.#flush("buffer" in bytes ? bytes : new Uint8Array(bytes)); + return; + } + this.ensureUnusedCapacity(len); + new Uint8Array(this.#view.buffer, this.#written) + .set(bytes instanceof ArrayBuffer ? new Uint8Array(bytes) : bytes); + this.#written += bytes.byteLength; + } + + varUint(int: number) { + ASSERT(int >= 0); + do { + const data = int % (1 << 7); + int = Math.floor(int / (1 << 7)); + const continuation = int > 0 ? (1 << 7) : 0; + this.u8(data + continuation); + } while (int > 0); + } + + u8(int: number) { + const size = 1; + this.ensureUnusedCapacity(size); + this.#view.setUint8(this.#written, int); + this.#written += size; + } + + u16(int: number) { + const size = 2; + this.ensureUnusedCapacity(size); + this.#view.setUint16(this.#written, int, true); + this.#written += size; + } + + u32(int: number) { + const size = 4; + this.ensureUnusedCapacity(size); + this.#view.setUint32(this.#written, int, true); + this.#written += size; + } + + u64(int: bigint) { + const size = 8; + this.ensureUnusedCapacity(size); + this.#view.setBigUint64(this.#written, int, true); + this.#written += size; + } + + i8(int: number) { + const size = 1; + this.ensureUnusedCapacity(size); + this.#view.setUint8(this.#written, int); + this.#written += size; + } + + i16(int: number) { + const size = 2; + this.ensureUnusedCapacity(size); + this.#view.setInt16(this.#written, int, true); + this.#written += size; + } + + i32(int: number) { + const size = 4; + this.ensureUnusedCapacity(size); + this.#view.setInt32(this.#written, int, true); + this.#written += size; + } + + i64(int: bigint) { + const size = 8; + this.ensureUnusedCapacity(size); + this.#view.setBigInt64(this.#written, int, true); + this.#written += size; + } + + f16(float: number) { + const size = 2; + this.ensureUnusedCapacity(size); + this.#view.setFloat16(this.#written, float, true); + this.#written += size; + } + + f32(float: number) { + const size = 4; + this.ensureUnusedCapacity(size); + this.#view.setFloat32(this.#written, float, true); + this.#written += size; + } + + f64(float: number) { + const size = 8; + this.ensureUnusedCapacity(size); + this.#view.setFloat64(this.#written, float, true); + this.#written += size; + } + + stringWithLength(text: string) { + const buf = string.encodeUtf8(text); + this.varUint(buf.byteLength); + this.write(buf); + } + + get written(): number { + return this.#written; + } + get size(): number { + return this.#size; + } + + ensureUnusedCapacity(extra: number) { + const size = this.size; + const written = this.#written; + if (written + extra > size) { + if (extra > size) { + this.resize(extra); + } else { + this.flush(); + } + } + } + + resize(newCapacity: number) { + newCapacity = Math.round(newCapacity); + const written = this.#written; + if (written > 0) { + // transfer + this.#flush(new Uint8Array(this.#view.buffer, 0, written)); + this.#written = 0; + } + this.#size = newCapacity; + this.#view = new DataView(new ArrayBuffer(newCapacity)); + } +} + +/** + * a reader for `ReadableStream` which allows the caller to specify the + * chunking size and read numbers at a time. + */ +export class BufferedReader implements Disposable { + #reader: ReadableStreamDefaultReader | null; + #offset = 0; + #buffer: Uint8Array[] = []; + #available = 0; + + constructor(reader: ReadableStreamDefaultReader) { + this.#reader = reader; + } + + get available() { + return this.#available; + } + + /** Releasing will discard the bytes remaining in the buffer. */ + releaseLock() { + UNWRAP(this.#reader, "reader disposed").releaseLock(); + this.#reader = null; + } + + cancel(error: unknown) { + UNWRAP(this.#reader, "reader disposed").cancel(error); + } + + async ensureAvailableOrFalse(size: number): Promise { + if (size <= 0 || this.#available >= size) return true; + const reader = UNWRAP(this.#reader, "reader disposed"); + const buffer = this.#buffer; + while (this.#available < size) { + const { value, done } = await reader.read(); + if (done) break; + buffer.push(value); + this.#available += value.byteLength; + } + return this.#available >= size; + } + async ensureAvailable(size: number): Promise { + ASSERT(await this.ensureAvailableOrFalse(size), "End of stream"); + } + + #collectLimitedView(length: number): DataView { + ASSERT(this.available >= length); + const buffer = this.#buffer; + let chunk = UNWRAP(buffer[0]); + let offset = this.#offset; + if (chunk.byteLength - offset >= length) { + this.#shiftOffset(length); + return new DataView(chunk.buffer, chunk.byteOffset + offset, length); + } + let x = 0; + while (true) { + for ( + ; + offset < chunk.byteLength && x < length; + offset += 1, x += 1 + ) shared[x] = chunk[offset]!; + if (x === length) { + if (offset === chunk.byteLength) { + ASSERT(buffer.shift()); + this.#offset = 0; + } else this.#offset = offset; + this.#available -= length; + return sharedView; + } + ASSERT(buffer.shift()); + chunk = UNWRAP(buffer[0]); + offset = 0; + } + } + #collectBuffer(length: number) { + ASSERT(this.available >= length); + this.#available -= length; + const buffer = this.#buffer; + let offset = this.#offset; + + // slice an existing data view + if (buffer[0] && buffer[0].byteLength - offset >= length) { + const first = buffer[0]; + if (first.byteLength - offset === length) { + this.#offset = 0; + ASSERT(buffer.shift()); + return offset ? first.subarray(offset) : first; + } + this.#offset += length; + return first.subarray(offset, offset + length); + } + + // join by allocating a new buffer + const joined = new Uint8Array(length); + let i = 0; + do { + const chunk = UNWRAP(buffer[0]); + const { byteLength } = chunk; + const copyLen = Math.min(byteLength, length); + joined.set( + copyLen === byteLength + ? chunk + : chunk.subarray(offset, offset + copyLen), + i, + ); + i += copyLen; + length -= copyLen; + if (copyLen === byteLength) ASSERT(buffer.shift()); + if (length === 0) this.#offset = copyLen; + } while (length > 0); + return joined; + } + #shiftOffset(offset: number) { + const newOffset = this.#offset += offset; + this.#available -= offset; + if (UNWRAP(this.#buffer[0]).byteLength === newOffset) { + this.#buffer.shift(); + this.#offset = 0; + } else { + this.#offset = newOffset; + } + } + + async u8() { + const int = await this.peekU8(); + this.#shiftOffset(1); + return int; + } + + async peekU8() { + await this.ensureAvailable(1); + return UNWRAP(UNWRAP(this.#buffer[0])[this.#offset]); + } + + async i8() { + await this.ensureAvailable(2); + return this.#collectLimitedView(1).getInt8(0); + } + + async u16() { + await this.ensureAvailable(2); + return this.#collectLimitedView(2).getUint16(0, true); + } + + async i16() { + await this.ensureAvailable(2); + return this.#collectLimitedView(2).getInt16(0, true); + } + + async u32() { + await this.ensureAvailable(4); + return this.#collectLimitedView(4).getUint32(0, true); + } + + async i32() { + await this.ensureAvailable(4); + return this.#collectLimitedView(4).getInt32(0, true); + } + + async u64() { + await this.ensureAvailable(8); + return this.#collectLimitedView(8).getBigUint64(0, true); + } + + async i64() { + await this.ensureAvailable(8); + return this.#collectLimitedView(8).getBigInt64(0, true); + } + + async f16() { + await this.ensureAvailable(2); + return this.#collectLimitedView(2).getFloat16(0, true); + } + + async f32() { + await this.ensureAvailable(4); + return this.#collectLimitedView(4).getFloat32(0, true); + } + + async f64() { + await this.ensureAvailable(8); + return this.#collectLimitedView(8).getFloat64(0, true); + } + + async varUint() { + let result = 0; + let shift = 1; + let byte; + do { + byte = await this.u8(); + result += (byte & 0x7F) * shift; + shift *= 128; + } while (byte & 0x80); + return result; + } + + async stringWithLength() { + const len = await this.varUint(); + return string.decodeUtf8(await this.readExactly(len)); + } + + /** Read up to a maximum number of bytes */ + async readUpTo(size: number): Promise { + await this.ensureAvailableOrFalse(size); + return this.#collectBuffer(Math.min(size, this.#available)); + } + async readExactly(size: number): Promise { + await this.ensureAvailable(size); + return this.#collectBuffer(size); + } + + [Symbol.dispose]() { + return this.releaseLock(); + } +} + +/** easily create a readable stream with a writer */ +export function bufferedReadableStream(): [ + ReadableStream, + BufferedWriter, + ReadableStreamDefaultController, +] { + let unbuffered: ReadableStreamDefaultController | null = null; + const reader = new ReadableStream({ + start: (controller) => unbuffered = controller, + }); + ASSERT(unbuffered); + return [reader, new BufferedWriter(unbuffered), unbuffered]; +} + +/** easily create a readable stream with a writer */ +export function readWritePair(): [ + ReadableStream, + ReadableStreamDefaultController, +] { + let unbuffered: ReadableStreamDefaultController | null = null; + const reader = new ReadableStream({ + start: (controller) => unbuffered = controller, + }); + ASSERT(unbuffered); + return [reader, unbuffered]; +} + +export function fromChunks(chunks: T[]) { + return new ReadableStream({ + start(c) { + for (const chunk of chunks) c.enqueue(chunk); + c.close(); + }, + }); +} + +import * as string from "./string.ts"; +import { ASSERT, UNWRAP } from "./assert.ts"; diff --git a/lib/ts.ts b/lib/ts.ts index 0ea59052641c7175ea88ade99a5859c42e7aad78..5633d88a53cf88e1c0e6ce7112d01d01803f050a 100644 --- a/lib/ts.ts +++ b/lib/ts.ts @@ -1,5 +1,13 @@ export type Timer = ReturnType; export type Interval = ReturnType; +/** opposite of the built-in `Readonly` type */ +export type Writeable = { -readonly [P in keyof T]: T[P] }; +/** opposite of the built-in `Readonly` type, but recursive */ +// TODO: tuples? +export type DeepWriteable = { + -readonly [P in keyof T]: T[P] extends ReadonlyArray ? Array + : DeepWriteable; +}; /** * redeclared here because it is only provided by `lib: ["dom"]` and not @@ -23,3 +31,51 @@ export function defer(fn: () => void): Dispose { f[Symbol.dispose] = f; return f; } + +export type Mixin = A & B; + +/** + * mutates `base` to implement all methods of `other`, using `other` as state. + */ +export function mixin( + base: A, + other: B, +): Mixin { + const prototype = Object.getPrototypeOf(other); + for ( + const key of [ + ...Object.getOwnPropertyNames(prototype), + ...Object.getOwnPropertySymbols(prototype), + ] + ) { + if (key === "constructor") continue; + if (key in base) continue; + const fn = prototype[key]; + if (typeof fn !== "function") continue; + (base as Record)[key] = fn.bind(other); + } + return base as A & B; +} + +export type EmptyObject = Record; + +/** TODO: this type has many subtle bugs */ +export type ToJson = T extends JsonValue ? T + : T extends { toJSON(): infer J } ? ToJson + : T extends Set | Map | Record + ? EmptyObject + : T extends unknown[] ? { + [K in keyof Omit]: T[K] extends undefined | void + ? null + : ToJson; + } + : T extends object ? { [K in keyof Omit]: ToJson } + : never; +export type Json = JsonValue | undefined | void; +export type JsonValue = + | string + | number + | boolean + | null + | JsonValue[] + | { [k: string]: Exclude };