forked from
mitchellh.com/tack
Stitch any CI into Tangled
1package main
2
3// Spindle-owned pipeline queries used by the sh.tangled.ci XRPC surface.
4
5import (
6 "context"
7 "database/sql"
8 "encoding/json"
9 "errors"
10 "fmt"
11 "strconv"
12 "strings"
13 "time"
14
15 "tangled.org/core/api/tangled"
16 "tangled.org/core/workflow"
17)
18
19// QueryPipelines returns newest-first pipeline summaries for one repository.
20// The filtering and cursor semantics mirror the stock Spindle implementation.
21func (s *store) QueryPipelines(
22 ctx context.Context,
23 repoDID string,
24 commits []string,
25 cursor string,
26 kinds []string,
27 limit int,
28) ([]*tangled.CiPipeline, string, int64, error) {
29 if limit <= 0 {
30 limit = 30
31 }
32 if limit > 250 {
33 limit = 250
34 }
35
36 // Build one parameterized base query, then wrap it for the commit-specific
37 // deduplication contract below. RepoDid is preferred, with Did retained only
38 // for pipeline records created before stable repository DIDs were universal.
39 query := `
40 SELECT rkey, event_json, created
41 FROM events
42 WHERE nsid = ?
43 AND COALESCE(
44 json_extract(event_json, '$.triggerMetadata.repo.repoDid'),
45 json_extract(event_json, '$.triggerMetadata.repo.did')
46 ) = ?`
47 args := []any{tangled.PipelineNSID, repoDID}
48
49 if len(commits) > 0 {
50 // Commit identity lives in a different union branch for each trigger kind.
51 // COALESCE gives filtering and later partitioning one common SQL key.
52 query += ` AND COALESCE(
53 json_extract(event_json, '$.triggerMetadata.push.newSha'),
54 json_extract(event_json, '$.triggerMetadata.pullRequest.sourceSha'),
55 json_extract(event_json, '$.triggerMetadata.manual.sha')
56 ) IN (` + sqlPlaceholders(len(commits)) + `)`
57 for _, commit := range commits {
58 args = append(args, commit)
59 }
60 }
61 if len(kinds) > 0 {
62 query += ` AND json_extract(event_json, '$.triggerMetadata.kind') IN (` +
63 sqlPlaceholders(len(kinds)) + `)`
64 for _, kind := range kinds {
65 args = append(args, kind)
66 }
67 }
68 var cursorValue *int64
69 if cursor != "" {
70 value, err := strconv.ParseInt(cursor, 10, 64)
71 if err != nil {
72 return nil, "", 0, fmt.Errorf("invalid cursor %q: %w", cursor, err)
73 }
74 cursorValue = &value
75 if len(commits) == 0 {
76 query += ` AND created < ?`
77 args = append(args, value)
78 }
79 }
80 if len(commits) > 0 {
81 // The query contract permits at most one pipeline per requested
82 // commit. Keep the newest retry so callers that index by SHA do not
83 // accidentally replace it with an older run.
84 query = `
85 SELECT rkey, event_json, created
86 FROM (
87 SELECT rkey, event_json, created,
88 ROW_NUMBER() OVER (
89 PARTITION BY COALESCE(
90 json_extract(event_json, '$.triggerMetadata.push.newSha'),
91 json_extract(event_json, '$.triggerMetadata.pullRequest.sourceSha'),
92 json_extract(event_json, '$.triggerMetadata.manual.sha')
93 )
94 ORDER BY created DESC
95 ) AS commit_rank
96 FROM (` + query + `)
97 )
98 WHERE commit_rank = 1`
99 if cursorValue != nil {
100 query += ` AND created < ?`
101 args = append(args, *cursorValue)
102 }
103 }
104
105 var total int64
106 if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM (`+query+`)`, args...).Scan(&total); err != nil {
107 return nil, "", 0, fmt.Errorf("count pipelines: %w", err)
108 }
109
110 rows, err := s.db.QueryContext(ctx, query+` ORDER BY created DESC LIMIT ?`, append(args, limit)...)
111 if err != nil {
112 return nil, "", 0, fmt.Errorf("query pipelines: %w", err)
113 }
114 defer rows.Close()
115 // Buffer raw rows before status projection. store deliberately permits only
116 // one SQLite connection; querying statuses while rows is open would deadlock
117 // waiting for that same connection.
118 type storedPipeline struct {
119 rkey string
120 raw tangled.Pipeline
121 created int64
122 }
123 stored := make([]storedPipeline, 0, limit)
124 for rows.Next() {
125 var rkey string
126 var rawJSON []byte
127 var created int64
128 if err := rows.Scan(&rkey, &rawJSON, &created); err != nil {
129 return nil, "", 0, fmt.Errorf("scan pipeline: %w", err)
130 }
131 var raw tangled.Pipeline
132 if err := json.Unmarshal(rawJSON, &raw); err != nil {
133 return nil, "", 0, fmt.Errorf("decode pipeline %q: %w", rkey, err)
134 }
135 stored = append(stored, storedPipeline{rkey: rkey, raw: raw, created: created})
136 }
137 if err := rows.Err(); err != nil {
138 rows.Close()
139 return nil, "", 0, fmt.Errorf("iterate pipelines: %w", err)
140 }
141 if err := rows.Close(); err != nil {
142 return nil, "", 0, fmt.Errorf("close pipeline rows: %w", err)
143 }
144
145 // mapPipeline performs status lookups. Close the outer rows first because
146 // Tack intentionally uses one SQLite connection to serialize writers.
147 pipelines := make([]*tangled.CiPipeline, 0, len(stored))
148 var lastCreated int64
149 for i := range stored {
150 pipeline, err := s.mapPipeline(ctx, stored[i].rkey, stored[i].created, &stored[i].raw)
151 if err != nil {
152 return nil, "", 0, err
153 }
154 pipelines = append(pipelines, pipeline)
155 lastCreated = stored[i].created
156 }
157
158 nextCursor := ""
159 if len(pipelines) == limit {
160 nextCursor = strconv.FormatInt(lastCreated, 10)
161 }
162 return pipelines, nextCursor, total, nil
163}
164
165// HasPullPipeline makes live repo.pull processing idempotent across
166// Jetstream's replay window and harmless metadata-only pull updates.
167func (s *store) HasPullPipeline(ctx context.Context, pullURI, sourceSHA string) (bool, error) {
168 var count int
169 err := s.db.QueryRowContext(ctx, `
170 SELECT COUNT(*)
171 FROM events
172 WHERE nsid = ?
173 AND json_extract(event_json, '$.triggerMetadata.pullRequest.pull') = ?
174 AND json_extract(event_json, '$.triggerMetadata.pullRequest.sourceSha') = ?`,
175 tangled.PipelineNSID, pullURI, sourceSHA,
176 ).Scan(&count)
177 if err != nil {
178 return false, fmt.Errorf("check pull pipeline: %w", err)
179 }
180 return count > 0, nil
181}
182
183// sqlPlaceholders returns a comma-separated placeholder list for a non-empty
184// filter. Callers append values separately, so user input never enters SQL text.
185func sqlPlaceholders(count int) string {
186 return strings.TrimSuffix(strings.Repeat("?,", count), ",")
187}
188
189// GetPipeline returns one spindle-local pipeline summary by TID.
190func (s *store) GetPipeline(ctx context.Context, rkey string) (*tangled.CiPipeline, error) {
191 raw, created, err := s.GetRawPipeline(ctx, rkey)
192 if err != nil {
193 return nil, err
194 }
195 return s.mapPipeline(ctx, rkey, created, raw)
196}
197
198// GetRawPipeline returns the retained legacy record used internally to launch
199// providers. The public API maps it to sh.tangled.ci.pipeline.
200func (s *store) GetRawPipeline(ctx context.Context, rkey string) (*tangled.Pipeline, int64, error) {
201 var rawJSON []byte
202 var created int64
203 err := s.db.QueryRowContext(ctx, `
204 SELECT event_json, created
205 FROM events
206 WHERE nsid = ? AND rkey = ?
207 LIMIT 1`, tangled.PipelineNSID, rkey,
208 ).Scan(&rawJSON, &created)
209 if err != nil {
210 return nil, 0, err
211 }
212 var raw tangled.Pipeline
213 if err := json.Unmarshal(rawJSON, &raw); err != nil {
214 return nil, 0, fmt.Errorf("decode pipeline %q: %w", rkey, err)
215 }
216 return &raw, created, nil
217}
218
219// mapPipeline projects the internal sh.tangled.pipeline record and retained
220// workflow statuses into the spindle-owned sh.tangled.ci.pipeline response.
221// Keeping the legacy record internally avoids duplicating Tangled's compiler
222// and provider trigger types at write time.
223func (s *store) mapPipeline(
224 ctx context.Context,
225 rkey string,
226 created int64,
227 raw *tangled.Pipeline,
228) (*tangled.CiPipeline, error) {
229 createdAt := time.Unix(0, created).UTC().Format(time.RFC3339)
230 result := &tangled.CiPipeline{
231 Id: rkey,
232 CreatedAt: &createdAt,
233 }
234
235 var knot string
236 if raw.TriggerMetadata != nil {
237 result.SourceRepo = raw.TriggerMetadata.SourceRepo
238 if repo := raw.TriggerMetadata.Repo; repo != nil {
239 knot = repo.Knot
240 repoDID := repo.Did
241 if repo.RepoDid != nil {
242 repoDID = *repo.RepoDid
243 }
244 if repoDID != "" {
245 result.Repo = &repoDID
246 }
247 }
248 result.Trigger, result.Commit = ciTrigger(raw.TriggerMetadata)
249 }
250
251 for _, wf := range raw.Workflows {
252 if wf == nil || wf.Name == "" {
253 continue
254 }
255 status, startedAt, finishedAt, statusErr, err := s.workflowStatus(
256 ctx, knot, rkey, wf.Name,
257 )
258 if err != nil {
259 return nil, fmt.Errorf("load status for pipeline %q workflow %q: %w", rkey, wf.Name, err)
260 }
261 result.Workflows = append(result.Workflows, &tangled.CiPipeline_Workflow{
262 Id: wf.Name,
263 Name: wf.Name,
264 Status: status,
265 StartedAt: startedAt,
266 FinishedAt: finishedAt,
267 Error: statusErr,
268 })
269 }
270 return result, nil
271}
272
273// ciTrigger converts legacy Pipeline_TriggerMetadata into the public CI trigger
274// union and returns the commit SHA used by CiPipeline.Commit. Unknown or
275// structurally incomplete trigger variants produce a nil trigger and empty SHA
276// rather than inventing metadata.
277func ciTrigger(metadata *tangled.Pipeline_TriggerMetadata) (*tangled.CiPipeline_Trigger, string) {
278 trigger := &tangled.CiPipeline_Trigger{}
279 switch workflow.TriggerKind(metadata.Kind) {
280 case workflow.TriggerKindPush:
281 if value := metadata.Push; value != nil {
282 trigger.CiTrigger_Push = &tangled.CiTrigger_Push{
283 NewSha: value.NewSha,
284 OldSha: value.OldSha,
285 Ref: value.Ref,
286 }
287 return trigger, value.NewSha
288 }
289 case workflow.TriggerKindPullRequest:
290 if value := metadata.PullRequest; value != nil {
291 sourceBranch := value.SourceBranch
292 trigger.CiTrigger_PullRequest = &tangled.CiTrigger_PullRequest{
293 Pull: value.Pull,
294 SourceBranch: &sourceBranch,
295 SourceRepo: metadata.SourceRepo,
296 SourceSha: value.SourceSha,
297 TargetBranch: value.TargetBranch,
298 }
299 return trigger, value.SourceSha
300 }
301 case workflow.TriggerKindManual:
302 if value := metadata.Manual; value != nil {
303 inputs := make([]*tangled.CiTrigger_Pair, 0, len(value.Inputs))
304 for _, pair := range value.Inputs {
305 if pair != nil {
306 inputs = append(inputs, &tangled.CiTrigger_Pair{Key: pair.Key, Value: pair.Value})
307 }
308 }
309 trigger.CiTrigger_Manual = &tangled.CiTrigger_Manual{
310 Inputs: inputs,
311 Ref: value.Ref,
312 Sha: value.Sha,
313 SourceRepo: metadata.SourceRepo,
314 }
315 return trigger, value.Sha
316 }
317 }
318 return nil, ""
319}
320
321// workflowStatus returns the latest status/error plus the first running and
322// last terminal timestamps for one workflow. Missing status rows mean pending:
323// a pipeline is persisted before asynchronous provider setup publishes its
324// first explicit status.
325func (s *store) workflowStatus(
326 ctx context.Context,
327 knot, pipelineRkey, workflowName string,
328) (string, *string, *string, *string, error) {
329 status := "pending"
330 if knot == "" {
331 return status, nil, nil, nil, nil
332 }
333 pipelineURI := pipelineATURI(knot, pipelineRkey)
334
335 var rawJSON []byte
336 // Status is an append-only event stream. Ordering by the monotonic created
337 // value gives the current state without maintaining a second mutable table.
338 err := s.db.QueryRowContext(ctx, `
339 SELECT event_json
340 FROM events
341 WHERE nsid = ?
342 AND json_extract(event_json, '$.pipeline') = ?
343 AND json_extract(event_json, '$.workflow') = ?
344 ORDER BY created DESC
345 LIMIT 1`, tangled.PipelineStatusNSID, pipelineURI, workflowName,
346 ).Scan(&rawJSON)
347 var statusErr *string
348 if err == nil {
349 var record tangled.PipelineStatus
350 if err := json.Unmarshal(rawJSON, &record); err != nil {
351 return "", nil, nil, nil, err
352 }
353 status = record.Status
354 statusErr = record.Error
355 } else if !errors.Is(err, sql.ErrNoRows) {
356 return "", nil, nil, nil, err
357 }
358
359 // Timestamps come from record payloads, not insertion time: provider webhook
360 // delivery may lag the actual transition we want to display.
361 var started, finished sql.NullString
362 err = s.db.QueryRowContext(ctx, `
363 SELECT
364 MIN(CASE
365 WHEN json_extract(event_json, '$.status') = 'running'
366 THEN json_extract(event_json, '$.createdAt')
367 END),
368 MAX(CASE
369 WHEN json_extract(event_json, '$.status') IN ('success', 'failed', 'timeout', 'cancelled')
370 THEN json_extract(event_json, '$.createdAt')
371 END)
372 FROM events
373 WHERE nsid = ?
374 AND json_extract(event_json, '$.pipeline') = ?
375 AND json_extract(event_json, '$.workflow') = ?`,
376 tangled.PipelineStatusNSID, pipelineURI, workflowName,
377 ).Scan(&started, &finished)
378 if err != nil {
379 return "", nil, nil, nil, err
380 }
381 return status, nullStringPointer(started), nullStringPointer(finished), statusErr, nil
382}
383
384// nullStringPointer converts SQLite aggregate NULLs to optional lexicon fields.
385func nullStringPointer(value sql.NullString) *string {
386 if !value.Valid {
387 return nil
388 }
389 return &value.String
390}