-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathfrom_chan.go
More file actions
47 lines (37 loc) · 812 Bytes
/
Copy pathfrom_chan.go
File metadata and controls
47 lines (37 loc) · 812 Bytes
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
package co
import (
syncx "go.tempura.ink/co/internal/syncx"
)
type AsyncChannel[R any] struct {
*asyncSequence[R]
sourceCh chan R
}
func FromChan[R any](ch chan R) *AsyncChannel[R] {
a := &AsyncChannel[R]{
sourceCh: ch,
}
a.asyncSequence = NewAsyncSequence[R](a)
return a
}
func (a *AsyncChannel[R]) Complete() *AsyncChannel[R] {
syncx.SafeClose(a.sourceCh)
return a
}
func (a *AsyncChannel[R]) iterator() Iterator[R] {
it := &asyncChannelIterator[R]{
AsyncChannel: a,
}
it.asyncSequenceIterator = NewAsyncSequenceIterator[R](it)
return it
}
type asyncChannelIterator[R any] struct {
*asyncSequenceIterator[R]
*AsyncChannel[R]
}
func (it *asyncChannelIterator[R]) next() *Optional[R] {
val, ok := <-it.sourceCh
if !ok {
return NewOptionalEmpty[R]()
}
return OptionalOf(val)
}