stream_logs.go 3.0 KB

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