tuil
ReferencePackages@mwillbanks/tuil-streamingAPI

BoundedBackpressureController

class exported by @mwillbanks/tuil-streaming.

View rawEdit

class

Public class exported by @mwillbanks/tuil-streaming.

export class BoundedBackpressureController implements BackpressureController {
  readonly #highWaterMark: number;
  readonly #waiters: BackpressureWaiter[] = [];
  #size = 0;

  constructor(highWaterMark = 64) {
    if (!Number.isSafeInteger(highWaterMark) || highWaterMark < 1) {
      throw new Error("Backpressure highWaterMark must be positive");
    }
    this.#highWaterMark = highWaterMark;
  }

  get desiredSize(): number {
    return this.#highWaterMark - this.#size;
  }

  async wait(signal?: AbortSignal): Promise<void> {
    if (signal?.aborted) throw signal.reason;
    if (this.#size < this.#highWaterMark) {
      this.#size += 1;
      return;
    }
    await new Promise<void>((resolve, reject) => {
      let waiter: BackpressureWaiter;
      const admit = () => {
        signal?.removeEventListener("abort", waiter.abort);
        this.#size += 1;
        resolve();
      };
      const abort = () => {
        const index = this.#waiters.indexOf(waiter);
        if (index >= 0) this.#waiters.splice(index, 1);
        reject(signal?.reason);
      };
      waiter = { admit, abort };
      this.#waiters.push(waiter);
      signal?.addEventListener("abort", waiter.abort, { once: true });
    });
  }

  release(count = 1): void {
    this.#size = Math.max(0, this.#size - Math.max(1, count));
    while (this.#size < this.#highWaterMark && this.#waiters.length > 0) {
      this.#waiters.shift()?.admit();
    }
  }
}

Members

MemberTypeRequiredDescriptionRelated types
#highWaterMarknumberYesThe #highWaterMark member uses the number contract.
#waitersBackpressureWaiter[]YesThe #waiters member uses the BackpressureWaiter[] contract.
#sizenumberYesThe #size member uses the number contract.
__constructoranyYesThe __constructor member uses the any contract.
desiredSizenumberYesThe desiredSize member uses the number contract.
wait(signal?: AbortSignal) => Promise<void>YesThe wait member uses the (signal?: AbortSignal) => Promise<void> contract.
release(count?: number) => voidYesThe release member uses the (count?: number) => void contract.

Parameters

This declaration has no public members.

Returns

This declaration does not return a value.

Throws

No thrown errors are documented for this declaration.

Source

View the secondary source reference

Package

@mwillbanks/tuil-streaming

On this page