MCPcopy Create free account
hub / github.com/coder/agentapi / subscribeEvents

Method subscribeEvents

lib/httpapi/server.go:541–578  ·  view source on GitHub ↗

subscribeEvents is an SSE endpoint that sends events to the client

(ctx context.Context, input *struct{}, send sse.Sender)

Source from the content-addressed store, hash-verified

539
540// subscribeEvents is an SSE endpoint that sends events to the client
541func (s *Server) subscribeEvents(ctx context.Context, input *struct{}, send sse.Sender) {
542 subscriberId, ch, stateEvents := s.emitter.Subscribe()
543 defer s.emitter.Unsubscribe(subscriberId)
544
545 s.logger.Info("New subscriber", "subscriberId", subscriberId)
546 for _, event := range stateEvents {
547 if event.Type == EventTypeScreenUpdate {
548 continue
549 }
550 if err := send.Data(event.Payload); err != nil {
551 s.logger.Error("Failed to send event", "subscriberId", subscriberId, "error", err)
552 return
553 }
554 }
555
556 for {
557 select {
558 case event, ok := <-ch:
559 if !ok {
560 s.logger.Info("Channel closed", "subscriberId", subscriberId)
561 return
562 }
563 if event.Type == EventTypeScreenUpdate {
564 continue
565 }
566 if err := send.Data(event.Payload); err != nil {
567 s.logger.Error("Failed to send event", "subscriberId", subscriberId, "error", err)
568 return
569 }
570 case <-s.shutdownCtx.Done():
571 s.logger.Info("Server stop initiated, unsubscribing.", "subscriberId", subscriberId)
572 return
573 case <-ctx.Done():
574 s.logger.Info("Context done", "subscriberId", subscriberId)
575 return
576 }
577 }
578}
579
580func (s *Server) subscribeScreen(ctx context.Context, input *struct{}, send sse.Sender) {
581 subscriberId, ch, stateEvents := s.emitter.Subscribe()

Callers

nothing calls this directly

Calls 2

SubscribeMethod · 0.80
UnsubscribeMethod · 0.80

Tested by

no test coverage detected