diff --git a/AdaptixServer/core/httphandler/auth.go b/AdaptixServer/core/httphandler/auth.go index f5659278..f353c635 100644 --- a/AdaptixServer/core/httphandler/auth.go +++ b/AdaptixServer/core/httphandler/auth.go @@ -72,8 +72,6 @@ func (th *TsHttpHandler) connect(ctx *gin.Context) { return } - //defer wsConn.Close() - go th.handleWsConnect(wsConn) } @@ -91,7 +89,8 @@ func (th *TsHttpHandler) handleWsConnect(wsConn *websocket.Conn) { return } - if err = json.Unmarshal(body, &structToken); err != nil { + err = json.Unmarshal(body, &structToken) + if err != nil { logs.Error("JSON Unmarshal error: " + err.Error()) return } @@ -111,6 +110,8 @@ func (th *TsHttpHandler) handleWsConnect(wsConn *websocket.Conn) { logs.Debug("User '%s' disconnected: %s\n", username, err.Error()) + th.teamserver.ClientDisconnect(username) + break } } diff --git a/AdaptixServer/core/httphandler/teamserver.go b/AdaptixServer/core/httphandler/teamserver.go index f1ad43f8..5e15df81 100644 --- a/AdaptixServer/core/httphandler/teamserver.go +++ b/AdaptixServer/core/httphandler/teamserver.go @@ -14,7 +14,7 @@ import ( type Teamserver interface { ClientConnect(username string, socket *websocket.Conn) - ClientDisconnect(username string) error + ClientDisconnect(username string) } type TsHttpHandler struct { diff --git a/AdaptixServer/core/server/server.go b/AdaptixServer/core/server/server.go index c552d7c0..cebd8918 100644 --- a/AdaptixServer/core/server/server.go +++ b/AdaptixServer/core/server/server.go @@ -13,6 +13,7 @@ func NewTeamserver() *Teamserver { ts := new(Teamserver) ts.Profile = profile.NewProfile() ts.clients = safe.NewMap() + ts.syncpackets = safe.NewMap() return ts } diff --git a/AdaptixServer/core/server/syncpacket.go b/AdaptixServer/core/server/syncpacket.go new file mode 100644 index 00000000..8519d15a --- /dev/null +++ b/AdaptixServer/core/server/syncpacket.go @@ -0,0 +1,53 @@ +package server + +import "time" + +const ( + storeLog = "log" + + typeSync = "sync" + typeSyncStart = "start" + typeSyncFinish = "finish" + + typeClient = "client" + typeClientConnect = "connect" + typeClientDisconnect = "disconnect" +) + +func CreateSpSyncStart(count int) SyncPackerStart { + return SyncPackerStart{ + SpType: typeSync, + SpSubType: typeSyncStart, + + Count: count, + } +} + +func CreateSpSyncFinish() SyncPackerFinish { + return SyncPackerFinish{ + SpType: typeSync, + SpSubType: typeSyncFinish, + } +} + +func CreateSpClientConnect(username string) SyncPackerClientConnect { + return SyncPackerClientConnect{ + store: storeLog, + SpCreateTime: time.Now().UTC().Unix(), + SpType: typeClient, + SpSubType: typeClientConnect, + + Username: username, + } +} + +func CreateSpClientDisconnect(username string) SyncPackerClientDisconnect { + return SyncPackerClientDisconnect{ + store: storeLog, + SpCreateTime: time.Now().UTC().Unix(), + SpType: typeClient, + SpSubType: typeClientDisconnect, + + Username: username, + } +} diff --git a/AdaptixServer/core/server/ts_clients.go b/AdaptixServer/core/server/ts_clients.go index c432c118..d89318fa 100644 --- a/AdaptixServer/core/server/ts_clients.go +++ b/AdaptixServer/core/server/ts_clients.go @@ -1,7 +1,6 @@ package server import ( - "errors" "github.com/gorilla/websocket" ) @@ -9,29 +8,23 @@ func (ts *Teamserver) ClientConnect(username string, socket *websocket.Conn) { ts.clients.Put(username, socket) - //TODO: Sync Data to Client + ts.SyncStored(socket) - //TODO: Broadcast LogEvent - UserConnected + packet := CreateSpClientConnect(username) + ts.SyncSavePacket(packet.store, packet) + ts.SyncAllClients(packet) } -func (ts *Teamserver) ClientDisconnect(username string) error { - var ( - val any - ok bool - wsConn *websocket.Conn - err error - ) +func (ts *Teamserver) ClientDisconnect(username string) { - val, ok = ts.clients.GetDelete(username) + value, ok := ts.clients.GetDelete(username) if !ok { - return errors.New("user not found") + return } - wsConn = val.(*websocket.Conn) + value.(*websocket.Conn).Close() - err = wsConn.Close() - - //TODO: Broadcast LogEvent - UserDisconnected - - return err + packet := CreateSpClientDisconnect(username) + ts.SyncSavePacket(packet.store, packet) + ts.SyncAllClients(packet) } diff --git a/AdaptixServer/core/server/ts_sync.go b/AdaptixServer/core/server/ts_sync.go new file mode 100644 index 00000000..c5104438 --- /dev/null +++ b/AdaptixServer/core/server/ts_sync.go @@ -0,0 +1,110 @@ +package server + +import ( + "AdaptixServer/core/utils/logs" + "AdaptixServer/core/utils/safe" + "bytes" + "encoding/json" + "github.com/gorilla/websocket" +) + +func (ts *Teamserver) SyncClient(username string, packet interface{}) { + var ( + buffer bytes.Buffer + err error + clientWS *websocket.Conn + ) + + err = json.NewEncoder(&buffer).Encode(packet) + if err != nil { + return + } + + logs.Debug("SyncClient to '%v': %v\n", username, buffer.String()) + + value, found := ts.clients.Get(username) + if found { + clientWS = value.(*websocket.Conn) + err = clientWS.WriteMessage(websocket.BinaryMessage, buffer.Bytes()) + if err != nil { + return + } + } +} + +func (ts *Teamserver) SyncAllClients(packet interface{}) { + var ( + buffer bytes.Buffer + err error + ) + + err = json.NewEncoder(&buffer).Encode(packet) + if err != nil { + return + } + + logs.Debug("SyncAllClients: %v\n", buffer.String()) + + ts.clients.ForEach(func(key string, value interface{}) { + clientWS := value.(*websocket.Conn) + _ = clientWS.WriteMessage(websocket.BinaryMessage, buffer.Bytes()) + }) +} + +func (ts *Teamserver) SyncSavePacket(store string, packet interface{}) { + if ts.syncpackets.Contains(store) == false { + ts.syncpackets.Put(store, safe.NewSlice()) + } + + value, found := ts.syncpackets.Get(store) + if found { + value.(*safe.Slice).Put(packet) + } +} + +func (ts *Teamserver) SyncStored(clientWS *websocket.Conn) { + var ( + buffer bytes.Buffer + packet interface{} + ) + + ts.syncpackets.DirectLock() + mapPackets := ts.syncpackets.DirectMap() + + sumLen := 0 + for _, value := range mapPackets { + sumLen += value.(*safe.Slice).Len() + } + + if sumLen != 0 { + packet = CreateSpSyncStart(sumLen) + _ = json.NewEncoder(&buffer).Encode(packet) + _ = clientWS.WriteMessage(websocket.BinaryMessage, buffer.Bytes()) + buffer.Reset() + + sliceLog := mapPackets[storeLog].(*safe.Slice) + for value := range sliceLog.Iterator() { + _ = json.NewEncoder(&buffer).Encode(value) + _ = clientWS.WriteMessage(websocket.BinaryMessage, buffer.Bytes()) + buffer.Reset() + } + + for mKey, mValue := range mapPackets { + if mKey != storeLog { + sliceVal := mValue.(*safe.Slice) + for value := range sliceVal.Iterator() { + _ = json.NewEncoder(&buffer).Encode(value) + _ = clientWS.WriteMessage(websocket.BinaryMessage, buffer.Bytes()) + buffer.Reset() + } + } + } + + packet = CreateSpSyncFinish() + _ = json.NewEncoder(&buffer).Encode(packet) + _ = clientWS.WriteMessage(websocket.BinaryMessage, buffer.Bytes()) + buffer.Reset() + } + + ts.syncpackets.DirectUnlock() +} diff --git a/AdaptixServer/core/server/utils.go b/AdaptixServer/core/server/utils.go index 1247394f..ed8ad764 100644 --- a/AdaptixServer/core/server/utils.go +++ b/AdaptixServer/core/server/utils.go @@ -6,9 +6,44 @@ import ( "AdaptixServer/core/utils/safe" ) +// TeamServer + type Teamserver struct { Profile *profile.AdaptixProfile AdaptixServer *httphandler.TsHttpHandler - clients safe.Map + clients safe.Map + syncpackets safe.Map +} + +// SyncPacket + +type SyncPackerStart struct { + SpType string `json:"type"` + SpSubType string `json:"subtype"` + + Count int `json:"count"` +} + +type SyncPackerFinish struct { + SpType string `json:"type"` + SpSubType string `json:"subtype"` +} + +type SyncPackerClientConnect struct { + store string + SpCreateTime int64 `json:"time"` + SpType string `json:"type"` + SpSubType string `json:"subtype"` + + Username string `json:"username"` +} + +type SyncPackerClientDisconnect struct { + store string + SpCreateTime int64 `json:"time"` + SpType string `json:"type"` + SpSubType string `json:"subtype"` + + Username string `json:"username"` } diff --git a/AdaptixServer/core/utils/safe/map.go b/AdaptixServer/core/utils/safe/map.go index 9b1fba7c..8ed3f2a9 100644 --- a/AdaptixServer/core/utils/safe/map.go +++ b/AdaptixServer/core/utils/safe/map.go @@ -56,3 +56,23 @@ func (s *Map) Len() int { defer s.mutex.RUnlock() return len(s.m) } + +func (s *Map) ForEach(f func(key string, value interface{})) { + s.mutex.RLock() + defer s.mutex.RUnlock() + for key, value := range s.m { + f(key, value) + } +} + +func (s *Map) DirectLock() { + s.mutex.RLock() +} + +func (s *Map) DirectUnlock() { + s.mutex.RUnlock() +} + +func (s *Map) DirectMap() map[string]interface{} { + return s.m +} diff --git a/AdaptixServer/core/utils/safe/slice.go b/AdaptixServer/core/utils/safe/slice.go new file mode 100644 index 00000000..9244cfd5 --- /dev/null +++ b/AdaptixServer/core/utils/safe/slice.go @@ -0,0 +1,61 @@ +package safe + +import ( + "sync" +) + +type Slice struct { + mutex sync.RWMutex + items []interface{} +} + +func NewSlice() *Slice { + return &Slice{} +} + +func (sl *Slice) Put(value interface{}) { + sl.mutex.Lock() + defer sl.mutex.Unlock() + + sl.items = append(sl.items, value) +} + +func (sl *Slice) Get(index int) (interface{}, bool) { + sl.mutex.RLock() + defer sl.mutex.RUnlock() + + if index < 0 || index >= len(sl.items) { + return nil, false + } + return sl.items[index], true +} + +func (sl *Slice) Delete(index int) { + sl.mutex.Lock() + defer sl.mutex.Unlock() + + if index < 0 || index >= len(sl.items) { + return + } + sl.items = append(sl.items[:index], sl.items[index+1:]...) +} + +func (sl *Slice) Len() int { + sl.mutex.RLock() + defer sl.mutex.RUnlock() + return len(sl.items) +} + +func (sl *Slice) Iterator() <-chan interface{} { + sl.mutex.RLock() + defer sl.mutex.RUnlock() + + ch := make(chan interface{}) + go func() { + defer close(ch) + for _, item := range sl.items { + ch <- item + } + }() + return ch +}