This repository has no description
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}