package api import ( "encoding/json" "fmt" "net/http" "strconv" "time" "github.com/syncserver/internal/syncengine" ) type SSEHandler struct { engine *syncengine.Engine } func NewSSEHandler(engine *syncengine.Engine) *SSEHandler { return &SSEHandler{engine: engine} } func (h *SSEHandler) Stream(w http.ResponseWriter, r *http.Request) { jobIDStr := r.URL.Query().Get("job_id") var filterJobID int64 if jobIDStr != "" { filterJobID, _ = strconv.ParseInt(jobIDStr, 10, 64) } flusher, ok := w.(http.Flusher) if !ok { http.Error(w, "SSE not supported", http.StatusInternalServerError) return } w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") w.Header().Set("X-Accel-Buffering", "no") flusher.Flush() if h.engine == nil { return } events := h.engine.Events() for { select { case evt := <-events: if filterJobID != 0 && evt.JobID != filterJobID { continue } data, _ := json.Marshal(evt) fmt.Fprintf(w, "event: %s\ndata: %s\n\n", evt.Type, data) flusher.Flush() case <-r.Context().Done(): return case <-time.After(30 * time.Second): fmt.Fprintf(w, ": keepalive\n\n") flusher.Flush() } } }