forked from Huawei/dockyard
/
httpsink.go
93 lines (75 loc) · 1.77 KB
/
httpsink.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
package notifications
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"sync"
"time"
)
type httpSink struct {
url string
mu sync.Mutex
closed bool
client *http.Client
}
func newHttpSink(u string, timeout time.Duration, headers http.Header) *httpSink {
return &httpSink{
url: u,
client: &http.Client{
Transport: &headerRoundTripper{
Transport: http.DefaultTransport.(*http.Transport),
headers: headers,
},
Timeout: timeout,
},
}
}
func (hs *httpSink) Write(events ...Event) Error {
hs.mu.Lock()
defer hs.mu.Unlock()
defer hs.client.Transport.(*headerRoundTripper).CloseIdleConnections()
if hs.closed {
return Error{ErrSinkClosed, http.StatusInternalServerError}
}
envelope := Envelope{
Events: events,
}
p, err := json.MarshalIndent(envelope, "", " ")
if err != nil {
return Error{fmt.Errorf("%v: error marshaling event envelope: %v", hs, err), http.StatusBadRequest}
}
body := bytes.NewReader(p)
resp, err := hs.client.Post(hs.url, EventsMediaType, body)
if err != nil {
return Error{fmt.Errorf("%v: error posting: %v", hs, err), http.StatusInternalServerError}
}
defer resp.Body.Close()
return Error{nil, resp.StatusCode}
}
func (hs *httpSink) Close() error {
hs.mu.Lock()
defer hs.mu.Unlock()
if hs.closed {
return fmt.Errorf("httpsink: already closed")
}
hs.closed = true
return nil
}
type headerRoundTripper struct {
*http.Transport
headers http.Header
}
func (hrt *headerRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
var nreq http.Request
nreq = *req
nreq.Header = make(http.Header)
merge := func(headers http.Header) {
for k, v := range headers {
nreq.Header[k] = append(nreq.Header[k], v...)
}
}
merge(req.Header)
merge(hrt.headers)
return hrt.Transport.RoundTrip(&nreq)
}