This repository has no description
1package spindle
2
3import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "log/slog"
10 "net"
11 "net/http"
12 "net/url"
13 "sync"
14 "time"
15
16 "github.com/bluesky-social/indigo/atproto/syntax"
17 indigoxrpc "github.com/bluesky-social/indigo/xrpc"
18 "tangled.org/core/api/tangled"
19 avmodels "tangled.org/core/appview/models"
20 "tangled.org/core/eventconsumer"
21 "tangled.org/core/log"
22 "tangled.org/core/rbac"
23 "tangled.org/core/spindle/db"
24 "tangled.org/core/spindle/git"
25 "tangled.org/core/spindle/models"
26 "tangled.org/core/spindle/netguard"
27 "tangled.org/core/tapc"
28 "tangled.org/core/tid"
29 "tangled.org/core/workflow"
30)
31
32const (
33 maxPendingPerRepo = 64
34 pendingCollabTTL = 10 * time.Minute
35)
36
37// blobs are fetched from user controlled PDSes so protect our transport
38// from dialing internal addresses
39var guardedBlobClient = &http.Client{
40 Transport: &http.Transport{
41 DialContext: (&net.Dialer{
42 Timeout: 30 * time.Second,
43 KeepAlive: 30 * time.Second,
44 Control: netguard.RefuseSpecialPurposeAddrs,
45 }).DialContext,
46 ForceAttemptHTTP2: true,
47 MaxIdleConns: 100,
48 IdleConnTimeout: 90 * time.Second,
49 TLSHandshakeTimeout: 10 * time.Second,
50 ExpectContinueTimeout: 1 * time.Second,
51 },
52}
53
54type pendingCollabEvent struct {
55 evt *tapc.RecordEventData
56 at time.Time
57}
58
59type Tap struct {
60 logger *slog.Logger
61 spindle *Spindle
62 tap tapc.Client
63 pendingMu sync.Mutex
64 pendingCollabs map[syntax.DID][]pendingCollabEvent
65}
66
67func NewTapClient(s *Spindle) *Tap {
68 return &Tap{
69 logger: log.SubLogger(s.l, "tapclient"),
70 spindle: s,
71 tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword),
72 pendingCollabs: make(map[syntax.DID][]pendingCollabEvent),
73 }
74}
75
76func (t *Tap) AddOwnerDIDs(ctx context.Context, dids []syntax.DID) error {
77 if len(dids) == 0 {
78 return nil
79 }
80 return t.tap.AddRepos(ctx, dids)
81}
82
83func (t *Tap) Start(connCtx context.Context) {
84 go t.tap.Connect(connCtx, &tapc.SimpleIndexer{
85 EventHandler: t.processEvent,
86 ConnectHandler: t.onConnect,
87 })
88 go t.purgePendingCollabsLoop(t.spindle.rootCtx)
89}
90
91func (t *Tap) onConnect(ctx context.Context) {
92 t.spindle.declareTapInterest(ctx)
93}
94
95func (t *Tap) processEvent(ctx context.Context, evt tapc.Event) error {
96 if evt.Type != tapc.EvtRecord || evt.Record == nil {
97 return nil
98 }
99 switch evt.Record.Collection.String() {
100 case tangled.RepoNSID:
101 return t.processRepo(ctx, evt.Record)
102 case tangled.RepoCollaboratorNSID:
103 return t.processCollaborator(ctx, evt.Record)
104 }
105 return nil
106}
107
108func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error {
109 l := t.logger.With("collection", tangled.RepoNSID, "did", evt.Did, "rkey", evt.Rkey)
110
111 ownerDid := evt.Did
112 rkey := evt.Rkey
113
114 switch evt.Action {
115 case tapc.RecordCreateAction, tapc.RecordUpdateAction:
116 record := tangled.Repo{}
117 if err := json.Unmarshal(evt.Record, &record); err != nil {
118 l.Warn("skipping invalid repo record", "err", err)
119 return nil
120 }
121
122 hostname := t.spindle.cfg.Server.Hostname
123 prior, priorErr := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey)
124 knownRepo := priorErr == nil
125
126 if record.Spindle == nil || *record.Spindle != hostname {
127 if knownRepo {
128 l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle)
129 return t.teardownRepo(l, prior, ownerDid, rkey)
130 }
131 return nil
132 }
133
134 if record.RepoDid == nil || *record.RepoDid == "" {
135 l.Warn("skipping repo record without repoDid")
136 return nil
137 }
138 repoDid, err := syntax.ParseDID(*record.RepoDid)
139 if err != nil {
140 l.Warn("skipping repo record with malformed repoDid", "value", *record.RepoDid, "err", err)
141 return nil
142 }
143
144 isMember, err := t.spindle.e.IsSpindleMember(ownerDid.String(), rbac.ThisServer)
145 if err != nil {
146 return fmt.Errorf("checking spindle membership: %w", err)
147 }
148 if !isMember {
149 l.Warn("rejecting repo record: owner is not a spindle member", "owner", ownerDid)
150 return nil
151 }
152
153 // check if this repo DID is already owned by someone else
154 existingRepo, err := t.spindle.db.GetRepoByDid(repoDid)
155 if err == nil {
156 if existingRepo.Owner != ownerDid {
157 l.Warn("rejecting repo record: repoDid already registered by another owner", "repoDid", repoDid, "existingOwner", existingRepo.Owner, "newOwner", ownerDid)
158 return nil
159 }
160 } else if !errors.Is(err, sql.ErrNoRows) {
161 return fmt.Errorf("lookup existing repo by DID: %w", err)
162 }
163
164 if err := t.spindle.e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()); err != nil {
165 l.Error("failed to add repo policy", "err", err)
166 return fmt.Errorf("add repo policy: %w", err)
167 }
168
169 src := eventconsumer.NewKnotSource(record.Knot)
170 t.spindle.ks.AddSource(t.spindle.rootCtx, src)
171
172 repo := db.Repo{
173 Knot: record.Knot,
174 Owner: ownerDid,
175 Rkey: rkey,
176 RepoDid: repoDid,
177 CreatedAt: record.CreatedAt,
178 }
179
180 if err := t.spindle.db.AddRepo(repo); err != nil {
181 l.Error("failed to add repo row", "err", err)
182 return fmt.Errorf("add repo: %w", err)
183 }
184
185 // setup sparse sync
186 repoCloneUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid)
187 repoPath := t.spindle.newRepoPath(repo.RepoDid)
188 if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, ""); err != nil {
189 return fmt.Errorf("setting up sparse-clone git repo: %w", err)
190 }
191
192 legacyName := ""
193 if record.Name != nil {
194 legacyName = *record.Name
195 }
196 migrateLegacyRepoSecrets(ctx, t.spindle.db, t.spindle.vault, l, ownerDid, legacyName, rkey, repoDid)
197 migrateLegacyRepoCasbin(ctx, t.spindle.db, t.spindle.e, l, ownerDid, legacyName, rkey, repoDid)
198
199 if removed, err := t.spindle.db.CollapseRepoSiblings(ownerDid, repoDid); err != nil {
200 l.Warn("collapse rename siblings failed", "err", err)
201 } else if removed > 0 {
202 l.Info("collapsed rename leftovers", "owner", ownerDid, "repo_did", repoDid, "removed", removed)
203 }
204
205 if e := t.spindle.embedTap; e == nil || !e.closed.Load() {
206 if err := t.tap.AddRepos(ctx, []syntax.DID{ownerDid}); err != nil {
207 l.Warn("tap AddRepos rejected", "did", ownerDid, "err", err)
208 }
209 }
210 t.spindle.jc.AddDid(ownerDid.String())
211
212 t.drainPendingCollabs(ctx, repoDid)
213
214 case tapc.RecordDeleteAction:
215 repo, err := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey)
216 if err != nil {
217 l.Info("skipping delete for unknown repo")
218 return nil
219 }
220 return t.teardownRepo(l, repo, ownerDid, rkey)
221 }
222 return nil
223}
224
225func (t *Tap) teardownRepo(l *slog.Logger, repo *db.Repo, ownerDid syntax.DID, rkey syntax.RecordKey) error {
226 if repo.RepoDid != "" {
227 collabs, err := t.spindle.db.ListCollaboratorsByRepoDid(repo.RepoDid)
228 if err != nil {
229 l.Error("failed to list collaborators for cleanup", "err", err)
230 return fmt.Errorf("list collaborators: %w", err)
231 }
232 for _, c := range collabs {
233 if err := t.spindle.e.RemoveCollaborator(c.Subject.String(), rbac.ThisServer, repo.RepoDid.String()); err != nil {
234 l.Error("failed to remove collaborator policy", "subject", c.Subject, "err", err)
235 return fmt.Errorf("remove collaborator policy: %w", err)
236 }
237 }
238 if err := t.spindle.db.DeleteRepoCollaboratorsByRepoDid(repo.RepoDid); err != nil {
239 l.Error("failed to clear collaborator rows", "err", err)
240 return err
241 }
242 if err := t.spindle.e.RemoveRepo(ownerDid.String(), rbac.ThisServer, repo.RepoDid.String()); err != nil {
243 l.Error("failed to remove repo policy", "err", err)
244 return fmt.Errorf("remove repo policy: %w", err)
245 }
246 }
247 if err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey); err != nil {
248 l.Error("failed to delete repo row", "err", err)
249 return fmt.Errorf("delete repo row: %w", err)
250 }
251 // TODO: clear sparse-synced git repo
252 return nil
253}
254
255func (t *Tap) processCollaborator(ctx context.Context, evt *tapc.RecordEventData) error {
256 l := t.logger.With("collection", tangled.RepoCollaboratorNSID, "did", evt.Did, "rkey", evt.Rkey)
257
258 switch evt.Action {
259 case tapc.RecordCreateAction, tapc.RecordUpdateAction:
260 record := tangled.RepoCollaborator{}
261 if err := json.Unmarshal(evt.Record, &record); err != nil {
262 l.Warn("skipping invalid collaborator record", "err", err)
263 return nil
264 }
265
266 actor := evt.Did
267 rkey := evt.Rkey
268
269 subjectDid, err := syntax.ParseDID(record.Subject)
270 if err != nil {
271 l.Info("skipping collaborator with malformed subject DID", "subject", record.Subject, "err", err)
272 return nil
273 }
274 if _, err := t.spindle.res.ResolveIdent(ctx, subjectDid.String()); err != nil {
275 l.Info("skipping unresolvable collaborator subject", "subject", subjectDid, "err", err)
276 return nil
277 }
278
279 repoRefDid, err := syntax.ParseDID(record.Repo)
280 if err != nil {
281 l.Info("skipping collaborator with non-DID repo ref", "repo", record.Repo, "err", err)
282 return nil
283 }
284 repo, lookupErr := t.spindle.db.GetRepoByDid(repoRefDid)
285 if errors.Is(lookupErr, sql.ErrNoRows) {
286 t.bufferCollab(repoRefDid, evt)
287 l.Info("buffering collaborator until repo arrives", "repo", repoRefDid)
288 return nil
289 }
290 if lookupErr != nil {
291 return fmt.Errorf("lookup repo %s: %w", repoRefDid, lookupErr)
292 }
293 repoDid := repo.RepoDid
294 ownerDid := repo.Owner
295
296 if actor != ownerDid {
297 l.Info("rejecting collaborator with non-owner actor", "actor", actor, "owner", ownerDid)
298 return nil
299 }
300
301 ok, err := t.spindle.e.IsCollaboratorInviteAllowed(ownerDid.String(), rbac.ThisServer, repoDid.String())
302 if err != nil {
303 l.Error("invite permission check failed", "err", err)
304 return fmt.Errorf("invite check: %w", err)
305 }
306 if !ok {
307 l.Info("rejecting collaborator invite", "owner", ownerDid, "repo", repoDid)
308 return nil
309 }
310
311 prior, priorErr := t.spindle.db.GetRepoCollaborator(actor, rkey)
312 staleSubject := priorErr == nil && (prior.Subject != subjectDid || prior.RepoDid != repoDid)
313
314 if err := t.spindle.e.AddCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()); err != nil {
315 l.Error("failed to add collaborator policy", "err", err)
316 return fmt.Errorf("add collaborator policy: %w", err)
317 }
318 if staleSubject {
319 if err := t.spindle.e.RemoveCollaborator(prior.Subject.String(), rbac.ThisServer, prior.RepoDid.String()); err != nil {
320 l.Error("failed to remove stale collaborator policy", "err", err)
321 return fmt.Errorf("remove stale collaborator: %w", err)
322 }
323 }
324 if err := t.spindle.db.AddRepoCollaborator(db.RepoCollaborator{
325 OwnerDid: actor,
326 Rkey: rkey,
327 Subject: subjectDid,
328 RepoDid: repoDid,
329 }); err != nil {
330 l.Error("failed to persist collaborator row", "err", err)
331 return fmt.Errorf("track collaborator: %w", err)
332 }
333
334 case tapc.RecordDeleteAction:
335 actor := evt.Did
336 rkey := evt.Rkey
337
338 tracked, err := t.spindle.db.GetRepoCollaborator(actor, rkey)
339 if err != nil {
340 l.Info("skipping delete for unknown collaborator record")
341 return nil
342 }
343 if err := t.spindle.e.RemoveCollaborator(tracked.Subject.String(), rbac.ThisServer, tracked.RepoDid.String()); err != nil {
344 l.Error("failed to remove collaborator policy", "err", err)
345 return fmt.Errorf("remove collaborator policy: %w", err)
346 }
347 if err := t.spindle.db.DeleteRepoCollaborator(actor, rkey); err != nil {
348 l.Error("failed to delete collaborator row", "err", err)
349 return fmt.Errorf("delete collaborator row: %w", err)
350 }
351 }
352 return nil
353}
354
355func (s *Spindle) processPull(ctx context.Context, evt *tapc.RecordEventData) error {
356 l := s.l.With("component", "ingester", "collection", evt.Collection, "did", evt.Did, "rkey", evt.Rkey)
357
358 // only listen to live events
359 if !evt.Live {
360 l.Info("skipping backfill event", "event", evt.AtUri())
361 return nil
362 }
363
364 switch evt.Action {
365 case tapc.RecordCreateAction, tapc.RecordUpdateAction:
366 record := tangled.RepoPull{}
367 if err := json.Unmarshal(evt.Record, &record); err != nil {
368 l.Error("invalid record", "err", err)
369 return fmt.Errorf("parsing record: %w", err)
370 }
371
372 // ignore legacy records
373 if record.Target == nil {
374 l.Info("ignoring pull record: target repo is nil")
375 return nil
376 }
377
378 // ignore patch-based and fork-based PRs
379 if record.Source == nil || record.Source.Repo != nil {
380 l.Info("ignoring pull record: not a branch-based pull request")
381 return nil
382 }
383
384 // skip if target repo is unknown
385 repo, err := s.db.GetRepoByDid(syntax.DID(record.Target.Repo))
386 if err != nil {
387 l.Warn("target repo is not ingested yet", "repo", record.Target.Repo, "err", err)
388 return fmt.Errorf("target repo is unknown")
389 }
390
391 // only accept branch-based PR (excluding patch-based and fork-based)
392 if record.Source == nil || record.Source.Repo != nil {
393 l.Warn("skipping non-branch-based PR")
394 return nil
395 }
396
397 // check if pull record author has push access to target repo
398 allowed, err := s.e.IsPushAllowed(evt.Did.String(), rbac.ThisServer, repo.RepoDid.String())
399 if err != nil {
400 return fmt.Errorf("checking push access for pull record author: %w", err)
401 }
402 if !allowed {
403 l.Warn("rejecting pull-triggered pipeline. author has no push access",
404 "author", evt.Did, "repo", repo.RepoDid)
405 return nil
406 }
407
408 latestSubmission, err := s.fetchLatestSubmission(ctx, evt.Did.String(), evt.Rkey.String(), &record)
409 if err != nil {
410 return err
411 }
412 sourceSha := latestSubmission.SourceRev
413
414 scheme := "https"
415 if s.cfg.Server.Dev {
416 scheme = "http"
417 }
418 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, repo.Knot)}
419
420 // fetch current default branch
421 defaultBranch, _ := func(repo syntax.DID) (string, error) {
422 defaultBranchOut, err := tangled.RepoGetDefaultBranch(ctx, client, repo.String())
423 if err != nil {
424 return "", err
425 }
426 return defaultBranchOut.Name, nil
427 }(repo.RepoDid)
428
429 compiler := workflow.Compiler{
430 Trigger: tangled.Pipeline_TriggerMetadata{
431 Kind: string(workflow.TriggerKindPullRequest),
432 PullRequest: &tangled.Pipeline_PullRequestTriggerData{
433 SourceBranch: record.Source.Branch,
434 SourceSha: sourceSha,
435 TargetBranch: record.Target.Branch,
436 },
437 Repo: &tangled.Pipeline_TriggerRepo{
438 Did: repo.Owner.String(),
439 Knot: repo.Knot,
440 Repo: (*string)(&repo.Rkey),
441 RepoDid: (*string)(&repo.RepoDid),
442 DefaultBranch: defaultBranch,
443 },
444 },
445 }
446
447 repoUri := s.newRepoCloneUrl(repo.Knot, repo.RepoDid)
448 repoPath := s.newRepoPath(repo.RepoDid)
449
450 // load workflow definitions from rev (without spindle context)
451 rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha)
452 if err != nil {
453 // don't retry
454 l.Error("failed loading pipeline", "err", err)
455 return nil
456 }
457 if len(rawPipeline) == 0 {
458 l.Info("no workflow definition find for the repo. skipping the event")
459 return nil
460 }
461 tpl := compiler.Compile(compiler.Parse(rawPipeline))
462 // TODO: pass compile error to workflow log
463 for _, w := range compiler.Diagnostics.Errors {
464 l.Error(w.String())
465 }
466 for _, w := range compiler.Diagnostics.Warnings {
467 l.Warn(w.String())
468 }
469 if len(tpl.Workflows) == 0 {
470 l.Info("no workflow matching trigger 'pull_request'. skipping the event")
471 return nil
472 }
473
474 pipelineId := models.PipelineId{
475 Knot: tpl.TriggerMetadata.Repo.Knot,
476 Rkey: tid.TID(),
477 }
478 if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil {
479 l.Error("failed to create pipeline event", "err", err)
480 return nil
481 }
482 sourceRepo, err := s.resolvePipelineSourceRepo(ctx, tpl.TriggerMetadata)
483 if err != nil {
484 l.Error("failed resolving pipeline source repo", "err", err)
485 return nil
486 }
487 err = s.processPipeline(repo.RepoDid, tpl, pipelineId, sourceRepo)
488 if err != nil {
489 // don't retry
490 l.Error("failed processing pipeline", "err", err)
491 return nil
492 }
493 case tapc.RecordDeleteAction:
494 // no-op
495 }
496 return nil
497}
498
499func (t *Tap) bufferCollab(repoDid syntax.DID, evt *tapc.RecordEventData) {
500 t.pendingMu.Lock()
501 defer t.pendingMu.Unlock()
502 list := t.pendingCollabs[repoDid]
503 list = append(list, pendingCollabEvent{evt: evt, at: time.Now()})
504 if len(list) > maxPendingPerRepo {
505 list = list[len(list)-maxPendingPerRepo:]
506 }
507 t.pendingCollabs[repoDid] = list
508}
509
510func (t *Tap) drainPendingCollabs(ctx context.Context, repoDid syntax.DID) {
511 t.pendingMu.Lock()
512 list := t.pendingCollabs[repoDid]
513 delete(t.pendingCollabs, repoDid)
514 t.pendingMu.Unlock()
515 if len(list) == 0 {
516 return
517 }
518 cutoff := time.Now().Add(-pendingCollabTTL)
519 for _, p := range list {
520 if p.at.Before(cutoff) {
521 continue
522 }
523 if err := t.processCollaborator(ctx, p.evt); err != nil {
524 t.logger.Warn("replaying buffered collaborator failed", "repo", repoDid, "rkey", p.evt.Rkey, "err", err)
525 }
526 }
527}
528
529func (t *Tap) purgePendingCollabsLoop(ctx context.Context) {
530 ticker := time.NewTicker(pendingCollabTTL / 2)
531 defer ticker.Stop()
532 for {
533 select {
534 case <-ctx.Done():
535 return
536 case <-ticker.C:
537 t.purgeStalePendingCollabs()
538 }
539 }
540}
541
542func (t *Tap) purgeStalePendingCollabs() {
543 cutoff := time.Now().Add(-pendingCollabTTL)
544 t.pendingMu.Lock()
545 defer t.pendingMu.Unlock()
546 expired := 0
547 for did, list := range t.pendingCollabs {
548 kept := list[:0]
549 for _, p := range list {
550 if !p.at.Before(cutoff) {
551 kept = append(kept, p)
552 } else {
553 expired++
554 }
555 }
556 if len(kept) == 0 {
557 delete(t.pendingCollabs, did)
558 } else {
559 t.pendingCollabs[did] = kept
560 }
561 }
562 if expired > 0 {
563 t.logger.Warn("expired buffered collaborator events without matching repo arrival", "count", expired, "ttl", pendingCollabTTL)
564 }
565}
566
567func (s *Spindle) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) {
568 // resolve the PR owner's identity to fetch the blob from their PDS
569 prOwnerIdent, err := s.res.ResolveIdent(ctx, did)
570 if err != nil || prOwnerIdent.Handle.IsInvalidHandle() {
571 return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err)
572 }
573
574 if len(record.Rounds) == 0 {
575 return nil, fmt.Errorf("failed to fetch latest submission, no rounds in record")
576 }
577
578 roundNumber := len(record.Rounds) - 1
579 round := record.Rounds[roundNumber]
580
581 // fetch the blob from the PR owner's PDS
582 prOwnerPds := prOwnerIdent.PDSEndpoint()
583 blobUrl, err := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", prOwnerPds))
584 if err != nil {
585 return nil, fmt.Errorf("failed to construct blob URL: %w", err)
586 }
587 q := blobUrl.Query()
588 q.Set("cid", round.PatchBlob.Ref.String())
589 q.Set("did", did)
590 blobUrl.RawQuery = q.Encode()
591
592 req, err := http.NewRequestWithContext(ctx, http.MethodGet, blobUrl.String(), nil)
593 if err != nil {
594 return nil, fmt.Errorf("failed to create blob request: %w", err)
595 }
596 req.Header.Set("Content-Type", "application/json")
597
598 blobResp, err := guardedBlobClient.Do(req)
599 if err != nil {
600 return nil, fmt.Errorf("failed to fetch blob: %w", err)
601 }
602 defer blobResp.Body.Close()
603
604 latestSubmission, err := avmodels.PullSubmissionFromRecord(did, rkey, roundNumber, round, blobResp.Body)
605 if err != nil {
606 return nil, fmt.Errorf("failed to parse submission: %w", err)
607 }
608
609 return latestSubmission, nil
610}