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
|
package http
import (
"bufio"
"bytes"
"context"
"fmt"
"io"
"net"
"net/http"
"github.com/go-chassis/go-chassis/core/common"
"github.com/go-chassis/go-chassis/core/handler"
"github.com/go-chassis/go-chassis/core/invocation"
"k8s.io/klog"
"github.com/kubeedge/kubeedge/edgemesh/pkg/config"
)
type HTTP struct {
Conn net.Conn
SvcNamespace string
SvcName string
Port int
req *http.Request
}
// Process handles http protocol
func (p *HTTP) Process() {
defer p.Conn.Close()
for {
// parse http request
req, err := http.ReadRequest(bufio.NewReader(p.Conn))
if err != nil {
if err == io.EOF {
klog.Infof("[EdgeMesh] http client disconnected.")
return
}
klog.Errorf("[EdgeMesh] parse http request err: %v", err)
return
}
// http: Request.RequestURI can't be set in client requests
// just reset it before transport
req.RequestURI = ""
// create invocation
inv := invocation.New(context.Background())
// set invocation
inv.MicroServiceName = req.Host
inv.SourceServiceID = ""
inv.Protocol = "rest"
inv.Strategy = config.Config.LBStrategy
inv.Args = req
inv.Reply = &http.Response{}
// create handlerchain
c, err := handler.CreateChain(common.Consumer, "http", handler.Loadbalance, handler.Transport)
if err != nil {
klog.Errorf("[EdgeMesh] create http handlerchain error: %v", err)
return
}
// start to handle
p.req = req
c.Next(inv, p.responseCallback)
}
}
// responseCallback implements http handlerchain callback
func (p *HTTP) responseCallback(data *invocation.Response) error {
var err error
var resp *http.Response
var respBytes []byte
// as a proxy server, make sure that edgemesh always response to a request
// send either the response of the real backend server or 503 back
if data.Err != nil {
klog.Errorf("[EdgeMesh] error in http handlerchain: %v", data.Err)
err = data.Err
} else {
if data.Result == nil {
klog.Errorf("[EdgeMesh] empty response from http handlerchain")
err = fmt.Errorf("empty response from http handlerchain")
} else {
var ok bool
resp, ok = data.Result.(*http.Response)
if !ok {
klog.Errorf("[EdgeMesh] http handlerchain result %+v not *http.Response type", data.Result)
err = fmt.Errorf("result not *http.Response type")
} else {
respBytes, err = httpResponseToBytes(resp)
if err != nil {
klog.Errorf("[EdgeMesh] convert http response to bytes err: %v", err)
} else {
// send response back
p.Conn.Write(respBytes)
return nil
}
}
}
}
// 503
resp = &http.Response{
Status: fmt.Sprintf("%d %s", http.StatusServiceUnavailable, err),
StatusCode: http.StatusServiceUnavailable,
Proto: p.req.Proto,
Request: p.req,
Header: make(http.Header, 0),
}
respBytes, _ = httpResponseToBytes(resp)
// send error response back
p.Conn.Write(respBytes)
return err
}
// httpResponseToBytes transforms http.Response to bytes
func httpResponseToBytes(resp *http.Response) ([]byte, error) {
buf := new(bytes.Buffer)
if resp == nil {
return nil, fmt.Errorf("http response nil")
}
err := resp.Write(buf)
if err != nil {
return nil, err
}
return buf.Bytes(), nil
}
|