-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathasync_merged_sequence.go
More file actions
53 lines (45 loc) · 1.09 KB
/
Copy pathasync_merged_sequence.go
File metadata and controls
53 lines (45 loc) · 1.09 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
package co
type AsyncMergedSequence[R any] struct {
*asyncSequence[R]
aSequenceables []AsyncSequenceable[R]
}
func NewAsyncMergedSequence[R any](as ...AsyncSequenceable[R]) *AsyncMergedSequence[R] {
a := &AsyncMergedSequence[R]{
aSequenceables: as,
}
a.asyncSequence = NewAsyncSequence[R](a)
return a
}
func (a *AsyncMergedSequence[R]) iterator() Iterator[R] {
it := &asyncMergedSequenceIterator[R]{}
it.asyncSequenceIterator = NewAsyncSequenceIterator[R](it)
for i := range a.aSequenceables {
it.its = append(it.its, a.aSequenceables[i].iterator())
}
return it
}
type asyncMergedSequenceIterator[R any] struct {
*asyncSequenceIterator[R]
its []Iterator[R]
currentIndex int
}
func (it *asyncMergedSequenceIterator[R]) nextIndex() int {
defer func() {
if it.currentIndex+1 >= len(it.its) {
it.currentIndex = 0
} else {
it.currentIndex++
}
}()
return it.currentIndex
}
func (it *asyncMergedSequenceIterator[R]) next() *Optional[R] {
for range it.its {
op := it.its[it.nextIndex()].next()
if !op.valid {
continue
}
return op
}
return NewOptionalEmpty[R]()
}