From b1d5ce3c98fb99d57d433b128244572dc6cb8434 Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 03 Mar 2026 12:53:03 +0000 Subject: [PATCH] appview: DID-based routing, state/handler/middleware updates Signed-off-by: Lewis Lewis: May this revision serve well! --- appview/ingester.go | 92 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------ appview/issues/issues.go | 4 ++-- appview/middleware/middleware.go | 14 +++++++++----- appview/models/issue.go | 6 +++++- appview/pulls/pulls.go | 46 ++++++++++++++++++++++++++-------------------- appview/repo/archive.go | 2 +- appview/repo/artifact.go | 23 ++++++++++++++++------- appview/repo/blob.go | 7 +++---- appview/repo/compare.go | 8 ++++---- appview/repo/index.go | 1 - appview/repo/log.go | 3 +-- appview/repo/repo.go | 175 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------------------------------------------- appview/repo/settings.go | 14 +++++++------- appview/reporesolver/resolver.go | 19 +++++++++++++------ appview/state/git_http.go | 29 +++++++---------------------- appview/state/knotstream.go | 103 +++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------------------------------------- appview/state/router.go | 28 ++++++++++++++++++++++++++-- appview/state/star.go | 18 +++++++++++++----- appview/state/state.go | 126 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------------------------- appview/validator/label.go | 2 +- 20 file(s) changed, 469 insertion(s)(+), 251 deletion(s)(-) diff --git a/appview/ingester.go b/appview/ingester.go --- a/appview/ingester.go +++ b/appview/ingester.go @@ -2,7 +2,9 @@ package appview import ( "context" + "database/sql" "encoding/json" + "errors" "fmt" "log/slog" "maps" @@ -116,16 +118,36 @@ l.Error("invalid record", "err", err) return err } - subjectUri, err = syntax.ParseATURI(record.Subject) - if err != nil { - l.Error("invalid record", "err", err) - return err + star := &models.Star{ + Did: did, + Rkey: e.Commit.RKey, } - err = db.AddStar(i.Db, &models.Star{ - Did: did, - RepoAt: subjectUri, - Rkey: e.Commit.RKey, - }) + + switch { + case record.SubjectDid != nil: + repo, repoErr := db.GetRepo(i.Db, orm.FilterEq("repo_did", *record.SubjectDid)) + if repoErr == nil { + subjectUri = repo.RepoAt() + star.RepoAt = subjectUri + } + case record.Subject != nil: + subjectUri, err = syntax.ParseATURI(*record.Subject) + if err != nil { + l.Error("invalid record", "err", err) + return err + } + star.RepoAt = subjectUri + repo, repoErr := db.GetRepoByAtUri(i.Db, subjectUri.String()) + if repoErr == nil && repo.RepoDid != "" { + if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.FeedStarNSID, e.Commit.RKey, *record.Subject); enqErr != nil { + l.Warn("failed to enqueue PDS rewrite for star", "err", enqErr, "did", did, "repoDid", repo.RepoDid) + } + } + default: + l.Error("star record has neither subject nor subjectDid") + return fmt.Errorf("star record has neither subject nor subjectDid") + } + err = db.AddStar(i.Db, star) case jmodels.CommitOperationDelete: err = db.DeleteStarByRkey(i.Db, did, e.Commit.RKey) } @@ -220,19 +242,40 @@ l.Error("invalid record", "err", err) return err } - repoAt, err := syntax.ParseATURI(record.Repo) - if err != nil { - return err + var repo *models.Repo + if record.RepoDid != nil && *record.RepoDid != "" { + repo, err = db.GetRepoByDid(i.Db, *record.RepoDid) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return fmt.Errorf("failed to look up repo by DID %s: %w", *record.RepoDid, err) + } + } + if repo == nil && record.Repo != nil { + repoAt, parseErr := syntax.ParseATURI(*record.Repo) + if parseErr != nil { + return parseErr + } + repo, err = db.GetRepoByAtUri(i.Db, repoAt.String()) + if err != nil { + return err + } + } + if repo == nil { + return fmt.Errorf("artifact record has neither valid repoDid nor repo field") } - repo, err := db.GetRepoByAtUri(i.Db, repoAt.String()) - if err != nil { + ok, err := i.Enforcer.E.Enforce(did, repo.Knot, repo.RepoIdentifier(), "repo:push") + if err != nil || !ok { return err } - ok, err := i.Enforcer.E.Enforce(did, repo.Knot, repo.DidSlashRepo(), "repo:push") - if err != nil || !ok { - return err + repoDid := repo.RepoDid + if repoDid == "" && record.RepoDid != nil { + repoDid = *record.RepoDid + } + if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil { + if enqErr := db.EnqueuePdsRewrite(i.Db, did, repoDid, tangled.RepoArtifactNSID, e.Commit.RKey, *record.Repo); enqErr != nil { + l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid) + } } createdAt, err := time.Parse(time.RFC3339, record.CreatedAt) @@ -243,7 +286,7 @@ artifact := models.Artifact{ Did: did, Rkey: e.Commit.RKey, - RepoAt: repoAt, + RepoAt: repo.RepoAt(), Tag: plumbing.Hash(record.Tag), CreatedAt: createdAt, BlobCid: cid.Cid(record.Artifact.Ref), @@ -834,8 +877,21 @@ } issue := models.IssueFromRecord(did, rkey, record) + if issue.RepoAt == "" { + return fmt.Errorf("issue record has no repo field") + } + if err := i.Validator.ValidateIssue(&issue); err != nil { return fmt.Errorf("failed to validate issue: %w", err) + } + + if record.Repo != nil { + repo, repoErr := db.GetRepoByAtUri(i.Db, *record.Repo) + if repoErr == nil && repo.RepoDid != "" { + if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.RepoIssueNSID, rkey, *record.Repo); enqErr != nil { + l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) + } + } } tx, err := ddb.BeginTx(ctx, nil) diff --git a/appview/issues/issues.go b/appview/issues/issues.go --- a/appview/issues/issues.go +++ b/appview/issues/issues.go @@ -309,7 +309,7 @@ rp.pages.Error404(w) return } - roles := repoinfo.RolesInRepo{Roles: rp.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} + roles := repoinfo.RolesInRepo{Roles: rp.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.RepoIdentifier())} isRepoOwner := roles.IsOwner() isCollaborator := roles.IsCollaborator() isIssueOwner := user.Active.Did == issue.Did @@ -357,7 +357,7 @@ rp.pages.Error404(w) return } - roles := repoinfo.RolesInRepo{Roles: rp.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} + roles := repoinfo.RolesInRepo{Roles: rp.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.RepoIdentifier())} isRepoOwner := roles.IsOwner() isCollaborator := roles.IsCollaborator() isIssueOwner := user.Active.Did == issue.Did diff --git a/appview/middleware/middleware.go b/appview/middleware/middleware.go --- a/appview/middleware/middleware.go +++ b/appview/middleware/middleware.go @@ -18,6 +18,7 @@ "tangled.org/core/appview/oauth" "tangled.org/core/appview/pages" "tangled.org/core/appview/pagination" "tangled.org/core/appview/reporesolver" + "tangled.org/core/appview/state/userutil" "tangled.org/core/idresolver" "tangled.org/core/orm" "tangled.org/core/rbac" @@ -162,9 +163,9 @@ http.Error(w, "malformed url", http.StatusBadRequest) return } - ok, err := mw.enforcer.E.Enforce(actor.Active.Did, f.Knot, f.DidSlashRepo(), requiredPerm) + ok, err := mw.enforcer.E.Enforce(actor.Active.Did, f.Knot, f.RepoIdentifier(), requiredPerm) if err != nil || !ok { - log.Printf("%s does not have perms of a %s in repo %s", actor.Active.Did, requiredPerm, f.DidSlashRepo()) + log.Printf("%s does not have perms of a %s in repo %s", actor.Active.Did, requiredPerm, f.RepoIdentifier()) http.Error(w, "Forbiden", http.StatusUnauthorized) return } @@ -195,7 +196,6 @@ id, err = mw.idResolver.ResolveIdent(req.Context(), string(did)) } } } - // invalid did or handle if err != nil { log.Printf("failed to resolve did/handle '%s': %s\n", didOrHandle, err) mw.pages.Error404(w) @@ -342,11 +342,15 @@ fullName := reporesolver.GetBaseRepoPath(r, f) if r.Header.Get("User-Agent") == "Go-http-client/1.1" { if r.URL.Query().Get("go-get") == "1" { + modulePath := userutil.FlattenDid(fullName) + if strings.Contains(modulePath, ":") { + modulePath = userutil.FlattenDid(f.Did) + "/" + f.Name + } html := fmt.Sprintf( ` `, - fullName, fullName, - fullName, fullName, + modulePath, fullName, + modulePath, fullName, ) w.Header().Set("Content-Type", "text/html") w.Write([]byte(html)) diff --git a/appview/models/issue.go b/appview/models/issue.go --- a/appview/models/issue.go +++ b/appview/models/issue.go @@ -45,7 +45,7 @@ for i, uri := range i.References { references[i] = string(uri) } repoAtStr := i.RepoAt.String() - return tangled.RepoIssue{ + rec := tangled.RepoIssue{ Repo: &repoAtStr, Title: i.Title, Body: &i.Body, @@ -53,6 +53,10 @@ Mentions: mentions, References: references, CreatedAt: i.Created.Format(time.RFC3339), } + if i.Repo != nil && i.Repo.RepoDid != "" { + rec.RepoDid = &i.Repo.RepoDid + } + return rec } func (i *Issue) State() string { diff --git a/appview/pulls/pulls.go b/appview/pulls/pulls.go --- a/appview/pulls/pulls.go +++ b/appview/pulls/pulls.go @@ -408,7 +408,7 @@ return nil } // user can only delete branch if they are a collaborator in the repo that the branch belongs to - perms := s.enforcer.GetPermissionsInRepo(user.Active.Did, repo.Knot, repo.DidSlashRepo()) + perms := s.enforcer.GetPermissionsInRepo(user.Active.Did, repo.Knot, repo.RepoIdentifier()) if !slices.Contains(perms, "repo:push") { return nil } @@ -432,10 +432,8 @@ } var sourceRepo syntax.ATURI if pull.PullSource.RepoAt != nil { - // fork-based pulls sourceRepo = *pull.PullSource.RepoAt } else { - // pulls within the same repo sourceRepo = repo.RepoAt() } @@ -929,7 +927,7 @@ return } // Determine PR type based on input parameters - roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} + roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.RepoIdentifier())} isPushAllowed := roles.IsPushAllowed() isBranchBased := isPushAllowed && sourceBranch != "" && fromFork == "" isForkBased := fromFork != "" && sourceBranch != "" @@ -1045,8 +1043,7 @@ xrpcc := &indigoxrpc.Client{ Host: host, } - didSlashRepo := fmt.Sprintf("%s/%s", repo.Did, repo.Name) - xrpcBytes, err := tangled.RepoCompare(r.Context(), xrpcc, didSlashRepo, targetBranch, sourceBranch) + xrpcBytes, err := tangled.RepoCompare(r.Context(), xrpcc, repo.RepoIdentifier(), targetBranch, sourceBranch) if err != nil { if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { s.logger.Error("failed to call XRPC repo.compare", "err", xrpcerr) @@ -1155,8 +1152,7 @@ forkXrpcc := &indigoxrpc.Client{ Host: forkHost, } - forkRepoId := fmt.Sprintf("%s/%s", fork.Did, fork.Name) - forkXrpcBytes, err := tangled.RepoCompare(r.Context(), forkXrpcc, forkRepoId, hiddenRef, sourceBranch) + forkXrpcBytes, err := tangled.RepoCompare(r.Context(), forkXrpcc, fork.RepoIdentifier(), 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) @@ -1197,6 +1193,9 @@ Branch: sourceBranch, Repo: &forkAtUriStr, Sha: sourceRev, } + if fork.RepoDid != "" { + recordPullSource.RepoDid = &fork.RepoDid + } s.createPullRequest(w, r, repo, user, title, body, targetBranch, patch, combined, sourceRev, pullSource, recordPullSource, isStacked) } @@ -1313,11 +1312,8 @@ Repo: user.Active.Did, Rkey: rkey, Record: &lexutil.LexiconTypeDecoder{ Val: &tangled.RepoPull{ - Title: title, - Target: &tangled.RepoPull_Target{ - Repo: string(repo.RepoAt()), - Branch: targetBranch, - }, + Title: title, + Target: repoPullTarget(repo, targetBranch), PatchBlob: blob.Blob, Source: recordPullSource, CreatedAt: time.Now().Format(time.RFC3339), @@ -1707,7 +1703,7 @@ w.WriteHeader(http.StatusUnauthorized) return } - roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} + roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.RepoIdentifier())} if !roles.IsPushAllowed() { s.logger.Warn("unauthorized user") w.WriteHeader(http.StatusUnauthorized) @@ -1723,8 +1719,7 @@ xrpcc := &indigoxrpc.Client{ Host: host, } - repo := fmt.Sprintf("%s/%s", f.Did, f.Name) - xrpcBytes, err := tangled.RepoCompare(r.Context(), xrpcc, repo, pull.TargetBranch, pull.PullSource.Branch) + xrpcBytes, err := tangled.RepoCompare(r.Context(), xrpcc, f.RepoIdentifier(), 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) @@ -1817,8 +1812,7 @@ if !s.config.Core.Dev { forkScheme = "https" } forkHost := fmt.Sprintf("%s://%s", forkScheme, forkRepo.Knot) - forkRepoId := fmt.Sprintf("%s/%s", forkRepo.Did, forkRepo.Name) - forkXrpcBytes, err := tangled.RepoCompare(r.Context(), &indigoxrpc.Client{Host: forkHost}, forkRepoId, hiddenRef, pull.PullSource.Branch) + forkXrpcBytes, err := tangled.RepoCompare(r.Context(), &indigoxrpc.Client{Host: forkHost}, forkRepo.RepoIdentifier(), 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) @@ -2296,7 +2290,7 @@ return } // auth filter: only owner or collaborators can close - roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} + roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.RepoIdentifier())} isOwner := roles.IsOwner() isCollaborator := roles.IsCollaborator() isPullAuthor := user.Active.Did == pull.OwnerDid @@ -2370,7 +2364,7 @@ return } // auth filter: only owner or collaborators can close - roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.DidSlashRepo())} + roles := repoinfo.RolesInRepo{Roles: s.enforcer.GetPermissionsInRepo(user.Active.Did, f.Knot, f.RepoIdentifier())} isOwner := roles.IsOwner() isCollaborator := roles.IsCollaborator() isPullAuthor := user.Active.Did == pull.OwnerDid @@ -2495,3 +2489,15 @@ return &b } func ptrPullState(s models.PullState) *models.PullState { return &s } + +func repoPullTarget(repo *models.Repo, branch string) *tangled.RepoPull_Target { + s := string(repo.RepoAt()) + t := &tangled.RepoPull_Target{ + Branch: branch, + Repo: &s, + } + if repo.RepoDid != "" { + t.RepoDid = &repo.RepoDid + } + return t +} diff --git a/appview/repo/archive.go b/appview/repo/archive.go --- a/appview/repo/archive.go +++ b/appview/repo/archive.go @@ -60,7 +60,7 @@ if link := resp.Header.Get("Link"); link != "" { if resolvedRef, err := extractImmutableLink(link); err == nil { newLink := fmt.Sprintf("<%s/%s/archive/%s.tar.gz>; rel=\"immutable\"", - rp.config.Core.BaseUrl(), f.DidSlashRepo(), resolvedRef) + rp.config.Core.BaseUrl(), f.RepoIdentifier(), resolvedRef) w.Header().Set("Link", newLink) } } diff --git a/appview/repo/artifact.go b/appview/repo/artifact.go --- a/appview/repo/artifact.go +++ b/appview/repo/artifact.go @@ -80,13 +80,7 @@ Collection: tangled.RepoArtifactNSID, Repo: user.Active.Did, Rkey: rkey, Record: &lexutil.LexiconTypeDecoder{ - Val: &tangled.RepoArtifact{ - Artifact: uploadBlobResp.Blob, - CreatedAt: createdAt.Format(time.RFC3339), - Name: header.Filename, - Repo: f.RepoAt().String(), - Tag: tag.Tag.Hash[:], - }, + Val: repoArtifactRecord(f, uploadBlobResp.Blob, createdAt, header.Filename, tag.Tag.Hash[:]), }, }) if err != nil { @@ -350,3 +344,18 @@ } return tag, nil } + +func repoArtifactRecord(f *models.Repo, blob *lexutil.LexBlob, createdAt time.Time, name string, tag []byte) *tangled.RepoArtifact { + rec := &tangled.RepoArtifact{ + Artifact: blob, + CreatedAt: createdAt.Format(time.RFC3339), + Name: name, + Tag: tag, + } + s := f.RepoAt().String() + rec.Repo = &s + if f.RepoDid != "" { + rec.RepoDid = &f.RepoDid + } + return rec +} diff --git a/appview/repo/blob.go b/appview/repo/blob.go --- a/appview/repo/blob.go +++ b/appview/repo/blob.go @@ -58,8 +58,7 @@ host := fmt.Sprintf("%s://%s", scheme, f.Knot) xrpcc := &indigoxrpc.Client{ Host: host, } - repo := fmt.Sprintf("%s/%s", f.Did, f.Name) - resp, err := tangled.RepoBlob(r.Context(), xrpcc, filePath, false, ref, repo) + resp, err := tangled.RepoBlob(r.Context(), xrpcc, filePath, false, ref, f.RepoIdentifier()) if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { l.Error("failed to call XRPC repo.blob", "err", xrpcerr) rp.pages.Error503(w) @@ -139,7 +138,7 @@ scheme := "http" if !rp.config.Core.Dev { scheme = "https" } - repo := f.DidSlashRepo() + repo := f.RepoIdentifier() baseURL := &url.URL{ Scheme: scheme, Host: f.Knot, @@ -290,7 +289,7 @@ if !config.Core.Dev { scheme = "https" } - repoName := fmt.Sprintf("%s/%s", repo.Did, repo.Name) + repoName := repo.RepoIdentifier() baseURL := &url.URL{ Scheme: scheme, Host: repo.Knot, diff --git a/appview/repo/compare.go b/appview/repo/compare.go --- a/appview/repo/compare.go +++ b/appview/repo/compare.go @@ -141,9 +141,9 @@ xrpcc := &indigoxrpc.Client{ Host: host, } - repo := fmt.Sprintf("%s/%s", f.Did, f.Name) + repoId := f.RepoIdentifier() - branchBytes, err := tangled.RepoBranches(r.Context(), xrpcc, "", 0, repo) + branchBytes, err := tangled.RepoBranches(r.Context(), xrpcc, "", 0, repoId) if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { l.Error("failed to call XRPC repo.branches", "err", xrpcerr) rp.pages.Error503(w) @@ -157,7 +157,7 @@ rp.pages.Notice(w, "compare-error", "Failed to produce comparison. Try again later.") return } - tagBytes, err := tangled.RepoTags(r.Context(), xrpcc, "", 0, repo) + tagBytes, err := tangled.RepoTags(r.Context(), xrpcc, "", 0, repoId) if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { l.Error("failed to call XRPC repo.tags", "err", xrpcerr) rp.pages.Error503(w) @@ -171,7 +171,7 @@ rp.pages.Notice(w, "compare-error", "Failed to produce comparison. Try again later.") return } - compareBytes, err := tangled.RepoCompare(r.Context(), xrpcc, repo, base, head) + compareBytes, err := tangled.RepoCompare(r.Context(), xrpcc, repoId, base, head) if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { l.Error("failed to call XRPC repo.compare", "err", xrpcerr) rp.pages.Error503(w) diff --git a/appview/repo/index.go b/appview/repo/index.go --- a/appview/repo/index.go +++ b/appview/repo/index.go @@ -239,7 +239,6 @@ // buildIndexResponse creates a RepoIndexResponse by combining multiple xrpc calls in parallel func (rp *Repo) buildIndexResponse(ctx context.Context, repo *models.Repo, ref string) (*types.RepoIndexResponse, error) { xrpcc := &indigoxrpc.Client{Host: rp.config.KnotMirror.Url} - // first get branches to determine the ref if not specified branchesBytes, err := tangled.GitTempListBranches(ctx, xrpcc, "", 0, repo.RepoAt().String()) if err != nil { return nil, fmt.Errorf("calling knotmirror git.listBranches: %w", err) diff --git a/appview/repo/log.go b/appview/repo/log.go --- a/appview/repo/log.go +++ b/appview/repo/log.go @@ -164,8 +164,7 @@ xrpcc := &indigoxrpc.Client{ Host: host, } - repo := fmt.Sprintf("%s/%s", f.Did, f.Name) - xrpcBytes, err := tangled.RepoDiff(r.Context(), xrpcc, ref, repo) + xrpcBytes, err := tangled.RepoDiff(r.Context(), xrpcc, ref, f.RepoIdentifier()) if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { l.Error("failed to call XRPC repo.diff", "err", xrpcerr) rp.pages.Error503(w) diff --git a/appview/repo/repo.go b/appview/repo/repo.go --- a/appview/repo/repo.go +++ b/appview/repo/repo.go @@ -36,7 +36,7 @@ comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/atproto/syntax" lexutil "github.com/bluesky-social/indigo/lex/util" - securejoin "github.com/cyphar/filepath-securejoin" + "github.com/go-chi/chi/v5" ) @@ -318,10 +318,13 @@ fail("Failed to add label.", err) return } - err = db.SubscribeLabel(tx, &models.RepoLabel{ + if err = db.SubscribeLabel(tx, &models.RepoLabel{ RepoAt: f.RepoAt(), LabelAt: label.AtUri(), - }) + }); err != nil { + fail("Failed to subscribe to label.", err) + return + } err = tx.Commit() if err != nil { @@ -755,11 +758,8 @@ Collection: tangled.RepoCollaboratorNSID, Repo: currentUser.Active.Did, Rkey: rkey, Record: &lexutil.LexiconTypeDecoder{ - Val: &tangled.RepoCollaborator{ - Subject: collaboratorIdent.DID.String(), - Repo: string(f.RepoAt()), - CreatedAt: createdAt.Format(time.RFC3339), - }}, + Val: repoCollaboratorRecord(f, collaboratorIdent.DID.String(), createdAt), + }, }) // invalid record if err != nil { @@ -794,7 +794,7 @@ } } defer rollback() - err = rp.enforcer.AddCollaborator(collaboratorIdent.DID.String(), f.Knot, f.DidSlashRepo()) + err = rp.enforcer.AddCollaborator(collaboratorIdent.DID.String(), f.Knot, f.RepoIdentifier()) if err != nil { fail("Failed to add collaborator permissions.", err) return @@ -900,19 +900,19 @@ } }() // remove collaborator RBAC - repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(f.DidSlashRepo(), f.Knot) + repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(f.RepoIdentifier(), f.Knot) if err != nil { rp.pages.Notice(w, noticeId, "Failed to remove collaborators") return } for _, c := range repoCollaborators { did := c[0] - rp.enforcer.RemoveCollaborator(did, f.Knot, f.DidSlashRepo()) + rp.enforcer.RemoveCollaborator(did, f.Knot, f.RepoIdentifier()) } l.Info("removed collaborators") // remove repo RBAC - err = rp.enforcer.RemoveRepo(f.Did, f.Knot, f.DidSlashRepo()) + err = rp.enforcer.RemoveRepo(f.Did, f.Knot, f.RepoIdentifier()) if err != nil { rp.pages.Notice(w, noticeId, "Failed to update RBAC rules") return @@ -1067,28 +1067,106 @@ if rp.config.Core.Dev { uri = "http" } - forkSourceUrl := fmt.Sprintf("%s://%s/%s/%s", uri, f.Knot, f.Did, f.Name) + forkSourceUrl := fmt.Sprintf("%s://%s/%s", uri, f.Knot, f.RepoIdentifier()) l = l.With("cloneUrl", forkSourceUrl) - sourceAt := f.RepoAt().String() - - // create an atproto record for this fork rkey := tid.TID() + + // TODO: this could coordinate better with the knot to recieve a clone status + client, err := rp.oauth.ServiceClient( + r, + oauth.WithService(targetKnot), + oauth.WithLxm(tangled.RepoCreateNSID), + oauth.WithDev(rp.config.Core.Dev), + oauth.WithTimeout(time.Second*20), + ) + if err != nil { + l.Error("could not create service client", "err", err) + rp.pages.Notice(w, "repo", "Failed to connect to knot server.") + return + } + + forkInput := &tangled.RepoCreate_Input{ + Rkey: rkey, + Name: forkName, + Source: &forkSourceUrl, + } + createResp, createErr := tangled.RepoCreate( + r.Context(), + client, + forkInput, + ) + if err := xrpcclient.HandleXrpcErr(createErr); err != nil { + rp.pages.Notice(w, "repo", err.Error()) + return + } + + var repoDid string + if createResp != nil && createResp.RepoDid != nil { + repoDid = *createResp.RepoDid + } + if repoDid == "" { + l.Error("knot returned empty repo DID for fork") + rp.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.") + return + } + + forkSource := f.RepoAt().String() + if f.RepoDid != "" { + forkSource = f.RepoDid + } + repo := &models.Repo{ Did: user.Active.Did, Name: forkName, Knot: targetKnot, Rkey: rkey, - Source: sourceAt, + Source: forkSource, Description: f.Description, Created: time.Now(), Labels: rp.config.Label.DefaultLabelDefs, + RepoDid: repoDid, } record := repo.AsRecord() + cleanupKnot := func() { + go func() { + delays := []time.Duration{0, 2 * time.Second, 5 * time.Second} + for attempt, delay := range delays { + time.Sleep(delay) + deleteClient, dErr := rp.oauth.ServiceClient( + r, + oauth.WithService(targetKnot), + oauth.WithLxm(tangled.RepoDeleteNSID), + oauth.WithDev(rp.config.Core.Dev), + ) + if dErr != nil { + l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr) + continue + } + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{ + Did: user.Active.Did, + Name: forkName, + Rkey: rkey, + }); dErr != nil { + cancel() + l.Error("failed to clean up fork on knot after rollback", "attempt", attempt+1, "err", dErr) + continue + } + cancel() + l.Info("successfully cleaned up fork on knot after rollback", "attempt", attempt+1) + return + } + l.Error("exhausted retries for knot cleanup, fork may be orphaned", + "did", user.Active.Did, "fork", forkName, "knot", targetKnot) + }() + } + atpClient, err := rp.oauth.AuthorizedClient(r) if err != nil { l.Error("failed to create xrpcclient", "err", err) + cleanupKnot() rp.pages.Notice(w, "repo", "Failed to fork repository.") return } @@ -1103,6 +1181,7 @@ }, }) if err != nil { l.Error("failed to write to PDS", "err", err) + cleanupKnot() rp.pages.Notice(w, "repo", "Failed to announce repository creation.") return } @@ -1118,54 +1197,25 @@ rp.pages.Notice(w, "repo", "Failed to save repository information.") return } - // The rollback function reverts a few things on failure: - // - the pending txn - // - the ACLs - // - the atproto record created rollback := func() { err1 := tx.Rollback() err2 := rp.enforcer.E.LoadPolicy() err3 := rollbackRecord(context.Background(), aturi, atpClient) - // ignore txn complete errors, this is okay if errors.Is(err1, sql.ErrTxDone) { err1 = nil } if errs := errors.Join(err1, err2, err3); errs != nil { l.Error("failed to rollback changes", "errs", errs) - return + } + + if aturi != "" { + cleanupKnot() } } defer rollback() - // TODO: this could coordinate better with the knot to recieve a clone status - client, err := rp.oauth.ServiceClient( - r, - oauth.WithService(targetKnot), - oauth.WithLxm(tangled.RepoCreateNSID), - oauth.WithDev(rp.config.Core.Dev), - oauth.WithTimeout(time.Second*20), // big repos take time to clone - ) - if err != nil { - l.Error("could not create service client", "err", err) - rp.pages.Notice(w, "repo", "Failed to connect to knot server.") - return - } - - err = tangled.RepoCreate( - r.Context(), - client, - &tangled.RepoCreate_Input{ - Rkey: rkey, - Source: &forkSourceUrl, - }, - ) - if err := xrpcclient.HandleXrpcErr(err); err != nil { - rp.pages.Notice(w, "repo", err.Error()) - return - } - err = db.AddRepo(tx, repo) if err != nil { l.Error("failed to AddRepo", "err", err) @@ -1173,9 +1223,8 @@ rp.pages.Notice(w, "repo", "Failed to save repository information.") return } - // acls - p, _ := securejoin.SecureJoin(user.Active.Did, forkName) - err = rp.enforcer.AddRepo(user.Active.Did, targetKnot, p) + rbacPath := repo.RepoIdentifier() + err = rp.enforcer.AddRepo(user.Active.Did, targetKnot, rbacPath) if err != nil { l.Error("failed to add ACLs", "err", err) rp.pages.Notice(w, "repo", "Failed to set up repository permissions.") @@ -1196,11 +1245,14 @@ http.Error(w, err.Error(), http.StatusInternalServerError) return } - // reset the ATURI because the transaction completed successfully aturi = "" rp.notifier.NewRepo(r.Context(), repo) - rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Active.Did, forkName)) + if repoDid != "" { + rp.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid)) + } else { + rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Active.Did, forkName)) + } } } @@ -1225,3 +1277,16 @@ Rkey: rkey, }) return err } + +func repoCollaboratorRecord(f *models.Repo, subject string, createdAt time.Time) *tangled.RepoCollaborator { + rec := &tangled.RepoCollaborator{ + Subject: subject, + CreatedAt: createdAt.Format(time.RFC3339), + } + s := string(f.RepoAt()) + rec.Repo = &s + if f.RepoDid != "" { + rec.RepoDid = &f.RepoDid + } + return rec +} diff --git a/appview/repo/settings.go b/appview/repo/settings.go --- a/appview/repo/settings.go +++ b/appview/repo/settings.go @@ -293,7 +293,7 @@ // Trigger an initial deploy asynchronously so the handler returns promptly. // Skip entirely if there is no active domain claim — the site cannot be served anyway. ownerClaim, _ := db.GetActiveDomainClaimForDid(rp.db, f.Did) if ownerClaim == nil { - rp.logger.Info("skipping deploy: no active domain claim", "repo", f.DidSlashRepo()) + rp.logger.Info("skipping deploy: no active domain claim", "repo", f.RepoIdentifier()) } else if rp.cfClient.Enabled() { scheme := "http" if !rp.config.Core.Dev { @@ -313,7 +313,7 @@ } deployErr := sites.Deploy(ctx, rp.cfClient, knotHost, f.Did, f.Name, branch, dir) if deployErr != nil { - l.Error("sites: initial R2 sync failed", "repo", f.DidSlashRepo(), "err", deployErr) + l.Error("sites: initial R2 sync failed", "repo", f.RepoIdentifier(), "err", deployErr) deploy.Status = models.SiteDeployStatusFailure deploy.Error = deployErr.Error() } else { @@ -321,18 +321,18 @@ deploy.Status = models.SiteDeployStatusSuccess } if err := db.AddSiteDeploy(rp.db, deploy); err != nil { - l.Error("sites: failed to record deploy", "repo", f.DidSlashRepo(), "err", err) + l.Error("sites: failed to record deploy", "repo", f.RepoIdentifier(), "err", err) } if deployErr == nil { if err := sites.PutDomainMapping(ctx, rp.cfClient, ownerClaim.Domain, f.Did, f.Name, isIndex); err != nil { l.Error("sites: KV write failed", "domain", ownerClaim.Domain, "err", err) } - rp.logger.Info("site deployed to r2", "repo", f.DidSlashRepo(), "is_index", isIndex) + rp.logger.Info("site deployed to r2", "repo", f.RepoIdentifier(), "is_index", isIndex) } }() } else { - rp.logger.Warn("cloudflare integration is disabled; site won't be deployed", "repo", f.DidSlashRepo()) + rp.logger.Warn("cloudflare integration is disabled; site won't be deployed", "repo", f.RepoIdentifier()) } rp.pages.HxRefresh(w) @@ -367,7 +367,7 @@ go func() { ctx := context.Background() if err := sites.Delete(ctx, rp.cfClient, f.Did, f.Name); err != nil { - l.Error("sites: R2 delete failed", "repo", f.DidSlashRepo(), "err", err) + l.Error("sites: R2 delete failed", "repo", f.RepoIdentifier(), "err", err) } if ownerClaim != nil { if err := sites.DeleteDomainMapping(ctx, rp.cfClient, ownerClaim.Domain, f.Name); err != nil { @@ -459,7 +459,7 @@ f, err := rp.repoResolver.Resolve(r) user := rp.oauth.GetMultiAccountUser(r) collaborators, err := func(repo *models.Repo) ([]pages.Collaborator, error) { - repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(repo.DidSlashRepo(), repo.Knot) + repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(repo.RepoIdentifier(), repo.Knot) if err != nil { return nil, err } diff --git a/appview/reporesolver/resolver.go b/appview/reporesolver/resolver.go --- a/appview/reporesolver/resolver.go +++ b/appview/reporesolver/resolver.go @@ -36,12 +36,15 @@ } // NOTE: this... should not even be here. the entire package will be removed in future refactor func GetBaseRepoPath(r *http.Request, repo *models.Repo) string { + if repo.RepoDid != "" { + return repo.RepoDid + } var ( user = chi.URLParam(r, "user") name = chi.URLParam(r, "repo") ) if user == "" || name == "" { - return repo.DidSlashRepo() + return repo.RepoIdentifier() } return path.Join(user, name) } @@ -77,13 +80,13 @@ isStarred := false roles := repoinfo.RolesInRepo{} if user != nil && user.Active != nil { isStarred = db.GetStarStatus(rr.execer, user.Active.Did, repoAt) - roles.Roles = rr.enforcer.GetPermissionsInRepo(user.Active.Did, repo.Knot, repo.DidSlashRepo()) + roles.Roles = rr.enforcer.GetPermissionsInRepo(user.Active.Did, repo.Knot, repo.RepoIdentifier()) } stats := repo.RepoStats if stats == nil { - starCount, err := db.GetStarCount(rr.execer, repoAt) - if err != nil { + starCount, starErr := db.GetStarCount(rr.execer, repoAt) + if starErr != nil { log.Println("failed to get star count for ", repoAt) } issueCount, err := db.GetIssueCount(rr.execer, repoAt) @@ -104,9 +107,13 @@ var sourceRepo *models.Repo var err error if repo.Source != "" { - sourceRepo, err = db.GetRepoByAtUri(rr.execer, repo.Source) + if strings.HasPrefix(repo.Source, "did:") { + sourceRepo, err = db.GetRepoByDid(rr.execer, repo.Source) + } else { + sourceRepo, err = db.GetRepoByAtUri(rr.execer, repo.Source) + } if err != nil { - log.Println("failed to get repo by at uri", err) + log.Println("failed to get source repo", err) } } diff --git a/appview/state/git_http.go b/appview/state/git_http.go --- a/appview/state/git_http.go +++ b/appview/state/git_http.go @@ -37,7 +37,6 @@ w.Header().Set("X-Content-Type-Options", "nosniff") } func (s *State) InfoRefs(w http.ResponseWriter, r *http.Request) { - user := r.Context().Value("resolvedId").(identity.Identity) repo := r.Context().Value("repo").(*models.Repo) scheme := "https" @@ -45,27 +44,20 @@ if s.config.Core.Dev { scheme = "http" } - // check for the 'service' url param service := r.URL.Query().Get("service") var contentType string switch service { case "git-receive-pack": contentType = "application/x-git-receive-pack-advertisement" default: - // git-upload-pack is the default service for git-clone / git-fetch. contentType = "application/x-git-upload-pack-advertisement" } - targetURL := fmt.Sprintf("%s://%s/%s/%s/info/refs?%s", scheme, repo.Knot, user.DID, repo.Name, r.URL.RawQuery) + targetURL := fmt.Sprintf("%s://%s/%s/info/refs?%s", scheme, repo.Knot, repo.RepoIdentifier(), r.URL.RawQuery) s.proxyRequest(w, r, targetURL, contentType) } func (s *State) UploadArchive(w http.ResponseWriter, r *http.Request) { - user, ok := r.Context().Value("resolvedId").(identity.Identity) - if !ok { - http.Error(w, "failed to resolve user", http.StatusInternalServerError) - return - } repo := r.Context().Value("repo").(*models.Repo) scheme := "https" @@ -73,16 +65,11 @@ if s.config.Core.Dev { scheme = "http" } - targetURL := fmt.Sprintf("%s://%s/%s/%s/git-upload-archive?%s", scheme, repo.Knot, user.DID, repo.Name, r.URL.RawQuery) + targetURL := fmt.Sprintf("%s://%s/%s/git-upload-archive?%s", scheme, repo.Knot, repo.RepoIdentifier(), r.URL.RawQuery) s.proxyRequest(w, r, targetURL, "application/x-git-upload-archive-result") } func (s *State) UploadPack(w http.ResponseWriter, r *http.Request) { - user, ok := r.Context().Value("resolvedId").(identity.Identity) - if !ok { - http.Error(w, "failed to resolve user", http.StatusInternalServerError) - return - } repo := r.Context().Value("repo").(*models.Repo) scheme := "https" @@ -90,16 +77,11 @@ if s.config.Core.Dev { scheme = "http" } - targetURL := fmt.Sprintf("%s://%s/%s/%s/git-upload-pack?%s", scheme, repo.Knot, user.DID, repo.Name, r.URL.RawQuery) + targetURL := fmt.Sprintf("%s://%s/%s/git-upload-pack?%s", scheme, repo.Knot, repo.RepoIdentifier(), r.URL.RawQuery) s.proxyRequest(w, r, targetURL, "application/x-git-upload-pack-result") } func (s *State) ReceivePack(w http.ResponseWriter, r *http.Request) { - user, ok := r.Context().Value("resolvedId").(identity.Identity) - if !ok { - http.Error(w, "failed to resolve user", http.StatusInternalServerError) - return - } repo := r.Context().Value("repo").(*models.Repo) scheme := "https" @@ -107,7 +89,7 @@ if s.config.Core.Dev { scheme = "http" } - targetURL := fmt.Sprintf("%s://%s/%s/%s/git-receive-pack?%s", scheme, repo.Knot, user.DID, repo.Name, r.URL.RawQuery) + targetURL := fmt.Sprintf("%s://%s/%s/git-receive-pack?%s", scheme, repo.Knot, repo.RepoIdentifier(), r.URL.RawQuery) s.proxyRequest(w, r, targetURL, "application/x-git-receive-pack-result") } @@ -123,6 +105,9 @@ proxyReq.Header = r.Header.Clone() repoOwnerHandle := chi.URLParam(r, "user") + if id, ok := r.Context().Value("resolvedId").(identity.Identity); ok && !id.Handle.IsInvalidHandle() { + repoOwnerHandle = id.Handle.String() + } proxyReq.Header.Set("x-tangled-repo-owner-handle", repoOwnerHandle) resp, err := client.Do(proxyReq) diff --git a/appview/state/knotstream.go b/appview/state/knotstream.go --- a/appview/state/knotstream.go +++ b/appview/state/knotstream.go @@ -2,6 +2,7 @@ package state import ( "context" + "database/sql" "encoding/json" "errors" "fmt" @@ -66,6 +67,20 @@ return ec.NewConsumer(cfg), nil } +func resolveRepo(d *db.DB, repoDid *string, ownerDid, repoName string) (*models.Repo, error) { + if repoDid != nil && *repoDid != "" { + return db.GetRepoByDid(d, *repoDid) + } + repos, err := db.GetRepos(d, orm.FilterEq("did", ownerDid), orm.FilterEq("name", repoName)) + if err != nil { + return nil, err + } + if len(repos) == 0 { + return nil, sql.ErrNoRows + } + return &repos[0], nil +} + func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client) ec.ProcessFunc { return func(ctx context.Context, source ec.Source, msg ec.Message) error { switch msg.Nsid { @@ -96,19 +111,18 @@ if !slices.Contains(knownKnots, source.Key()) { return fmt.Errorf("%s does not belong to %s, something is fishy", record.CommitterDid, source.Key()) } - repo, err := db.GetRepo( - d, - orm.FilterEq("did", record.RepoDid), - orm.FilterEq("name", record.RepoName), - orm.FilterEq("knot", source.Key()), - ) - if err != nil { - return fmt.Errorf("repo %s/%s on knot %s not found", record.RepoDid, record.RepoName, source.Key()) + ownerDid := "" + if record.OwnerDid != nil { + ownerDid = *record.OwnerDid + } + + repo, lookupErr := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName) + if lookupErr != nil { + return fmt.Errorf("failed to look up repo: %w", lookupErr) } logger.Info("processing gitRefUpdate event", - "repo_did", record.RepoDid, - "repo_name", record.RepoName, + "repo", repo.RepoIdentifier(), "ref", record.Ref, "old_sha", record.OldSha, "new_sha", record.NewSha) @@ -145,15 +159,15 @@ return } pushedBranch := ref.Short() - repos, err := db.GetRepos( - d, - orm.FilterEq("did", record.RepoDid), - orm.FilterEq("name", record.RepoName), - ) - if err != nil || len(repos) != 1 { + ownerDid := "" + if record.OwnerDid != nil { + ownerDid = *record.OwnerDid + } + + repo, err := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName) + if err != nil { return } - repo := repos[0] siteConfig, err := db.GetRepoSiteConfig(d, repo.RepoAt().String()) if err != nil || siteConfig == nil { @@ -177,9 +191,9 @@ CommitSHA: record.NewSha, Trigger: models.SiteDeployTriggerPush, } - deployErr := sites.Deploy(ctx, cfClient, knotHost, record.RepoDid, record.RepoName, siteConfig.Branch, siteConfig.Dir) + deployErr := sites.Deploy(ctx, cfClient, knotHost, repo.RepoIdentifier(), record.RepoName, siteConfig.Branch, siteConfig.Dir) if deployErr != nil { - logger.Error("sites: R2 sync failed on push", "repo", record.RepoDid+"/"+record.RepoName, "err", deployErr) + logger.Error("sites: R2 sync failed on push", "repo", repo.RepoIdentifier(), "err", deployErr) deploy.Status = models.SiteDeployStatusFailure deploy.Error = deployErr.Error() } else { @@ -187,11 +201,11 @@ deploy.Status = models.SiteDeployStatusSuccess } if err := db.AddSiteDeploy(d, deploy); err != nil { - logger.Error("sites: failed to record deploy", "repo", record.RepoDid+"/"+record.RepoName, "err", err) + logger.Error("sites: failed to record deploy", "repo", repo.RepoIdentifier(), "err", err) } if deployErr == nil { - logger.Info("site deployed to r2", "repo", record.RepoDid+"/"+record.RepoName) + logger.Info("site deployed to r2", "repo", repo.RepoIdentifier()) } } @@ -233,21 +247,19 @@ } func updateRepoLanguages(d *db.DB, record tangled.GitRefUpdate) error { if record.Meta == nil || record.Meta.LangBreakdown == nil || record.Meta.LangBreakdown.Inputs == nil { - return fmt.Errorf("empty language data for repo: %s/%s", record.RepoDid, record.RepoName) + return fmt.Errorf("empty language data for repo: %v/%s", record.OwnerDid, record.RepoName) } - repos, err := db.GetRepos( - d, - orm.FilterEq("did", record.RepoDid), - orm.FilterEq("name", record.RepoName), - ) - if err != nil { - return fmt.Errorf("failed to look for repo in DB (%s/%s): %w", record.RepoDid, record.RepoName, err) + ownerDid := "" + if record.OwnerDid != nil { + ownerDid = *record.OwnerDid } - if len(repos) != 1 { - return fmt.Errorf("incorrect number of repos returned: %d (expected 1)", len(repos)) + + r, lookupErr := resolveRepo(d, record.RepoDid, ownerDid, record.RepoName) + if lookupErr != nil { + return fmt.Errorf("failed to look up repo: %w", lookupErr) } - repo := repos[0] + repo := *r ref := plumbing.ReferenceName(record.Ref) if !ref.IsBranch() { @@ -300,25 +312,15 @@ if record.TriggerMetadata.Repo == nil { return fmt.Errorf("empty repo: nsid %s, rkey %s", msg.Nsid, msg.Rkey) } - repo, err := db.GetRepo( - d, - orm.FilterEq("did", record.TriggerMetadata.Repo.Did), - orm.FilterEq("name", record.TriggerMetadata.Repo.Repo), - orm.FilterEq("knot", source.Key()), - ) - if err != nil { - return fmt.Errorf( - "failed to look for repo in DB: nsid %s, rkey %s, %s/%s, knot %s, %w", - msg.Nsid, - msg.Rkey, - record.TriggerMetadata.Repo.Did, - record.TriggerMetadata.Repo.Did, - source.Key(), - err, - ) + repoName := "" + if record.TriggerMetadata.Repo.Repo != nil { + repoName = *record.TriggerMetadata.Repo.Repo } - // does this repo have a spindle configured? + repo, lookupErr := resolveRepo(d, record.TriggerMetadata.Repo.RepoDid, record.TriggerMetadata.Repo.Did, repoName) + if lookupErr != nil { + return fmt.Errorf("failed to look up repo: %w", lookupErr) + } if repo.Spindle == "" { return fmt.Errorf("repo does not have a spindle configured yet: nsid %s, rkey %s", msg.Nsid, msg.Rkey) } @@ -355,7 +357,8 @@ pipeline := models.Pipeline{ Rkey: msg.Rkey, Knot: source.Key(), RepoOwner: syntax.DID(record.TriggerMetadata.Repo.Did), - RepoName: record.TriggerMetadata.Repo.Repo, + RepoName: repoName, + RepoDid: repo.RepoDid, TriggerId: int(triggerId), Sha: sha, } diff --git a/appview/state/router.go b/appview/state/router.go --- a/appview/state/router.go +++ b/appview/state/router.go @@ -1,10 +1,13 @@ package state import ( + "database/sql" + "errors" "net/http" "strings" "github.com/go-chi/chi/v5" + "tangled.org/core/appview/db" "tangled.org/core/appview/issues" "tangled.org/core/appview/knots" "tangled.org/core/appview/labels" @@ -46,8 +49,29 @@ if len(pathParts) > 0 { firstPart := pathParts[0] - // if using a DID or handle, just continue as per usual - if userutil.IsDid(firstPart) || userutil.IsHandle(firstPart) { + if userutil.IsDid(firstPart) { + repo, err := db.GetRepoByDid(s.db, firstPart) + switch { + case err == nil: + remaining := "" + if len(pathParts) > 1 { + remaining = "/" + pathParts[1] + } + rewritten := "/" + repo.Did + "/" + repo.Name + remaining + r2 := r.Clone(r.Context()) + r2.URL.Path = rewritten + r2.URL.RawPath = rewritten + userRouter.ServeHTTP(w, r2) + case errors.Is(err, sql.ErrNoRows): + userRouter.ServeHTTP(w, r) + default: + s.logger.Error("db error looking up repo DID", "repoDid", firstPart, "err", err) + http.Error(w, "internal server error", http.StatusInternalServerError) + } + return + } + + if userutil.IsHandle(firstPart) { userRouter.ServeHTTP(w, r) return } diff --git a/appview/state/star.go b/appview/state/star.go --- a/appview/state/star.go +++ b/appview/state/star.go @@ -12,6 +12,7 @@ "tangled.org/core/api/tangled" "tangled.org/core/appview/db" "tangled.org/core/appview/models" "tangled.org/core/appview/pages" + "tangled.org/core/orm" "tangled.org/core/tid" ) @@ -40,15 +41,22 @@ switch r.Method { case http.MethodPost: createdAt := time.Now().Format(time.RFC3339) rkey := tid.TID() + + subjectStr := subjectUri.String() + starRecord := &tangled.FeedStar{ + CreatedAt: createdAt, + Subject: &subjectStr, + } + repo, err := db.GetRepo(s.db, orm.FilterEq("at_uri", subjectUri.String())) + if err == nil && repo.RepoDid != "" { + starRecord.SubjectDid = &repo.RepoDid + } + resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ Collection: tangled.FeedStarNSID, Repo: currentUser.Active.Did, Rkey: rkey, - Record: &lexutil.LexiconTypeDecoder{ - Val: &tangled.FeedStar{ - Subject: subjectUri.String(), - CreatedAt: createdAt, - }}, + Record: &lexutil.LexiconTypeDecoder{Val: starRecord}, }) if err != nil { log.Println("failed to create atproto record", err) diff --git a/appview/state/state.go b/appview/state/state.go --- a/appview/state/state.go +++ b/appview/state/state.go @@ -42,7 +42,7 @@ "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/atproto/syntax" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/xrpc" - securejoin "github.com/cyphar/filepath-securejoin" + "github.com/go-chi/chi/v5" "github.com/posthog/posthog-go" ) @@ -457,8 +457,46 @@ s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository by this name on %s", existingRepo.Knot)) return } - // create atproto record for this repo rkey := tid.TID() + + client, err := s.oauth.ServiceClient( + r, + oauth.WithService(domain), + oauth.WithLxm(tangled.RepoCreateNSID), + oauth.WithDev(s.config.Core.Dev), + ) + if err != nil { + l.Error("service auth failed", "err", err) + s.pages.Notice(w, "repo", "Failed to reach knot server.") + return + } + + input := &tangled.RepoCreate_Input{ + Rkey: rkey, + Name: repoName, + DefaultBranch: &defaultBranch, + } + createResp, xe := tangled.RepoCreate( + r.Context(), + client, + input, + ) + if err := xrpcclient.HandleXrpcErr(xe); err != nil { + l.Error("xrpc error", "xe", xe) + s.pages.Notice(w, "repo", err.Error()) + return + } + + var repoDid string + if createResp != nil && createResp.RepoDid != nil { + repoDid = *createResp.RepoDid + } + if repoDid == "" { + l.Error("knot returned empty repo DID") + s.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.") + return + } + repo := &models.Repo{ Did: user.Active.Did, Name: repoName, @@ -467,12 +505,48 @@ Rkey: rkey, Description: description, Created: time.Now(), Labels: s.config.Label.DefaultLabelDefs, + RepoDid: repoDid, } record := repo.AsRecord() + cleanupKnot := func() { + go func() { + delays := []time.Duration{0, 2 * time.Second, 5 * time.Second} + for attempt, delay := range delays { + time.Sleep(delay) + deleteClient, dErr := s.oauth.ServiceClient( + r, + oauth.WithService(domain), + oauth.WithLxm(tangled.RepoDeleteNSID), + oauth.WithDev(s.config.Core.Dev), + ) + if dErr != nil { + l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr) + continue + } + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{ + Did: user.Active.Did, + Name: repoName, + Rkey: rkey, + }); dErr != nil { + cancel() + l.Error("failed to clean up repo on knot after rollback", "attempt", attempt+1, "err", dErr) + continue + } + cancel() + l.Info("successfully cleaned up repo on knot after rollback", "attempt", attempt+1) + return + } + l.Error("exhausted retries for knot cleanup, repo may be orphaned", + "did", user.Active.Did, "repo", repoName, "knot", domain) + }() + } + atpClient, err := s.oauth.AuthorizedClient(r) if err != nil { l.Info("PDS write failed", "err", err) + cleanupKnot() s.pages.Notice(w, "repo", "Failed to write record to PDS.") return } @@ -487,6 +561,7 @@ }, }) if err != nil { l.Info("PDS write failed", "err", err) + cleanupKnot() s.pages.Notice(w, "repo", "Failed to announce repository creation.") return } @@ -502,51 +577,24 @@ s.pages.Notice(w, "repo", "Failed to save repository information.") return } - // The rollback function reverts a few things on failure: - // - the pending txn - // - the ACLs - // - the atproto record created rollback := func() { err1 := tx.Rollback() err2 := s.enforcer.E.LoadPolicy() err3 := rollbackRecord(context.Background(), aturi, atpClient) - // ignore txn complete errors, this is okay if errors.Is(err1, sql.ErrTxDone) { err1 = nil } if errs := errors.Join(err1, err2, err3); errs != nil { l.Error("failed to rollback changes", "errs", errs) - return } - } - defer rollback() - client, err := s.oauth.ServiceClient( - r, - oauth.WithService(domain), - oauth.WithLxm(tangled.RepoCreateNSID), - oauth.WithDev(s.config.Core.Dev), - ) - if err != nil { - l.Error("service auth failed", "err", err) - s.pages.Notice(w, "repo", "Failed to reach PDS.") - return + if aturi != "" { + cleanupKnot() + } } - - xe := tangled.RepoCreate( - r.Context(), - client, - &tangled.RepoCreate_Input{ - Rkey: rkey, - }, - ) - if err := xrpcclient.HandleXrpcErr(xe); err != nil { - l.Error("xrpc error", "xe", xe) - s.pages.Notice(w, "repo", err.Error()) - return - } + defer rollback() err = db.AddRepo(tx, repo) if err != nil { @@ -555,9 +603,8 @@ s.pages.Notice(w, "repo", "Failed to save repository information.") return } - // acls - p, _ := securejoin.SecureJoin(user.Active.Did, repoName) - err = s.enforcer.AddRepo(user.Active.Did, domain, p) + rbacPath := repo.RepoIdentifier() + err = s.enforcer.AddRepo(user.Active.Did, domain, rbacPath) if err != nil { l.Error("acl setup failed", "err", err) s.pages.Notice(w, "repo", "Failed to set up repository permissions.") @@ -578,11 +625,14 @@ http.Error(w, err.Error(), http.StatusInternalServerError) return } - // reset the ATURI because the transaction completed successfully aturi = "" s.notifier.NewRepo(r.Context(), repo) - s.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Active.Did, repoName)) + if repoDid != "" { + s.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid)) + } else { + s.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Active.Did, repoName)) + } } } diff --git a/appview/validator/label.go b/appview/validator/label.go --- a/appview/validator/label.go +++ b/appview/validator/label.go @@ -109,7 +109,7 @@ // validate permissions: only collaborators can apply labels currently // // TODO: introduce a repo:triage permission - ok, err := v.enforcer.IsPushAllowed(labelOp.Did, repo.Knot, repo.DidSlashRepo()) + ok, err := v.enforcer.IsPushAllowed(labelOp.Did, repo.Knot, repo.RepoIdentifier()) if err != nil { return fmt.Errorf("failed to enforce permissions: %w", err) } -- tangled.sh