/
subscriber.go
83 lines (69 loc) · 1.55 KB
/
subscriber.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
package guppeteer
import (
"errors"
"sync"
"sync/atomic"
"github.com/jacexh/guppeteer/cdp/target"
)
type (
Subscriber struct {
receivers sync.Map
wg sync.WaitGroup
}
Receiver struct {
received int32
event string
f func([]byte)
}
eventloop struct {
sessions sync.Map
}
)
var (
defaultEventloop = &eventloop{}
)
func (el *eventloop) Register(sid target.SessionID, sub *Subscriber) {
el.sessions.Store(sid, sub)
}
func (el *eventloop) Cancel(sid target.SessionID) {
el.sessions.Delete(sid)
}
func (el *eventloop) Handle(sid target.SessionID, event string, d []byte) {
if val, loaded := el.sessions.Load(sid); loaded {
sub := val.(*Subscriber)
sub.Handle(event, d)
}
}
func (sub *Subscriber) Subscribe(event string, f func([]byte)) {
sub.wg.Add(1)
sub.receivers.Store(event, NewReceiver(event, f))
}
func (sub *Subscriber) Handle(event string, d []byte) {
if val, ok := sub.receivers.Load(event); ok {
go func(r *Receiver) {
defer sub.wg.Done()
r.Receive(d)
}(val.(*Receiver))
}
}
func (sub *Subscriber) WaitUtilPublished() map[string]*Receiver {
sub.wg.Wait()
ret := map[string]*Receiver{}
sub.receivers.Range(func(key, value interface{}) bool {
ret[key.(string)] = value.(*Receiver)
return true
})
return ret
}
func NewReceiver(event string, f func([]byte)) *Receiver {
return &Receiver{event: event, f: f}
}
func (rc *Receiver) Receive(d []byte) error {
if atomic.CompareAndSwapInt32(&rc.received, 0, 1) {
if rc.f != nil {
rc.f(d)
}
return nil
}
return errors.New("received message")
}