From 1493f351c6402978cf3f6bfdcaa603d3e42f2324 Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Thu, 30 Jul 2026 20:58:52 -0400 Subject: [PATCH] pipelines: add trigger --- internal/app/dependencies.go | 2 + internal/app/pipelines.go | 62 +++++++++++++++++++++++++++ internal/app/pipelines_test.go | 42 +++++++++++++++--- internal/app/service_test.go | 4 ++ internal/app/types.go | 7 +++ internal/cli/pipeline_trigger.go | 41 ++++++++++++++++++ internal/cli/pipeline_trigger_test.go | 19 ++++++++ internal/cli/root.go | 2 +- internal/gitutil/branch.go | 9 ++++ spindle/client.go | 28 ++++++++++++ spindle/client_test.go | 36 ++++++++++++++++ 11 files changed, 246 insertions(+), 6 deletions(-) create mode 100644 internal/cli/pipeline_trigger.go create mode 100644 internal/cli/pipeline_trigger_test.go diff --git a/internal/app/dependencies.go b/internal/app/dependencies.go index 2e2cbd5..95645aa 100644 --- a/internal/app/dependencies.go +++ b/internal/app/dependencies.go @@ -51,6 +51,7 @@ type gitClient interface { CheckoutPatch(context.Context, gitutil.CheckoutPatchParams) error GeneratePatch(context.Context, string, string, string) ([]byte, error) CurrentBranch(context.Context, string) (string, error) + ResolveCommit(context.Context, string, string) (string, error) DefaultBranch(context.Context, string) (string, error) DetectRepoCandidatesFromCWD(context.Context) ([]gitutil.RepoContext, error) } @@ -71,6 +72,7 @@ type pipelineClient interface { QueryLatestPipeline(context.Context, string) (*spindle.QueryPipelinesOutput, error) GetPipeline(context.Context, string) (*spindle.Pipeline, error) CancelPipeline(context.Context, spindle.CancelPipelineInput) error + TriggerPipeline(context.Context, spindle.TriggerPipelineInput) (*spindle.TriggerPipelineOutput, error) } type spindleClientFactory interface { diff --git a/internal/app/pipelines.go b/internal/app/pipelines.go index be0cd40..43bc282 100644 --- a/internal/app/pipelines.go +++ b/internal/app/pipelines.go @@ -3,6 +3,7 @@ package app import ( "context" "fmt" + "strings" "github.com/alyraffauf/tg/spindle" ) @@ -18,6 +19,67 @@ func (s *Service) ListPipelines(ctx context.Context, target Target) ([]Pipeline, return listPipelinePages(ctx, client, repoDID) } +// TriggerPipeline starts a manual pipeline for revision. A full commit SHA can be +// used without a local Git checkout; other revisions are resolved locally. +func (s *Service) TriggerPipeline(ctx context.Context, target Target, revision string, workflows []string) (*PipelineTriggerResult, error) { + commit, ref, err := s.resolvePipelineRevision(ctx, revision) + if err != nil { + return nil, err + } + spindleHost, repoDID, err := s.pipelineTarget(ctx, target) + if err != nil { + return nil, err + } + pds, _, err := s.authenticatedPDS(ctx) + if err != nil { + return nil, err + } + audience, err := spindle.ServiceDID(spindleHost) + if err != nil { + return nil, err + } + token, err := pds.GetServiceAuth(ctx, audience, "sh.tangled.ci.triggerPipeline") + if err != nil { + return nil, fmt.Errorf("mint pipeline trigger token: %w", err) + } + client, err := s.spindle.NewWithToken(spindleHost, token) + if err != nil { + return nil, fmt.Errorf("connect to pipeline spindle: %w", err) + } + response, err := client.TriggerPipeline(ctx, spindle.TriggerPipelineInput{ + Repo: repoDID, Workflows: workflows, + Trigger: spindle.ManualTrigger{LexiconTypeID: "sh.tangled.ci.trigger#manual", SHA: commit, Ref: ref}, + }) + if err != nil { + return nil, err + } + return &PipelineTriggerResult{Pipeline: extractRKey(response.Pipeline), Commit: commit, Workflows: workflows}, nil +} + +func (s *Service) resolvePipelineRevision(ctx context.Context, revision string) (commit, ref string, err error) { + if isFullCommitSHA(revision) { + return revision, "", nil + } + commit, err = s.git.ResolveCommit(ctx, "", revision) + if err != nil { + return "", "", err + } + if revision == "HEAD" { + if branch, branchErr := s.git.CurrentBranch(ctx, ""); branchErr == nil { + return commit, "refs/heads/" + branch, nil + } + return commit, "", nil + } + return commit, revision, nil +} + +func isFullCommitSHA(revision string) bool { + if len(revision) != 40 { + return false + } + return strings.Trim(revision, "0123456789abcdefABCDEF") == "" +} + // CancelPipeline cancels every workflow in a pipeline, or only the selected workflows. func (s *Service) CancelPipeline(ctx context.Context, target Target, pipelineID string, workflows []string) (*PipelineCancelResult, error) { spindleHost, repoDID, err := s.pipelineTarget(ctx, target) diff --git a/internal/app/pipelines_test.go b/internal/app/pipelines_test.go index 8a09d49..6f1a20b 100644 --- a/internal/app/pipelines_test.go +++ b/internal/app/pipelines_test.go @@ -130,12 +130,39 @@ func TestCancelPipelineMintsSpindleToken(t *testing.T) { } } +func TestTriggerPipelineUsesFullSHAWithoutGitResolution(t *testing.T) { + commit := "0123456789abcdef0123456789abcdef01234567" + client := &testPipelineClient{triggerOutput: &spindle.TriggerPipelineOutput{Pipeline: "at://did:plc:spindle/sh.tangled.ci.pipeline/3mrvk5dbnep22"}} + pds := &testPDS{} + service := testService(pds, &testGit{}, &testKnot{}) + service.appview = testAppview{repo: &tangled.Repo{Value: tangledlex.Repo{ + Knot: "knot.example", Spindle: optionalString("spindle.example"), RepoDid: optionalString("did:plc:repo"), + }}} + service.spindle = testSpindleFactory{client: client} + + result, err := service.TriggerPipeline(context.Background(), Target{Handle: "owner.test", Repo: "example"}, commit, []string{"test.yml"}) + if err != nil { + t.Fatalf("TriggerPipeline() error = %v", err) + } + if result.Pipeline != "3mrvk5dbnep22" || result.Commit != commit { + t.Fatalf("TriggerPipeline() = %+v", result) + } + if pds.serviceAuthLexiconMethods[0] != "sh.tangled.ci.triggerPipeline" { + t.Fatalf("service auth method = %q", pds.serviceAuthLexiconMethods[0]) + } + if client.triggerInput.Trigger.SHA != commit || client.triggerInput.Trigger.Ref != "" || client.triggerInput.Repo != "did:plc:repo" { + t.Fatalf("trigger input = %+v", client.triggerInput) + } +} + type testPipelineClient struct { - responses []*spindle.QueryPipelinesOutput - cursors []string - err error - cancelInput spindle.CancelPipelineInput - pipeline *spindle.Pipeline + responses []*spindle.QueryPipelinesOutput + cursors []string + err error + cancelInput spindle.CancelPipelineInput + pipeline *spindle.Pipeline + triggerInput spindle.TriggerPipelineInput + triggerOutput *spindle.TriggerPipelineOutput } func (c *testPipelineClient) QueryLatestPipeline(_ context.Context, _ string) (*spindle.QueryPipelinesOutput, error) { @@ -151,6 +178,11 @@ func (c *testPipelineClient) CancelPipeline(_ context.Context, input spindle.Can return c.err } +func (c *testPipelineClient) TriggerPipeline(_ context.Context, input spindle.TriggerPipelineInput) (*spindle.TriggerPipelineOutput, error) { + c.triggerInput = input + return c.triggerOutput, c.err +} + type testSpindleFactory struct { client pipelineClient } diff --git a/internal/app/service_test.go b/internal/app/service_test.go index d606aab..121657b 100644 --- a/internal/app/service_test.go +++ b/internal/app/service_test.go @@ -350,6 +350,7 @@ func (p *testPDS) GetServiceAuth(_ context.Context, audience, lexiconMethod stri type testGit struct { branch string + commit string patch []byte patchErr error clones []gitutil.CloneRepoParams @@ -378,6 +379,9 @@ func (g *testGit) GeneratePatch(context.Context, string, string, string) ([]byte return nil, errors.New("not implemented") } func (g *testGit) CurrentBranch(context.Context, string) (string, error) { return g.branch, nil } +func (g *testGit) ResolveCommit(context.Context, string, string) (string, error) { + return g.commit, nil +} func (g *testGit) DefaultBranch(context.Context, string) (string, error) { return "", errors.New("not implemented") } diff --git a/internal/app/types.go b/internal/app/types.go index 54e90a4..662f595 100644 --- a/internal/app/types.go +++ b/internal/app/types.go @@ -67,6 +67,13 @@ type PipelineCancelResult struct { CancellationRequested bool `json:"cancellationRequested"` } +// PipelineTriggerResult describes a manually triggered pipeline. +type PipelineTriggerResult struct { + Pipeline string `json:"pipeline"` + Commit string `json:"commit"` + Workflows []string `json:"workflows,omitempty"` +} + // SSHKeyItem is one SSH public key in a listing. type SSHKeyItem struct { Name string `json:"name"` diff --git a/internal/cli/pipeline_trigger.go b/internal/cli/pipeline_trigger.go new file mode 100644 index 0000000..bef4ad5 --- /dev/null +++ b/internal/cli/pipeline_trigger.go @@ -0,0 +1,41 @@ +package cli + +import ( + "fmt" + + "github.com/alyraffauf/tg/internal/app" + "github.com/spf13/cobra" +) + +func newPipelineTriggerCommand(service *app.Service) *cobra.Command { + var repository string + var workflows []string + + command := &cobra.Command{ + Use: "trigger ", + Short: "Trigger pipelines for a commit or branch", + Long: `Trigger pipelines for a commit or branch. + +A full commit SHA works without a local checkout. Other Git revisions, such +as HEAD or a branch name, are resolved from the current checkout. If --repo +is not set, the repository is detected from the current directory's git +origin remote.`, + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + target, err := resolveTargetFlag(cmd.Context(), repository, service) + if err != nil { + return err + } + result, err := service.TriggerPipeline(cmd.Context(), target, args[0], workflows) + if err != nil { + return err + } + return output(cmd, result, func(result *app.PipelineTriggerResult) { + fmt.Fprintf(cmd.OutOrStdout(), "Triggered pipeline %s.\n", result.Pipeline) + }) + }, + } + command.Flags().StringVarP(&repository, "repo", "R", "", "Target repository as handle/repo") + command.Flags().StringSliceVarP(&workflows, "workflow", "w", nil, "Workflow name to trigger (repeatable)") + return command +} diff --git a/internal/cli/pipeline_trigger_test.go b/internal/cli/pipeline_trigger_test.go new file mode 100644 index 0000000..bacc877 --- /dev/null +++ b/internal/cli/pipeline_trigger_test.go @@ -0,0 +1,19 @@ +package cli + +import ( + "testing" + + "github.com/alyraffauf/tg/internal/app" +) + +func TestPipelineTriggerCommandFlags(t *testing.T) { + command := newPipelineTriggerCommand(&app.Service{}) + if command.Use != "trigger " { + t.Fatalf("command use = %q", command.Use) + } + for _, name := range []string{"repo", "workflow"} { + if command.Flags().Lookup(name) == nil { + t.Errorf("pipeline trigger has no --%s flag", name) + } + } +} diff --git a/internal/cli/root.go b/internal/cli/root.go index 4e40864..bedeea6 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)) + pipeline.AddCommand(newPipelineListCommand(service), newPipelineViewCommand(service), newPipelineStatusCommand(service), newPipelineCancelCommand(service), newPipelineTriggerCommand(service)) rootCmd.AddCommand(pipeline) keys := newSSHKeyCommand(service) diff --git a/internal/gitutil/branch.go b/internal/gitutil/branch.go index df4ede7..c69787b 100644 --- a/internal/gitutil/branch.go +++ b/internal/gitutil/branch.go @@ -32,3 +32,12 @@ func (c *Client) CurrentBranch(ctx context.Context, dir string) (string, error) func CurrentBranch(ctx context.Context, dir string) (string, error) { return defaultClient.CurrentBranch(ctx, dir) } + +// ResolveCommit resolves revision to its full commit SHA. +func (c *Client) ResolveCommit(ctx context.Context, dir, revision string) (string, error) { + commit, err := c.gitOutput(ctx, dir, "rev-parse", "--verify", revision+"^{commit}") + if err != nil { + return "", fmt.Errorf("resolve commit %q in %q: %w", revision, dir, err) + } + return strings.TrimSpace(string(commit)), nil +} diff --git a/spindle/client.go b/spindle/client.go index a25c538..5c74ce9 100644 --- a/spindle/client.go +++ b/spindle/client.go @@ -102,6 +102,34 @@ type CancelPipelineInput struct { Workflows []string `json:"workflows,omitempty"` } +// TriggerPipelineInput is the argument to sh.tangled.ci.triggerPipeline. +type TriggerPipelineInput struct { + Repo string `json:"repo"` + Trigger ManualTrigger `json:"trigger"` + Workflows []string `json:"workflows,omitempty"` +} + +// ManualTrigger describes a manually requested pipeline trigger. +type ManualTrigger struct { + LexiconTypeID string `json:"$type"` + SHA string `json:"sha"` + Ref string `json:"ref,omitempty"` +} + +// TriggerPipelineOutput is returned after creating a manual pipeline. +type TriggerPipelineOutput struct { + Pipeline string `json:"pipeline"` +} + +// TriggerPipeline starts a manual pipeline for a commit. +func (c *Client) TriggerPipeline(ctx context.Context, input TriggerPipelineInput) (*TriggerPipelineOutput, error) { + var output TriggerPipelineOutput + if err := c.Post(ctx, syntax.NSID("sh.tangled.ci.triggerPipeline"), input, &output); err != nil { + return nil, fmt.Errorf("trigger pipeline: %w", err) + } + return &output, nil +} + // GetPipeline fetches one pipeline by its spindle-local ID. func (c *Client) GetPipeline(ctx context.Context, pipelineID string) (*Pipeline, error) { var pipeline Pipeline diff --git a/spindle/client_test.go b/spindle/client_test.go index 32e5fba..a2299a0 100644 --- a/spindle/client_test.go +++ b/spindle/client_test.go @@ -64,6 +64,42 @@ func TestCancelPipelineAuthenticatesAndPostsInput(t *testing.T) { } } +func TestTriggerPipelineAuthenticatesAndReturnsPipeline(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) { + if request.URL.Path != "/xrpc/sh.tangled.ci.triggerPipeline" || request.Method != http.MethodPost { + t.Fatalf("request = %s %s", request.Method, request.URL.Path) + } + if request.Header.Get("Authorization") != "Bearer token" { + t.Fatalf("authorization = %q", request.Header.Get("Authorization")) + } + var input TriggerPipelineInput + if err := json.NewDecoder(request.Body).Decode(&input); err != nil { + t.Fatalf("decode input: %v", err) + } + if input.Trigger.SHA != "0123456789abcdef0123456789abcdef01234567" || input.Trigger.Ref != "refs/heads/feature" { + t.Fatalf("input = %+v", input) + } + writer.Header().Set("Content-Type", "application/json") + _, _ = writer.Write([]byte(`{"pipeline":"at://did:plc:spindle/sh.tangled.ci.pipeline/3mrvk5dbnep22"}`)) + })) + defer server.Close() + + client, err := NewWithToken(server.URL, "token", server.Client()) + if err != nil { + t.Fatalf("NewWithToken() error = %v", err) + } + output, err := client.TriggerPipeline(context.Background(), TriggerPipelineInput{ + Repo: "did:plc:repo", + Trigger: ManualTrigger{LexiconTypeID: "sh.tangled.ci.trigger#manual", SHA: "0123456789abcdef0123456789abcdef01234567", Ref: "refs/heads/feature"}, + }) + if err != nil { + t.Fatalf("TriggerPipeline() error = %v", err) + } + if output.Pipeline != "at://did:plc:spindle/sh.tangled.ci.pipeline/3mrvk5dbnep22" { + t.Fatalf("output = %+v", output) + } +} + func TestQueryLatestPipelineLimitsResults(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) { if request.URL.Query().Get("repo") != "did:plc:repo" { -- 2.51.2