diff --git a/.changeset/fast-fibers-drain.md b/.changeset/fast-fibers-drain.md new file mode 100644 index 00000000000..7e8433b4e68 --- /dev/null +++ b/.changeset/fast-fibers-drain.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Use a constant-time queue when draining fiber messages. diff --git a/packages/effect/src/internal/fiberRuntime.ts b/packages/effect/src/internal/fiberRuntime.ts index f93e75c92f9..df76ad90936 100644 --- a/packages/effect/src/internal/fiberRuntime.ts +++ b/packages/effect/src/internal/fiberRuntime.ts @@ -28,6 +28,7 @@ import type { Logger } from "../Logger.js" import * as LogLevel from "../LogLevel.js" import type * as MetricLabel from "../MetricLabel.js" import * as Micro from "../Micro.js" +import * as MutableQueue from "../MutableQueue.js" import * as MRef from "../MutableRef.js" import * as Option from "../Option.js" import { pipeArguments } from "../Pipeable.js" @@ -297,7 +298,7 @@ export class FiberRuntime extends Effectable.Class() + private _queue = MutableQueue.unbounded() private _children: Set> | null = null private _observers = new Array<(exit: Exit.Exit) => void>() private _running = false @@ -441,7 +442,7 @@ export class FiberRuntime extends Effectable.Class extends Effectable.Class extends Effectable.Class 0 && !this._running) { + if (!MutableQueue.isEmpty(this._queue) && !this._running) { this._running = true if (evaluationSignal === EvaluationSignalYieldNow) { this.drainQueueLaterOnExecutor() @@ -725,8 +726,8 @@ export class FiberRuntime extends Effectable.Class ) { let cur = cur0 - while (this._queue.length > 0) { - const message = this._queue.splice(0, 1)[0] + while (!MutableQueue.isEmpty(this._queue)) { + const message = MutableQueue.poll(this._queue, undefined)! // @ts-expect-error cur = drainQueueWhileRunningTable[message._tag](this, runtimeFlags, cur, message) } @@ -969,7 +970,7 @@ export class FiberRuntime extends Effectable.Class exit) } else { - if (this._queue.length === 0) { + if (MutableQueue.isEmpty(this._queue)) { // No more messages to process, so we will allow the fiber to end life: this.setExitValue(exit) } else { @@ -1009,7 +1010,7 @@ export class FiberRuntime extends Effectable.Class 0) { + if (!MutableQueue.isEmpty(this._queue)) { this.drainQueueLaterOnExecutor() } } @@ -1366,7 +1367,7 @@ export class FiberRuntime extends Effectable.Class 0) { + if (!MutableQueue.isEmpty(this._queue)) { cur = this.drainQueueWhileRunning(this.currentRuntimeFlags, cur) } if (!this._isYielding) {