MCPcopy Create free account
hub / github.com/Effect-TS/effect / offerAllUnsafe

Function offerAllUnsafe

packages/effect/src/Queue.ts:812–844  ·  view source on GitHub ↗
(self: Enqueue<A, E>, messages: Iterable<A>)

Source from the content-addressed store, hash-verified

810 * @since 4.0.0
811 */
812export const offerAllUnsafe = <A, E>(self: Enqueue<A, E>, messages: Iterable<A>): Array<A> => {
813 if (self.state._tag !== "Open") {
814 return Arr.fromIterable(messages)
815 } else if (
816 self.capacity === Number.POSITIVE_INFINITY ||
817 self.strategy === "sliding"
818 ) {
819 MutableList.appendAll(self.messages, messages)
820 if (self.strategy === "sliding") {
821 MutableList.takeN(self.messages, self.messages.length - self.capacity)
822 }
823 scheduleReleaseTaker(self as Queue<A, E>)
824 return []
825 }
826 const free = self.capacity <= 0
827 ? self.state.takers.size
828 : self.capacity - self.messages.length
829 if (free === 0) {
830 return Arr.fromIterable(messages)
831 }
832 const remaining: Array<A> = []
833 let i = 0
834 for (const message of messages) {
835 if (i < free) {
836 MutableList.append(self.messages, message)
837 } else {
838 remaining.push(message)
839 }
840 i++
841 }
842 scheduleReleaseTaker(self as Queue<A, E>)
843 return remaining
844}
845
846/**
847 * Fails the queue with an error. If the queue is already done, `false` is

Callers 1

offerAllFunction · 0.85

Calls 3

scheduleReleaseTakerFunction · 0.85
pushMethod · 0.80
takeNMethod · 0.65

Tested by

no test coverage detected