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)
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}