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)
}