This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / tapclient.go
16 kB 515 lines
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}