Monorepo for Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / stream.go
3.6 kB 156 lines
1package spindle 2 3import ( 4 "context" 5 "errors" 6 "fmt" 7 "io" 8 "net/http" 9 "time" 10 11 "tangled.org/core/eventstream" 12 "tangled.org/core/log" 13 "tangled.org/core/spindle/models" 14 15 "github.com/go-chi/chi/v5" 16 "github.com/gorilla/websocket" 17 "github.com/hpcloud/tail" 18) 19 20var upgrader = websocket.Upgrader{ 21 ReadBufferSize: 1024, 22 WriteBufferSize: 1024, 23} 24 25func (s *Spindle) Events(w http.ResponseWriter, r *http.Request) { 26 l := log.SubLogger(s.l, "eventstream") 27 l.Debug("received new connection") 28 29 err := eventstream.Stream(w, r, eventstream.StreamConfig{ 30 Backend: s.db, 31 Notifier: s.n, 32 Logger: l, 33 }) 34 if err != nil && !errors.Is(err, eventstream.ErrDrainCap) { 35 l.Error("event stream ended with error", "err", err) 36 } 37} 38 39func (s *Spindle) Logs(w http.ResponseWriter, r *http.Request) { 40 wid, err := getWorkflowID(r) 41 if err != nil { 42 http.Error(w, err.Error(), http.StatusBadRequest) 43 return 44 } 45 46 l := s.l.With("handler", "Logs") 47 l = s.l.With("wid", wid) 48 49 conn, err := upgrader.Upgrade(w, r, nil) 50 if err != nil { 51 l.Error("websocket upgrade failed", "err", err) 52 http.Error(w, "failed to upgrade", http.StatusInternalServerError) 53 return 54 } 55 defer func() { 56 _ = conn.WriteControl( 57 websocket.CloseMessage, 58 websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), 59 time.Now().Add(time.Second), 60 ) 61 conn.Close() 62 }() 63 l.Debug("upgraded http to wss") 64 65 ctx, cancel := context.WithCancel(r.Context()) 66 defer cancel() 67 68 go func() { 69 for { 70 if _, _, err := conn.NextReader(); err != nil { 71 l.Debug("client disconnected", "err", err) 72 cancel() 73 return 74 } 75 } 76 }() 77 78 if err := s.streamLogsFromDisk(ctx, conn, wid); err != nil { 79 l.Info("log stream ended", "err", err) 80 } 81 82 l.Info("logs connection closed") 83} 84 85func (s *Spindle) streamLogsFromDisk(ctx context.Context, conn *websocket.Conn, wid models.WorkflowId) error { 86 status, err := s.db.GetStatus(wid) 87 if err != nil { 88 return err 89 } 90 isFinished := models.StatusKind(status.Status).IsFinish() 91 92 filePath := models.LogFilePath(s.cfg.Server.LogDir, wid) 93 94 config := tail.Config{ 95 Follow: !isFinished, 96 ReOpen: !isFinished, 97 MustExist: false, 98 Location: &tail.SeekInfo{ 99 Offset: 0, 100 Whence: io.SeekStart, 101 }, 102 // Logger: tail.DiscardingLogger, 103 } 104 105 t, err := tail.TailFile(filePath, config) 106 if err != nil { 107 return fmt.Errorf("failed to tail log file: %w", err) 108 } 109 defer t.Stop() 110 111 for { 112 select { 113 case <-ctx.Done(): 114 return ctx.Err() 115 case line := <-t.Lines: 116 if line == nil && isFinished { 117 return fmt.Errorf("tail completed") 118 } 119 120 if line == nil { 121 return fmt.Errorf("tail channel closed unexpectedly") 122 } 123 124 if line.Err != nil { 125 return fmt.Errorf("error tailing log file: %w", line.Err) 126 } 127 128 if err := conn.WriteMessage(websocket.TextMessage, []byte(line.Text)); err != nil { 129 return fmt.Errorf("failed to write to websocket: %w", err) 130 } 131 case <-time.After(30 * time.Second): 132 // send a keep-alive 133 if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { 134 return fmt.Errorf("failed to write control: %w", err) 135 } 136 } 137 } 138} 139 140func getWorkflowID(r *http.Request) (models.WorkflowId, error) { 141 knot := chi.URLParam(r, "knot") 142 rkey := chi.URLParam(r, "rkey") 143 name := chi.URLParam(r, "name") 144 145 if knot == "" || rkey == "" || name == "" { 146 return models.WorkflowId{}, fmt.Errorf("missing required parameters") 147 } 148 149 return models.WorkflowId{ 150 PipelineId: models.PipelineId{ 151 Knot: knot, 152 Rkey: rkey, 153 }, 154 Name: name, 155 }, nil 156}