From 0b30a23c2a9573a1e728c6e0e34563144c8970ef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=C5=81ukasz=20Mierzwa?= Date: Sat, 20 Jan 2018 14:17:45 -0800 Subject: [PATCH 1/5] Refactor mapper package responsibilities Mapper should only take input, decode it and map onto internal model instances --- internal/mapper/mapper.go | 26 +++++++++++++------------- internal/mapper/v04/alerts.go | 15 +++++---------- internal/mapper/v04/silences.go | 19 +++++-------------- internal/mapper/v05/alerts.go | 15 +++++---------- internal/mapper/v05/silences.go | 15 +++++---------- internal/mapper/v061/alerts.go | 15 +++++---------- internal/mapper/v062/alerts.go | 15 +++++---------- 7 files changed, 43 insertions(+), 77 deletions(-) diff --git a/internal/mapper/mapper.go b/internal/mapper/mapper.go index 2d1bcb778..84761409e 100644 --- a/internal/mapper/mapper.go +++ b/internal/mapper/mapper.go @@ -2,7 +2,7 @@ package mapper import ( "fmt" - "time" + "io" "github.com/cloudflare/unsee/internal/models" ) @@ -12,11 +12,19 @@ var ( silenceMappers = []SilenceMapper{} ) -// AlertMapper implements Alertmanager -> unsee alert data mapping that works -// for a specific range of Alertmanager versions -type AlertMapper interface { +type Mapper interface { IsSupported(version string) bool - GetAlerts(uri string, timeout time.Duration) ([]models.AlertGroup, error) +} + +type AlertMapper interface { + Mapper + Decode(io.ReadCloser) ([]models.AlertGroup, error) +} + +type SilenceMapper interface { + Mapper + Decode(io.ReadCloser) ([]models.Silence, error) + Release() string } // RegisterAlertMapper allows to register mapper implementing alert data @@ -35,14 +43,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) { diff --git a/internal/mapper/v04/alerts.go b/internal/mapper/v04/alerts.go index eb03db9e3..756206910 100644 --- a/internal/mapper/v04/alerts.go +++ b/internal/mapper/v04/alerts.go @@ -5,7 +5,9 @@ package v04 import ( + "encoding/json" "errors" + "io" "sort" "strconv" "time" @@ -13,7 +15,6 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" - "github.com/cloudflare/unsee/internal/transport" ) type alert struct { @@ -58,19 +59,13 @@ func (m AlertMapper) IsSupported(version string) bool { 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) { +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 } diff --git a/internal/mapper/v04/silences.go b/internal/mapper/v04/silences.go index 9439c6b4d..7476fbe9f 100644 --- a/internal/mapper/v04/silences.go +++ b/internal/mapper/v04/silences.go @@ -5,16 +5,15 @@ package v04 import ( + "encoding/json" "errors" - "fmt" - "math" + "io" "strconv" "time" "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" - "github.com/cloudflare/unsee/internal/transport" ) // Alertmanager 0.4 silence format @@ -53,20 +52,12 @@ func (m SilenceMapper) IsSupported(version string) bool { 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) { +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 } diff --git a/internal/mapper/v05/alerts.go b/internal/mapper/v05/alerts.go index 6c6b422a3..ad48648fd 100644 --- a/internal/mapper/v05/alerts.go +++ b/internal/mapper/v05/alerts.go @@ -5,14 +5,15 @@ package v05 import ( + "encoding/json" "errors" + "io" "sort" "time" "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" - "github.com/cloudflare/unsee/internal/transport" ) type alert struct { @@ -57,19 +58,13 @@ func (m AlertMapper) IsSupported(version string) bool { 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) { +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 } diff --git a/internal/mapper/v05/silences.go b/internal/mapper/v05/silences.go index 118ded6b0..54ac07838 100644 --- a/internal/mapper/v05/silences.go +++ b/internal/mapper/v05/silences.go @@ -5,13 +5,14 @@ package v05 import ( + "encoding/json" "errors" + "io" "time" "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" - "github.com/cloudflare/unsee/internal/transport" ) type silence struct { @@ -45,18 +46,12 @@ func (m SilenceMapper) IsSupported(version string) bool { 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) { +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 } diff --git a/internal/mapper/v061/alerts.go b/internal/mapper/v061/alerts.go index 32b4f8865..70e4881e3 100644 --- a/internal/mapper/v061/alerts.go +++ b/internal/mapper/v061/alerts.go @@ -6,14 +6,15 @@ package v061 import ( + "encoding/json" "errors" + "io" "sort" "time" "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" - "github.com/cloudflare/unsee/internal/transport" ) type alert struct { @@ -59,19 +60,13 @@ func (m AlertMapper) IsSupported(version string) bool { 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) { +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 } diff --git a/internal/mapper/v062/alerts.go b/internal/mapper/v062/alerts.go index af0337fcf..80ea4b836 100644 --- a/internal/mapper/v062/alerts.go +++ b/internal/mapper/v062/alerts.go @@ -6,14 +6,15 @@ package v062 import ( + "encoding/json" "errors" + "io" "sort" "time" "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" - "github.com/cloudflare/unsee/internal/transport" ) type alertStatus struct { @@ -63,19 +64,13 @@ func (m AlertMapper) IsSupported(version string) bool { 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) { +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 } From 2cf9253d3c333c1163ce9aec886a3653d1d4979d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=C5=81ukasz=20Mierzwa?= Date: Sat, 20 Jan 2018 16:16:28 -0800 Subject: [PATCH 2/5] Refactor transport package With this change we'll initialize Transport object for each Alertmanager and just call Read() on it when we need to use this transport to read from upstream Alertmanager --- internal/transport/file.go | 42 ++++++++++++++++++++++++---- internal/transport/http.go | 26 +++++++---------- internal/transport/transport.go | 41 ++++++++++----------------- internal/transport/transport_test.go | 22 +++++++++++++-- 4 files changed, 79 insertions(+), 52 deletions(-) diff --git a/internal/transport/file.go b/internal/transport/file.go index c1241e3f9..6850be4dc 100644 --- a/internal/transport/file.go +++ b/internal/transport/file.go @@ -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 } diff --git a/internal/transport/http.go b/internal/transport/http.go index e55467347..a30fa7fca 100644 --- a/internal/transport/http.go +++ b/internal/transport/http.go @@ -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 diff --git a/internal/transport/transport.go b/internal/transport/transport.go index e1f386a26..e57d4aa55 100644 --- a/internal/transport/transport.go +++ b/internal/transport/transport.go @@ -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) } diff --git a/internal/transport/transport_test.go b/internal/transport/transport_test.go index 6f92ed510..981673f1a 100644 --- a/internal/transport/transport_test.go +++ b/internal/transport/transport_test.go @@ -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) } From 6124196b0f11f81a1551c303964e3f3a58e07d23 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=C5=81ukasz=20Mierzwa?= Date: Sat, 20 Jan 2018 16:18:04 -0800 Subject: [PATCH 3/5] Finish refactoring mapper package Fix bugs, add docstrings and let each mapper give us full url since it doesn't handle any http requests now (it just maps text to objects now) --- internal/mapper/mapper.go | 4 ++++ internal/mapper/v04/alerts.go | 9 ++++++++- internal/mapper/v04/silences.go | 9 ++++++++- internal/mapper/v05/alerts.go | 9 ++++++++- internal/mapper/v05/silences.go | 9 ++++++++- internal/mapper/v061/alerts.go | 9 ++++++++- internal/mapper/v062/alerts.go | 9 ++++++++- 7 files changed, 52 insertions(+), 6 deletions(-) diff --git a/internal/mapper/mapper.go b/internal/mapper/mapper.go index 84761409e..8d7ef6924 100644 --- a/internal/mapper/mapper.go +++ b/internal/mapper/mapper.go @@ -12,15 +12,19 @@ var ( silenceMappers = []SilenceMapper{} ) +// Mapper converts Alertmanager response body and maps to unsee data structures type Mapper interface { IsSupported(version string) bool + 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) diff --git a/internal/mapper/v04/alerts.go b/internal/mapper/v04/alerts.go index 756206910..becbffa37 100644 --- a/internal/mapper/v04/alerts.go +++ b/internal/mapper/v04/alerts.go @@ -15,6 +15,7 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" + "github.com/cloudflare/unsee/internal/transport" ) type alert struct { @@ -53,19 +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)) } +// 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{} defer source.Close() - err := json.NewDecoder(source).Decode(resp) + err := json.NewDecoder(source).Decode(&resp) if err != nil { return groups, err } diff --git a/internal/mapper/v04/silences.go b/internal/mapper/v04/silences.go index 7476fbe9f..0362ba4cf 100644 --- a/internal/mapper/v04/silences.go +++ b/internal/mapper/v04/silences.go @@ -14,6 +14,7 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" + "github.com/cloudflare/unsee/internal/transport" ) // Alertmanager 0.4 silence format @@ -46,18 +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)) } +// 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{} defer source.Close() - err := json.NewDecoder(source).Decode(resp) + err := json.NewDecoder(source).Decode(&resp) if err != nil { return silences, err } diff --git a/internal/mapper/v05/alerts.go b/internal/mapper/v05/alerts.go index ad48648fd..efc1b1f06 100644 --- a/internal/mapper/v05/alerts.go +++ b/internal/mapper/v05/alerts.go @@ -14,6 +14,7 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" + "github.com/cloudflare/unsee/internal/transport" ) type alert struct { @@ -52,19 +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)) } +// 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{} defer source.Close() - err := json.NewDecoder(source).Decode(resp) + err := json.NewDecoder(source).Decode(&resp) if err != nil { return groups, err } diff --git a/internal/mapper/v05/silences.go b/internal/mapper/v05/silences.go index 54ac07838..eeac798d8 100644 --- a/internal/mapper/v05/silences.go +++ b/internal/mapper/v05/silences.go @@ -13,6 +13,7 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" + "github.com/cloudflare/unsee/internal/transport" ) type silence struct { @@ -40,18 +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)) } +// 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{} defer source.Close() - err := json.NewDecoder(source).Decode(resp) + err := json.NewDecoder(source).Decode(&resp) if err != nil { return silences, err } diff --git a/internal/mapper/v061/alerts.go b/internal/mapper/v061/alerts.go index 70e4881e3..799759603 100644 --- a/internal/mapper/v061/alerts.go +++ b/internal/mapper/v061/alerts.go @@ -15,6 +15,7 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" + "github.com/cloudflare/unsee/internal/transport" ) type alert struct { @@ -54,19 +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)) } +// 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{} defer source.Close() - err := json.NewDecoder(source).Decode(resp) + err := json.NewDecoder(source).Decode(&resp) if err != nil { return groups, err } diff --git a/internal/mapper/v062/alerts.go b/internal/mapper/v062/alerts.go index 80ea4b836..4d9ba687a 100644 --- a/internal/mapper/v062/alerts.go +++ b/internal/mapper/v062/alerts.go @@ -15,6 +15,7 @@ import ( "github.com/blang/semver" "github.com/cloudflare/unsee/internal/mapper" "github.com/cloudflare/unsee/internal/models" + "github.com/cloudflare/unsee/internal/transport" ) type alertStatus struct { @@ -58,19 +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)) } +// 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{} defer source.Close() - err := json.NewDecoder(source).Decode(resp) + err := json.NewDecoder(source).Decode(&resp) if err != nil { return groups, err } From 1f89ba05fe4a134f4f6d4dfc51fdab5c4bbf4a7a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=C5=81ukasz=20Mierzwa?= Date: Sat, 20 Jan 2018 16:18:21 -0800 Subject: [PATCH 4/5] Migrate rest of the code to new mapper and transport packages --- internal/alertmanager/dedup_test.go | 5 ++- internal/alertmanager/models.go | 53 ++++++++++++++++++++++++++--- internal/alertmanager/upstream.go | 11 ++++-- internal/alertmanager/version.go | 11 +++++- internal/filters/filter_test.go | 5 ++- main.go | 9 +++-- proxy_test.go | 5 ++- 7 files changed, 86 insertions(+), 13 deletions(-) diff --git a/internal/alertmanager/dedup_test.go b/internal/alertmanager/dedup_test.go index 1d7e64430..59e6484b1 100644 --- a/internal/alertmanager/dedup_test.go +++ b/internal/alertmanager/dedup_test.go @@ -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) } } diff --git a/internal/alertmanager/models.go b/internal/alertmanager/models.go index ce5399119..9bf41f6a6 100644 --- a/internal/alertmanager/models.go +++ b/internal/alertmanager/models.go @@ -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 } diff --git a/internal/alertmanager/upstream.go b/internal/alertmanager/upstream.go index 285b973c5..a67018811 100644 --- a/internal/alertmanager/upstream.go +++ b/internal/alertmanager/upstream.go @@ -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 diff --git a/internal/alertmanager/version.go b/internal/alertmanager/version.go index 0ff58e26a..0fd021d68 100644 --- a/internal/alertmanager/version.go +++ b/internal/alertmanager/version.go @@ -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 diff --git a/internal/filters/filter_test.go b/internal/filters/filter_test.go index c329afb2d..3b7f85476 100644 --- a/internal/filters/filter_test.go +++ b/internal/filters/filter_test.go @@ -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 { diff --git a/main.go b/main.go index b517c54dd..50f8fc2ce 100644 --- a/main.go +++ b/main.go @@ -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) } } } diff --git a/proxy_test.go b/proxy_test.go index 18ed8c9f1..330bcc279 100644 --- a/proxy_test.go +++ b/proxy_test.go @@ -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() From c7fb8db98d922729fab229edd2af240ae45a1ace Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=C5=81ukasz=20Mierzwa?= Date: Sat, 20 Jan 2018 16:27:00 -0800 Subject: [PATCH 5/5] Drop Release() from SilenceMapper Not used --- internal/mapper/mapper.go | 1 - 1 file changed, 1 deletion(-) diff --git a/internal/mapper/mapper.go b/internal/mapper/mapper.go index 8d7ef6924..8315e97d3 100644 --- a/internal/mapper/mapper.go +++ b/internal/mapper/mapper.go @@ -28,7 +28,6 @@ type AlertMapper interface { type SilenceMapper interface { Mapper Decode(io.ReadCloser) ([]models.Silence, error) - Release() string } // RegisterAlertMapper allows to register mapper implementing alert data