This repository has no description
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}