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