forked from
tangled.org/core
Monorepo for Tangled
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}