From 58e7f51ba7d91ef2b437180fc0cc59fe82225f18 Mon Sep 17 00:00:00 2001 From: Daniel Arroyo Date: Wed, 8 Jul 2026 22:38:56 -0400 Subject: [PATCH] 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 --- internal/api/handlers_jobs.go | 12 ++++++++- internal/api/handlers_machines.go | 10 ++++++- internal/api/handlers_ws.go | 10 +++++-- internal/config/config.go | 20 +++++++++----- internal/scheduler/scheduler.go | 5 ++++ internal/syncengine/engine.go | 7 ++++- internal/syncengine/eventbus.go | 45 +++++++++++++++++++++++++------ 7 files changed, 89 insertions(+), 20 deletions(-) diff --git a/internal/api/handlers_jobs.go b/internal/api/handlers_jobs.go index 9a21cd1..0d80c08 100644 --- a/internal/api/handlers_jobs.go +++ b/internal/api/handlers_jobs.go @@ -4,6 +4,7 @@ import ( "database/sql" "encoding/json" "fmt" + "log/slog" "net/http" "os" "strconv" @@ -150,11 +151,20 @@ func (h *JobHandler) TriggerRun(w http.ResponseWriter, r *http.Request) { } go func() { + defer func() { + if r := recover(); r != nil { + slog.Error("job run goroutine panicked", "job_id", jobID, "panic", r) + } + }() h.engine.Run(r.Context(), jobID, pairID) }() jobRepo := models.NewJobRepository(h.db) - j, _ := jobRepo.GetByID(jobID) + j, err := jobRepo.GetByID(jobID) + if err != nil { + writeError(w, http.StatusInternalServerError, "failed to fetch created job") + return + } writeJSON(w, jobToResp(*j), http.StatusCreated) } diff --git a/internal/api/handlers_machines.go b/internal/api/handlers_machines.go index ca14648..ec9cfb2 100644 --- a/internal/api/handlers_machines.go +++ b/internal/api/handlers_machines.go @@ -3,6 +3,7 @@ package api import ( "database/sql" "encoding/json" + "log/slog" "net/http" "regexp" "strconv" @@ -216,7 +217,14 @@ func (h *MachineHandler) Refresh(w http.ResponseWriter, r *http.Request) { writeError(w, http.StatusInternalServerError, "engine not available") return } - go h.engine.ProbeAllMachines() + go func() { + defer func() { + if r := recover(); r != nil { + slog.Error("ProbeAllMachines panicked", "panic", r) + } + }() + h.engine.ProbeAllMachines() + }() w.WriteHeader(http.StatusAccepted) writeJSON(w, map[string]string{"status": "probing"}) } diff --git a/internal/api/handlers_ws.go b/internal/api/handlers_ws.go index 97211f8..1372707 100644 --- a/internal/api/handlers_ws.go +++ b/internal/api/handlers_ws.go @@ -38,6 +38,9 @@ func (h *SSEHandler) StreamAll(w http.ResponseWriter, r *http.Request) { events, unsub := h.engine.SubscribeGlobal() defer unsub() + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + for { select { case evt := <-events: @@ -46,7 +49,7 @@ func (h *SSEHandler) StreamAll(w http.ResponseWriter, r *http.Request) { flusher.Flush() case <-r.Context().Done(): return - case <-time.After(30 * time.Second): + case <-ticker.C: fmt.Fprintf(w, ": keepalive\n\n") flusher.Flush() } @@ -84,6 +87,9 @@ func (h *SSEHandler) StreamJob(w http.ResponseWriter, r *http.Request) { events, unsub := h.engine.SubscribeJob(jobID) defer unsub() + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + for { select { case evt := <-events: @@ -92,7 +98,7 @@ func (h *SSEHandler) StreamJob(w http.ResponseWriter, r *http.Request) { flusher.Flush() case <-r.Context().Done(): return - case <-time.After(30 * time.Second): + case <-ticker.C: fmt.Fprintf(w, ": keepalive\n\n") flusher.Flush() } diff --git a/internal/config/config.go b/internal/config/config.go index 1435afd..3813ce2 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -1,6 +1,8 @@ package config import ( + "crypto/rand" + "encoding/hex" "fmt" "os" "path/filepath" @@ -80,19 +82,23 @@ func Load(configPath, dataDir, addr string) (*Config, error) { } } + secretPath := filepath.Join(cfg.DataDir, ".jwt_secret") + if cfg.Auth.JWTSecret == "" { + if data, err := os.ReadFile(secretPath); err == nil && len(data) >= 32 { + cfg.Auth.JWTSecret = strings.TrimSpace(string(data)) + } + } if cfg.Auth.JWTSecret == "" { b := make([]byte, 32) - f, err := os.Open("/dev/urandom") - if err == nil { - defer f.Close() - n, _ := f.Read(b) - if n == 32 { - cfg.Auth.JWTSecret = fmt.Sprintf("%x", b) - } + if _, err := rand.Read(b); err == nil { + cfg.Auth.JWTSecret = hex.EncodeToString(b) } if cfg.Auth.JWTSecret == "" { cfg.Auth.JWTSecret = "insecure-dev-secret-change-in-production" } + if dirErr := os.MkdirAll(cfg.DataDir, 0700); dirErr == nil { + _ = os.WriteFile(secretPath, []byte(cfg.Auth.JWTSecret+"\n"), 0600) + } } if dataDir := os.Getenv("SYNCSERVER_DATA_DIR"); dataDir != "" { diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index 4d9f136..11c03ab 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -83,6 +83,11 @@ func (s *Scheduler) tick() { ctx := context.Background() go func(jobID int64, pairID int64, schID int64) { + defer func() { + if r := recover(); r != nil { + slog.Error("scheduler: job run panicked", "job_id", jobID, "panic", r) + } + }() if err := s.engine.Run(ctx, jobID, pairID); err != nil { slog.Warn("scheduler: job failed", "job_id", jobID, "error", err) } diff --git a/internal/syncengine/engine.go b/internal/syncengine/engine.go index b4a4915..57b854d 100644 --- a/internal/syncengine/engine.go +++ b/internal/syncengine/engine.go @@ -368,7 +368,12 @@ func (e *Engine) ProbeAllMachines() { sem <- struct{}{} go func(m *models.Machine) { defer wg.Done() - defer func() { <-sem }() + defer func() { + if r := recover(); r != nil { + slog.Error("probe goroutine panicked", "machine_id", m.ID, "panic", r) + } + <-sem + }() ctx, cancel := context.WithTimeout(context.Background(), probeTimeout) defer cancel() diff --git a/internal/syncengine/eventbus.go b/internal/syncengine/eventbus.go index f5acfa4..486850e 100644 --- a/internal/syncengine/eventbus.go +++ b/internal/syncengine/eventbus.go @@ -10,6 +10,12 @@ type EventBus 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 { @@ -17,6 +23,7 @@ func NewEventBus(bufferSize int) *EventBus { subscribers: make(map[int64]map[chan Event]struct{}), global: make(chan Event, bufferSize), bufferSize: bufferSize, + globalSubs: nil, } } @@ -43,20 +50,42 @@ func (eb *EventBus) Subscribe(jobID int64) (chan Event, func()) { } func (eb *EventBus) SubscribeGlobal() (chan Event, func()) { - eb.mu.RLock() ch := make(chan Event, eb.bufferSize) - eb.mu.RUnlock() + done := make(chan struct{}) + eb.mu.Lock() + eb.globalSubs = append(eb.globalSubs, globalSub{ch: ch, done: done}) + eb.mu.Unlock() go func() { - for evt := range eb.global { + defer func() { + if r := recover(); r != nil { + slog.Error("SubscribeGlobal goroutine panicked", "reason", r) + } + close(ch) + }() + for { select { - case ch <- evt: - default: - slog.Warn("global event subscriber buffer full, dropping event", "type", evt.Type) + case evt := <-eb.global: + select { + case ch <- evt: + default: + slog.Warn("global event subscriber buffer full, dropping event", "type", evt.Type) + } + case <-done: + return } } - close(ch) }() - return ch, func() { close(ch) } + 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) {