stream_logs.go 2.7 KB

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