From 44a5763cdcaed07be3ea22a643885a0dc84d609e Mon Sep 17 00:00:00 2001 From: dawn Date: Wed, 26 Aug 2026 17:18:49 +0900 Subject: [PATCH] spindle/xrpc: implement org.tangled CI and secret endpoints Signed-off-by: dawn --- spindle/db/db.go | 14 +- spindle/engine/engine_test.go | 59 ++- ...i_pipeline_describe_workflow_definition.go | 4 +- spindle/xrpc/ci_pipeline_trigger_pipeline.go | 19 +- spindle/xrpc/org_tangled_ci_get_pipeline.go | 33 ++ .../org_tangled_ci_get_workflow_definition.go | 69 +++ spindle/xrpc/org_tangled_ci_pipeline.go | 136 +++++ .../xrpc/org_tangled_ci_query_pipelines.go | 43 ++ .../org_tangled_ci_subscribe_pipeline_logs.go | 41 ++ spindle/xrpc/org_tangled_secret_add_secret.go | 48 ++ spindle/xrpc/org_tangled_secret_common.go | 40 ++ .../xrpc/org_tangled_secret_list_secrets.go | 29 ++ .../xrpc/org_tangled_secret_remove_secret.go | 39 ++ spindle/xrpc/org_tangled_test.go | 491 ++++++++++++++++++ spindle/xrpc/validation_test.go | 8 +- spindle/xrpc/xrpc.go | 25 +- spindle/xrpc/xrpc_test.go | 110 +++- workflow/def.go | 1 + 18 files changed, 1182 insertions(+), 27 deletions(-) create mode 100644 spindle/xrpc/org_tangled_ci_get_pipeline.go create mode 100644 spindle/xrpc/org_tangled_ci_get_workflow_definition.go create mode 100644 spindle/xrpc/org_tangled_ci_pipeline.go create mode 100644 spindle/xrpc/org_tangled_ci_query_pipelines.go create mode 100644 spindle/xrpc/org_tangled_ci_subscribe_pipeline_logs.go create mode 100644 spindle/xrpc/org_tangled_secret_add_secret.go create mode 100644 spindle/xrpc/org_tangled_secret_common.go create mode 100644 spindle/xrpc/org_tangled_secret_list_secrets.go create mode 100644 spindle/xrpc/org_tangled_secret_remove_secret.go create mode 100644 spindle/xrpc/org_tangled_test.go diff --git a/spindle/db/db.go b/spindle/db/db.go index 388f9500b..c604f049b 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -107,13 +107,6 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { unique (did, rkey) ); - create table if not exists events ( - rkey text not null, - nsid text not null, - event text not null, - created integer not null - ); - create table if not exists nixos_toplevel_cache ( config_key text primary key, toplevel text not null, @@ -343,6 +336,13 @@ func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error if err := orm.RunMigration(conn, logger, "events-pipeline-index", func(tx *sql.Tx) error { _, err := tx.Exec(` + create table if not exists events ( + rkey text not null, + nsid text not null, + event text not null, + created integer not null + ); + create index if not exists idx_events_pipeline_lookup on events( coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), json_extract(event, '$.triggerMetadata.repo.did')), diff --git a/spindle/engine/engine_test.go b/spindle/engine/engine_test.go index 39ba955ae..f58670ccf 100644 --- a/spindle/engine/engine_test.go +++ b/spindle/engine/engine_test.go @@ -6,6 +6,7 @@ import ( "log/slog" "os" "path/filepath" + "slices" "strings" "sync" "testing" @@ -34,13 +35,14 @@ func (m mockStep) Command() string { return m.command } func (m mockStep) Kind() models.StepKind { return models.StepKindUser } type mockEngine struct { - mu sync.Mutex - setupCalls []models.WorkflowId - runStepCalls []models.WorkflowId - setupFunc func(ctx context.Context, wid models.WorkflowId) error - runStepFunc func(ctx context.Context, wid models.WorkflowId, idx int, wfLogger models.WorkflowLogger) error - destroyFunc func(ctx context.Context, wid models.WorkflowId) error - timeout time.Duration + mu sync.Mutex + setupCalls []models.WorkflowId + runStepCalls []models.WorkflowId + runStepSecrets [][]secrets.UnlockedSecret + setupFunc func(ctx context.Context, wid models.WorkflowId) error + runStepFunc func(ctx context.Context, wid models.WorkflowId, idx int, wfLogger models.WorkflowLogger) error + destroyFunc func(ctx context.Context, wid models.WorkflowId) error + timeout time.Duration } func (m *mockEngine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) { @@ -86,6 +88,7 @@ func (m *mockEngine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) func (m *mockEngine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, secrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { m.mu.Lock() m.runStepCalls = append(m.runStepCalls, wid) + m.runStepSecrets = append(m.runStepSecrets, slices.Clone(secrets)) fn := m.runStepFunc m.mu.Unlock() @@ -242,6 +245,48 @@ func TestCancelWorkflow_NotOverwritten(t *testing.T) { } } +func TestStartWorkflowsScopesSecretsToTrustedPipelineRepo(t *testing.T) { + vault, err := secrets.NewSQLiteManager(":memory:") + if err != nil { + t.Fatal(err) + } + repoA := syntax.DID("did:plc:workflowrepoa") + repoB := syntax.DID("did:plc:workflowrepob") + for _, secret := range []secrets.UnlockedSecret{ + {Repo: secrets.RepoIdentifier(repoA), Key: "A_SECRET", Value: "alpha"}, + {Repo: secrets.RepoIdentifier(repoB), Key: "B_SECRET", Value: "bravo"}, + } { + if err := vault.AddSecret(context.Background(), secret); err != nil { + t.Fatal(err) + } + } + run := func(trusted bool) [][]secrets.UnlockedSecret { + t.Helper() + database := newTestDB(t) + defer database.Close() + engine := &mockEngine{} + pipeline := &models.Pipeline{ + RepoDid: repoA, + TrustedSource: trusted, + Workflows: map[models.Engine][]models.Workflow{ + engine: {{Name: "ci.yml", Steps: []models.Step{mockStep{name: "test"}}}}, + }, + } + StartWorkflows(slog.Default(), vault, &config.Config{Server: config.Server{LogDir: t.TempDir()}}, nil, nil, database, nil, context.Background(), pipeline, models.PipelineId("3mu2xwiorc2xl")) + engine.mu.Lock() + defer engine.mu.Unlock() + return slices.Clone(engine.runStepSecrets) + } + + trustedSecrets := run(true) + if len(trustedSecrets) != 1 || len(trustedSecrets[0]) != 1 || trustedSecrets[0][0].Key != "A_SECRET" || trustedSecrets[0][0].Value != "alpha" { + t.Fatalf("trusted secrets = %#v", trustedSecrets) + } + if untrustedSecrets := run(false); len(untrustedSecrets) != 1 || len(untrustedSecrets[0]) != 0 { + t.Fatalf("untrusted secrets = %#v", untrustedSecrets) + } +} + func TestStartWorkflows_FlushesLogBeforeTerminalStatus(t *testing.T) { t.Parallel() diff --git a/spindle/xrpc/ci_pipeline_describe_workflow_definition.go b/spindle/xrpc/ci_pipeline_describe_workflow_definition.go index 3bb988537..20d434692 100644 --- a/spindle/xrpc/ci_pipeline_describe_workflow_definition.go +++ b/spindle/xrpc/ci_pipeline_describe_workflow_definition.go @@ -22,14 +22,14 @@ func (x *Xrpc) DescribeWorkflowDefinition(w http.ResponseWriter, r *http.Request sha := q.Get("sha") if err := requireSha(sha); err != nil { - fail(xrpcerr.NewXrpcError(xrpcerr.WithTag("InvalidRequest"), xrpcerr.WithError(err))) + fail(invalidRequest(err)) return } sourceRepoParam := q.Get("sourceRepo") sourceRepo, err := parseOptionalDID("sourceRepo", &sourceRepoParam) if err != nil { - fail(xrpcerr.NewXrpcError(xrpcerr.WithTag("InvalidRequest"), xrpcerr.WithError(err))) + fail(invalidRequest(err)) return } diff --git a/spindle/xrpc/ci_pipeline_trigger_pipeline.go b/spindle/xrpc/ci_pipeline_trigger_pipeline.go index 207875c7a..c1089278f 100644 --- a/spindle/xrpc/ci_pipeline_trigger_pipeline.go +++ b/spindle/xrpc/ci_pipeline_trigger_pipeline.go @@ -2,6 +2,7 @@ package xrpc import ( "context" + "database/sql" "encoding/json" "errors" "fmt" @@ -170,14 +171,28 @@ func (x *Xrpc) resolveOwnedRepo(ctx context.Context, actorDid syntax.DID, repoDi return repoDid, xrpcerr.XrpcError{}, true } +func writeKnownRepoError(w http.ResponseWriter, xerr xrpcerr.XrpcError) { + status := http.StatusInternalServerError + switch xerr.Tag { + case "InvalidRequest": + status = http.StatusBadRequest + case "RepoNotFound": + status = http.StatusNotFound + } + writeError(w, xerr, status) +} + func (x *Xrpc) resolveKnownRepoDid(repoDidStr string) (syntax.DID, xrpcerr.XrpcError, bool) { repoDid, err := syntax.ParseDID(repoDidStr) if err != nil { - return "", xrpcerr.GenericError(fmt.Errorf("invalid repo DID %q: %w", repoDidStr, err)), false + return "", invalidRequest(fmt.Errorf("invalid repo DID %q: %w", repoDidStr, err)), false } if _, err := x.Db.GetRepoByDid(repoDid); err != nil { - return "", xrpcerr.RepoNotFoundError, false + if errors.Is(err, sql.ErrNoRows) { + return "", xrpcerr.RepoNotFoundError, false + } + return "", xrpcerr.GenericError(err), false } return repoDid, xrpcerr.XrpcError{}, true diff --git a/spindle/xrpc/org_tangled_ci_get_pipeline.go b/spindle/xrpc/org_tangled_ci_get_pipeline.go new file mode 100644 index 000000000..8910a6f72 --- /dev/null +++ b/spindle/xrpc/org_tangled_ci_get_pipeline.go @@ -0,0 +1,33 @@ +package xrpc + +import ( + "database/sql" + "fmt" + "net/http" + + "tangled.org/core/spindle/models" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledCiGetPipeline(w http.ResponseWriter, r *http.Request) { + pipeline := r.URL.Query().Get("pipeline") + if pipeline == "" { + writeError(w, invalidRequest(fmt.Errorf("pipeline is required")), http.StatusBadRequest) + return + } + p, err := x.Db.GetPipeline(r.Context(), models.PipelineId(pipeline)) + if err != nil { + if err == sql.ErrNoRows { + writeError(w, xrpcerr.NewXrpcError(xrpcerr.WithTag("PipelineNotFound"), xrpcerr.WithError(err)), http.StatusNotFound) + return + } + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + out, err := toOrgPipeline(p) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + x.writeResponse(w, r, http.StatusOK, out) +} diff --git a/spindle/xrpc/org_tangled_ci_get_workflow_definition.go b/spindle/xrpc/org_tangled_ci_get_workflow_definition.go new file mode 100644 index 000000000..5519abbf8 --- /dev/null +++ b/spindle/xrpc/org_tangled_ci_get_workflow_definition.go @@ -0,0 +1,69 @@ +package xrpc + +import ( + "fmt" + "net/http" + + "tangled.org/core/api/org_tangled" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledCiGetWorkflowDefinition(w http.ResponseWriter, r *http.Request) { + q := r.URL.Query() + repo, xerr, ok := x.resolveKnownRepoDid(q.Get("repo")) + if !ok { + writeKnownRepoError(w, xerr) + return + } + workflowID := q.Get("workflow") + if len(workflowID) < 1 || len(workflowID) > 128 { + writeError(w, invalidRequest(fmt.Errorf("workflow length must be between 1 and 128")), http.StatusBadRequest) + return + } + workflows, err := x.Trigger.ListWorkflowDefinitions(r.Context(), repo, "HEAD") + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + for _, workflow := range workflows { + if workflow != nil && workflow.ID == workflowID { + mapped, err := toOrgWorkflowDefinition(workflow) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + x.writeResponse(w, r, http.StatusOK, org_tangled.CiGetWorkflowDefinition_Output{Workflow: mapped}) + return + } + } + writeError(w, xrpcerr.NewXrpcError(xrpcerr.WithTag("WorkflowNotFound"), xrpcerr.WithError(fmt.Errorf("workflow %q not found", workflowID))), http.StatusNotFound) +} + +func (x *Xrpc) handleOrgTangledCiListWorkflowDefinitions(w http.ResponseWriter, r *http.Request) { + q := r.URL.Query() + repo, xerr, ok := x.resolveKnownRepoDid(q.Get("repo")) + if !ok { + writeKnownRepoError(w, xerr) + return + } + ref := q.Get("ref") + if len(ref) > 256 { + writeError(w, invalidRequest(fmt.Errorf("ref is too long")), http.StatusBadRequest) + return + } + workflows, err := x.Trigger.ListWorkflowDefinitions(r.Context(), repo, ref) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + mapped := make([]*org_tangled.CiWorkflow, 0, len(workflows)) + for _, workflow := range workflows { + definition, err := toOrgWorkflowDefinition(workflow) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + mapped = append(mapped, definition) + } + x.writeResponse(w, r, http.StatusOK, org_tangled.CiListWorkflowDefinitions_Output{Workflows: mapped}) +} diff --git a/spindle/xrpc/org_tangled_ci_pipeline.go b/spindle/xrpc/org_tangled_ci_pipeline.go new file mode 100644 index 000000000..b482b799b --- /dev/null +++ b/spindle/xrpc/org_tangled_ci_pipeline.go @@ -0,0 +1,136 @@ +package xrpc + +import ( + "fmt" + + "tangled.org/core/api/org_tangled" + "tangled.org/core/spindle/models" + pipelinecodec "tangled.org/core/spindle/pipeline" + "tangled.org/core/workflow" +) + +func normalizeOrgTriggerKinds(triggers []string) []string { + out := make([]string, 0, len(triggers)) + for _, trigger := range triggers { + switch trigger { + case pipelinecodec.TriggerNSIDPush: + trigger = string(workflow.TriggerKindPush) + case pipelinecodec.TriggerNSIDPullRequest: + trigger = string(workflow.TriggerKindPullRequest) + case pipelinecodec.TriggerNSIDManual: + trigger = string(workflow.TriggerKindManual) + default: + trigger = "\x00" + trigger + } + out = append(out, trigger) + } + return out +} + +func toOrgPipeline(record *models.PipelineRecord) (*org_tangled.CiPipeline, error) { + if record == nil { + return nil, nil + } + trigger, err := orgPipelineTrigger(record) + if err != nil { + return nil, err + } + workflows := make([]*org_tangled.CiPipeline_Workflow, 0, len(record.Workflows)) + for _, workflow := range record.Workflows { + if workflow == nil { + continue + } + definition, err := toOrgWorkflowDefinition(workflow.Definition) + if err != nil { + return nil, fmt.Errorf("pipeline %s workflow %q: %w", record.ID, workflow.ID, err) + } + name := workflow.Name + workflows = append(workflows, &org_tangled.CiPipeline_Workflow{ + Definition: definition, + Error: workflow.Error, + FinishedAt: workflow.FinishedAt, + Id: workflow.ID, + Name: &name, + StartedAt: workflow.StartedAt, + Status: workflow.Status, + }) + } + createdAt := record.CreatedAt + return &org_tangled.CiPipeline{ + Commit: record.Commit, + CreatedAt: &createdAt, + Id: string(record.ID), + Repo: record.RepoDID, + SourceRepo: record.SourceRepo, + Trigger: trigger, + Workflows: workflows, + }, nil +} + +func toOrgPipelines(records []*models.PipelineRecord) []*org_tangled.CiPipeline { + out := make([]*org_tangled.CiPipeline, 0, len(records)) + for _, record := range records { + pipeline, err := toOrgPipeline(record) + if err == nil && pipeline != nil { + out = append(out, pipeline) + } + } + return out +} + +func toOrgWorkflowDefinition(definition *models.WorkflowDefinition) (*org_tangled.CiWorkflow, error) { + if definition == nil { + return nil, fmt.Errorf("missing definition") + } + var source org_tangled.CiWorkflow_Source + switch { + case definition.Source.File != nil: + file := definition.Source.File + mapped := &org_tangled.CiWorkflow_FileSource{ + Repo: file.Repo, Commit: file.Commit, Path: file.Path, + } + if file.Lines != nil { + mapped.Lines = &org_tangled.CiWorkflow_FileSource_Lines{Start: file.Lines.Start, End: file.Lines.End} + } + source.CiWorkflow_FileSource = mapped + case definition.Source.External != nil: + external := definition.Source.External + mapped := &org_tangled.CiWorkflow_ExternalSource{Link: external.Link} + if external.Name != "" { + mapped.Name = &external.Name + } + source.CiWorkflow_ExternalSource = mapped + default: + return nil, fmt.Errorf("missing definition source") + } + mapped := &org_tangled.CiWorkflow{ + Hash: definition.Hash, Id: definition.ID, Source: &source, + Triggers: append([]string{}, definition.Triggers...), + } + if definition.Name != "" { + mapped.Name = &definition.Name + } + return mapped, nil +} + +func orgPipelineTrigger(record *models.PipelineRecord) (*org_tangled.CiPipeline_Trigger, error) { + switch { + case record.Trigger.Push != nil: + return &org_tangled.CiPipeline_Trigger{EventPush: &org_tangled.EventPush{}}, nil + case record.Trigger.PullRequest != nil: + return &org_tangled.CiPipeline_Trigger{EventPullRequest: &org_tangled.EventPullRequest{}}, nil + case record.Trigger.Manual != nil: + manual := record.Trigger.Manual + inputs := make([]*org_tangled.CiTriggerManual_Pair, 0, len(manual.Inputs)) + for _, input := range manual.Inputs { + if input != nil { + inputs = append(inputs, &org_tangled.CiTriggerManual_Pair{Key: input.Key, Value: input.Value}) + } + } + return &org_tangled.CiPipeline_Trigger{CiTriggerManual: &org_tangled.CiTriggerManual{ + Repo: record.RepoDID, Commit: record.Commit, Ref: manual.Ref, Inputs: inputs, + }}, nil + default: + return nil, fmt.Errorf("pipeline %s has no supported trigger", record.ID) + } +} diff --git a/spindle/xrpc/org_tangled_ci_query_pipelines.go b/spindle/xrpc/org_tangled_ci_query_pipelines.go new file mode 100644 index 000000000..2eaff03c6 --- /dev/null +++ b/spindle/xrpc/org_tangled_ci_query_pipelines.go @@ -0,0 +1,43 @@ +package xrpc + +import ( + "fmt" + "net/http" + "strconv" + + "tangled.org/core/api/org_tangled" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledCiQueryPipelines(w http.ResponseWriter, r *http.Request) { + fail := func(err error) { + writeError(w, invalidRequest(err), http.StatusBadRequest) + } + q := r.URL.Query() + repo, xerr, ok := x.resolveKnownRepoDid(q.Get("repo")) + if !ok { + writeKnownRepoError(w, xerr) + return + } + limit := 50 + if value := q.Get("limit"); value != "" { + parsed, err := strconv.Atoi(value) + if err != nil || parsed < 1 || parsed > 250 { + fail(fmt.Errorf("limit must be between 1 and 250")) + return + } + limit = parsed + } + cursor := q.Get("cursor") + pipelines, nextCursor, total, err := x.Db.QueryPipelines(r.Context(), repo.String(), q["commits"], cursor, normalizeOrgTriggerKinds(q["triggers"]), limit) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + orgPipelines := toOrgPipelines(pipelines) + out := org_tangled.CiQueryPipelines_Output{Pipelines: orgPipelines, Total: total} + if nextCursor != "" { + out.Cursor = &nextCursor + } + x.writeResponse(w, r, http.StatusOK, out) +} diff --git a/spindle/xrpc/org_tangled_ci_subscribe_pipeline_logs.go b/spindle/xrpc/org_tangled_ci_subscribe_pipeline_logs.go new file mode 100644 index 000000000..0d3bf5ade --- /dev/null +++ b/spindle/xrpc/org_tangled_ci_subscribe_pipeline_logs.go @@ -0,0 +1,41 @@ +package xrpc + +import ( + "database/sql" + "errors" + "fmt" + "net/http" + + "tangled.org/core/spindle/models" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledCiSubscribePipelineLogs(w http.ResponseWriter, r *http.Request) { + pipeline := r.URL.Query().Get("pipeline") + pipelineID := models.PipelineId(pipeline) + stored, err := x.Db.GetPipeline(r.Context(), pipelineID) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + writeError(w, xrpcerr.NewXrpcError(xrpcerr.WithTag("PipelineNotFound"), xrpcerr.WithError(err)), http.StatusNotFound) + return + } + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + workflows := r.URL.Query()["workflows"] + if len(workflows) != 0 { + known := make(map[string]struct{}, len(stored.Workflows)) + for _, workflow := range stored.Workflows { + if workflow != nil { + known[workflow.Name] = struct{}{} + } + } + for _, workflow := range workflows { + if _, ok := known[workflow]; !ok { + writeError(w, xrpcerr.NewXrpcError(xrpcerr.WithTag("WorkflowNotFound"), xrpcerr.WithError(fmt.Errorf("workflow %q is not part of pipeline %s", workflow, pipeline))), http.StatusNotFound) + return + } + } + } + x.handleSubscribeLogs(w, r, pipelineID, workflows) +} diff --git a/spindle/xrpc/org_tangled_secret_add_secret.go b/spindle/xrpc/org_tangled_secret_add_secret.go new file mode 100644 index 000000000..6a6b52260 --- /dev/null +++ b/spindle/xrpc/org_tangled_secret_add_secret.go @@ -0,0 +1,48 @@ +package xrpc + +import ( + "encoding/json" + "errors" + "fmt" + "net/http" + "time" + + "tangled.org/core/api/org_tangled" + "tangled.org/core/spindle/secrets" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledSecretAddSecret(w http.ResponseWriter, r *http.Request) { + var input org_tangled.SecretAddSecret_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, invalidRequest(err), http.StatusBadRequest) + return + } + if err := validateOrgSecretKey(input.Key); err != nil { + writeError(w, invalidRequest(err), http.StatusBadRequest) + return + } + if len(input.Value) < 1 || len(input.Value) > 200 { + writeError(w, invalidRequest(fmt.Errorf("secret value length must be between 1 and 200")), http.StatusBadRequest) + return + } + actor, repo, ok := x.authorizeOrgSecretRepo(w, r, input.Repo) + if !ok { + return + } + if err := x.Vault.AddSecret(r.Context(), secrets.UnlockedSecret{ + Repo: secrets.RepoIdentifier(repo.String()), Key: input.Key, Value: input.Value, + CreatedAt: time.Now(), CreatedBy: actor, + }); err != nil { + switch { + case errors.Is(err, secrets.ErrInvalidKeyIdent): + writeError(w, invalidRequest(err), http.StatusBadRequest) + case errors.Is(err, secrets.ErrKeyAlreadyPresent): + writeError(w, xrpcerr.NewXrpcError(xrpcerr.WithTag("SecretAlreadyExists"), xrpcerr.WithError(err)), http.StatusConflict) + default: + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + } + return + } + w.WriteHeader(http.StatusOK) +} diff --git a/spindle/xrpc/org_tangled_secret_common.go b/spindle/xrpc/org_tangled_secret_common.go new file mode 100644 index 000000000..da535f0b7 --- /dev/null +++ b/spindle/xrpc/org_tangled_secret_common.go @@ -0,0 +1,40 @@ +package xrpc + +import ( + "fmt" + "net/http" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/rbac" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func validateOrgSecretKey(key string) error { + if len(key) < 1 || len(key) > 50 { + return fmt.Errorf("secret key length must be between 1 and 50") + } + return nil +} + +func (x *Xrpc) authorizeOrgSecretRepo(w http.ResponseWriter, r *http.Request, repo string) (syntax.DID, syntax.DID, bool) { + actor, ok := r.Context().Value(ActorDid).(syntax.DID) + if !ok { + writeError(w, xrpcerr.MissingActorDidError, http.StatusUnauthorized) + return "", "", false + } + repoDid, xerr, ok := x.resolveKnownRepoDid(repo) + if !ok { + writeKnownRepoError(w, xerr) + return "", "", false + } + allowed, err := x.Enforcer.IsSettingsAllowed(actor.String(), rbac.ThisServer, repoDid.String()) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return "", "", false + } + if !allowed { + writeError(w, xrpcerr.AccessControlError(actor.String()), http.StatusUnauthorized) + return "", "", false + } + return actor, repoDid, true +} diff --git a/spindle/xrpc/org_tangled_secret_list_secrets.go b/spindle/xrpc/org_tangled_secret_list_secrets.go new file mode 100644 index 000000000..17da6e02a --- /dev/null +++ b/spindle/xrpc/org_tangled_secret_list_secrets.go @@ -0,0 +1,29 @@ +package xrpc + +import ( + "net/http" + "time" + + "tangled.org/core/api/org_tangled" + "tangled.org/core/spindle/secrets" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledSecretListSecrets(w http.ResponseWriter, r *http.Request) { + _, repo, ok := x.authorizeOrgSecretRepo(w, r, r.URL.Query().Get("repo")) + if !ok { + return + } + locked, err := x.Vault.GetSecretsLocked(r.Context(), secrets.RepoIdentifier(repo.String())) + if err != nil { + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } + out := org_tangled.SecretListSecrets_Output{Secrets: make([]*org_tangled.SecretListSecrets_Secret, 0, len(locked))} + for _, secret := range locked { + out.Secrets = append(out.Secrets, &org_tangled.SecretListSecrets_Secret{ + CreatedAt: secret.CreatedAt.Format(time.RFC3339), CreatedBy: secret.CreatedBy.String(), Key: secret.Key, + }) + } + x.writeResponse(w, r, http.StatusOK, out) +} diff --git a/spindle/xrpc/org_tangled_secret_remove_secret.go b/spindle/xrpc/org_tangled_secret_remove_secret.go new file mode 100644 index 000000000..80d9115f7 --- /dev/null +++ b/spindle/xrpc/org_tangled_secret_remove_secret.go @@ -0,0 +1,39 @@ +package xrpc + +import ( + "encoding/json" + "errors" + "net/http" + + "tangled.org/core/api/org_tangled" + "tangled.org/core/spindle/secrets" + xrpcerr "tangled.org/core/xrpc/errors" +) + +func (x *Xrpc) handleOrgTangledSecretRemoveSecret(w http.ResponseWriter, r *http.Request) { + var input org_tangled.SecretRemoveSecret_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeError(w, invalidRequest(err), http.StatusBadRequest) + return + } + if err := validateOrgSecretKey(input.Key); err != nil { + writeError(w, invalidRequest(err), http.StatusBadRequest) + return + } + _, repo, ok := x.authorizeOrgSecretRepo(w, r, input.Repo) + if !ok { + return + } + if err := x.Vault.RemoveSecret(r.Context(), secrets.Secret[any]{Repo: secrets.RepoIdentifier(repo.String()), Key: input.Key}); err != nil { + switch { + case errors.Is(err, secrets.ErrInvalidKeyIdent): + writeError(w, invalidRequest(err), http.StatusBadRequest) + case errors.Is(err, secrets.ErrKeyNotFound): + writeError(w, xrpcerr.NewXrpcError(xrpcerr.WithTag("SecretNotFound"), xrpcerr.WithError(err)), http.StatusNotFound) + default: + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + } + return + } + w.WriteHeader(http.StatusOK) +} diff --git a/spindle/xrpc/org_tangled_test.go b/spindle/xrpc/org_tangled_test.go new file mode 100644 index 000000000..c4927cc68 --- /dev/null +++ b/spindle/xrpc/org_tangled_test.go @@ -0,0 +1,491 @@ +package xrpc + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "log/slog" + "net/http" + "net/http/httptest" + "slices" + "strings" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/org_tangled" + "tangled.org/core/api/tangled" + "tangled.org/core/rbac" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" + "tangled.org/core/spindle/secrets" +) + +func newOrgPipelineXrpc(t *testing.T) (*Xrpc, models.PipelineId, string) { + t.Helper() + d, e := newTestXrpcDB(t) + repoDid := syntax.DID("did:plc:orgpipelinerepo") + if err := d.AddRepo(db.Repo{Knot: "knot.test", Owner: "did:plc:owner", Rkey: "repo", RepoDid: repoDid}); err != nil { + t.Fatal(err) + } + id := models.PipelineId("pipeline_custom") + repo := repoDid.String() + createTestPipeline(t, d, id, tangled.Pipeline{ + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{Kind: "push", Repo: &tangled.Pipeline_TriggerRepo{Knot: "knot.test", RepoDid: &repo}, Push: &tangled.Pipeline_PushTriggerData{NewSha: "1111111111111111111111111111111111111111", Ref: "refs/heads/main"}}, + Workflows: []*tangled.Pipeline_Workflow{{Name: "ci.yml", Raw: "when:\n - event: pull_request\n"}}, + }) + return &Xrpc{Logger: slog.Default(), Db: d, Enforcer: e}, id, repo +} + +func TestOrgCiGetAndQueryPipelines(t *testing.T) { + x, id, repo := newOrgPipelineXrpc(t) + + w := httptest.NewRecorder() + r := httptest.NewRequest(http.MethodGet, "/?pipeline="+id.String(), nil) + x.handleOrgTangledCiGetPipeline(w, r) + if w.Code != http.StatusOK { + t.Fatalf("get pipeline = %d: %s", w.Code, w.Body.String()) + } + var got org_tangled.CiPipeline + if err := json.Unmarshal(w.Body.Bytes(), &got); err != nil || got.Id != id.String() || len(got.Workflows) != 1 { + t.Fatalf("pipeline = %+v, %v", got, err) + } + if !slices.Equal(got.Workflows[0].Definition.Triggers, []string{"org.tangled.event.pullRequest"}) { + t.Fatalf("persisted workflow triggers = %v", got.Workflows[0].Definition.Triggers) + } + + w = httptest.NewRecorder() + r = httptest.NewRequest(http.MethodGet, "/?pipeline=missing_pipeline", nil) + x.handleOrgTangledCiGetPipeline(w, r) + if w.Code != http.StatusNotFound { + t.Fatalf("missing pipeline = %d", w.Code) + } + + w = httptest.NewRecorder() + r = httptest.NewRequest(http.MethodGet, "/?repo="+repo+"&triggers=org.tangled.event.push&limit=50", nil) + x.handleOrgTangledCiQueryPipelines(w, r) + if w.Code != http.StatusOK { + t.Fatalf("query pipelines = %d: %s", w.Code, w.Body.String()) + } + var queried org_tangled.CiQueryPipelines_Output + if err := json.Unmarshal(w.Body.Bytes(), &queried); err != nil || queried.Total != 1 || len(queried.Pipelines) != 1 { + t.Fatalf("query = %+v, %v", queried, err) + } + + w = httptest.NewRecorder() + r = httptest.NewRequest(http.MethodGet, "/?repo="+repo+"&triggers=push&limit=50", nil) + x.handleOrgTangledCiQueryPipelines(w, r) + if w.Code != http.StatusOK { + t.Fatalf("plain trigger alias = %d: %s", w.Code, w.Body.String()) + } + queried = org_tangled.CiQueryPipelines_Output{} + if err := json.Unmarshal(w.Body.Bytes(), &queried); err != nil || queried.Total != 0 || len(queried.Pipelines) != 0 { + t.Fatalf("plain trigger alias matched = %+v, %v", queried, err) + } +} + +func TestOrgSecretInputValidation(t *testing.T) { + d, enforcer := newTestXrpcDB(t) + actor := syntax.DID("did:plc:secretvalidator") + repo := syntax.DID("did:plc:secretvalidationrepo") + if err := d.AddRepo(db.Repo{Knot: "knot.test", Owner: actor, Rkey: "repo", RepoDid: repo}); err != nil { + t.Fatal(err) + } + if err := enforcer.AddRepo(actor.String(), rbac.ThisServer, repo.String()); err != nil { + t.Fatal(err) + } + vault, err := secrets.NewSQLiteManager(":memory:") + if err != nil { + t.Fatal(err) + } + x := &Xrpc{Logger: slog.Default(), Db: d, Enforcer: enforcer, Vault: vault} + valid := org_tangled.SecretAddSecret_Input{Repo: repo.String(), Key: "VALID_KEY", Value: "value"} + cases := map[string][]byte{ + "malformed json": []byte("{"), + "empty key": mustJSON(t, org_tangled.SecretAddSecret_Input{Repo: repo.String(), Value: "value"}), + "long key": mustJSON(t, org_tangled.SecretAddSecret_Input{Repo: repo.String(), Key: strings.Repeat("k", 51), Value: "value"}), + "empty value": mustJSON(t, org_tangled.SecretAddSecret_Input{Repo: repo.String(), Key: valid.Key}), + "long value": mustJSON(t, org_tangled.SecretAddSecret_Input{Repo: repo.String(), Key: valid.Key, Value: strings.Repeat("v", 201)}), + } + for name, body := range cases { + t.Run(name, func(t *testing.T) { + r := httptest.NewRequest(http.MethodPost, "/", bytes.NewReader(body)) + r = r.WithContext(context.WithValue(r.Context(), ActorDid, actor)) + w := httptest.NewRecorder() + x.handleOrgTangledSecretAddSecret(w, r) + if w.Code != http.StatusBadRequest || !strings.Contains(w.Body.String(), "InvalidRequest") { + t.Fatalf("response = %d: %s", w.Code, w.Body.String()) + } + }) + } + stored, err := vault.GetSecretsUnlocked(context.Background(), secrets.RepoIdentifier(repo)) + if err != nil || len(stored) != 0 { + t.Fatalf("invalid input changed vault: %+v, %v", stored, err) + } + + r := httptest.NewRequest(http.MethodPost, "/", bytes.NewReader(mustJSON(t, org_tangled.SecretAddSecret_Input{Repo: repo.String(), Key: "key with spaces", Value: "value"}))) + r = r.WithContext(context.WithValue(r.Context(), ActorDid, actor)) + w := httptest.NewRecorder() + x.handleOrgTangledSecretAddSecret(w, r) + if w.Code != http.StatusOK { + t.Fatalf("schema-valid key response = %d: %s", w.Code, w.Body.String()) + } +} + +func TestOrgCiInputBounds(t *testing.T) { + x, _, repo := newOrgPipelineXrpc(t) + for name, path := range map[string]string{ + "zero limit": "/?repo=" + repo + "&limit=0", + "large limit": "/?repo=" + repo + "&limit=251", + } { + t.Run(name, func(t *testing.T) { + w := httptest.NewRecorder() + x.handleOrgTangledCiQueryPipelines(w, httptest.NewRequest(http.MethodGet, path, nil)) + if w.Code != http.StatusBadRequest || !strings.Contains(w.Body.String(), "InvalidRequest") { + t.Fatalf("response = %d: %s", w.Code, w.Body.String()) + } + }) + } + for name, path := range map[string]string{ + "long commit": "/?repo=" + repo + "&commits=" + strings.Repeat("a", 65), + "short commit": "/?repo=" + repo + "&commits=" + strings.Repeat("a", 39), + "non-hex commit": "/?repo=" + repo + "&commits=" + strings.Repeat("z", 40), + "unknown trigger": "/?repo=" + repo + "&triggers=scheduled", + "invalid cursor": "/?repo=" + repo + "&cursor=nope", + } { + t.Run(name, func(t *testing.T) { + w := httptest.NewRecorder() + x.handleOrgTangledCiQueryPipelines(w, httptest.NewRequest(http.MethodGet, path, nil)) + if w.Code != http.StatusOK { + t.Fatalf("response = %d: %s", w.Code, w.Body.String()) + } + }) + } + w := httptest.NewRecorder() + x.handleOrgTangledCiGetPipeline(w, httptest.NewRequest(http.MethodGet, "/", nil)) + if w.Code != http.StatusBadRequest || !strings.Contains(w.Body.String(), "InvalidRequest") { + t.Fatalf("missing pipeline = %d: %s", w.Code, w.Body.String()) + } + + for name, handlerPath := range map[string]struct { + handler func(http.ResponseWriter, *http.Request) + path string + }{ + "empty workflow": {x.handleOrgTangledCiGetWorkflowDefinition, "/?repo=" + repo}, + "long workflow": {x.handleOrgTangledCiGetWorkflowDefinition, "/?repo=" + repo + "&workflow=" + strings.Repeat("w", 129)}, + "long ref": {x.handleOrgTangledCiListWorkflowDefinitions, "/?repo=" + repo + "&ref=" + strings.Repeat("r", 257)}, + } { + t.Run(name, func(t *testing.T) { + w := httptest.NewRecorder() + handlerPath.handler(w, httptest.NewRequest(http.MethodGet, handlerPath.path, nil)) + if w.Code != http.StatusBadRequest || !strings.Contains(w.Body.String(), "InvalidRequest") { + t.Fatalf("response = %d: %s", w.Code, w.Body.String()) + } + }) + } +} + +func mustJSON(t *testing.T, value any) []byte { + t.Helper() + body, err := json.Marshal(value) + if err != nil { + t.Fatal(err) + } + return body +} + +func TestOrgSecretsAreIsolatedByResolvedRepo(t *testing.T) { + d, enforcer := newTestXrpcDB(t) + actor := syntax.DID("did:plc:settingsactor") + repoA := syntax.DID("did:plc:secretrepoa") + repoB := syntax.DID("did:plc:secretrepob") + for i, repo := range []syntax.DID{repoA, repoB} { + if err := d.AddRepo(db.Repo{Knot: "knot.test", Owner: actor, Rkey: syntax.RecordKey(fmt.Sprintf("repo-%d", i)), RepoDid: repo}); err != nil { + t.Fatal(err) + } + } + if err := enforcer.AddRepo(actor.String(), rbac.ThisServer, repoA.String()); err != nil { + t.Fatal(err) + } + vault, err := secrets.NewSQLiteManager(":memory:") + if err != nil { + t.Fatal(err) + } + for _, repo := range []syntax.DID{repoA, repoB} { + if err := vault.AddSecret(context.Background(), secrets.UnlockedSecret{Repo: secrets.RepoIdentifier(repo), Key: "SAME_KEY", Value: repo.String(), CreatedAt: time.Now(), CreatedBy: actor}); err != nil { + t.Fatal(err) + } + } + x := &Xrpc{Logger: slog.Default(), Db: d, Enforcer: enforcer, Vault: vault} + request := func(repo syntax.DID) *httptest.ResponseRecorder { + r := httptest.NewRequest(http.MethodGet, "/?repo="+repo.String(), nil) + r = r.WithContext(context.WithValue(r.Context(), ActorDid, actor)) + w := httptest.NewRecorder() + x.handleOrgTangledSecretListSecrets(w, r) + return w + } + if w := request(repoA); w.Code != http.StatusOK { + t.Fatalf("repo A list = %d: %s", w.Code, w.Body.String()) + } else { + var listed org_tangled.SecretListSecrets_Output + if err := json.Unmarshal(w.Body.Bytes(), &listed); err != nil || len(listed.Secrets) != 1 || listed.Secrets[0].Key != "SAME_KEY" { + t.Fatalf("repo A secrets = %+v, %v", listed.Secrets, err) + } + } + if w := request(repoB); w.Code != http.StatusUnauthorized || !strings.Contains(w.Body.String(), "AccessControl") { + t.Fatalf("repo B list = %d: %s", w.Code, w.Body.String()) + } + removeBody, _ := json.Marshal(org_tangled.SecretRemoveSecret_Input{Repo: repoB.String(), Key: "SAME_KEY"}) + r := httptest.NewRequest(http.MethodPost, "/", strings.NewReader(string(removeBody))) + r = r.WithContext(context.WithValue(r.Context(), ActorDid, actor)) + w := httptest.NewRecorder() + x.handleOrgTangledSecretRemoveSecret(w, r) + if w.Code != http.StatusUnauthorized { + t.Fatalf("repo B remove = %d: %s", w.Code, w.Body.String()) + } + remaining, err := vault.GetSecretsUnlocked(context.Background(), secrets.RepoIdentifier(repoB)) + if err != nil || len(remaining) != 1 || remaining[0].Value != repoB.String() { + t.Fatalf("repo B secret changed: %+v, %v", remaining, err) + } +} + +func TestOrgSubscribePipelineLogsValidatesPipelineAndWorkflows(t *testing.T) { + x, id, _ := newOrgPipelineXrpc(t) + + w := httptest.NewRecorder() + r := httptest.NewRequest(http.MethodGet, "/?pipeline=3mu2xwiorc2xm", nil) + x.handleOrgTangledCiSubscribePipelineLogs(w, r) + if w.Code != http.StatusNotFound || !strings.Contains(w.Body.String(), "PipelineNotFound") { + t.Fatalf("missing pipeline = %d: %s", w.Code, w.Body.String()) + } + + otherID := models.PipelineId("3mu2xwiorc2xn") + repo := "did:plc:orgpipelinerepo" + createTestPipeline(t, x.Db, otherID, tangled.Pipeline{ + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{Kind: "push", Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &repo}, Push: &tangled.Pipeline_PushTriggerData{NewSha: strings.Repeat("2", 40)}}, + Workflows: []*tangled.Pipeline_Workflow{{Name: "other.yml"}}, + }) + w = httptest.NewRecorder() + r = httptest.NewRequest(http.MethodGet, "/?pipeline="+id.String()+"&workflows=other.yml", nil) + x.handleOrgTangledCiSubscribePipelineLogs(w, r) + if w.Code != http.StatusNotFound || !strings.Contains(w.Body.String(), "WorkflowNotFound") { + t.Fatalf("missing workflow = %d: %s", w.Code, w.Body.String()) + } +} + +func TestResolveKnownRepoDidContract(t *testing.T) { + x, _, knownRepo := newOrgPipelineXrpc(t) + if got, xerr, ok := x.resolveKnownRepoDid(knownRepo); !ok || got.String() != knownRepo || xerr.Tag != "" { + t.Fatalf("known repo = %q, %+v, %v", got, xerr, ok) + } + for _, test := range []struct { + name, repo, tag string + }{ + {"unknown", "did:plc:unknownrepo", "RepoNotFound"}, + {"malformed", "not-a-did", "InvalidRequest"}, + {"empty", "", "InvalidRequest"}, + } { + t.Run(test.name, func(t *testing.T) { + if _, xerr, ok := x.resolveKnownRepoDid(test.repo); ok || xerr.Tag != test.tag { + t.Fatalf("resolve = %+v, %v; want %s", xerr, ok, test.tag) + } + }) + } +} + +func TestAuthorizeOrgSecretRepoRequiresTypedActor(t *testing.T) { + x, _, repo := newOrgPipelineXrpc(t) + r := httptest.NewRequest(http.MethodGet, "/", nil) + r = r.WithContext(context.WithValue(r.Context(), ActorDid, "did:plc:not-typed")) + w := httptest.NewRecorder() + if _, _, ok := x.authorizeOrgSecretRepo(w, r, repo); ok { + t.Fatal("string actor unexpectedly authorized") + } + if w.Code != http.StatusUnauthorized || !strings.Contains(w.Body.String(), "MissingActorDid") { + t.Fatalf("response = %d: %s", w.Code, w.Body.String()) + } +} + +func TestOrgRepoScopedEndpointsReturnRepoNotFound(t *testing.T) { + x, _, _ := newOrgPipelineXrpc(t) + for name, handler := range map[string]func(http.ResponseWriter, *http.Request){ + "query pipelines": x.handleOrgTangledCiQueryPipelines, + "get workflow": x.handleOrgTangledCiGetWorkflowDefinition, + "list workflows": x.handleOrgTangledCiListWorkflowDefinitions, + } { + t.Run(name, func(t *testing.T) { + w := httptest.NewRecorder() + r := httptest.NewRequest(http.MethodGet, "/?repo=did:plc:unknownrepo&workflow=ci.yml", nil) + handler(w, r) + if w.Code != http.StatusNotFound || !strings.Contains(w.Body.String(), "RepoNotFound") { + t.Fatalf("response = %d: %s", w.Code, w.Body.String()) + } + }) + } +} + +func TestOrgWorkflowDefinitionHandlers(t *testing.T) { + d, e := newTestXrpcDB(t) + repo := syntax.DID("did:plc:workflowrepo") + if err := d.AddRepo(db.Repo{Knot: "knot.test", Owner: "did:plc:owner", Rkey: "repo", RepoDid: repo}); err != nil { + t.Fatal(err) + } + name := "ci.yml" + definition := &models.WorkflowDefinition{ + ID: name, Name: name, + Source: models.WorkflowDefinitionSource{File: &models.WorkflowFileSource{ + Repo: repo.String(), Commit: "1111111111111111111111111111111111111111", Path: ".tangled/workflows/ci.yml", + }}, + Triggers: []string{"org.tangled.event.push"}, + } + x := &Xrpc{Logger: slog.Default(), Db: d, Enforcer: e, Trigger: &mockTrigger{workflowDefinitions: []*models.WorkflowDefinition{definition}}} + + w := httptest.NewRecorder() + r := httptest.NewRequest(http.MethodGet, "/?repo="+repo.String()+"&workflow="+name, nil) + x.handleOrgTangledCiGetWorkflowDefinition(w, r) + if w.Code != http.StatusOK { + t.Fatalf("get workflow = %d: %s", w.Code, w.Body.String()) + } + + w = httptest.NewRecorder() + r = httptest.NewRequest(http.MethodGet, "/?repo="+repo.String(), nil) + x.handleOrgTangledCiListWorkflowDefinitions(w, r) + if w.Code != http.StatusOK { + t.Fatalf("list workflows = %d: %s", w.Code, w.Body.String()) + } +} + +func TestOrgPipelineAdapter(t *testing.T) { + source := "did:plc:source" + record := &models.PipelineRecord{ + ID: "3mu2xwiorc2xl", RepoDID: "did:plc:target", SourceRepo: &source, + Commit: "1111111111111111111111111111111111111111", CreatedAt: time.Now().Format(time.RFC3339), + Trigger: models.PipelineTrigger{Kind: "pull_request", PullRequest: &models.PipelinePullRequestTrigger{ + SourceCommit: "1111111111111111111111111111111111111111", TargetBranch: "main", + }}, + Workflows: []*models.PipelineWorkflow{ + nil, + {ID: "ci.yml", Name: "ci.yml", Status: "running", Definition: &models.WorkflowDefinition{ + ID: "ci.yml", Name: "ci.yml", Triggers: []string{"org.tangled.event.push"}, + Source: models.WorkflowDefinitionSource{File: &models.WorkflowFileSource{ + Repo: source, Commit: "1111111111111111111111111111111111111111", Path: ".tangled/workflows/ci.yml", + }}, + }}, + }, + } + pipeline, err := toOrgPipeline(record) + if err != nil { + t.Fatal(err) + } + if len(pipeline.Workflows) != 1 { + t.Fatalf("workflows = %+v", pipeline.Workflows) + } + workflow := pipeline.Workflows[0] + sourceFile := workflow.Definition.Source.CiWorkflow_FileSource + if workflow.Id != "ci.yml" || sourceFile.Repo != source || !slices.Equal(workflow.Definition.Triggers, []string{"org.tangled.event.push"}) { + t.Fatalf("workflow = %+v", workflow) + } + if pipeline.Trigger.EventPullRequest == nil { + t.Fatalf("trigger = %+v", pipeline.Trigger) + } + encoded, err := json.Marshal(pipeline) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(encoded), `"definition"`) || strings.Contains(string(encoded), `"defintion"`) || !strings.Contains(string(encoded), `"$type":"org.tangled.event.pullRequest"`) { + t.Fatalf("pipeline JSON = %s", encoded) + } +} + +func TestNormalizeOrgTriggerKinds(t *testing.T) { + got := normalizeOrgTriggerKinds([]string{"org.tangled.event.push", "org.tangled.event.pullRequest", "org.tangled.ci.trigger.manual", "push"}) + want := []string{"push", "pull_request", "manual", "\x00push"} + if !slices.Equal(got, want) { + t.Fatalf("normalized trigger kinds = %v, want %v", got, want) + } +} + +func TestOrgSecretConflictErrors(t *testing.T) { + d, enforcer := newTestXrpcDB(t) + actor := syntax.DID("did:plc:secret-errors") + repo := syntax.DID("did:plc:secret-errors-repo") + if err := d.AddRepo(db.Repo{Knot: "knot.test", Owner: actor, Rkey: "repo", RepoDid: repo}); err != nil { + t.Fatal(err) + } + if err := enforcer.AddRepo(actor.String(), rbac.ThisServer, repo.String()); err != nil { + t.Fatal(err) + } + vault, err := secrets.NewSQLiteManager(":memory:") + if err != nil { + t.Fatal(err) + } + x := &Xrpc{Logger: slog.Default(), Db: d, Enforcer: enforcer, Vault: vault} + request := func(handler func(http.ResponseWriter, *http.Request), value any) *httptest.ResponseRecorder { + request := httptest.NewRequest(http.MethodPost, "/", bytes.NewReader(mustJSON(t, value))) + request = request.WithContext(context.WithValue(request.Context(), ActorDid, actor)) + response := httptest.NewRecorder() + handler(response, request) + return response + } + input := org_tangled.SecretAddSecret_Input{Repo: repo.String(), Key: "KEY", Value: "value"} + if response := request(x.handleOrgTangledSecretAddSecret, input); response.Code != http.StatusOK { + t.Fatalf("first add = %d: %s", response.Code, response.Body.String()) + } + if response := request(x.handleOrgTangledSecretAddSecret, input); response.Code != http.StatusConflict || !strings.Contains(response.Body.String(), "SecretAlreadyExists") { + t.Fatalf("duplicate add = %d: %s", response.Code, response.Body.String()) + } + missing := org_tangled.SecretRemoveSecret_Input{Repo: repo.String(), Key: "MISSING"} + if response := request(x.handleOrgTangledSecretRemoveSecret, missing); response.Code != http.StatusNotFound || !strings.Contains(response.Body.String(), "SecretNotFound") { + t.Fatalf("missing remove = %d: %s", response.Code, response.Body.String()) + } +} + +func TestOrgPipelineTriggerUnion(t *testing.T) { + base := models.PipelineRecord{ + ID: "pipeline", RepoDID: "did:plc:repo", Commit: strings.Repeat("1", 40), CreatedAt: time.Now().Format(time.RFC3339), + Workflows: []*models.PipelineWorkflow{{ + ID: "ci.yml", Name: "ci.yml", Status: "pending", + Definition: &models.WorkflowDefinition{ + ID: "ci.yml", Name: "ci.yml", Triggers: []string{}, + Source: models.WorkflowDefinitionSource{External: &models.WorkflowExternalSource{Name: "generated", Link: "https://example.com/workflow"}}, + }, + }}, + } + ref := "refs/heads/main" + cases := []struct { + name string + set func(*models.PipelineRecord) + typeID string + }{ + {"push", func(record *models.PipelineRecord) { + record.Trigger = models.PipelineTrigger{Kind: "push", Push: &models.PipelinePushTrigger{}} + }, "org.tangled.event.push"}, + {"pull_request", func(record *models.PipelineRecord) { + record.Trigger = models.PipelineTrigger{Kind: "pull_request", PullRequest: &models.PipelinePullRequestTrigger{}} + }, "org.tangled.event.pullRequest"}, + {"manual", func(record *models.PipelineRecord) { + record.Trigger = models.PipelineTrigger{Kind: "manual", Manual: &models.PipelineManualTrigger{Ref: &ref, Inputs: []*models.PipelineInputPair{{Key: "target", Value: "all"}}}} + }, "org.tangled.ci.trigger.manual"}, + } + for _, test := range cases { + t.Run(test.name, func(t *testing.T) { + record := base + test.set(&record) + pipeline, err := toOrgPipeline(&record) + if err != nil { + t.Fatal(err) + } + encoded, err := json.Marshal(pipeline) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(encoded), `"$type":"`+test.typeID+`"`) { + t.Fatalf("pipeline JSON = %s", encoded) + } + if pipeline.Workflows[0].Definition.Source.CiWorkflow_ExternalSource == nil { + t.Fatalf("external source was not preserved: %+v", pipeline.Workflows[0].Definition) + } + }) + } +} diff --git a/spindle/xrpc/validation_test.go b/spindle/xrpc/validation_test.go index 1450b9a8b..452ee170c 100644 --- a/spindle/xrpc/validation_test.go +++ b/spindle/xrpc/validation_test.go @@ -9,10 +9,12 @@ func TestRequireSha(t *testing.T) { if err := requireSha(strings.Repeat("a", 40)); err != nil { t.Fatalf("valid SHA rejected: %v", err) } + if err := requireSha(strings.Repeat("z", 40)); err != nil { + t.Fatalf("40-character non-hex SHA rejected: %v", err) + } for name, sha := range map[string]string{ - "short": strings.Repeat("a", 39), - "long": strings.Repeat("a", 41), - "non-hex": strings.Repeat("z", 40), + "short": strings.Repeat("a", 39), + "long": strings.Repeat("a", 41), } { t.Run(name, func(t *testing.T) { if err := requireSha(sha); err == nil { diff --git a/spindle/xrpc/xrpc.go b/spindle/xrpc/xrpc.go index a24aafa20..acab27809 100644 --- a/spindle/xrpc/xrpc.go +++ b/spindle/xrpc/xrpc.go @@ -3,7 +3,6 @@ package xrpc import ( "context" _ "embed" - "encoding/hex" "encoding/json" "errors" "fmt" @@ -36,9 +35,6 @@ func requireSha(sha string) error { if len(sha) != 40 { return fmt.Errorf("sha must be a 40-character commit hash") } - if _, err := hex.DecodeString(sha); err != nil { - return fmt.Errorf("sha must be hexadecimal: %w", err) - } return nil } @@ -47,6 +43,7 @@ func requireSha(sha string) error { type PipelineTrigger interface { TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (models.PipelineId, error) DescribeWorkflowDefinition(ctx context.Context, repoDid syntax.DID, sha string, sourceRepo syntax.DID) (*tangled.CiDescribeWorkflowDefinition_Output, error) + ListWorkflowDefinitions(ctx context.Context, repoDid syntax.DID, ref string) ([]*models.WorkflowDefinition, error) } type RepoWiper interface { @@ -80,9 +77,19 @@ type Xrpc struct { func (x *Xrpc) Router() http.Handler { r := chi.NewRouter() + r.Get("/"+org_tangled.CiGetWorkflowDefinitionNSID, x.handleOrgTangledCiGetWorkflowDefinition) + r.Get("/"+org_tangled.CiListWorkflowDefinitionsNSID, x.handleOrgTangledCiListWorkflowDefinitions) + r.Get("/"+org_tangled.CiSubscribePipelineLogsNSID, x.handleOrgTangledCiSubscribePipelineLogs) + r.Get("/"+org_tangled.CiQueryPipelinesNSID, x.handleOrgTangledCiQueryPipelines) + r.Get("/"+org_tangled.CiGetPipelineNSID, x.handleOrgTangledCiGetPipeline) + + // authenticated endpoints r.Group(func(r chi.Router) { r.Use(x.ServiceAuth.VerifyServiceAuth) + r.Post("/"+org_tangled.SecretAddSecretNSID, x.handleOrgTangledSecretAddSecret) + r.Post("/"+org_tangled.SecretRemoveSecretNSID, x.handleOrgTangledSecretRemoveSecret) + r.Get("/"+org_tangled.SecretListSecretsNSID, x.handleOrgTangledSecretListSecrets) r.Post("/"+tangled.RepoAddSecretNSID, x.AddSecret) r.Post("/"+tangled.RepoRemoveSecretNSID, x.RemoveSecret) r.Get("/"+tangled.RepoListSecretsNSID, x.ListSecrets) @@ -109,6 +116,10 @@ func (x *Xrpc) Router() http.Handler { return r } +func invalidRequest(err error) xrpcerr.XrpcError { + return xrpcerr.NewXrpcError(xrpcerr.WithTag("InvalidRequest"), xrpcerr.WithError(err)) +} + // this is slightly different from http_util::write_error to follow the spec: // // the json object returned must include an "error" and a "message" @@ -123,3 +134,9 @@ func writeJson(w http.ResponseWriter, status int, response any) error { w.WriteHeader(status) return json.NewEncoder(w).Encode(response) } + +func (x *Xrpc) writeResponse(w http.ResponseWriter, r *http.Request, status int, response any) { + if err := writeJson(w, status, response); err != nil { + x.Logger.ErrorContext(r.Context(), "failed to write response", "err", err) + } +} diff --git a/spindle/xrpc/xrpc_test.go b/spindle/xrpc/xrpc_test.go index 3dd54e828..776a4d316 100644 --- a/spindle/xrpc/xrpc_test.go +++ b/spindle/xrpc/xrpc_test.go @@ -4,7 +4,6 @@ import ( "bytes" "context" "encoding/json" - "github.com/bluesky-social/indigo/atproto/identity" "log/slog" "net/http" "net/http/httptest" @@ -13,7 +12,11 @@ import ( "testing" "time" + "github.com/bluesky-social/indigo/atproto/atcrypto" + "github.com/bluesky-social/indigo/atproto/auth" + "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/org_tangled" "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/rbac" @@ -22,10 +25,12 @@ import ( "tangled.org/core/spindle/models" pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/spindle/secrets" + "tangled.org/core/xrpc/serviceauth" ) type mockTrigger struct { - triggered bool + triggered bool + workflowDefinitions []*models.WorkflowDefinition } func (m *mockTrigger) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (models.PipelineId, error) { @@ -37,6 +42,10 @@ func (m *mockTrigger) DescribeWorkflowDefinition(context.Context, syntax.DID, st return &tangled.CiDescribeWorkflowDefinition_Output{}, nil } +func (m *mockTrigger) ListWorkflowDefinitions(context.Context, syntax.DID, string) ([]*models.WorkflowDefinition, error) { + return m.workflowDefinitions, nil +} + func createTestPipeline(t *testing.T, database *db.DB, id models.PipelineId, raw tangled.Pipeline) { t.Helper() record, err := pipelinecodec.FromTangled(id, time.Now(), raw) @@ -407,4 +416,101 @@ func TestSecrets_RBAC(t *testing.T) { if len(listOut.Secrets) != 0 { t.Fatal("secret was not removed") } + + orgAdd := org_tangled.SecretAddSecret_Input{Repo: repoDid.String(), Key: "ORG_SECRET", Value: "value"} + body, _ := json.Marshal(orgAdd) + req := httptest.NewRequest(http.MethodPost, "/"+org_tangled.SecretAddSecretNSID, bytes.NewReader(body)) + req = req.WithContext(context.WithValue(req.Context(), ActorDid, pusherDid)) + w = httptest.NewRecorder() + x.handleOrgTangledSecretAddSecret(w, req) + if w.Code != http.StatusOK { + t.Fatalf("org add secret = %d: %s", w.Code, w.Body.String()) + } + + req = httptest.NewRequest(http.MethodGet, "/"+org_tangled.SecretListSecretsNSID+"?repo="+repoDid.String(), nil) + req = req.WithContext(context.WithValue(req.Context(), ActorDid, pusherDid)) + w = httptest.NewRecorder() + x.handleOrgTangledSecretListSecrets(w, req) + var orgList org_tangled.SecretListSecrets_Output + if w.Code != http.StatusOK || json.Unmarshal(w.Body.Bytes(), &orgList) != nil || len(orgList.Secrets) != 1 || orgList.Secrets[0].Key != "ORG_SECRET" { + t.Fatalf("org list secrets = %d, %+v: %s", w.Code, orgList, w.Body.String()) + } + + body, _ = json.Marshal(orgAdd) + req = httptest.NewRequest(http.MethodPost, "/"+org_tangled.SecretAddSecretNSID, bytes.NewReader(body)) + req = req.WithContext(context.WithValue(req.Context(), ActorDid, nonPusherDid)) + w = httptest.NewRecorder() + x.handleOrgTangledSecretAddSecret(w, req) + if w.Code != http.StatusUnauthorized { + t.Fatalf("unauthorized org add secret = %d", w.Code) + } + + orgRemove := org_tangled.SecretRemoveSecret_Input{Repo: repoDid.String(), Key: "ORG_SECRET"} + body, _ = json.Marshal(orgRemove) + req = httptest.NewRequest(http.MethodPost, "/"+org_tangled.SecretRemoveSecretNSID, bytes.NewReader(body)) + req = req.WithContext(context.WithValue(req.Context(), ActorDid, pusherDid)) + w = httptest.NewRecorder() + x.handleOrgTangledSecretRemoveSecret(w, req) + if w.Code != http.StatusOK { + t.Fatalf("org remove secret = %d: %s", w.Code, w.Body.String()) + } + + privateKey, err := atcrypto.GeneratePrivateKeyP256() + if err != nil { + t.Fatal(err) + } + publicKey, err := privateKey.PublicKey() + if err != nil { + t.Fatal(err) + } + serviceIdentity := &identity.Identity{ + DID: pusherDid, + Keys: map[string]identity.VerificationMethod{ + "atproto": {Type: "Multikey", PublicKeyMultibase: publicKey.Multibase()}, + }, + } + x.ServiceAuth = serviceauth.NewServiceAuth(slog.Default(), idresolver.MockDirectory{Ident: serviceIdentity}, "did:web:spindle.test") + signedToken := func(nsid string) string { + t.Helper() + lxm := syntax.NSID(nsid) + token, err := auth.SignServiceAuth(pusherDid, "did:web:spindle.test", time.Minute, &lxm, privateKey) + if err != nil { + t.Fatal(err) + } + return token + } + routerRequest := func(method, nsid string, body []byte, query string, token string) *httptest.ResponseRecorder { + t.Helper() + req := httptest.NewRequest(method, "/"+nsid+query, bytes.NewReader(body)) + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + w := httptest.NewRecorder() + x.Router().ServeHTTP(w, req) + return w + } + + w = routerRequest(http.MethodGet, org_tangled.SecretListSecretsNSID, nil, "?repo="+repoDid.String(), "") + if w.Code != http.StatusForbidden { + t.Fatalf("org secret route without service auth = %d, want 403", w.Code) + } + + body, _ = json.Marshal(orgAdd) + w = routerRequest(http.MethodPost, org_tangled.SecretAddSecretNSID, body, "", signedToken(org_tangled.SecretAddSecretNSID)) + if w.Code != http.StatusOK { + t.Fatalf("service-auth org add = %d: %s", w.Code, w.Body.String()) + } + w = routerRequest(http.MethodGet, org_tangled.SecretListSecretsNSID, nil, "?repo="+repoDid.String(), signedToken(org_tangled.SecretListSecretsNSID)) + if w.Code != http.StatusOK || json.Unmarshal(w.Body.Bytes(), &orgList) != nil || len(orgList.Secrets) != 1 || orgList.Secrets[0].CreatedBy != pusherDid.String() { + t.Fatalf("service-auth org list = %d, %+v: %s", w.Code, orgList, w.Body.String()) + } + w = routerRequest(http.MethodGet, org_tangled.SecretListSecretsNSID, nil, "?repo="+repoDid.String(), signedToken(org_tangled.SecretAddSecretNSID)) + if w.Code != http.StatusForbidden { + t.Fatalf("mismatched service-auth lxm = %d, want 403", w.Code) + } + body, _ = json.Marshal(orgRemove) + w = routerRequest(http.MethodPost, org_tangled.SecretRemoveSecretNSID, body, "", signedToken(org_tangled.SecretRemoveSecretNSID)) + if w.Code != http.StatusOK { + t.Fatalf("service-auth org remove = %d: %s", w.Code, w.Body.String()) + } } diff --git a/workflow/def.go b/workflow/def.go index 14467b187..be6ace864 100644 --- a/workflow/def.go +++ b/workflow/def.go @@ -115,6 +115,7 @@ func FromFile(name string, contents []byte) (Workflow, error) { } // if any of the constraints on a workflow is true, return true + func (w *Workflow) Match(trigger tangled.Pipeline_TriggerMetadata, changedFiles []string) (bool, error) { // manual dispatch skips matching constraints since selection is done by the caller if trigger.Manual != nil { -- 2.51.2