This repository has been archived by the owner on Dec 6, 2022. It is now read-only.
/
server.go
118 lines (93 loc) · 2.49 KB
/
server.go
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
package nrt
import (
"bufio"
"bytes"
"context"
"fmt"
"io/ioutil"
"log"
"net/http"
"strconv"
"strings"
"sync"
"github.com/nats-io/nats.go"
"github.com/nats-io/nuid"
)
// implements http.ResponseWriter
type response struct {
gotStatus bool
resp *http.Response
mu sync.Mutex
}
func (r *NatsRoundTripper) ListenAndServ(ctx context.Context, hostname string, handler http.Handler) error {
if handler == nil {
handler = http.DefaultServeMux
}
hostname = strings.ToLower(hostname)
if !strings.HasSuffix(hostname, ".nats") {
return fmt.Errorf("invalid hostname, hostnames should end in .nats")
}
subject := fmt.Sprintf("%s.requests.%s.*", r.opts.prefix, hostname)
sub, err := r.opts.conn.QueueSubscribe(subject, "nrt", func(m *nats.Msg) {
reqr := bufio.NewReader(bytes.NewBuffer(m.Data))
req, err := http.ReadRequest(reqr)
if err != nil {
return
}
if req.Host != hostname {
log.Printf("Ignoring request for host %s while serving %s", req.Host, hostname)
return
}
log.Printf("[%s] Serving %s", hostname, req.URL.String())
httpResp := &http.Response{Header: map[string][]string{}}
handler.ServeHTTP(&response{resp: httpResp}, req)
httpResp.Header.Set("NRT-ID", nuid.Next())
httpResp.Header.Set("NRT-Connected-Server", r.opts.conn.ConnectedServerName())
cid, err := r.opts.conn.GetClientID()
if err == nil {
httpResp.Header.Set("NRT-Client", strconv.Itoa(int(cid)))
}
cluster := r.opts.conn.ConnectedClusterName()
if cluster != "" {
httpResp.Header.Set("NRT-Connected-Cluster", r.opts.conn.ConnectedClusterName())
}
buf := bytes.NewBuffer([]byte{})
err = httpResp.Write(buf)
if err != nil {
return
}
data := buf.Bytes()
m.Respond(data)
})
if err != nil {
return err
}
defer sub.Unsubscribe()
log.Printf("Listening on %s", sub.Subject)
<-ctx.Done()
return nil
}
func (r *response) Header() http.Header {
return r.resp.Header
}
func (r *response) Write(body []byte) (int, error) {
bc := make([]byte, len(body))
copy(bc, body)
r.resp.Body = ioutil.NopCloser(bytes.NewBuffer(bc))
if !r.gotStatus {
r.WriteHeader(http.StatusOK)
r.resp.Proto = "HTTP/1.1"
r.resp.ProtoMinor = 1
r.resp.ProtoMajor = 1
}
if r.resp.Header.Get("Content-Type") == "" {
r.resp.Header.Set("Content-Type", http.DetectContentType(body))
}
r.resp.ContentLength = int64(len(body))
return len(body), nil
}
func (r *response) WriteHeader(status int) {
r.gotStatus = true
r.resp.Status = http.StatusText(status)
r.resp.StatusCode = status
}