stream_state.go 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133
  1. package infra
  2. import (
  3. "context"
  4. "errors"
  5. "io"
  6. "net/http"
  7. "sync"
  8. "github.com/porter-dev/porter/api/server/handlers"
  9. "github.com/porter-dev/porter/api/server/shared"
  10. "github.com/porter-dev/porter/api/server/shared/apierrors"
  11. "github.com/porter-dev/porter/api/server/shared/config"
  12. "github.com/porter-dev/porter/api/server/shared/websocket"
  13. "github.com/porter-dev/porter/api/types"
  14. "github.com/porter-dev/porter/internal/models"
  15. "github.com/porter-dev/porter/provisioner/pb"
  16. "google.golang.org/grpc"
  17. "google.golang.org/grpc/metadata"
  18. )
  19. type InfraStreamStateHandler struct {
  20. handlers.PorterHandlerWriter
  21. }
  22. func NewInfraStreamStateHandler(
  23. config *config.Config,
  24. writer shared.ResultWriter,
  25. ) *InfraStreamStateHandler {
  26. return &InfraStreamStateHandler{
  27. PorterHandlerWriter: handlers.NewDefaultPorterHandler(config, nil, writer),
  28. }
  29. }
  30. func (c *InfraStreamStateHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
  31. safeRW := r.Context().Value(types.RequestCtxWebsocketKey).(*websocket.WebsocketSafeReadWriter)
  32. infra, _ := r.Context().Value(types.InfraScope).(*models.Infra)
  33. operation, _ := r.Context().Value(types.OperationScope).(*models.Operation)
  34. workspaceID := models.GetWorkspaceID(infra, operation)
  35. conn, err := grpc.Dial("localhost:8082", grpc.WithInsecure())
  36. if err != nil {
  37. c.HandleAPIError(w, r, apierrors.NewErrInternal(err))
  38. return
  39. }
  40. client := pb.NewProvisionerClient(conn)
  41. header := metadata.New(map[string]string{
  42. "workspace_id": workspaceID,
  43. })
  44. ctx := metadata.NewOutgoingContext(context.Background(), header)
  45. ctx, cancel := context.WithCancel(ctx)
  46. defer cancel()
  47. stream, err := client.GetStateUpdate(ctx, &pb.Infra{
  48. ProjectId: int64(infra.ProjectID),
  49. Id: int64(infra.ID),
  50. Suffix: infra.Suffix,
  51. })
  52. if err != nil {
  53. c.HandleAPIError(w, r, apierrors.NewErrInternal(err))
  54. return
  55. }
  56. errorchan := make(chan error)
  57. var wg sync.WaitGroup
  58. wg.Add(2)
  59. go func() {
  60. wg.Wait()
  61. close(errorchan)
  62. }()
  63. go func() {
  64. defer wg.Done()
  65. for {
  66. if _, _, err := safeRW.ReadMessage(); err != nil {
  67. errorchan <- nil
  68. return
  69. }
  70. }
  71. }()
  72. go func() {
  73. defer wg.Done()
  74. for {
  75. stateUpdate, err := stream.Recv()
  76. if err != nil {
  77. if err == io.EOF || errors.Is(ctx.Err(), context.Canceled) {
  78. errorchan <- nil
  79. } else {
  80. errorchan <- err
  81. }
  82. return
  83. }
  84. err = safeRW.WriteJSON(stateUpdate)
  85. if err != nil {
  86. errorchan <- err
  87. }
  88. }
  89. }()
  90. for err = range errorchan {
  91. if err != nil {
  92. c.HandleAPIErrorNoWrite(w, r, apierrors.NewErrInternal(err))
  93. }
  94. // close the grpc stream: do not check for error case since the stream could already be
  95. // closed
  96. stream.CloseSend()
  97. // close the websocket stream: do not check for error case since the WS could already be
  98. // closed
  99. safeRW.Close()
  100. // cancel the context set for the grpc stream to ensure that Recv is unblocked
  101. cancel()
  102. }
  103. }