MCPcopy Create free account
hub / github.com/coder/aibridge / Start

Method Start

intercept/eventstream/eventstream.go:90–143  ·  view source on GitHub ↗

Start handles sending Server-Sent Event to the client.

(w http.ResponseWriter, r *http.Request)

Source from the content-addressed store, hash-verified

88
89// Start handles sending Server-Sent Event to the client.
90func (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.

Callers 15

ProcessRequestMethod · 0.95
ProcessRequestMethod · 0.95
ProcessRequestMethod · 0.95
newInterceptionProcessorFunction · 0.80
newPassthroughRouterFunction · 0.80
newStreamMethod · 0.80
ProcessRequestMethod · 0.80
newResponseMethod · 0.80
newStreamMethod · 0.80
ProcessRequestMethod · 0.80
newMessageMethod · 0.80
newStreamMethod · 0.80

Calls 5

InitiateStreamMethod · 0.95
IsConnErrorFunction · 0.85
flushFunction · 0.85
ErrorMethod · 0.45
WriteMethod · 0.45

Tested by

no test coverage detected