Monorepo for Tangled
tangled.org
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}