This repository has no description
4.7 kB
157 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.GitTempListLanguagesNSID: tangled.RepoLanguagesNSID,
29}
30
31var hopByHopHeaders = map[string]bool{
32 "Connection": true,
33 "Keep-Alive": true,
34 "Transfer-Encoding": true,
35 "Te": true,
36 "Trailer": true,
37 "Upgrade": true,
38 "Proxy-Authorization": true,
39 "Proxy-Authenticate": true,
40}
41
42type knotInfo struct {
43 baseURL string
44 didSlashRepo string
45}
46
47func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, error) {
48 repo, err := db.GetRepoByAtUri(ctx, x.db, repoAt)
49 if err == nil && repo != nil {
50 if repo.State != models.RepoStatePending && repo.State != models.RepoStateResyncing {
51 go func() {
52 if err := db.UpdateRepoState(context.Background(), x.db, repo.Did, repo.Rkey, models.RepoStatePending); err != nil {
53 x.logger.Error("failed to mark repo for resync after proxy", "err", err)
54 }
55 }()
56 }
57 return &knotInfo{baseURL: repo.KnotDomain, didSlashRepo: repo.DidSlashRepo()}, nil
58 }
59
60 owner, err := x.resolver.ResolveIdent(ctx, repoAt.Authority().String())
61 if err != nil {
62 return nil, fmt.Errorf("resolving repo owner: %w", err)
63 }
64
65 xrpcc := indigoxrpc.Client{Host: owner.PDSEndpoint()}
66 out, err := atproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String())
67 if err != nil {
68 return nil, fmt.Errorf("fetching repo record from PDS: %w", err)
69 }
70
71 record := out.Value.Val.(*tangled.Repo)
72 knotURL := record.Knot
73 if !strings.Contains(record.Knot, "://") {
74 if host, _ := db.GetHost(ctx, x.db, record.Knot); host != nil {
75 knotURL = host.URL()
76 } else {
77 x.logger.Warn("repo is from unknown knot")
78 if x.cfg.KnotUseSSL {
79 knotURL = "https://" + knotURL
80 } else {
81 knotURL = "http://" + knotURL
82 }
83 }
84 }
85
86 go func() {
87 bgCtx := context.Background()
88 pending := &models.Repo{
89 Did: owner.DID,
90 Rkey: repoAt.RecordKey(),
91 Cid: (*syntax.CID)(out.Cid),
92 Name: record.Name,
93 KnotDomain: knotURL,
94 State: models.RepoStatePending,
95 }
96 x.logger.Debug("pending: upserting repo with knot", "knot", pending.KnotDomain)
97 if upsertErr := db.UpsertRepo(bgCtx, x.db, pending); upsertErr != nil {
98 x.logger.Error("failed to upsert repo after proxy resolution", "err", upsertErr)
99 }
100 }()
101
102 return &knotInfo{
103 baseURL: knotURL,
104 didSlashRepo: fmt.Sprintf("%s/%s", owner.DID, record.Name),
105 }, nil
106}
107
108func (x *Xrpc) proxyToKnot(w http.ResponseWriter, r *http.Request, repoAt syntax.ATURI) bool {
109 mirrorNSID := strings.TrimPrefix(r.URL.Path, "/xrpc/")
110 knotNSID, ok := mirrorToKnotNSID[mirrorNSID]
111 if !ok {
112 return false
113 }
114
115 knot, err := x.resolveKnot(r.Context(), repoAt)
116 if err != nil {
117 x.logger.Warn("proxy: failed to resolve knot", "repo", repoAt, "err", err)
118 return false
119 }
120
121 params := make(url.Values)
122 for k, v := range r.URL.Query() {
123 params[k] = v
124 }
125 params.Set("repo", knot.didSlashRepo)
126
127 target := fmt.Sprintf("%s/xrpc/%s?%s", knot.baseURL, knotNSID, params.Encode())
128
129 req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, target, nil)
130 if err != nil {
131 x.logger.Warn("proxy: failed to build request", "target", target, "err", err)
132 return false
133 }
134
135 resp, err := x.httpClient.Do(req)
136 if err != nil {
137 x.logger.Warn("proxy: knot request failed", "target", target, "err", err)
138 return false
139 }
140 defer resp.Body.Close()
141
142 for k, vv := range resp.Header {
143 if hopByHopHeaders[k] {
144 continue
145 }
146 for _, v := range vv {
147 w.Header().Add(k, v)
148 }
149 }
150 w.WriteHeader(resp.StatusCode)
151 if _, err := io.Copy(w, resp.Body); err != nil {
152 x.logger.Warn("proxy: response copy interrupted", "target", target, "err", err)
153 }
154
155 x.logger.Info("proxy: served from knot", "repo", repoAt, "knot", knot.baseURL, "status", resp.StatusCode)
156 return true
157}