(self: Enqueue<A, E>, messages: Iterable<A>)
| 810 | * @since 4.0.0 |
| 811 | */ |
| 812 | export 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 |
no test coverage detected