aboutsummaryrefslogtreecommitdiff
path: root/backend/internal/api/sse.go
diff options
context:
space:
mode:
Diffstat (limited to 'backend/internal/api/sse.go')
-rw-r--r--backend/internal/api/sse.go57
1 files changed, 57 insertions, 0 deletions
diff --git a/backend/internal/api/sse.go b/backend/internal/api/sse.go
new file mode 100644
index 0000000..7c2332c
--- /dev/null
+++ b/backend/internal/api/sse.go
@@ -0,0 +1,57 @@
+package api
+
+import (
+ "fmt"
+ "log"
+ "net/http"
+
+ "yaum/internal/sse"
+)
+
+type StreamHandler struct {
+ broadcaster *sse.Broadcaster
+}
+
+func NewStreamHandler(b *sse.Broadcaster) *StreamHandler {
+ return &StreamHandler{broadcaster: b}
+}
+
+func (sh *StreamHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
+ flusher, ok := w.(http.Flusher)
+ if !ok {
+ http.Error(w, "streaming 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")
+
+ ch := sh.broadcaster.Subscribe()
+ defer sh.broadcaster.Unsubscribe(ch)
+
+ _, err := fmt.Fprintf(w, "event: connected\ndata: {}\n\n")
+ if err != nil {
+ return
+ }
+ flusher.Flush()
+
+ ctx := r.Context()
+ for {
+ select {
+ case data, ok := <-ch:
+ if !ok {
+ return
+ }
+ _, err := fmt.Fprintf(w, "event: heartbeat\ndata: %s\n\n", data)
+ if err != nil {
+ log.Printf("[sse] cliente desconectado: %v", err)
+ return
+ }
+ flusher.Flush()
+ case <-ctx.Done():
+ return
+ }
+ }
+}