package spindle import ( "context" "encoding/json" "fmt" "log/slog" "net/http" "net/url" "slices" "strings" "time" "github.com/bluesky-social/indigo/atproto/syntax" indigoxrpc "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" avmodels "tangled.org/core/appview/models" "tangled.org/core/eventconsumer" "tangled.org/core/log" "tangled.org/core/spindle/db" "tangled.org/core/spindle/git" "tangled.org/core/spindle/models" "tangled.org/core/tapc" "tangled.org/core/tid" "tangled.org/core/workflow" ) const ( collaboratorPageLimit = 100 maxCollaboratorPages = 50 ) type Tap struct { logger *slog.Logger spindle *Spindle tap tapc.Client } func NewTapClient(s *Spindle) *Tap { return &Tap{ logger: log.SubLogger(s.l, "tapclient"), spindle: s, tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword), } } func (t *Tap) Start(connCtx context.Context) { go t.tap.Connect(connCtx, &tapc.SimpleIndexer{ EventHandler: t.processEvent, ConnectHandler: t.onConnect, }) } func (t *Tap) onConnect(ctx context.Context) { l := t.logger if t.spindle.cfg.Server.InviteOnly { // listen to owners of registered repositories owners, err := t.spindle.db.RepoOwners() if err != nil { l.Warn("tap declare: failed to load known repos", "err", err) return } if err := t.tap.AddRepos(ctx, owners); err != nil { l.Warn("tap declare: AddRepos rejected", "count", len(owners), "err", err) return } l.Info("tap declare: known owner DIDs registered", "count", len(owners)) } else { // public spindle. listen to full network l.Info("tap declare: listening to full network") } } func (t *Tap) processEvent(ctx context.Context, evt tapc.Event) error { if evt.Type != tapc.EvtRecord || evt.Record == nil { return nil } switch evt.Record.Collection.String() { case tangled.RepoNSID: return t.processRepo(ctx, evt.Record) } return nil } // processRepo ingests repo declaration record: `sh.tangled.repo`. It skips alias records. func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error { l := t.logger.With("collection", tangled.RepoNSID, "did", evt.Did, "rkey", evt.Rkey) ownerDid := evt.Did rkey := evt.Rkey switch evt.Action { case tapc.RecordCreateAction, tapc.RecordUpdateAction: record := tangled.Repo{} if err := json.Unmarshal(evt.Record, &record); err != nil { l.Warn("skipping invalid repo record", "err", err) return nil } if record.RepoDid == nil || *record.RepoDid == "" { l.Warn("skipping repo record without repoDid") return nil } repoDid, err := syntax.ParseDID(*record.RepoDid) if err != nil { l.Warn("skipping repo record with malformed repoDid", "value", *record.RepoDid, "err", err) return nil } isMember, err := t.spindle.db.IsAllowedMember(ctx, ownerDid, !t.spindle.cfg.Server.InviteOnly) if err != nil { return fmt.Errorf("checking spindle membership: %w", err) } if !isMember { l.Warn("rejecting repo record: owner is not a spindle member", "owner", ownerDid) return nil } // ignore repos not pointing this spindle. hostname := t.spindle.cfg.Server.Hostname if record.Spindle == nil || *record.Spindle != hostname { // teardown existing repo prior, err := t.spindle.db.GetRepoByDid(repoDid) if err != nil { return nil } if prior.Owner == ownerDid && prior.Rkey == rkey { l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle) return t.teardownRepo(l, prior.RepoDid) } l.Warn("ignoring reassignment from non-registering record", "owner", prior.Owner, "rkey", prior.Rkey) return nil } // verify repo declaration // NOTE: we are ignoring repo declaration records that *might* become correct pointer in // future with repository rename or ownership transfer. On repo rename, new pointer record // should always be recreated even when it already exists. verified, err := t.verifyRepoDeclaration(ctx, l, repoDid, ownerDid, rkey.String(), record.Knot) if err != nil { l.Warn("failed to verify repo declaration", "err", err) return nil } if !verified { // ignore alias records return nil } if err := t.spindle.e.SetRepoOwner(ownerDid, repoDid); err != nil { l.Error("failed to add repo policy", "err", err) return fmt.Errorf("add repo policy: %w", err) } repo := db.Repo{ Knot: record.Knot, Owner: ownerDid, Rkey: rkey, RepoDid: repoDid, CreatedAt: record.CreatedAt, } if err := t.spindle.db.UpsertRepo(repo); err != nil { l.Error("failed to add repo row", "err", err) return fmt.Errorf("add repo: %w", err) } t.reconcileCollaborators(ctx, l, record.Knot, repoDid, ownerDid) t.spindle.ks.AddSource(t.spindle.rootCtx, eventconsumer.NewKnotSource(record.Knot)) // setup sparse sync repoCloneUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid) repoPath := t.spindle.newRepoPath(repo.RepoDid) if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, ""); err != nil { return fmt.Errorf("setting up sparse-clone git repo: %w", err) } legacyName := "" if record.Name != nil { legacyName = *record.Name } migrateLegacyRepoSecrets(ctx, t.spindle.db, t.spindle.vault, l, ownerDid, legacyName, rkey, repoDid) if e := t.spindle.embedTap; e == nil || !e.closed.Load() { if err := t.tap.AddRepos(ctx, []syntax.DID{ownerDid}); err != nil { l.Warn("tap AddRepos rejected", "did", ownerDid, "err", err) } } t.spindle.jc.AddDid(ownerDid.String()) case tapc.RecordDeleteAction: repoDid, err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey) if err != nil { return fmt.Errorf("deleting repo record: %w", err) } if repoDid == "" { // no record is deleted. record was not pointing this spindle or was an alias record return nil } return t.teardownRepo(l, repoDid) } return nil } func (t *Tap) verifyRepoDeclaration(ctx context.Context, l *slog.Logger, repo, owner syntax.DID, rkey, knot string) (bool, error) { result, err := t.spindle.VerifyRepo(ctx, repo) if err != nil { return false, err } l = l.With( "repoDid", result.RepoDid, "repoOwner", result.OwnerDid, "repoKnot", result.KnotURL.Host, "claimedOwner", owner, "claimedRkey", rkey, "claimedKnot", knot, ) if syntax.DID(result.OwnerDid) != owner { l.Warn("rejecting repo event: owner mismatch") return false, nil } if result.Rkey != rkey { l.Warn("rejecting repo event: rkey mismatch") return false, nil } if !strings.EqualFold(knot, result.KnotURL.Host) { l.Warn("rejecting repo event: record knot does not match DID-doc endpoint") return false, nil } return true, nil } func (t *Tap) reconcileCollaborators(ctx context.Context, l *slog.Logger, knot string, repo, owner syntax.DID) { wanted, err := t.fetchKnotCollaborators(ctx, knot, repo) if err != nil { l.Warn("collaborator reconcile: failed to fetch roster from knot", "knot", knot, "err", err) return } have, err := t.spindle.e.GetRepoCollaborators(repo) if err != nil { l.Warn("collaborator reconcile: failed to read current grants", "err", err) return } // the owner is a collaborator by role inheritance, never by an explicit grant delete(wanted, owner) for did := range wanted { if slices.Contains(have, did) { continue } if err := t.spindle.grantCollaborator(did, repo); err != nil { l.Error("collaborator reconcile: failed to add", "subject", did, "err", err) return } l.Info("collaborator reconcile: added", "subject", did) } for _, did := range have { if did == owner { continue } if _, ok := wanted[did]; ok { continue } if err := t.spindle.e.RemoveRepoCollaborator(did, repo); err != nil { l.Error("collaborator reconcile: failed to remove", "subject", did, "err", err) return } l.Info("collaborator reconcile: removed", "subject", did) } } func (t *Tap) fetchKnotCollaborators(ctx context.Context, knot string, repo syntax.DID) (map[syntax.DID]struct{}, error) { scheme := "https" if t.spindle.cfg.Server.Dev { scheme = "http" } xc := &indigoxrpc.Client{ Host: fmt.Sprintf("%s://%s", scheme, knot), Client: &http.Client{Timeout: 30 * time.Second}, } subjects := make(map[syntax.DID]struct{}) cursor := "" for range maxCollaboratorPages { out, err := tangled.RepoListCollaborators(ctx, xc, cursor, collaboratorPageLimit, "", repo.String()) if err != nil { return nil, err } for _, item := range out.Items { did, err := syntax.ParseDID(item.Subject) if err != nil { continue } subjects[did] = struct{}{} } if out.Cursor == nil || *out.Cursor == "" { return subjects, nil } cursor = *out.Cursor } return nil, fmt.Errorf("collaborator roster exceeded %d pages", maxCollaboratorPages) } func (t *Tap) teardownRepo(l *slog.Logger, repo syntax.DID) error { if repo == "" { return nil } if err := t.spindle.db.DeleteRepo(repo); err != nil { l.Error("failed to remove repo", "err", err) return fmt.Errorf("remove repo: %w", err) } if err := t.spindle.e.DeleteRepo(repo); err != nil { l.Error("failed to remove repo policy", "err", err) return fmt.Errorf("remove repo policy: %w", err) } // TODO: clear sparse-synced git repo return nil } func (s *Spindle) processPull(ctx context.Context, evt *tapc.RecordEventData) error { l := s.l.With("component", "ingester", "collection", evt.Collection, "did", evt.Did, "rkey", evt.Rkey) // only listen to live events if !evt.Live { l.Info("skipping backfill event", "event", evt.AtUri()) return nil } switch evt.Action { case tapc.RecordCreateAction, tapc.RecordUpdateAction: record := tangled.RepoPull{} if err := json.Unmarshal(evt.Record, &record); err != nil { l.Error("invalid record", "err", err) return fmt.Errorf("parsing record: %w", err) } // ignore legacy records if record.Target == nil { l.Info("ignoring pull record: target repo is nil") return nil } // ignore patch-based and fork-based PRs if record.Source == nil || record.Source.Repo != nil { l.Info("ignoring pull record: not a branch-based pull request") return nil } // skip if target repo is unknown repo, err := s.db.GetRepoByDid(syntax.DID(record.Target.Repo)) if err != nil { l.Warn("target repo is not ingested yet", "repo", record.Target.Repo, "err", err) return fmt.Errorf("target repo is unknown") } // only accept branch-based PR (excluding patch-based and fork-based) if record.Source == nil || record.Source.Repo != nil { l.Warn("skipping non-branch-based PR") return nil } // check if pull record author can trigger CI in target repo allowed, err := s.e.IsRepoCiTriggerAllowed(evt.Did, repo.RepoDid) if err != nil { return fmt.Errorf("checking push access for pull record author: %w", err) } if !allowed { l.Warn("rejecting pull-triggered pipeline. author has no push access", "author", evt.Did, "repo", repo.RepoDid) return nil } latestSubmission, err := s.fetchLatestSubmission(ctx, evt.Did.String(), evt.Rkey.String(), &record) if err != nil { return err } sourceSha := latestSubmission.SourceRev scheme := "https" if s.cfg.Server.Dev { scheme = "http" } client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, repo.Knot)} // fetch current default branch defaultBranch, _ := func(repo syntax.DID) (string, error) { defaultBranchOut, err := tangled.RepoGetDefaultBranch(ctx, client, repo.String()) if err != nil { return "", err } return defaultBranchOut.Name, nil }(repo.RepoDid) compiler := workflow.Compiler{ Trigger: tangled.Pipeline_TriggerMetadata{ Kind: string(workflow.TriggerKindPullRequest), PullRequest: &tangled.Pipeline_PullRequestTriggerData{ SourceBranch: record.Source.Branch, SourceSha: sourceSha, TargetBranch: record.Target.Branch, }, Repo: &tangled.Pipeline_TriggerRepo{ Did: repo.Owner.String(), Knot: repo.Knot, Repo: (*string)(&repo.Rkey), RepoDid: (*string)(&repo.RepoDid), DefaultBranch: defaultBranch, }, }, } repoUri := s.newRepoCloneUrl(repo.Knot, repo.RepoDid) repoPath := s.newRepoPath(repo.RepoDid) // load workflow definitions from rev (without spindle context) rawPipeline, err := s.loadPipeline(ctx, repoUri, repoPath, sourceSha) if err != nil { // don't retry l.Error("failed loading pipeline", "err", err) return nil } if len(rawPipeline) == 0 { l.Info("no workflow definition find for the repo. skipping the event") return nil } tpl := compiler.Compile(compiler.Parse(rawPipeline)) // TODO: pass compile error to workflow log for _, w := range compiler.Diagnostics.Errors { l.Error(w.String()) } for _, w := range compiler.Diagnostics.Warnings { l.Warn(w.String()) } if len(tpl.Workflows) == 0 { l.Info("no workflow matching trigger 'pull_request'. skipping the event") return nil } pipelineId := models.PipelineId{ Knot: tpl.TriggerMetadata.Repo.Knot, Rkey: tid.TID(), } if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { l.Error("failed to create pipeline event", "err", err) return nil } sourceRepo, err := s.resolvePipelineSourceRepo(ctx, tpl.TriggerMetadata) if err != nil { l.Error("failed resolving pipeline source repo", "err", err) return nil } err = s.processPipeline(repo.RepoDid, tpl, pipelineId, sourceRepo) if err != nil { // don't retry l.Error("failed processing pipeline", "err", err) return nil } case tapc.RecordDeleteAction: // no-op } return nil } func (s *Spindle) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*avmodels.PullSubmission, error) { // resolve the PR owner's identity to fetch the blob from their PDS prOwnerIdent, err := s.res.ResolveIdent(ctx, did) if err != nil || prOwnerIdent.Handle.IsInvalidHandle() { return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err) } if len(record.Rounds) == 0 { return nil, fmt.Errorf("failed to fetch latest submission, no rounds in record") } roundNumber := len(record.Rounds) - 1 round := record.Rounds[roundNumber] // fetch the blob from the PR owner's PDS prOwnerPds := prOwnerIdent.PDSEndpoint() blobUrl, err := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", prOwnerPds)) if err != nil { return nil, fmt.Errorf("failed to construct blob URL: %w", err) } q := blobUrl.Query() q.Set("cid", round.PatchBlob.Ref.String()) q.Set("did", did) blobUrl.RawQuery = q.Encode() req, err := http.NewRequestWithContext(ctx, http.MethodGet, blobUrl.String(), nil) if err != nil { return nil, fmt.Errorf("failed to create blob request: %w", err) } req.Header.Set("Content-Type", "application/json") blobResp, err := http.DefaultClient.Do(req) if err != nil { return nil, fmt.Errorf("failed to fetch blob: %w", err) } defer blobResp.Body.Close() latestSubmission, err := avmodels.PullSubmissionFromRecord(did, rkey, roundNumber, round, blobResp.Body) if err != nil { return nil, fmt.Errorf("failed to parse submission: %w", err) } return latestSubmission, nil }