This repository has no description
0

Configure Feed

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

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