fix(backend): fix memory leaks and resource issues in logs system
CI / build (push) Successful in 1m18s
CI / build (push) Successful in 1m18s
- Fix Subscribe overwriting channel without closing old one (goroutine leak) - Add context.Context to Stream() for cancellation support - Add WebSocket read pump to detect dead clients (StreamLogs, StreamSteamCMDLogs, StreamRPTLogs) - Use defer f.Close() pattern in StreamRPTLogs file handling - Fix findLatestRPT nil dereference on failed os.Stat - Optimize findLatestRPT to collect stat results before sorting - Apply gofmt formatting to modlists.go, settings.go
This commit is contained in:
@@ -19,6 +19,25 @@ var upgrader = websocket.Upgrader{
|
||||
CheckOrigin: func(r *http.Request) bool { return true },
|
||||
}
|
||||
|
||||
const wsReadTimeout = 60 * time.Second
|
||||
|
||||
func readPump(conn *websocket.Conn, done chan struct{}) {
|
||||
defer close(done)
|
||||
conn.SetReadLimit(4096)
|
||||
conn.SetReadDeadline(time.Now().Add(wsReadTimeout))
|
||||
conn.SetPongHandler(func(string) error {
|
||||
conn.SetReadDeadline(time.Now().Add(wsReadTimeout))
|
||||
return nil
|
||||
})
|
||||
for {
|
||||
_, _, err := conn.ReadMessage()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
conn.SetReadDeadline(time.Now().Add(wsReadTimeout))
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) StreamLogs(c *gin.Context) {
|
||||
conn, err := upgrader.Upgrade(c.Writer, c.Request, nil)
|
||||
if err != nil {
|
||||
@@ -27,13 +46,24 @@ func (h *Handler) StreamLogs(c *gin.Context) {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
done := make(chan struct{})
|
||||
go readPump(conn, done)
|
||||
|
||||
clientID := conn.RemoteAddr().String()
|
||||
ch := h.streamer.Subscribe("server", clientID)
|
||||
defer h.streamer.Unsubscribe("server", clientID)
|
||||
|
||||
for line := range ch {
|
||||
for {
|
||||
select {
|
||||
case line, ok := <-ch:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if err := conn.WriteMessage(websocket.TextMessage, []byte(line)); err != nil {
|
||||
break
|
||||
return
|
||||
}
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -83,13 +113,24 @@ func (h *Handler) StreamSteamCMDLogs(c *gin.Context) {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
done := make(chan struct{})
|
||||
go readPump(conn, done)
|
||||
|
||||
clientID := conn.RemoteAddr().String()
|
||||
ch := h.streamer.Subscribe("steamcmd", clientID)
|
||||
defer h.streamer.Unsubscribe("steamcmd", clientID)
|
||||
|
||||
for line := range ch {
|
||||
for {
|
||||
select {
|
||||
case line, ok := <-ch:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if err := conn.WriteMessage(websocket.TextMessage, []byte(line)); err != nil {
|
||||
break
|
||||
return
|
||||
}
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -102,13 +143,22 @@ func (h *Handler) StreamRPTLogs(c *gin.Context) {
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
done := make(chan struct{})
|
||||
go readPump(conn, done)
|
||||
|
||||
profilesDir := h.process.ProfilesDir()
|
||||
var currentPath string
|
||||
var currentOffset int64
|
||||
ticker := time.NewTicker(250 * time.Millisecond)
|
||||
defer ticker.Stop()
|
||||
|
||||
for range ticker.C {
|
||||
for {
|
||||
select {
|
||||
case <-done:
|
||||
return
|
||||
case <-ticker.C:
|
||||
}
|
||||
|
||||
path := findLatestRPT(profilesDir)
|
||||
|
||||
if path != currentPath {
|
||||
@@ -134,15 +184,19 @@ func (h *Handler) StreamRPTLogs(c *gin.Context) {
|
||||
continue
|
||||
}
|
||||
|
||||
if err := func() error {
|
||||
f, err := os.Open(currentPath)
|
||||
if err != nil {
|
||||
currentPath = ""
|
||||
continue
|
||||
return nil
|
||||
}
|
||||
defer f.Close()
|
||||
|
||||
if _, err := f.Seek(currentOffset, io.SeekStart); err != nil {
|
||||
return nil
|
||||
}
|
||||
f.Seek(currentOffset, io.SeekStart)
|
||||
buf := make([]byte, fi.Size()-currentOffset)
|
||||
n, _ := io.ReadFull(f, buf)
|
||||
f.Close()
|
||||
|
||||
currentOffset += int64(n)
|
||||
for _, line := range strings.Split(string(buf[:n]), "\n") {
|
||||
@@ -150,9 +204,13 @@ func (h *Handler) StreamRPTLogs(c *gin.Context) {
|
||||
continue
|
||||
}
|
||||
if err := conn.WriteMessage(websocket.TextMessage, []byte(line)); err != nil {
|
||||
return
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}(); err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -162,12 +220,25 @@ func findLatestRPT(dir string) string {
|
||||
if err != nil || len(matches) == 0 {
|
||||
return ""
|
||||
}
|
||||
sort.Slice(matches, func(i, j int) bool {
|
||||
fi, _ := os.Stat(matches[i])
|
||||
fj, _ := os.Stat(matches[j])
|
||||
return fi.ModTime().After(fj.ModTime())
|
||||
type rptInfo struct {
|
||||
path string
|
||||
time time.Time
|
||||
}
|
||||
var infos []rptInfo
|
||||
for _, m := range matches {
|
||||
fi, err := os.Stat(m)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
infos = append(infos, rptInfo{path: m, time: fi.ModTime()})
|
||||
}
|
||||
if len(infos) == 0 {
|
||||
return ""
|
||||
}
|
||||
sort.Slice(infos, func(i, j int) bool {
|
||||
return infos[i].time.After(infos[j].time)
|
||||
})
|
||||
return matches[0]
|
||||
return infos[0].path
|
||||
}
|
||||
|
||||
func (h *Handler) GetLog(c *gin.Context) {
|
||||
|
||||
@@ -45,20 +45,48 @@ func (h *Handler) UpdateSettings(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
if input.IPPort != nil { s.IPPort = *input.IPPort }
|
||||
if input.ServerParameters != nil { s.ServerParameters = *input.ServerParameters }
|
||||
if input.SteamBranch != nil { s.SteamBranch = *input.SteamBranch }
|
||||
if input.SteamUser != nil { s.SteamUser = *input.SteamUser }
|
||||
if input.Platform != nil { s.Platform = *input.Platform }
|
||||
if input.CBASettings != nil { s.CBASettings = *input.CBASettings }
|
||||
if input.AILevelPresets != nil { s.AILevelPresets = *input.AILevelPresets }
|
||||
if input.DifficultyPresets != nil { s.DifficultyPresets = *input.DifficultyPresets }
|
||||
if input.ActiveConfig != nil { s.ActiveConfig = *input.ActiveConfig }
|
||||
if input.ActiveModlist != nil { s.ActiveModlist = *input.ActiveModlist }
|
||||
if input.AutoUpdateOnStartup != nil { s.AutoUpdateOnStartup = *input.AutoUpdateOnStartup }
|
||||
if input.AutoStartOnStartup != nil { s.AutoStartOnStartup = *input.AutoStartOnStartup }
|
||||
if input.AutoUpdateModsOnStartup != nil { s.AutoUpdateModsOnStartup = *input.AutoUpdateModsOnStartup }
|
||||
if input.ScheduledUpdate != nil { s.ScheduledUpdate = *input.ScheduledUpdate }
|
||||
if input.IPPort != nil {
|
||||
s.IPPort = *input.IPPort
|
||||
}
|
||||
if input.ServerParameters != nil {
|
||||
s.ServerParameters = *input.ServerParameters
|
||||
}
|
||||
if input.SteamBranch != nil {
|
||||
s.SteamBranch = *input.SteamBranch
|
||||
}
|
||||
if input.SteamUser != nil {
|
||||
s.SteamUser = *input.SteamUser
|
||||
}
|
||||
if input.Platform != nil {
|
||||
s.Platform = *input.Platform
|
||||
}
|
||||
if input.CBASettings != nil {
|
||||
s.CBASettings = *input.CBASettings
|
||||
}
|
||||
if input.AILevelPresets != nil {
|
||||
s.AILevelPresets = *input.AILevelPresets
|
||||
}
|
||||
if input.DifficultyPresets != nil {
|
||||
s.DifficultyPresets = *input.DifficultyPresets
|
||||
}
|
||||
if input.ActiveConfig != nil {
|
||||
s.ActiveConfig = *input.ActiveConfig
|
||||
}
|
||||
if input.ActiveModlist != nil {
|
||||
s.ActiveModlist = *input.ActiveModlist
|
||||
}
|
||||
if input.AutoUpdateOnStartup != nil {
|
||||
s.AutoUpdateOnStartup = *input.AutoUpdateOnStartup
|
||||
}
|
||||
if input.AutoStartOnStartup != nil {
|
||||
s.AutoStartOnStartup = *input.AutoStartOnStartup
|
||||
}
|
||||
if input.AutoUpdateModsOnStartup != nil {
|
||||
s.AutoUpdateModsOnStartup = *input.AutoUpdateModsOnStartup
|
||||
}
|
||||
if input.ScheduledUpdate != nil {
|
||||
s.ScheduledUpdate = *input.ScheduledUpdate
|
||||
}
|
||||
|
||||
if err := h.process.WriteUserconfigFiles(s); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "write userconfig: " + err.Error()})
|
||||
|
||||
@@ -2,6 +2,7 @@ package services
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"io"
|
||||
"log"
|
||||
"sync"
|
||||
@@ -26,6 +27,10 @@ func (ls *LogStreamer) Subscribe(serverID, clientID string) chan string {
|
||||
ls.subs[serverID] = make(map[string]chan string)
|
||||
}
|
||||
|
||||
if old, ok := ls.subs[serverID][clientID]; ok {
|
||||
close(old)
|
||||
}
|
||||
|
||||
ch := make(chan string, 256)
|
||||
ls.subs[serverID][clientID] = ch
|
||||
return ch
|
||||
@@ -46,15 +51,24 @@ func (ls *LogStreamer) Unsubscribe(serverID, clientID string) {
|
||||
}
|
||||
}
|
||||
|
||||
func (ls *LogStreamer) Stream(serverID string, reader io.Reader, closeMsg string) {
|
||||
func (ls *LogStreamer) Stream(ctx context.Context, serverID string, reader io.Reader, closeMsg string) {
|
||||
scanner := bufio.NewScanner(reader)
|
||||
for scanner.Scan() {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
line := scanner.Text()
|
||||
ls.broadcast(serverID, line)
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
default:
|
||||
log.Printf("log stream error for server %s: %v", serverID, err)
|
||||
}
|
||||
}
|
||||
if closeMsg != "" {
|
||||
ls.broadcast(serverID, closeMsg)
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -95,7 +96,7 @@ func TestLogStreamer_Stream(t *testing.T) {
|
||||
ch := ls.Subscribe("server", "client1")
|
||||
|
||||
reader := strings.NewReader("line1\nline2\nline3\n")
|
||||
go ls.Stream("server", reader, "DONE")
|
||||
go ls.Stream(context.Background(), "server", reader, "DONE")
|
||||
|
||||
lines := []string{}
|
||||
for i := 0; i < 4; i++ { // 3 lines + DONE close message
|
||||
@@ -174,3 +175,68 @@ func TestLogStreamer_ConcurrentSubscribeUnsubscribe(t *testing.T) {
|
||||
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
func TestLogStreamer_SubscribeOverwriteClosesOld(t *testing.T) {
|
||||
ls := NewLogStreamer()
|
||||
|
||||
ch1 := ls.Subscribe("server", "client1")
|
||||
ls.Broadcast("server", "first")
|
||||
<-ch1
|
||||
|
||||
// Re-subscribe with same clientID — old channel should be closed
|
||||
ch2 := ls.Subscribe("server", "client1")
|
||||
|
||||
// Old channel should be closed
|
||||
select {
|
||||
case _, ok := <-ch1:
|
||||
if ok {
|
||||
t.Error("old channel should be closed after re-subscribe")
|
||||
}
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
t.Error("old channel not closed within timeout")
|
||||
}
|
||||
|
||||
// New channel should work
|
||||
ls.Broadcast("server", "second")
|
||||
select {
|
||||
case line := <-ch2:
|
||||
if line != "second" {
|
||||
t.Errorf("new channel received %q, want %q", line, "second")
|
||||
}
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
t.Error("timeout waiting on new channel")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLogStreamer_StreamContextCancel(t *testing.T) {
|
||||
ls := NewLogStreamer()
|
||||
ch := ls.Subscribe("server", "client1")
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
reader := contextReader{ctx: ctx}
|
||||
go ls.Stream(ctx, "server", reader, "")
|
||||
|
||||
// Broadcast should still work while stream is running
|
||||
ls.Broadcast("server", "live")
|
||||
select {
|
||||
case line := <-ch:
|
||||
if line != "live" {
|
||||
t.Errorf("received %q, want %q", line, "live")
|
||||
}
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
t.Error("timeout waiting for broadcast")
|
||||
}
|
||||
|
||||
// Cancel context — stream should stop
|
||||
cancel()
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
}
|
||||
|
||||
type contextReader struct {
|
||||
ctx context.Context
|
||||
}
|
||||
|
||||
func (r contextReader) Read(p []byte) (int, error) {
|
||||
<-r.ctx.Done()
|
||||
return 0, r.ctx.Err()
|
||||
}
|
||||
|
||||
@@ -148,8 +148,8 @@ func (pm *ProcessManager) Start() error {
|
||||
pm.mu.Unlock()
|
||||
pm.state.Store(int32(procRunning))
|
||||
|
||||
go pm.streamer.Stream("server", stdout, "")
|
||||
go pm.streamer.Stream("server", stderr, "")
|
||||
go pm.streamer.Stream(ctx, "server", stdout, "")
|
||||
go pm.streamer.Stream(ctx, "server", stderr, "")
|
||||
|
||||
go func() {
|
||||
cmd.Wait()
|
||||
|
||||
@@ -117,8 +117,8 @@ func (s *SteamCmdManager) run(label string, args []string) error {
|
||||
|
||||
s.cancel = cancel
|
||||
|
||||
go s.streamer.Stream("steamcmd", stdout, "")
|
||||
go s.streamer.Stream("steamcmd", stderr, "")
|
||||
go s.streamer.Stream(ctx, "steamcmd", stdout, "")
|
||||
go s.streamer.Stream(ctx, "steamcmd", stderr, "")
|
||||
|
||||
go func() {
|
||||
err := cmd.Wait()
|
||||
|
||||
Reference in New Issue
Block a user