From 0df0d697af6899e7c3f59617e11eeadd18334983 Mon Sep 17 00:00:00 2001 From: oppiliappan Date: Fri, 20 Jun 2025 14:19:45 +0100 Subject: [PATCH] appview: stream logs from workflow endpoint Signed-off-by: oppiliappan --- appview/db/pipeline.go | 32 ++-- appview/pages/pages.go | 1 + appview/pages/repoinfo/repoinfo.go | 1 + appview/pages/templates/layouts/base.html | 1 + .../pipelines/fragments/pipelineSymbol.html | 17 +- .../repo/pipelines/fragments/tooltip.html | 5 + .../templates/repo/pipelines/pipelines.html | 32 ++-- .../templates/repo/pipelines/workflow.html | 53 +++--- appview/pipelines/pipelines.go | 152 +++++++++++++++++- appview/pipelines/router.go | 1 + appview/reporesolver/resolver.go | 1 + flake.nix | 2 +- 12 files changed, 254 insertions(+), 44 deletions(-) diff --git a/appview/db/pipeline.go b/appview/db/pipeline.go index 9c76806f..142af306 100644 --- a/appview/db/pipeline.go +++ b/appview/db/pipeline.go @@ -27,17 +27,17 @@ type Pipeline struct { } type WorkflowStatus struct { - data []PipelineStatus + Data []PipelineStatus } func (w WorkflowStatus) Latest() PipelineStatus { - return w.data[len(w.data)-1] + return w.Data[len(w.Data)-1] } // time taken by this workflow to reach an "end state" func (w WorkflowStatus) TimeTaken() time.Duration { var start, end *time.Time - for _, s := range w.data { + for _, s := range w.Data { if s.Status.IsStart() { start = &s.Created } @@ -78,6 +78,11 @@ func (p Pipeline) Workflows() []string { return ws } +// if we know that a spindle has picked up this pipeline, then it is Responding +func (p Pipeline) IsResponding() bool { + return len(p.Statuses) != 0 +} + type Trigger struct { Id int Kind string @@ -256,6 +261,7 @@ func AddPipelineStatus(e Execer, status PipelineStatus) error { status.Status, status.Error, status.ExitCode, + status.Created.Format(time.RFC3339), } placeholders := make([]string, len(args)) @@ -272,7 +278,8 @@ func AddPipelineStatus(e Execer, status PipelineStatus) error { workflow, status, error, - exit_code + exit_code, + created ) values (%s) `, strings.Join(placeholders, ",")) @@ -355,13 +362,11 @@ func GetPipelineStatuses(e Execer, filters ...filter) ([]Pipeline, error) { return nil, err } - // Parse created time manually p.Created, err = time.Parse(time.RFC3339, created) if err != nil { return nil, fmt.Errorf("invalid pipeline created timestamp %q: %w", created, err) } - // Link trigger to pipeline t.Id = p.TriggerId p.Trigger = &t p.Statuses = make(map[string]WorkflowStatus) @@ -440,7 +445,7 @@ func GetPipelineStatuses(e Execer, filters ...filter) ([]Pipeline, error) { } // append - statuses.data = append(statuses.data, ps) + statuses.Data = append(statuses.Data, ps) // reassign pipeline.Statuses[ps.Workflow] = statuses @@ -450,11 +455,20 @@ func GetPipelineStatuses(e Execer, filters ...filter) ([]Pipeline, error) { var all []Pipeline for _, p := range pipelines { for _, s := range p.Statuses { - slices.SortFunc(s.data, func(a, b PipelineStatus) int { + slices.SortFunc(s.Data, func(a, b PipelineStatus) int { if a.Created.After(b.Created) { return 1 } - return -1 + if a.Created.Before(b.Created) { + return -1 + } + if a.ID > b.ID { + return 1 + } + if a.ID < b.ID { + return -1 + } + return 0 }) } all = append(all, p) diff --git a/appview/pages/pages.go b/appview/pages/pages.go index cecad95d..8ef647eb 100644 --- a/appview/pages/pages.go +++ b/appview/pages/pages.go @@ -968,6 +968,7 @@ type WorkflowParams struct { RepoInfo repoinfo.RepoInfo Pipeline db.Pipeline Workflow string + LogUrl string Active string } diff --git a/appview/pages/repoinfo/repoinfo.go b/appview/pages/repoinfo/repoinfo.go index 0c9dcb31..40cb5ad8 100644 --- a/appview/pages/repoinfo/repoinfo.go +++ b/appview/pages/repoinfo/repoinfo.go @@ -56,6 +56,7 @@ type RepoInfo struct { OwnerHandle string Description string Knot string + Spindle string RepoAt syntax.ATURI IsStarred bool Stats db.RepoStats diff --git a/appview/pages/templates/layouts/base.html b/appview/pages/templates/layouts/base.html index 1501f5f9..a0ce2b21 100644 --- a/appview/pages/templates/layouts/base.html +++ b/appview/pages/templates/layouts/base.html @@ -9,6 +9,7 @@ /> + {{ block "title" . }}{{ end }} ยท tangled {{ block "extrameta" . }}{{ end }} diff --git a/appview/pages/templates/repo/pipelines/fragments/pipelineSymbol.html b/appview/pages/templates/repo/pipelines/fragments/pipelineSymbol.html index fab91ddc..1b2e621d 100644 --- a/appview/pages/templates/repo/pipelines/fragments/pipelineSymbol.html +++ b/appview/pages/templates/repo/pipelines/fragments/pipelineSymbol.html @@ -4,13 +4,26 @@ {{ $statuses := .Statuses }} {{ $total := len $statuses }} {{ $success := index $c "success" }} + {{ $fail := index $c "failed" }} + {{ $empty := eq $total 0 }} {{ $allPass := eq $success $total }} + {{ $allFail := eq $fail $total }} - {{ if $allPass }} + {{ if $empty }}
- {{ i "check" "size-4 text-green-600 dark:text-green-400 " }} + {{ i "hourglass" "size-4 text-gray-600 dark:text-gray-400 " }} + 0/{{ $total }} +
+ {{ else if $allPass }} +
+ {{ i "check" "size-4 text-green-600" }} {{ $total }}/{{ $total }}
+ {{ else if $allFail }} +
+ {{ i "x" "size-4 text-red-600" }} + 0/{{ $total }} +
{{ else }} {{ $radius := f64 8 }} {{ $circumference := mulf64 2.0 (mulf64 3.1416 $radius) }} diff --git a/appview/pages/templates/repo/pipelines/fragments/tooltip.html b/appview/pages/templates/repo/pipelines/fragments/tooltip.html index c313f87f..ce3ca10a 100644 --- a/appview/pages/templates/repo/pipelines/fragments/tooltip.html +++ b/appview/pages/templates/repo/pipelines/fragments/tooltip.html @@ -23,6 +23,11 @@ + {{ else }} +
+ {{ i "hourglass" "size-4" }} + Waiting for spindle ... +
{{ end }} diff --git a/appview/pages/templates/repo/pipelines/pipelines.html b/appview/pages/templates/repo/pipelines/pipelines.html index 69811420..06cb12b8 100644 --- a/appview/pages/templates/repo/pipelines/pipelines.html +++ b/appview/pages/templates/repo/pipelines/pipelines.html @@ -41,15 +41,17 @@ {{ $root := index . 0 }} {{ $p := index . 1 }} {{ with $p }} -
-
+
+
{{ $target := .Trigger.TargetRef }} {{ $workflows := .Workflows }} + {{ $link := "" }} + {{ if .IsResponding }} + {{ $link = printf "/%s/pipelines/%s/workflow/%d" $root.RepoInfo.FullName .Id (index $workflows 0) }} + {{ end }} {{ if .Trigger.IsPush }} - - {{ $target }} - push - + {{ $target }} + push
-
+
{{ template "repo/pipelines/fragments/pipelineSymbolLong" . }}
-
+
{{ $t := .TimeTaken }} -
+
{{ if $t }} {{ else }} {{ end }}
+ +
+ {{ if $link }} + + {{ i "arrow-up-right" "size-4" }} + + + {{ end }} +
+
{{ end }} {{ end }} diff --git a/appview/pages/templates/repo/pipelines/workflow.html b/appview/pages/templates/repo/pipelines/workflow.html index ddc0ee84..1f0d3872 100644 --- a/appview/pages/templates/repo/pipelines/workflow.html +++ b/appview/pages/templates/repo/pipelines/workflow.html @@ -23,30 +23,47 @@ {{ define "sidebar" }} {{ $active := .Workflow }} {{ with .Pipeline }} -
+ {{ $id := .Id }} +
{{ range $name, $all := .Statuses }} - {{ end }} {{ end }} + +{{ define "logs" }} +
+
+ + +
+
+{{ end }} diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go index 58c7a591..8919b1fa 100644 --- a/appview/pipelines/pipelines.go +++ b/appview/pipelines/pipelines.go @@ -1,8 +1,13 @@ package pipelines import ( + "context" + "encoding/json" + "fmt" "log/slog" "net/http" + "strings" + "time" "tangled.sh/tangled.sh/core/appview/config" "tangled.sh/tangled.sh/core/appview/db" @@ -13,8 +18,10 @@ import ( "tangled.sh/tangled.sh/core/eventconsumer" "tangled.sh/tangled.sh/core/log" "tangled.sh/tangled.sh/core/rbac" + spindlemodel "tangled.sh/tangled.sh/core/spindle/models" "github.com/go-chi/chi/v5" + "github.com/gorilla/websocket" "github.com/posthog/posthog-go" ) @@ -28,7 +35,7 @@ type Pipelines struct { db *db.DB enforcer *rbac.Enforcer posthog posthog.Client - Logger *slog.Logger + logger *slog.Logger } func New( @@ -53,13 +60,13 @@ func New( db: db, posthog: posthog, enforcer: enforcer, - Logger: logger, + logger: logger, } } func (p *Pipelines) Index(w http.ResponseWriter, r *http.Request) { user := p.oauth.GetUser(r) - l := p.Logger.With("handler", "Index") + l := p.logger.With("handler", "Index") f, err := p.repoResolver.Resolve(r) if err != nil { @@ -89,7 +96,7 @@ func (p *Pipelines) Index(w http.ResponseWriter, r *http.Request) { func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { user := p.oauth.GetUser(r) - l := p.Logger.With("handler", "Workflow") + l := p.logger.With("handler", "Workflow") f, err := p.repoResolver.Resolve(r) if err != nil { @@ -106,7 +113,7 @@ func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { } workflow := chi.URLParam(r, "workflow") - if pipelineId == "" { + if workflow == "" { l.Error("empty workflow name") return } @@ -137,3 +144,138 @@ func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { Workflow: workflow, }) } + +var upgrader = websocket.Upgrader{ + ReadBufferSize: 1024, + WriteBufferSize: 1024, +} + +func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { + l := p.logger.With("handler", "logs") + + clientConn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + l.Error("websocket upgrade failed", "err", err) + return + } + defer clientConn.Close() + + ctx, cancel := context.WithCancel(r.Context()) + defer cancel() + go func() { + for { + if _, _, err := clientConn.NextReader(); err != nil { + l.Error("failed to read", "err", err) + cancel() + return + } + } + }() + + user := p.oauth.GetUser(r) + f, err := p.repoResolver.Resolve(r) + if err != nil { + l.Error("failed to get repo and knot", "err", err) + http.Error(w, "bad repo/knot", http.StatusBadRequest) + return + } + + repoInfo := f.RepoInfo(user) + + pipelineId := chi.URLParam(r, "pipeline") + workflow := chi.URLParam(r, "workflow") + if pipelineId == "" || workflow == "" { + http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) + return + } + + ps, err := db.GetPipelineStatuses( + p.db, + db.FilterEq("repo_owner", repoInfo.OwnerDid), + db.FilterEq("repo_name", repoInfo.Name), + db.FilterEq("knot", repoInfo.Knot), + db.FilterEq("id", pipelineId), + ) + if err != nil || len(ps) != 1 { + l.Error("pipeline query failed", "err", err, "count", len(ps)) + http.Error(w, "pipeline not found", http.StatusNotFound) + return + } + + singlePipeline := ps[0] + spindle := repoInfo.Spindle + knot := repoInfo.Knot + rkey := singlePipeline.Rkey + + if spindle == "" || knot == "" || rkey == "" { + http.Error(w, "invalid repo info", http.StatusBadRequest) + return + } + + scheme := "wss" + if p.config.Core.Dev { + scheme = "ws" + } + + url := scheme + "://" + strings.Join([]string{spindle, "logs", knot, rkey, workflow}, "/") + l = l.With("url", url) + l.Info("logs endpoint hit") + + spindleConn, _, err := websocket.DefaultDialer.Dial(url, nil) + if err != nil { + l.Error("websocket dial failed", "err", err) + http.Error(w, "failed to connect to log stream", http.StatusBadGateway) + return + } + defer spindleConn.Close() + + // create a channel for incoming messages + msgChan := make(chan []byte, 10) + errChan := make(chan error, 1) + + // start a goroutine to read from spindle + go func() { + defer close(msgChan) + for { + _, msg, err := spindleConn.ReadMessage() + if err != nil { + errChan <- err + return + } + msgChan <- msg + } + }() + + for { + select { + case <-ctx.Done(): + l.Info("client disconnected") + return + case err := <-errChan: + l.Error("error reading from spindle", "err", err) + return + case msg := <-msgChan: + var logLine spindlemodel.LogLine + if err = json.Unmarshal(msg, &logLine); err != nil { + l.Error("failed to parse logline", "err", err) + continue + } + + html := fmt.Appendf(nil, ` +
+

%s: %s

+
+ `, logLine.Stream, logLine.Data) + + if err = clientConn.WriteMessage(websocket.TextMessage, html); err != nil { + 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) + } + } + } +} diff --git a/appview/pipelines/router.go b/appview/pipelines/router.go index e0a0947f..a54e7117 100644 --- a/appview/pipelines/router.go +++ b/appview/pipelines/router.go @@ -11,6 +11,7 @@ func (p *Pipelines) Router(mw *middleware.Middleware) http.Handler { r := chi.NewRouter() r.Get("/", p.Index) r.Get("/{pipeline}/workflow/{workflow}", p.Workflow) + r.Get("/{pipeline}/workflow/{workflow}/logs", p.Logs) return r } diff --git a/appview/reporesolver/resolver.go b/appview/reporesolver/resolver.go index 1d692239..8efdced4 100644 --- a/appview/reporesolver/resolver.go +++ b/appview/reporesolver/resolver.go @@ -251,6 +251,7 @@ func (f *ResolvedRepo) RepoInfo(user *oauth.User) repoinfo.RepoInfo { Ref: f.Ref, IsStarred: isStarred, Knot: knot, + Spindle: f.Spindle, Roles: f.RolesInRepo(user), Stats: db.RepoStats{ StarCount: starCount, diff --git a/flake.nix b/flake.nix index 4530767c..b3a99161 100644 --- a/flake.nix +++ b/flake.nix @@ -59,7 +59,7 @@ inherit (gitignore.lib) gitignoreSource; in { overlays.default = final: prev: let - goModHash = "sha256-G+59ZwQwBbnO9ZjAB5zMEmWZbeG4k7ko/lPz+ceqYKs="; + goModHash = "sha256-2RUwj16RNaZ/gCOcd7b3LRCHiROCRj9HuzbBdLdgWGo="; appviewDeps = { inherit htmx-src htmx-ws-src lucide-src inter-fonts-src ibm-plex-mono-src goModHash gitignoreSource; }; -- 2.51.2