This repository has no description
0

Configure Feed

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

core / spindle / stream.go
3.6 kB 158 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 95 96 config := tail.Config{ 97 Follow: !isFinished, 98 ReOpen: !isFinished, 99 MustExist: false, 100 Location: &tail.SeekInfo{ 101 Offset: 0, 102 Whence: io.SeekStart, 103 }, 104 // Logger: tail.DiscardingLogger, 105 } 106 107 t, err := tail.TailFile(filePath, config) 108 if err != nil { 109 return fmt.Errorf("failed to tail log file: %w", err) 110 } 111 defer t.Stop() 112 113 for { 114 select { 115 case <-ctx.Done(): 116 return ctx.Err() 117 case line := <-t.Lines: 118 if line == nil && isFinished { 119 return fmt.Errorf("tail completed") 120 } 121 122 if line == nil { 123 return fmt.Errorf("tail channel closed unexpectedly") 124 } 125 126 if line.Err != nil { 127 return fmt.Errorf("error tailing log file: %w", line.Err) 128 } 129 130 if err := conn.WriteMessage(websocket.TextMessage, []byte(line.Text)); err != nil { 131 return fmt.Errorf("failed to write to websocket: %w", err) 132 } 133 case <-time.After(30 * time.Second): 134 // send a keep-alive 135 if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { 136 return fmt.Errorf("failed to write control: %w", err) 137 } 138 } 139 } 140} 141 142func getWorkflowID(r *http.Request) (models.WorkflowId, error) { 143 knot := chi.URLParam(r, "knot") 144 rkey := chi.URLParam(r, "rkey") 145 name := chi.URLParam(r, "name") 146 147 if knot == "" || rkey == "" || name == "" { 148 return models.WorkflowId{}, fmt.Errorf("missing required parameters") 149 } 150 151 return models.WorkflowId{ 152 PipelineId: models.PipelineId{ 153 Knot: knot, 154 Rkey: rkey, 155 }, 156 Name: name, 157 }, nil 158}