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