Stitch any CI into Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

tack / store_pipeline.go
12 kB 390 lines
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}