Monorepo for Tangled tangled.org
3

Configure Feed

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

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