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