Files
move-data-nas/internal/syncengine/eventbus.go
T
darroyo 58e7f51ba7 fix: EventBus panic on SSE disconnect + JWT secret persistence + recover() guards
- eventbus.go: Fix send-on-closed-channel panic in SubscribeGlobal by
  using a done channel; add recover() in fan-out goroutine; track active
  global subs for proper cleanup on unsubscribe
- config.go: Persist JWT secret to $DATA_DIR/.jwt_secret instead of
  regenerating a random one on every restart (which invalidated all sessions)
- handlers_ws.go: Replace time.After with time.Ticker to fix timer leak in
  SSE keepalive loop
- handlers_jobs.go: Add recover() in fire-and-forget job goroutine; fix nil
  pointer deref when GetByID fails after job creation
- handlers_machines.go: Add recover() in ProbeAllMachines goroutine
- scheduler.go: Add recover() in scheduled job run goroutine
- engine.go: Add recover() in per-machine probe goroutines
2026-07-08 22:38:56 -04:00

122 lines
2.4 KiB
Go

package syncengine
import (
"log/slog"
"sync"
)
type EventBus struct {
subscribers map[int64]map[chan Event]struct{}
mu sync.RWMutex
global chan Event
bufferSize int
globalSubs []globalSub
}
type globalSub struct {
ch chan Event
done chan struct{}
}
func NewEventBus(bufferSize int) *EventBus {
return &EventBus{
subscribers: make(map[int64]map[chan Event]struct{}),
global: make(chan Event, bufferSize),
bufferSize: bufferSize,
globalSubs: nil,
}
}
func (eb *EventBus) Subscribe(jobID int64) (chan Event, func()) {
eb.mu.Lock()
defer eb.mu.Unlock()
if eb.subscribers[jobID] == nil {
eb.subscribers[jobID] = make(map[chan Event]struct{})
}
ch := make(chan Event, eb.bufferSize)
eb.subscribers[jobID][ch] = struct{}{}
unsubscribe := func() {
eb.mu.Lock()
defer eb.mu.Unlock()
if subs, ok := eb.subscribers[jobID]; ok {
delete(subs, ch)
if len(subs) == 0 {
delete(eb.subscribers, jobID)
}
}
close(ch)
}
return ch, unsubscribe
}
func (eb *EventBus) SubscribeGlobal() (chan Event, func()) {
ch := make(chan Event, eb.bufferSize)
done := make(chan struct{})
eb.mu.Lock()
eb.globalSubs = append(eb.globalSubs, globalSub{ch: ch, done: done})
eb.mu.Unlock()
go func() {
defer func() {
if r := recover(); r != nil {
slog.Error("SubscribeGlobal goroutine panicked", "reason", r)
}
close(ch)
}()
for {
select {
case evt := <-eb.global:
select {
case ch <- evt:
default:
slog.Warn("global event subscriber buffer full, dropping event", "type", evt.Type)
}
case <-done:
return
}
}
}()
return ch, func() {
close(done)
eb.mu.Lock()
for i, s := range eb.globalSubs {
if s.ch == ch {
eb.globalSubs = append(eb.globalSubs[:i], eb.globalSubs[i+1:]...)
break
}
}
eb.mu.Unlock()
}
}
func (eb *EventBus) Publish(evt Event) {
eb.mu.RLock()
defer eb.mu.RUnlock()
if subs, ok := eb.subscribers[evt.JobID]; ok {
for ch := range subs {
select {
case ch <- evt:
default:
slog.Warn("job event subscriber buffer full, dropping event", "job_id", evt.JobID)
}
}
}
select {
case eb.global <- evt:
default:
slog.Warn("global event bus full, dropping event", "type", evt.Type)
}
}
func (eb *EventBus) CloseJobChannels(jobID int64) {
eb.mu.Lock()
defer eb.mu.Unlock()
if subs, ok := eb.subscribers[jobID]; ok {
for ch := range subs {
close(ch)
}
delete(eb.subscribers, jobID)
}
}