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
96 lines
1.7 KiB
Go
96 lines
1.7 KiB
Go
package services
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"io"
|
|
"log"
|
|
"sync"
|
|
)
|
|
|
|
type LogStreamer struct {
|
|
mu sync.RWMutex
|
|
subs map[string]map[string]chan string
|
|
}
|
|
|
|
func NewLogStreamer() *LogStreamer {
|
|
return &LogStreamer{
|
|
subs: make(map[string]map[string]chan string),
|
|
}
|
|
}
|
|
|
|
func (ls *LogStreamer) Subscribe(serverID, clientID string) chan string {
|
|
ls.mu.Lock()
|
|
defer ls.mu.Unlock()
|
|
|
|
if ls.subs[serverID] == nil {
|
|
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
|
|
}
|
|
|
|
func (ls *LogStreamer) Unsubscribe(serverID, clientID string) {
|
|
ls.mu.Lock()
|
|
defer ls.mu.Unlock()
|
|
|
|
if ls.subs[serverID] != nil {
|
|
if ch, ok := ls.subs[serverID][clientID]; ok {
|
|
close(ch)
|
|
delete(ls.subs[serverID], clientID)
|
|
}
|
|
if len(ls.subs[serverID]) == 0 {
|
|
delete(ls.subs, serverID)
|
|
}
|
|
}
|
|
}
|
|
|
|
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)
|
|
}
|
|
}
|
|
|
|
func (ls *LogStreamer) Broadcast(serverID, line string) {
|
|
ls.broadcast(serverID, line)
|
|
}
|
|
|
|
func (ls *LogStreamer) broadcast(serverID, line string) {
|
|
ls.mu.RLock()
|
|
defer ls.mu.RUnlock()
|
|
|
|
if ls.subs[serverID] == nil {
|
|
return
|
|
}
|
|
|
|
for _, ch := range ls.subs[serverID] {
|
|
select {
|
|
case ch <- line:
|
|
default:
|
|
}
|
|
}
|
|
}
|