Merge pull request #216 from cloudflare/mapper-refactor

Mapper package refactoring
This commit is contained in:
Łukasz Mierzwa
2018-01-22 11:39:19 -08:00
committed by GitHub
18 changed files with 247 additions and 136 deletions
+4 -1
View File
@@ -17,7 +17,10 @@ func init() {
log.SetLevel(log.ErrorLevel)
for i, uri := range mock.ListAllMockURIs() {
name := fmt.Sprintf("dedup-mock-%d", i)
am := alertmanager.NewAlertmanager(name, uri, alertmanager.WithRequestTimeout(time.Second))
am, err := alertmanager.NewAlertmanager(name, uri, alertmanager.WithRequestTimeout(time.Second))
if err != nil {
log.Fatal(err)
}
alertmanager.RegisterAlertmanager(am)
}
}
+49 -4
View File
@@ -1,6 +1,7 @@
package alertmanager
import (
"encoding/json"
"fmt"
"path"
"sort"
@@ -34,6 +35,8 @@ type Alertmanager struct {
Name string `json:"name"`
// whenever this instance should be proxied
ProxyRequests bool
// transport instances are specific to URI scheme we collect from
transport transport.Transport
// lock protects data access while updating
lock sync.RWMutex
// fields for storing pulled data
@@ -56,9 +59,19 @@ func (am *Alertmanager) detectVersion() string {
return defaultVersion
}
ver := alertmanagerVersion{}
err = transport.ReadJSON(url, am.RequestTimeout, &ver)
// read raw body from the source
source, err := am.transport.Read(url)
defer source.Close()
if err != nil {
log.Errorf("[%s] %s request failed: %s", am.Name, url, err.Error())
log.Errorf("[%s] %s request failed: %s", am.Name, url, err)
return defaultVersion
}
// decode body as JSON
err = json.NewDecoder(source).Decode(&ver)
if err != nil {
log.Errorf("[%s] %s failed to decode as JSON: %s", am.Name, url, err)
return defaultVersion
}
@@ -91,8 +104,24 @@ func (am *Alertmanager) pullSilences(version string) error {
return err
}
// generate full URL to collect silences from
url, err := mapper.AbsoluteURL(am.URI)
if err != nil {
log.Errorf("[%s] Failed to generate silences endpoint URL: %s", am.Name, err)
return err
}
start := time.Now()
silences, err := mapper.GetSilences(am.URI, am.RequestTimeout)
// read raw body from the source
source, err := am.transport.Read(url)
defer source.Close()
if err != nil {
log.Errorf("[%s] %s request failed: %s", am.Name, url, err)
return err
}
// decode body text
silences, err := mapper.Decode(source)
if err != nil {
return err
}
@@ -134,8 +163,24 @@ func (am *Alertmanager) pullAlerts(version string) error {
return err
}
// generate full URL to collect alerts from
url, err := mapper.AbsoluteURL(am.URI)
if err != nil {
log.Errorf("[%s] Failed to generate alerts endpoint URL: %s", am.Name, err)
return err
}
start := time.Now()
groups, err := mapper.GetAlerts(am.URI, am.RequestTimeout)
// read raw body from the source
source, err := am.transport.Read(url)
defer source.Close()
if err != nil {
log.Errorf("[%s] %s request failed: %s", am.Name, url, err)
return err
}
// decode body text
groups, err := mapper.Decode(source)
if err != nil {
return err
}
+9 -2
View File
@@ -6,6 +6,7 @@ import (
"time"
"github.com/cloudflare/unsee/internal/models"
"github.com/cloudflare/unsee/internal/transport"
log "github.com/sirupsen/logrus"
)
@@ -18,7 +19,7 @@ var (
)
// NewAlertmanager creates a new Alertmanager instance
func NewAlertmanager(name, uri string, opts ...Option) *Alertmanager {
func NewAlertmanager(name, uri string, opts ...Option) (*Alertmanager, error) {
am := &Alertmanager{
URI: uri,
RequestTimeout: time.Second * 10,
@@ -40,7 +41,13 @@ func NewAlertmanager(name, uri string, opts ...Option) *Alertmanager {
opt(am)
}
return am
var err error
am.transport, err = transport.NewTransport(am.URI, am.RequestTimeout)
if err != nil {
return am, err
}
return am, nil
}
// RegisterAlertmanager will add an Alertmanager instance to the list of
+10 -1
View File
@@ -1,6 +1,7 @@
package alertmanager
import (
"encoding/json"
"time"
"github.com/cloudflare/unsee/internal/transport"
@@ -30,7 +31,15 @@ func GetVersion(uri string, timeout time.Duration) string {
return defaultVersion
}
ver := alertmanagerVersion{}
err = transport.ReadJSON(url, timeout, &ver)
t, err := transport.NewTransport(uri, timeout)
if err != nil {
log.Errorf("Unable to get the version information from %s", url)
return defaultVersion
}
source, err := t.Read(url)
err = json.NewDecoder(source).Decode(&ver)
if err != nil {
log.Errorf("%s request failed: %s", url, err.Error())
return defaultVersion
+4 -1
View File
@@ -485,7 +485,10 @@ var tests = []filterTest{
func TestFilters(t *testing.T) {
log.SetLevel(log.ErrorLevel)
am := alertmanager.NewAlertmanager("test", "http://localhost", alertmanager.WithRequestTimeout(time.Second))
am, err := alertmanager.NewAlertmanager("test", "http://localhost", alertmanager.WithRequestTimeout(time.Second))
if err != nil {
t.Error(err)
}
for _, ft := range tests {
alert := models.Alert(ft.Alert)
if &ft.Silence != nil {
+16 -13
View File
@@ -2,7 +2,7 @@ package mapper
import (
"fmt"
"time"
"io"
"github.com/cloudflare/unsee/internal/models"
)
@@ -12,11 +12,22 @@ var (
silenceMappers = []SilenceMapper{}
)
// AlertMapper implements Alertmanager -> unsee alert data mapping that works
// for a specific range of Alertmanager versions
type AlertMapper interface {
// Mapper converts Alertmanager response body and maps to unsee data structures
type Mapper interface {
IsSupported(version string) bool
GetAlerts(uri string, timeout time.Duration) ([]models.AlertGroup, error)
AbsoluteURL(baseURI string) (string, error)
}
// AlertMapper handles mapping of Alertmanager alert information to unsee AlertGroup models
type AlertMapper interface {
Mapper
Decode(io.ReadCloser) ([]models.AlertGroup, error)
}
// SilenceMapper handles mapping of Alertmanager silence information to unsee Silence models
type SilenceMapper interface {
Mapper
Decode(io.ReadCloser) ([]models.Silence, error)
}
// RegisterAlertMapper allows to register mapper implementing alert data
@@ -35,14 +46,6 @@ func GetAlertMapper(version string) (AlertMapper, error) {
return nil, fmt.Errorf("Can't find alert mapper for Alertmanager %s", version)
}
// SilenceMapper implements Alertmanager -> unsee silence data mapping that
// works for a specific range of Alertmanager versions
type SilenceMapper interface {
Release() string
IsSupported(version string) bool
GetSilences(uri string, timeout time.Duration) ([]models.Silence, error)
}
// RegisterSilenceMapper allows to register mapper implementing silence data
// handling for specific Alertmanager versions
func RegisterSilenceMapper(m SilenceMapper) {
+11 -9
View File
@@ -5,7 +5,9 @@
package v04
import (
"encoding/json"
"errors"
"io"
"sort"
"strconv"
"time"
@@ -52,25 +54,25 @@ type AlertMapper struct {
mapper.AlertMapper
}
// AbsoluteURL for alerts API endpoint this mapper supports
func (m AlertMapper) AbsoluteURL(baseURI string) (string, error) {
return transport.JoinURL(baseURI, "api/v1/alerts/groups")
}
// IsSupported returns true if given version string is supported
func (m AlertMapper) IsSupported(version string) bool {
versionRange := semver.MustParseRange(">=0.4.0 <0.5.0")
return versionRange(semver.MustParse(version))
}
// GetAlerts will make a request to Alertmanager API and parse the response
// It will only return alerts or error (if any)
func (m AlertMapper) GetAlerts(uri string, timeout time.Duration) ([]models.AlertGroup, error) {
// Decode Alertmanager API response body and return unsee model instances
func (m AlertMapper) Decode(source io.ReadCloser) ([]models.AlertGroup, error) {
groups := []models.AlertGroup{}
receivers := map[string]alertsGroupReceiver{}
resp := alertsGroupsAPISchema{}
url, err := transport.JoinURL(uri, "api/v1/alerts/groups")
if err != nil {
return groups, err
}
err = transport.ReadJSON(url, timeout, &resp)
defer source.Close()
err := json.NewDecoder(source).Decode(&resp)
if err != nil {
return groups, err
}
+11 -13
View File
@@ -5,9 +5,9 @@
package v04
import (
"encoding/json"
"errors"
"fmt"
"math"
"io"
"strconv"
"time"
@@ -47,26 +47,24 @@ type SilenceMapper struct {
mapper.SilenceMapper
}
// AbsoluteURL for silences API endpoint this mapper supports
func (m SilenceMapper) AbsoluteURL(baseURI string) (string, error) {
return transport.JoinURL(baseURI, "api/v1/silences")
}
// IsSupported returns true if given version string is supported
func (m SilenceMapper) IsSupported(version string) bool {
versionRange := semver.MustParseRange(">=0.4.0 <0.5.0")
return versionRange(semver.MustParse(version))
}
// GetSilences will make a request to Alertmanager API and parse the response
// It will only return silences or error (if any)
func (m SilenceMapper) GetSilences(uri string, timeout time.Duration) ([]models.Silence, error) {
// Decode Alertmanager API response body and return unsee model instances
func (m SilenceMapper) Decode(source io.ReadCloser) ([]models.Silence, error) {
silences := []models.Silence{}
resp := silenceAPISchema{}
url, err := transport.JoinURL(uri, "api/v1/silences")
if err != nil {
return silences, err
}
// Alertmanager 0.4 uses pagination for silences
url = fmt.Sprintf("%s?limit=%d", url, math.MaxUint32)
err = transport.ReadJSON(url, timeout, &resp)
defer source.Close()
err := json.NewDecoder(source).Decode(&resp)
if err != nil {
return silences, err
}
+11 -9
View File
@@ -5,7 +5,9 @@
package v05
import (
"encoding/json"
"errors"
"io"
"sort"
"time"
@@ -51,25 +53,25 @@ type AlertMapper struct {
mapper.AlertMapper
}
// AbsoluteURL for alerts API endpoint this mapper supports
func (m AlertMapper) AbsoluteURL(baseURI string) (string, error) {
return transport.JoinURL(baseURI, "api/v1/alerts/groups")
}
// IsSupported returns true if given version string is supported
func (m AlertMapper) IsSupported(version string) bool {
versionRange := semver.MustParseRange(">=0.5.0 <=0.6.0")
return versionRange(semver.MustParse(version))
}
// GetAlerts will make a request to Alertmanager API and parse the response
// It will only return alerts or error (if any)
func (m AlertMapper) GetAlerts(uri string, timeout time.Duration) ([]models.AlertGroup, error) {
// Decode Alertmanager API response body and return unsee model instances
func (m AlertMapper) Decode(source io.ReadCloser) ([]models.AlertGroup, error) {
groups := []models.AlertGroup{}
receivers := map[string]alertsGroupReceiver{}
resp := alertsGroupsAPISchema{}
url, err := transport.JoinURL(uri, "api/v1/alerts/groups")
if err != nil {
return groups, err
}
err = transport.ReadJSON(url, timeout, &resp)
defer source.Close()
err := json.NewDecoder(source).Decode(&resp)
if err != nil {
return groups, err
}
+11 -9
View File
@@ -5,7 +5,9 @@
package v05
import (
"encoding/json"
"errors"
"io"
"time"
"github.com/blang/semver"
@@ -39,24 +41,24 @@ type SilenceMapper struct {
mapper.SilenceMapper
}
// AbsoluteURL for silences API endpoint this mapper supports
func (m SilenceMapper) AbsoluteURL(baseURI string) (string, error) {
return transport.JoinURL(baseURI, "api/v1/silences")
}
// IsSupported returns true if given version string is supported
func (m SilenceMapper) IsSupported(version string) bool {
versionRange := semver.MustParseRange(">=0.5.0")
return versionRange(semver.MustParse(version))
}
// GetSilences will make a request to Alertmanager API and parse the response
// It will only return silences or error (if any)
func (m SilenceMapper) GetSilences(uri string, timeout time.Duration) ([]models.Silence, error) {
// Decode Alertmanager API response body and return unsee model instances
func (m SilenceMapper) Decode(source io.ReadCloser) ([]models.Silence, error) {
silences := []models.Silence{}
resp := silenceAPISchema{}
url, err := transport.JoinURL(uri, "api/v1/silences")
if err != nil {
return silences, err
}
err = transport.ReadJSON(url, timeout, &resp)
defer source.Close()
err := json.NewDecoder(source).Decode(&resp)
if err != nil {
return silences, err
}
+11 -9
View File
@@ -6,7 +6,9 @@
package v061
import (
"encoding/json"
"errors"
"io"
"sort"
"time"
@@ -53,25 +55,25 @@ type AlertMapper struct {
mapper.AlertMapper
}
// AbsoluteURL for alerts API endpoint this mapper supports
func (m AlertMapper) AbsoluteURL(baseURI string) (string, error) {
return transport.JoinURL(baseURI, "api/v1/alerts/groups")
}
// IsSupported returns true if given version string is supported
func (m AlertMapper) IsSupported(version string) bool {
versionRange := semver.MustParseRange("=0.6.1")
return versionRange(semver.MustParse(version))
}
// GetAlerts will make a request to Alertmanager API and parse the response
// It will only return alerts or error (if any)
func (m AlertMapper) GetAlerts(uri string, timeout time.Duration) ([]models.AlertGroup, error) {
// Decode Alertmanager API response body and return unsee model instances
func (m AlertMapper) Decode(source io.ReadCloser) ([]models.AlertGroup, error) {
groups := []models.AlertGroup{}
receivers := map[string]alertsGroupReceiver{}
resp := alertsGroupsAPISchema{}
url, err := transport.JoinURL(uri, "api/v1/alerts/groups")
if err != nil {
return groups, err
}
err = transport.ReadJSON(url, timeout, &resp)
defer source.Close()
err := json.NewDecoder(source).Decode(&resp)
if err != nil {
return groups, err
}
+11 -9
View File
@@ -6,7 +6,9 @@
package v062
import (
"encoding/json"
"errors"
"io"
"sort"
"time"
@@ -57,25 +59,25 @@ type AlertMapper struct {
mapper.AlertMapper
}
// AbsoluteURL for alerts API endpoint this mapper supports
func (m AlertMapper) AbsoluteURL(baseURI string) (string, error) {
return transport.JoinURL(baseURI, "api/v1/alerts/groups")
}
// IsSupported returns true if given version string is supported
func (m AlertMapper) IsSupported(version string) bool {
versionRange := semver.MustParseRange(">=0.6.2")
return versionRange(semver.MustParse(version))
}
// GetAlerts will make a request to Alertmanager API and parse the response
// It will only return alerts or error (if any)
func (m AlertMapper) GetAlerts(uri string, timeout time.Duration) ([]models.AlertGroup, error) {
// Decode Alertmanager API response body and return unsee model instances
func (m AlertMapper) Decode(source io.ReadCloser) ([]models.AlertGroup, error) {
groups := []models.AlertGroup{}
receivers := map[string]alertsGroupReceiver{}
resp := alertsGroupsAPISchema{}
url, err := transport.JoinURL(uri, "api/v1/alerts/groups")
if err != nil {
return groups, err
}
err = transport.ReadJSON(url, timeout, &resp)
defer source.Close()
err := json.NewDecoder(source).Decode(&resp)
if err != nil {
return groups, err
}
+36 -6
View File
@@ -2,14 +2,16 @@ package transport
import (
"io"
"net/url"
"os"
"path"
"strings"
log "github.com/sirupsen/logrus"
)
type fileReader struct {
filename string
fd *os.File
fd *os.File
}
func (fr *fileReader) Read(b []byte) (n int, err error) {
@@ -20,9 +22,37 @@ func (fr *fileReader) Close() error {
return fr.fd.Close()
}
func newFileReader(filname string) (io.ReadCloser, error) {
log.Infof("Reading file '%s'", filname)
fd, err := os.Open(filname)
fr := fileReader{filename: filname, fd: fd}
// FileTransport can read data from file:// URIs
type FileTransport struct {
}
func (t *FileTransport) pathFromURI(uri string) (string, error) {
u, err := url.Parse(uri)
if err != nil {
return "", err
}
// if we have a file URI with relative path we need to expand it into an
// absolute path, url.Parse doesn't support relative file paths
if strings.HasPrefix(uri, "file:///") {
return u.Path, nil
}
wd, err := os.Getwd()
if err != nil {
return "", err
}
absolutePath := path.Join(wd, strings.TrimPrefix(uri, "file://"))
return absolutePath, nil
}
func (t *FileTransport) Read(uri string) (io.ReadCloser, error) {
filename, err := t.pathFromURI(uri)
if err != nil {
return nil, err
}
log.Infof("Reading file '%s'", filename)
fd, err := os.Open(filename)
fr := fileReader{fd: fd}
return &fr, err
}
+10 -16
View File
@@ -5,37 +5,31 @@ import (
"fmt"
"io"
"net/http"
"time"
log "github.com/sirupsen/logrus"
)
type httpReader struct {
URL string
Timeout time.Duration
// HTTPTransport can read data from http:// and https:// URIs
type HTTPTransport struct {
client http.Client
}
func newHTTPReader(url string, timeout time.Duration) (io.ReadCloser, error) {
hr := httpReader{URL: url, Timeout: timeout}
func (t *HTTPTransport) Read(uri string) (io.ReadCloser, error) {
log.Infof("GET %s timeout=%s", uri, t.client.Timeout)
log.Infof("GET %s timeout=%s", hr.URL, hr.Timeout)
c := &http.Client{
Timeout: timeout,
}
req, err := http.NewRequest("GET", hr.URL, nil)
request, err := http.NewRequest("GET", uri, nil)
if err != nil {
return nil, err
}
req.Header.Add("Accept-Encoding", "gzip")
resp, err := c.Do(req)
request.Header.Add("Accept-Encoding", "gzip")
resp, err := t.client.Do(request)
if err != nil {
return nil, err
}
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("Request to Alertmanager failed with %s", resp.Status)
return nil, fmt.Errorf("Request to %s failed with %s", uri, resp.Status)
}
var reader io.ReadCloser
+14 -27
View File
@@ -1,45 +1,32 @@
package transport
import (
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"os"
"path"
"strings"
"time"
)
// ReadJSON using one of supported transports (file:// http://)
func ReadJSON(uri string, timeout time.Duration, target interface{}) error {
// Transport reads from a specific URI schema
type Transport interface {
Read(string) (io.ReadCloser, error)
}
// NewTransport creates an instance of Transport that can handle URI schema
// for the passed uri string
func NewTransport(uri string, timeout time.Duration) (Transport, error) {
u, err := url.Parse(uri)
if err != nil {
return err
return nil, err
}
var reader io.ReadCloser
switch u.Scheme {
case "http", "https":
reader, err = newHTTPReader(u.String(), timeout)
return &HTTPTransport{client: http.Client{Timeout: timeout}}, nil
case "file":
// if we have a file URI with relative path we need to expand it into an
// absolute path, url.Parse doesn't support relative file paths
if strings.HasPrefix(uri, "file:///") {
reader, err = newFileReader(u.Path)
} else {
wd, e := os.Getwd()
if e != nil {
return e
}
absolutePath := path.Join(wd, strings.TrimPrefix(uri, "file://"))
reader, err = newFileReader(absolutePath)
}
return &FileTransport{}, nil
default:
return fmt.Errorf("Unsupported URI scheme '%s' in '%s'", u.Scheme, u)
return nil, fmt.Errorf("Unsupported URI scheme '%s' in '%s'", u.Scheme, u)
}
if err != nil {
return err
}
defer reader.Close()
return json.NewDecoder(reader).Decode(target)
}
+19 -3
View File
@@ -1,6 +1,7 @@
package transport_test
import (
"encoding/json"
"fmt"
"testing"
"time"
@@ -61,8 +62,8 @@ type mockStatus struct {
no bool
}
func TestFileReader(t *testing.T) {
log.SetLevel(log.ErrorLevel)
func TestTransport(t *testing.T) {
log.SetLevel(log.FatalLevel)
httpmock.Activate()
defer httpmock.DeactivateAndReset()
mockJSON := `{
@@ -79,8 +80,23 @@ func TestFileReader(t *testing.T) {
httpmock.RegisterResponder("GET", "https://localhost/invalid", httpmock.NewStringResponder(200, "bad json}{}"))
for _, testCase := range transportTests {
tr, err := transport.NewTransport(testCase.uri, testCase.timeout)
if err != nil {
t.Error(err)
}
source, err := tr.Read(testCase.uri)
if err != nil {
if !testCase.failed {
t.Errorf("[%s] transport Read() failed with: %s", testCase.uri, err)
}
continue
}
r := mockStatus{}
err := transport.ReadJSON(testCase.uri, testCase.timeout, &r)
err = json.NewDecoder(source).Decode(&r)
source.Close()
if (err != nil) != testCase.failed {
t.Errorf("[%s] Expected failure: %v, Read() failed: %v, error: %s", testCase.uri, testCase.failed, (err != nil), err)
}
+6 -3
View File
@@ -60,10 +60,13 @@ func setupRouter(router *gin.Engine) {
func setupUpstreams() {
for _, s := range config.Config.Alertmanager.Servers {
am := alertmanager.NewAlertmanager(s.Name, s.URI, alertmanager.WithRequestTimeout(s.Timeout), alertmanager.WithProxy(s.Proxy))
err := alertmanager.RegisterAlertmanager(am)
am, err := alertmanager.NewAlertmanager(s.Name, s.URI, alertmanager.WithRequestTimeout(s.Timeout), alertmanager.WithProxy(s.Proxy))
if err != nil {
log.Fatalf("Failed to configure Alertmanager '%s' with URI '%s': %s", s.Name, s.URI, err)
log.Fatalf("Failed to create Alertmanager '%s' with URI '%s': %s", s.Name, s.URI, err)
}
err = alertmanager.RegisterAlertmanager(am)
if err != nil {
log.Fatalf("Failed to register Alertmanager '%s' with URI '%s': %s", s.Name, s.URI, err)
}
}
}
+4 -1
View File
@@ -90,12 +90,15 @@ var proxyTests = []proxyTest{
func TestProxy(t *testing.T) {
r := ginTestEngine()
am := alertmanager.NewAlertmanager(
am, err := alertmanager.NewAlertmanager(
"dummy",
"http://localhost:9093",
alertmanager.WithRequestTimeout(time.Second*5),
alertmanager.WithProxy(true),
)
if err != nil {
t.Error(err)
}
setupRouterProxyHandlers(r, am)
httpmock.Activate()