This repository has no description
4.9 kB
169 lines
1package xrpc
2
3import (
4 "context"
5 "fmt"
6 "io"
7 "net/http"
8 "net/url"
9 "strings"
10
11 "github.com/bluesky-social/indigo/api/atproto"
12 "github.com/bluesky-social/indigo/atproto/syntax"
13 indigoxrpc "github.com/bluesky-social/indigo/xrpc"
14 "tangled.org/core/api/tangled"
15 "tangled.org/core/knotmirror/db"
16 "tangled.org/core/knotmirror/models"
17)
18
19var mirrorToKnotNSID = map[string]string{
20 tangled.GitTempListBranchesNSID: tangled.RepoBranchesNSID,
21 tangled.GitTempListTagsNSID: tangled.RepoTagsNSID,
22 tangled.GitTempListCommitsNSID: tangled.RepoLogNSID,
23 tangled.GitTempGetTreeNSID: tangled.RepoTreeNSID,
24 tangled.GitTempGetBranchNSID: tangled.RepoBranchNSID,
25 tangled.GitTempGetBlobNSID: tangled.RepoBlobNSID,
26 tangled.GitTempGetTagNSID: tangled.RepoTagNSID,
27 tangled.GitTempGetArchiveNSID: tangled.RepoArchiveNSID,
28 tangled.RepoBlobNSID: tangled.RepoBlobNSID,
29 tangled.GitTempListLanguagesNSID: tangled.RepoLanguagesNSID,
30}
31
32var hopByHopHeaders = map[string]bool{
33 "Connection": true,
34 "Keep-Alive": true,
35 "Transfer-Encoding": true,
36 "Te": true,
37 "Trailer": true,
38 "Upgrade": true,
39 "Proxy-Authorization": true,
40 "Proxy-Authenticate": true,
41}
42
43type knotInfo struct {
44 baseURL string
45 repoIdentifier string
46}
47
48func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, error) {
49 repo, err := db.GetRepoByAtUri(ctx, x.db, repoAt)
50 if err == nil && repo != nil {
51 knotURL := repo.KnotDomain
52 if !strings.Contains(repo.KnotDomain, "://") {
53 if host, _ := db.GetHost(ctx, x.db, repo.KnotDomain); host != nil {
54 knotURL = host.URL()
55 } else {
56 x.logger.Warn("repo is from unknown knot")
57 if x.cfg.KnotUseSSL {
58 knotURL = "https://" + knotURL
59 } else {
60 knotURL = "http://" + knotURL
61 }
62 }
63 }
64 return &knotInfo{baseURL: knotURL, repoIdentifier: repo.RepoIdentifier()}, nil
65 }
66
67 owner, err := x.resolver.ResolveIdent(ctx, repoAt.Authority().String())
68 if err != nil {
69 return nil, fmt.Errorf("resolving repo owner: %w", err)
70 }
71
72 xrpcc := indigoxrpc.Client{Host: owner.PDSEndpoint()}
73 out, err := atproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String())
74 if err != nil {
75 return nil, fmt.Errorf("fetching repo record from PDS: %w", err)
76 }
77
78 record := out.Value.Val.(*tangled.Repo)
79 if record.RepoDid == nil || *record.RepoDid == "" {
80 return nil, fmt.Errorf("repo record has no repo_did")
81 }
82 knotURL := record.Knot
83 if !strings.Contains(record.Knot, "://") {
84 if host, _ := db.GetHost(ctx, x.db, record.Knot); host != nil {
85 knotURL = host.URL()
86 } else {
87 x.logger.Warn("repo is from unknown knot")
88 if x.cfg.KnotUseSSL {
89 knotURL = "https://" + knotURL
90 } else {
91 knotURL = "http://" + knotURL
92 }
93 }
94 }
95
96 rkey := repoAt.RecordKey().String()
97 repoDid := syntax.DID(*record.RepoDid)
98 go func() {
99 bgCtx := context.Background()
100 pending := &models.Repo{
101 Did: owner.DID,
102 Rkey: repoAt.RecordKey(),
103 Cid: (*syntax.CID)(out.Cid),
104 Name: rkey,
105 KnotDomain: knotURL,
106 RepoDid: repoDid,
107 State: models.RepoStatePending,
108 }
109 if upsertErr := db.UpsertRepo(bgCtx, x.db, pending); upsertErr != nil {
110 x.logger.Error("failed to upsert repo after proxy resolution", "err", upsertErr)
111 }
112 }()
113
114 return &knotInfo{
115 baseURL: knotURL,
116 repoIdentifier: repoDid.String(),
117 }, nil
118}
119
120func (x *Xrpc) proxyToKnot(w http.ResponseWriter, r *http.Request, repoAt syntax.ATURI) bool {
121 mirrorNSID := strings.TrimPrefix(r.URL.Path, "/xrpc/")
122 knotNSID, ok := mirrorToKnotNSID[mirrorNSID]
123 if !ok {
124 return false
125 }
126
127 knot, err := x.resolveKnot(r.Context(), repoAt)
128 if err != nil {
129 x.logger.Warn("proxy: failed to resolve knot", "repo", repoAt, "err", err)
130 return false
131 }
132
133 params := make(url.Values)
134 for k, v := range r.URL.Query() {
135 params[k] = v
136 }
137 params.Set("repo", knot.repoIdentifier)
138
139 target := fmt.Sprintf("%s/xrpc/%s?%s", knot.baseURL, knotNSID, params.Encode())
140
141 req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, target, nil)
142 if err != nil {
143 x.logger.Warn("proxy: failed to build request", "target", target, "err", err)
144 return false
145 }
146
147 resp, err := x.httpClient.Do(req)
148 if err != nil {
149 x.logger.Warn("proxy: knot request failed", "target", target, "err", err)
150 return false
151 }
152 defer resp.Body.Close()
153
154 for k, vv := range resp.Header {
155 if hopByHopHeaders[k] {
156 continue
157 }
158 for _, v := range vv {
159 w.Header().Add(k, v)
160 }
161 }
162 w.WriteHeader(resp.StatusCode)
163 if _, err := io.Copy(w, resp.Body); err != nil {
164 x.logger.Warn("proxy: response copy interrupted", "target", target, "err", err)
165 }
166
167 x.logger.Info("proxy: served from knot", "repo", repoAt, "knot", knot.baseURL, "status", resp.StatusCode)
168 return true
169}