(ctx context.Context, input *struct{}, send sse.Sender)
| 578 | } |
| 579 | |
| 580 | func (s *Server) subscribeScreen(ctx context.Context, input *struct{}, send sse.Sender) { |
| 581 | subscriberId, ch, stateEvents := s.emitter.Subscribe() |
| 582 | defer s.emitter.Unsubscribe(subscriberId) |
| 583 | s.logger.Info("New screen subscriber", "subscriberId", subscriberId) |
| 584 | for _, event := range stateEvents { |
| 585 | if event.Type != EventTypeScreenUpdate { |
| 586 | continue |
| 587 | } |
| 588 | if err := send.Data(event.Payload); err != nil { |
| 589 | s.logger.Error("Failed to send screen event", "subscriberId", subscriberId, "error", err) |
| 590 | return |
| 591 | } |
| 592 | } |
| 593 | for { |
| 594 | select { |
| 595 | case event, ok := <-ch: |
| 596 | if !ok { |
| 597 | s.logger.Info("Screen channel closed", "subscriberId", subscriberId) |
| 598 | return |
| 599 | } |
| 600 | if event.Type != EventTypeScreenUpdate { |
| 601 | continue |
| 602 | } |
| 603 | if err := send.Data(event.Payload); err != nil { |
| 604 | s.logger.Error("Failed to send screen event", "subscriberId", subscriberId, "error", err) |
| 605 | return |
| 606 | } |
| 607 | case <-s.shutdownCtx.Done(): |
| 608 | s.logger.Info("Server stop initiated, unsubscribing.", "subscriberId", subscriberId) |
| 609 | return |
| 610 | case <-ctx.Done(): |
| 611 | s.logger.Info("Screen context done", "subscriberId", subscriberId) |
| 612 | return |
| 613 | } |
| 614 | } |
| 615 | } |
| 616 | |
| 617 | // Start starts the HTTP server |
| 618 | func (s *Server) Start() error { |
nothing calls this directly
no test coverage detected