-
Notifications
You must be signed in to change notification settings - Fork 0
/
observable.go
101 lines (90 loc) · 2.26 KB
/
observable.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
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
package observable
import (
"errors"
"io"
"os"
"reflect"
"strings"
"sync"
)
/*
*/
type observableInterface interface {
Next(interface{})
Subscribe(func(*chan interface{}))
}
type observable struct {
buffer chan interface{}
subscribers sync.Map
i *int
}
// Observable provides methods and hooks for transmitting data.
type Observable struct {
observable
}
// New returns a pointer to a new Observable
func New() *Observable {
var i int
return &Observable{observable: observable{i: &i}}
}
// Next - pass a new value to be broadcasted.
func (o *Observable) Next(sl ...interface{}) {
for k := range sl {
/*o.subscribers.Range(func(key interface{}, value interface{}) bool {
*value.(*chan interface{}) <- sl[k]
return true
})*/
for i := 0; i < *o.i; i++ {
if ch, ok := o.subscribers.Load(i); ok {
*ch.(*chan interface{}) <- sl[k]
}
}
}
}
func (o *Observable) allocate(id int) {
io.Copy(os.Stderr, strings.NewReader("Unimplemented"))
}
// Subscribe with your own channel to the events of the observable. The function returns an id that can be used to unsubscribe.
func (o *Observable) Subscribe(ch *chan interface{}) (int, error) {
if ch == nil {
return 0, errors.New("Argument to subscribe is nil not chan")
}
k := *o.i
*o.i++
o.subscribers.Store(k, ch)
return k, nil
}
// On provides a hook to execute a callback when a certain value is passed through the channel.
func (o *Observable) On(value interface{}, ch *chan interface{}, fn func()) {
for true {
select {
case v := <-*ch:
if reflect.DeepEqual(v, value) {
fn()
}
}
}
}
// Once provides a hook to execute a callback when a certain value is passed through the channel.
func (o *Observable) Once(value interface{}, ch *chan interface{}, fn func()) {
select {
case v := <-*ch:
if reflect.DeepEqual(v, value) {
fn()
}
}
}
// Close effectivelly calls close() on all channels used by subscribers to this observable.
func (o *Observable) Close() {
for i := 0; i < *o.i; i++ {
ch, ok := o.subscribers.Load(i)
if ok {
close(*ch.(*chan interface{}))
}
}
}
// Unsubscribe removes the channel from the list of active subscribers, does not call close on them.
func (o *Observable) Unsubscribe(id int) (ok bool) {
_, ok = o.subscribers.LoadAndDelete(id)
return
}