Monorepo for Tangled tangled.org
1

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / db / pipelines.go
5.9 kB 225 lines
1package db 2 3import ( 4 "context" 5 "encoding/json" 6 "strconv" 7 "strings" 8 "time" 9 10 "tangled.org/core/api/tangled" 11 "tangled.org/core/spindle/models" 12) 13 14func (d *DB) QueryPipelines(ctx context.Context, repoDid string, commits []string, cursor string, limit int) ([]*tangled.CiDefs_Pipeline, string, int64, error) { 15 if limit <= 0 { 16 limit = 30 17 } 18 19 var query string 20 var args []interface{} 21 query = ` 22 select 23 rkey, event, created from events 24 where 25 nsid = 'sh.tangled.pipeline' 26 and coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), json_extract(event, '$.triggerMetadata.repo.did')) = ? 27 ` 28 args = append(args, repoDid) 29 30 if len(commits) > 0 { 31 placeholders := make([]string, len(commits)) 32 for i := range commits { 33 placeholders[i] = "?" 34 args = append(args, commits[i]) 35 } 36 query += " and json_extract(event, '$.triggerMetadata.push.newSha') in (" + strings.Join(placeholders, ",") + ")" 37 } 38 39 if cursor != "" { 40 if cVal, err := strconv.ParseInt(cursor, 10, 64); err == nil { 41 query += " and created < ?" 42 args = append(args, cVal) 43 } 44 } 45 46 // First get total count 47 var total int64 48 countQuery := "select count(*) from (" + query + ")" 49 if err := d.QueryRowContext(ctx, countQuery, args...).Scan(&total); err != nil { 50 return nil, "", 0, err 51 } 52 53 query += " order by created desc limit ?" 54 args = append(args, limit) 55 56 rows, err := d.QueryContext(ctx, query, args...) 57 if err != nil { 58 return nil, "", 0, err 59 } 60 defer rows.Close() 61 62 var pipelines []*tangled.CiDefs_Pipeline 63 var lastCreated int64 64 65 for rows.Next() { 66 var rkey, eventJson string 67 var created int64 68 if err := rows.Scan(&rkey, &eventJson, &created); err != nil { 69 return nil, "", 0, err 70 } 71 lastCreated = created 72 73 var rawPipeline tangled.Pipeline 74 if err := json.Unmarshal([]byte(eventJson), &rawPipeline); err != nil { 75 continue 76 } 77 78 p, err := d.mapToCiDefsPipeline(ctx, rkey, created, rawPipeline) 79 if err != nil { 80 return nil, "", 0, err 81 } 82 pipelines = append(pipelines, p) 83 } 84 85 nextCursor := "" 86 if len(pipelines) == limit { 87 nextCursor = strconv.FormatInt(lastCreated, 10) 88 } 89 90 return pipelines, nextCursor, total, nil 91} 92 93func (d *DB) GetPipeline(ctx context.Context, rkey string) (*tangled.CiDefs_Pipeline, error) { 94 var eventJson string 95 var created int64 96 err := d.QueryRowContext(ctx, 97 ` 98 select 99 event, created from events 100 where 101 nsid = 'sh.tangled.pipeline' 102 and rkey = ? 103 `, 104 rkey, 105 ).Scan(&eventJson, &created) 106 107 if err != nil { 108 return nil, err 109 } 110 111 var rawPipeline tangled.Pipeline 112 if err := json.Unmarshal([]byte(eventJson), &rawPipeline); err != nil { 113 return nil, err 114 } 115 116 return d.mapToCiDefsPipeline(ctx, rkey, created, rawPipeline) 117} 118 119func (d *DB) mapToCiDefsPipeline(ctx context.Context, rkey string, created int64, raw tangled.Pipeline) (*tangled.CiDefs_Pipeline, error) { 120 createdAtStr := time.Unix(0, created).Format(time.RFC3339) 121 122 var repoDidStr string 123 if raw.TriggerMetadata != nil && raw.TriggerMetadata.Repo != nil { 124 if raw.TriggerMetadata.Repo.RepoDid != nil { 125 repoDidStr = *raw.TriggerMetadata.Repo.RepoDid 126 } else { 127 repoDidStr = raw.TriggerMetadata.Repo.Did 128 } 129 } 130 131 commitSha := "" 132 var trigger tangled.CiDefs_Pipeline_Trigger 133 134 if raw.TriggerMetadata != nil { 135 switch raw.TriggerMetadata.Kind { 136 case "push": 137 if raw.TriggerMetadata.Push != nil { 138 commitSha = raw.TriggerMetadata.Push.NewSha 139 trigger.CiTrigger_Push = &tangled.CiTrigger_Push{ 140 NewSha: raw.TriggerMetadata.Push.NewSha, 141 OldSha: raw.TriggerMetadata.Push.OldSha, 142 Ref: raw.TriggerMetadata.Push.Ref, 143 } 144 } 145 case "pullRequest": 146 if raw.TriggerMetadata.PullRequest != nil { 147 commitSha = raw.TriggerMetadata.PullRequest.SourceSha 148 trigger.CiTrigger_PullRequest = &tangled.CiTrigger_PullRequest{ 149 Action: raw.TriggerMetadata.PullRequest.Action, 150 SourceBranch: &raw.TriggerMetadata.PullRequest.SourceBranch, 151 SourceSha: raw.TriggerMetadata.PullRequest.SourceSha, 152 TargetBranch: raw.TriggerMetadata.PullRequest.TargetBranch, 153 } 154 } 155 case "manual": 156 if raw.TriggerMetadata.Manual != nil { 157 commitSha = raw.TriggerMetadata.Manual.Sha 158 trigger.CiTrigger_Manual = &tangled.CiTrigger_Manual{} 159 } 160 } 161 } 162 163 var workflows []*tangled.CiDefs_Workflow 164 for _, wf := range raw.Workflows { 165 status := "pending" 166 var startedAt, finishedAt, wfError *string 167 168 if raw.TriggerMetadata != nil && raw.TriggerMetadata.Repo != nil { 169 wfId := models.WorkflowId{ 170 PipelineId: models.PipelineId{ 171 Knot: raw.TriggerMetadata.Repo.Knot, 172 Rkey: rkey, 173 }, 174 Name: wf.Name, 175 } 176 177 wfStatus, err := d.GetStatus(wfId) 178 if err == nil && wfStatus != nil { 179 status = wfStatus.Status 180 startedAt, finishedAt = d.GetWorkflowTimes(wfId) 181 wfError = wfStatus.Error 182 } 183 } 184 185 workflows = append(workflows, &tangled.CiDefs_Workflow{ 186 Id: wf.Name, 187 Name: wf.Name, 188 Status: status, 189 StartedAt: startedAt, 190 FinishedAt: finishedAt, 191 Error: wfError, 192 }) 193 } 194 195 return &tangled.CiDefs_Pipeline{ 196 Id: rkey, 197 Commit: commitSha, 198 Repo: &repoDidStr, 199 CreatedAt: &createdAtStr, 200 Trigger: &trigger, 201 Workflows: workflows, 202 }, nil 203} 204 205func (d *DB) GetWorkflowTimes(workflowId models.WorkflowId) (startedAt, finishedAt *string) { 206 pipelineAtUri := workflowId.PipelineId.AtUri() 207 208 _ = d.QueryRow( 209 ` 210 select 211 min(case when json_extract(event, '$.status') = 'running' then json_extract(event, '$.createdAt') end), 212 max(case when json_extract(event, '$.status') in ('success', 'failed', 'timeout', 'cancelled') then json_extract(event, '$.createdAt') end) 213 from events 214 where 215 nsid = ? 216 and json_extract(event, '$.pipeline') = ? 217 and json_extract(event, '$.workflow') = ? 218 `, 219 tangled.PipelineStatusNSID, 220 string(pipelineAtUri), 221 workflowId.Name, 222 ).Scan(&startedAt, &finishedAt) 223 224 return 225}