-
Notifications
You must be signed in to change notification settings - Fork 88
Expand file tree
/
Copy pathlistener.go
More file actions
121 lines (101 loc) · 2.63 KB
/
Copy pathlistener.go
File metadata and controls
121 lines (101 loc) · 2.63 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
110
111
112
113
114
115
116
117
118
119
120
121
package relayer
import "github.com/nbd-wtf/go-nostr"
type Listener struct {
filters nostr.Filters
}
func GetListeningFilters() nostr.Filters {
serversMutex.RLock()
defer serversMutex.RUnlock()
respfilters := make(nostr.Filters, 0, len(servers)*2)
for srv := range servers {
respfilters = appendDistinctFilters(respfilters, srv.GetListeningFilters())
}
return respfilters
}
func (s *Server) GetListeningFilters() nostr.Filters {
s.listenersMu.RLock()
defer s.listenersMu.RUnlock()
respfilters := make(nostr.Filters, 0, len(s.listeners)*2)
for _, connlisteners := range s.listeners {
for _, listener := range connlisteners {
respfilters = appendDistinctFilters(respfilters, listener.filters)
}
}
return respfilters
}
func appendDistinctFilters(dst nostr.Filters, src nostr.Filters) nostr.Filters {
for _, listenerfilter := range src {
duplicate := false
for _, respfilter := range dst {
if nostr.FilterEqual(listenerfilter, respfilter) {
duplicate = true
break
}
}
if !duplicate {
dst = append(dst, listenerfilter)
}
}
return dst
}
func (s *Server) setListener(id string, ws *WebSocket, filters nostr.Filters) {
s.listenersMu.Lock()
defer s.listenersMu.Unlock()
subs, ok := s.listeners[ws]
if !ok {
subs = make(map[string]*Listener)
s.listeners[ws] = subs
}
subs[id] = &Listener{filters: filters}
}
// Remove a specific subscription id from listeners for a given ws client
func (s *Server) removeListenerId(ws *WebSocket, id string) {
s.listenersMu.Lock()
defer s.listenersMu.Unlock()
if subs, ok := s.listeners[ws]; ok {
delete(s.listeners[ws], id)
if len(subs) == 0 {
delete(s.listeners, ws)
}
}
}
// Remove WebSocket conn from listeners
func (s *Server) removeListener(ws *WebSocket) {
s.listenersMu.Lock()
defer s.listenersMu.Unlock()
clear(s.listeners[ws])
delete(s.listeners, ws)
}
type listenerDelivery struct {
ws *WebSocket
subID string
event nostr.Event
}
func (s *Server) notifyListeners(event *nostr.Event) {
s.listenersMu.RLock()
deliveries := make([]listenerDelivery, 0, len(s.listeners))
for ws, subs := range s.listeners {
for id, listener := range subs {
if !listener.filters.Match(event) {
continue
}
deliveries = append(deliveries, listenerDelivery{
ws: ws,
subID: id,
event: *event,
})
}
}
s.listenersMu.RUnlock()
for _, delivery := range deliveries {
delivery := delivery
delivery.ws.WriteJSON(nostr.EventEnvelope{SubscriptionID: &delivery.subID, Event: delivery.event})
}
}
func BroadcastEvent(evt *nostr.Event) {
serversMutex.RLock()
defer serversMutex.RUnlock()
for srv := range servers {
srv.notifyListeners(evt)
}
}