From 8a3a3f4aa331537f454beb20cdc0287dcd43741b Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Sun, 14 Jun 2026 05:39:12 +0900 Subject: [PATCH] wip: appview/pipelines: use new ci models Signed-off-by: Seongmin Lee --- appview/pipelines/logs.go | 60 ++++------- appview/pipelines/pipelines.go | 185 ++++++++++++++++++++------------- 2 files changed, 134 insertions(+), 111 deletions(-) diff --git a/appview/pipelines/logs.go b/appview/pipelines/logs.go index 0c035d31..c96902d8 100644 --- a/appview/pipelines/logs.go +++ b/appview/pipelines/logs.go @@ -1,18 +1,15 @@ package pipelines import ( - "fmt" + "errors" "html/template" - "net/url" "regexp" "strings" + "time" - "github.com/bluesky-social/indigo/atproto/syntax" terminal "github.com/buildkite/terminal-to-html/v3" "github.com/gorilla/websocket" - "tangled.org/core/api/tangled" "tangled.org/core/appview/pages/markup" - "tangled.org/core/hostutil" ) // matches any ANSI escape sequence: ESC [ m @@ -55,43 +52,32 @@ func (a *ansiState) Render(line string) template.HTML { return template.HTML(sanitized) } -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 +// isExpectedClose reports whether err is a clean websocket close (or nil). +func isExpectedClose(err error) bool { + if err == nil { + return true + } + var ce *websocket.CloseError + if errors.As(err, &ce) { + switch ce.Code { + case websocket.CloseNormalClosure, websocket.CloseGoingAway, websocket.CloseAbnormalClosure: + return true } - ch <- LogEvent{Msg: msg} } + return false } -func SubscribeLogsUrl(spindle string, pipeline syntax.TID, workflow string) string { - u, err := hostutil.EnsureWsScheme(spindle) - if err != nil { +func derefStr(s *string) string { + if s == nil { return "" } + return *s +} - query := url.Values{} - query.Set("pipeline", pipeline.String()) - query.Set("workflow", workflow) - return fmt.Sprintf("%s/xrpc/%s?%s", u, tangled.CiPipelineSubscribeLogsNSID, query.Encode()) +func parseRFC3339(s string) time.Time { + t, err := time.Parse(time.RFC3339, s) + if err != nil { + return time.Time{} + } + return t } diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go index 2387bc2d..269ca733 100644 --- a/appview/pipelines/pipelines.go +++ b/appview/pipelines/pipelines.go @@ -3,9 +3,9 @@ package pipelines import ( "bytes" "context" - "encoding/json" "log/slog" "net/http" + "sync" "time" "tangled.org/core/api/tangled" @@ -17,9 +17,9 @@ import ( "tangled.org/core/appview/reporesolver" "tangled.org/core/hostutil" "tangled.org/core/idresolver" + "tangled.org/core/lexutil" "tangled.org/core/orm" "tangled.org/core/rbac" - spindlemodel "tangled.org/core/spindle/models" "github.com/bluesky-social/indigo/atproto/syntax" indigoxrpc "github.com/bluesky-social/indigo/xrpc" @@ -180,6 +180,25 @@ var upgrader = websocket.Upgrader{ WriteBufferSize: 1024, } +type webLogScheduler struct { + ch chan *tangled.CiPipelineSubscribeLogs_Event +} + +var _ lexutil.Scheduler[tangled.CiPipelineSubscribeLogs_Event] = (*webLogScheduler)(nil) + +// AddWork implements [lexutil.Scheduler]. +func (w *webLogScheduler) AddWork(ctx context.Context, _ string, val *tangled.CiPipelineSubscribeLogs_Event) error { + select { + case w.ch <- val: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +// Shutdown implements [lexutil.Scheduler]. +func (w *webLogScheduler) Shutdown() { close(w.ch) } + func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { l := p.logger.With("handler", "logs") @@ -209,44 +228,72 @@ func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { return } - logsUrl := SubscribeLogsUrl(f.Spindle, pipelineId, workflowName) - if logsUrl == "" { - http.Error(w, "invalid spindle hostname", http.StatusBadRequest) - return - } - l = l.With("url", logsUrl) - clientConn, err := upgrader.Upgrade(w, r, nil) if err != nil { l.Error("websocket upgrade failed", "err", err) return } - defer func() { - _ = clientConn.WriteControl( - websocket.CloseMessage, - websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), - time.Now().Add(time.Second), - ) - clientConn.Close() - }() + defer clientConn.Close() ctx, cancel := context.WithCancel(r.Context()) defer cancel() - l.Info("logs endpoint hit") + evChan := make(chan *tangled.CiPipelineSubscribeLogs_Event, 100) + done := make(chan error, 1) + sched := &webLogScheduler{ch: evChan} + xrpcc := &lexutil.Client{Client: indigoxrpc.Client{Host: f.Spindle}} + go func() { + done <- tangled.CiPipelineSubscribeLogs(ctx, xrpcc, pipelineId.String(), []string{workflowName}, sched) + }() - spindleConn, _, err := websocket.DefaultDialer.Dial(logsUrl, nil) - if err != nil { - l.Error("websocket dial failed", "err", err) - return - } - defer spindleConn.Close() + var lastWriteLk sync.Mutex + lastWrite := time.Now() + + // Start a goroutine to ping the client periodically to check if it's still + // alive. If the client doesn't respond to a ping within 5 seconds, we'll + // close the connection and teardown the consumer. + go func() { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + for { + select { + case <-ticker.C: + lastWriteLk.Lock() + lw := lastWrite + lastWriteLk.Unlock() + if time.Since(lw) < 30*time.Second { + continue + } + if err := clientConn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil { + l.Warn("failed to ping client", "err", err) + cancel() + return + } + case <-ctx.Done(): + return + } + } + }() - // create a channel for incoming messages - evChan := make(chan LogEvent, 100) - // start a goroutine to read from spindle - go ReadLogs(spindleConn, evChan) + clientConn.SetPingHandler(func(message string) error { + err := clientConn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(60*time.Second)) + if err == websocket.ErrCloseSent { + return nil + } + return err + }) + + // Start a goroutine to read messages from the client and discard them. + go func() { + for { + if _, _, err := clientConn.ReadMessage(); err != nil { + cancel() + return + } + } + }() + // Main loop: sole writer of data frames to the client. stepStartTimes := make(map[int]time.Time) stepAnsi := make(map[int]*ansiState) var fragment bytes.Buffer @@ -258,61 +305,55 @@ func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { case ev, ok := <-evChan: if !ok { - continue - } - - if ev.Err != nil && ev.IsCloseError() { - l.Debug("graceful shutdown, tail complete", "err", err) - return - } - if ev.Err != nil { - l.Error("error reading from spindle", "err", err) + // Stream ended: Shutdown closed the upstream channel. + if err := <-done; !isExpectedClose(err) { + l.Error("spindle stream error", "err", err) + } return } - var logLine spindlemodel.LogLine - if err = json.Unmarshal(ev.Msg, &logLine); err != nil { - l.Error("failed to parse logline", "err", err) - continue - } - fragment.Reset() - switch logLine.Kind { - case spindlemodel.LogKindControl: - switch logLine.StepStatus { - case spindlemodel.StepStatusStart: - stepStartTimes[logLine.StepId] = logLine.Time - collapsed := false - if logLine.StepKind == spindlemodel.StepKindSystem { - collapsed = true - } + switch { + case ev.Error != nil: + l.Error("spindle error frame", "err", ev.Error.Error, "msg", ev.Error.Message) + return + + case ev.Control != nil: + c := ev.Control + step := int(c.Step) + switch derefStr(c.Status) { + case "start": + t := parseRFC3339(c.Time) + stepStartTimes[step] = t + // "system" steps are injected by the CI runner; collapse them. + collapsed := derefStr(c.Kind) == "system" err = p.pages.LogBlock(&fragment, pages.LogBlockParams{ - Id: logLine.StepId, - Name: logLine.Content, - Command: logLine.StepCommand, + Id: step, + Name: c.Content, + Command: derefStr(c.Command), Collapsed: collapsed, - StartTime: logLine.Time, + StartTime: t, }) - case spindlemodel.StepStatusEnd: - startTime := stepStartTimes[logLine.StepId] - endTime := logLine.Time + case "end": err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{ - Id: logLine.StepId, - StartTime: startTime, - EndTime: endTime, + Id: step, + StartTime: stepStartTimes[step], + EndTime: parseRFC3339(c.Time), }) } - case spindlemodel.LogKindData: - ansi, ok := stepAnsi[logLine.StepId] + case ev.Data != nil: + d := ev.Data + step := int(d.Step) + ansi, ok := stepAnsi[step] if !ok { ansi = NewAnsiState() - stepAnsi[logLine.StepId] = ansi + stepAnsi[step] = ansi } err = p.pages.LogLine(&fragment, pages.LogLineParams{ - Id: logLine.StepId, - Content: ansi.Render(logLine.Content), + Id: step, + Content: ansi.Render(d.Content), }) } if err != nil { @@ -324,13 +365,9 @@ func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { l.Error("error writing to client", "err", err) return } - - case <-time.After(30 * time.Second): - l.Debug("sent keepalive") - if err = clientConn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { - l.Error("failed to write control", "err", err) - return - } + lastWriteLk.Lock() + lastWrite = time.Now() + lastWriteLk.Unlock() } } } -- 2.51.2