subscribeEvents is an SSE endpoint that sends events to the client
(ctx context.Context, input *struct{}, send sse.Sender)
| 539 | |
| 540 | // subscribeEvents is an SSE endpoint that sends events to the client |
| 541 | func (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 | |
| 580 | func (s *Server) subscribeScreen(ctx context.Context, input *struct{}, send sse.Sender) { |
| 581 | subscriberId, ch, stateEvents := s.emitter.Subscribe() |
nothing calls this directly
no test coverage detected