Server: Sync stored packets

This commit is contained in:
Ralf
2024-09-07 19:31:27 +03:00
parent c40d93168b
commit 299e17c693
9 changed files with 297 additions and 23 deletions
+4 -3
View File
@@ -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
}
}
+1 -1
View File
@@ -14,7 +14,7 @@ import (
type Teamserver interface {
ClientConnect(username string, socket *websocket.Conn)
ClientDisconnect(username string) error
ClientDisconnect(username string)
}
type TsHttpHandler struct {
+1
View File
@@ -13,6 +13,7 @@ func NewTeamserver() *Teamserver {
ts := new(Teamserver)
ts.Profile = profile.NewProfile()
ts.clients = safe.NewMap()
ts.syncpackets = safe.NewMap()
return ts
}
+53
View File
@@ -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,
}
}
+11 -18
View File
@@ -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)
}
+110
View File
@@ -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()
}
+36 -1
View File
@@ -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"`
}
+20
View File
@@ -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
}
+61
View File
@@ -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
}