package db import ( "context" "encoding/json" "strconv" "strings" "tangled.org/core/orm" "tangled.org/core/spindle/models" ) 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)} if len(commits) > 0 { filters = append(filters, orm.FilterIn("commit_sha", commits)) } if len(kinds) > 0 { filters = append(filters, orm.FilterIn("kind", kinds)) } var conditions []string var args []any for _, filter := range filters { conditions = append(conditions, filter.Condition()) args = append(args, filter.Arg()...) } whereClause := " where " + strings.Join(conditions, " and ") var total int64 if err := d.QueryRowContext(ctx, `select count(*) from pipelines`+whereClause, args...).Scan(&total); err != nil { return nil, "", 0, err } if cursor != "" { 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 ") } } args = append(args, limit) rows, err := d.QueryContext(ctx, `select id, payload from pipelines`+whereClause+` order by id desc limit ?`, args...) if err != nil { return nil, "", 0, err } defer rows.Close() 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 } lastID = id scanned++ var record models.PipelineRecord if json.Unmarshal([]byte(payload), &record) != nil { continue } records = append(records, &record) } if err := rows.Err(); err != nil { return nil, "", 0, err } if err := d.applyStatuses(ctx, records); err != nil { return nil, "", 0, err } nextCursor := "" if scanned == limit { nextCursor = strconv.FormatInt(lastID, 10) } return records, nextCursor, total, nil } 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, ).Scan(&payload); err != nil { return nil, err } var record models.PipelineRecord if err := json.Unmarshal([]byte(payload), &record); err != nil { return nil, err } if err := d.applyStatuses(ctx, []*models.PipelineRecord{&record}); err != nil { return nil, err } return &record, nil } 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 (?, ?, ?, ?, ?)`, record.ID, record.RepoDID, record.Commit, record.Trigger.Kind, string(payload), ) return err } func (d *DB) applyStatuses(ctx context.Context, pipelines []*models.PipelineRecord) error { if len(pipelines) == 0 { return nil } pipelineIDs := make([]string, 0, len(pipelines)) for _, pipeline := range pipelines { pipelineIDs = append(pipelineIDs, string(pipeline.ID)) } statuses, err := d.workflowStatuses(ctx, pipelineIDs) if err != nil { return err } for _, pipeline := range pipelines { for _, workflow := range pipeline.Workflows { if workflow == nil { continue } status, ok := statuses[wfKey{string(pipeline.ID), workflow.Name}] if !ok { continue } workflow.Status = status.Status workflow.Error = status.Error workflow.StartedAt = status.StartedAt workflow.FinishedAt = status.FinishedAt } } return nil } type wfKey struct { PipelineID string Workflow string } type wfStatus struct { Status string Error *string StartedAt *string FinishedAt *string } func (d *DB) workflowStatuses(ctx context.Context, pipelineIDs []string) (map[wfKey]wfStatus, error) { filter := orm.FilterIn("pipeline_id", pipelineIDs) rows, err := d.QueryContext(ctx, `select pipeline_id, workflow, status, error, created_at from workflow_statuses where `+filter.Condition()+` order by id asc`, filter.Arg()...) if err != nil { return nil, err } defer rows.Close() out := make(map[wfKey]wfStatus) for rows.Next() { var pipelineID, wfName, status, createdAt string var wfError *string if err := rows.Scan(&pipelineID, &wfName, &status, &wfError, &createdAt); err != nil { return nil, err } k := wfKey{pipelineID, wfName} st := out[k] // rows arrive in insertion order, so the last one wins for the current status st.Status = status st.Error = wfError switch kind := models.StatusKind(status); { case kind == models.StatusKindRunning: // the first running row is when the workflow actually started if st.StartedAt == nil { at := createdAt st.StartedAt = &at } case kind.IsFinish(): at := createdAt st.FinishedAt = &at } out[k] = st } return out, rows.Err() }