From 12ebd85199b872c871de0d967fe8eb35ec1456a1 Mon Sep 17 00:00:00 2001 From: oppiliappan Date: Wed, 20 May 2026 17:10:09 +0000 Subject: [PATCH] appview/pipelines: separate utilities out of web handler Signed-off-by: oppiliappan --- appview/pipelines/logs.go | 44 ++++++++++++++++++++++++++++++++++++++++++++ appview/pipelines/pipelines.go | 51 ++++++--------------------------------------------- 2 file(s) changed, 50 insertion(s)(+), 45 deletion(s)(-) diff --git a/appview/pipelines/logs.go b/appview/pipelines/logs.go new file mode 100644 --- /dev/null +++ b/appview/pipelines/logs.go @@ -0,0 +1,44 @@ +package pipelines + +import ( + "strings" + + "github.com/gorilla/websocket" +) + +type LogEvent struct { + Msg []byte + Err error +} + +func (ev *LogEvent) IsCloseError() bool { + return websocket.IsCloseError( + ev.Err, + websocket.CloseNormalClosure, + websocket.CloseGoingAway, + websocket.CloseAbnormalClosure, + ) +} + +func ReadLogs(conn *websocket.Conn, ch chan LogEvent) { + defer close(ch) + for { + if conn == nil { + return + } + _, msg, err := conn.ReadMessage() + if err != nil { + ch <- LogEvent{Err: err} + return + } + ch <- LogEvent{Msg: msg} + } +} + +func SpindleURL(dev bool, spindle, knot, rkey, workflow string) string { + scheme := "wss" + if dev { + scheme = "ws" + } + return scheme + "://" + strings.Join([]string{spindle, "logs", knot, rkey, workflow}, "/") +} diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go --- a/appview/pipelines/pipelines.go +++ b/appview/pipelines/pipelines.go @@ -7,7 +7,6 @@ "fmt" "log/slog" "net/http" - "strings" "time" "tangled.org/core/api/tangled" @@ -223,12 +222,7 @@ return } - scheme := "wss" - if p.config.Core.Dev { - scheme = "ws" - } - - url := scheme + "://" + strings.Join([]string{spindle, "logs", knot, rkey, workflow}, "/") + url := SpindleURL(p.config.Core.Dev, spindle, knot, rkey, workflow) l = l.With("url", url) clientConn, err := upgrader.Upgrade(w, r, nil) @@ -258,9 +252,9 @@ defer spindleConn.Close() // create a channel for incoming messages - evChan := make(chan logEvent, 100) + evChan := make(chan LogEvent, 100) // start a goroutine to read from spindle - go readLogs(spindleConn, evChan) + go ReadLogs(spindleConn, evChan) stepStartTimes := make(map[int]time.Time) var fragment bytes.Buffer @@ -275,17 +269,17 @@ continue } - if ev.err != nil && ev.isCloseError() { + if ev.Err != nil && ev.IsCloseError() { l.Debug("graceful shutdown, tail complete", "err", err) return } - if ev.err != nil { + if ev.Err != nil { l.Error("error reading from spindle", "err", err) return } var logLine spindlemodel.LogLine - if err = json.Unmarshal(ev.msg, &logLine); err != nil { + if err = json.Unmarshal(ev.Msg, &logLine); err != nil { l.Error("failed to parse logline", "err", err) continue } @@ -446,37 +440,4 @@ return } l.Debug("canceled pipeline", "uri", pipeline.AtUri()) -} - -// either a message or an error -type logEvent struct { - msg []byte - err error -} - -func (ev *logEvent) isCloseError() bool { - return websocket.IsCloseError( - ev.err, - websocket.CloseNormalClosure, - websocket.CloseGoingAway, - websocket.CloseAbnormalClosure, - ) -} - -// read logs from spindle and pass through to chan -func readLogs(conn *websocket.Conn, ch chan logEvent) { - defer close(ch) - - for { - if conn == nil { - return - } - - _, msg, err := conn.ReadMessage() - if err != nil { - ch <- logEvent{err: err} - return - } - ch <- logEvent{msg: msg} - } } -- tangled.sh