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