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
5.8 kB 224 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 trigger.CiTrigger_Manual = &tangled.CiTrigger_Manual{} 158 } 159 } 160 } 161 162 var workflows []*tangled.CiDefs_Workflow 163 for _, wf := range raw.Workflows { 164 status := "pending" 165 var startedAt, finishedAt, wfError *string 166 167 if raw.TriggerMetadata != nil && raw.TriggerMetadata.Repo != nil { 168 wfId := models.WorkflowId{ 169 PipelineId: models.PipelineId{ 170 Knot: raw.TriggerMetadata.Repo.Knot, 171 Rkey: rkey, 172 }, 173 Name: wf.Name, 174 } 175 176 wfStatus, err := d.GetStatus(wfId) 177 if err == nil && wfStatus != nil { 178 status = wfStatus.Status 179 startedAt, finishedAt = d.GetWorkflowTimes(wfId) 180 wfError = wfStatus.Error 181 } 182 } 183 184 workflows = append(workflows, &tangled.CiDefs_Workflow{ 185 Id: wf.Name, 186 Name: wf.Name, 187 Status: status, 188 StartedAt: startedAt, 189 FinishedAt: finishedAt, 190 Error: wfError, 191 }) 192 } 193 194 return &tangled.CiDefs_Pipeline{ 195 Id: rkey, 196 Commit: commitSha, 197 Repo: &repoDidStr, 198 CreatedAt: &createdAtStr, 199 Trigger: &trigger, 200 Workflows: workflows, 201 }, nil 202} 203 204func (d *DB) GetWorkflowTimes(workflowId models.WorkflowId) (startedAt, finishedAt *string) { 205 pipelineAtUri := workflowId.PipelineId.AtUri() 206 207 _ = d.QueryRow( 208 ` 209 select 210 min(case when json_extract(event, '$.status') = 'running' then json_extract(event, '$.createdAt') end), 211 max(case when json_extract(event, '$.status') in ('success', 'failed', 'timeout', 'cancelled') then json_extract(event, '$.createdAt') end) 212 from events 213 where 214 nsid = ? 215 and json_extract(event, '$.pipeline') = ? 216 and json_extract(event, '$.workflow') = ? 217 `, 218 tangled.PipelineStatusNSID, 219 string(pipelineAtUri), 220 workflowId.Name, 221 ).Scan(&startedAt, &finishedAt) 222 223 return 224}