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