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, + }, } }