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