-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathhttp.go
More file actions
314 lines (264 loc) · 8.28 KB
/
Copy pathhttp.go
File metadata and controls
314 lines (264 loc) · 8.28 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
package geecache
import (
"fmt"
"geecache/consistenthash"
pb "geecache/geecachepb"
"io"
"log"
"math/rand"
"net/http"
"net/url"
"strings"
"sync"
"crypto/tls"
"crypto/x509"
"github.com/golang/protobuf/proto"
)
const (
defaultBasePath = "/_geecache/" // 默认路径
defaultReplicas = 50
)
// HTTPPool implements PeerPicker for a pool of HTTP peers.
type HTTPPool struct {
// this peer's base URL, e.g. "https://example.net:8000"
self string
basePath string
mu sync.Mutex // guards peers and httpGetters
peers *consistenthash.Map
httpGetters map[string]*httpGetter // keyed by e.g. "http://10.0.0.2:8008"
client *http.Client // 支持自定义 TLS Client
}
// NewHTTPPool initializes an HTTP pool of peers.
func NewHTTPPool(self string) *HTTPPool {
return &HTTPPool{
self: self, // 本机地址
basePath: defaultBasePath, // 默认路径
}
}
func NewHTTPPoolWithTLS(self, caFile string) *HTTPPool {
client := newTLSClient(caFile)
return &HTTPPool{self: self, basePath: defaultBasePath, client: client}
}
func newTLSClient(caFile string) *http.Client {
pool := x509.NewCertPool()
pem, err := ioutil.ReadFile(caFile)
pool.AppendCertsFromPEM(pem)
tr := &http.Transport{
TLSClientConfig: &tls.Config{
RootCAs: pool,
MinVersion: tls.VersionTLS12,
},
}
return &http.Client{Transport: tr}
}
// Log info with server name
func (p *HTTPPool) Log(format string, v ...interface{}) {
log.Printf("[Server %s] %s", p.self, fmt.Sprintf(format, v...))
}
// ServeHTTP handle all http requests
func (p *HTTPPool) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if !strings.HasPrefix(r.URL.Path, p.basePath) {
panic("HTTPPool serving unexpected path: " + r.URL.Path)
}
p.Log("%s %s", r.Method, r.URL.Path) // 打印请求方法和路径
// /<basepath>/<groupname>/<key> required
parts := strings.SplitN(r.URL.Path[len(p.basePath):], "/", 2)
if len(parts) != 2 {
http.Error(w, "bad request", http.StatusBadRequest)
return
}
groupName := parts[0]
key := parts[1]
group := GetGroup(groupName) // 获取指定组
if group == nil {
http.Error(w, "no such group: "+groupName, http.StatusNotFound)
return
}
// 处理不同的HTTP方法
switch r.Method {
case http.MethodGet:
// 处理GET请求,获取缓存数据
view, err := group.Get(key) // 从指定组中获取指定值
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Write the value to the response body as a proto message.
body, err := proto.Marshal(&pb.Response{Value: view.ByteSlice()})
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/octet-stream")
w.Write(body)
case http.MethodPut:
// 处理PUT请求,存储热点数据
body, err := io.ReadAll(r.Body)
if err != nil {
http.Error(w, fmt.Sprintf("reading request body: %v", err), http.StatusBadRequest)
return
}
defer r.Body.Close()
// 解析请求体中的protobuf数据
res := &pb.Response{}
if err = proto.Unmarshal(body, res); err != nil {
http.Error(w, fmt.Sprintf("decoding request body: %v", err), http.StatusBadRequest)
return
}
// 将数据添加到本地缓存
value := ByteView{b: cloneBytes(res.Value)}
group.mainCache.add(key, value)
p.Log("Stored hot spot data for group=%s, key=%s", groupName, key)
w.WriteHeader(http.StatusOK)
default:
w.Header().Set("Allow", "GET, PUT")
http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
}
}
// Set updates the pool's list of peers.
func (p *HTTPPool) Set(peers ...string) {
p.mu.Lock()
defer p.mu.Unlock()
p.peers = consistenthash.New(defaultReplicas, nil) // 创建一个一致性哈希
p.peers.Add(peers...) // 添加节点到一致性哈希
p.httpGetters = make(map[string]*httpGetter, len(peers))
for _, peer := range peers {
// 为每个节点创建一个 HTTP 客户端,地址:http://10.0.0.2:8008/_geecache/
// p.httpGetters[peer] = &httpGetter{baseURL: peer + p.basePath}
p.httpGetters[peer] = &httpGetter{baseURL: peer + p.basePath, client: p.client}
}
}
// PickPeer picks a peer according to key
func (p *HTTPPool) PickPeer(key string) (PeerGetter, bool) {
p.mu.Lock()
defer p.mu.Unlock()
// 确保p.peers已经初始化
if p.peers == nil {
log.Println("PickPeer() called but peers not properly initialized")
return nil, false
}
if peer := p.peers.Get(key); peer != "" && peer != p.self {
// 通过一致性哈希(节点负荷平衡)找到该值(应该)存储的节点
p.Log("Pick peer %s", peer)
return p.httpGetters[peer], true
} else if peer == p.self {
p.Log("Pick self %s", peer)
return nil, false
}
return nil, false
}
// PickPeers picks multiple peers for hot spot data backup
func (p *HTTPPool) PickPeers(key string, count int) ([]PeerGetter, bool) {
p.mu.Lock()
defer p.mu.Unlock()
// 确保p.peers已经初始化
if p.peers == nil {
log.Println("PickPeers() called but peers not properly initialized")
return nil, false
}
// 获取主节点
mainPeer := p.peers.Get(key)
if mainPeer == "" || mainPeer == p.self {
return nil, false
}
// 收集所有可用的节点(除了自己)
var availablePeers []string
for peer := range p.httpGetters {
if peer != p.self {
availablePeers = append(availablePeers, peer)
}
}
// 如果可用节点数量不足,返回所有可用节点
if len(availablePeers) <= count {
peers := make([]PeerGetter, 0, len(availablePeers))
for _, peer := range availablePeers {
peers = append(peers, p.httpGetters[peer])
}
p.Log("Pick all %d available peers for hot spot data", len(peers))
return peers, len(peers) > 0
}
// 否则,选择指定数量的节点(包括主节点)
peers := make([]PeerGetter, 0, count)
// 首先添加主节点
peers = append(peers, p.httpGetters[mainPeer])
// 随机选择其他节点
for i := 0; i < count-1 && len(availablePeers) > 0; i++ {
// 简单随机选择一个节点
index := rand.Intn(len(availablePeers))
peer := availablePeers[index]
// 如果不是主节点,添加到结果中
if peer != mainPeer {
peers = append(peers, p.httpGetters[peer])
}
// 从可用节点中移除已选择的节点
availablePeers = append(availablePeers[:index], availablePeers[index+1:]...)
}
p.Log("Pick %d peers for hot spot data", len(peers))
return peers, true
}
var _ PeerPicker = (*HTTPPool)(nil)
type httpGetter struct {
baseURL string
client *http.Client
}
func (h *httpGetter) Get(in *pb.Request, out *pb.Response) error {
// 构建请求的 URL:http://10.0.0.2:8008/_geecache/<groupname>/<key>
u := fmt.Sprintf(
"%v%v/%v",
h.baseURL,
url.QueryEscape(in.GetGroup()),
url.QueryEscape(in.GetKey()),
)
// res, err := http.Get(u)
// 发送 HTTP 请求给该地址的 HTTP 服务端,由ServeHttp来处理
res, err := h.client.Get(u)
if err != nil {
return err
}
defer res.Body.Close()
if res.StatusCode != http.StatusOK {
return fmt.Errorf("server returned: %v", res.Status)
}
bytes, err := io.ReadAll(res.Body)
if err != nil {
return fmt.Errorf("reading response body: %v", err)
}
if err = proto.Unmarshal(bytes, out); err != nil {
return fmt.Errorf("decoding response body: %v", err)
}
return nil
}
// Set sends a PUT request to store a value for a key in remote peer
func (h *httpGetter) Set(in *pb.Request, out *pb.Response) error {
// 构建请求的 URL:http://10.0.0.2:8008/_geecache/<groupname>/<key>
u := fmt.Sprintf(
"%v%v/%v",
h.baseURL,
url.QueryEscape(in.GetGroup()),
url.QueryEscape(in.GetKey()),
)
// 将响应数据序列化为protobuf
body, err := proto.Marshal(out)
if err != nil {
return fmt.Errorf("encoding request body: %v", err)
}
// 创建PUT请求
req, err := http.NewRequest(http.MethodPut, u, strings.NewReader(string(body)))
if err != nil {
return fmt.Errorf("creating request: %v", err)
}
req.Header.Set("Content-Type", "application/octet-stream")
// 发送请求
// client := &http.Client{}
// res, err := client.Do(req)
res, err := h.client.Do(req)
if err != nil {
return fmt.Errorf("sending request: %v", err)
}
defer res.Body.Close()
if res.StatusCode != http.StatusOK {
return fmt.Errorf("server returned: %v", res.Status)
}
return nil
}
var _ PeerGetter = (*httpGetter)(nil)