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