diff --git a/appview/pages/pages.go b/appview/pages/pages.go
--- a/appview/pages/pages.go
+++ b/appview/pages/pages.go
@@ -974,6 +974,26 @@
return p.executeRepo("repo/pipelines/pipelines", w, params)
}
+type LogBlockParams struct {
+ Id int
+ Name string
+ Command string
+ Collapsed bool
+}
+
+func (p *Pages) LogBlock(w io.Writer, params LogBlockParams) error {
+ return p.executePlain("repo/pipelines/fragments/logBlock", w, params)
+}
+
+type LogLineParams struct {
+ Id int
+ Content string
+}
+
+func (p *Pages) LogLine(w io.Writer, params LogLineParams) error {
+ return p.executePlain("repo/pipelines/fragments/logLine", w, params)
+}
+
type WorkflowParams struct {
LoggedInUser *oauth.User
RepoInfo repoinfo.RepoInfo
diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go
--- a/appview/pipelines/pipelines.go
+++ b/appview/pipelines/pipelines.go
@@ -1,10 +1,9 @@
package pipelines
import (
+ "bytes"
"context"
"encoding/json"
- "fmt"
- "html"
"log/slog"
"net/http"
"strings"
@@ -170,15 +169,6 @@
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)
@@ -238,86 +228,110 @@
defer spindleConn.Close()
// create a channel for incoming messages
- msgChan := make(chan []byte, 10)
- errChan := make(chan error, 1)
-
+ evChan := make(chan logEvent, 100)
// start a goroutine to read from spindle
- go func() {
- defer close(msgChan)
- defer close(errChan)
-
- for {
- _, msg, err := spindleConn.ReadMessage()
- if err != nil {
- if websocket.IsCloseError(err,
- websocket.CloseNormalClosure,
- websocket.CloseGoingAway,
- websocket.CloseAbnormalClosure) {
- errChan <- nil // signal graceful end
- } else {
- errChan <- err
- }
- return
- }
- msgChan <- msg
- }
- }()
+ go readLogs(spindleConn, evChan)
stepIdx := 0
+ var fragment bytes.Buffer
for {
select {
case <-ctx.Done():
l.Info("client disconnected")
return
- case err := <-errChan:
- if err != nil {
+
+ 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)
+ return
}
- if err == nil {
- l.Info("log tail complete")
- }
-
- return
- case msg := <-msgChan:
var logLine spindlemodel.LogLine
- if err = json.Unmarshal(msg, &logLine); err != nil {
+ if err = json.Unmarshal(ev.msg, &logLine); err != nil {
l.Error("failed to parse logline", "err", err)
continue
}
- var fragment []byte
+ fragment.Reset()
+
switch logLine.Kind {
case spindlemodel.LogKindControl:
// control messages create a new step block
stepIdx++
- fragment = fmt.Appendf(nil, `
-
- `, stepIdx, logLine.Content, stepIdx)
+ collapsed := false
+ if logLine.StepKind == spindlemodel.StepKindSystem {
+ collapsed = true
+ }
+ err = p.pages.LogBlock(&fragment, pages.LogBlockParams{
+ Id: stepIdx,
+ Name: logLine.Content,
+ Command: logLine.StepCommand,
+ Collapsed: collapsed,
+ })
case spindlemodel.LogKindData:
// data messages simply insert new log lines into current step
- escaped := html.EscapeString(logLine.Content)
- fragment = fmt.Appendf(nil, `
-
- `, stepIdx, escaped)
+ err = p.pages.LogLine(&fragment, pages.LogLineParams{
+ Id: stepIdx,
+ Content: logLine.Content,
+ })
+ }
+ if err != nil {
+ l.Error("failed to render log line", "err", err)
+ return
}
- if err = clientConn.WriteMessage(websocket.TextMessage, fragment); err != nil {
+ if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); 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)
+ return
}
}
+ }
+}
+
+// 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}
}
}
diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go
--- a/spindle/engine/engine.go
+++ b/spindle/engine/engine.go
@@ -315,8 +315,8 @@
}
defer wfLogger.Close()
- ctl := wfLogger.ControlWriter()
- ctl.Write([]byte(step.Command))
+ ctl := wfLogger.ControlWriter(stepIdx, step)
+ ctl.Write([]byte(step.Name))
logs, err := e.docker.ContainerLogs(ctx, containerID, container.LogsOptions{
Follow: true,
diff --git a/spindle/engine/logger.go b/spindle/engine/logger.go
--- a/spindle/engine/logger.go
+++ b/spindle/engine/logger.go
@@ -41,29 +41,44 @@
func (l *WorkflowLogger) DataWriter(stream string) io.Writer {
// TODO: emit stream
- return &jsonWriter{logger: l, kind: models.LogKindData}
-}
-
-func (l *WorkflowLogger) ControlWriter() io.Writer {
- return &jsonWriter{logger: l, kind: models.LogKindControl}
-}
-
-type jsonWriter struct {
- logger *WorkflowLogger
- kind models.LogKind
-}
-
-func (w *jsonWriter) Write(p []byte) (int, error) {
- line := strings.TrimRight(string(p), "\r\n")
-
- entry := models.LogLine{
- Kind: w.kind,
- Content: line,
+ return &dataWriter{
+ logger: l,
+ stream: stream,
}
+}
+func (l *WorkflowLogger) ControlWriter(idx int, step models.Step) io.Writer {
+ return &controlWriter{
+ logger: l,
+ idx: idx,
+ step: step,
+ }
+}
+
+type dataWriter struct {
+ logger *WorkflowLogger
+ stream string
+}
+
+func (w *dataWriter) Write(p []byte) (int, error) {
+ line := strings.TrimRight(string(p), "\r\n")
+ entry := models.NewDataLogLine(line, w.stream)
if err := w.logger.encoder.Encode(entry); err != nil {
return 0, err
}
-
return len(p), nil
+}
+
+type controlWriter struct {
+ logger *WorkflowLogger
+ idx int
+ step models.Step
+}
+
+func (w *controlWriter) Write(_ []byte) (int, error) {
+ entry := models.NewControlLogLine(w.idx, w.step)
+ if err := w.logger.encoder.Encode(entry); err != nil {
+ return 0, err
+ }
+ return len(w.step.Name), nil
}
diff --git a/spindle/models/models.go b/spindle/models/models.go
--- a/spindle/models/models.go
+++ b/spindle/models/models.go
@@ -88,4 +88,25 @@
Stream string `json:"stream,omitempty"`
// fields if kind is "control"
+ StepId int `json:"step_id,omitempty"`
+ StepKind StepKind `json:"step_kind,omitempty"`
+ StepCommand string `json:"step_command,omitempty"`
+}
+
+func NewDataLogLine(content, stream string) LogLine {
+ return LogLine{
+ Kind: LogKindData,
+ Content: content,
+ Stream: stream,
+ }
+}
+
+func NewControlLogLine(idx int, step Step) LogLine {
+ return LogLine{
+ Kind: LogKindControl,
+ Content: step.Name,
+ StepId: idx,
+ StepKind: step.Kind,
+ StepCommand: step.Command,
+ }
}
diff --git a/spindle/models/pipeline.go b/spindle/models/pipeline.go
--- a/spindle/models/pipeline.go
+++ b/spindle/models/pipeline.go
@@ -15,7 +15,17 @@
Command string
Name string
Environment map[string]string
+ Kind StepKind
}
+
+type StepKind int
+
+const (
+ // steps injected by the CI runner
+ StepKindSystem StepKind = iota
+ // steps defined by the user in the original pipeline
+ StepKindUser
+)
type Workflow struct {
Steps []Step
@@ -46,6 +56,7 @@
sstep.Environment = stepEnvToMap(tstep.Environment)
sstep.Command = tstep.Command
sstep.Name = tstep.Name
+ sstep.Kind = StepKindUser
swf.Steps = append(swf.Steps, sstep)
}
swf.Name = twf.Name
@@ -59,7 +70,10 @@
setup.addStep(nixConfStep())
setup.addStep(cloneStep(*twf, *pl.TriggerMetadata.Repo, cfg.Server.Dev))
setup.addStep(checkoutStep(*twf, *pl.TriggerMetadata))
- setup.addStep(dependencyStep(*twf))
+ // this step could be empty
+ if s := dependencyStep(*twf); s != nil {
+ setup.addStep(*s)
+ }
// append setup steps in order to the start of workflow steps
swf.Steps = append(*setup, swf.Steps...)
diff --git a/spindle/models/setup_steps.go b/spindle/models/setup_steps.go
--- a/spindle/models/setup_steps.go
+++ b/spindle/models/setup_steps.go
@@ -83,7 +83,7 @@
// For dependencies using a custom registry (i.e. not nixpkgs), it collects
// all packages and adds a single 'nix profile install' step to the
// beginning of the workflow's step list.
-func dependencyStep(twf tangled.Pipeline_Workflow) Step {
+func dependencyStep(twf tangled.Pipeline_Workflow) *Step {
var customPackages []string
for _, d := range twf.Dependencies {
@@ -111,7 +111,7 @@
"NIX_SHOW_DOWNLOAD_PROGRESS": "0",
},
}
- return installStep
+ return &installStep
}
- return Step{}
+ return nil
}
diff --git a/appview/pages/templates/repo/pipelines/pipelines.html b/appview/pages/templates/repo/pipelines/pipelines.html
--- a/appview/pages/templates/repo/pipelines/pipelines.html
+++ b/appview/pages/templates/repo/pipelines/pipelines.html
@@ -8,17 +8,7 @@
{{ define "repoContent" }}
-
-{{ end }}
-
-{{ define "repoAfter" }}
-
+
{{ range .Pipelines }}
{{ block "pipeline" (list $ .) }} {{ end }}
{{ else }}
@@ -26,13 +16,15 @@
No pipelines run for this repository.
{{ end }}
-
+
+
{{ end }}
+
{{ define "pipeline" }}
{{ $root := index . 0 }}
{{ $p := index . 1 }}
-
+
{{ block "pipelineHeader" $ }} {{ end }}
{{ end }}
diff --git a/appview/pages/templates/repo/pipelines/workflow.html b/appview/pages/templates/repo/pipelines/workflow.html
--- a/appview/pages/templates/repo/pipelines/workflow.html
+++ b/appview/pages/templates/repo/pipelines/workflow.html
@@ -24,7 +24,7 @@
{{ $active := .Workflow }}
{{ with .Pipeline }}
{{ $id := .Id }}
-