stream_logs.go 2.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114
  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. )
  17. type InfraStreamLogHandler struct {
  18. handlers.PorterHandlerWriter
  19. }
  20. func NewInfraStreamLogHandler(
  21. config *config.Config,
  22. writer shared.ResultWriter,
  23. ) *InfraStreamLogHandler {
  24. return &InfraStreamLogHandler{
  25. PorterHandlerWriter: handlers.NewDefaultPorterHandler(config, nil, writer),
  26. }
  27. }
  28. func (c *InfraStreamLogHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
  29. safeRW := r.Context().Value(types.RequestCtxWebsocketKey).(*websocket.WebsocketSafeReadWriter)
  30. infra, _ := r.Context().Value(types.InfraScope).(*models.Infra)
  31. operation, _ := r.Context().Value(types.OperationScope).(*models.Operation)
  32. workspaceID := models.GetWorkspaceID(infra, operation)
  33. ctx, cancel := c.Config().ProvisionerClient.NewGRPCContext(workspaceID)
  34. defer cancel()
  35. stream, err := c.Config().ProvisionerClient.GRPCClient.GetLog(ctx, &pb.Infra{
  36. ProjectId: int64(infra.ProjectID),
  37. Id: int64(infra.ID),
  38. Suffix: infra.Suffix,
  39. })
  40. if err != nil {
  41. c.HandleAPIError(w, r, apierrors.NewErrInternal(err))
  42. return
  43. }
  44. errorchan := make(chan error)
  45. var wg sync.WaitGroup
  46. wg.Add(2)
  47. go func() {
  48. wg.Wait()
  49. close(errorchan)
  50. }()
  51. go func() {
  52. defer wg.Done()
  53. for {
  54. if _, _, err := safeRW.ReadMessage(); err != nil {
  55. errorchan <- nil
  56. return
  57. }
  58. }
  59. }()
  60. go func() {
  61. defer wg.Done()
  62. for {
  63. tfLog, err := stream.Recv()
  64. if err != nil {
  65. if err == io.EOF || errors.Is(ctx.Err(), context.Canceled) {
  66. errorchan <- nil
  67. } else {
  68. errorchan <- err
  69. }
  70. return
  71. }
  72. _, err = safeRW.Write([]byte(tfLog.Log))
  73. if err != nil {
  74. errorchan <- nil
  75. return
  76. }
  77. }
  78. }()
  79. for err = range errorchan {
  80. if err != nil {
  81. c.HandleAPIErrorNoWrite(w, r, apierrors.NewErrInternal(err))
  82. }
  83. // close the grpc stream: do not check for error case since the stream could already be
  84. // closed
  85. stream.CloseSend()
  86. // close the websocket stream: do not check for error case since the WS could already be
  87. // closed
  88. safeRW.Close()
  89. // cancel the context set for the grpc stream to ensure that Recv is unblocked
  90. cancel()
  91. }
  92. }