From 9243b43ccb14a1993113901fbf44871a994c16bf Mon Sep 17 00:00:00 2001 From: oppiliappan Date: Mon, 23 Jun 2025 09:24:08 +0000 Subject: [PATCH] appview,spindle: handle graceful shutdown of websockets Signed-off-by: oppiliappan --- spindle/stream.go | 36 +++++++++++++++++++++++++++++------- appview/pipelines/pipelines.go | 29 ++++++++++++++++++++++++++--- spindle/db/events.go | 34 ++++++++++++++++++++++++++++++++++ spindle/engine/engine.go | 6 +++++- spindle/engine/logger.go | 28 +++++----------------------- spindle/models/models.go | 1 + 6 file(s) changed, 100 insertion(s)(+), 34 deletion(s)(-) diff --git a/spindle/stream.go b/spindle/stream.go --- a/spindle/stream.go +++ b/spindle/stream.go @@ -4,6 +4,7 @@ "context" "encoding/json" "fmt" + "io" "net/http" "strconv" "time" @@ -105,7 +106,14 @@ http.Error(w, "failed to upgrade", http.StatusInternalServerError) return } - defer conn.Close() + defer func() { + _ = conn.WriteControl( + websocket.CloseMessage, + websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), + time.Now().Add(time.Second), + ) + conn.Close() + }() l.Debug("upgraded http to wss") ctx, cancel := context.WithCancel(r.Context()) @@ -122,20 +130,30 @@ }() if err := s.streamLogsFromDisk(ctx, conn, wid); err != nil { - l.Error("log stream failed", "err", err) + l.Info("log stream ended", "err", err) } - l.Debug("logs connection closed") + + l.Info("logs connection closed") } func (s *Spindle) streamLogsFromDisk(ctx context.Context, conn *websocket.Conn, wid models.WorkflowId) error { + status, err := s.db.GetStatus(wid) + if err != nil { + return err + } + isFinished := models.StatusKind(status.Status).IsFinish() + filePath := engine.LogFilePath(s.cfg.Pipelines.LogDir, wid) config := tail.Config{ - Follow: true, - ReOpen: true, + Follow: !isFinished, + ReOpen: !isFinished, MustExist: false, - Location: &tail.SeekInfo{Offset: 0, Whence: 0}, - Logger: tail.DiscardingLogger, + Location: &tail.SeekInfo{ + Offset: 0, + Whence: io.SeekStart, + }, + // Logger: tail.DiscardingLogger, } t, err := tail.TailFile(filePath, config) @@ -149,6 +167,10 @@ case <-ctx.Done(): return ctx.Err() case line := <-t.Lines: + if line == nil && isFinished { + return fmt.Errorf("tail completed") + } + if line == nil { return fmt.Errorf("tail channel closed unexpectedly") } diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go --- a/appview/pipelines/pipelines.go +++ b/appview/pipelines/pipelines.go @@ -158,7 +158,14 @@ l.Error("websocket upgrade failed", "err", err) return } - defer clientConn.Close() + defer func() { + _ = clientConn.WriteControl( + websocket.CloseMessage, + websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), + time.Now().Add(time.Second), + ) + clientConn.Close() + }() ctx, cancel := context.WithCancel(r.Context()) defer cancel() @@ -236,10 +243,19 @@ // start a goroutine to read from spindle go func() { defer close(msgChan) + defer close(errChan) + for { _, msg, err := spindleConn.ReadMessage() if err != nil { - errChan <- err + if websocket.IsCloseError(err, + websocket.CloseNormalClosure, + websocket.CloseGoingAway, + websocket.CloseAbnormalClosure) { + errChan <- nil // signal graceful end + } else { + errChan <- err + } return } msgChan <- msg @@ -252,7 +268,14 @@ l.Info("client disconnected") return case err := <-errChan: - l.Error("error reading from spindle", "err", err) + if err != nil { + l.Error("error reading from spindle", "err", err) + } + + if err == nil { + l.Info("log tail complete") + } + return case msg := <-msgChan: var logLine spindlemodel.LogLine diff --git a/spindle/db/events.go b/spindle/db/events.go --- a/spindle/db/events.go +++ b/spindle/db/events.go @@ -120,6 +120,40 @@ } +func (d *DB) GetStatus(workflowId models.WorkflowId) (*tangled.PipelineStatus, error) { + pipelineAtUri := workflowId.PipelineId.AtUri() + + var eventJson string + err := d.QueryRow( + ` + select + event from events + where + nsid = ? + and json_extract(event, '$.pipeline') = ? + and json_extract(event, '$.workflow') = ? + order by + created desc + limit + 1 + `, + tangled.PipelineStatusNSID, + string(pipelineAtUri), + workflowId.Name, + ).Scan(&eventJson) + + if err != nil { + return nil, err + } + + var status tangled.PipelineStatus + if err := json.Unmarshal([]byte(eventJson), &status); err != nil { + return nil, err + } + + return &status, nil +} + func (d *DB) StatusPending(workflowId models.WorkflowId, n *notifier.Notifier) error { return d.createStatusEvent(workflowId, models.StatusKindPending, nil, nil, n) } diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -326,7 +326,11 @@ } defer wfLogger.Close() - _, err = stdcopy.StdCopy(wfLogger.Stdout(), wfLogger.Stderr(), logs) + _, err = stdcopy.StdCopy( + wfLogger.Writer("stdout", stepIdx), + wfLogger.Writer("stderr", stepIdx), + logs, + ) if err != nil && err != io.EOF && !errors.Is(err, context.DeadlineExceeded) { return fmt.Errorf("failed to copy logs: %w", err) } diff --git a/spindle/engine/logger.go b/spindle/engine/logger.go --- a/spindle/engine/logger.go +++ b/spindle/engine/logger.go @@ -17,14 +17,9 @@ } func NewWorkflowLogger(baseDir string, wid models.WorkflowId) (*WorkflowLogger, error) { - dir := filepath.Join(baseDir, wid.String()) - if err := os.MkdirAll(dir, 0755); err != nil { - return nil, fmt.Errorf("creating log dir: %w", err) - } - path := LogFilePath(baseDir, wid) - file, err := os.Create(path) + file, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644) if err != nil { return nil, fmt.Errorf("creating log file: %w", err) } @@ -43,33 +38,19 @@ return l.file.Close() } -func OpenLogFile(baseDir string, workflowID models.WorkflowId) (*os.File, error) { - logPath := LogFilePath(baseDir, workflowID) - - file, err := os.Open(logPath) - if err != nil { - return nil, fmt.Errorf("error opening log file: %w", err) - } - - return file, nil -} - func LogFilePath(baseDir string, workflowID models.WorkflowId) string { logFilePath := filepath.Join(baseDir, fmt.Sprintf("%s.log", workflowID.String())) return logFilePath } -func (l *WorkflowLogger) Stdout() io.Writer { - return &jsonWriter{logger: l, stream: "stdout"} -} - -func (l *WorkflowLogger) Stderr() io.Writer { - return &jsonWriter{logger: l, stream: "stderr"} +func (l *WorkflowLogger) Writer(stream string, stepId int) io.Writer { + return &jsonWriter{logger: l, stream: stream, stepId: stepId} } type jsonWriter struct { logger *WorkflowLogger stream string + stepId int } func (w *jsonWriter) Write(p []byte) (int, error) { @@ -78,6 +59,7 @@ entry := models.LogLine{ Stream: w.stream, Data: line, + StepId: w.stepId, } if err := w.logger.encoder.Encode(entry); err != nil { diff --git a/spindle/models/models.go b/spindle/models/models.go --- a/spindle/models/models.go +++ b/spindle/models/models.go @@ -74,4 +74,5 @@ type LogLine struct { Stream string `json:"s"` Data string `json:"d"` + StepId int `json:"i"` } -- tangled.sh