This repository has no description
0

Configure Feed

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

core / spindle / tapclient.go
15 kB 509 lines
1package spindle 2 3import ( 4 "context" 5 "encoding/json" 6 "fmt" 7 "log/slog" 8 "net/http" 9 "net/url" 10 "slices" 11 "strings" 12 "time" 13 14 "github.com/bluesky-social/indigo/atproto/syntax" 15 indigoxrpc "github.com/bluesky-social/indigo/xrpc" 16 "tangled.org/core/api/tangled" 17 avmodels "tangled.org/core/appview/models" 18 "tangled.org/core/eventconsumer" 19 "tangled.org/core/log" 20 "tangled.org/core/spindle/db" 21 "tangled.org/core/spindle/git" 22 "tangled.org/core/spindle/models" 23 "tangled.org/core/tapc" 24 "tangled.org/core/tid" 25 "tangled.org/core/workflow" 26) 27 28const ( 29 collaboratorPageLimit = 100 30 maxCollaboratorPages = 50 31) 32 33type Tap struct { 34 logger *slog.Logger 35 spindle *Spindle 36 tap tapc.Client 37} 38 39func NewTapClient(s *Spindle) *Tap { 40 return &Tap{ 41 logger: log.SubLogger(s.l, "tapclient"), 42 spindle: s, 43 tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword), 44 } 45} 46 47func (t *Tap) Start(connCtx context.Context) { 48 go t.tap.Connect(connCtx, &tapc.SimpleIndexer{ 49 EventHandler: t.processEvent, 50 ConnectHandler: t.onConnect, 51 }) 52} 53 54func (t *Tap) onConnect(ctx context.Context) { 55 l := t.logger 56 if t.spindle.cfg.Server.InviteOnly { 57 // listen to owners of registered repositories 58 owners, err := t.spindle.db.RepoOwners() 59 if err != nil { 60 l.Warn("tap declare: failed to load known repos", "err", err) 61 return 62 } 63 if err := t.tap.AddRepos(ctx, owners); err != nil { 64 l.Warn("tap declare: AddRepos rejected", "count", len(owners), "err", err) 65 return 66 } 67 l.Info("tap declare: known owner DIDs registered", "count", len(owners)) 68 } else { 69 // public spindle. listen to full network 70 l.Info("tap declare: listening to full network") 71 } 72} 73 74func (t *Tap) processEvent(ctx context.Context, evt tapc.Event) error { 75 if evt.Type != tapc.EvtRecord || evt.Record == nil { 76 return nil 77 } 78 switch evt.Record.Collection.String() { 79 case tangled.RepoNSID: 80 return t.processRepo(ctx, evt.Record) 81 } 82 return nil 83} 84 85// processRepo ingests repo declaration record: `sh.tangled.repo`. It skips alias records. 86func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error { 87 l := t.logger.With("collection", tangled.RepoNSID, "did", evt.Did, "rkey", evt.Rkey) 88 89 ownerDid := evt.Did 90 rkey := evt.Rkey 91 92 switch evt.Action { 93 case tapc.RecordCreateAction, tapc.RecordUpdateAction: 94 record := tangled.Repo{} 95 if err := json.Unmarshal(evt.Record, &record); err != nil { 96 l.Warn("skipping invalid repo record", "err", err) 97 return nil 98 } 99 100 if record.RepoDid == nil || *record.RepoDid == "" { 101 l.Warn("skipping repo record without repoDid") 102 return nil 103 } 104 repoDid, err := syntax.ParseDID(*record.RepoDid) 105 if err != nil { 106 l.Warn("skipping repo record with malformed repoDid", "value", *record.RepoDid, "err", err) 107 return nil 108 } 109 110 isMember, err := t.spindle.db.IsAllowedMember(ctx, ownerDid, !t.spindle.cfg.Server.InviteOnly) 111 if err != nil { 112 return fmt.Errorf("checking spindle membership: %w", err) 113 } 114 if !isMember { 115 l.Warn("rejecting repo record: owner is not a spindle member", "owner", ownerDid) 116 return nil 117 } 118 119 // ignore repos not pointing this spindle. 120 hostname := t.spindle.cfg.Server.Hostname 121 if record.Spindle == nil || *record.Spindle != hostname { 122 // teardown existing repo 123 prior, err := t.spindle.db.GetRepoByDid(repoDid) 124 if err != nil { 125 return nil 126 } 127 if prior.Owner == ownerDid && prior.Rkey == rkey { 128 l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle) 129 return t.teardownRepo(l, prior.RepoDid) 130 } 131 l.Warn("ignoring reassignment from non-registering record", "owner", prior.Owner, "rkey", prior.Rkey) 132 return nil 133 } 134 135 // verify repo declaration 136 // NOTE: we are ignoring repo declaration records that *might* become correct pointer in 137 // future with repository rename or ownership transfer. On repo rename, new pointer record 138 // should always be recreated even when it already exists. 139 verified, err := t.verifyRepoDeclaration(ctx, l, repoDid, ownerDid, rkey.String(), record.Knot) 140 if err != nil { 141 l.Warn("failed to verify repo declaration", "err", err) 142 return nil 143 } 144 if !verified { 145 // ignore alias records 146 return nil 147 } 148 149 if err := t.spindle.e.SetRepoOwner(ownerDid, repoDid); err != nil { 150 l.Error("failed to add repo policy", "err", err) 151 return fmt.Errorf("add repo policy: %w", err) 152 } 153 154 repo := db.Repo{ 155 Knot: record.Knot, 156 Owner: ownerDid, 157 Rkey: rkey, 158 RepoDid: repoDid, 159 CreatedAt: record.CreatedAt, 160 } 161 162 if err := t.spindle.db.UpsertRepo(repo); err != nil { 163 l.Error("failed to add repo row", "err", err) 164 return fmt.Errorf("add repo: %w", err) 165 } 166 167 t.reconcileCollaborators(ctx, l, record.Knot, repoDid, ownerDid) 168 t.spindle.ks.AddSource(t.spindle.rootCtx, eventconsumer.NewKnotSource(record.Knot)) 169 170 // setup sparse sync 171 repoCloneUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid) 172 repoPath := t.spindle.newRepoPath(repo.RepoDid) 173 if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, ""); err != nil { 174 return fmt.Errorf("setting up sparse-clone git repo: %w", err) 175 } 176 177 legacyName := "" 178 if record.Name != nil { 179 legacyName = *record.Name 180 } 181 migrateLegacyRepoSecrets(ctx, t.spindle.db, t.spindle.vault, l, ownerDid, legacyName, rkey, repoDid) 182 183 if e := t.spindle.embedTap; e == nil || !e.closed.Load() { 184 if err := t.tap.AddRepos(ctx, []syntax.DID{ownerDid}); err != nil { 185 l.Warn("tap AddRepos rejected", "did", ownerDid, "err", err) 186 } 187 } 188 t.spindle.jc.AddDid(ownerDid.String()) 189 190 case tapc.RecordDeleteAction: 191 repoDid, err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey) 192 if err != nil { 193 return fmt.Errorf("deleting repo record: %w", err) 194 } 195 196 if repoDid == "" { 197 // no record is deleted. record was not pointing this spindle or was an alias record 198 return nil 199 } 200 return t.teardownRepo(l, repoDid) 201 } 202 return nil 203} 204 205func (t *Tap) verifyRepoDeclaration(ctx context.Context, l *slog.Logger, repo, owner syntax.DID, rkey, knot string) (bool, error) { 206 result, err := t.spindle.VerifyRepo(ctx, repo) 207 if err != nil { 208 return false, err 209 } 210 l = l.With( 211 "repoDid", result.RepoDid, 212 "repoOwner", result.OwnerDid, 213 "repoKnot", result.KnotURL.Host, 214 "claimedOwner", owner, 215 "claimedRkey", rkey, 216 "claimedKnot", knot, 217 ) 218 if syntax.DID(result.OwnerDid) != owner { 219 l.Warn("rejecting repo event: owner mismatch") 220 return false, nil 221 } 222 if result.Rkey != rkey { 223 l.Warn("rejecting repo event: rkey mismatch") 224 return false, nil 225 } 226 if !strings.EqualFold(knot, result.KnotURL.Host) { 227 l.Warn("rejecting repo event: record knot does not match DID-doc endpoint") 228 return false, nil 229 } 230 return true, nil 231} 232 233func (t *Tap) reconcileCollaborators(ctx context.Context, l *slog.Logger, knot string, repo, owner syntax.DID) { 234 wanted, err := t.fetchKnotCollaborators(ctx, knot, repo) 235 if err != nil { 236 l.Warn("collaborator reconcile: failed to fetch roster from knot", "knot", knot, "err", err) 237 return 238 } 239 240 have, err := t.spindle.e.GetRepoCollaborators(repo) 241 if err != nil { 242 l.Warn("collaborator reconcile: failed to read current grants", "err", err) 243 return 244 } 245 246 // the owner is a collaborator by role inheritance, never by an explicit grant 247 delete(wanted, owner) 248 249 for did := range wanted { 250 if slices.Contains(have, did) { 251 continue 252 } 253 if err := t.spindle.grantCollaborator(did, repo); err != nil { 254 l.Error("collaborator reconcile: failed to add", "subject", did, "err", err) 255 return 256 } 257 l.Info("collaborator reconcile: added", "subject", did) 258 } 259 for _, did := range have { 260 if did == owner { 261 continue 262 } 263 if _, ok := wanted[did]; ok { 264 continue 265 } 266 if err := t.spindle.e.RemoveRepoCollaborator(did, repo); err != nil { 267 l.Error("collaborator reconcile: failed to remove", "subject", did, "err", err) 268 return 269 } 270 l.Info("collaborator reconcile: removed", "subject", did) 271 } 272} 273 274func (t *Tap) fetchKnotCollaborators(ctx context.Context, knot string, repo syntax.DID) (map[syntax.DID]struct{}, error) { 275 scheme := "https" 276 if t.spindle.cfg.Server.Dev { 277 scheme = "http" 278 } 279 xc := &indigoxrpc.Client{ 280 Host: fmt.Sprintf("%s://%s", scheme, knot), 281 Client: &http.Client{Timeout: 30 * time.Second}, 282 } 283 284 subjects := make(map[syntax.DID]struct{}) 285 cursor := "" 286 for range maxCollaboratorPages { 287 out, err := tangled.RepoListCollaborators(ctx, xc, cursor, collaboratorPageLimit, "", repo.String()) 288 if err != nil { 289 return nil, err 290 } 291 for _, item := range out.Items { 292 did, err := syntax.ParseDID(item.Subject) 293 if err != nil { 294 continue 295 } 296 subjects[did] = struct{}{} 297 } 298 if out.Cursor == nil || *out.Cursor == "" { 299 return subjects, nil 300 } 301 cursor = *out.Cursor 302 } 303 return nil, fmt.Errorf("collaborator roster exceeded %d pages", maxCollaboratorPages) 304} 305 306func (t *Tap) teardownRepo(l *slog.Logger, repo syntax.DID) error { 307 if repo == "" { 308 return nil 309 } 310 if err := t.spindle.db.DeleteRepo(repo); err != nil { 311 l.Error("failed to remove repo", "err", err) 312 return fmt.Errorf("remove repo: %w", err) 313 } 314 if err := t.spindle.e.DeleteRepo(repo); err != nil { 315 l.Error("failed to remove repo policy", "err", err) 316 return fmt.Errorf("remove repo policy: %w", err) 317 } 318 // TODO: clear sparse-synced git repo 319 return nil 320} 321 322func (s *Spindle) processPull(ctx context.Context, evt *tapc.RecordEventData) error { 323 l := s.l.With("component", "ingester", "collection", evt.Collection, "did", evt.Did, "rkey", evt.Rkey) 324 325 // only listen to live events 326 if !evt.Live { 327 l.Info("skipping backfill event", "event", evt.AtUri()) 328 return nil 329 } 330 331 switch evt.Action { 332 case tapc.RecordCreateAction, tapc.RecordUpdateAction: 333 record := tangled.RepoPull{} 334 if err := json.Unmarshal(evt.Record, &record); err != nil { 335 l.Error("invalid record", "err", err) 336 return fmt.Errorf("parsing record: %w", err) 337 } 338 339 // ignore legacy records 340 if record.Target == nil { 341 l.Info("ignoring pull record: target repo is nil") 342 return nil 343 } 344 345 // ignore patch-based and fork-based PRs 346 if record.Source == nil || record.Source.Repo != nil { 347 l.Info("ignoring pull record: not a branch-based pull request") 348 return nil 349 } 350 351 // skip if target repo is unknown 352 repo, err := s.db.GetRepoByDid(syntax.DID(record.Target.Repo)) 353 if err != nil { 354 l.Warn("target repo is not ingested yet", "repo", record.Target.Repo, "err", err) 355 return fmt.Errorf("target repo is unknown") 356 } 357 358 // only accept branch-based PR (excluding patch-based and fork-based) 359 if record.Source == nil || record.Source.Repo != nil { 360 l.Warn("skipping non-branch-based PR") 361 return nil 362 } 363 364 // check if pull record author can trigger CI in target repo 365 allowed, err := s.e.IsRepoCiTriggerAllowed(evt.Did, repo.RepoDid) 366 if err != nil { 367 return fmt.Errorf("checking push access for pull record author: %w", err) 368 } 369 if !allowed { 370 l.Warn("rejecting pull-triggered pipeline. author has no push access", 371 "author", evt.Did, "repo", repo.RepoDid) 372 return nil 373 } 374 375 latestSubmission, err := s.fetchLatestSubmission(ctx, evt.Did.String(), evt.Rkey.String(), &record) 376 if err != nil { 377 return err 378 } 379 sourceSha := latestSubmission.SourceRev 380 381 scheme := "https" 382 if s.cfg.Server.Dev { 383 scheme = "http" 384 } 385 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, repo.Knot)} 386 387 // fetch current default branch 388 defaultBranch, _ := func(repo syntax.DID) (string, error) { 389 defaultBranchOut, err := tangled.RepoGetDefaultBranch(ctx, client, repo.String()) 390 if err != nil { 391 return "", err 392 } 393 return defaultBranchOut.Name, nil 394 }(repo.RepoDid) 395 396 compiler := workflow.Compiler{ 397 Trigger: tangled.Pipeline_TriggerMetadata{ 398 Kind: string(workflow.TriggerKindPullRequest), 399 PullRequest: &tangled.Pipeline_PullRequestTriggerData{ 400 SourceBranch: record.Source.Branch, 401 SourceSha: sourceSha, 402 TargetBranch: record.Target.Branch, 403 }, 404 Repo: &tangled.Pipeline_TriggerRepo{ 405 Did: repo.Owner.String(), 406 Knot: repo.Knot, 407 Repo: (*string)(&repo.Rkey), 408 RepoDid: (*string)(&repo.RepoDid), 409 DefaultBranch: defaultBranch, 410 }, 411 }, 412 } 413 414 repoUri := s.newRepoCloneUrl(repo.Knot, repo.RepoDid) 415 repoPath := s.newRepoPath(repo.RepoDid) 416 417 // load workflow definitions from rev (without spindle context) 418 rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha) 419 if err != nil { 420 // don't retry 421 l.Error("failed loading pipeline", "err", err) 422 return nil 423 } 424 if len(rawPipeline) == 0 { 425 l.Info("no workflow definition find for the repo. skipping the event") 426 return nil 427 } 428 tpl := compiler.Compile(compiler.Parse(rawPipeline)) 429 // TODO: pass compile error to workflow log 430 for _, w := range compiler.Diagnostics.Errors { 431 l.Error(w.String()) 432 } 433 for _, w := range compiler.Diagnostics.Warnings { 434 l.Warn(w.String()) 435 } 436 if len(tpl.Workflows) == 0 { 437 l.Info("no workflow matching trigger 'pull_request'. skipping the event") 438 return nil 439 } 440 441 pipelineId := models.PipelineId{ 442 Knot: tpl.TriggerMetadata.Repo.Knot, 443 Rkey: tid.TID(), 444 } 445 if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { 446 l.Error("failed to create pipeline event", "err", err) 447 return nil 448 } 449 sourceRepo, err := s.resolvePipelineSourceRepo(ctx, tpl.TriggerMetadata) 450 if err != nil { 451 l.Error("failed resolving pipeline source repo", "err", err) 452 return nil 453 } 454 err = s.processPipeline(repo.RepoDid, tpl, pipelineId, sourceRepo) 455 if err != nil { 456 // don't retry 457 l.Error("failed processing pipeline", "err", err) 458 return nil 459 } 460 case tapc.RecordDeleteAction: 461 // no-op 462 } 463 return nil 464} 465 466func (s *Spindle) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) { 467 // resolve the PR owner's identity to fetch the blob from their PDS 468 prOwnerIdent, err := s.res.ResolveIdent(ctx, did) 469 if err != nil || prOwnerIdent.Handle.IsInvalidHandle() { 470 return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err) 471 } 472 473 if len(record.Rounds) == 0 { 474 return nil, fmt.Errorf("failed to fetch latest submission, no rounds in record") 475 } 476 477 roundNumber := len(record.Rounds) - 1 478 round := record.Rounds[roundNumber] 479 480 // fetch the blob from the PR owner's PDS 481 prOwnerPds := prOwnerIdent.PDSEndpoint() 482 blobUrl, err := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", prOwnerPds)) 483 if err != nil { 484 return nil, fmt.Errorf("failed to construct blob URL: %w", err) 485 } 486 q := blobUrl.Query() 487 q.Set("cid", round.PatchBlob.Ref.String()) 488 q.Set("did", did) 489 blobUrl.RawQuery = q.Encode() 490 491 req, err := http.NewRequestWithContext(ctx, http.MethodGet, blobUrl.String(), nil) 492 if err != nil { 493 return nil, fmt.Errorf("failed to create blob request: %w", err) 494 } 495 req.Header.Set("Content-Type", "application/json") 496 497 blobResp, err := http.DefaultClient.Do(req) 498 if err != nil { 499 return nil, fmt.Errorf("failed to fetch blob: %w", err) 500 } 501 defer blobResp.Body.Close() 502 503 latestSubmission, err := avmodels.PullSubmissionFromRecord(did, rkey, roundNumber, round, blobResp.Body) 504 if err != nil { 505 return nil, fmt.Errorf("failed to parse submission: %w", err) 506 } 507 508 return latestSubmission, nil 509}