This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knotmirror / xrpc / proxy.go
8.1 kB 278 lines
1package xrpc 2 3import ( 4 "cmp" 5 "context" 6 "errors" 7 "fmt" 8 "io" 9 "maps" 10 "net/http" 11 "net/url" 12 "path" 13 "strings" 14 15 "github.com/bluesky-social/indigo/atproto/atclient" 16 "github.com/bluesky-social/indigo/atproto/syntax" 17 indigoxrpc "github.com/bluesky-social/indigo/xrpc" 18 "github.com/go-git/go-git/v5/plumbing/filemode" 19 "github.com/samber/lo" 20 "tangled.org/core/api/tangled" 21 "tangled.org/core/knotmirror/db" 22 "tangled.org/core/knotmirror/models" 23 "tangled.org/core/repoident" 24 "tangled.org/core/repoverify" 25) 26 27var mirrorToKnotNSID = map[string]string{ 28 tangled.GitTempListBranchesNSID: tangled.RepoBranchesNSID, 29 tangled.GitTempListTagsNSID: tangled.RepoTagsNSID, 30 tangled.GitTempListCommitsNSID: tangled.RepoLogNSID, 31 tangled.GitTempGetTreeNSID: tangled.RepoTreeNSID, 32 tangled.GitTempGetBranchNSID: tangled.RepoBranchNSID, 33 tangled.GitTempGetTagNSID: tangled.RepoTagNSID, 34 tangled.GitTempGetArchiveNSID: tangled.RepoArchiveNSID, 35 tangled.GitTempListLanguagesNSID: tangled.RepoLanguagesNSID, 36 tangled.GitTempGetBlobNSID: tangled.RepoBlobNSID, 37} 38 39var hopByHopHeaders = map[string]bool{ 40 "Connection": true, 41 "Keep-Alive": true, 42 "Transfer-Encoding": true, 43 "Te": true, 44 "Trailer": true, 45 "Upgrade": true, 46 "Proxy-Authorization": true, 47 "Proxy-Authenticate": true, 48} 49 50type knotInfo struct { 51 baseURL string 52 repoIdentifier string 53} 54 55func (x *Xrpc) resolveKnot(ctx context.Context, repoDid syntax.DID) (*knotInfo, error) { 56 policy := repoident.SchemeFor(!x.cfg.KnotSSRF) 57 58 if repo, err := db.GetRepoByRepoDid(ctx, x.db, repoDid); err == nil && repo != nil { 59 knotURL := repo.KnotDomain 60 if !strings.Contains(repo.KnotDomain, "://") { 61 if host, _ := db.GetHost(ctx, x.db, repo.KnotDomain); host != nil { 62 knotURL = host.URL() 63 } else { 64 x.logger.Warn("repo is from unknown knot") 65 knotURL = lo.Ternary(x.cfg.KnotUseSSL, "https://", "http://") + knotURL 66 } 67 } 68 base, err := repoident.ParseKnotURL(knotURL, policy) 69 if err != nil { 70 return nil, err 71 } 72 return &knotInfo{baseURL: base.String(), repoIdentifier: repo.RepoIdentifier()}, nil 73 } 74 75 ident, err := x.resolver.ResolveIdent(ctx, repoDid.String()) 76 if err != nil { 77 return nil, fmt.Errorf("resolving repoDid %s: %w", repoDid, err) 78 } 79 base, err := repoident.KnotURLFromIdentity(ident, policy) 80 if err != nil { 81 return nil, fmt.Errorf("repoDid %s: %w", repoDid, err) 82 } 83 knotURL := base.String() 84 85 described, err := repoverify.Describe(ctx, x.httpClient, base, repoident.RepoDid(repoDid)) 86 if errors.Is(err, repoverify.ErrKnotAnswer) { 87 return nil, err 88 } 89 if err != nil { 90 x.logger.Warn("describeRepo failed; serving without metadata upsert", "knot", knotURL, "repo", repoDid, "err", err) 91 return &knotInfo{baseURL: knotURL, repoIdentifier: repoDid.String()}, nil 92 } 93 94 go func() { 95 pending := &models.Repo{ 96 Did: syntax.DID(described.OwnerDid), 97 Rkey: described.Rkey, 98 Name: string(described.Rkey), 99 KnotDomain: knotURL, 100 RepoDid: repoDid, 101 State: models.RepoStatePending, 102 } 103 if err := db.UpsertRepo(context.Background(), x.db, pending); err != nil { 104 x.logger.Error("failed to upsert repo after directory resolution", "err", err) 105 } 106 }() 107 108 return &knotInfo{baseURL: knotURL, repoIdentifier: repoDid.String()}, nil 109} 110 111func (x *Xrpc) proxyToKnot(w http.ResponseWriter, r *http.Request, repoDid syntax.DID) bool { 112 mirrorNSID := strings.TrimPrefix(r.URL.Path, "/xrpc/") 113 knotNSID, ok := mirrorToKnotNSID[mirrorNSID] 114 if !ok { 115 return false 116 } 117 118 knot, err := x.resolveKnot(r.Context(), repoDid) 119 if err != nil { 120 x.logger.Warn("proxy: failed to resolve knot", "repo", repoDid, "err", err) 121 return false 122 } 123 124 params := make(url.Values) 125 maps.Copy(params, r.URL.Query()) 126 params.Set("repo", knot.repoIdentifier) 127 128 target := fmt.Sprintf("%s/xrpc/%s?%s", knot.baseURL, knotNSID, params.Encode()) 129 130 req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, target, nil) 131 if err != nil { 132 x.logger.Warn("proxy: failed to build request", "target", target, "err", err) 133 return false 134 } 135 136 resp, err := x.httpClient.Do(req) 137 if err != nil { 138 x.logger.Warn("proxy: knot request failed", "target", target, "err", err) 139 return false 140 } 141 defer resp.Body.Close() 142 143 for k, vv := range resp.Header { 144 if hopByHopHeaders[k] { 145 continue 146 } 147 for _, v := range vv { 148 w.Header().Add(k, v) 149 } 150 } 151 w.WriteHeader(resp.StatusCode) 152 if _, err := io.Copy(w, resp.Body); err != nil { 153 x.logger.Warn("proxy: response copy interrupted", "target", target, "err", err) 154 } 155 156 x.logger.Info("proxy: served from knot", "repo", repoDid, "knot", knot.baseURL, "status", resp.StatusCode) 157 return true 158} 159 160func (x *Xrpc) forwardSuspended(next http.Handler) http.Handler { 161 return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { 162 repoDid, err := syntax.ParseDID(r.URL.Query().Get("repo")) 163 if err != nil { 164 next.ServeHTTP(w, r) 165 return 166 } 167 168 repo, err := db.GetRepoByRepoDid(r.Context(), x.db, repoDid) 169 if err != nil || repo == nil || repo.State != models.RepoStateSuspended { 170 next.ServeHTTP(w, r) 171 return 172 } 173 174 nsid := strings.TrimPrefix(r.URL.Path, "/xrpc/") 175 switch nsid { 176 case tangled.GitTempGetEntryNSID: 177 x.serveSuspendedEntry(w, r, repoDid) 178 case tangled.GitTempGetBlobNSID: 179 q := r.URL.Query() 180 q.Set("raw", "true") 181 r.URL.RawQuery = q.Encode() 182 x.forwardOrFail(w, r, repoDid) 183 default: 184 if _, ok := mirrorToKnotNSID[nsid]; !ok { 185 next.ServeHTTP(w, r) 186 return 187 } 188 x.forwardOrFail(w, r, repoDid) 189 } 190 }) 191} 192 193func (x *Xrpc) forwardOrFail(w http.ResponseWriter, r *http.Request, repoDid syntax.DID) { 194 if x.proxyToKnot(w, r, repoDid) { 195 return 196 } 197 writeJson(w, http.StatusBadGateway, atclient.ErrorBody{Name: "BadGateway", Message: "failed to reach knot for suspended repo"}) 198} 199 200func (x *Xrpc) serveSuspendedEntry(w http.ResponseWriter, r *http.Request, repoDid syntax.DID) { 201 ref := cmp.Or(r.URL.Query().Get("ref"), "HEAD") 202 filePath := r.URL.Query().Get("path") 203 if filePath == "" { 204 writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: "missing path parameter"}) 205 return 206 } 207 208 knot, err := x.resolveKnot(r.Context(), repoDid) 209 if err != nil { 210 x.logger.Warn("suspended entry: failed to resolve knot", "repo", repoDid, "err", err) 211 writeJson(w, http.StatusBadGateway, atclient.ErrorBody{Name: "BadGateway", Message: "failed to resolve knot for suspended repo"}) 212 return 213 } 214 215 client := &indigoxrpc.Client{Host: knot.baseURL, Client: x.httpClient} 216 out, err := tangled.RepoBlob(r.Context(), client, filePath, false, ref, knot.repoIdentifier) 217 if err != nil { 218 x.logger.Warn("suspended entry: knot repo.blob failed", "repo", repoDid, "err", err) 219 writeJson(w, http.StatusBadGateway, atclient.ErrorBody{Name: "BadGateway", Message: "failed to read entry from knot"}) 220 return 221 } 222 223 mode := filemode.Regular 224 if out.Submodule != nil { 225 mode = filemode.Submodule 226 } 227 228 writeJson(w, http.StatusOK, tangled.GitTempGetEntry_Output{ 229 Name: path.Base(filePath), 230 Mode: mode.String(), 231 Size: derefInt64(out.Size), 232 LastCommit: suspendedLastCommit(out.LastCommit), 233 Submodule: suspendedSubmodule(out.Submodule), 234 }) 235} 236 237func suspendedLastCommit(c *tangled.RepoBlob_LastCommit) *tangled.GitTempDefs_Commit { 238 if c == nil || c.Author == nil { 239 return nil 240 } 241 sig := suspendedSignature(c.Author) 242 hash := c.Hash 243 return &tangled.GitTempDefs_Commit{ 244 Author: sig, 245 Committer: sig, 246 Hash: &hash, 247 Message: c.Message, 248 } 249} 250 251func suspendedSignature(s *tangled.RepoBlob_Signature) *tangled.GitTempDefs_Signature { 252 if s == nil { 253 return nil 254 } 255 return &tangled.GitTempDefs_Signature{ 256 Name: s.Name, 257 Email: s.Email, 258 When: s.When, 259 } 260} 261 262func suspendedSubmodule(s *tangled.RepoBlob_Submodule) *tangled.GitTempDefs_Submodule { 263 if s == nil { 264 return nil 265 } 266 return &tangled.GitTempDefs_Submodule{ 267 Name: s.Name, 268 Url: s.Url, 269 Branch: s.Branch, 270 } 271} 272 273func derefInt64(v *int64) int64 { 274 if v == nil { 275 return 0 276 } 277 return *v 278}