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