forked from rancher/remotedialer
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathsession_manager.go
146 lines (119 loc) · 3.09 KB
/
session_manager.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
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
package remotedialer
import (
"fmt"
"math/rand"
"net"
"sync"
"time"
"github.com/gorilla/websocket"
"github.com/rancher/norman/metrics"
)
type sessionListener interface {
sessionAdded(clientKey string, sessionKey int64)
sessionRemoved(clientKey string, sessionKey int64)
}
type sessionManager struct {
sync.Mutex
clients map[string][]*Session
peers map[string][]*Session
listeners map[sessionListener]bool
}
func newSessionManager() *sessionManager {
return &sessionManager{
clients: map[string][]*Session{},
peers: map[string][]*Session{},
listeners: map[sessionListener]bool{},
}
}
func toDialer(s *Session, prefix string, deadline time.Duration) Dialer {
return func(proto, address string) (net.Conn, error) {
if prefix == "" {
return s.serverConnect(deadline, proto, address)
}
return s.serverConnect(deadline, prefix+"::"+proto, address)
}
}
func (sm *sessionManager) removeListener(listener sessionListener) {
sm.Lock()
defer sm.Unlock()
delete(sm.listeners, listener)
}
func (sm *sessionManager) addListener(listener sessionListener) {
sm.Lock()
defer sm.Unlock()
sm.listeners[listener] = true
for k, sessions := range sm.clients {
for _, session := range sessions {
listener.sessionAdded(k, session.sessionKey)
}
}
for k, sessions := range sm.peers {
for _, session := range sessions {
listener.sessionAdded(k, session.sessionKey)
}
}
}
func (sm *sessionManager) getDialer(clientKey string, deadline time.Duration) (Dialer, error) {
sm.Lock()
defer sm.Unlock()
sessions := sm.clients[clientKey]
if len(sessions) > 0 {
return toDialer(sessions[0], "", deadline), nil
}
for _, sessions := range sm.peers {
for _, session := range sessions {
session.Lock()
keys := session.remoteClientKeys[clientKey]
session.Unlock()
if len(keys) > 0 {
return toDialer(session, clientKey, deadline), nil
}
}
}
return nil, fmt.Errorf("failed to find Session for client %s", clientKey)
}
func (sm *sessionManager) add(clientKey string, conn *websocket.Conn, peer bool) *Session {
sessionKey := rand.Int63()
session := newSession(sessionKey, clientKey, conn)
sm.Lock()
defer sm.Unlock()
if peer {
sm.peers[clientKey] = append(sm.peers[clientKey], session)
} else {
sm.clients[clientKey] = append(sm.clients[clientKey], session)
}
metrics.IncSMTotalAddWS(clientKey, peer)
for l := range sm.listeners {
l.sessionAdded(clientKey, session.sessionKey)
}
return session
}
func (sm *sessionManager) remove(s *Session) {
var isPeer bool
sm.Lock()
defer sm.Unlock()
for i, store := range []map[string][]*Session{sm.clients, sm.peers} {
var newSessions []*Session
for _, v := range store[s.clientKey] {
if v.sessionKey == s.sessionKey {
if i == 0 {
isPeer = false
} else {
isPeer = true
}
metrics.IncSMTotalRemoveWS(s.clientKey, isPeer)
continue
}
newSessions = append(newSessions, v)
}
if len(newSessions) == 0 {
delete(store, s.clientKey)
} else {
store[s.clientKey] = newSessions
}
}
for l := range sm.listeners {
l.sessionRemoved(s.clientKey, s.sessionKey)
}
s.Close()
}