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