Files
Foundry-Scrapping-API/internal/foundry/foundry.go

155 lines
4.1 KiB
Go

package foundry
import (
"errors"
"time"
"gitea.local.lab/Lbenedar/foundry_helper_service/internal/foundry/actions"
"gitea.local.lab/Lbenedar/foundry_helper_service/internal/foundry/transport"
"gitea.local.lab/Lbenedar/foundry_helper_service/internal/foundry/types"
)
var (
ListenIsDone = errors.New("Listen for websocket data in foundry is stopped")
)
type FoundryApi struct {
//TODO: make check of admin's authentication
Transport *transport.FoundryTransport
IsReadReady bool
Status types.FoundryStatus
}
func NewFoundry() *FoundryApi {
return &FoundryApi{}
}
func (foundry *FoundryApi) StartListen() error {
foundry.Transport.ReadChan = *types.InitWsChannels()
wsChannels := &foundry.Transport.ReadChan
defer foundry.Transport.ReadChan.Close()
go foundry.ListenAndServeWS()
var err error
for {
select {
// case msg, ok := <-foundry.ws.channels.Msg():
// if !ok {
// time.Sleep(5 * time.Microsecond)
// continue
// }
// var foundryState *json_model.FoundryState
// foundryState, err = json_model.ParseSetupModel(msg.ToByteSlice())
// if err != nil {
// return nil
// }
// foundry.models.FoundryState.Insert(foundryState)
case err = <-wsChannels.Err():
var foundryErr *types.FoundryError
if errors.As(err, &foundryErr) {
if foundryErr.IsFatal {
foundry.Transport.CloseWebSocketConn()
return foundryErr
}
foundry.Transport.Logger.Warn("Got error when listening or served", "err", err.Error(), "type", foundryErr.Type, "direction", foundryErr.Direction)
}
case <-wsChannels.Done():
foundry.Transport.CloseWebSocketConn()
return ListenIsDone
}
}
}
func (foundry *FoundryApi) ServeWebSocket() {
for {
select {
case message, ok := <-foundry.Transport.ReadChan.Msg():
if !ok {
foundry.Transport.Logger.Debug("Read channel is closed", "type", types.WebSocketCode)
return
}
foundry.Transport.Logger.Debug("Data has been received\n", "msg", message.Code, "type", types.WebSocketCode) //, "msg", message.MsgJson)
switch message.Code {
case types.RespPingCode, types.RespSessionDataCode:
err := foundry.Transport.SendOnlyCodeRequest(message.Code)
if err != nil {
foundry.IsReadReady = false
foundry.Transport.ReadChan.Err() <- &types.FoundryError{Direction: types.WriterCode, Err: err, Type: types.WebSocketCode}
continue
}
case types.RespServerChangeCode:
foundry.IsReadReady = true
data, err := actions.SplitToTypeAndData([]byte(message.MsgJson))
if err != nil {
foundry.Transport.ReadChan.Err() <- &types.FoundryError{Direction: types.ReaderCode, Err: err, Type: types.WebSocketCode}
continue
}
if data == nil {
continue
}
err = data.Action(foundry.Transport, &foundry.Status)
if err != nil {
foundry.Transport.ReadChan.Err() <- &types.FoundryError{Direction: types.ReaderCode, Err: err, Type: types.WebSocketCode}
continue
}
case types.RespDataCode:
foundry.Transport.ExchangeChan.Msgs[message.Id] = make(chan []byte, 1)
foundry.Transport.ExchangeChan.Msgs[message.Id] <- []byte(message.MsgJson)
go foundry.Transport.CloseMsgChannel(message.Id, 5*time.Second)
default:
}
case <-foundry.Transport.ReadChan.Done():
return
}
}
}
func (foundry *FoundryApi) ListenAndServeWS() {
go foundry.Transport.ListenWebSocket()
foundry.ServeWebSocket()
}
func (foundry *FoundryApi) HandleWSRequest(msgType string) ([]byte, error) {
msg := types.NewWsMessageByPage(msgType, foundry.Transport.CurrWsId)
return foundry.Transport.HandleWebsocketRequest(msg)
}
func (foundry *FoundryApi) ConnectToWebSocket() error {
var err error
if !foundry.Transport.HasSessionId() {
err = foundry.Transport.InitSessionId()
if err != nil {
return err
}
}
err = foundry.Transport.ConnectToFoundry()
if err != nil {
return err
}
err = foundry.Transport.InitWebSocketConnection()
if err != nil {
return err
}
status, err := foundry.Transport.Http.GetStatus()
if err != nil {
return err
}
foundry.Status = *types.NewFoundryStatus(status)
err = foundry.StartListen()
if err != nil && err != ListenIsDone {
return err
}
return nil
}