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