/
channel.go
71 lines (61 loc) · 1.48 KB
/
channel.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
package robin
import "sync"
// Channel is a struct that has a member variable to store subscribers
type Channel struct {
sync.Map
}
// NewChannel new a Channel instance
func NewChannel() *Channel {
c := &Channel{}
return c
}
// Subscribe to register a receiver to receive the Channel's message
func (c *Channel) Subscribe(taskFunc any, params ...any) *Subscriber {
s := &Subscriber{channel: c, receiver: newTask(taskFunc, params...)}
c.Store(s, s)
return s
}
// Publish a message to all subscribers
func (c *Channel) Publish(msg ...any) {
fiber.Enqueue(func(c *Channel, message []any) {
c.Range(func(k, v any) bool {
if s, ok := v.(*Subscriber); ok {
s.locker.Lock()
s.receiver.params(msg...)
fiber.EnqueueWithTask(s.receiver)
s.locker.Unlock()
}
return true
})
}, c, msg)
}
// Clear empty the subscribers
func (c *Channel) Clear() {
c.Range(func(k, v any) bool {
c.Delete(k)
return true
})
}
// Count returns a number that how many subscribers in the Channel.
func (c *Channel) Count() int {
count := 0
c.Range(func(k, v any) bool {
count++
return true
})
return count
}
// Unsubscribe remove the subscriber from the channel
func (c *Channel) Unsubscribe(subscriber any) {
c.Delete(subscriber)
}
// Subscriber is a struct for register to a channel
type Subscriber struct {
channel *Channel
receiver Task
locker sync.Mutex
}
// Unsubscribe remove the subscriber from the channel
func (c *Subscriber) Unsubscribe() {
c.channel.Unsubscribe(c)
}