This repository has no description
0

Configure Feed

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

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