Start handles sending Server-Sent Event to the client.
(w http.ResponseWriter, r *http.Request)
| 88 | |
| 89 | // Start handles sending Server-Sent Event to the client. |
| 90 | func (s *EventStream) Start(w http.ResponseWriter, r *http.Request) { |
| 91 | // Signal completion on exit so senders don't block indefinitely after closure. |
| 92 | defer close(s.doneCh) |
| 93 | |
| 94 | ctx := r.Context() |
| 95 | |
| 96 | defer s.tick.Stop() |
| 97 | |
| 98 | for { |
| 99 | var ( |
| 100 | ev event |
| 101 | open bool |
| 102 | ) |
| 103 | |
| 104 | select { |
| 105 | case <-s.ctx.Done(): |
| 106 | return |
| 107 | case <-ctx.Done(): |
| 108 | s.logger.Debug(ctx, "request context canceled", slog.Error(ctx.Err())) |
| 109 | return |
| 110 | case ev, open = <-s.eventsCh: // Once closed, the buffered channel will drain all buffered values before showing as closed. |
| 111 | if !open { |
| 112 | s.logger.Debug(ctx, "events channel closed") |
| 113 | return |
| 114 | } |
| 115 | |
| 116 | // Initiate the stream on first event (if not already initiated). |
| 117 | s.InitiateStream(w) |
| 118 | case <-s.tick.C: |
| 119 | ev = s.pingPayload |
| 120 | if ev == nil { |
| 121 | continue |
| 122 | } |
| 123 | } |
| 124 | |
| 125 | _, err := w.Write(ev) |
| 126 | if err != nil { |
| 127 | if IsConnError(err) { |
| 128 | s.logger.Debug(ctx, "client disconnected during SSE write", slog.Error(err)) |
| 129 | } else { |
| 130 | s.logger.Warn(ctx, "failed to write SSE event", slog.Error(err)) |
| 131 | } |
| 132 | return |
| 133 | } |
| 134 | if err := flush(w); err != nil { |
| 135 | s.logger.Warn(ctx, "failed to flush event stream", slog.Error(err)) |
| 136 | return |
| 137 | } |
| 138 | |
| 139 | // Reset the timer once we've flushed some data to the stream, since it's already fresh. |
| 140 | // No need to ping in that case. |
| 141 | s.tick.Reset(pingInterval) |
| 142 | } |
| 143 | } |
| 144 | |
| 145 | // Send enqueues an event in a non-blocking fashion, but if the channel is full |
| 146 | // then it will block. |
no test coverage detected