From 7b3842c84ede3b02150839d65ad63d0ba3c139ba Mon Sep 17 00:00:00 2001 From: Lewis Date: Thu, 2 Apr 2026 01:44:47 +0300 Subject: [PATCH] knotmirror/proxy: proxy in case of knotmirror err We ought to harden the knotmirror in case we get errors from tap or anywhere else, in situations such that the appview knows about a git repo somewhere but the knot- mirror doesn't. I would call this good practice in general so that we have robust infrastructure and other possible future appviews would also be able to trust that the knotmirror will serve them repos that in fact exist. Lewis: May this revision serve well! --- knotmirror/xrpc/git_get_archive.go | 15 ++- knotmirror/xrpc/git_get_blob.go | 6 +- knotmirror/xrpc/git_get_branch.go | 8 +- knotmirror/xrpc/git_get_tag.go | 8 +- knotmirror/xrpc/git_get_tree.go | 8 +- knotmirror/xrpc/git_list_branches.go | 8 +- knotmirror/xrpc/git_list_commits.go | 8 +- knotmirror/xrpc/git_list_languages.go | 5 +- knotmirror/xrpc/git_list_tags.go | 8 +- knotmirror/xrpc/proxy.go | 151 ++++++++++++++++++++++++++ knotmirror/xrpc/xrpc.go | 25 +++-- 11 files changed, 210 insertions(+), 40 deletions(-) create mode 100644 knotmirror/xrpc/proxy.go diff --git a/knotmirror/xrpc/git_get_archive.go b/knotmirror/xrpc/git_get_archive.go index 16ee7f14..9541c19d 100644 --- a/knotmirror/xrpc/git_get_archive.go +++ b/knotmirror/xrpc/git_get_archive.go @@ -42,14 +42,20 @@ func (x *Xrpc) GetArchive(w http.ResponseWriter, r *http.Request) { repoPath, err := x.makeRepoPath(ctx, repo) if err != nil { - l.Error("failed to resolve repo at-uri", "err", err) + l.Warn("local mirror failed, trying proxy", "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to resolve repo"}) return } gr, err := git.Open(repoPath, ref) if err != nil { - l.Error("failed to open git repo", "err", err) + l.Warn("local mirror failed, trying proxy", "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to open git repo"}) return } @@ -65,7 +71,10 @@ func (x *Xrpc) GetArchive(w http.ResponseWriter, r *http.Request) { return r.Name, nil }() if err != nil { - l.Error("failed to get repo name", "err", err) + l.Warn("local mirror failed, trying proxy", "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to retrieve repo name"}) return } diff --git a/knotmirror/xrpc/git_get_blob.go b/knotmirror/xrpc/git_get_blob.go index a1ce934b..cb72816d 100644 --- a/knotmirror/xrpc/git_get_blob.go +++ b/knotmirror/xrpc/git_get_blob.go @@ -35,8 +35,10 @@ func (x *Xrpc) GetBlob(w http.ResponseWriter, r *http.Request) { file, err := x.getFile(r.Context(), repo, ref, path) if err != nil { - // TODO: better error return - l.Error("failed to get blob", "err", err) + l.Warn("local mirror failed, trying proxy", "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to get blob"}) return } diff --git a/knotmirror/xrpc/git_get_branch.go b/knotmirror/xrpc/git_get_branch.go index b66d2bbf..2a30bd76 100644 --- a/knotmirror/xrpc/git_get_branch.go +++ b/knotmirror/xrpc/git_get_branch.go @@ -33,12 +33,12 @@ func (x *Xrpc) GetBranch(w http.ResponseWriter, r *http.Request) { } branchName, _ := url.PathUnescape(nameQuery) - l := x.logger.With("repo", repo, "branch", branchName) - out, err := x.getBranch(r.Context(), repo, branchName) if err != nil { - // TODO: better error return - l.Error("failed to get branch", "err", err) + x.logger.Warn("local mirror failed, trying proxy", "repo", repo, "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to get branch"}) return } diff --git a/knotmirror/xrpc/git_get_tag.go b/knotmirror/xrpc/git_get_tag.go index 00cf0193..698a1f8c 100644 --- a/knotmirror/xrpc/git_get_tag.go +++ b/knotmirror/xrpc/git_get_tag.go @@ -30,12 +30,12 @@ func (x *Xrpc) GetTag(w http.ResponseWriter, r *http.Request) { return } - l := x.logger.With("repo", repo, "tag", tagName) - out, err := x.getTag(r.Context(), repo, tagName) if err != nil { - // TODO: better error return - l.Error("failed to get tag", "err", err) + x.logger.Warn("local mirror failed, trying proxy", "repo", repo, "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to get tag"}) return } diff --git a/knotmirror/xrpc/git_get_tree.go b/knotmirror/xrpc/git_get_tree.go index a2742116..6ef142f9 100644 --- a/knotmirror/xrpc/git_get_tree.go +++ b/knotmirror/xrpc/git_get_tree.go @@ -28,12 +28,12 @@ func (x *Xrpc) GetTree(w http.ResponseWriter, r *http.Request) { return } - l := x.logger.With("repo", repo, "ref", ref, "path", path) - out, err := x.getTree(r.Context(), repo, ref, path) if err != nil { - // TODO: better error return - l.Error("failed to get tree", "err", err) + x.logger.Warn("local mirror failed, trying proxy", "repo", repo, "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to get tree"}) return } diff --git a/knotmirror/xrpc/git_list_branches.go b/knotmirror/xrpc/git_list_branches.go index b7b224d4..28604972 100644 --- a/knotmirror/xrpc/git_list_branches.go +++ b/knotmirror/xrpc/git_list_branches.go @@ -44,12 +44,12 @@ func (x *Xrpc) ListBranches(w http.ResponseWriter, r *http.Request) { } } - l := x.logger.With("repo", repoQuery, "limit", limit, "cursor", cursor) - out, err := x.listBranches(r.Context(), repo, limit, cursor) if err != nil { - // TODO: better error return - l.Error("failed to list branches", "err", err) + x.logger.Warn("local mirror failed, trying proxy", "repo", repo, "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to list branches"}) return } diff --git a/knotmirror/xrpc/git_list_commits.go b/knotmirror/xrpc/git_list_commits.go index 5efcb744..87a9d7c7 100644 --- a/knotmirror/xrpc/git_list_commits.go +++ b/knotmirror/xrpc/git_list_commits.go @@ -44,12 +44,12 @@ func (x *Xrpc) ListCommits(w http.ResponseWriter, r *http.Request) { } } - l := x.logger.With("repo", repo, "ref", ref) - out, err := x.listCommits(r.Context(), repo, ref, limit, cursor) if err != nil { - // TODO: better error return - l.Error("failed to list commits", "err", err) + x.logger.Warn("local mirror failed, trying proxy", "repo", repo, "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to list commits"}) return } diff --git a/knotmirror/xrpc/git_list_languages.go b/knotmirror/xrpc/git_list_languages.go index 040a422f..e6a60539 100644 --- a/knotmirror/xrpc/git_list_languages.go +++ b/knotmirror/xrpc/git_list_languages.go @@ -29,7 +29,10 @@ func (x *Xrpc) ListLanguages(w http.ResponseWriter, r *http.Request) { out, err := x.listLanguages(r.Context(), repo, ref) if err != nil { - l.Error("failed to list languages", "err", err) + l.Warn("local mirror failed, trying proxy", "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeErr(w, err) return } diff --git a/knotmirror/xrpc/git_list_tags.go b/knotmirror/xrpc/git_list_tags.go index 553352ab..8661e4b5 100644 --- a/knotmirror/xrpc/git_list_tags.go +++ b/knotmirror/xrpc/git_list_tags.go @@ -45,12 +45,12 @@ func (x *Xrpc) ListTags(w http.ResponseWriter, r *http.Request) { } } - l := x.logger.With("repo", repo, "limit", limit, "cursor", cursor) - out, err := x.listTags(r.Context(), repo, limit, cursor) if err != nil { - // TODO: better error return - l.Error("failed to list tags", "err", err) + x.logger.Warn("local mirror failed, trying proxy", "repo", repo, "err", err) + if x.proxyToKnot(w, r, repo) { + return + } writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to list tags"}) return } diff --git a/knotmirror/xrpc/proxy.go b/knotmirror/xrpc/proxy.go new file mode 100644 index 00000000..bc2b60aa --- /dev/null +++ b/knotmirror/xrpc/proxy.go @@ -0,0 +1,151 @@ +package xrpc + +import ( + "context" + "fmt" + "io" + "net/http" + "net/url" + "strings" + + "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" + indigoxrpc "github.com/bluesky-social/indigo/xrpc" + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/models" +) + +var mirrorToKnotNSID = map[string]string{ + tangled.GitTempListBranchesNSID: tangled.RepoBranchesNSID, + tangled.GitTempListTagsNSID: tangled.RepoTagsNSID, + tangled.GitTempListCommitsNSID: tangled.RepoLogNSID, + tangled.GitTempGetTreeNSID: tangled.RepoTreeNSID, + tangled.GitTempGetBranchNSID: tangled.RepoBranchNSID, + tangled.GitTempGetBlobNSID: tangled.RepoBlobNSID, + tangled.GitTempGetTagNSID: tangled.RepoTagNSID, + tangled.GitTempGetArchiveNSID: tangled.RepoArchiveNSID, + tangled.GitTempListLanguagesNSID: tangled.RepoLanguagesNSID, +} + +var hopByHopHeaders = map[string]bool{ + "Connection": true, + "Keep-Alive": true, + "Transfer-Encoding": true, + "Te": true, + "Trailer": true, + "Upgrade": true, + "Proxy-Authorization": true, + "Proxy-Authenticate": true, +} + +type knotInfo struct { + baseURL string + didSlashRepo string +} + +func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, error) { + repo, err := db.GetRepoByAtUri(ctx, x.db, repoAt) + if err == nil && repo != nil { + if repo.State != models.RepoStatePending && repo.State != models.RepoStateResyncing { + go func() { + if err := db.UpdateRepoState(context.Background(), x.db, repo.Did, repo.Rkey, models.RepoStatePending); err != nil { + x.logger.Error("failed to mark repo for resync after proxy", "err", err) + } + }() + } + return &knotInfo{baseURL: repo.KnotDomain, didSlashRepo: repo.DidSlashRepo()}, nil + } + + owner, err := x.resolver.ResolveIdent(ctx, repoAt.Authority().String()) + if err != nil { + return nil, fmt.Errorf("resolving repo owner: %w", err) + } + + xrpcc := indigoxrpc.Client{Host: owner.PDSEndpoint()} + out, err := atproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) + if err != nil { + return nil, fmt.Errorf("fetching repo record from PDS: %w", err) + } + + record := out.Value.Val.(*tangled.Repo) + knotURL := record.Knot + if !strings.Contains(knotURL, "://") { + scheme := "http" + if x.cfg.KnotUseSSL { + scheme = "https" + } + knotURL = scheme + "://" + knotURL + } + + go func() { + bgCtx := context.Background() + pending := &models.Repo{ + Did: owner.DID, + Rkey: repoAt.RecordKey(), + Cid: (*syntax.CID)(out.Cid), + Name: record.Name, + KnotDomain: knotURL, + State: models.RepoStatePending, + } + if upsertErr := db.UpsertRepo(bgCtx, x.db, pending); upsertErr != nil { + x.logger.Error("failed to upsert repo after proxy resolution", "err", upsertErr) + } + }() + + return &knotInfo{ + baseURL: knotURL, + didSlashRepo: fmt.Sprintf("%s/%s", owner.DID, record.Name), + }, nil +} + +func (x *Xrpc) proxyToKnot(w http.ResponseWriter, r *http.Request, repoAt syntax.ATURI) bool { + mirrorNSID := strings.TrimPrefix(r.URL.Path, "/xrpc/") + knotNSID, ok := mirrorToKnotNSID[mirrorNSID] + if !ok { + return false + } + + knot, err := x.resolveKnot(r.Context(), repoAt) + if err != nil { + x.logger.Warn("proxy: failed to resolve knot", "repo", repoAt, "err", err) + return false + } + + params := make(url.Values) + for k, v := range r.URL.Query() { + params[k] = v + } + params.Set("repo", knot.didSlashRepo) + + target := fmt.Sprintf("%s/xrpc/%s?%s", knot.baseURL, knotNSID, params.Encode()) + + req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, target, nil) + if err != nil { + x.logger.Warn("proxy: failed to build request", "target", target, "err", err) + return false + } + + resp, err := x.httpClient.Do(req) + if err != nil { + x.logger.Warn("proxy: knot request failed", "target", target, "err", err) + return false + } + defer resp.Body.Close() + + for k, vv := range resp.Header { + if hopByHopHeaders[k] { + continue + } + for _, v := range vv { + w.Header().Add(k, v) + } + } + w.WriteHeader(resp.StatusCode) + if _, err := io.Copy(w, resp.Body); err != nil { + x.logger.Warn("proxy: response copy interrupted", "target", target, "err", err) + } + + x.logger.Info("proxy: served from knot", "repo", repoAt, "knot", knot.baseURL, "status", resp.StatusCode) + return true +} diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go index 5239ccef..c5da80ca 100644 --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -6,6 +6,7 @@ import ( "errors" "log/slog" "net/http" + "time" "github.com/bluesky-social/indigo/atproto/atclient" "github.com/go-chi/chi/v5" @@ -17,20 +18,24 @@ import ( ) type Xrpc struct { - cfg *config.Config - db *sql.DB - resolver *idresolver.Resolver - ks *knotstream.KnotStream - logger *slog.Logger + cfg *config.Config + db *sql.DB + resolver *idresolver.Resolver + ks *knotstream.KnotStream + logger *slog.Logger + httpClient *http.Client } func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, resolver *idresolver.Resolver, ks *knotstream.KnotStream) *Xrpc { return &Xrpc{ - cfg, - db, - resolver, - ks, - log.SubLogger(logger, "xrpc"), + cfg: cfg, + db: db, + resolver: resolver, + ks: ks, + logger: log.SubLogger(logger, "xrpc"), + httpClient: &http.Client{ + Timeout: 30 * time.Second, + }, } } -- 2.51.2