-
Notifications
You must be signed in to change notification settings - Fork 0
/
bytes_stream_test.go
68 lines (60 loc) · 1.18 KB
/
bytes_stream_test.go
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
package stream
import (
"log"
"strconv"
"testing"
)
func TestBytesStream(t *testing.T) {
obv := NewObserver(nil)
stream := NewBytesStream(obv)
stream.Target = func() {
for i := 0; i <= 10; i++ {
stream.Send([]byte{byte(i)})
}
stream.OnComplete()
}
stream.Publish(nil).Subscribe(func(data []byte) {
log.Print(data)
})
}
func TestCancel(t *testing.T) {
handler := &ObservHandler{
AtCancel: func() {},
AtComplete: func() {},
AtError: func(error) {},
}
obv := NewObserver(handler)
stream := NewBytesStream(obv)
stream.Target = func() {
for i := 0; i <= 100; i++ {
select {
case <-stream.AfterCancel():
return
default:
stream.Send([]byte(strconv.Itoa(i)))
}
}
stream.OnComplete()
}
stream.Cancel()
stream.Publish(nil).Subscribe(func(data []byte) {
t.Fail()
})
}
func TestChunk(t *testing.T) {
stream := NewBytesStream(NewObserver(nil))
handler := func() {
for i := 0; i <= 100; i++ {
select {
case <-stream.AfterCancel():
return
default:
stream.Send([]byte(strconv.Itoa(i)))
}
}
stream.OnComplete()
}
stream.Publish(handler).Chunk(5).Subscribe(func(data []byte) {
log.Printf("received data: %s", data)
})
}