diff --git a/go.mod b/go.mod index fee6462..e2bc059 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require ( charm.land/lipgloss/v2 v2.0.5 github.com/bluesky-social/indigo v0.0.0-20260629160527-dfe5578fd537 github.com/charmbracelet/x/term v0.2.2 + github.com/gorilla/websocket v1.5.1 github.com/ipfs/go-cid v0.6.2 github.com/spf13/cobra v1.10.2 github.com/spf13/viper v1.21.0 diff --git a/go.sum b/go.sum index 4172b7c..1e69a07 100644 --- a/go.sum +++ b/go.sum @@ -66,6 +66,8 @@ github.com/google/go-querystring v1.2.0 h1:yhqkPbu2/OH+V9BfpCVPZkNmUXhb2gBxJArfh github.com/google/go-querystring v1.2.0/go.mod h1:8IFJqpSRITyJ8QhQ13bmbeMBDfmeEJZD5A0egEOmkqU= github.com/gorilla/css v1.0.1 h1:ntNaBIghp6JmvWnxbZKANoLyuXTPZ4cAMlo6RyhlbO8= github.com/gorilla/css v1.0.1/go.mod h1:BvnYkspnSzMmwRK+b8/xgNPLiIuNZr6vbZBTPQ2A3b0= +github.com/gorilla/websocket v1.5.1 h1:gmztn0JnHVt9JZquRuzLw3g4wouNVzKL15iLr/zn/QY= +github.com/gorilla/websocket v1.5.1/go.mod h1:x3kM2JMyaluk02fnUJpQuwD2dCS5NDG2ZHL0uE0tcaY= github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/hexops/gotextdiff v1.0.3 h1:gitA9+qJrrTCsiCl7+kh75nPqQt1cx4ZkudSTLoUqJM= diff --git a/internal/app/dependencies.go b/internal/app/dependencies.go index 87c825c..13a7815 100644 --- a/internal/app/dependencies.go +++ b/internal/app/dependencies.go @@ -77,6 +77,7 @@ type pipelineClient interface { GetPipeline(context.Context, string) (*spindle.Pipeline, error) CancelPipeline(context.Context, spindle.CancelPipelineInput) error TriggerPipeline(context.Context, spindle.TriggerPipelineInput) (*spindle.TriggerPipelineOutput, error) + SubscribePipelineLogs(context.Context, string, []string, func(spindle.PipelineLogEvent) error) error } type spindleClientFactory interface { diff --git a/internal/app/domain_types.go b/internal/app/domain_types.go index b19a4d7..c63a1f5 100644 --- a/internal/app/domain_types.go +++ b/internal/app/domain_types.go @@ -229,3 +229,30 @@ type GitCredentialResult struct { Handle string MatchesRequestedHost bool } + +// PipelineLogControl marks the start or end of a workflow step. +type PipelineLogControl struct { + Kind string `json:"kind"` + Step int `json:"step"` + Time string `json:"time"` + Status string `json:"status,omitempty"` + Content string `json:"content"` + Workflow string `json:"workflow"` + Command *string `json:"command,omitempty"` +} + +// PipelineLogData is one line of workflow output. +type PipelineLogData struct { + Step int `json:"step"` + Time string `json:"time"` + Stream string `json:"stream"` + Content string `json:"content"` + Workflow string `json:"workflow"` +} + +// PipelineLogEvent is one event from a log subscription. +type PipelineLogEvent struct { + Type string `json:"type"` + Control *PipelineLogControl `json:"control,omitempty"` + Data *PipelineLogData `json:"data,omitempty"` +} diff --git a/internal/app/pipeline_logs.go b/internal/app/pipeline_logs.go new file mode 100644 index 0000000..d992f8d --- /dev/null +++ b/internal/app/pipeline_logs.go @@ -0,0 +1,48 @@ +package app + +import ( + "context" + "fmt" + + "github.com/alyraffauf/tg/spindle" +) + +// PipelineLogs streams log events from a pipeline. A non-empty workflows list +// filters to the named workflows. +func (s *Service) PipelineLogs(ctx context.Context, target Target, pipelineID string, workflows []string, onEvent func(PipelineLogEvent) error) error { + spindleHost, _, err := s.pipelineTarget(ctx, target) + if err != nil { + return err + } + client, err := s.spindle.New(spindleHost) + if err != nil { + return fmt.Errorf("connect to pipeline spindle: %w", err) + } + return client.SubscribePipelineLogs(ctx, pipelineID, workflows, func(event spindle.PipelineLogEvent) error { + return onEvent(PipelineLogEvent{ + Type: event.Type, + Control: convertControl(event.Control), + Data: convertData(event.Data), + }) + }) +} + +func convertControl(c *spindle.PipelineLogControl) *PipelineLogControl { + if c == nil { + return nil + } + return &PipelineLogControl{ + Kind: c.Kind, Step: c.Step, Time: c.Time, Status: c.Status, + Content: c.Content, Workflow: c.Workflow, Command: c.Command, + } +} + +func convertData(d *spindle.PipelineLogData) *PipelineLogData { + if d == nil { + return nil + } + return &PipelineLogData{ + Step: d.Step, Time: d.Time, Stream: d.Stream, + Content: d.Content, Workflow: d.Workflow, + } +} diff --git a/internal/app/pipelines_test.go b/internal/app/pipelines_test.go index 46a8377..7437ba7 100644 --- a/internal/app/pipelines_test.go +++ b/internal/app/pipelines_test.go @@ -188,6 +188,7 @@ type testPipelineClient struct { pipelineID string triggerInput spindle.TriggerPipelineInput triggerOutput *spindle.TriggerPipelineOutput + logEvents []spindle.PipelineLogEvent } func (c *testPipelineClient) QueryLatestPipeline(_ context.Context, _ string) (*spindle.QueryPipelinesOutput, error) { @@ -230,3 +231,12 @@ func (c *testPipelineClient) QueryPipelines(_ context.Context, _ string, cursor c.responses = c.responses[1:] return response, nil } + +func (c *testPipelineClient) SubscribePipelineLogs(_ context.Context, _ string, _ []string, onEvent func(spindle.PipelineLogEvent) error) error { + for _, event := range c.logEvents { + if err := onEvent(event); err != nil { + return err + } + } + return nil +} diff --git a/internal/cli/pipeline_logs.go b/internal/cli/pipeline_logs.go new file mode 100644 index 0000000..109a19a --- /dev/null +++ b/internal/cli/pipeline_logs.go @@ -0,0 +1,111 @@ +package cli + +import ( + "encoding/json" + "fmt" + "image/color" + "io" + "strings" + + "charm.land/lipgloss/v2" + "github.com/alyraffauf/tg/internal/app" + "github.com/spf13/cobra" +) + +var workflowColors = []color.Color{ + lipgloss.Cyan, lipgloss.Magenta, lipgloss.Green, + lipgloss.Yellow, lipgloss.Blue, lipgloss.Red, +} + +func newPipelineLogsCommand(service *app.Service) *cobra.Command { + var repository string + var workflows []string + + command := &cobra.Command{ + Use: "logs ", + Short: "Stream pipeline logs", + Long: `Stream log output from a pipeline in real time. + +If --repo is not set, the repository is detected from the current directory's +git origin remote. Use --workflow to filter to specific workflows.`, + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + target, err := resolveTargetFlag(cmd.Context(), repository, service) + if err != nil { + return err + } + jsonOutput, _ := cmd.Flags().GetBool("json") + out := cmd.OutOrStdout() + terminal := isTerminal(out) + return service.PipelineLogs(cmd.Context(), target, args[0], workflows, func(event app.PipelineLogEvent) error { + if jsonOutput { + return json.NewEncoder(out).Encode(event) + } + return renderPipelineLogEvent(out, cmd.ErrOrStderr(), event, terminal) + }) + }, + } + command.Flags().StringVarP(&repository, "repo", "R", "", "Target repository as handle/repo") + command.Flags().StringSliceVarP(&workflows, "workflow", "w", nil, "Workflow name to stream (repeatable)") + return command +} + +func renderPipelineLogEvent(stdout, stderr io.Writer, event app.PipelineLogEvent, terminal bool) error { + switch { + case event.Control != nil: + if event.Control.Status == "start" { + renderStepHeader(stdout, event.Control, terminal) + } + case event.Data != nil: + w := stderr + if event.Data.Stream == "stdout" { + w = stdout + } + return writeLogLines(w, event.Data, terminal) + } + return nil +} + +func renderStepHeader(w io.Writer, control *app.PipelineLogControl, terminal bool) { + sep := "──" + tag := control.Workflow + if terminal { + sep = lipgloss.NewStyle().Faint(true).Render(sep) + tag = lipgloss.NewStyle().Foreground(workflowColor(control.Workflow)).Render(tag) + } + fmt.Fprintf(w, "\n%s %s: %s ", sep, tag, control.Content) + if control.Command != nil { + command := "$ " + *control.Command + if terminal { + command = lipgloss.NewStyle().Faint(true).Render(command) + } + fmt.Fprintf(w, "(%s) ", command) + } + fmt.Fprintf(w, "%s\n", sep) +} + +func writeLogLines(w io.Writer, data *app.PipelineLogData, terminal bool) error { + content := strings.TrimRight(data.Content, "\n") + if content == "" { + return nil + } + tag := data.Workflow + if terminal { + tag = lipgloss.NewStyle().Foreground(workflowColor(data.Workflow)).Render(tag) + } + for _, line := range strings.Split(content, "\n") { + if terminal && data.Stream == "stderr" { + line = lipgloss.NewStyle().Foreground(lipgloss.Yellow).Render(line) + } + fmt.Fprintf(w, "[%s] %s\n", tag, line) + } + return nil +} + +func workflowColor(workflow string) color.Color { + var hash int + for _, r := range workflow { + hash += int(r) + } + return workflowColors[hash%len(workflowColors)] +} diff --git a/internal/cli/root.go b/internal/cli/root.go index 8cdc845..a246131 100644 --- a/internal/cli/root.go +++ b/internal/cli/root.go @@ -49,7 +49,7 @@ func newRoot(service *app.Service, defaultKnot, defaultSSHPort, defaultProtocol rootCmd.AddCommand(repo) pipeline := newPipelineCommand(service) - pipeline.AddCommand(newPipelineListCommand(service), newPipelineViewCommand(service), newPipelineStatusCommand(service), newPipelineCancelCommand(service), newPipelineTriggerCommand(service)) + pipeline.AddCommand(newPipelineListCommand(service), newPipelineViewCommand(service), newPipelineStatusCommand(service), newPipelineCancelCommand(service), newPipelineTriggerCommand(service), newPipelineLogsCommand(service)) rootCmd.AddCommand(pipeline) keys := newSSHKeyCommand(service) diff --git a/nix/tg.nix b/nix/tg.nix index a367eac..bd48ae9 100644 --- a/nix/tg.nix +++ b/nix/tg.nix @@ -10,7 +10,7 @@ buildGoModule { version = "dev"; src = ../.; proxyVendor = true; - vendorHash = "sha256-xN7gzhN7ljxIdOhgmAKsMGg3lORIGtcI6ScfNgupBWA="; + vendorHash = "sha256-NXjakn2F/FC+KiaoNs7kei1dELKyy5FqSAlo/MxbtaM="; subPackages = ["cmd/tg"]; nativeBuildInputs = [installShellFiles]; diff --git a/spindle/cbor.go b/spindle/cbor.go new file mode 100644 index 0000000..ad21081 --- /dev/null +++ b/spindle/cbor.go @@ -0,0 +1,107 @@ +package spindle + +import ( + "bytes" + "fmt" + "io" + + cborgen "github.com/whyrusleeping/cbor-gen" +) + +// decodeCBORMaps decodes consecutive CBOR maps from a WebSocket binary frame. +// Each spindle log frame is a header map followed by a body map. +func decodeCBORMaps(data []byte) ([]map[string]any, error) { + r := bytes.NewReader(data) + var maps []map[string]any + for r.Len() > 0 { + m, err := decodeCBORMap(r) + if err != nil { + return nil, err + } + maps = append(maps, m) + } + return maps, nil +} + +func decodeCBORMap(r io.Reader) (map[string]any, error) { + majorType, length, err := cborgen.CborReadHeader(r) + if err != nil { + return nil, fmt.Errorf("read map header: %w", err) + } + if majorType != 5 { + return nil, fmt.Errorf("expected CBOR map, got major type %d", majorType) + } + result := make(map[string]any, length) + for i := uint64(0); i < length; i++ { + key, err := decodeCBORValue(r) + if err != nil { + return nil, fmt.Errorf("read map key %d: %w", i, err) + } + strKey, ok := key.(string) + if !ok { + return nil, fmt.Errorf("expected string map key, got %T", key) + } + value, err := decodeCBORValue(r) + if err != nil { + return nil, fmt.Errorf("read map value for %q: %w", strKey, err) + } + result[strKey] = value + } + return result, nil +} + +func decodeCBORValue(r io.Reader) (any, error) { + majorType, value, err := cborgen.CborReadHeader(r) + if err != nil { + return nil, err + } + switch majorType { + case 0: + return value, nil + case 1: + return -int64(value) - 1, nil + case 3: + buf := make([]byte, value) + if _, err := io.ReadFull(r, buf); err != nil { + return nil, fmt.Errorf("read string: %w", err) + } + return string(buf), nil + case 7: + switch value { + case 20: + return false, nil + case 21: + return true, nil + case 22: + return nil, nil + default: + return nil, fmt.Errorf("unsupported simple value %d", value) + } + default: + return nil, fmt.Errorf("unsupported CBOR major type %d", majorType) + } +} + +func mapString(m map[string]any, key string) string { + if s, ok := m[key].(string); ok { + return s + } + return "" +} + +func mapInt(m map[string]any, key string) int { + switch n := m[key].(type) { + case int64: + return int(n) + case int: + return n + } + return 0 +} + +func mapStringOpt(m map[string]any, key string) *string { + if s, ok := m[key].(string); ok && s != "" { + return &s + } + return nil +} diff --git a/spindle/ci_subscribe_pipeline_logs.go b/spindle/ci_subscribe_pipeline_logs.go new file mode 100644 index 0000000..bab96ac --- /dev/null +++ b/spindle/ci_subscribe_pipeline_logs.go @@ -0,0 +1,103 @@ +package spindle + +import ( + "context" + "fmt" + "net/url" + "strings" + "time" + + "github.com/gorilla/websocket" +) + +// SubscribePipelineLogs streams log events from a pipeline over WebSocket. +// The spindle closes the connection when logs are fully delivered. +func (c *Client) SubscribePipelineLogs(ctx context.Context, pipelineID string, workflows []string, onEvent func(PipelineLogEvent) error) error { + u, err := url.Parse(c.Host) + if err != nil { + return fmt.Errorf("parse spindle host: %w", err) + } + switch strings.ToLower(u.Scheme) { + case "https": + u.Scheme = "wss" + case "http": + u.Scheme = "ws" + default: + return fmt.Errorf("spindle host must be http(s), got %q", u.Scheme) + } + u.Path = "/xrpc/sh.tangled.ci.subscribePipelineLogs" + query := url.Values{"pipeline": []string{pipelineID}} + if len(workflows) > 0 { + query["workflows"] = workflows + } + u.RawQuery = query.Encode() + + dialer := websocket.Dialer{HandshakeTimeout: 30 * time.Second} + conn, _, err := dialer.DialContext(ctx, u.String(), nil) + if err != nil { + return fmt.Errorf("connect pipeline log subscription: %w", err) + } + defer conn.Close() + + for { + if err := ctx.Err(); err != nil { + return err + } + _, data, err := conn.ReadMessage() + if err != nil { + if websocket.IsCloseError(err, websocket.CloseNormalClosure, websocket.CloseAbnormalClosure) { + return nil + } + return fmt.Errorf("read pipeline log: %w", err) + } + event, err := decodeLogEvent(data) + if err != nil { + return err + } + if event == nil { + continue + } + if err := onEvent(*event); err != nil { + return err + } + } +} + +func decodeLogEvent(data []byte) (*PipelineLogEvent, error) { + maps, err := decodeCBORMaps(data) + if err != nil { + return nil, fmt.Errorf("decode pipeline log frame: %w", err) + } + if len(maps) < 2 { + return nil, nil + } + header, body := maps[0], maps[1] + switch header["t"] { + case "#control": + return &PipelineLogEvent{ + Type: "control", + Control: &PipelineLogControl{ + Kind: mapString(body, "kind"), + Step: mapInt(body, "step"), + Time: mapString(body, "time"), + Status: mapString(body, "status"), + Content: mapString(body, "content"), + Workflow: mapString(body, "workflow"), + Command: mapStringOpt(body, "command"), + }, + }, nil + case "#data": + return &PipelineLogEvent{ + Type: "data", + Data: &PipelineLogData{ + Step: mapInt(body, "step"), + Time: mapString(body, "time"), + Stream: mapString(body, "stream"), + Content: mapString(body, "content"), + Workflow: mapString(body, "workflow"), + }, + }, nil + default: + return nil, nil + } +} diff --git a/spindle/ci_subscribe_pipeline_logs_test.go b/spindle/ci_subscribe_pipeline_logs_test.go new file mode 100644 index 0000000..66151fa --- /dev/null +++ b/spindle/ci_subscribe_pipeline_logs_test.go @@ -0,0 +1,105 @@ +package spindle + +import ( + "bytes" + "testing" +) + +func cborString(s string) []byte { + var buf bytes.Buffer + l := len(s) + if l < 24 { + buf.WriteByte(0x60 | byte(l)) + } else { + buf.WriteByte(0x78) + buf.WriteByte(byte(l)) + } + buf.WriteString(s) + return buf.Bytes() +} + +func cborInt(n int) []byte { + if n >= 0 { + return []byte{byte(n)} + } + return []byte{0x20 | byte(-n-1)} +} + +func cborMap(pairs ...[2][]byte) []byte { + var buf bytes.Buffer + buf.WriteByte(0xa0 | byte(len(pairs))) + for _, p := range pairs { + buf.Write(p[0]) + buf.Write(p[1]) + } + return buf.Bytes() +} + +func TestDecodeLogEventControl(t *testing.T) { + header := cborMap( + [2][]byte{cborString("t"), cborString("#control")}, + [2][]byte{cborString("op"), cborInt(1)}, + ) + body := cborMap( + [2][]byte{cborString("kind"), cborString("system")}, + [2][]byte{cborString("step"), cborInt(-1)}, + [2][]byte{cborString("time"), cborString("2026-07-31T06:34:08+03:00")}, + [2][]byte{cborString("status"), cborString("start")}, + [2][]byte{cborString("content"), cborString("Pull image")}, + [2][]byte{cborString("workflow"), cborString("build.yml")}, + ) + event, err := decodeLogEvent(append(header, body...)) + if err != nil { + t.Fatalf("decodeLogEvent() error = %v", err) + } + if event.Type != "control" || event.Control == nil { + t.Fatalf("decodeLogEvent() = %+v, want control event", event) + } + if event.Control.Workflow != "build.yml" || event.Control.Step != -1 { + t.Fatalf("control = %+v", event.Control) + } + if event.Data != nil { + t.Fatalf("unexpected data event") + } +} + +func TestDecodeLogEventData(t *testing.T) { + header := cborMap( + [2][]byte{cborString("t"), cborString("#data")}, + [2][]byte{cborString("op"), cborInt(1)}, + ) + body := cborMap( + [2][]byte{cborString("step"), cborInt(0)}, + [2][]byte{cborString("time"), cborString("2026-07-31T06:34:10+03:00")}, + [2][]byte{cborString("stream"), cborString("stderr")}, + [2][]byte{cborString("content"), cborString("hint: Using 'master'")}, + [2][]byte{cborString("workflow"), cborString("build.yml")}, + ) + event, err := decodeLogEvent(append(header, body...)) + if err != nil { + t.Fatalf("decodeLogEvent() error = %v", err) + } + if event.Type != "data" || event.Data == nil { + t.Fatalf("decodeLogEvent() = %+v, want data event", event) + } + if event.Data.Stream != "stderr" || event.Data.Workflow != "build.yml" { + t.Fatalf("data = %+v", event.Data) + } + if event.Control != nil { + t.Fatalf("unexpected control event") + } +} + +func TestDecodeLogEventUnknownType(t *testing.T) { + header := cborMap( + [2][]byte{cborString("t"), cborString("#unknown")}, + [2][]byte{cborString("op"), cborInt(1)}, + ) + event, err := decodeLogEvent(append(header, cborMap()...)) + if err != nil { + t.Fatalf("decodeLogEvent() error = %v", err) + } + if event != nil { + t.Fatalf("decodeLogEvent() = %+v, want nil for unknown type", event) + } +} diff --git a/spindle/ci_types.go b/spindle/ci_types.go index a7085fd..588d9af 100644 --- a/spindle/ci_types.go +++ b/spindle/ci_types.go @@ -48,3 +48,30 @@ type ManualTrigger struct { type TriggerPipelineOutput struct { Pipeline string `json:"pipeline"` } + +// PipelineLogControl marks the start or end of a workflow step. +type PipelineLogControl struct { + Kind string `json:"kind"` + Step int `json:"step"` + Time string `json:"time"` + Status string `json:"status,omitempty"` + Content string `json:"content"` + Workflow string `json:"workflow"` + Command *string `json:"command,omitempty"` +} + +// PipelineLogData is one line of workflow output. +type PipelineLogData struct { + Step int `json:"step"` + Time string `json:"time"` + Stream string `json:"stream"` + Content string `json:"content"` + Workflow string `json:"workflow"` +} + +// PipelineLogEvent is one event from a log subscription. +type PipelineLogEvent struct { + Type string `json:"type"` + Control *PipelineLogControl `json:"control,omitempty"` + Data *PipelineLogData `json:"data,omitempty"` +}