Repository navigation
Expand file tree
/
Copy pathaction.go
More file actions
109 lines (88 loc) · 1.86 KB
/
Copy pathaction.go
File metadata and controls
109 lines (88 loc) · 1.86 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
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
package co
import (
"sync"
syncx "go.tempura.ink/co/internal/syncx"
)
type Action[E any] struct {
emitChs []chan E
bufferedData *List[E]
ifWaitData bool
firstChan chan bool
lastChan chan bool
rwmux sync.RWMutex
}
func NewAction[E any]() *Action[E] {
return &Action[E]{
emitChs: make([]chan E, 0),
bufferedData: NewList[E](),
firstChan: make(chan bool),
lastChan: make(chan bool),
ifWaitData: true,
}
}
func (a *Action[E]) Iter() chan E {
ch := make(chan E)
a.emitChs = append(a.emitChs, ch)
return ch
}
func (a *Action[E]) DiscardData() *Action[E] {
a.ifWaitData = false
return a
}
func (a *Action[E]) wait() {
<-a.lastChan
}
func (a *Action[E]) GetData() []E {
a.wait()
a.rwmux.RLock()
defer a.rwmux.RUnlock()
return a.bufferedData.list
}
func (a *Action[E]) PeakData() E {
if a.bufferedData.len() == 0 {
<-a.firstChan
}
a.rwmux.RLock()
defer a.rwmux.RUnlock()
if a.bufferedData.len() == 0 {
return *new(E)
}
return a.bufferedData.getAt(0)
}
func (a *Action[E]) listen(el ...E) {
a.rwmux.RLock()
defer a.rwmux.RUnlock()
for _, e := range el {
for _, ch := range a.emitChs {
syncx.SafeNSend(ch, e)
}
}
sendFirstCh := a.bufferedData.len() == 0
if a.ifWaitData {
a.bufferedData.add(el...)
}
if sendFirstCh {
go func() {
syncx.SafeSend(a.firstChan, true)
syncx.SafeClose(a.firstChan)
}()
}
}
func (a *Action[E]) done() {
for i := range a.emitChs {
syncx.SafeClose(a.emitChs[i])
}
syncx.SafeSend(a.lastChan, true)
syncx.SafeClose(a.lastChan)
}
func MapAction[T1, T2 any](a1 *Action[T1], fn func(T1) T2) *Action[T2] {
a2 := NewAction[T2]()
a2.firstChan = a1.firstChan
a2.lastChan = a1.lastChan
a2.bufferedData = NewList[T2]()
a2.bufferedData.resizeTo(a1.bufferedData.len())
for i := range a1.bufferedData.list {
a2.bufferedData.list[i] = fn(a1.bufferedData.list[i])
}
return a2
}