import type { InternalWaitConstraintsEntry, InternalWaitConstraintsOptions, } from "./internal/wait-constraints.ts" import { internalAssertTimeoutOption, internalBindWaitConstraints, internalCleanupWaitConstraints, internalThrowIfAborted, } from "./internal/wait-constraints.ts" import type { InternalWaitQueueEntry } from "./internal/wait-queue.ts" import { InternalWaitQueue } from "./internal/wait-queue.ts" /** * @description 断言闭锁初始计数为合法的非负整数。 */ const internalAssertInitialCount = (count: number): void => { if (Number.isFinite(count) === false || Number.isInteger(count) === false || count < 0) { throw new RangeError( "CountDownLatch initialCount must be a finite integer greater than or equal to 0.", ) } } /** * @description 断言单次递减步长为合法的正整数。 */ const internalAssertDelta = (count: number): void => { if (Number.isFinite(count) === false || Number.isInteger(count) === false || count <= 0) { throw new RangeError("CountDownLatch count must be a finite integer greater than 0.") } } /** * @description 表示 `CountDownLatch` 等待者在内部队列中的节点。 */ interface InternalCountDownLatchQueueEntry extends InternalWaitQueueEntry, InternalWaitConstraintsEntry { resolve: () => void reject: (reason: unknown) => void isSettled: boolean } /** * @description 表示调用 `CountDownLatch.wait()` 时可选的等待约束。 */ export interface CountDownLatchWaitOptions extends InternalWaitConstraintsOptions { /** * @description 等待闭锁打开的最长时间,单位为毫秒。 */ timeout?: number | undefined /** * @description 用于主动取消等待的信号。 */ abortSignal?: AbortSignal | undefined } /** * @description 表示一个倒计时闭锁。 * * 闭锁在剩余计数降到 0 前保持关闭,所有等待者都会被挂起; * 一旦打开,当前及后续等待都会立即通过,且该实例不会再次关闭。 */ export class CountDownLatch { private readonly queue: InternalWaitQueue private internalRemainingCount: number /** * @description 使用给定初始计数创建闭锁。 */ constructor(initialCount: number) { internalAssertInitialCount(initialCount) this.queue = new InternalWaitQueue() this.internalRemainingCount = initialCount } /** * @description 判断闭锁是否已经打开。 */ isOpen(): boolean { return this.internalRemainingCount === 0 } /** * @description 以同步方式检测当前是否可以继续执行而无需等待。 */ tryWait(): boolean { return this.internalRemainingCount === 0 } /** * @description 读取距离打开闭锁还差多少次 `countDown()`。 */ getRemainingCount(): number { return this.internalRemainingCount } /** * @description 读取当前挂起的等待者数量。 */ getPendingCount(): number { return this.queue.getCount() } /** * @description 在闭锁打开时一次性唤醒全部等待者。 */ private resolveAll(): void { 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() } } /** * @description 以失败状态结束指定等待者,并从队列中摘除。 */ private rejectEntry(entry: InternalCountDownLatchQueueEntry, reason: unknown): void { if (entry.isSettled === true) { return } entry.isSettled = true internalCleanupWaitConstraints(entry) this.queue.remove(entry) entry.reject(reason) } /** * @description 将闭锁剩余计数减少指定值。 * * 当剩余计数降到 0 时,所有等待者会立即被放行。 */ countDown(count: number = 1): number { internalAssertDelta(count) if (this.internalRemainingCount === 0) { return 0 } this.internalRemainingCount = Math.max(0, this.internalRemainingCount - count) if (this.internalRemainingCount === 0) { this.resolveAll() } return this.internalRemainingCount } /** * @description `countDown()` 的语义化别名,用于强调“到达一次事件”。 */ arrive(count: number = 1): number { return this.countDown(count) } /** * @description 等待直到闭锁打开。 */ async wait(options: CountDownLatchWaitOptions = {}): Promise { internalAssertTimeoutOption("CountDownLatch wait", options.timeout) internalThrowIfAborted("CountDownLatch wait", options.abortSignal) if (this.internalRemainingCount === 0) { return } await new Promise((resolve, reject) => { const entry: InternalCountDownLatchQueueEntry = { resolve, reject, isSettled: false, isQueued: false, previous: undefined, next: undefined, timeoutId: undefined, abortSignal: undefined, abortListener: undefined, } internalBindWaitConstraints(entry, options, { abortOperation: "CountDownLatch wait", timeoutOperation: "CountDownLatch wait", onAbort: (error) => { this.rejectEntry(entry, error) }, onTimeout: (error) => { this.rejectEntry(entry, error) }, }) this.queue.enqueue(entry) if (this.internalRemainingCount === 0) { this.resolveAll() } }) } }