This repository has no description
0

Configure Feed

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

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