summaryrefslogtreecommitdiff
path: root/edgemesh/pkg/protocol/http/http.go
blob: a38738c2970a21845ea14de46d6cbe08b60e3132 (about) (plain)
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
}