tuil
ReferencePackages@mwillbanks/tuil-operationsAPI

OperationExecutor

class exported by @mwillbanks/tuil-operations.

View rawEdit

class

Public class exported by @mwillbanks/tuil-operations.

export class OperationExecutor<TResult = unknown> {
  readonly #store;
  readonly #observers = new Set<(event: OperationEvent) => void>();
  #controller = new AbortController();
  readonly #children: ChildOperationHandle[] = [];
  readonly #childAttempts = new Map<ChildOperationHandle, number>();
  readonly #childUnsubscribes = new Map<ChildOperationHandle, () => void>();
  readonly #logBatcher: Batcher<{
    readonly generation: number;
    readonly line: string;
  }>;
  #lastResult?: TResult;
  #generation = 0;
  #attemptGeneration = 0;

  constructor(
    readonly definition: OperationDefinition<TResult>,
    readonly options: OperationExecutorOptions = {},
  ) {
    this.#store = createNusmStore<OperationSnapshot<TResult>>(
      Object.freeze({
        id: definition.id,
        title: definition.title,
        description: definition.description,
        status: "idle",
        attempt: 0,
        children: Object.freeze([]),
        metadata: Object.freeze({ ...(definition.metadata ?? {}) }),
        logs: Object.freeze([]),
      }),
    );
    this.#logBatcher = new Batcher(
      (entries) => {
        const lines = entries
          .filter((entry) => entry.generation === this.#generation)
          .map((entry) => entry.line);
        if (lines.length === 0) return;
        const logs = [...this.state.logs, ...lines].slice(
          -(this.options.maxLogs ?? 1_000),
        );
        this.#update({ logs: Object.freeze(logs) });
      },
      { maxSize: 50, wait: 16 },
    );
  }

  get state(): OperationSnapshot<TResult> {
    return this.#store.state;
  }

  subscribe(observer: () => void): () => void {
    const subscription = this.#store.subscribe(observer);
    return () => subscription.unsubscribe();
  }

  observe(observer: (event: OperationEvent) => void): () => void {
    this.#observers.add(observer);
    return () => this.#observers.delete(observer);
  }

  restore(snapshot: OperationSnapshot<TResult>): void {
    if (
      ["running", "queued", "waiting", "retrying"].includes(this.state.status)
    ) {
      throw new Error("Cannot restore while the operation is running");
    }
    if (
      ["running", "queued", "waiting", "retrying"].includes(snapshot.status)
    ) {
      throw new Error("Cannot restore an in-flight operation snapshot");
    }
    this.#generation += 1;
    this.#attemptGeneration += 1;
    this.#lastResult = snapshot.result;
    this.#update(
      Object.freeze({
        ...snapshot,
        children: Object.freeze([...snapshot.children]),
        metadata: Object.freeze({ ...snapshot.metadata }),
        logs: Object.freeze([...snapshot.logs]),
      }),
    );
  }

  async execute(signal?: AbortSignal): Promise<TResult> {
    if (["running", "waiting", "retrying"].includes(this.state.status)) {
      throw new Error(`Operation "${this.definition.id}" is already running`);
    }
    if (this.#controller.signal.aborted) {
      this.#controller = new AbortController();
    }
    this.#generation += 1;
    this.#attemptGeneration += 1;
    const generation = this.#generation;
    this.#disposeChildren();
    this.#lastResult = undefined;
    this.#update({
      status: "idle",
      progress: undefined,
      attempt: 0,
      startedAt: undefined,
      completedAt: undefined,
      result: undefined,
      error: undefined,
      children: Object.freeze([]),
      logs: Object.freeze([]),
    });
    const combined = signal
      ? AbortSignal.any([signal, this.#controller.signal])
      : this.#controller.signal;
    const retries = Math.max(0, this.definition.retries ?? 0);
    let lastError: unknown;
    this.#transition("queued");
    for (let attempt = 1; attempt <= retries + 1; attempt += 1) {
      combined.throwIfAborted();
      this.#update({
        attempt,
        startedAt: this.state.startedAt ?? Date.now(),
        error: undefined,
      });
      this.#transition(attempt === 1 ? "running" : "retrying");
      const attemptController = new AbortController();
      const attemptGeneration = ++this.#attemptGeneration;
      const attemptSignal = AbortSignal.any([
        combined,
        attemptController.signal,
      ]);
      let attemptSucceeded = false;
      let attemptChildrenCleaned = false;
      const context = this.#context(
        attemptSignal,
        attempt,
        generation,
        attemptGeneration,
      );
      const execution = Promise.resolve().then(() =>
        this.definition.run(context),
      );
      let settled = false;
      void execution.then(
        () => {
          settled = true;
        },
        () => {
          settled = true;
        },
      );
      const timeout = this.definition.timeout;
      let timer: ReturnType<typeof setTimeout> | undefined;
      let timedOut = false;
      let removeAbortListener: (() => void) | undefined;
      const aborted = new Promise<never>((_resolve, reject) => {
        const rejectAborted = () => reject(attemptSignal.reason);
        attemptSignal.addEventListener("abort", rejectAborted, { once: true });
        removeAbortListener = () =>
          attemptSignal.removeEventListener("abort", rejectAborted);
      });
      try {
        const attempts: Promise<TResult>[] = [execution, aborted];
        if (timeout && timeout > 0) {
          attempts.push(
            new Promise<never>((_, reject) => {
              timer = setTimeout(() => {
                timedOut = true;
                const error = new OperationTimeoutError(
                  `Operation timed out after ${timeout}ms`,
                );
                attemptController.abort(error);
                reject(error);
              }, timeout);
            }),
          );
        }
        const result = await Promise.race(attempts);
        combined.throwIfAborted();
        this.#lastResult = result;
        this.#logBatcher.flush();
        this.#update({ result, completedAt: Date.now() });
        this.#transition("succeeded");
        attemptSucceeded = true;
        return result;
      } catch (error) {
        this.#cleanupAttemptChildren(attemptGeneration, false);
        attemptChildrenCleaned = true;
        lastError = error;
        if (combined.aborted) {
          this.#logBatcher.flush();
          this.#update({
            error: operationError(combined.reason ?? error),
            completedAt: Date.now(),
          });
          this.#transition("cancelled");
          throw combined.reason ?? error;
        }
        if (error instanceof OperationBlockedError) {
          this.#logBatcher.flush();
          this.#update({ error: operationError(error) });
          this.#transition("blocked");
          throw error;
        }
        if (timedOut) {
          await Promise.resolve();
          if (!settled) break;
        }
        if (attempt <= retries) {
          this.#update({ error: operationError(error) });
          if ((this.definition.retryDelay ?? 0) > 0) {
            this.#transition("waiting");
            try {
              await abortableDelay(this.definition.retryDelay ?? 0, combined);
            } catch (delayError) {
              this.#logBatcher.flush();
              this.#update({
                error: operationError(delayError),
                completedAt: Date.now(),
              });
              this.#transition("cancelled");
              throw delayError;
            }
          }
        }
      } finally {
        if (timer) clearTimeout(timer);
        removeAbortListener?.();
        if (!attemptChildrenCleaned) {
          this.#cleanupAttemptChildren(attemptGeneration, attemptSucceeded);
        }
        if (this.#attemptGeneration === attemptGeneration) {
          this.#attemptGeneration += 1;
        }
      }
    }
    this.#logBatcher.flush();
    this.#update({
      error: operationError(lastError),
      completedAt: Date.now(),
    });
    this.#transition("failed");
    throw lastError;
  }

  cancel(reason: unknown = new DOMException("Cancelled", "AbortError")): void {
    this.#attemptGeneration += 1;
    this.#controller.abort(reason);
    for (const child of this.#children) child.cancel(reason);
  }

  async rollback(signal?: AbortSignal): Promise<void> {
    const activeSignal = signal ?? new AbortController().signal;
    activeSignal.throwIfAborted();
    const context = this.#context(
      activeSignal,
      this.state.attempt,
      this.#generation,
      this.#attemptGeneration,
    );
    for (const child of [...this.#children].reverse()) {
      await child.rollback(activeSignal);
      activeSignal.throwIfAborted();
    }
    await this.definition.rollback?.(this.#lastResult, context);
    activeSignal.throwIfAborted();
  }

  skip(): void {
    if (this.state.status !== "idle" && this.state.status !== "queued") {
      throw new Error("Only idle or queued operations can be skipped");
    }
    this.#update({ completedAt: Date.now() });
    this.#transition("skipped");
  }

  dispose(): void {
    this.cancel(new DOMException("Disposed", "AbortError"));
    this.#disposeChildren();
    this.#logBatcher.cancel();
    this.#observers.clear();
  }

  #context(
    signal: AbortSignal,
    attempt: number,
    generation: number,
    attemptGeneration: number,
  ): OperationContext {
    const active = () =>
      generation === this.#generation &&
      attemptGeneration === this.#attemptGeneration &&
      !signal.aborted;
    return {
      signal,
      attempt,
      updateProgress: (progress) => {
        if (active()) {
          this.#update({ progress: Object.freeze({ ...progress }) });
        }
      },
      log: (line) => {
        if (active()) this.#logBatcher.addItem({ generation, line });
      },
      block: (message) => {
        throw new OperationBlockedError(message);
      },
      waitFor: async <T>(work: Promise<T>) => {
        if (active()) this.#transition("waiting");
        try {
          return await work;
        } finally {
          if (active()) this.#transition("running");
        }
      },
      runChild: async <T>(definition: OperationDefinition<T>) => {
        signal.throwIfAborted();
        if (!active()) throw new DOMException("Stale operation", "AbortError");
        const child = new OperationExecutor(definition, this.options);
        this.#children.push(child);
        this.#childAttempts.set(child, attemptGeneration);
        this.#childUnsubscribes.set(
          child,
          child.subscribe(() => {
            if (active())
              this.#update({
                children: Object.freeze(
                  this.#children.map((candidate) => candidate.state),
                ),
              });
          }),
        );
        if (active())
          this.#update({
            children: Object.freeze(
              this.#children.map((candidate) => candidate.state),
            ),
          });
        return child.execute(signal);
      },
    };
  }

  #cleanupAttemptChildren(
    attemptGeneration: number,
    preserveSucceeded: boolean,
  ): void {
    const retained: ChildOperationHandle[] = [];
    for (const child of this.#children) {
      const belongsToAttempt =
        this.#childAttempts.get(child) === attemptGeneration;
      const shouldPreserve =
        preserveSucceeded && child.state.status === "succeeded";
      if (!belongsToAttempt || shouldPreserve) {
        retained.push(child);
        continue;
      }
      this.#disposeChild(child);
    }
    this.#children.splice(0, this.#children.length, ...retained);
    this.#update({
      children: Object.freeze(this.#children.map((child) => child.state)),
    });
  }

  #disposeChild(child: ChildOperationHandle): void {
    this.#childUnsubscribes.get(child)?.();
    this.#childUnsubscribes.delete(child);
    this.#childAttempts.delete(child);
    child.dispose();
  }

  #disposeChildren(): void {
    for (const child of this.#children) this.#disposeChild(child);
    this.#children.length = 0;
    this.#childAttempts.clear();
    this.#childUnsubscribes.clear();
  }

  #transition(status: OperationStatus): void {
    const previousStatus = this.state.status;
    this.#update({ status });
    const event = Object.freeze({
      operation: this.state,
      previousStatus,
      at: Date.now(),
    });
    for (const observer of this.#observers) {
      try {
        observer(event);
      } catch (error) {
        try {
          this.options.onObserverError?.(error);
        } catch {
          // Observer error reporting must not corrupt operation state.
        }
      }
    }
  }

  #update(patch: Partial<OperationSnapshot<TResult>>): void {
    this.#store.setState((state) => Object.freeze({ ...state, ...patch }));
  }
}

Members

MemberTypeRequiredDescriptionRelated types
#storeimport("nusm").NusmStore<OperationSnapshot<TResult>>YesThe #store member uses the import("nusm").NusmStore<OperationSnapshot<TResult>> contract.OperationSnapshot
#observersSet<(event: OperationEvent) => void>YesThe #observers member uses the Set<(event: OperationEvent) => void> contract.OperationEvent
#controllerAbortControllerYesThe #controller member uses the AbortController contract.
#childrenChildOperationHandle[]YesThe #children member uses the ChildOperationHandle[] contract.
#childAttemptsMap<ChildOperationHandle, number>YesThe #childAttempts member uses the Map<ChildOperationHandle, number> contract.
#childUnsubscribesMap<ChildOperationHandle, () => void>YesThe #childUnsubscribes member uses the Map<ChildOperationHandle, () => void> contract.
#logBatcherBatcher<&#123; readonly generation: number; readonly line: string; &#125;>YesThe #logBatcher member uses the Batcher<&#123; readonly generation: number; readonly line: string; &#125;> contract.
#lastResultTResult | undefinedNoThe #lastResult member uses the TResult | undefined contract.
#generationnumberYesThe #generation member uses the number contract.
#attemptGenerationnumberYesThe #attemptGeneration member uses the number contract.
__constructoranyYesThe __constructor member uses the any contract.
stateOperationSnapshot<TResult>YesThe state member uses the OperationSnapshot<TResult> contract.OperationSnapshot
subscribe(observer: () => void) => () => voidYesThe subscribe member uses the (observer: () => void) => () => void contract.
observe(observer: (event: OperationEvent) => void) => () => voidYesThe observe member uses the (observer: (event: OperationEvent) => void) => () => void contract.OperationEvent
restore(snapshot: OperationSnapshot<TResult>) => voidYesThe restore member uses the (snapshot: OperationSnapshot<TResult>) => void contract.OperationSnapshot
execute(signal?: AbortSignal) => Promise<TResult>YesThe execute member uses the (signal?: AbortSignal) => Promise<TResult> contract.
cancel(reason?: unknown) => voidYesThe cancel member uses the (reason?: unknown) => void contract.
rollback(signal?: AbortSignal) => Promise<void>YesThe rollback member uses the (signal?: AbortSignal) => Promise<void> contract.
skip() => voidYesThe skip member uses the () => void contract.
dispose() => voidYesThe dispose member uses the () => void contract.
#context(signal: AbortSignal, attempt: number, generation: number, attemptGeneration: number) => OperationContextYesThe #context member uses the (signal: AbortSignal, attempt: number, generation: number, attemptGeneration: number) => OperationContext contract.OperationContext
#cleanupAttemptChildren(attemptGeneration: number, preserveSucceeded: boolean) => voidYesThe #cleanupAttemptChildren member uses the (attemptGeneration: number, preserveSucceeded: boolean) => void contract.
#disposeChild(child: ChildOperationHandle) => voidYesThe #disposeChild member uses the (child: ChildOperationHandle) => void contract.
#disposeChildren() => voidYesThe #disposeChildren member uses the () => void contract.
#transition(status: OperationStatus) => voidYesThe #transition member uses the (status: OperationStatus) => void contract.OperationStatus
#update(patch: Partial<OperationSnapshot<TResult>>) => voidYesThe #update member uses the (patch: Partial<OperationSnapshot<TResult>>) => void contract.OperationSnapshot

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-operations

On this page