This repository has no description
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}