This repository has no description
0

Configure Feed

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

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