import { BrokenBarrierError, CoordinationAbortError } from "./errors.ts" import type { InternalWaitConstraintsEntry, InternalWaitConstraintsOptions, } from "./internal/wait-constraints.ts" import { internalAssertTimeoutOption, internalBindWaitConstraints, internalCleanupWaitConstraints, } from "./internal/wait-constraints.ts" import type { InternalWaitQueueEntry } from "./internal/wait-queue.ts" import { InternalWaitQueue } from "./internal/wait-queue.ts" /** * @description 断言栅栏的参与者数量是一个有效的正整数。 */ const internalAssertParticipantCount = (participantCount: number): void => { if ( Number.isFinite(participantCount) === false || Number.isInteger(participantCount) === false || participantCount <= 0 ) { throw new RangeError("Barrier participantCount must be a finite integer greater than 0.") } } /** * @description 表示某一代栅栏中单个等待者对应的内部节点。 * * 每个节点既要能挂入等待队列,也要能绑定超时与中止约束,最终在整代完成或整代破坏时统一结算。 */ interface InternalBarrierQueueEntry extends InternalWaitQueueEntry, InternalWaitConstraintsEntry { resolve: (generation: number) => void reject: (reason: unknown) => void isSettled: boolean } /** * @description 表示调用 `Barrier.signalAndWait()` 时可选的等待约束。 */ export interface BarrierWaitOptions extends InternalWaitConstraintsOptions { /** * @description 单次等待允许持续的最长时间,单位为毫秒。 */ timeout?: number | undefined /** * @description 用于主动中止当前等待的信号。 */ abortSignal?: AbortSignal | undefined } /** * @description 表示一个循环栅栏,用于让固定数量的参与者在同一同步点会合。 * * 每次当所有参与者都调用 `signalAndWait()` 后,当前这一代会被一次性完成, * 所有等待者都会收到相同的代号,然后栅栏自动进入下一代继续复用。 */ export class Barrier { private readonly queue: InternalWaitQueue private internalArrivedCount: number private internalGeneration: number private readonly participantCount: number /** * @description 创建一个需要固定参与者数量才能完成一代的栅栏。 */ constructor(participantCount: number) { internalAssertParticipantCount(participantCount) this.queue = new InternalWaitQueue() this.internalArrivedCount = 0 this.internalGeneration = 0 this.participantCount = participantCount } /** * @description 读取每一代完成所需的参与者总数。 */ getParticipantCount(): number { return this.participantCount } /** * @description 读取当前所处的代号。 * * 该值会在一代正常完成或被破坏后递增。 */ getGeneration(): number { return this.internalGeneration } /** * @description 读取当前这一代中仍在等待其余参与者的任务数量。 */ getPendingCount(): number { return this.queue.getCount() } /** * @description 读取当前这一代距离完成还差多少参与者。 */ getRemainingCount(): number { return this.participantCount - this.internalArrivedCount } /** * @description 以成功状态完成当前代。 * * 完成时会先重置计数并推进代号,再唤醒所有尚未结算的等待者,确保下一代可以立即开始接收新的参与者。 */ private completeGeneration(completedGeneration: number): void { this.internalArrivedCount = 0 this.internalGeneration = this.internalGeneration + 1 while (this.queue.isEmpty() === false) { const entry = this.queue.shift() if (entry === undefined) { return } if (entry.isSettled === true) { continue } entry.isSettled = true internalCleanupWaitConstraints(entry) entry.resolve(completedGeneration) } } /** * @description 以失败状态破坏当前代。 * * 触发 abort 或 timeout 的参与者收到其原始错误,其余同代等待者统一收到 `BrokenBarrierError`, * 这样调用方可以区分“我自己失败了”与“别人导致这一代失效”。 */ private breakGeneration(failingEntry: InternalBarrierQueueEntry, reason: unknown): void { if (failingEntry.isSettled === true) { return } const brokenReason = new BrokenBarrierError(reason) this.internalArrivedCount = 0 this.internalGeneration = this.internalGeneration + 1 while (this.queue.isEmpty() === false) { const entry = this.queue.shift() if (entry === undefined) { return } if (entry.isSettled === true) { continue } entry.isSettled = true internalCleanupWaitConstraints(entry) if (entry === failingEntry) { entry.reject(reason) } else { entry.reject(brokenReason) } } } /** * @description 声明当前参与者已经到达同步点,并等待整代完成。 * * 返回值是本次完成时对应的代号;如果本调用导致最后一个参与者到达,则会立即完成整代并同步返回。 */ async signalAndWait(options: BarrierWaitOptions = {}): Promise { internalAssertTimeoutOption("Barrier wait", options.timeout) if (options.abortSignal?.aborted === true) { throw new CoordinationAbortError("Barrier wait", options.abortSignal.reason) } const currentGeneration = this.internalGeneration this.internalArrivedCount = this.internalArrivedCount + 1 if (this.internalArrivedCount === this.participantCount) { this.completeGeneration(currentGeneration) return currentGeneration } const generation = await new Promise((resolve, reject) => { const entry: InternalBarrierQueueEntry = { resolve, reject, isSettled: false, isQueued: false, previous: undefined, next: undefined, timeoutId: undefined, abortSignal: undefined, abortListener: undefined, } internalBindWaitConstraints(entry, options, { abortOperation: "Barrier wait", timeoutOperation: "Barrier wait", onAbort: (error) => { this.breakGeneration(entry, error) }, onTimeout: (error) => { this.breakGeneration(entry, error) }, }) this.queue.enqueue(entry) }) return generation } }