| 1 | /** |
| 2 | * this module contains helpers for creating and consuming `ReadableStream`, |
| 3 | * particularly with binary payloads. |
| 4 | * |
| 5 | * @module |
| 6 | */ |
| 7 | |
| 8 | const shared = new Uint8Array(8); |
| 9 | const sharedView = new DataView(shared.buffer); |
| 10 | |
| 11 | /** write bytes `data` to a destination */ |
| 12 | export type WriteBytesFn = (data: Uint8Array) => void; |
| 13 | |
| 14 | /** something that can be written to */ |
| 15 | export type WriteTarget = |
| 16 | | WriteBytesFn |
| 17 | | { enqueue(chunk: Uint8Array): void; close?(): void }; |
| 18 | |
| 19 | /** |
| 20 | * `ReadableStreamDefaultController.enqueue` does not buffer data in any way. |
| 21 | * for use cases such as a streaming write of many small binary chunks, such as |
| 22 | * differently sized integers in the case of `progress.encodeByteStream`, this |
| 23 | * batches those many small writes into their own `ArrayBuffer` objects. |
| 24 | * |
| 25 | * when passing a `ReadableStreamDefaultController`, the `BufferedWriter` |
| 26 | * "owns" it, and will close the stream when disposed. to prevent that, pass |
| 27 | * a `WriteBytesFn` instead. |
| 28 | * |
| 29 | * ```ts |
| 30 | * const w = new BufferedWriter(unbufferedWrite); |
| 31 | * w.write(string.encodeUtf8("a null-terminated string")); |
| 32 | * w.u8(0); |
| 33 | * w.write(string.encodeUtf8("a second null-terminated string")); |
| 34 | * w.u8(0); |
| 35 | * |
| 36 | * w.flush(); // don't forget to flush! |
| 37 | * // `unbufferedWrite` will be called once with 57 bytes |
| 38 | * ``` |
| 39 | */ |
| 40 | export class BufferedWriter { |
| 41 | #view: DataView<ArrayBuffer>; |
| 42 | #written: number = 0; |
| 43 | #flush: WriteBytesFn; |
| 44 | #size: number; |
| 45 | #onClose?: (() => void) | null; |
| 46 | |
| 47 | constructor( |
| 48 | flush: WriteTarget, |
| 49 | size: number = 8192, |
| 50 | ) { |
| 51 | this.#flush = typeof flush === "function" |
| 52 | ? flush |
| 53 | : flush.enqueue.bind(flush); |
| 54 | this.#view = new DataView(new ArrayBuffer(size)); |
| 55 | this.#size = size; |
| 56 | this.#onClose = typeof flush === "function" |
| 57 | ? null |
| 58 | : flush.close?.bind(flush); |
| 59 | } |
| 60 | |
| 61 | flush() { |
| 62 | const written = this.#written; |
| 63 | if (written === 0) return; |
| 64 | const size = this.#size; |
| 65 | if (written === size) { |
| 66 | // transfer |
| 67 | this.#flush(new Uint8Array(this.#view.buffer)); |
| 68 | } else { |
| 69 | // copy |
| 70 | this.#flush(new Uint8Array(this.#view.buffer.slice(0, written))); |
| 71 | } |
| 72 | this.#view = new DataView(new ArrayBuffer(size)); |
| 73 | this.#written = 0; |
| 74 | } |
| 75 | |
| 76 | close() { |
| 77 | const written = this.#written; |
| 78 | if (written === 0) return; |
| 79 | const size = this.#size; |
| 80 | if (written === size) { |
| 81 | // transfer |
| 82 | this.#flush(new Uint8Array(this.#view.buffer)); |
| 83 | } else { |
| 84 | // copy |
| 85 | this.#flush(new Uint8Array(this.#view.buffer.slice(0, written))); |
| 86 | } |
| 87 | this.#written = 0; |
| 88 | this.#size = 0; |
| 89 | this.#view = new DataView(new ArrayBuffer(0)); |
| 90 | this.#onClose?.(); |
| 91 | } |
| 92 | |
| 93 | [Symbol.dispose]() { |
| 94 | this.close(); |
| 95 | } |
| 96 | |
| 97 | write(bytes: Uint8Array<ArrayBuffer> | ArrayBuffer) { |
| 98 | const len = bytes.byteLength; |
| 99 | if (len >= this.size) { |
| 100 | this.flush(); |
| 101 | this.#flush("buffer" in bytes ? bytes : new Uint8Array(bytes)); |
| 102 | return; |
| 103 | } |
| 104 | this.ensureUnusedCapacity(len); |
| 105 | new Uint8Array(this.#view.buffer, this.#written) |
| 106 | .set(bytes instanceof ArrayBuffer ? new Uint8Array(bytes) : bytes); |
| 107 | this.#written += bytes.byteLength; |
| 108 | } |
| 109 | |
| 110 | varUint(int: number) { |
| 111 | ASSERT(int >= 0); |
| 112 | do { |
| 113 | const data = int % (1 << 7); |
| 114 | int = Math.floor(int / (1 << 7)); |
| 115 | const continuation = int > 0 ? (1 << 7) : 0; |
| 116 | this.u8(data + continuation); |
| 117 | } while (int > 0); |
| 118 | } |
| 119 | |
| 120 | u8(int: number) { |
| 121 | const size = 1; |
| 122 | this.ensureUnusedCapacity(size); |
| 123 | this.#view.setUint8(this.#written, int); |
| 124 | this.#written += size; |
| 125 | } |
| 126 | |
| 127 | u16(int: number) { |
| 128 | const size = 2; |
| 129 | this.ensureUnusedCapacity(size); |
| 130 | this.#view.setUint16(this.#written, int, true); |
| 131 | this.#written += size; |
| 132 | } |
| 133 | |
| 134 | u32(int: number) { |
| 135 | const size = 4; |
| 136 | this.ensureUnusedCapacity(size); |
| 137 | this.#view.setUint32(this.#written, int, true); |
| 138 | this.#written += size; |
| 139 | } |
| 140 | |
| 141 | u64(int: bigint) { |
| 142 | const size = 8; |
| 143 | this.ensureUnusedCapacity(size); |
| 144 | this.#view.setBigUint64(this.#written, int, true); |
| 145 | this.#written += size; |
| 146 | } |
| 147 | |
| 148 | i8(int: number) { |
| 149 | const size = 1; |
| 150 | this.ensureUnusedCapacity(size); |
| 151 | this.#view.setUint8(this.#written, int); |
| 152 | this.#written += size; |
| 153 | } |
| 154 | |
| 155 | i16(int: number) { |
| 156 | const size = 2; |
| 157 | this.ensureUnusedCapacity(size); |
| 158 | this.#view.setInt16(this.#written, int, true); |
| 159 | this.#written += size; |
| 160 | } |
| 161 | |
| 162 | i32(int: number) { |
| 163 | const size = 4; |
| 164 | this.ensureUnusedCapacity(size); |
| 165 | this.#view.setInt32(this.#written, int, true); |
| 166 | this.#written += size; |
| 167 | } |
| 168 | |
| 169 | i64(int: bigint) { |
| 170 | const size = 8; |
| 171 | this.ensureUnusedCapacity(size); |
| 172 | this.#view.setBigInt64(this.#written, int, true); |
| 173 | this.#written += size; |
| 174 | } |
| 175 | |
| 176 | f16(float: number) { |
| 177 | const size = 2; |
| 178 | this.ensureUnusedCapacity(size); |
| 179 | this.#view.setFloat16(this.#written, float, true); |
| 180 | this.#written += size; |
| 181 | } |
| 182 | |
| 183 | f32(float: number) { |
| 184 | const size = 4; |
| 185 | this.ensureUnusedCapacity(size); |
| 186 | this.#view.setFloat32(this.#written, float, true); |
| 187 | this.#written += size; |
| 188 | } |
| 189 | |
| 190 | f64(float: number) { |
| 191 | const size = 8; |
| 192 | this.ensureUnusedCapacity(size); |
| 193 | this.#view.setFloat64(this.#written, float, true); |
| 194 | this.#written += size; |
| 195 | } |
| 196 | |
| 197 | stringWithLength(text: string) { |
| 198 | const buf = string.encodeUtf8(text); |
| 199 | this.varUint(buf.byteLength); |
| 200 | this.write(buf); |
| 201 | } |
| 202 | |
| 203 | get written(): number { |
| 204 | return this.#written; |
| 205 | } |
| 206 | get size(): number { |
| 207 | return this.#size; |
| 208 | } |
| 209 | |
| 210 | ensureUnusedCapacity(extra: number) { |
| 211 | const size = this.size; |
| 212 | const written = this.#written; |
| 213 | if (written + extra > size) { |
| 214 | if (extra > size) { |
| 215 | this.resize(extra); |
| 216 | } else { |
| 217 | this.flush(); |
| 218 | } |
| 219 | } |
| 220 | } |
| 221 | |
| 222 | resize(newCapacity: number) { |
| 223 | newCapacity = Math.round(newCapacity); |
| 224 | const written = this.#written; |
| 225 | if (written > 0) { |
| 226 | // transfer |
| 227 | this.#flush(new Uint8Array(this.#view.buffer, 0, written)); |
| 228 | this.#written = 0; |
| 229 | } |
| 230 | this.#size = newCapacity; |
| 231 | this.#view = new DataView(new ArrayBuffer(newCapacity)); |
| 232 | } |
| 233 | } |
| 234 | |
| 235 | /** |
| 236 | * a reader for `ReadableStream` which allows the caller to specify the |
| 237 | * chunking size and read numbers at a time. |
| 238 | */ |
| 239 | export class BufferedReader implements Disposable { |
| 240 | #reader: ReadableStreamDefaultReader<Uint8Array> | null; |
| 241 | #offset = 0; |
| 242 | #buffer: Uint8Array[] = []; |
| 243 | #available = 0; |
| 244 | |
| 245 | constructor(reader: ReadableStreamDefaultReader<Uint8Array>) { |
| 246 | this.#reader = reader; |
| 247 | } |
| 248 | |
| 249 | get available(): number { |
| 250 | return this.#available; |
| 251 | } |
| 252 | |
| 253 | /** Releasing will discard the bytes remaining in the buffer. */ |
| 254 | releaseLock() { |
| 255 | UNWRAP(this.#reader, "reader disposed").releaseLock(); |
| 256 | this.#reader = null; |
| 257 | } |
| 258 | |
| 259 | cancel(error: unknown) { |
| 260 | UNWRAP(this.#reader, "reader disposed").cancel(error); |
| 261 | } |
| 262 | |
| 263 | async ensureAvailableOrFalse(size: number): Promise<boolean> { |
| 264 | if (size <= 0 || this.#available >= size) return true; |
| 265 | const reader = UNWRAP(this.#reader, "reader disposed"); |
| 266 | const buffer = this.#buffer; |
| 267 | while (this.#available < size) { |
| 268 | const { value, done } = await reader.read(); |
| 269 | if (done) break; |
| 270 | buffer.push(value); |
| 271 | this.#available += value.byteLength; |
| 272 | } |
| 273 | return this.#available >= size; |
| 274 | } |
| 275 | async ensureAvailable(size: number): Promise<void> { |
| 276 | ASSERT(await this.ensureAvailableOrFalse(size), "End of stream"); |
| 277 | } |
| 278 | |
| 279 | #collectLimitedView(length: number): DataView { |
| 280 | ASSERT(this.available >= length); |
| 281 | const buffer = this.#buffer; |
| 282 | let chunk = UNWRAP(buffer[0]); |
| 283 | let offset = this.#offset; |
| 284 | if (chunk.byteLength - offset >= length) { |
| 285 | this.#shiftOffset(length); |
| 286 | return new DataView(chunk.buffer, chunk.byteOffset + offset, length); |
| 287 | } |
| 288 | let x = 0; |
| 289 | while (true) { |
| 290 | for ( |
| 291 | ; |
| 292 | offset < chunk.byteLength && x < length; |
| 293 | offset += 1, x += 1 |
| 294 | ) shared[x] = chunk[offset]!; |
| 295 | if (x === length) { |
| 296 | if (offset === chunk.byteLength) { |
| 297 | ASSERT(buffer.shift()); |
| 298 | this.#offset = 0; |
| 299 | } else this.#offset = offset; |
| 300 | this.#available -= length; |
| 301 | return sharedView; |
| 302 | } |
| 303 | ASSERT(buffer.shift()); |
| 304 | chunk = UNWRAP(buffer[0]); |
| 305 | offset = 0; |
| 306 | } |
| 307 | } |
| 308 | #collectBuffer(length: number) { |
| 309 | ASSERT(this.available >= length); |
| 310 | this.#available -= length; |
| 311 | const buffer = this.#buffer; |
| 312 | let offset = this.#offset; |
| 313 | |
| 314 | // slice an existing data view |
| 315 | if (buffer[0] && buffer[0].byteLength - offset >= length) { |
| 316 | const first = buffer[0]; |
| 317 | if (first.byteLength - offset === length) { |
| 318 | this.#offset = 0; |
| 319 | ASSERT(buffer.shift()); |
| 320 | return offset ? first.subarray(offset) : first; |
| 321 | } |
| 322 | this.#offset += length; |
| 323 | return first.subarray(offset, offset + length); |
| 324 | } |
| 325 | |
| 326 | // join by allocating a new buffer |
| 327 | const joined = new Uint8Array(length); |
| 328 | let i = 0; |
| 329 | do { |
| 330 | const chunk = UNWRAP(buffer[0]); |
| 331 | const { byteLength } = chunk; |
| 332 | const copyLen = Math.min(byteLength, length); |
| 333 | joined.set( |
| 334 | copyLen === byteLength |
| 335 | ? chunk |
| 336 | : chunk.subarray(offset, offset + copyLen), |
| 337 | i, |
| 338 | ); |
| 339 | i += copyLen; |
| 340 | length -= copyLen; |
| 341 | if (copyLen === byteLength) ASSERT(buffer.shift()); |
| 342 | if (length === 0) this.#offset = copyLen; |
| 343 | } while (length > 0); |
| 344 | return joined; |
| 345 | } |
| 346 | #shiftOffset(offset: number) { |
| 347 | const newOffset = this.#offset += offset; |
| 348 | this.#available -= offset; |
| 349 | if (UNWRAP(this.#buffer[0]).byteLength === newOffset) { |
| 350 | this.#buffer.shift(); |
| 351 | this.#offset = 0; |
| 352 | } else { |
| 353 | this.#offset = newOffset; |
| 354 | } |
| 355 | } |
| 356 | |
| 357 | async u8(): Promise<number> { |
| 358 | const int = await this.peekU8(); |
| 359 | this.#shiftOffset(1); |
| 360 | return int; |
| 361 | } |
| 362 | |
| 363 | async peekU8(): Promise<number> { |
| 364 | await this.ensureAvailable(1); |
| 365 | return UNWRAP(UNWRAP(this.#buffer[0])[this.#offset]); |
| 366 | } |
| 367 | |
| 368 | async i8(): Promise<number> { |
| 369 | await this.ensureAvailable(2); |
| 370 | return this.#collectLimitedView(1).getInt8(0); |
| 371 | } |
| 372 | |
| 373 | async u16(): Promise<number> { |
| 374 | await this.ensureAvailable(2); |
| 375 | return this.#collectLimitedView(2).getUint16(0, true); |
| 376 | } |
| 377 | |
| 378 | async i16(): Promise<number> { |
| 379 | await this.ensureAvailable(2); |
| 380 | return this.#collectLimitedView(2).getInt16(0, true); |
| 381 | } |
| 382 | |
| 383 | async u32(): Promise<number> { |
| 384 | await this.ensureAvailable(4); |
| 385 | return this.#collectLimitedView(4).getUint32(0, true); |
| 386 | } |
| 387 | |
| 388 | async i32(): Promise<number> { |
| 389 | await this.ensureAvailable(4); |
| 390 | return this.#collectLimitedView(4).getInt32(0, true); |
| 391 | } |
| 392 | |
| 393 | async u64(): Promise<bigint> { |
| 394 | await this.ensureAvailable(8); |
| 395 | return this.#collectLimitedView(8).getBigUint64(0, true); |
| 396 | } |
| 397 | |
| 398 | async i64(): Promise<bigint> { |
| 399 | await this.ensureAvailable(8); |
| 400 | return this.#collectLimitedView(8).getBigInt64(0, true); |
| 401 | } |
| 402 | |
| 403 | async f16(): Promise<number> { |
| 404 | await this.ensureAvailable(2); |
| 405 | return this.#collectLimitedView(2).getFloat16(0, true); |
| 406 | } |
| 407 | |
| 408 | async f32(): Promise<number> { |
| 409 | await this.ensureAvailable(4); |
| 410 | return this.#collectLimitedView(4).getFloat32(0, true); |
| 411 | } |
| 412 | |
| 413 | async f64(): Promise<number> { |
| 414 | await this.ensureAvailable(8); |
| 415 | return this.#collectLimitedView(8).getFloat64(0, true); |
| 416 | } |
| 417 | |
| 418 | async varUint(): Promise<number> { |
| 419 | let result = 0; |
| 420 | let shift = 1; |
| 421 | let byte; |
| 422 | do { |
| 423 | byte = await this.u8(); |
| 424 | result += (byte & 0x7F) * shift; |
| 425 | shift *= 128; |
| 426 | } while (byte & 0x80); |
| 427 | return result; |
| 428 | } |
| 429 | |
| 430 | async stringWithLength(): Promise<string> { |
| 431 | const len = await this.varUint(); |
| 432 | return string.decodeUtf8(await this.readExactly(len)); |
| 433 | } |
| 434 | |
| 435 | /** Read up to a maximum number of bytes */ |
| 436 | async readUpTo(size: number): Promise<Uint8Array> { |
| 437 | await this.ensureAvailableOrFalse(size); |
| 438 | return this.#collectBuffer(Math.min(size, this.#available)); |
| 439 | } |
| 440 | async readExactly(size: number): Promise<Uint8Array> { |
| 441 | await this.ensureAvailable(size); |
| 442 | return this.#collectBuffer(size); |
| 443 | } |
| 444 | |
| 445 | [Symbol.dispose]() { |
| 446 | this.releaseLock(); |
| 447 | } |
| 448 | } |
| 449 | |
| 450 | /** easily create a readable stream with a writer */ |
| 451 | export function bufferedReadableStream(): [ |
| 452 | ReadableStream, |
| 453 | BufferedWriter, |
| 454 | ReadableStreamDefaultController<Uint8Array>, |
| 455 | ] { |
| 456 | let unbuffered: ReadableStreamDefaultController<Uint8Array> | null = null; |
| 457 | const reader = new ReadableStream<Uint8Array>({ |
| 458 | start: (controller) => unbuffered = controller, |
| 459 | }); |
| 460 | ASSERT(unbuffered); |
| 461 | return [reader, new BufferedWriter(unbuffered), unbuffered]; |
| 462 | } |
| 463 | |
| 464 | /** easily create a readable stream with a writer */ |
| 465 | export function readWritePair(): [ |
| 466 | ReadableStream, |
| 467 | ReadableStreamDefaultController, |
| 468 | ] { |
| 469 | let unbuffered: ReadableStreamDefaultController<Uint8Array> | null = null; |
| 470 | const reader = new ReadableStream<Uint8Array>({ |
| 471 | start: (controller) => unbuffered = controller, |
| 472 | }); |
| 473 | ASSERT(unbuffered); |
| 474 | return [reader, unbuffered]; |
| 475 | } |
| 476 | |
| 477 | import { ASSERT, UNWRAP } from "./assert.ts"; |
| 478 | import * as string from "./string.ts"; |