diff --git a/spindle/db/db.go b/spindle/db/db.go index 1c84da1f6..475056ed4 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -14,6 +14,7 @@ import ( "tangled.org/core/log" "tangled.org/core/orm" "tangled.org/core/spindle/models" + pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/sqlite" ) @@ -1139,9 +1140,8 @@ func (d *DB) GetRepoOwnerAndDid(knot, rkey string) (owner string, repoDid string return } -// migratePipelines converts the legacy sh.tangled.pipeline event rows into -// sh.tangled.ci.pipeline payloads. ordered by created so the new rowids follow -// historical order, which is what the pagination cursor reads. +// migratePipelines orders converted rows by their original creation time so +// pagination keeps the historical order. func migratePipelines(tx *sql.Tx, logger *slog.Logger) error { rows, err := tx.Query( `select rkey, event, created from events where nsid = ? order by created asc`, @@ -1176,8 +1176,12 @@ func migratePipelines(tx *sql.Tx, logger *slog.Logger) error { continue } - p, kind := mapToCiPipeline(models.PipelineId(rkey), time.Unix(0, created), raw) - payload, err := json.Marshal(p) + record, err := pipelinecodec.FromTangled(models.PipelineId(rkey), time.Unix(0, created), raw) + if err != nil { + skipped++ + continue + } + payload, err := json.Marshal(record) if err != nil { skipped++ continue @@ -1186,7 +1190,7 @@ func migratePipelines(tx *sql.Tx, logger *slog.Logger) error { if raw.TriggerMetadata.Repo != nil { knot = raw.TriggerMetadata.Repo.Knot } - out = append(out, converted{rkey, knot, p.Repo, p.Commit, string(kind), payload}) + out = append(out, converted{rkey, knot, record.RepoDID, record.Commit, record.Trigger.Kind, payload}) } if err := rows.Err(); err != nil { return err @@ -1249,16 +1253,28 @@ func migrateWorkflowStatuses(tx *sql.Tx, logger *slog.Logger) error { return err } + convertedCount := 0 for _, c := range out { - if _, err := tx.Exec( + result, err := tx.Exec( `insert into workflow_statuses (rkey, workflow, status, error, exit_code, created_at) - values (?, ?, ?, ?, ?, ?)`, - c.rkey, c.workflow, c.status, c.wfError, c.exitCode, c.createdAt, - ); err != nil { + select ?, ?, ?, ?, ?, ? + where exists (select 1 from pipelines where rkey = ?)`, + c.rkey, c.workflow, c.status, c.wfError, c.exitCode, c.createdAt, c.rkey, + ) + if err != nil { + return err + } + inserted, err := result.RowsAffected() + if err != nil { return err } + if inserted == 0 { + skipped++ + continue + } + convertedCount++ } - logger.Info("backfilled workflow statuses", "converted", len(out), "skipped", skipped) + logger.Info("backfilled workflow statuses", "converted", convertedCount, "skipped", skipped) return nil } diff --git a/spindle/db/events.go b/spindle/db/events.go index f8860cd44..d6bb93467 100644 --- a/spindle/db/events.go +++ b/spindle/db/events.go @@ -6,7 +6,6 @@ import ( "encoding/json" "time" - "tangled.org/core/api/tangled" "tangled.org/core/notifier" "tangled.org/core/spindle/models" ) @@ -219,7 +218,7 @@ func (d *DB) ListPipelineWorkflows(repoDid string) ([]PipelineWorkflow, error) { if err := rows.Scan(&pipelineID, &raw); err != nil { return nil, err } - var p tangled.CiPipeline + var p models.PipelineRecord if err := json.Unmarshal([]byte(raw), &p); err != nil { continue } diff --git a/spindle/db/identity_migration_test.go b/spindle/db/identity_migration_test.go new file mode 100644 index 000000000..05d87f1f4 --- /dev/null +++ b/spindle/db/identity_migration_test.go @@ -0,0 +1,270 @@ +package db + +import ( + "context" + "database/sql" + "encoding/json" + "os" + "path/filepath" + "slices" + "testing" + + "tangled.org/core/api/tangled" + "tangled.org/core/spindle/models" +) + +func TestPipelineIdentityFinalSchemaHasNoLegacyState(t *testing.T) { + ctx := context.Background() + database, err := Make(ctx, filepath.Join(t.TempDir(), "final-schema.db")) + if err != nil { + t.Fatal(err) + } + defer database.Close() + if err := database.MigratePipelineLogFiles(""); err != nil { + t.Fatal(err) + } + + for _, table := range []string{"events", "events_legacy", "pipeline_log_renames"} { + var count int + if err := database.QueryRow(`select count(*) from sqlite_master where type = 'table' and name = ?`, table).Scan(&count); err != nil || count != 0 { + t.Fatalf("legacy table %s count = %d, %v", table, count, err) + } + } + legacyColumns := map[string][]string{ + "pipelines": {"knot", "rkey", "definitions"}, + "workflow_statuses": {"rkey"}, + "jobs": {"pipeline_id_knot", "pipeline_id_rkey"}, + "mill_leases": {"knot", "rkey"}, + "mill_artifacts": {"knot", "rkey"}, + "executor_pending_artifacts": {"knot", "rkey"}, + } + for table, columns := range legacyColumns { + for _, column := range columns { + var count int + if err := database.QueryRow(`select count(*) from pragma_table_info(?) where name = ?`, table, column).Scan(&count); err != nil || count != 0 { + t.Fatalf("legacy column %s.%s count = %d, %v", table, column, count, err) + } + } + } +} + +func TestPipelineMigrationSkipsMalformedLegacyRows(t *testing.T) { + ctx := context.Background() + dbPath := filepath.Join(t.TempDir(), "malformed.db") + database, err := Make(ctx, dbPath) + if err != nil { + t.Fatal(err) + } + if err := database.Close(); err != nil { + t.Fatal(err) + } + + rawDB, err := sql.Open("sqlite3", dbPath) + if err != nil { + t.Fatal(err) + } + if _, err := rawDB.Exec(` + delete from migrations where name in ('pipelines-and-workflow-statuses', 'pipeline-id-identity', 'drop-legacy-pipeline-events'); + + alter table jobs rename column pipeline_id to pipeline_id_rkey; + alter table jobs add column pipeline_id_knot text not null default ''; + + alter table mill_leases rename column pipeline_id to rkey; + alter table mill_leases add column knot text not null default ''; + + drop index idx_mill_artifacts_workflow_identity; + alter table mill_artifacts rename column pipeline_id to rkey; + alter table mill_artifacts add column knot text not null default ''; + create index idx_mill_artifacts_workflow_identity on mill_artifacts(knot, rkey, workflow, id desc); + + alter table executor_pending_artifacts rename column pipeline_id to rkey; + alter table executor_pending_artifacts add column knot text not null default ''; + + drop table pipeline_log_renames; + drop table if exists events_legacy; + drop table pipelines; + drop table workflow_statuses; + create table events ( + rkey text not null, + nsid text not null, + event text not null, + created integer not null + ); + insert into events (rkey, nsid, event, created) values ('3mu2xwiorc2xl', 'sh.tangled.pipeline', '{', 1); + `); err != nil { + rawDB.Close() + t.Fatal(err) + } + repo := "did:plc:validrepo" + validPipeline := 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: push\n"}}, + } + validJSON, err := json.Marshal(validPipeline) + if err != nil { + t.Fatal(err) + } + if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "3mu2xwiorc2xm", tangled.PipelineNSID, validJSON, int64(2)); err != nil { + t.Fatal(err) + } + statusJSON, err := json.Marshal(tangled.PipelineStatus{ + Pipeline: "at://did:plc:validrepo/sh.tangled.pipeline/3mu2xwiorc2xm", + Workflow: "ci.yml", Status: "success", CreatedAt: "2026-08-31T08:00:00Z", + }) + if err != nil { + t.Fatal(err) + } + if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "status-valid", tangled.PipelineStatusNSID, statusJSON, int64(3)); err != nil { + t.Fatal(err) + } + if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "status-broken", tangled.PipelineStatusNSID, "{", int64(4)); err != nil { + t.Fatal(err) + } + invalidDefinitionPipeline := validPipeline + invalidDefinitionPipeline.Workflows = []*tangled.Pipeline_Workflow{ + {Name: "valid.yml", Raw: "when:\n - event: push\n"}, + {Name: "broken.yml", Raw: "when: ["}, + } + invalidDefinitionJSON, err := json.Marshal(invalidDefinitionPipeline) + if err != nil { + t.Fatal(err) + } + if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "3mu2xwiorc2xn", tangled.PipelineNSID, invalidDefinitionJSON, int64(5)); err != nil { + t.Fatal(err) + } + + orphanStatusJSON, err := json.Marshal(tangled.PipelineStatus{ + Pipeline: "at://did:plc:validrepo/sh.tangled.pipeline/3mu2xwiorc2xn", + Workflow: "broken.yml", Status: "failed", CreatedAt: "2026-08-31T08:01:00Z", + }) + if err != nil { + t.Fatal(err) + } + if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "status-orphan", tangled.PipelineStatusNSID, orphanStatusJSON, int64(6)); err != nil { + t.Fatal(err) + } + if err := rawDB.Close(); err != nil { + t.Fatal(err) + } + logDir := filepath.Join(t.TempDir(), "logs") + if err := os.MkdirAll(logDir, 0o755); err != nil { + t.Fatal(err) + } + pipelineID := models.PipelineId("3mu2xwiorc2xm") + oldLog := models.LegacyLogFilePath(logDir, "knot.test", pipelineID, "ci.yml") + if err := os.WriteFile(oldLog, []byte("historical log"), 0o600); err != nil { + t.Fatal(err) + } + + upgraded, err := Make(ctx, dbPath) + if err != nil { + t.Fatal(err) + } + defer upgraded.Close() + if err := upgraded.MigratePipelineLogFiles(logDir); err != nil { + t.Fatal(err) + } + var archivedEventsTable int + if err := upgraded.QueryRow(`select count(*) from sqlite_master where type = 'table' and name = 'events_legacy'`).Scan(&archivedEventsTable); err != nil || archivedEventsTable != 0 { + t.Fatalf("events_legacy table count = %d, %v", archivedEventsTable, err) + } + var converted int + if err := upgraded.QueryRow(`select count(*) from pipelines`).Scan(&converted); err != nil { + t.Fatal(err) + } + if converted != 1 { + t.Fatalf("valid legacy pipelines converted = %d", converted) + } + + var statuses int + if err := upgraded.QueryRow(`select count(*) from workflow_statuses`).Scan(&statuses); err != nil || statuses != 1 { + t.Fatalf("workflow statuses = %d, %v", statuses, err) + } + record, err := upgraded.GetPipeline(ctx, pipelineID) + if err != nil { + t.Fatal(err) + } + if len(record.Workflows) != 1 || record.Workflows[0].Definition == nil { + t.Fatalf("canonical workflow = %#v", record.Workflows) + } + workflow := record.Workflows[0] + definition := workflow.Definition + if definition.Hash != nil || !slices.Equal(definition.Triggers, []string{"org.tangled.event.push"}) { + t.Fatalf("canonical definition = %#v", definition) + } + file := definition.Source.File + if file == nil || file.Repo != repo || file.Commit != record.Commit || file.Path != ".tangled/workflows/ci.yml" { + t.Fatalf("canonical source = %#v", file) + } + if record.ID != pipelineID || record.Trigger.Push == nil || record.Trigger.Push.Ref != "refs/heads/main" || workflow.Status != "success" || workflow.FinishedAt == nil { + t.Fatalf("canonical run = %#v", record) + } + newLog := models.LogFilePath(logDir, models.WorkflowId{PipelineId: pipelineID, Name: "ci.yml"}) + if contents, err := os.ReadFile(newLog); err != nil || string(contents) != "historical log" { + t.Fatalf("migrated log = %q, %v", contents, err) + } + if _, err := upgraded.GetPipeline(ctx, models.PipelineId("3mu2xwiorc2xn")); err != sql.ErrNoRows { + t.Fatalf("pipeline with malformed definition survived migration: %v", err) + } +} + +func TestMigratePipelineLogFilesDoesNotClobberAndResumesHardLink(t *testing.T) { + ctx := context.Background() + database, err := Make(ctx, filepath.Join(t.TempDir(), "logs.db")) + if err != nil { + t.Fatal(err) + } + defer database.Close() + pipelineID := models.PipelineId("3mu2xwiorc2xl") + if _, err := database.Exec(`insert into pipeline_log_renames (knot, pipeline_id, workflow) values (?, ?, ?)`, "knot.test", pipelineID, "ci.yml"); err != nil { + t.Fatal(err) + } + logDir := filepath.Join(t.TempDir(), "logs") + if err := os.MkdirAll(logDir, 0o755); err != nil { + t.Fatal(err) + } + oldPath := models.LegacyLogFilePath(logDir, "knot.test", pipelineID, "ci.yml") + newPath := models.LogFilePath(logDir, models.WorkflowId{PipelineId: pipelineID, Name: "ci.yml"}) + if err := os.WriteFile(oldPath, []byte("old"), 0o600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(newPath, []byte("new"), 0o600); err != nil { + t.Fatal(err) + } + if err := database.MigratePipelineLogFiles(logDir); err == nil { + t.Fatal("expected existing destination to reject migration") + } + if oldData, _ := os.ReadFile(oldPath); string(oldData) != "old" { + t.Fatalf("old path changed: %q", oldData) + } + if newData, _ := os.ReadFile(newPath); string(newData) != "new" { + t.Fatalf("new path changed: %q", newData) + } + + if err := os.Remove(newPath); err != nil { + t.Fatal(err) + } + if err := os.Link(oldPath, newPath); err != nil { + t.Fatal(err) + } + if err := database.MigratePipelineLogFiles(logDir); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(oldPath); !os.IsNotExist(err) { + t.Fatalf("old hard link still exists: %v", err) + } + if data, err := os.ReadFile(newPath); err != nil || string(data) != "old" { + t.Fatalf("new hard link = %q, %v", data, err) + } + var renameTable int + if err := database.QueryRow(`select count(*) from sqlite_master where type = 'table' and name = 'pipeline_log_renames'`).Scan(&renameTable); err != nil || renameTable != 0 { + t.Fatalf("pipeline_log_renames table count = %d, %v", renameTable, err) + } + if err := database.MigratePipelineLogFiles(logDir); err != nil { + t.Fatalf("repeat log migration: %v", err) + } +} diff --git a/spindle/db/log_migration.go b/spindle/db/log_migration.go index 992749524..ff16e7729 100644 --- a/spindle/db/log_migration.go +++ b/spindle/db/log_migration.go @@ -6,7 +6,6 @@ import ( "fmt" "os" - "tangled.org/core/api/tangled" "tangled.org/core/spindle/models" ) @@ -33,7 +32,7 @@ func stagePipelineLogRenames(tx *sql.Tx) error { rows.Close() return err } - var pipeline tangled.CiPipeline + var pipeline models.PipelineRecord if err := json.Unmarshal([]byte(payload), &pipeline); err != nil { rows.Close() return fmt.Errorf("decode pipeline %s while staging log migration: %w", rkey, err) diff --git a/spindle/db/pipelines.go b/spindle/db/pipelines.go index cfe921adc..509dcafff 100644 --- a/spindle/db/pipelines.go +++ b/spindle/db/pipelines.go @@ -5,21 +5,17 @@ import ( "encoding/json" "strconv" "strings" - "time" - "tangled.org/core/api/tangled" "tangled.org/core/orm" "tangled.org/core/spindle/models" - "tangled.org/core/workflow" ) -func (d *DB) QueryPipelines(ctx context.Context, repoDid string, commits []string, cursor string, kinds []string, limit int) ([]*tangled.CiPipeline, string, int64, error) { +func (d *DB) QueryPipelines(ctx context.Context, repoDID string, commits []string, cursor string, kinds []string, limit int) ([]*models.PipelineRecord, string, int64, error) { if limit <= 0 { limit = 30 } - filters := []orm.Filter{orm.FilterEq("repo_did", repoDid)} - // only filter when asked: FilterIn compiles an empty slice to `1 = 0` + filters := []orm.Filter{orm.FilterEq("repo_did", repoDID)} if len(commits) > 0 { filters = append(filters, orm.FilterIn("commit_sha", commits)) } @@ -41,10 +37,9 @@ func (d *DB) QueryPipelines(ctx context.Context, repoDid string, commits []strin return nil, "", 0, err } - // the cursor bounds the page but not the total, so it joins after the count if cursor != "" { - if cVal, err := strconv.ParseInt(cursor, 10, 64); err == nil { - filter := orm.FilterLt("id", cVal) + if value, err := strconv.ParseInt(cursor, 10, 64); err == nil { + filter := orm.FilterLt("id", value) conditions = append(conditions, filter.Condition()) args = append(args, filter.Arg()...) whereClause = " where " + strings.Join(conditions, " and ") @@ -59,39 +54,38 @@ func (d *DB) QueryPipelines(ctx context.Context, repoDid string, commits []strin } defer rows.Close() - var pipelines []*tangled.CiPipeline - var lastId int64 - + records := make([]*models.PipelineRecord, 0, limit) + var lastID int64 + var scanned int for rows.Next() { var id int64 var payload string if err := rows.Scan(&id, &payload); err != nil { return nil, "", 0, err } - var p tangled.CiPipeline - if err := json.Unmarshal([]byte(payload), &p); err != nil { + lastID = id + scanned++ + var record models.PipelineRecord + if json.Unmarshal([]byte(payload), &record) != nil { continue } - lastId = id - pipelines = append(pipelines, &p) + records = append(records, &record) } if err := rows.Err(); err != nil { return nil, "", 0, err } - - if err := d.applyStatuses(ctx, pipelines); err != nil { + if err := d.applyStatuses(ctx, records); err != nil { return nil, "", 0, err } nextCursor := "" - if len(pipelines) == limit { - nextCursor = strconv.FormatInt(lastId, 10) + if scanned == limit { + nextCursor = strconv.FormatInt(lastID, 10) } - - return pipelines, nextCursor, total, nil + return records, nextCursor, total, nil } -func (d *DB) GetPipeline(ctx context.Context, id models.PipelineId) (*tangled.CiPipeline, error) { +func (d *DB) GetPipeline(ctx context.Context, id models.PipelineId) (*models.PipelineRecord, error) { var payload string if err := d.QueryRowContext(ctx, `select payload from pipelines where pipeline_id = ?`, id, @@ -99,160 +93,56 @@ func (d *DB) GetPipeline(ctx context.Context, id models.PipelineId) (*tangled.Ci return nil, err } - var p tangled.CiPipeline - if err := json.Unmarshal([]byte(payload), &p); err != nil { + var record models.PipelineRecord + if err := json.Unmarshal([]byte(payload), &record); err != nil { return nil, err } - - if err := d.applyStatuses(ctx, []*tangled.CiPipeline{&p}); err != nil { + if err := d.applyStatuses(ctx, []*models.PipelineRecord{&record}); err != nil { return nil, err } - return &p, nil -} - -// mapToCiPipeline stores the initial status; reads overlay live workflow_statuses rows. -func mapToCiPipeline(id models.PipelineId, createdAt time.Time, raw tangled.Pipeline) (*tangled.CiPipeline, workflow.TriggerKind) { - createdAtStr := createdAt.Format(time.RFC3339) - metadata := raw.TriggerMetadata - - var repoDidStr string - if metadata != nil && metadata.Repo != nil { - if rd := metadata.Repo.RepoDid; rd != nil && *rd != "" { - repoDidStr = *rd - } else { - repoDidStr = metadata.Repo.Did - } - } - - commitSha := "" - var trigger tangled.CiPipeline_Trigger - var kind workflow.TriggerKind - - if metadata != nil { - kind = workflow.TriggerKind(metadata.Kind) - switch kind { - case workflow.TriggerKindPush: - if metadata.Push != nil { - commitSha = metadata.Push.NewSha - trigger.CiTrigger_Push = &tangled.CiTrigger_Push{ - NewSha: metadata.Push.NewSha, - OldSha: metadata.Push.OldSha, - Ref: metadata.Push.Ref, - } - } - case workflow.TriggerKindPullRequest: - if metadata.PullRequest != nil { - commitSha = metadata.PullRequest.SourceSha - trigger.CiTrigger_PullRequest = &tangled.CiTrigger_PullRequest{ - Action: metadata.PullRequest.Action, - SourceBranch: &metadata.PullRequest.SourceBranch, - SourceRepo: metadata.SourceRepo, - SourceSha: metadata.PullRequest.SourceSha, - TargetBranch: metadata.PullRequest.TargetBranch, - Pull: metadata.PullRequest.Pull, - } - } - case workflow.TriggerKindManual: - if metadata.Manual != nil { - commitSha = metadata.Manual.Sha - trigger.CiTrigger_Manual = &tangled.CiTrigger_Manual{ - Inputs: pipelinePairsToCiTriggerPairs(metadata.Manual.Inputs), - Ref: metadata.Manual.Ref, - Sha: metadata.Manual.Sha, - SourceRepo: metadata.SourceRepo, - } - } - } - } - - workflows := make([]*tangled.CiPipeline_Workflow, 0, len(raw.Workflows)) - for _, wf := range raw.Workflows { - if wf == nil { - continue - } - - workflows = append(workflows, &tangled.CiPipeline_Workflow{ - Id: wf.Name, - Name: wf.Name, - Status: string(models.StatusKindPending), - }) - } - - var sourceRepo *string - if metadata != nil { - sourceRepo = metadata.SourceRepo - } - - return &tangled.CiPipeline{ - Id: string(id), - Commit: commitSha, - Repo: repoDidStr, - CreatedAt: &createdAtStr, - Trigger: &trigger, - Workflows: workflows, - SourceRepo: sourceRepo, - }, kind -} - -func pipelinePairsToCiTriggerPairs(inputs []*tangled.Pipeline_Pair) []*tangled.CiTrigger_Pair { - if len(inputs) == 0 { - return nil - } - pairs := make([]*tangled.CiTrigger_Pair, 0, len(inputs)) - for _, input := range inputs { - if input == nil { - continue - } - pairs = append(pairs, &tangled.CiTrigger_Pair{ - Key: input.Key, - Value: input.Value, - }) - } - return pairs + return &record, nil } -func (d *DB) CreatePipeline(id models.PipelineId, raw tangled.Pipeline) error { - p, kind := mapToCiPipeline(id, time.Now(), raw) - payload, err := json.Marshal(p) +func (d *DB) CreatePipeline(record *models.PipelineRecord) error { + payload, err := json.Marshal(record) if err != nil { return err } _, err = d.Exec( `insert into pipelines (pipeline_id, repo_did, commit_sha, kind, payload) values (?, ?, ?, ?, ?)`, - id, p.Repo, p.Commit, string(kind), string(payload), + record.ID, record.RepoDID, record.Commit, record.Trigger.Kind, string(payload), ) return err } -// applyStatuses overlays the latest workflow statuses onto stored pipelines. -func (d *DB) applyStatuses(ctx context.Context, pipelines []*tangled.CiPipeline) error { +func (d *DB) applyStatuses(ctx context.Context, pipelines []*models.PipelineRecord) error { if len(pipelines) == 0 { return nil } - rkeys := make([]string, 0, len(pipelines)) - for _, p := range pipelines { - rkeys = append(rkeys, p.Id) + pipelineIDs := make([]string, 0, len(pipelines)) + for _, pipeline := range pipelines { + pipelineIDs = append(pipelineIDs, string(pipeline.ID)) } - statuses, err := d.workflowStatuses(ctx, rkeys) + statuses, err := d.workflowStatuses(ctx, pipelineIDs) if err != nil { return err } - for _, p := range pipelines { - for _, wf := range p.Workflows { - if wf == nil { + for _, pipeline := range pipelines { + for _, workflow := range pipeline.Workflows { + if workflow == nil { continue } - st, ok := statuses[wfKey{p.Id, wf.Name}] + status, ok := statuses[wfKey{string(pipeline.ID), workflow.Name}] if !ok { continue } - wf.Status = st.Status - wf.Error = st.Error - wf.StartedAt = st.StartedAt - wf.FinishedAt = st.FinishedAt + workflow.Status = status.Status + workflow.Error = status.Error + workflow.StartedAt = status.StartedAt + workflow.FinishedAt = status.FinishedAt } } return nil diff --git a/spindle/db/pipelines_test.go b/spindle/db/pipelines_test.go index c18688532..f25b489f0 100644 --- a/spindle/db/pipelines_test.go +++ b/spindle/db/pipelines_test.go @@ -10,6 +10,7 @@ import ( "tangled.org/core/api/tangled" "tangled.org/core/notifier" "tangled.org/core/spindle/models" + pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/workflow" ) @@ -32,7 +33,11 @@ func seedPipeline(t *testing.T, d *DB, rkey, repoDid, kind string) { TriggerMetadata: tm, Workflows: []*tangled.Pipeline_Workflow{{Name: "ci.yml"}}, } - if err := d.CreatePipeline(models.PipelineId(rkey), raw); err != nil { + record, err := pipelinecodec.FromTangled(models.PipelineId(rkey), time.Now(), raw) + if err != nil { + t.Fatalf("map pipeline %s: %v", rkey, err) + } + if err := d.CreatePipeline(record); err != nil { t.Fatalf("seed pipeline %s: %v", rkey, err) } } @@ -91,25 +96,13 @@ func TestQueryPipelines_KindScopedToRepo(t *testing.T) { if total != 1 || len(pipelines) != 1 { t.Fatalf("total=%d len=%d, want exactly alice's single push pipeline", total, len(pipelines)) } - if pipelines[0].Repo != "did:plc:alice" { - t.Errorf("returned pipeline repo = %q, want did:plc:alice", pipelines[0].Repo) - + if pipelines[0].RepoDID != "did:plc:alice" { + t.Errorf("returned pipeline repo = %q, want did:plc:alice", pipelines[0].RepoDID) } } -func triggerKindOf(p *tangled.CiPipeline) string { - if p.Trigger == nil { - return "" - } - switch { - case p.Trigger.CiTrigger_Push != nil: - return "push" - case p.Trigger.CiTrigger_PullRequest != nil: - return "pull_request" - case p.Trigger.CiTrigger_Manual != nil: - return "manual" - } - return "" +func triggerKindOf(p *models.PipelineRecord) string { + return p.Trigger.Kind } func TestQueryPipelines_WorkflowStatuses(t *testing.T) { @@ -126,7 +119,11 @@ func TestQueryPipelines_WorkflowStatuses(t *testing.T) { }, Workflows: []*tangled.Pipeline_Workflow{{Name: "a"}, {Name: "b"}}, } - if err := d.CreatePipeline(models.PipelineId("pl1"), raw); err != nil { + record, err := pipelinecodec.FromTangled(models.PipelineId("pl1"), time.Now(), raw) + if err != nil { + t.Fatalf("map pipeline: %v", err) + } + if err := d.CreatePipeline(record); err != nil { t.Fatalf("CreatePipeline: %v", err) } @@ -152,7 +149,7 @@ func TestQueryPipelines_WorkflowStatuses(t *testing.T) { t.Fatalf("total = %d, len = %d; want 1, 1", total, len(pipelines)) } - byName := map[string]*tangled.CiPipeline_Workflow{} + byName := map[string]*models.PipelineWorkflow{} for _, wf := range pipelines[0].Workflows { byName[wf.Name] = wf } @@ -232,27 +229,56 @@ func TestToCiPipeline_TriggerKinds(t *testing.T) { } tt.meta(meta) - p, kind := mapToCiPipeline("rk1", time.Now(), tangled.Pipeline{ + record, err := pipelinecodec.FromTangled("rk1", time.Now(), tangled.Pipeline{ TriggerMetadata: meta, Workflows: []*tangled.Pipeline_Workflow{{Name: "ci.yml"}}, }) - - if kind != tt.kind { - t.Errorf("kind = %q, want %q", kind, tt.kind) + if err != nil { + t.Fatalf("map pipeline: %v", err) } - if p.Repo != repo { - t.Errorf("p.Repo = %q, want %q", p.Repo, repo) + if record.Trigger.Kind != string(tt.kind) { + t.Errorf("kind = %q, want %q", record.Trigger.Kind, tt.kind) } - if p.Commit != sha { - t.Errorf("p.Commit = %q, want %q", p.Commit, sha) + if record.RepoDID != repo { + t.Errorf("record.RepoDID = %q, want %q", record.RepoDID, repo) } - if !tt.wantUnion(p.Trigger) { + if record.Commit != sha { + t.Errorf("record.Commit = %q, want %q", record.Commit, sha) + } + wire := pipelinecodec.ToTangled(record) + if !tt.wantUnion(wire.Trigger) { t.Errorf("wrong trigger union populated for kind %q", tt.kind) } - // the payload must be marshalable, CreatePipeline stores it as JSON - if _, err := json.Marshal(p); err != nil { + if _, err := json.Marshal(record); err != nil { t.Errorf("marshal payload: %v", err) } }) } } + +func TestQueryPipelinesSkipsMalformedRowsAndAdvancesCursor(t *testing.T) { + d := newTestDB(t) + ctx := context.Background() + repo := "did:plc:malformed-page" + seedPipeline(t, d, "oldest", repo, "push") + seedPipeline(t, d, "broken", repo, "push") + seedPipeline(t, d, "newest", repo, "push") + if _, err := d.Exec(`update pipelines set payload = '{' where pipeline_id = 'broken'`); err != nil { + t.Fatal(err) + } + + first, cursor, total, err := d.QueryPipelines(ctx, repo, nil, "", nil, 2) + if err != nil { + t.Fatal(err) + } + if total != 3 || len(first) != 1 || first[0].ID != "newest" || cursor == "" { + t.Fatalf("first page = (%+v, %q, %d)", first, cursor, total) + } + second, _, _, err := d.QueryPipelines(ctx, repo, nil, cursor, nil, 2) + if err != nil { + t.Fatal(err) + } + if len(second) != 1 || second[0].ID != "oldest" { + t.Fatalf("second page = %+v", second) + } +} diff --git a/spindle/models/pipeline_record.go b/spindle/models/pipeline_record.go new file mode 100644 index 000000000..f0ba2c7b0 --- /dev/null +++ b/spindle/models/pipeline_record.go @@ -0,0 +1,84 @@ +package models + +type PipelineRecord struct { + ID PipelineId `json:"id"` + RepoDID string `json:"repoDid"` + SourceRepo *string `json:"sourceRepo,omitempty"` + Commit string `json:"commit"` + CreatedAt string `json:"createdAt"` + Trigger PipelineTrigger `json:"trigger"` + Workflows []*PipelineWorkflow `json:"workflows"` +} + +type PipelineTrigger struct { + Kind string `json:"kind"` + Push *PipelinePushTrigger `json:"push,omitempty"` + PullRequest *PipelinePullRequestTrigger `json:"pullRequest,omitempty"` + Manual *PipelineManualTrigger `json:"manual,omitempty"` +} + +type PipelinePushTrigger struct { + Ref string `json:"ref"` + NewCommit string `json:"newCommit"` + OldCommit string `json:"oldCommit"` +} + +type PipelinePullRequestTrigger struct { + Action *string `json:"action,omitempty"` + SourceRepo *string `json:"sourceRepo,omitempty"` + SourceBranch *string `json:"sourceBranch,omitempty"` + SourceCommit string `json:"sourceCommit"` + TargetBranch string `json:"targetBranch"` + Pull *string `json:"pull,omitempty"` +} + +type PipelineManualTrigger struct { + Ref *string `json:"ref,omitempty"` + SourceRepo *string `json:"sourceRepo,omitempty"` + Inputs []*PipelineInputPair `json:"inputs,omitempty"` +} + +type PipelineInputPair struct { + Key string `json:"key"` + Value string `json:"value"` +} + +type PipelineWorkflow struct { + ID string `json:"id"` + Name string `json:"name"` + Definition *WorkflowDefinition `json:"definition"` + Status string `json:"status"` + Error *string `json:"error,omitempty"` + StartedAt *string `json:"startedAt,omitempty"` + FinishedAt *string `json:"finishedAt,omitempty"` +} + +type WorkflowDefinition struct { + ID string `json:"id"` + Name string `json:"name"` + Source WorkflowDefinitionSource `json:"source"` + Triggers []string `json:"triggers"` + Hash *string `json:"hash,omitempty"` +} + +type WorkflowDefinitionSource struct { + File *WorkflowFileSource `json:"file,omitempty"` + External *WorkflowExternalSource `json:"external,omitempty"` +} + +type WorkflowFileSource struct { + Repo string `json:"repo"` + Commit string `json:"commit"` + Path string `json:"path"` + Lines *WorkflowLineRange `json:"lines,omitempty"` +} + +type WorkflowLineRange struct { + Start int64 `json:"start"` + End int64 `json:"end"` +} + +type WorkflowExternalSource struct { + Name string `json:"name"` + Link string `json:"link"` +} diff --git a/spindle/pipeline/codec.go b/spindle/pipeline/codec.go new file mode 100644 index 000000000..71779dd73 --- /dev/null +++ b/spindle/pipeline/codec.go @@ -0,0 +1,200 @@ +package pipeline + +import ( + "fmt" + "time" + + "tangled.org/core/api/tangled" + "tangled.org/core/spindle/models" + "tangled.org/core/workflow" +) + +const ( + TriggerNSIDPush = "org.tangled.event.push" + TriggerNSIDPullRequest = "org.tangled.event.pullRequest" + TriggerNSIDManual = "org.tangled.ci.trigger.manual" +) + +func FromTangled(id models.PipelineId, createdAt time.Time, raw tangled.Pipeline) (*models.PipelineRecord, error) { + metadata := raw.TriggerMetadata + if metadata == nil { + return nil, fmt.Errorf("pipeline %s has no trigger metadata", id) + } + + repoDID := "" + if metadata.Repo != nil { + if metadata.Repo.RepoDid != nil && *metadata.Repo.RepoDid != "" { + repoDID = *metadata.Repo.RepoDid + } else { + repoDID = metadata.Repo.Did + } + } + + record := &models.PipelineRecord{ + ID: id, + RepoDID: repoDID, + SourceRepo: metadata.SourceRepo, + CreatedAt: createdAt.Format(time.RFC3339), + Trigger: models.PipelineTrigger{ + Kind: metadata.Kind, + }, + Workflows: make([]*models.PipelineWorkflow, 0, len(raw.Workflows)), + } + + switch workflow.TriggerKind(metadata.Kind) { + case workflow.TriggerKindPush: + if metadata.Push == nil { + return nil, fmt.Errorf("pipeline %s has no push trigger", id) + } + record.Commit = metadata.Push.NewSha + record.Trigger.Push = &models.PipelinePushTrigger{ + Ref: metadata.Push.Ref, + NewCommit: metadata.Push.NewSha, + OldCommit: metadata.Push.OldSha, + } + case workflow.TriggerKindPullRequest: + if metadata.PullRequest == nil { + return nil, fmt.Errorf("pipeline %s has no pull request trigger", id) + } + record.Commit = metadata.PullRequest.SourceSha + record.Trigger.PullRequest = &models.PipelinePullRequestTrigger{ + Action: metadata.PullRequest.Action, + SourceRepo: metadata.SourceRepo, + SourceBranch: &metadata.PullRequest.SourceBranch, + SourceCommit: metadata.PullRequest.SourceSha, + TargetBranch: metadata.PullRequest.TargetBranch, + Pull: metadata.PullRequest.Pull, + } + case workflow.TriggerKindManual: + if metadata.Manual == nil { + return nil, fmt.Errorf("pipeline %s has no manual trigger", id) + } + record.Commit = metadata.Manual.Sha + record.Trigger.Manual = &models.PipelineManualTrigger{ + Ref: metadata.Manual.Ref, + SourceRepo: metadata.SourceRepo, + Inputs: inputPairs(metadata.Manual.Inputs), + } + default: + return nil, fmt.Errorf("pipeline %s has unknown trigger kind %q", id, metadata.Kind) + } + + sourceRepo := record.RepoDID + if record.SourceRepo != nil && *record.SourceRepo != "" { + sourceRepo = *record.SourceRepo + } + for _, rawWorkflow := range raw.Workflows { + if rawWorkflow == nil { + continue + } + definition, err := FileDefinition(rawWorkflow.Name, []byte(rawWorkflow.Raw), sourceRepo, record.Commit) + if err != nil { + return nil, fmt.Errorf("pipeline %s workflow %q: %w", id, rawWorkflow.Name, err) + } + record.Workflows = append(record.Workflows, &models.PipelineWorkflow{ + ID: rawWorkflow.Name, + Name: rawWorkflow.Name, + Definition: definition, + Status: string(models.StatusKindPending), + }) + } + return record, nil +} + +func FileDefinition(name string, contents []byte, repo, commit string) (*models.WorkflowDefinition, error) { + parsed, err := workflow.FromFile(name, contents) + if err != nil { + return nil, err + } + return &models.WorkflowDefinition{ + ID: name, + Name: name, + Source: models.WorkflowDefinitionSource{File: &models.WorkflowFileSource{ + Repo: repo, + Commit: commit, + Path: workflow.WorkflowDir + "/" + name, + }}, + Triggers: declaredTriggerNSIDs(parsed), + }, nil +} + +func declaredTriggerNSIDs(definition workflow.Workflow) []string { + out := make([]string, 0) + for _, constraint := range definition.When { + for _, event := range constraint.Event { + switch event { + case string(workflow.TriggerKindPush): + event = TriggerNSIDPush + case string(workflow.TriggerKindPullRequest): + event = TriggerNSIDPullRequest + case string(workflow.TriggerKindManual): + event = TriggerNSIDManual + } + out = append(out, event) + } + } + return out +} + +func inputPairs(inputs []*tangled.Pipeline_Pair) []*models.PipelineInputPair { + out := make([]*models.PipelineInputPair, 0, len(inputs)) + for _, input := range inputs { + if input != nil { + out = append(out, &models.PipelineInputPair{Key: input.Key, Value: input.Value}) + } + } + return out +} + +func ToTangled(record *models.PipelineRecord) *tangled.CiPipeline { + if record == nil { + return nil + } + trigger := &tangled.CiPipeline_Trigger{} + switch { + case record.Trigger.Push != nil: + push := record.Trigger.Push + trigger.CiTrigger_Push = &tangled.CiTrigger_Push{ + Ref: push.Ref, NewSha: push.NewCommit, OldSha: push.OldCommit, + } + case record.Trigger.PullRequest != nil: + pullRequest := record.Trigger.PullRequest + trigger.CiTrigger_PullRequest = &tangled.CiTrigger_PullRequest{ + Action: pullRequest.Action, SourceRepo: pullRequest.SourceRepo, + SourceBranch: pullRequest.SourceBranch, SourceSha: pullRequest.SourceCommit, + TargetBranch: pullRequest.TargetBranch, Pull: pullRequest.Pull, + } + case record.Trigger.Manual != nil: + manual := record.Trigger.Manual + trigger.CiTrigger_Manual = &tangled.CiTrigger_Manual{ + Ref: manual.Ref, Sha: record.Commit, SourceRepo: manual.SourceRepo, + Inputs: tangledInputPairs(manual.Inputs), + } + } + + createdAt := record.CreatedAt + workflows := make([]*tangled.CiPipeline_Workflow, 0, len(record.Workflows)) + for _, workflow := range record.Workflows { + if workflow == nil { + continue + } + workflows = append(workflows, &tangled.CiPipeline_Workflow{ + Id: workflow.ID, Name: workflow.Name, Status: workflow.Status, + Error: workflow.Error, StartedAt: workflow.StartedAt, FinishedAt: workflow.FinishedAt, + }) + } + return &tangled.CiPipeline{ + Id: string(record.ID), Repo: record.RepoDID, SourceRepo: record.SourceRepo, + Commit: record.Commit, CreatedAt: &createdAt, Trigger: trigger, Workflows: workflows, + } +} + +func tangledInputPairs(inputs []*models.PipelineInputPair) []*tangled.CiTrigger_Pair { + out := make([]*tangled.CiTrigger_Pair, 0, len(inputs)) + for _, input := range inputs { + if input != nil { + out = append(out, &tangled.CiTrigger_Pair{Key: input.Key, Value: input.Value}) + } + } + return out +} diff --git a/spindle/pipeline/codec_test.go b/spindle/pipeline/codec_test.go new file mode 100644 index 000000000..325eb4e50 --- /dev/null +++ b/spindle/pipeline/codec_test.go @@ -0,0 +1,62 @@ +package pipeline + +import ( + "slices" + "testing" + "time" + + "tangled.org/core/api/tangled" + "tangled.org/core/spindle/models" +) + +func TestFromTangledPreservesDeclaredDefinitionMetadata(t *testing.T) { + target := "did:plc:target" + source := "did:plc:source" + raw := tangled.Pipeline{ + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ + Kind: "push", Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &target}, SourceRepo: &source, + Push: &tangled.Pipeline_PushTriggerData{NewSha: "1111111111111111111111111111111111111111", Ref: "refs/heads/main"}, + }, + Workflows: []*tangled.Pipeline_Workflow{{ + Name: "ci.yml", + Raw: `when: + - event: [pull_request, push, custom.example.trigger] +`, + }}, + } + record, err := FromTangled(models.PipelineId("pipeline"), time.Unix(1, 0), raw) + if err != nil { + t.Fatal(err) + } + definition := record.Workflows[0].Definition + wantTriggers := []string{TriggerNSIDPullRequest, TriggerNSIDPush, "custom.example.trigger"} + if !slices.Equal(definition.Triggers, wantTriggers) { + t.Fatalf("triggers = %v, want %v", definition.Triggers, wantTriggers) + } + if definition.Hash != nil { + t.Fatalf("hash was invented: %q", *definition.Hash) + } + file := definition.Source.File + if file == nil || file.Repo != source || file.Commit != record.Commit || file.Path != ".tangled/workflows/ci.yml" { + t.Fatalf("file source = %+v", file) + } +} + +func TestFromTangledDoesNotInferDefinitionTriggers(t *testing.T) { + repo := "did:plc:repo" + record, err := FromTangled(models.PipelineId("pipeline"), time.Now(), tangled.Pipeline{ + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ + Kind: "manual", Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &repo}, + Manual: &tangled.Pipeline_ManualTriggerData{Sha: "1111111111111111111111111111111111111111"}, + }, + Workflows: []*tangled.Pipeline_Workflow{{Name: "ci.yml", Raw: `engine: nixery +`}}, + }) + if err != nil { + t.Fatal(err) + } + triggers := record.Workflows[0].Definition.Triggers + if triggers == nil || len(triggers) != 0 { + t.Fatalf("implicit triggers = %#v", triggers) + } +} diff --git a/spindle/server.go b/spindle/server.go index a546df305..38fa39417 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -45,6 +45,7 @@ import ( "tangled.org/core/spindle/mill/executor" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" + pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/storage" @@ -873,6 +874,14 @@ func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host(), ownership.OwnerDid.String(), ownership.Rkey.String(), repoDid.String()), nil } +func (s *Spindle) createPipeline(id models.PipelineId, raw tangled.Pipeline) error { + record, err := pipelinecodec.FromTangled(id, time.Now(), raw) + if err != nil { + return err + } + return s.db.CreatePipeline(record) +} + func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger tangled.Pipeline_TriggerMetadata, changedFiles []string, repoCloneUri, repoPath, rev string, only []string, sourceRepo *tangled.Pipeline_TriggerRepo) (models.PipelineId, error) { l := log.FromContext(ctx) @@ -881,7 +890,7 @@ func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger t Trigger: trigger, } - rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) + rawPipeline, _, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) if err != nil { return "", fmt.Errorf("loading pipeline: %w", err) } @@ -906,7 +915,7 @@ func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger t } pipelineId := models.PipelineId(tid.TID()) - if err := s.db.CreatePipeline(pipelineId, tpl); err != nil { + if err := s.createPipeline(pipelineId, tpl); err != nil { return "", fmt.Errorf("creating pipeline: %w", err) } err = s.processPipeline(ctx, repoDid, tpl, pipelineId, sourceRepo) @@ -1026,13 +1035,36 @@ func (s *Spindle) resolveCheckout(ctx context.Context, repoDid syntax.DID, sourc // resolves the workflow definition at sha without executing it // returns a deterministic fingerprint over the resolved files. +func (s *Spindle) ListWorkflowDefinitions(ctx context.Context, repoDid syntax.DID, ref string) ([]*models.WorkflowDefinition, error) { + if ref == "" { + ref = "HEAD" + } + repoCloneURI, repoPath, _, err := s.resolveCheckout(ctx, repoDid, "") + if err != nil { + return nil, err + } + rawPipeline, commit, err := s.loadPipeline(ctx, repoCloneURI, repoPath, ref) + if err != nil { + return nil, fmt.Errorf("loading pipeline: %w", err) + } + out := make([]*models.WorkflowDefinition, 0, len(rawPipeline)) + for _, raw := range rawPipeline { + definition, err := pipelinecodec.FileDefinition(raw.Name, raw.Contents, repoDid.String(), commit) + if err != nil { + return nil, fmt.Errorf("parsing workflow %s: %w", raw.Name, err) + } + out = append(out, definition) + } + return out, nil +} + func (s *Spindle) DescribeWorkflowDefinition(ctx context.Context, repoDid syntax.DID, sha string, sourceRepo syntax.DID) (*tangled.CiDescribeWorkflowDefinition_Output, error) { repoCloneUri, repoPath, _, err := s.resolveCheckout(ctx, repoDid, sourceRepo) if err != nil { return nil, err } - rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, sha) + rawPipeline, _, err := s.loadPipeline(ctx, repoCloneUri, repoPath, sha) if err != nil { return nil, fmt.Errorf("loading pipeline: %w", err) } @@ -1068,21 +1100,20 @@ func fingerprintWorkflowDefinition(rawPipeline workflow.RawPipeline) string { return fmt.Sprintf("sha256:%x", h.Sum(nil)) } -func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { +func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, string, error) { if err := gitutil.SparseSync(ctx, repoUri, repoPath, rev, sparseWorkflowDir); err != nil { - return nil, fmt.Errorf("syncing git repo: %w", err) + return nil, "", fmt.Errorf("syncing git repo: %w", err) } gr, err := kgit.Open(repoPath, rev) if err != nil { - return nil, fmt.Errorf("opening git repo: %w", err) + return nil, "", fmt.Errorf("opening git repo: %w", err) } workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) if errors.Is(err, object.ErrDirectoryNotFound) { - // return empty RawPipeline when directory doesn't exist - return nil, nil + return nil, gr.Hash().String(), nil } else if err != nil { - return nil, fmt.Errorf("loading file tree: %w", err) + return nil, "", fmt.Errorf("loading file tree: %w", err) } var rawPipeline workflow.RawPipeline @@ -1094,7 +1125,7 @@ func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev strin fpath := filepath.Join(workflow.WorkflowDir, e.Name) contents, err := gr.RawContent(fpath) if err != nil { - return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) + return nil, "", fmt.Errorf("reading raw content of '%s': %w", fpath, err) } rawPipeline = append(rawPipeline, workflow.RawWorkflow{ @@ -1103,7 +1134,7 @@ func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev strin }) } - return rawPipeline, nil + return rawPipeline, gr.Hash().String(), nil } func (s *Spindle) newRepoPath(repo syntax.DID) string { diff --git a/spindle/tapclient.go b/spindle/tapclient.go index 64104ba84..29b81880e 100644 --- a/spindle/tapclient.go +++ b/spindle/tapclient.go @@ -680,7 +680,7 @@ func (s *Spindle) triggerPullRequestPipeline(ctx context.Context, l *slog.Logger repoPath := s.newRepoPath(repo.RepoDid) // load workflow definitions from rev (without spindle context) - rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha) + rawPipeline, _, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha) if err != nil { // don't retry l.Error("failed loading pipeline", "err", err) @@ -704,7 +704,7 @@ func (s *Spindle) triggerPullRequestPipeline(ctx context.Context, l *slog.Logger } pipelineId := models.PipelineId(tid.TID()) - if err := s.db.CreatePipeline(pipelineId, tpl); err != nil { + if err := s.createPipeline(pipelineId, tpl); err != nil { l.Error("failed to create pipeline event", "err", err) return nil } diff --git a/spindle/xrpc/ci_query_pipelines.go b/spindle/xrpc/ci_query_pipelines.go index c738a5bc4..24727c845 100644 --- a/spindle/xrpc/ci_query_pipelines.go +++ b/spindle/xrpc/ci_query_pipelines.go @@ -7,6 +7,7 @@ import ( "tangled.org/core/api/tangled" "tangled.org/core/spindle/models" + pipelinecodec "tangled.org/core/spindle/pipeline" xrpcerr "tangled.org/core/xrpc/errors" ) @@ -35,11 +36,15 @@ func (x *Xrpc) HandleCiQueryPipelines(w http.ResponseWriter, r *http.Request) { } } - pipelines, nextCursor, total, err := x.Db.QueryPipelines(r.Context(), repo, commits, cursor, kinds, limit) + records, nextCursor, total, err := x.Db.QueryPipelines(r.Context(), repo, commits, cursor, kinds, limit) if err != nil { fail(xrpcerr.GenericError(err), http.StatusInternalServerError) return } + pipelines := make([]*tangled.CiPipeline, 0, len(records)) + for _, record := range records { + pipelines = append(pipelines, pipelinecodec.ToTangled(record)) + } output := tangled.CiQueryPipelines_Output{ Pipelines: pipelines, @@ -67,13 +72,13 @@ func (x *Xrpc) HandleCiGetPipeline(w http.ResponseWriter, r *http.Request) { return } - p, err := x.Db.GetPipeline(r.Context(), models.PipelineId(pipeline)) + record, err := x.Db.GetPipeline(r.Context(), models.PipelineId(pipeline)) if err != nil { fail(xrpcerr.GenericError(err), http.StatusInternalServerError) return } - if err := writeJson(w, http.StatusOK, p); err != nil { + if err := writeJson(w, http.StatusOK, pipelinecodec.ToTangled(record)); err != nil { fail(xrpcerr.GenericError(err), http.StatusInternalServerError) } } diff --git a/spindle/xrpc/pipeline_cancel_pipeline.go b/spindle/xrpc/pipeline_cancel_pipeline.go index 83b662f51..b59e27344 100644 --- a/spindle/xrpc/pipeline_cancel_pipeline.go +++ b/spindle/xrpc/pipeline_cancel_pipeline.go @@ -51,7 +51,7 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { fail(xrpcerr.GenericError(fmt.Errorf("failed to get pipeline: %w", err))) return } - if p.Repo != repoDid.String() { + if p.RepoDID != repoDid.String() { fail(xrpcerr.AccessControlError(actorDid.String())) return } @@ -75,7 +75,7 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { for _, wName := range workflows { wid := models.WorkflowId{ PipelineId: pipelineId, - Name: wName, + Name: wName, } l.DebugContext(r.Context(), "cancel pipeline", "wid", wid) diff --git a/spindle/xrpc/xrpc_test.go b/spindle/xrpc/xrpc_test.go index 7778b9464..3dd54e828 100644 --- a/spindle/xrpc/xrpc_test.go +++ b/spindle/xrpc/xrpc_test.go @@ -20,6 +20,7 @@ import ( "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" + pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/spindle/secrets" ) @@ -36,6 +37,17 @@ func (m *mockTrigger) DescribeWorkflowDefinition(context.Context, syntax.DID, st return &tangled.CiDescribeWorkflowDefinition_Output{}, 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) + if err != nil { + t.Fatal(err) + } + if err := database.CreatePipeline(record); err != nil { + t.Fatal(err) + } +} + func newTestXrpcDB(t *testing.T) (*db.DB, *rbac.Enforcer) { t.Helper() p := filepath.Join(t.TempDir(), "spindle_xrpc.db") @@ -191,10 +203,7 @@ func TestCancelPipeline_RBAC(t *testing.T) { {Name: "test-workflow"}, }, } - err = d.CreatePipeline(models.PipelineId(pipelineTid), tpl) - if err != nil { - t.Fatalf("CreatePipeline: %v", err) - } + createTestPipeline(t, d, models.PipelineId(pipelineTid), tpl) x := &Xrpc{ Logger: slog.Default(),