This repository has no description
1package spindle
2
3import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "log/slog"
8 "net/http"
9 "net/url"
10 "slices"
11 "strings"
12 "time"
13
14 "github.com/bluesky-social/indigo/atproto/syntax"
15 indigoxrpc "github.com/bluesky-social/indigo/xrpc"
16 "tangled.org/core/api/tangled"
17 avmodels "tangled.org/core/appview/models"
18 "tangled.org/core/eventconsumer"
19 "tangled.org/core/log"
20 "tangled.org/core/spindle/db"
21 "tangled.org/core/spindle/git"
22 "tangled.org/core/spindle/models"
23 "tangled.org/core/tapc"
24 "tangled.org/core/tid"
25 "tangled.org/core/workflow"
26)
27
28const (
29 collaboratorPageLimit = 100
30 maxCollaboratorPages = 50
31)
32
33type Tap struct {
34 logger *slog.Logger
35 spindle *Spindle
36 tap tapc.Client
37}
38
39func NewTapClient(s *Spindle) *Tap {
40 return &Tap{
41 logger: log.SubLogger(s.l, "tapclient"),
42 spindle: s,
43 tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword),
44 }
45}
46
47func (t *Tap) Start(connCtx context.Context) {
48 go t.tap.Connect(connCtx, &tapc.SimpleIndexer{
49 EventHandler: t.processEvent,
50 ConnectHandler: t.onConnect,
51 })
52}
53
54func (t *Tap) onConnect(ctx context.Context) {
55 l := t.logger
56 if t.spindle.cfg.Server.InviteOnly {
57 // listen to owners of registered repositories
58 owners, err := t.spindle.db.RepoOwners()
59 if err != nil {
60 l.Warn("tap declare: failed to load known repos", "err", err)
61 return
62 }
63 if err := t.tap.AddRepos(ctx, owners); err != nil {
64 l.Warn("tap declare: AddRepos rejected", "count", len(owners), "err", err)
65 return
66 }
67 l.Info("tap declare: known owner DIDs registered", "count", len(owners))
68 } else {
69 // public spindle. listen to full network
70 l.Info("tap declare: listening to full network")
71 }
72}
73
74func (t *Tap) processEvent(ctx context.Context, evt tapc.Event) error {
75 if evt.Type != tapc.EvtRecord || evt.Record == nil {
76 return nil
77 }
78 switch evt.Record.Collection.String() {
79 case tangled.RepoNSID:
80 return t.processRepo(ctx, evt.Record)
81 }
82 return nil
83}
84
85// processRepo ingests repo declaration record: `sh.tangled.repo`. It skips alias records.
86func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error {
87 l := t.logger.With("collection", tangled.RepoNSID, "did", evt.Did, "rkey", evt.Rkey)
88
89 ownerDid := evt.Did
90 rkey := evt.Rkey
91
92 switch evt.Action {
93 case tapc.RecordCreateAction, tapc.RecordUpdateAction:
94 record := tangled.Repo{}
95 if err := json.Unmarshal(evt.Record, &record); err != nil {
96 l.Warn("skipping invalid repo record", "err", err)
97 return nil
98 }
99
100 if record.RepoDid == nil || *record.RepoDid == "" {
101 l.Warn("skipping repo record without repoDid")
102 return nil
103 }
104 repoDid, err := syntax.ParseDID(*record.RepoDid)
105 if err != nil {
106 l.Warn("skipping repo record with malformed repoDid", "value", *record.RepoDid, "err", err)
107 return nil
108 }
109
110 isMember, err := t.spindle.db.IsAllowedMember(ctx, ownerDid, !t.spindle.cfg.Server.InviteOnly)
111 if err != nil {
112 return fmt.Errorf("checking spindle membership: %w", err)
113 }
114 if !isMember {
115 l.Warn("rejecting repo record: owner is not a spindle member", "owner", ownerDid)
116 return nil
117 }
118
119 // ignore repos not pointing this spindle.
120 hostname := t.spindle.cfg.Server.Hostname
121 if record.Spindle == nil || *record.Spindle != hostname {
122 // teardown existing repo
123 prior, err := t.spindle.db.GetRepoByDid(repoDid)
124 if err != nil {
125 return nil
126 }
127 if prior.Owner == ownerDid && prior.Rkey == rkey {
128 l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle)
129 return t.teardownRepo(l, prior.RepoDid)
130 }
131 l.Warn("ignoring reassignment from non-registering record", "owner", prior.Owner, "rkey", prior.Rkey)
132 return nil
133 }
134
135 // verify repo declaration
136 // NOTE: we are ignoring repo declaration records that *might* become correct pointer in
137 // future with repository rename or ownership transfer. On repo rename, new pointer record
138 // should always be recreated even when it already exists.
139 verified, err := t.verifyRepoDeclaration(ctx, l, repoDid, ownerDid, rkey.String(), record.Knot)
140 if err != nil {
141 l.Warn("failed to verify repo declaration", "err", err)
142 return nil
143 }
144 if !verified {
145 // ignore alias records
146 return nil
147 }
148
149 if err := t.spindle.e.SetRepoOwner(ownerDid, repoDid); err != nil {
150 l.Error("failed to add repo policy", "err", err)
151 return fmt.Errorf("add repo policy: %w", err)
152 }
153
154 repo := db.Repo{
155 Knot: record.Knot,
156 Owner: ownerDid,
157 Rkey: rkey,
158 RepoDid: repoDid,
159 CreatedAt: record.CreatedAt,
160 }
161
162 if err := t.spindle.db.UpsertRepo(repo); err != nil {
163 l.Error("failed to add repo row", "err", err)
164 return fmt.Errorf("add repo: %w", err)
165 }
166
167 t.reconcileCollaborators(ctx, l, record.Knot, repoDid, ownerDid)
168 t.spindle.ks.AddSource(t.spindle.rootCtx, eventconsumer.NewKnotSource(record.Knot))
169
170 // setup sparse sync
171 repoCloneUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid)
172 repoPath := t.spindle.newRepoPath(repo.RepoDid)
173 if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, ""); err != nil {
174 return fmt.Errorf("setting up sparse-clone git repo: %w", err)
175 }
176
177 legacyName := ""
178 if record.Name != nil {
179 legacyName = *record.Name
180 }
181 migrateLegacyRepoSecrets(ctx, t.spindle.db, t.spindle.vault, l, ownerDid, legacyName, rkey, repoDid)
182
183 if e := t.spindle.embedTap; e == nil || !e.closed.Load() {
184 if err := t.tap.AddRepos(ctx, []syntax.DID{ownerDid}); err != nil {
185 l.Warn("tap AddRepos rejected", "did", ownerDid, "err", err)
186 }
187 }
188 t.spindle.jc.AddDid(ownerDid.String())
189
190 case tapc.RecordDeleteAction:
191 repoDid, err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey)
192 if err != nil {
193 return fmt.Errorf("deleting repo record: %w", err)
194 }
195
196 if repoDid == "" {
197 // no record is deleted. record was not pointing this spindle or was an alias record
198 return nil
199 }
200 return t.teardownRepo(l, repoDid)
201 }
202 return nil
203}
204
205func (t *Tap) verifyRepoDeclaration(ctx context.Context, l *slog.Logger, repo, owner syntax.DID, rkey, knot string) (bool, error) {
206 result, err := t.spindle.VerifyRepo(ctx, repo)
207 if err != nil {
208 return false, err
209 }
210 l = l.With(
211 "repoDid", result.RepoDid,
212 "repoOwner", result.OwnerDid,
213 "repoKnot", result.KnotURL.Host,
214 "claimedOwner", owner,
215 "claimedRkey", rkey,
216 "claimedKnot", knot,
217 )
218 if syntax.DID(result.OwnerDid) != owner {
219 l.Warn("rejecting repo event: owner mismatch")
220 return false, nil
221 }
222 if result.Rkey != rkey {
223 l.Warn("rejecting repo event: rkey mismatch")
224 return false, nil
225 }
226 if !strings.EqualFold(knot, result.KnotURL.Host) {
227 l.Warn("rejecting repo event: record knot does not match DID-doc endpoint")
228 return false, nil
229 }
230 return true, nil
231}
232
233func (t *Tap) reconcileCollaborators(ctx context.Context, l *slog.Logger, knot string, repo, owner syntax.DID) {
234 wanted, err := t.fetchKnotCollaborators(ctx, knot, repo)
235 if err != nil {
236 l.Warn("collaborator reconcile: failed to fetch roster from knot", "knot", knot, "err", err)
237 return
238 }
239
240 have, err := t.spindle.e.GetRepoCollaborators(repo)
241 if err != nil {
242 l.Warn("collaborator reconcile: failed to read current grants", "err", err)
243 return
244 }
245
246 // the owner is a collaborator by role inheritance, never by an explicit grant
247 delete(wanted, owner)
248
249 for did := range wanted {
250 if slices.Contains(have, did) {
251 continue
252 }
253 if err := t.spindle.grantCollaborator(did, repo); err != nil {
254 l.Error("collaborator reconcile: failed to add", "subject", did, "err", err)
255 return
256 }
257 l.Info("collaborator reconcile: added", "subject", did)
258 }
259 for _, did := range have {
260 if did == owner {
261 continue
262 }
263 if _, ok := wanted[did]; ok {
264 continue
265 }
266 if err := t.spindle.e.RemoveRepoCollaborator(did, repo); err != nil {
267 l.Error("collaborator reconcile: failed to remove", "subject", did, "err", err)
268 return
269 }
270 l.Info("collaborator reconcile: removed", "subject", did)
271 }
272}
273
274func (t *Tap) fetchKnotCollaborators(ctx context.Context, knot string, repo syntax.DID) (map[syntax.DID]struct{}, error) {
275 scheme := "https"
276 if t.spindle.cfg.Server.Dev {
277 scheme = "http"
278 }
279 xc := &indigoxrpc.Client{
280 Host: fmt.Sprintf("%s://%s", scheme, knot),
281 Client: &http.Client{Timeout: 30 * time.Second},
282 }
283
284 subjects := make(map[syntax.DID]struct{})
285 cursor := ""
286 for range maxCollaboratorPages {
287 out, err := tangled.RepoListCollaborators(ctx, xc, cursor, collaboratorPageLimit, "", repo.String())
288 if err != nil {
289 return nil, err
290 }
291 for _, item := range out.Items {
292 did, err := syntax.ParseDID(item.Subject)
293 if err != nil {
294 continue
295 }
296 subjects[did] = struct{}{}
297 }
298 if out.Cursor == nil || *out.Cursor == "" {
299 return subjects, nil
300 }
301 cursor = *out.Cursor
302 }
303 return nil, fmt.Errorf("collaborator roster exceeded %d pages", maxCollaboratorPages)
304}
305
306func (t *Tap) teardownRepo(l *slog.Logger, repo syntax.DID) error {
307 if repo == "" {
308 return nil
309 }
310 if err := t.spindle.db.DeleteRepo(repo); err != nil {
311 l.Error("failed to remove repo", "err", err)
312 return fmt.Errorf("remove repo: %w", err)
313 }
314 if err := t.spindle.e.DeleteRepo(repo); err != nil {
315 l.Error("failed to remove repo policy", "err", err)
316 return fmt.Errorf("remove repo policy: %w", err)
317 }
318 // TODO: clear sparse-synced git repo
319 return nil
320}
321
322func (s *Spindle) processPull(ctx context.Context, evt *tapc.RecordEventData) error {
323 l := s.l.With("component", "ingester", "collection", evt.Collection, "did", evt.Did, "rkey", evt.Rkey)
324
325 // only listen to live events
326 if !evt.Live {
327 l.Info("skipping backfill event", "event", evt.AtUri())
328 return nil
329 }
330
331 switch evt.Action {
332 case tapc.RecordCreateAction, tapc.RecordUpdateAction:
333 record := tangled.RepoPull{}
334 if err := json.Unmarshal(evt.Record, &record); err != nil {
335 l.Error("invalid record", "err", err)
336 return fmt.Errorf("parsing record: %w", err)
337 }
338
339 // ignore legacy records
340 if record.Target == nil {
341 l.Info("ignoring pull record: target repo is nil")
342 return nil
343 }
344
345 // ignore patch-based and fork-based PRs
346 if record.Source == nil || record.Source.Repo != nil {
347 l.Info("ignoring pull record: not a branch-based pull request")
348 return nil
349 }
350
351 // skip if target repo is unknown
352 repo, err := s.db.GetRepoByDid(syntax.DID(record.Target.Repo))
353 if err != nil {
354 l.Warn("target repo is not ingested yet", "repo", record.Target.Repo, "err", err)
355 return fmt.Errorf("target repo is unknown")
356 }
357
358 // only accept branch-based PR (excluding patch-based and fork-based)
359 if record.Source == nil || record.Source.Repo != nil {
360 l.Warn("skipping non-branch-based PR")
361 return nil
362 }
363
364 // check if pull record author can trigger CI in target repo
365 allowed, err := s.e.IsRepoCiTriggerAllowed(evt.Did, repo.RepoDid)
366 if err != nil {
367 return fmt.Errorf("checking push access for pull record author: %w", err)
368 }
369 if !allowed {
370 l.Warn("rejecting pull-triggered pipeline. author has no push access",
371 "author", evt.Did, "repo", repo.RepoDid)
372 return nil
373 }
374
375 latestSubmission, err := s.fetchLatestSubmission(ctx, evt.Did.String(), evt.Rkey.String(), &record)
376 if err != nil {
377 return err
378 }
379 sourceSha := latestSubmission.SourceRev
380
381 scheme := "https"
382 if s.cfg.Server.Dev {
383 scheme = "http"
384 }
385 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, repo.Knot)}
386
387 // fetch current default branch
388 defaultBranch, _ := func(repo syntax.DID) (string, error) {
389 defaultBranchOut, err := tangled.RepoGetDefaultBranch(ctx, client, repo.String())
390 if err != nil {
391 return "", err
392 }
393 return defaultBranchOut.Name, nil
394 }(repo.RepoDid)
395
396 compiler := workflow.Compiler{
397 Trigger: tangled.Pipeline_TriggerMetadata{
398 Kind: string(workflow.TriggerKindPullRequest),
399 PullRequest: &tangled.Pipeline_PullRequestTriggerData{
400 SourceBranch: record.Source.Branch,
401 SourceSha: sourceSha,
402 TargetBranch: record.Target.Branch,
403 },
404 Repo: &tangled.Pipeline_TriggerRepo{
405 Did: repo.Owner.String(),
406 Knot: repo.Knot,
407 Repo: (*string)(&repo.Rkey),
408 RepoDid: (*string)(&repo.RepoDid),
409 DefaultBranch: defaultBranch,
410 },
411 },
412 }
413
414 repoUri := s.newRepoCloneUrl(repo.Knot, repo.RepoDid)
415 repoPath := s.newRepoPath(repo.RepoDid)
416
417 // load workflow definitions from rev (without spindle context)
418 rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha)
419 if err != nil {
420 // don't retry
421 l.Error("failed loading pipeline", "err", err)
422 return nil
423 }
424 if len(rawPipeline) == 0 {
425 l.Info("no workflow definition find for the repo. skipping the event")
426 return nil
427 }
428 tpl := compiler.Compile(compiler.Parse(rawPipeline))
429 // TODO: pass compile error to workflow log
430 for _, w := range compiler.Diagnostics.Errors {
431 l.Error(w.String())
432 }
433 for _, w := range compiler.Diagnostics.Warnings {
434 l.Warn(w.String())
435 }
436 if len(tpl.Workflows) == 0 {
437 l.Info("no workflow matching trigger 'pull_request'. skipping the event")
438 return nil
439 }
440
441 pipelineId := models.PipelineId{
442 Knot: tpl.TriggerMetadata.Repo.Knot,
443 Rkey: tid.TID(),
444 }
445 if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil {
446 l.Error("failed to create pipeline event", "err", err)
447 return nil
448 }
449 sourceRepo, err := s.resolvePipelineSourceRepo(ctx, tpl.TriggerMetadata)
450 if err != nil {
451 l.Error("failed resolving pipeline source repo", "err", err)
452 return nil
453 }
454 err = s.processPipeline(repo.RepoDid, tpl, pipelineId, sourceRepo)
455 if err != nil {
456 // don't retry
457 l.Error("failed processing pipeline", "err", err)
458 return nil
459 }
460 case tapc.RecordDeleteAction:
461 // no-op
462 }
463 return nil
464}
465
466func (s *Spindle) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) {
467 // resolve the PR owner's identity to fetch the blob from their PDS
468 prOwnerIdent, err := s.res.ResolveIdent(ctx, did)
469 if err != nil || prOwnerIdent.Handle.IsInvalidHandle() {
470 return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err)
471 }
472
473 if len(record.Rounds) == 0 {
474 return nil, fmt.Errorf("failed to fetch latest submission, no rounds in record")
475 }
476
477 roundNumber := len(record.Rounds) - 1
478 round := record.Rounds[roundNumber]
479
480 // fetch the blob from the PR owner's PDS
481 prOwnerPds := prOwnerIdent.PDSEndpoint()
482 blobUrl, err := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", prOwnerPds))
483 if err != nil {
484 return nil, fmt.Errorf("failed to construct blob URL: %w", err)
485 }
486 q := blobUrl.Query()
487 q.Set("cid", round.PatchBlob.Ref.String())
488 q.Set("did", did)
489 blobUrl.RawQuery = q.Encode()
490
491 req, err := http.NewRequestWithContext(ctx, http.MethodGet, blobUrl.String(), nil)
492 if err != nil {
493 return nil, fmt.Errorf("failed to create blob request: %w", err)
494 }
495 req.Header.Set("Content-Type", "application/json")
496
497 blobResp, err := http.DefaultClient.Do(req)
498 if err != nil {
499 return nil, fmt.Errorf("failed to fetch blob: %w", err)
500 }
501 defer blobResp.Body.Close()
502
503 latestSubmission, err := avmodels.PullSubmissionFromRecord(did, rkey, roundNumber, round, blobResp.Body)
504 if err != nil {
505 return nil, fmt.Errorf("failed to parse submission: %w", err)
506 }
507
508 return latestSubmission, nil
509}