diff --git a/appview/middleware/middleware.go b/appview/middleware/middleware.go index 46e0f307..1583fad8 100644 --- a/appview/middleware/middleware.go +++ b/appview/middleware/middleware.go @@ -252,7 +252,7 @@ func (mw Middleware) ResolvePull() middlewareFunc { return } - pr, err := db.GetPull(mw.db, f.RepoAt(), prIdInt) + pr, err := db.GetPull(mw.db, orm.FilterEq("repo_at", f.RepoAt()), orm.FilterEq("pull_id", prIdInt)) if err != nil { log.Println("failed to get pull and comments", err) mw.pages.Error404(w) @@ -261,22 +261,13 @@ func (mw Middleware) ResolvePull() middlewareFunc { ctx := context.WithValue(r.Context(), "pull", pr) - if pr.IsStacked() { - stack, err := db.GetStack(mw.db, pr.StackId) - if err != nil { - log.Println("failed to get stack", err) - return - } - abandonedPulls, err := db.GetAbandonedPulls(mw.db, pr.StackId) - if err != nil { - log.Println("failed to get abandoned pulls", err) - return - } - - ctx = context.WithValue(ctx, "stack", stack) - ctx = context.WithValue(ctx, "abandonedPulls", abandonedPulls) + stack, err := db.GetStack(mw.db, pr.AtUri()) + if err != nil { + log.Println("failed to get stack", err) } + ctx = context.WithValue(ctx, "stack", stack) + next.ServeHTTP(w, r.WithContext(ctx)) }) } diff --git a/appview/notify/db/db.go b/appview/notify/db/db.go index 429cba56..10ebef72 100644 --- a/appview/notify/db/db.go +++ b/appview/notify/db/db.go @@ -282,8 +282,8 @@ func (n *databaseNotifier) NewPullComment(ctx context.Context, comment *models.P l := log.FromContext(ctx) pull, err := db.GetPull(n.db, - syntax.ATURI(comment.RepoAt), - comment.PullId, + orm.FilterEq("repo_at", syntax.ATURI(comment.RepoAt)), + orm.FilterEq("pull_id", comment.PullId), ) if err != nil { l.Error("failed to get pulls", "err", err) diff --git a/appview/pulls/pulls.go b/appview/pulls/pulls.go index b708f750..08e11fb0 100644 --- a/appview/pulls/pulls.go +++ b/appview/pulls/pulls.go @@ -9,6 +9,7 @@ import ( "errors" "fmt" "io" + "log" "log/slog" "net/http" "slices" @@ -47,7 +48,6 @@ import ( lexutil "github.com/bluesky-social/indigo/lex/util" indigoxrpc "github.com/bluesky-social/indigo/xrpc" "github.com/go-chi/chi/v5" - "github.com/google/uuid" ) const ApplicationGzip = "application/gzip" @@ -106,13 +106,13 @@ func (s *Pulls) PullActions(w http.ResponseWriter, r *http.Request) { user := s.oauth.GetMultiAccountUser(r) f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -127,7 +127,7 @@ func (s *Pulls) PullActions(w http.ResponseWriter, r *http.Request) { } if roundNumber >= len(pull.Submissions) { http.Error(w, "bad round id", http.StatusBadRequest) - s.logger.Error("failed to parse round id", "err", err) + log.Println("failed to parse round id", err) return } @@ -156,20 +156,20 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff user := s.oauth.GetMultiAccountUser(r) f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } backlinks, err := db.GetBacklinks(s.db, pull.AtUri()) if err != nil { - s.logger.Error("failed to get pull backlinks", "err", err) + log.Println("failed to get pull backlinks", err) s.pages.Notice(w, "pull-error", "Failed to get pull. Try again later.") return } @@ -181,7 +181,7 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff } if roundIdInt >= len(pull.Submissions) { http.Error(w, "bad round id", http.StatusBadRequest) - s.logger.Error("failed to parse round id", "err", err) + log.Println("failed to parse round id", err) return } @@ -192,7 +192,6 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff // can be nil if this pull is not stacked stack, _ := r.Context().Value("stack").(models.Stack) - abandonedPulls, _ := r.Context().Value("abandonedPulls").([]*models.Pull) mergeCheckResponse := s.mergeCheck(r, f, pull, stack) branchDeleteStatus := s.branchDeleteStatus(r, f, pull) @@ -210,9 +209,6 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff for _, p := range stack { shas = append(shas, p.LatestSha()) } - for _, p := range abandonedPulls { - shas = append(shas, p.LatestSha()) - } ps, err := db.GetPipelineStatuses( s.db, @@ -223,7 +219,7 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff orm.FilterIn("p.sha", shas), ) if err != nil { - s.logger.Error("failed to fetch pipeline statuses", "err", err) + log.Printf("failed to fetch pipeline statuses: %s", err) // non-fatal } @@ -233,7 +229,7 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff reactionMap, err := db.GetReactionMap(s.db, 20, pull.AtUri()) if err != nil { - s.logger.Error("failed to get pull reactions", "err", err) + log.Println("failed to get pull reactions") } userReactions := map[models.ReactionKind]bool{} @@ -247,7 +243,7 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff orm.FilterContains("scope", tangled.RepoPullNSID), ) if err != nil { - s.logger.Error("failed to fetch labels", "err", err) + log.Println("failed to fetch labels", err) s.pages.Error503(w) return } @@ -264,14 +260,14 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff if interdiff { currentPatch, err := patchutil.AsDiff(pull.Submissions[roundIdInt].CombinedPatch()) if err != nil { - s.logger.Error("failed to interdiff; current patch malformed", "err", err) + log.Println("failed to interdiff; current patch malformed") s.pages.Notice(w, fmt.Sprintf("interdiff-error-%d", roundIdInt), "Failed to calculate interdiff; current patch is invalid.") return } previousPatch, err := patchutil.AsDiff(pull.Submissions[roundIdInt-1].CombinedPatch()) if err != nil { - s.logger.Error("failed to interdiff; previous patch malformed", "err", err) + log.Println("failed to interdiff; previous patch malformed") s.pages.Notice(w, fmt.Sprintf("interdiff-error-%d", roundIdInt), "Failed to calculate interdiff; previous patch is invalid.") return } @@ -284,7 +280,6 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff RepoInfo: s.repoResolver.GetRepoInfo(r, user), Pull: pull, Stack: stack, - AbandonedPulls: abandonedPulls, Backlinks: backlinks, BranchDeleteStatus: branchDeleteStatus, MergeCheck: mergeCheckResponse, @@ -305,7 +300,7 @@ func (s *Pulls) repoPullHelper(w http.ResponseWriter, r *http.Request, interdiff func (s *Pulls) RepoSinglePull(w http.ResponseWriter, r *http.Request) { pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -328,15 +323,12 @@ func (s *Pulls) mergeCheck(r *http.Request, f *models.Repo, pull *models.Pull, s Host: host, } - patch := pull.LatestPatch() - if pull.IsStacked() { - // combine patches of substack - subStack := stack.Below(pull) - // collect the portion of the stack that is mergeable - mergeable := subStack.Mergeable() - // combine each patch - patch = mergeable.CombinedPatch() - } + // combine patches of substack + subStack := stack.Below(pull) + // collect the portion of the stack that is mergeable + mergeable := subStack.Mergeable() + // combine each patch + patch := mergeable.CombinedPatch() resp, xe := tangled.RepoMergeCheck( r.Context(), @@ -349,7 +341,7 @@ func (s *Pulls) mergeCheck(r *http.Request, f *models.Repo, pull *models.Pull, s }, ) if err := xrpcclient.HandleXrpcErr(xe); err != nil { - s.logger.Error("failed to check for mergeability", "err", err) + log.Println("failed to check for mergeability", "err", err) return types.MergeCheckResponse{ Error: fmt.Sprintf("failed to check merge status: %s", err.Error()), } @@ -426,7 +418,7 @@ func (s *Pulls) branchDeleteStatus(r *http.Request, repo *models.Repo, pull *mod } func (s *Pulls) resubmitCheck(r *http.Request, repo *models.Repo, pull *models.Pull, stack models.Stack) pages.ResubmitResult { - if pull.State == models.PullMerged || pull.State == models.PullDeleted || pull.PullSource == nil { + if pull.State == models.PullMerged || pull.State == models.PullAbandoned || pull.PullSource == nil { return pages.Unknown } @@ -443,21 +435,17 @@ func (s *Pulls) resubmitCheck(r *http.Request, repo *models.Repo, pull *models.P branchResp, err := tangled.GitTempGetBranch(r.Context(), xrpcc, pull.PullSource.Branch, sourceRepo.String()) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.branches", "err", xrpcerr) + log.Println("failed to call XRPC repo.branches", xrpcerr) return pages.Unknown } - s.logger.Error("failed to reach knotserver", "err", err) + log.Println("failed to reach knotserver", err) return pages.Unknown } targetBranch := branchResp - latestSourceRev := pull.LatestSha() - - if pull.IsStacked() && stack != nil { - top := stack[0] - latestSourceRev = top.LatestSha() - } + top := stack[0] + latestSourceRev := top.LatestSha() if latestSourceRev != targetBranch.Hash { return pages.ShouldResubmit @@ -477,7 +465,7 @@ func (s *Pulls) RepoPullInterdiff(w http.ResponseWriter, r *http.Request) { func (s *Pulls) RepoPullPatchRaw(w http.ResponseWriter, r *http.Request) { pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -486,7 +474,7 @@ func (s *Pulls) RepoPullPatchRaw(w http.ResponseWriter, r *http.Request) { roundIdInt, err := strconv.Atoi(roundId) if err != nil || roundIdInt >= len(pull.Submissions) { http.Error(w, "bad round id", http.StatusBadRequest) - s.logger.Error("failed to parse round id", "err", err) + log.Println("failed to parse round id", err) return } @@ -503,7 +491,7 @@ func (s *Pulls) RepoPulls(w http.ResponseWriter, r *http.Request) { f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } @@ -620,7 +608,6 @@ func (s *Pulls) RepoPulls(w http.ResponseWriter, r *http.Request) { countOpts := searchOpts countOpts.Page = pagination.Page{Limit: 1} for _, ps := range []models.PullState{models.PullOpen, models.PullMerged, models.PullClosed} { - ps := ps countOpts.State = &ps countRes, err := s.indexer.Search(r.Context(), countOpts) if err != nil { @@ -672,7 +659,7 @@ func (s *Pulls) RepoPulls(w http.ResponseWriter, r *http.Request) { if p.PullSource.RepoAt != nil { pullSourceRepo, err = db.GetRepoByAtUri(s.db, p.PullSource.RepoAt.String()) if err != nil { - s.logger.Error("failed to get repo by at uri", "err", err) + log.Printf("failed to get repo by at uri: %v", err) continue } else { p.PullSource.Repo = pullSourceRepo @@ -681,30 +668,59 @@ func (s *Pulls) RepoPulls(w http.ResponseWriter, r *http.Request) { } } - // we want to group all stacked PRs into just one list - stacks := make(map[string]models.Stack) + var stacks []models.Stack var shas []string - n := 0 + + pullMap := make(map[string]*models.Pull) for _, p := range pulls { - // store the sha for later shas = append(shas, p.LatestSha()) - // this PR is stacked - if p.StackId != "" { - // we have already seen this PR stack - if _, seen := stacks[p.StackId]; seen { - stacks[p.StackId] = append(stacks[p.StackId], p) - // skip this PR + pullMap[p.AtUri().String()] = p + } + + // track which PRs have been added to stacks + visited := make(map[string]bool) + + // group stacked PRs together using dependent_on relationships + for _, p := range pulls { + if visited[p.AtUri().String()] { + continue + } + + root := p + for root.DependentOn != nil { + if parent, ok := pullMap[root.DependentOn.String()]; ok { + root = parent } else { - stacks[p.StackId] = nil - pulls[n] = p - n++ + break // parent not in current page + } + } + + var stack models.Stack + current := root + for { + if visited[current.AtUri().String()] { + break + } + stack = append(stack, current) + visited[current.AtUri().String()] = true + + found := false + for _, candidate := range pulls { + if candidate.DependentOn != nil && + candidate.DependentOn.String() == current.AtUri().String() { + current = candidate + found = true + break + } + } + if !found { + break } - } else { - pulls[n] = p - n++ } + + slices.Reverse(stack) + stacks = append(stacks, stack) } - pulls = pulls[:n] ps, err := db.GetPipelineStatuses( s.db, @@ -715,7 +731,7 @@ func (s *Pulls) RepoPulls(w http.ResponseWriter, r *http.Request) { orm.FilterIn("p.sha", shas), ) if err != nil { - s.logger.Error("failed to fetch pipeline statuses", "err", err) + log.Printf("failed to fetch pipeline statuses: %s", err) // non-fatal } m := make(map[string]models.Pipeline) @@ -762,13 +778,13 @@ func (s *Pulls) PullComment(w http.ResponseWriter, r *http.Request) { user := s.oauth.GetMultiAccountUser(r) f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -777,7 +793,7 @@ func (s *Pulls) PullComment(w http.ResponseWriter, r *http.Request) { roundNumber, err := strconv.Atoi(roundNumberStr) if err != nil || roundNumber >= len(pull.Submissions) { http.Error(w, "bad round id", http.StatusBadRequest) - s.logger.Error("failed to parse round id", "err", err) + log.Println("failed to parse round id", err) return } @@ -802,7 +818,7 @@ func (s *Pulls) PullComment(w http.ResponseWriter, r *http.Request) { // Start a transaction tx, err := s.db.BeginTx(r.Context(), nil) if err != nil { - s.logger.Error("failed to start transaction", "err", err) + log.Println("failed to start transaction", err) s.pages.Notice(w, "pull-comment", "Failed to create comment.") return } @@ -812,7 +828,7 @@ func (s *Pulls) PullComment(w http.ResponseWriter, r *http.Request) { client, err := s.oauth.AuthorizedClient(r) if err != nil { - s.logger.Error("failed to get authorized client", "err", err) + log.Println("failed to get authorized client", err) s.pages.Notice(w, "pull-comment", "Failed to create comment.") return } @@ -829,7 +845,7 @@ func (s *Pulls) PullComment(w http.ResponseWriter, r *http.Request) { }, }) if err != nil { - s.logger.Error("failed to create pull comment", "err", err) + log.Println("failed to create pull comment", err) s.pages.Notice(w, "pull-comment", "Failed to create comment.") return } @@ -848,14 +864,14 @@ func (s *Pulls) PullComment(w http.ResponseWriter, r *http.Request) { // Create the pull comment in the database with the commentAt field commentId, err := db.NewPullComment(tx, comment) if err != nil { - s.logger.Error("failed to create pull comment", "err", err) + log.Println("failed to create pull comment", err) s.pages.Notice(w, "pull-comment", "Failed to create comment.") return } // Commit the transaction if err = tx.Commit(); err != nil { - s.logger.Error("failed to commit transaction", "err", err) + log.Println("failed to commit transaction", err) s.pages.Notice(w, "pull-comment", "Failed to create comment.") return } @@ -872,7 +888,7 @@ func (s *Pulls) NewPull(w http.ResponseWriter, r *http.Request) { user := s.oauth.GetMultiAccountUser(r) f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } @@ -883,17 +899,17 @@ func (s *Pulls) NewPull(w http.ResponseWriter, r *http.Request) { xrpcBytes, err := tangled.GitTempListBranches(r.Context(), xrpcc, "", 0, f.RepoAt().String()) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.branches", "err", xrpcerr) + log.Println("failed to call XRPC repo.branches", xrpcerr) s.pages.Error503(w) return } - s.logger.Error("failed to fetch branches", "err", err) + log.Println("failed to fetch branches", err) return } var result types.RepoBranchesResponse if err := json.Unmarshal(xrpcBytes, &result); err != nil { - s.logger.Error("failed to decode XRPC response", "err", err) + log.Println("failed to decode XRPC response", err) s.pages.Error503(w) return } @@ -1049,18 +1065,18 @@ func (s *Pulls) handleBranchBasedPull( xrpcBytes, err := tangled.RepoCompare(r.Context(), xrpcc, didSlashRepo, targetBranch, sourceBranch) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.compare", "err", xrpcerr) + log.Println("failed to call XRPC repo.compare", xrpcerr) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } - s.logger.Error("failed to compare", "err", err) + log.Println("failed to compare", err) s.pages.Notice(w, "pull", err.Error()) return } var comparison types.RepoFormatPatchResponse if err := json.Unmarshal(xrpcBytes, &comparison); err != nil { - s.logger.Error("failed to decode XRPC compare response", "err", err) + log.Println("failed to decode XRPC compare response", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } @@ -1080,7 +1096,6 @@ func (s *Pulls) handleBranchBasedPull( } recordPullSource := &tangled.RepoPull_Source{ Branch: sourceBranch, - Sha: comparison.Rev2, } s.createPullRequest(w, r, repo, user, title, body, targetBranch, patch, combined, sourceRev, pullSource, recordPullSource, isStacked) @@ -1105,7 +1120,7 @@ func (s *Pulls) handleForkBasedPull(w http.ResponseWriter, r *http.Request, repo s.pages.Notice(w, "pull", "No such fork.") return } else if err != nil { - s.logger.Error("failed to fetch fork:", "err", err) + log.Println("failed to fetch fork:", err) s.pages.Notice(w, "pull", "Failed to fetch fork.") return } @@ -1159,18 +1174,18 @@ func (s *Pulls) handleForkBasedPull(w http.ResponseWriter, r *http.Request, repo forkXrpcBytes, err := tangled.RepoCompare(r.Context(), forkXrpcc, forkRepoId, hiddenRef, sourceBranch) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.compare for fork", "err", xrpcerr) + log.Println("failed to call XRPC repo.compare for fork", xrpcerr) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } - s.logger.Error("failed to compare across branches", "err", err) + log.Println("failed to compare across branches", err) s.pages.Notice(w, "pull", err.Error()) return } var comparison types.RepoFormatPatchResponse if err := json.Unmarshal(forkXrpcBytes, &comparison); err != nil { - s.logger.Error("failed to decode XRPC compare response for fork", "err", err) + log.Println("failed to decode XRPC compare response for fork", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } @@ -1195,7 +1210,6 @@ func (s *Pulls) handleForkBasedPull(w http.ResponseWriter, r *http.Request, repo recordPullSource := &tangled.RepoPull_Source{ Branch: sourceBranch, Repo: &forkAtUriStr, - Sha: sourceRev, } s.createPullRequest(w, r, repo, user, title, body, targetBranch, patch, combined, sourceRev, pullSource, recordPullSource, isStacked) @@ -1231,14 +1245,14 @@ func (s *Pulls) createPullRequest( client, err := s.oauth.AuthorizedClient(r) if err != nil { - s.logger.Error("failed to get authorized client", "err", err) + log.Println("failed to get authorized client", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } tx, err := s.db.BeginTx(r.Context(), nil) if err != nil { - s.logger.Error("failed to start tx", "err", err) + log.Println("failed to start tx") s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } @@ -1268,10 +1282,35 @@ func (s *Pulls) createPullRequest( mentions, references := s.mentionsResolver.Resolve(r.Context(), body) rkey := tid.TID() + + blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(patch), ApplicationGzip) + if err != nil { + log.Println("failed to upload patch", err) + s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") + return + } + + record := tangled.RepoPull{ + Title: title, + Body: &body, + Target: &tangled.RepoPull_Target{ + Repo: string(repo.RepoAt()), + Branch: targetBranch, + }, + Source: recordPullSource, + CreatedAt: time.Now().Format(time.RFC3339), + Rounds: []*tangled.RepoPull_Round{ + { + CreatedAt: time.Now().Format(time.RFC3339), + PatchBlob: blob.Blob, + }, + }, + } initialSubmission := models.PullSubmission{ Patch: patch, Combined: combined, SourceRev: sourceRev, + Blob: *blob.Blob, } pull := &models.Pull{ Title: title, @@ -1286,52 +1325,38 @@ func (s *Pulls) createPullRequest( &initialSubmission, }, PullSource: pullSource, + State: models.PullOpen, } - err = db.NewPull(tx, pull) - if err != nil { - s.logger.Error("failed to create pull request", "err", err) - s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") - return - } - pullId, err := db.NextPullId(tx, repo.RepoAt()) + + _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ + Collection: tangled.RepoPullNSID, + Repo: user.Active.Did, + Rkey: rkey, + Record: &lexutil.LexiconTypeDecoder{ + Val: &record, + }, + }) if err != nil { - s.logger.Error("failed to get pull id", "err", err) + log.Println("failed to create pull request", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } - blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(patch), ApplicationGzip) + err = db.PutPull(tx, pull) if err != nil { - s.logger.Error("failed to upload patch", "err", err) + log.Println("failed to create pull request", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } - - _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ - Collection: tangled.RepoPullNSID, - Repo: user.Active.Did, - Rkey: rkey, - Record: &lexutil.LexiconTypeDecoder{ - Val: &tangled.RepoPull{ - Title: title, - Target: &tangled.RepoPull_Target{ - Repo: string(repo.RepoAt()), - Branch: targetBranch, - }, - PatchBlob: blob.Blob, - Source: recordPullSource, - CreatedAt: time.Now().Format(time.RFC3339), - }, - }, - }) + pullId, err := db.NextPullId(tx, repo.RepoAt()) if err != nil { - s.logger.Error("failed to create pull request", "err", err) + log.Println("failed to get pull id", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } if err = tx.Commit(); err != nil { - s.logger.Error("failed to create pull request", "err", err) + log.Println("failed to create pull request", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } @@ -1356,7 +1381,7 @@ func (s *Pulls) createStackedPullRequest( // must be branch or fork based if sourceRev == "" { - s.logger.Warn("stacked PR from patch-based pull") + s.logger.Error("stacked PR from patch-based pull") s.pages.Notice(w, "pull", "Stacking is only supported on branch and fork based pull-requests.") return } @@ -1375,15 +1400,6 @@ func (s *Pulls) createStackedPullRequest( return } - // build a stack out of this patch - stackId := uuid.New() - stack, err := s.newStack(r.Context(), repo, user, targetBranch, patch, pullSource, stackId.String()) - if err != nil { - s.logger.Error("failed to create stack", "err", err) - s.pages.Notice(w, "pull", fmt.Sprintf("Failed to create stack: %v", err)) - return - } - client, err := s.oauth.AuthorizedClient(r) if err != nil { s.logger.Error("failed to get authorized client", "err", err) @@ -1391,18 +1407,31 @@ func (s *Pulls) createStackedPullRequest( return } - // apply all record creations at once - var writes []*comatproto.RepoApplyWrites_Input_Writes_Elem - for _, p := range stack { - blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(p.LatestPatch()), ApplicationGzip) + // first upload all blobs + blobs := make([]*lexutil.LexBlob, len(formatPatches)) + for i, p := range formatPatches { + blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(p.Raw), ApplicationGzip) if err != nil { s.logger.Error("failed to upload patch blob", "err", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } + s.logger.Info("uploaded blob", "idx", i+1, "total", len(formatPatches)) + blobs[i] = blob.Blob + } + // build a stack out of this patch + stack, err := s.newStack(r.Context(), repo, user, targetBranch, pullSource, formatPatches, blobs) + if err != nil { + s.logger.Error("failed to create stack", "err", err) + s.pages.Notice(w, "pull", fmt.Sprintf("Failed to create stack: %v", err)) + return + } + + // apply all record creations at once + var writes []*comatproto.RepoApplyWrites_Input_Writes_Elem + for _, p := range stack { record := p.AsRecord() - record.PatchBlob = blob.Blob writes = append(writes, &comatproto.RepoApplyWrites_Input_Writes_Elem{ RepoApplyWrites_Create: &comatproto.RepoApplyWrites_Create{ Collection: tangled.RepoPullNSID, @@ -1426,14 +1455,14 @@ func (s *Pulls) createStackedPullRequest( // create all pulls at once tx, err := s.db.BeginTx(r.Context(), nil) if err != nil { - s.logger.Error("failed to start tx", "err", err) + s.logger.Error("failed to start tx") s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") return } defer tx.Rollback() for _, p := range stack { - err = db.NewPull(tx, p) + err = db.PutPull(tx, p) if err != nil { s.logger.Error("failed to create pull request", "err", err) s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") @@ -1497,7 +1526,7 @@ func (s *Pulls) CompareBranchesFragment(w http.ResponseWriter, r *http.Request) user := s.oauth.GetMultiAccountUser(r) f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } @@ -1505,14 +1534,14 @@ func (s *Pulls) CompareBranchesFragment(w http.ResponseWriter, r *http.Request) xrpcBytes, err := tangled.GitTempListBranches(r.Context(), xrpcc, "", 0, f.RepoAt().String()) if err != nil { - s.logger.Error("failed to fetch branches", "err", err) + log.Println("failed to fetch branches", err) s.pages.Error503(w) return } var result types.RepoBranchesResponse if err := json.Unmarshal(xrpcBytes, &result); err != nil { - s.logger.Error("failed to decode XRPC response", "err", err) + log.Println("failed to decode XRPC response", err) s.pages.Error503(w) return } @@ -1541,7 +1570,7 @@ func (s *Pulls) CompareForksFragment(w http.ResponseWriter, r *http.Request) { forks, err := db.GetForksByDid(s.db, user.Active.Did) if err != nil { - s.logger.Error("failed to get forks", "err", err) + log.Println("failed to get forks", err) return } @@ -1557,7 +1586,7 @@ func (s *Pulls) CompareForksBranchesFragment(w http.ResponseWriter, r *http.Requ f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } @@ -1574,25 +1603,25 @@ func (s *Pulls) CompareForksBranchesFragment(w http.ResponseWriter, r *http.Requ orm.FilterEq("name", forkName), ) if err != nil { - s.logger.Error("failed to get repo", "did", forkOwnerDid, "name", forkName, "err", err) + log.Println("failed to get repo", "did", forkOwnerDid, "name", forkName, "err", err) return } sourceXrpcBytes, err := tangled.GitTempListBranches(r.Context(), xrpcc, "", 0, repo.RepoAt().String()) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.branches for source", "err", xrpcerr) + log.Println("failed to call XRPC repo.branches for source", xrpcerr) s.pages.Error503(w) return } - s.logger.Error("failed to fetch source branches", "err", err) + log.Println("failed to fetch source branches", err) return } // Decode source branches var sourceBranches types.RepoBranchesResponse if err := json.Unmarshal(sourceXrpcBytes, &sourceBranches); err != nil { - s.logger.Error("failed to decode source branches XRPC response", "err", err) + log.Println("failed to decode source branches XRPC response", err) s.pages.Error503(w) return } @@ -1600,18 +1629,18 @@ func (s *Pulls) CompareForksBranchesFragment(w http.ResponseWriter, r *http.Requ targetXrpcBytes, err := tangled.GitTempListBranches(r.Context(), xrpcc, "", 0, f.RepoAt().String()) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.branches for target", "err", xrpcerr) + log.Println("failed to call XRPC repo.branches for target", xrpcerr) s.pages.Error503(w) return } - s.logger.Error("failed to fetch target branches", "err", err) + log.Println("failed to fetch target branches", err) return } // Decode target branches var targetBranches types.RepoBranchesResponse if err := json.Unmarshal(targetXrpcBytes, &targetBranches); err != nil { - s.logger.Error("failed to decode target branches XRPC response", "err", err) + log.Println("failed to decode target branches XRPC response", err) s.pages.Error503(w) return } @@ -1632,7 +1661,7 @@ func (s *Pulls) ResubmitPull(w http.ResponseWriter, r *http.Request) { pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -1663,19 +1692,19 @@ func (s *Pulls) resubmitPatch(w http.ResponseWriter, r *http.Request) { pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } if user.Active.Did != pull.OwnerDid { - s.logger.Warn("unauthorized user") + log.Println("unauthorized user") w.WriteHeader(http.StatusUnauthorized) return } @@ -1690,26 +1719,26 @@ func (s *Pulls) resubmitBranch(w http.ResponseWriter, r *http.Request) { pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "resubmit-error", "Failed to edit patch. Try again later.") return } f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } if user.Active.Did != pull.OwnerDid { - s.logger.Warn("unauthorized user") + log.Println("unauthorized user") w.WriteHeader(http.StatusUnauthorized) return } roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} if !roles.IsPushAllowed() { - s.logger.Warn("unauthorized user") + log.Println("unauthorized user") w.WriteHeader(http.StatusUnauthorized) return } @@ -1727,18 +1756,18 @@ func (s *Pulls) resubmitBranch(w http.ResponseWriter, r *http.Request) { xrpcBytes, err := tangled.RepoCompare(r.Context(), xrpcc, repo, pull.TargetBranch, pull.PullSource.Branch) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.compare", "err", xrpcerr) + log.Println("failed to call XRPC repo.compare", xrpcerr) s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } - s.logger.Error("compare request failed", "err", err) + log.Printf("compare request failed: %s", err) s.pages.Notice(w, "resubmit-error", err.Error()) return } var comparison types.RepoFormatPatchResponse if err := json.Unmarshal(xrpcBytes, &comparison); err != nil { - s.logger.Error("failed to decode XRPC compare response", "err", err) + log.Println("failed to decode XRPC compare response", err) s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } @@ -1755,26 +1784,26 @@ func (s *Pulls) resubmitFork(w http.ResponseWriter, r *http.Request) { pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "resubmit-error", "Failed to edit patch. Try again later.") return } f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to get repo and knot", "err", err) + log.Println("failed to get repo and knot", err) return } if user.Active.Did != pull.OwnerDid { - s.logger.Warn("unauthorized user") + log.Println("unauthorized user") w.WriteHeader(http.StatusUnauthorized) return } forkRepo, err := db.GetRepoByAtUri(s.db, pull.PullSource.RepoAt.String()) if err != nil { - s.logger.Error("failed to get source repo", "err", err) + log.Println("failed to get source repo", err) s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } @@ -1787,7 +1816,7 @@ func (s *Pulls) resubmitFork(w http.ResponseWriter, r *http.Request) { oauth.WithDev(s.config.Core.Dev), ) if err != nil { - s.logger.Error("failed to connect to knot server", "err", err) + log.Printf("failed to connect to knot server: %v", err) return } @@ -1805,7 +1834,7 @@ func (s *Pulls) resubmitFork(w http.ResponseWriter, r *http.Request) { return } if !resp.Success { - s.logger.Warn("failed to update tracking ref", "err", resp.Error) + log.Println("Failed to update tracking ref.", "err", resp.Error) s.pages.Notice(w, "resubmit-error", "Failed to update tracking ref.") return } @@ -1821,18 +1850,18 @@ func (s *Pulls) resubmitFork(w http.ResponseWriter, r *http.Request) { forkXrpcBytes, err := tangled.RepoCompare(r.Context(), &indigoxrpc.Client{Host: forkHost}, forkRepoId, hiddenRef, pull.PullSource.Branch) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { - s.logger.Error("failed to call XRPC repo.compare for fork", "err", xrpcerr) + log.Println("failed to call XRPC repo.compare for fork", xrpcerr) s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } - s.logger.Error("failed to compare branches", "err", err) + log.Printf("failed to compare branches: %s", err) s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } var forkComparison types.RepoFormatPatchResponse if err := json.Unmarshal(forkXrpcBytes, &forkComparison); err != nil { - s.logger.Error("failed to decode XRPC compare response for fork", "err", err) + log.Println("failed to decode XRPC compare response for fork", err) s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } @@ -1857,9 +1886,10 @@ func (s *Pulls) resubmitPullHelper( combined string, sourceRev string, ) { - if pull.IsStacked() { - s.logger.Info("resubmitting stacked PR") - s.resubmitStackedPullHelper(w, r, repo, user, pull, patch, pull.StackId) + stack := r.Context().Value("stack").(models.Stack) + if stack != nil && len(stack) != 1 { + log.Println("resubmitting stacked PR") + s.resubmitStackedPullHelper(w, r, repo, user, pull, patch) return } @@ -1881,28 +1911,15 @@ func (s *Pulls) resubmitPullHelper( } } - tx, err := s.db.BeginTx(r.Context(), nil) - if err != nil { - s.logger.Error("failed to start tx", "err", err) - s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") - return - } - defer tx.Rollback() - pullAt := pull.AtUri() newRoundNumber := len(pull.Submissions) newPatch := patch newSourceRev := sourceRev combinedPatch := combined - err = db.ResubmitPull(tx, pullAt, newRoundNumber, newPatch, combinedPatch, newSourceRev) - if err != nil { - s.logger.Error("failed to create pull request", "err", err) - s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") - return - } + client, err := s.oauth.AuthorizedClient(r) if err != nil { - s.logger.Error("failed to authorize client", "err", err) + log.Println("failed to authorize client") s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } @@ -1916,18 +1933,17 @@ func (s *Pulls) resubmitPullHelper( blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(patch), ApplicationGzip) if err != nil { - s.logger.Error("failed to upload patch blob", "err", err) + log.Println("failed to upload patch blob", err) s.pages.Notice(w, "resubmit-error", "Failed to update pull request on the PDS. Try again later.") return } record := pull.AsRecord() - record.PatchBlob = blob.Blob + record.Rounds = append(record.Rounds, &tangled.RepoPull_Round{ + CreatedAt: time.Now().Format(time.RFC3339), + PatchBlob: blob.Blob, + }) record.CreatedAt = time.Now().Format(time.RFC3339) - if record.Source != nil { - record.Source.Sha = newSourceRev - } - _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ Collection: tangled.RepoPullNSID, Repo: user.Active.Did, @@ -1938,14 +1954,15 @@ func (s *Pulls) resubmitPullHelper( }, }) if err != nil { - s.logger.Error("failed to update record", "err", err) + log.Println("failed to update record", err) s.pages.Notice(w, "resubmit-error", "Failed to update pull request on the PDS. Try again later.") return } - if err = tx.Commit(); err != nil { - s.logger.Error("failed to commit transaction", "err", err) - s.pages.Notice(w, "resubmit-error", "Failed to resubmit pull.") + err = db.ResubmitPull(s.db, pullAt, newRoundNumber, newPatch, combinedPatch, newSourceRev, blob.Blob) + if err != nil { + log.Println("failed to create pull request", err) + s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") return } @@ -1960,14 +1977,48 @@ func (s *Pulls) resubmitStackedPullHelper( user *oauth.MultiAccountUser, pull *models.Pull, patch string, - stackId string, ) { targetBranch := pull.TargetBranch origStack, _ := r.Context().Value("stack").(models.Stack) - newStack, err := s.newStack(r.Context(), repo, user, targetBranch, patch, pull.PullSource, stackId) + + formatPatches, err := patchutil.ExtractPatches(patch) + if err != nil { + s.logger.Error("Failed to extract patches", "err", err) + s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Failed to parse patches.") + return + } + + // must have atleast 1 patch to begin with + if len(formatPatches) == 0 { + s.logger.Error("No patches found in the generated format-patch.") + s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request: No patches found in the generated patch.") + return + } + + client, err := s.oauth.AuthorizedClient(r) + if err != nil { + s.logger.Error("failed to get authorized client", "err", err) + s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") + return + } + + // first upload all blobs + blobs := make([]*lexutil.LexBlob, len(formatPatches)) + for i, p := range formatPatches { + blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(p.Raw), ApplicationGzip) + if err != nil { + s.logger.Error("failed to upload patch blob", "err", err) + s.pages.Notice(w, "pull", "Failed to create pull request. Try again later.") + return + } + s.logger.Info("uploaded blob", "idx", i+1, "total", len(formatPatches)) + blobs[i] = blob.Blob + } + + newStack, err := s.newStack(r.Context(), repo, user, targetBranch, pull.PullSource, formatPatches, blobs) if err != nil { - s.logger.Error("failed to create resubmitted stack", "err", err) + log.Println("failed to create resubmitted stack", err) s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") return } @@ -1976,10 +2027,10 @@ func (s *Pulls) resubmitStackedPullHelper( origById := make(map[string]*models.Pull) newById := make(map[string]*models.Pull) for _, p := range origStack { - origById[p.ChangeId] = p + origById[p.LatestSubmission().ChangeId()] = p } for _, p := range newStack { - newById[p.ChangeId] = p + newById[p.LatestSubmission().ChangeId()] = p } // commits that got deleted: corresponding pull is closed @@ -1991,42 +2042,49 @@ func (s *Pulls) resubmitStackedPullHelper( // pulls in orignal stack but not in new one for _, op := range origStack { - if _, ok := newById[op.ChangeId]; !ok { - deletions[op.ChangeId] = op + if _, ok := newById[op.LatestSubmission().ChangeId()]; !ok { + deletions[op.LatestSubmission().ChangeId()] = op } } // pulls in new stack but not in original one for _, np := range newStack { - if _, ok := origById[np.ChangeId]; !ok { - additions[np.ChangeId] = np + if _, ok := origById[np.LatestSubmission().ChangeId()]; !ok { + additions[np.LatestSubmission().ChangeId()] = np } } // NOTE: this loop can be written in any of above blocks, // but is written separately in the interest of simpler code for _, np := range newStack { - if op, ok := origById[np.ChangeId]; ok { + if op, ok := origById[np.LatestSubmission().ChangeId()]; ok { // pull exists in both stacks - updated[op.ChangeId] = struct{}{} + updated[op.LatestSubmission().ChangeId()] = struct{}{} } } + // NOTE: we can go through the newStack and update dependent relations and + // rkeys now that we know which ones have been updated + // update dependentOn relations for the entire stack + var parentAt *syntax.ATURI + for _, np := range newStack { + if op, ok := origById[np.LatestSubmission().ChangeId()]; ok { + // pull exists in both stacks + np.Rkey = op.Rkey + } + np.DependentOn = parentAt + x := np.AtUri() + parentAt = &x + } + tx, err := s.db.Begin() if err != nil { - s.logger.Error("failed to start transaction", "err", err) + log.Println("failed to start transaction", err) s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") return } defer tx.Rollback() - client, err := s.oauth.AuthorizedClient(r) - if err != nil { - s.logger.Error("failed to authorize client", "err", err) - s.pages.Notice(w, "resubmit-error", "Failed to create pull request. Try again later.") - return - } - // pds updates to make var writes []*comatproto.RepoApplyWrites_Input_Writes_Elem @@ -2037,9 +2095,9 @@ func (s *Pulls) resubmitStackedPullHelper( continue } - err := db.DeletePull(tx, p.RepoAt, p.PullId) + err := db.AbandonPulls(tx, orm.FilterEq("repo_at", p.RepoAt), orm.FilterEq("at_uri", p.AtUri())) if err != nil { - s.logger.Error("failed to delete pull", "err", err, "pull_id", p.PullId) + log.Println("failed to delete pull", err, p.PullId) s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") return } @@ -2053,21 +2111,27 @@ func (s *Pulls) resubmitStackedPullHelper( // new pulls are created for _, p := range additions { - err := db.NewPull(tx, p) + blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(p.LatestPatch()), ApplicationGzip) if err != nil { - s.logger.Error("failed to create pull", "err", err, "pull_id", p.PullId) - s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") + log.Println("failed to upload patch blob", err) + s.pages.Notice(w, "resubmit-error", "Failed to update pull request on the PDS. Try again later.") return } + p.Submissions[0].Blob = *blob.Blob - blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(patch), ApplicationGzip) - if err != nil { - s.logger.Error("failed to upload patch blob", "err", err) - s.pages.Notice(w, "resubmit-error", "Failed to update pull request on the PDS. Try again later.") + if err = db.PutPull(tx, p); err != nil { + log.Println("failed to create pull", err, p.PullId) + s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") return } + record := p.AsRecord() - record.PatchBlob = blob.Blob + record.Rounds = []*tangled.RepoPull_Round{ + { + CreatedAt: time.Now().Format(time.RFC3339), + PatchBlob: blob.Blob, + }, + } writes = append(writes, &comatproto.RepoApplyWrites_Input_Writes_Elem{ RepoApplyWrites_Create: &comatproto.RepoApplyWrites_Create{ Collection: tangled.RepoPullNSID, @@ -2090,26 +2154,33 @@ func (s *Pulls) resubmitStackedPullHelper( } // resubmit the new pull + np.Rkey = op.Rkey pullAt := op.AtUri() newRoundNumber := len(op.Submissions) newPatch := np.LatestPatch() combinedPatch := np.LatestSubmission().Combined newSourceRev := np.LatestSha() - err := db.ResubmitPull(tx, pullAt, newRoundNumber, newPatch, combinedPatch, newSourceRev) + + blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(newPatch), ApplicationGzip) if err != nil { - s.logger.Error("failed to update pull", "err", err, "pull_id", op.PullId) - s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") + log.Println("failed to upload patch blob", err) + s.pages.Notice(w, "resubmit-error", "Failed to update pull request on the PDS. Try again later.") return } - blob, err := xrpc.RepoUploadBlob(r.Context(), client, gz(patch), ApplicationGzip) + err = db.ResubmitPull(tx, pullAt, newRoundNumber, newPatch, combinedPatch, newSourceRev, blob.Blob) if err != nil { - s.logger.Error("failed to upload patch blob", "err", err) - s.pages.Notice(w, "resubmit-error", "Failed to update pull request on the PDS. Try again later.") + log.Println("failed to update pull", err, op.PullId) + s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") return } + record := np.AsRecord() - record.PatchBlob = blob.Blob + record.Rounds = op.AsRecord().Rounds + record.Rounds = append(record.Rounds, &tangled.RepoPull_Round{ + CreatedAt: time.Now().Format(time.RFC3339), + PatchBlob: blob.Blob, + }) writes = append(writes, &comatproto.RepoApplyWrites_Input_Writes_Elem{ RepoApplyWrites_Update: &comatproto.RepoApplyWrites_Update{ Collection: tangled.RepoPullNSID, @@ -2121,41 +2192,23 @@ func (s *Pulls) resubmitStackedPullHelper( }) } - // update parent-change-id relations for the entire stack - for _, p := range newStack { - err := db.SetPullParentChangeId( - tx, - p.ParentChangeId, - // these should be enough filters to be unique per-stack - orm.FilterEq("repo_at", p.RepoAt.String()), - orm.FilterEq("owner_did", p.OwnerDid), - orm.FilterEq("change_id", p.ChangeId), - ) - - if err != nil { - s.logger.Error("failed to update pull", "err", err, "pull_id", p.PullId) - s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") - return - } - } - - err = tx.Commit() - if err != nil { - s.logger.Error("failed to resubmit pull", "err", err) - s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") - return - } - _, err = comatproto.RepoApplyWrites(r.Context(), client, &comatproto.RepoApplyWrites_Input{ Repo: user.Active.Did, Writes: writes, }) if err != nil { - s.logger.Error("failed to create stacked pull request", "err", err) + log.Println("failed to create stacked pull request", err) s.pages.Notice(w, "pull", "Failed to create stacked pull request. Try again later.") return } + err = tx.Commit() + if err != nil { + log.Println("failed to resubmit pull", err) + s.pages.Notice(w, "pull-resubmit-error", "Failed to resubmit pull request. Try again later.") + return + } + ownerSlashRepo := reporesolver.GetBaseRepoPath(r, repo) s.pages.HxLocation(w, fmt.Sprintf("/%s/pulls/%d", ownerSlashRepo, pull.PullId)) } @@ -2164,48 +2217,42 @@ func (s *Pulls) MergePull(w http.ResponseWriter, r *http.Request) { user := s.oauth.GetMultiAccountUser(r) f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to resolve repo:", "err", err) + log.Println("failed to resolve repo:", err) s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") return } pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-merge-error", "Failed to merge patch. Try again later.") return } - var pullsToMerge models.Stack - pullsToMerge = append(pullsToMerge, pull) - if pull.IsStacked() { - stack, ok := r.Context().Value("stack").(models.Stack) - if !ok { - s.logger.Error("failed to get stack") - s.pages.Notice(w, "pull-merge-error", "Failed to merge patch. Try again later.") - return - } - - // combine patches of substack - subStack := stack.StrictlyBelow(pull) - // collect the portion of the stack that is mergeable - mergeable := subStack.Mergeable() - // add to total patch - pullsToMerge = append(pullsToMerge, mergeable...) + stack, ok := r.Context().Value("stack").(models.Stack) + if !ok { + log.Println("failed to get stack") + s.pages.Notice(w, "pull-merge-error", "Failed to merge patch. Try again later.") + return } + // combine patches of substack + subStack := stack.Below(pull) + // collect the portion of the stack that is mergeable + pullsToMerge := subStack.Mergeable() + patch := pullsToMerge.CombinedPatch() ident, err := s.idResolver.ResolveIdent(r.Context(), pull.OwnerDid) if err != nil { - s.logger.Error("resolving identity", "err", err) + log.Printf("resolving identity: %s", err) w.WriteHeader(http.StatusNotFound) return } email, err := db.GetPrimaryEmail(s.db, pull.OwnerDid) if err != nil { - s.logger.Error("failed to get primary email", "err", err) + log.Printf("failed to get primary email: %s", err) } authorName := ident.Handle.String() @@ -2233,7 +2280,7 @@ func (s *Pulls) MergePull(w http.ResponseWriter, r *http.Request) { oauth.WithDev(s.config.Core.Dev), ) if err != nil { - s.logger.Error("failed to connect to knot server", "err", err) + log.Printf("failed to connect to knot server: %v", err) s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") return } @@ -2246,26 +2293,28 @@ func (s *Pulls) MergePull(w http.ResponseWriter, r *http.Request) { tx, err := s.db.Begin() if err != nil { - s.logger.Error("failed to start transaction", "err", err) + log.Println("failed to start transcation", err) s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") return } defer tx.Rollback() + var atUris []syntax.ATURI for _, p := range pullsToMerge { - err := db.MergePull(tx, f.RepoAt(), p.PullId) - if err != nil { - s.logger.Error("failed to update pull request status in database", "err", err) - s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") - return - } + atUris = append(atUris, p.AtUri()) p.State = models.PullMerged } + err = db.MergePulls(tx, orm.FilterEq("repo_at", f.RepoAt()), orm.FilterIn("at_uri", atUris)) + if err != nil { + log.Printf("failed to update pull request status in database: %s", err) + s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") + return + } err = tx.Commit() if err != nil { // TODO: this is unsound, we should also revert the merge from the knotserver here - s.logger.Error("failed to update pull request status in database", "err", err) + log.Printf("failed to update pull request status in database: %s", err) s.pages.Notice(w, "pull-merge-error", "Failed to merge pull request. Try again later.") return } @@ -2284,13 +2333,13 @@ func (s *Pulls) ClosePull(w http.ResponseWriter, r *http.Request) { f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("malformed middleware", "err", err) + log.Println("malformed middleware") return } pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -2302,7 +2351,7 @@ func (s *Pulls) ClosePull(w http.ResponseWriter, r *http.Request) { isPullAuthor := user.Active.Did == pull.OwnerDid isCloseAllowed := isOwner || isCollaborator || isPullAuthor if !isCloseAllowed { - s.logger.Warn("failed to close pull: unauthorized") + log.Println("failed to close pull") s.pages.Notice(w, "pull-close", "You are unauthorized to close this pull.") return } @@ -2310,36 +2359,33 @@ func (s *Pulls) ClosePull(w http.ResponseWriter, r *http.Request) { // Start a transaction tx, err := s.db.BeginTx(r.Context(), nil) if err != nil { - s.logger.Error("failed to start transaction", "err", err) + log.Println("failed to start transaction", err) s.pages.Notice(w, "pull-close", "Failed to close pull.") return } defer tx.Rollback() - var pullsToClose []*models.Pull - pullsToClose = append(pullsToClose, pull) - - // if this PR is stacked, then we want to close all PRs below this one on the stack - if pull.IsStacked() { - stack := r.Context().Value("stack").(models.Stack) - subStack := stack.StrictlyBelow(pull) - pullsToClose = append(pullsToClose, subStack...) - } - + // if this PR is stacked, then we want to close all PRs above this one on the stack + stack := r.Context().Value("stack").(models.Stack) + pullsToClose := stack.Above(pull) + var atUris []syntax.ATURI for _, p := range pullsToClose { - // Close the pull in the database - err = db.ClosePull(tx, f.RepoAt(), p.PullId) - if err != nil { - s.logger.Error("failed to close pull", "err", err) - s.pages.Notice(w, "pull-close", "Failed to close pull.") - return - } + atUris = append(atUris, p.AtUri()) p.State = models.PullClosed } + err = db.ClosePulls( + tx, + orm.FilterEq("repo_at", f.RepoAt()), + orm.FilterIn("at_uri", atUris), + ) + if err != nil { + log.Println("failed to close pulls", err) + s.pages.Notice(w, "pull-close", "Failed to close pull.") + } // Commit the transaction if err = tx.Commit(); err != nil { - s.logger.Error("failed to commit transaction", "err", err) + log.Println("failed to commit transaction", err) s.pages.Notice(w, "pull-close", "Failed to close pull.") return } @@ -2357,14 +2403,14 @@ func (s *Pulls) ReopenPull(w http.ResponseWriter, r *http.Request) { f, err := s.repoResolver.Resolve(r) if err != nil { - s.logger.Error("failed to resolve repo", "err", err) + log.Println("failed to resolve repo", err) s.pages.Notice(w, "pull-reopen", "Failed to reopen pull.") return } pull, ok := r.Context().Value("pull").(*models.Pull) if !ok { - s.logger.Error("failed to get pull") + log.Println("failed to get pull") s.pages.Notice(w, "pull-error", "Failed to edit patch. Try again later.") return } @@ -2376,7 +2422,7 @@ func (s *Pulls) ReopenPull(w http.ResponseWriter, r *http.Request) { isPullAuthor := user.Active.Did == pull.OwnerDid isCloseAllowed := isOwner || isCollaborator || isPullAuthor if !isCloseAllowed { - s.logger.Warn("failed to close pull: unauthorized") + log.Println("failed to close pull") s.pages.Notice(w, "pull-close", "You are unauthorized to close this pull.") return } @@ -2384,36 +2430,33 @@ func (s *Pulls) ReopenPull(w http.ResponseWriter, r *http.Request) { // Start a transaction tx, err := s.db.BeginTx(r.Context(), nil) if err != nil { - s.logger.Error("failed to start transaction", "err", err) + log.Println("failed to start transaction", err) s.pages.Notice(w, "pull-reopen", "Failed to reopen pull.") return } defer tx.Rollback() - var pullsToReopen []*models.Pull - pullsToReopen = append(pullsToReopen, pull) - // if this PR is stacked, then we want to reopen all PRs above this one on the stack - if pull.IsStacked() { - stack := r.Context().Value("stack").(models.Stack) - subStack := stack.StrictlyAbove(pull) - pullsToReopen = append(pullsToReopen, subStack...) - } - + stack := r.Context().Value("stack").(models.Stack) + pullsToReopen := stack.Below(pull) + var atUris []syntax.ATURI for _, p := range pullsToReopen { - // Close the pull in the database - err = db.ReopenPull(tx, f.RepoAt(), p.PullId) - if err != nil { - s.logger.Error("failed to close pull", "err", err) - s.pages.Notice(w, "pull-close", "Failed to close pull.") - return - } + atUris = append(atUris, p.AtUri()) p.State = models.PullOpen } + err = db.ReopenPulls( + tx, + orm.FilterEq("repo_at", f.RepoAt()), + orm.FilterIn("at_uri", atUris), + ) + if err != nil { + log.Println("failed to reopen pulls", err) + s.pages.Notice(w, "pull-close", "Failed to reopen pull.") + } // Commit the transaction if err = tx.Commit(); err != nil { - s.logger.Error("failed to commit transaction", "err", err) + log.Println("failed to commit transaction", err) s.pages.Notice(w, "pull-reopen", "Failed to reopen pull.") return } @@ -2426,23 +2469,20 @@ func (s *Pulls) ReopenPull(w http.ResponseWriter, r *http.Request) { s.pages.HxLocation(w, fmt.Sprintf("/%s/pulls/%d", ownerSlashRepo, pull.PullId)) } -func (s *Pulls) newStack(ctx context.Context, repo *models.Repo, user *oauth.MultiAccountUser, targetBranch, patch string, pullSource *models.PullSource, stackId string) (models.Stack, error) { - formatPatches, err := patchutil.ExtractPatches(patch) - if err != nil { - return nil, fmt.Errorf("Failed to extract patches: %v", err) - } - - // must have atleast 1 patch to begin with - if len(formatPatches) == 0 { - return nil, fmt.Errorf("No patches found in the generated format-patch.") - } - - // the stack is identified by a UUID +func (s *Pulls) newStack( + ctx context.Context, + repo *models.Repo, + user *oauth.MultiAccountUser, + targetBranch string, + pullSource *models.PullSource, + formatPatches []types.FormatPatch, + blobs []*lexutil.LexBlob, +) (models.Stack, error) { var stack models.Stack - parentChangeId := "" - for _, fp := range formatPatches { + var parentAtUri *syntax.ATURI + for i, fp := range formatPatches { // all patches must have a jj change-id - changeId, err := fp.ChangeId() + _, err := fp.ChangeId() if err != nil { return nil, fmt.Errorf("Stacking is only supported if all patches contain a change-id commit header.") } @@ -2457,6 +2497,7 @@ func (s *Pulls) newStack(ctx context.Context, repo *models.Repo, user *oauth.Mul Patch: fp.Raw, SourceRev: fp.SHA, Combined: fp.Raw, + Blob: *blobs[i], } pull := models.Pull{ Title: title, @@ -2472,15 +2513,15 @@ func (s *Pulls) newStack(ctx context.Context, repo *models.Repo, user *oauth.Mul }, PullSource: pullSource, Created: time.Now(), + State: models.PullOpen, - StackId: stackId, - ChangeId: changeId, - ParentChangeId: parentChangeId, + DependentOn: parentAtUri, } stack = append(stack, &pull) - parentChangeId = changeId + parent := pull.AtUri() + parentAtUri = &parent } return stack, nil diff --git a/knotserver/ingester.go b/knotserver/ingester.go index 5d2cc6e7..a336fd34 100644 --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -11,11 +11,13 @@ import ( "strings" comatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - "github.com/bluesky-social/jetstream/pkg/models" + jmodels "github.com/bluesky-social/jetstream/pkg/models" securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/api/tangled" + "tangled.org/core/appview/models" "tangled.org/core/knotserver/db" "tangled.org/core/knotserver/git" "tangled.org/core/log" @@ -23,7 +25,7 @@ import ( "tangled.org/core/workflow" ) -func (h *Knot) processPublicKey(ctx context.Context, event *models.Event) error { +func (h *Knot) processPublicKey(ctx context.Context, event *jmodels.Event) error { l := log.FromContext(ctx) raw := json.RawMessage(event.Commit.Record) did := event.Did @@ -45,7 +47,7 @@ func (h *Knot) processPublicKey(ctx context.Context, event *models.Event) error return nil } -func (h *Knot) processKnotMember(ctx context.Context, event *models.Event) error { +func (h *Knot) processKnotMember(ctx context.Context, event *jmodels.Event) error { l := log.FromContext(ctx) raw := json.RawMessage(event.Commit.Record) did := event.Did @@ -85,26 +87,11 @@ func (h *Knot) processKnotMember(ctx context.Context, event *models.Event) error return nil } -func (h *Knot) processPull(ctx context.Context, event *models.Event) error { - raw := json.RawMessage(event.Commit.Record) - did := event.Did - - var record tangled.RepoPull - if err := json.Unmarshal(raw, &record); err != nil { - return fmt.Errorf("failed to unmarshal record: %w", err) - } - - l := log.FromContext(ctx) - l = l.With("handler", "processPull") - l = l.With("did", did) - +func (h *Knot) validatePullRecord(record *tangled.RepoPull) error { if record.Target == nil { return fmt.Errorf("ignoring pull record: target repo is nil") } - l = l.With("target_repo", record.Target.Repo) - l = l.With("target_branch", record.Target.Branch) - if record.Source == nil { return fmt.Errorf("ignoring pull record: not a branch-based pull request") } @@ -113,50 +100,85 @@ func (h *Knot) processPull(ctx context.Context, event *models.Event) error { return fmt.Errorf("ignoring pull record: fork based pull") } - repoAt, err := syntax.ParseATURI(record.Target.Repo) + return nil +} + +func (h *Knot) resolveTargetRepo(ctx context.Context, targetRepoUri string) (*identity.Identity, *tangled.Repo, error) { + repoAt, err := syntax.ParseATURI(targetRepoUri) if err != nil { - return fmt.Errorf("failed to parse ATURI: %w", err) + return nil, nil, fmt.Errorf("failed to parse ATURI: %w", err) } - // resolve this aturi to extract the repo record - ident, err := h.resolver.ResolveIdent(ctx, repoAt.Authority().String()) - if err != nil || ident.Handle.IsInvalidHandle() { - return fmt.Errorf("failed to resolve handle: %w", err) + // resolve the repo owner to extract the repo record + repoOwnerIdent, err := h.resolver.ResolveIdent(ctx, repoAt.Authority().String()) + if err != nil || repoOwnerIdent.Handle.IsInvalidHandle() { + return nil, nil, fmt.Errorf("failed to resolve repo owner handle: %w", err) } xrpcc := xrpc.Client{ - Host: ident.PDSEndpoint(), + Host: repoOwnerIdent.PDSEndpoint(), } resp, err := comatproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) if err != nil { - return fmt.Errorf("failed to resolver repo: %w", err) + return nil, nil, fmt.Errorf("failed to resolve repo: %w", err) } repo := resp.Value.Val.(*tangled.Repo) + return repoOwnerIdent, repo, nil +} - if repo.Knot != h.c.Server.Hostname { - return fmt.Errorf("rejected pull record: not this knot, %s != %s", repo.Knot, h.c.Server.Hostname) +func (h *Knot) fetchLatestSubmission(ctx context.Context, did, rkey string, record *tangled.RepoPull) (*models.PullSubmission, error) { + // resolve the PR owner's identity to fetch the blob from their PDS + prOwnerIdent, err := h.resolver.ResolveIdent(ctx, did) + if err != nil || prOwnerIdent.Handle.IsInvalidHandle() { + return nil, fmt.Errorf("failed to resolve PR owner handle: %w", err) } - didSlashRepo, err := securejoin.SecureJoin(ident.DID.String(), repo.Name) + 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 fmt.Errorf("failed to construct relative repo path: %w", err) + 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() - repoPath, err := securejoin.SecureJoin(h.c.Repo.ScanPath, didSlashRepo) + req, err := http.NewRequestWithContext(ctx, http.MethodGet, blobUrl.String(), nil) if err != nil { - return fmt.Errorf("failed to construct absolute repo path: %w", err) + 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() - gr, err := git.Open(repoPath, record.Source.Sha) + blob := io.ReadCloser(blobResp.Body) + latestSubmission, err := models.PullSubmissionFromRecord(did, rkey, roundNumber, round, &blob) if err != nil { - return fmt.Errorf("failed to open git repository: %w", err) + return nil, fmt.Errorf("failed to parse submission: %w", err) + } + + return latestSubmission, nil +} + +func (h *Knot) discoverWorkflows(ctx context.Context, repoPath, sha string) (workflow.RawPipeline, error) { + gr, err := git.Open(repoPath, sha) + if err != nil { + return nil, fmt.Errorf("failed to open git repository: %w", err) } workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) if err != nil { - return fmt.Errorf("failed to open workflow directory: %w", err) + return nil, fmt.Errorf("failed to open workflow directory: %w", err) } var pipeline workflow.RawPipeline @@ -177,11 +199,17 @@ func (h *Knot) processPull(ctx context.Context, event *models.Event) error { }) } + return pipeline, nil +} + +func (h *Knot) compilePipeline(ctx context.Context, repoOwner *identity.Identity, repo *tangled.Repo, sourceBranch, sourceSha, targetBranch string, rawPipeline workflow.RawPipeline) tangled.Pipeline { + l := log.FromContext(ctx) + trigger := tangled.Pipeline_PullRequestTriggerData{ Action: "create", - SourceBranch: record.Source.Branch, - SourceSha: record.Source.Sha, - TargetBranch: record.Target.Branch, + SourceBranch: sourceBranch, + SourceSha: sourceSha, + TargetBranch: targetBranch, } compiler := workflow.Compiler{ @@ -189,35 +217,111 @@ func (h *Knot) processPull(ctx context.Context, event *models.Event) error { Kind: string(workflow.TriggerKindPullRequest), PullRequest: &trigger, Repo: &tangled.Pipeline_TriggerRepo{ - Did: ident.DID.String(), + Did: repoOwner.DID.String(), Knot: repo.Knot, Repo: repo.Name, }, }, } - cp := compiler.Compile(compiler.Parse(pipeline)) - eventJson, err := json.Marshal(cp) + l.Info("raw", "raw", rawPipeline) + parsed := compiler.Parse(rawPipeline) + l.Info("parsed", "parsed", parsed) + compiled := compiler.Compile(parsed) + + l.Info("compiler diagnostics", "diagnostics", compiler.Diagnostics) + + return compiled +} + +func (h *Knot) processPull(ctx context.Context, event *jmodels.Event) error { + raw := json.RawMessage(event.Commit.Record) + rkey := event.Commit.RKey + did := event.Did + + var record tangled.RepoPull + if err := json.Unmarshal(raw, &record); err != nil { + return fmt.Errorf("failed to unmarshal record: %w", err) + } + + l := log.FromContext(ctx) + l = l.With("handler", "processPull") + l = l.With("did", did) + + l.Info("validating pull record") + if err := h.validatePullRecord(&record); err != nil { + return err + } + + l = l.With("target_repo", record.Target.Repo) + l = l.With("target_branch", record.Target.Branch) + + l.Info("resolving target repo") + repoOwnerIdent, repo, err := h.resolveTargetRepo(ctx, record.Target.Repo) if err != nil { - return fmt.Errorf("failed to marshal pipeline event: %w", err) + return err + } + + if repo.Knot != h.c.Server.Hostname { + return fmt.Errorf("rejected pull record: not this knot, %s != %s", repo.Knot, h.c.Server.Hostname) + } + + l.Info("fetching latest submission") + latestSubmission, err := h.fetchLatestSubmission(ctx, did, rkey, &record) + if err != nil { + return err + } + + sha := latestSubmission.SourceRev + if sha == "" { + return fmt.Errorf("failed to extract source SHA from pull submission") + } + l = l.With("sha", sha) + + l.Info("constructing repo path") + didSlashRepo, err := securejoin.SecureJoin(repoOwnerIdent.DID.String(), repo.Name) + if err != nil { + return fmt.Errorf("failed to construct relative repo path: %w", err) } + repoPath, err := securejoin.SecureJoin(h.c.Repo.ScanPath, didSlashRepo) + if err != nil { + return fmt.Errorf("failed to construct absolute repo path: %w", err) + } + + l.Info("discovering workflows", "repo_path", repoPath) + pipeline, err := h.discoverWorkflows(ctx, repoPath, sha) + if err != nil { + return err + } + + l.Info("compiling pipeline", "workflow_count", len(pipeline)) + cp := h.compilePipeline(ctx, repoOwnerIdent, repo, record.Source.Branch, sha, record.Target.Branch, pipeline) + // do not run empty pipelines if cp.Workflows == nil { + l.Info("skipping empty pipeline") return nil } + l.Info("marshaling pipeline event") + eventJson, err := json.Marshal(cp) + if err != nil { + return fmt.Errorf("failed to marshal pipeline event: %w", err) + } + ev := db.Event{ Rkey: TID(), Nsid: tangled.PipelineNSID, EventJson: string(eventJson), } + l.Info("inserting pipeline event") return h.db.InsertEvent(ev, h.n) } // duplicated from add collaborator -func (h *Knot) processCollaborator(ctx context.Context, event *models.Event) error { +func (h *Knot) processCollaborator(ctx context.Context, event *jmodels.Event) error { raw := json.RawMessage(event.Commit.Record) did := event.Did @@ -319,8 +423,8 @@ func (h *Knot) fetchAndAddKeys(ctx context.Context, did string) error { return nil } -func (h *Knot) processMessages(ctx context.Context, event *models.Event) error { - if event.Kind != models.EventKindCommit { +func (h *Knot) processMessages(ctx context.Context, event *jmodels.Event) error { + if event.Kind != jmodels.EventKindCommit { return nil } diff --git a/patchutil/patchutil.go b/patchutil/patchutil.go index 90c81ad7..29150de3 100644 --- a/patchutil/patchutil.go +++ b/patchutil/patchutil.go @@ -17,9 +17,8 @@ import ( func ExtractPatches(formatPatch string) ([]types.FormatPatch, error) { patches := splitFormatPatch(formatPatch) - result := []types.FormatPatch{} - - for _, patch := range patches { + result := make([]types.FormatPatch, len(patches)) + for i, patch := range patches { files, headerStr, err := gitdiff.Parse(strings.NewReader(patch)) if err != nil { return nil, fmt.Errorf("failed to parse patch: %w", err) @@ -30,11 +29,11 @@ func ExtractPatches(formatPatch string) ([]types.FormatPatch, error) { return nil, fmt.Errorf("failed to parse patch header: %w", err) } - result = append(result, types.FormatPatch{ + result[i] = types.FormatPatch{ Files: files, PatchHeader: header, Raw: patch, - }) + } } return result, nil