add graceful shutdown, add reconnect on connection loss
This commit is contained in:
2
Makefile
2
Makefile
@@ -1,7 +1,7 @@
|
|||||||
-include .env
|
-include .env
|
||||||
|
|
||||||
api/run:
|
api/run:
|
||||||
go run ./cmd/api -foundry_pass=${FOUNDRY_PASS} -foundry_host=${FOUNDRY_HOST} -port=${LISTEN_PORT} -log_level=${LOG_LEVEL} -db-dsn=${DB_DSN}
|
go run ./cmd/api -foundry_pass=${FOUNDRY_PASS} -foundry_host=${FOUNDRY_HOST} -port=${LISTEN_PORT} -env=development -log_level=${LOG_LEVEL} -db-dsn=${DB_DSN}
|
||||||
|
|
||||||
db/migration/new:
|
db/migration/new:
|
||||||
@echo 'Creating migration files for ${name}...'
|
@echo 'Creating migration files for ${name}...'
|
||||||
|
|||||||
@@ -6,7 +6,6 @@ import (
|
|||||||
"flag"
|
"flag"
|
||||||
"log"
|
"log"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
|
||||||
"os"
|
"os"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -28,6 +27,7 @@ const (
|
|||||||
type config struct {
|
type config struct {
|
||||||
mode service_mode
|
mode service_mode
|
||||||
port string
|
port string
|
||||||
|
env string
|
||||||
|
|
||||||
db struct {
|
db struct {
|
||||||
dsn string
|
dsn string
|
||||||
@@ -72,6 +72,7 @@ func openDB(cfg config) (*sql.DB, error) {
|
|||||||
// TODO: change default values
|
// TODO: change default values
|
||||||
func (app *application) parseFlags() *requests.FoundryHttpRequest {
|
func (app *application) parseFlags() *requests.FoundryHttpRequest {
|
||||||
flag.StringVar(&app.cfg.port, "port", ":9090", "Service port")
|
flag.StringVar(&app.cfg.port, "port", ":9090", "Service port")
|
||||||
|
flag.StringVar(&app.cfg.env, "env", "development", "Enivornment (development|staging|production)")
|
||||||
|
|
||||||
flag.StringVar(&app.cfg.db.dsn, "db-dsn", "file:db/default.db?cache=shared", "PostgreSQL DSN")
|
flag.StringVar(&app.cfg.db.dsn, "db-dsn", "file:db/default.db?cache=shared", "PostgreSQL DSN")
|
||||||
flag.IntVar(&app.cfg.db.maxOpenConns, "db-max-open-conns", 25, "PostgreSQL max open connections")
|
flag.IntVar(&app.cfg.db.maxOpenConns, "db-max-open-conns", 25, "PostgreSQL max open connections")
|
||||||
@@ -91,7 +92,6 @@ func (app *application) parseFlags() *requests.FoundryHttpRequest {
|
|||||||
flag.Parse()
|
flag.Parse()
|
||||||
|
|
||||||
app.cfg.mode = service_mode(mode)
|
app.cfg.mode = service_mode(mode)
|
||||||
// app.foundryApp.SetConfig(&foundryConfig)
|
|
||||||
|
|
||||||
logLevel, ok := StringToLogLevel[logLevelStr]
|
logLevel, ok := StringToLogLevel[logLevelStr]
|
||||||
if !ok {
|
if !ok {
|
||||||
@@ -105,19 +105,7 @@ func (app *application) parseFlags() *requests.FoundryHttpRequest {
|
|||||||
return &foundryHttpData
|
return &foundryHttpData
|
||||||
}
|
}
|
||||||
|
|
||||||
// TODO: Handle application exit if wrong data
|
// TODO: Handle grace application exit
|
||||||
func (app *application) StartListenFoundry() {
|
|
||||||
var err error
|
|
||||||
|
|
||||||
for {
|
|
||||||
err = app.foundryApp.ConnectToWebSocket()
|
|
||||||
if err != nil {
|
|
||||||
app.slogger.Error("Error raised", "err", err.Error())
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
app := application{foundryApp: &foundry.FoundryApi{}}
|
app := application{foundryApp: &foundry.FoundryApi{}}
|
||||||
|
|
||||||
@@ -133,29 +121,8 @@ func main() {
|
|||||||
defer db.Close()
|
defer db.Close()
|
||||||
app.slogger.Info("Database connection pool established")
|
app.slogger.Info("Database connection pool established")
|
||||||
|
|
||||||
app.foundryApp.Transport = transport.NewFoundryTransport(db, app.slogger, foundryHttpData)
|
app.foundryApp.SetTransport(transport.NewFoundryTransport(db, app.slogger, foundryHttpData))
|
||||||
go app.StartListenFoundry()
|
go app.foundryApp.StartListenFoundry()
|
||||||
|
|
||||||
errLog := slog.NewLogLogger(app.slogger.Handler(), slog.LevelError)
|
app.serve()
|
||||||
server := &http.Server{
|
|
||||||
Addr: app.cfg.port,
|
|
||||||
Handler: app.routes(),
|
|
||||||
ErrorLog: errLog,
|
|
||||||
// TLSConfig: tlsConfig,
|
|
||||||
IdleTimeout: time.Minute,
|
|
||||||
ReadTimeout: 5 * time.Second,
|
|
||||||
WriteTimeout: 10 * time.Second,
|
|
||||||
}
|
|
||||||
|
|
||||||
switch app.cfg.mode {
|
|
||||||
case API_MODE:
|
|
||||||
app.slogger.Info("Server is started", "port", app.cfg.port)
|
|
||||||
server.ListenAndServe()
|
|
||||||
case DISCORD_BOT_MODE:
|
|
||||||
app.slogger.Warn("TODO: implement this mode")
|
|
||||||
case TG_BOT_MODE:
|
|
||||||
app.slogger.Warn("TODO: implement this mode")
|
|
||||||
default:
|
|
||||||
app.slogger.Error("Wrong Mode. Should be one of: 'api', 'discord', 'tg'", "mode", app.cfg.mode)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
70
cmd/api/server.go
Normal file
70
cmd/api/server.go
Normal file
@@ -0,0 +1,70 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"log/slog"
|
||||||
|
"net/http"
|
||||||
|
"os"
|
||||||
|
"os/signal"
|
||||||
|
"syscall"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (app *application) serve() error {
|
||||||
|
srv := &http.Server{
|
||||||
|
Addr: app.cfg.port,
|
||||||
|
Handler: app.routes(),
|
||||||
|
ErrorLog: slog.NewLogLogger(app.slogger.Handler(), slog.LevelError),
|
||||||
|
// TLSConfig: tlsConfig,
|
||||||
|
IdleTimeout: time.Minute,
|
||||||
|
ReadTimeout: 5 * time.Second,
|
||||||
|
WriteTimeout: 10 * time.Second,
|
||||||
|
}
|
||||||
|
|
||||||
|
shutdownError := make(chan error)
|
||||||
|
go func() {
|
||||||
|
quit := make(chan os.Signal, 1)
|
||||||
|
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
|
||||||
|
s := <-quit
|
||||||
|
|
||||||
|
app.slogger.Info("Caught signal", "signal", s.String())
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
err := srv.Shutdown(ctx)
|
||||||
|
if err != nil {
|
||||||
|
shutdownError <- err
|
||||||
|
}
|
||||||
|
|
||||||
|
err = app.foundryApp.Shutdown()
|
||||||
|
if err != nil {
|
||||||
|
shutdownError <- err
|
||||||
|
}
|
||||||
|
|
||||||
|
shutdownError <- nil
|
||||||
|
}()
|
||||||
|
|
||||||
|
switch app.cfg.mode {
|
||||||
|
case API_MODE:
|
||||||
|
app.slogger.Info("Starting server", "addr", srv.Addr, "env", app.cfg.env)
|
||||||
|
err := srv.ListenAndServe()
|
||||||
|
if !errors.Is(err, http.ErrServerClosed) {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
app.slogger.Info("Stopped server", "addr", srv.Addr)
|
||||||
|
err = <-shutdownError
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
case DISCORD_BOT_MODE:
|
||||||
|
app.slogger.Warn("TODO: implement this mode")
|
||||||
|
case TG_BOT_MODE:
|
||||||
|
app.slogger.Warn("TODO: implement this mode")
|
||||||
|
default:
|
||||||
|
app.slogger.Error("Wrong Mode. Should be one of: 'api', 'discord', 'tg'", "mode", app.cfg.mode)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -1,154 +1,200 @@
|
|||||||
package foundry
|
package foundry
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"log/slog"
|
||||||
|
"sync"
|
||||||
"time"
|
"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/transport"
|
||||||
"gitea.local.lab/Lbenedar/foundry_helper_service/internal/foundry/types"
|
"gitea.local.lab/Lbenedar/foundry_helper_service/internal/foundry/types"
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
ListenIsDone = errors.New("Listen for websocket data in foundry is stopped")
|
ListenIsDone = errors.New("Listen for websocket data in foundry is stopped")
|
||||||
|
ChannelIsClosed = errors.New("Channel is closed")
|
||||||
|
|
||||||
|
CloseTimeoutExceed = errors.New("Timeout of websocket close is exceed")
|
||||||
)
|
)
|
||||||
|
|
||||||
type FoundryApi struct {
|
type FoundryApi struct {
|
||||||
//TODO: make check of admin's authentication
|
//TODO: make check of admin's authentication
|
||||||
Transport *transport.FoundryTransport
|
transport *transport.FoundryTransport
|
||||||
IsReadReady bool
|
IsAvailable bool
|
||||||
|
|
||||||
|
Logger *slog.Logger
|
||||||
Status types.FoundryStatus
|
Status types.FoundryStatus
|
||||||
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewFoundry() *FoundryApi {
|
func NewFoundry() *FoundryApi {
|
||||||
return &FoundryApi{}
|
return &FoundryApi{}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (foundry *FoundryApi) StartListen() error {
|
func (foundry *FoundryApi) SetTransport(tr *transport.FoundryTransport) {
|
||||||
foundry.Transport.ReadChan = *types.InitWsChannels()
|
foundry.transport = tr
|
||||||
|
}
|
||||||
wsChannels := &foundry.Transport.ReadChan
|
|
||||||
defer foundry.Transport.ReadChan.Close()
|
|
||||||
go foundry.ListenAndServeWS()
|
|
||||||
|
|
||||||
|
func (foundry *FoundryApi) ListenAndServeWS() error {
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
|
foundry.transport.ReadChan = *types.InitWsChannels()
|
||||||
|
defer foundry.transport.ReadChan.Close()
|
||||||
|
|
||||||
|
foundry.background(foundry.transport.ListenWebSocket)
|
||||||
|
|
||||||
|
err = foundry.ServeWebSocket()
|
||||||
|
if err != nil {
|
||||||
|
foundry.IsAvailable = false
|
||||||
|
}
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (foundry *FoundryApi) ServeWebSocket() error {
|
||||||
|
var err error
|
||||||
|
wsChannels := &foundry.transport.ReadChan
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
// case msg, ok := <-foundry.ws.channels.Msg():
|
case message, ok := <-wsChannels.Msg():
|
||||||
// if !ok {
|
if !ok {
|
||||||
// time.Sleep(5 * time.Microsecond)
|
foundry.transport.Logger.Debug("Read channel is closed", "type", types.WebSocketCode)
|
||||||
// continue
|
return ChannelIsClosed
|
||||||
// }
|
}
|
||||||
// var foundryState *json_model.FoundryState
|
foundry.transport.Logger.Debug("Data has been received\n", "msg", message.Code, "type", types.WebSocketCode) //, "msg", message.MsgJson)
|
||||||
// foundryState, err = json_model.ParseSetupModel(msg.ToByteSlice())
|
|
||||||
// if err != nil {
|
switch message.Code {
|
||||||
// return nil
|
case types.RespPingCode, types.RespSessionDataCode:
|
||||||
// }
|
err := foundry.transport.SendOnlyCodeRequest(message.Code)
|
||||||
// foundry.models.FoundryState.Insert(foundryState)
|
if err != nil {
|
||||||
|
wsChannels.Err() <- &types.FoundryError{Direction: types.WriterCode, Err: err, Type: types.WebSocketCode}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
case types.RespServerChangeCode:
|
||||||
|
foundry.IsAvailable = true
|
||||||
|
// data, err := actions.SplitToTypeAndData([]byte(message.MsgJson))
|
||||||
|
// if err != nil {
|
||||||
|
// wsChannels.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 {
|
||||||
|
// wsChannels.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 err = <-wsChannels.Err():
|
case err = <-wsChannels.Err():
|
||||||
var foundryErr *types.FoundryError
|
var foundryErr *types.FoundryError
|
||||||
if errors.As(err, &foundryErr) {
|
if errors.As(err, &foundryErr) {
|
||||||
if foundryErr.IsFatal {
|
if foundryErr.IsFatal {
|
||||||
foundry.Transport.CloseWebSocketConn()
|
foundry.transport.CloseWebSocketConn()
|
||||||
return foundryErr
|
return foundryErr
|
||||||
}
|
}
|
||||||
foundry.Transport.Logger.Warn("Got error when listening or served", "err", err.Error(), "type", foundryErr.Type, "direction", foundryErr.Direction)
|
foundry.transport.Logger.Warn("Got error when listening or served", "err", err.Error(), "type", foundryErr.Type, "direction", foundryErr.Direction)
|
||||||
}
|
}
|
||||||
case <-wsChannels.Done():
|
case <-wsChannels.Done():
|
||||||
foundry.Transport.CloseWebSocketConn()
|
foundry.transport.CloseWebSocketConn()
|
||||||
return ListenIsDone
|
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) {
|
func (foundry *FoundryApi) HandleWSRequest(msgType string) ([]byte, error) {
|
||||||
msg := types.NewWsMessageByPage(msgType, foundry.Transport.CurrWsId)
|
msg := types.NewWsMessageByPage(msgType, foundry.transport.CurrWsId)
|
||||||
|
|
||||||
return foundry.Transport.HandleWebsocketRequest(msg)
|
return foundry.transport.HandleWebsocketRequest(msg)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (foundry *FoundryApi) ConnectToWebSocket() error {
|
func (foundry *FoundryApi) Shutdown() error {
|
||||||
var err error
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
||||||
if !foundry.Transport.HasSessionId() {
|
defer cancel()
|
||||||
err = foundry.Transport.InitSessionId()
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
err = foundry.Transport.ConnectToFoundry()
|
foundry.transport.Logger.Info("Completing foundry background tasks")
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
err = foundry.Transport.InitWebSocketConnection()
|
foundryClosed := make(chan struct{})
|
||||||
if err != nil {
|
go func() {
|
||||||
return err
|
types.CloseChannel(foundry.transport.ReadChan.Done())
|
||||||
}
|
|
||||||
|
|
||||||
status, err := foundry.Transport.Http.GetStatus()
|
foundry.wg.Wait()
|
||||||
if err != nil {
|
foundry.transport.ExchangeChan.Close()
|
||||||
return err
|
foundry.transport.ReadChan.Close()
|
||||||
}
|
|
||||||
foundry.Status = *types.NewFoundryStatus(status)
|
|
||||||
|
|
||||||
err = foundry.StartListen()
|
types.CloseChannel(foundryClosed)
|
||||||
if err != nil && err != ListenIsDone {
|
}()
|
||||||
return err
|
|
||||||
|
select {
|
||||||
|
case <-foundryClosed:
|
||||||
|
foundry.transport.Logger.Info("Stopped foundry server")
|
||||||
|
case <-ctx.Done():
|
||||||
|
return CloseTimeoutExceed
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (foundry *FoundryApi) ConnectToWebSocket() (bool, error) {
|
||||||
|
var err error
|
||||||
|
if !foundry.transport.HasSessionId() {
|
||||||
|
err = foundry.transport.InitSessionId()
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
err = foundry.transport.ConnectToFoundry()
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
|
||||||
|
err = foundry.transport.InitWebSocketConnection()
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
|
||||||
|
status, err := foundry.transport.Http.GetStatus()
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
foundry.Status = *types.NewFoundryStatus(status)
|
||||||
|
|
||||||
|
err = foundry.ListenAndServeWS()
|
||||||
|
if err != nil && err != ListenIsDone {
|
||||||
|
return true, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return true, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (foundry *FoundryApi) StartListenFoundry() {
|
||||||
|
foundry.background(func() {
|
||||||
|
timeInterval := 1 * time.Second
|
||||||
|
|
||||||
|
for {
|
||||||
|
ok, err := foundry.ConnectToWebSocket()
|
||||||
|
if err != nil {
|
||||||
|
foundry.transport.Logger.Error("Error raised", "err", err.Error())
|
||||||
|
|
||||||
|
if ok {
|
||||||
|
timeInterval = 1 * time.Second
|
||||||
|
}
|
||||||
|
foundry.transport.Logger.Info("Trying to reconnect", "timer", timeInterval.String())
|
||||||
|
|
||||||
|
time.Sleep(timeInterval)
|
||||||
|
timeInterval = min(timeInterval*2, 15*time.Second)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|||||||
18
internal/foundry/helpers.go
Normal file
18
internal/foundry/helpers.go
Normal file
@@ -0,0 +1,18 @@
|
|||||||
|
package foundry
|
||||||
|
|
||||||
|
import "fmt"
|
||||||
|
|
||||||
|
func (f *FoundryApi) background(fn func()) {
|
||||||
|
f.wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer f.wg.Done()
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
if err := recover(); err != nil {
|
||||||
|
f.transport.Logger.Error(fmt.Sprintf("%s", err))
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
fn()
|
||||||
|
}()
|
||||||
|
}
|
||||||
@@ -80,7 +80,12 @@ func (tr *FoundryTransport) ListenWebSocket() {
|
|||||||
for {
|
for {
|
||||||
_, message, err := tr.WsConn.ReadMessage()
|
_, message, err := tr.WsConn.ReadMessage()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
tr.ReadChan.Err() <- &types.FoundryError{Direction: types.ReaderCode, Err: err, Type: types.WebSocketCode}
|
select {
|
||||||
|
case <-tr.ReadChan.Done():
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
tr.ReadChan.Err() <- &types.FoundryError{Direction: types.ReaderCode, Err: err, Type: types.WebSocketCode, IsFatal: true}
|
||||||
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
data, err := types.ParseWsRespMessage(string(message), types.RequestCodes)
|
data, err := types.ParseWsRespMessage(string(message), types.RequestCodes)
|
||||||
@@ -105,7 +110,7 @@ func (tr *FoundryTransport) GetJsonData(stateType string) ([]byte, error) {
|
|||||||
return tr.HandleWebsocketRequest(msg)
|
return tr.HandleWebsocketRequest(msg)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (tr *FoundryTransport) CloseWebSocketConn() {
|
func (tr *FoundryTransport) CloseWebSocketConn() error {
|
||||||
tr.Logger.Debug("Websocket has been closed")
|
tr.Logger.Debug("Websocket has been closed")
|
||||||
tr.WsConn.Close()
|
return tr.WsConn.Close()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,6 +8,27 @@ type ExchangeChannels struct {
|
|||||||
ProgressMutex sync.Mutex
|
ProgressMutex sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func CloseMapChannel[K comparable, V any](channel map[K]chan V) {
|
||||||
|
for _, v := range channel {
|
||||||
|
select {
|
||||||
|
case _, ok := <-v:
|
||||||
|
if ok {
|
||||||
|
close(v)
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
close(v)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for k := range channel {
|
||||||
|
delete(channel, k)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (channels *ExchangeChannels) Close() {
|
||||||
|
CloseMapChannel(channels.Msgs)
|
||||||
|
CloseMapChannel(channels.ProgressMsg)
|
||||||
|
}
|
||||||
|
|
||||||
type ReadChannels struct {
|
type ReadChannels struct {
|
||||||
done chan struct{}
|
done chan struct{}
|
||||||
msg chan *WsMessage
|
msg chan *WsMessage
|
||||||
@@ -26,10 +47,21 @@ func (channels ReadChannels) Done() chan struct{} {
|
|||||||
return channels.done
|
return channels.done
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func CloseChannel[V any](channel chan V) {
|
||||||
|
select {
|
||||||
|
case _, ok := <-channel:
|
||||||
|
if ok {
|
||||||
|
close(channel)
|
||||||
|
}
|
||||||
|
default:
|
||||||
|
close(channel)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (channels *ReadChannels) Close() {
|
func (channels *ReadChannels) Close() {
|
||||||
close(channels.done)
|
CloseChannel(channels.done)
|
||||||
close(channels.err)
|
CloseChannel(channels.err)
|
||||||
close(channels.msg)
|
CloseChannel(channels.msg)
|
||||||
}
|
}
|
||||||
|
|
||||||
func InitWsChannels() *ReadChannels {
|
func InitWsChannels() *ReadChannels {
|
||||||
Reference in New Issue
Block a user