iterx.go (3353B)
1 package iterx 2 3 import ( 4 "iter" 5 "runtime" 6 "sync" 7 8 "git.sr.ht/~jackmordaunt/jackhammer/channel" 9 "git.sr.ht/~jackmordaunt/jackhammer/errorsx" 10 ) 11 12 // Fanout processes a sequence concurrently. 13 func Fanout[V any](seq iter.Seq[V], fn func(value V)) { 14 wg := sync.WaitGroup{} 15 defer wg.Wait() 16 17 for v := range seq { 18 wg.Add(1) 19 go func() { 20 defer wg.Done() 21 fn(v) 22 }() 23 } 24 } 25 26 // Fanout2 processes a sequence concurrently. 27 func Fanout2[K any, V any](seq iter.Seq2[K, V], fn func(key K, value V)) { 28 type pair[K any, V any] struct { 29 key K 30 value V 31 } 32 33 queue := make(chan pair[K, V]) 34 wg := sync.WaitGroup{} 35 36 for ii := 0; ii < runtime.NumCPU(); ii += 1 { 37 wg.Add(1) 38 go func() { 39 defer wg.Done() 40 for p := range queue { 41 fn(p.key, p.value) 42 } 43 }() 44 } 45 46 for k, v := range seq { 47 queue <- pair[K, V]{key: k, value: v} 48 } 49 50 close(queue) 51 wg.Wait() 52 } 53 54 // FanoutErr2 processes a sequence cuncurrently, returning an aggregate error. 55 func FanoutErr2[K any, V any](seq iter.Seq2[K, V], fn func(key K, value V) error) error { 56 errs := []error{} 57 failures := make(chan error) 58 done := make(chan any) 59 60 Fanout2(seq, func(key K, value V) { 61 if err := fn(key, value); err != nil { 62 failures <- err 63 } 64 }) 65 66 go func() { 67 errs = channel.Collect(failures) 68 close(done) 69 }() 70 71 close(failures) 72 <-done 73 74 if len(errs) > 0 { 75 return errorsx.PolyError{Errs: errs} 76 } 77 78 return nil 79 } 80 81 // For2 processes a sequence with the given function. 82 func For2[K any, V any](seq iter.Seq2[K, V], fn func(key K, value V)) { 83 for k, v := range seq { 84 fn(k, v) 85 } 86 } 87 88 // ForErr2 processes a sequence with the given function. 89 func ForErr2[K any, V any](seq iter.Seq2[K, V], fn func(key K, value V) error) error { 90 var errs []error 91 92 for k, v := range seq { 93 if err := fn(k, v); err != nil { 94 errs = append(errs, err) 95 } 96 } 97 if len(errs) > 0 { 98 return errorsx.PolyError{Errs: errs} 99 } 100 101 return nil 102 } 103 104 // Map transforms an iterator by applying a function to each element and returning the result. 105 func Map[S iter.Seq[In], In any, Out any](s S, fn func(In) Out) iter.Seq[Out] { 106 return func(yield func(Out) bool) { 107 for v := range s { 108 if !yield(fn(v)) { 109 break 110 } 111 } 112 } 113 } 114 115 // MapIndex is like [Map] but provides the iteration index to the transform function. 116 func MapIndex[S iter.Seq[In], In any, Out any](s S, fn func(int, In) Out) iter.Seq[Out] { 117 return func(yield func(Out) bool) { 118 ii := -1 119 for v := range s { 120 ii += 1 121 if !yield(fn(ii, v)) { 122 break 123 } 124 } 125 } 126 } 127 128 // Map2 transforms an iterator by applying a function to each element and returning the result. 129 func Map2[S iter.Seq2[In1, In2], In1 any, In2 any, Out any](s S, fn func(In1, In2) Out) iter.Seq[Out] { 130 return func(yield func(Out) bool) { 131 for v1, v2 := range s { 132 if !yield(fn(v1, v2)) { 133 break 134 } 135 } 136 } 137 } 138 139 // ToMap collects a sequence into a map, using [keyFor] to extract the key from each element. 140 func ToMap[S iter.Seq2[int, E], E any, K comparable](s S, keyFor func(int, E) K) map[K]E { 141 out := map[K]E{} 142 for ii, v := range s { 143 out[keyFor(ii, v)] = v 144 } 145 return out 146 } 147 148 // WithIndex adapts a single-value iterator to a two-value iterator that yeilds the index. 149 func WithIndex[T any](s iter.Seq[T]) iter.Seq2[int, T] { 150 return iter.Seq2[int, T](func(yield func(int, T) bool) { 151 i := 0 152 for v := range s { 153 if !yield(i, v) { 154 break 155 } 156 i++ 157 } 158 }) 159 }