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
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}