Alerting: Support context.Context in Loki interface (#61979)
This commit adds support for canceleable contexts in the Loki interface.
This commit is contained in:
@@ -23,8 +23,8 @@ const (
|
||||
)
|
||||
|
||||
type remoteLokiClient interface {
|
||||
ping() error
|
||||
push([]stream) error
|
||||
ping(context.Context) error
|
||||
push(context.Context, []stream) error
|
||||
}
|
||||
|
||||
type RemoteLokiBackend struct {
|
||||
@@ -40,8 +40,8 @@ func NewRemoteLokiBackend(cfg LokiConfig) *RemoteLokiBackend {
|
||||
}
|
||||
}
|
||||
|
||||
func (h *RemoteLokiBackend) TestConnection() error {
|
||||
return h.client.ping()
|
||||
func (h *RemoteLokiBackend) TestConnection(ctx context.Context) error {
|
||||
return h.client.ping(ctx)
|
||||
}
|
||||
|
||||
func (h *RemoteLokiBackend) RecordStatesAsync(ctx context.Context, rule history_model.RuleMeta, states []state.StateTransition) <-chan error {
|
||||
@@ -116,7 +116,7 @@ func (h *RemoteLokiBackend) recordStreamsAsync(ctx context.Context, streams []st
|
||||
}
|
||||
|
||||
func (h *RemoteLokiBackend) recordStreams(ctx context.Context, streams []stream, logger log.Logger) error {
|
||||
if err := h.client.push(streams); err != nil {
|
||||
if err := h.client.push(ctx, streams); err != nil {
|
||||
return err
|
||||
}
|
||||
logger.Debug("Done saving alert state history batch")
|
||||
|
||||
@@ -2,6 +2,7 @@ package historian
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -37,7 +38,7 @@ func newLokiClient(cfg LokiConfig, logger log.Logger) *httpLokiClient {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *httpLokiClient) ping() error {
|
||||
func (c *httpLokiClient) ping(ctx context.Context) error {
|
||||
uri := c.cfg.Url.JoinPath("/loki/api/v1/labels")
|
||||
req, err := http.NewRequest(http.MethodGet, uri.String(), nil)
|
||||
if err != nil {
|
||||
@@ -45,6 +46,7 @@ func (c *httpLokiClient) ping() error {
|
||||
}
|
||||
c.setAuthAndTenantHeaders(req)
|
||||
|
||||
req = req.WithContext(ctx)
|
||||
res, err := c.client.Do(req)
|
||||
if res != nil {
|
||||
defer func() {
|
||||
@@ -80,7 +82,7 @@ func (r *row) MarshalJSON() ([]byte, error) {
|
||||
})
|
||||
}
|
||||
|
||||
func (c *httpLokiClient) push(s []stream) error {
|
||||
func (c *httpLokiClient) push(ctx context.Context, s []stream) error {
|
||||
body := struct {
|
||||
Streams []stream `json:"streams"`
|
||||
}{Streams: s}
|
||||
@@ -98,6 +100,7 @@ func (c *httpLokiClient) push(s []stream) error {
|
||||
c.setAuthAndTenantHeaders(req)
|
||||
req.Header.Add("content-type", "application/json")
|
||||
|
||||
req = req.WithContext(ctx)
|
||||
resp, err := c.client.Do(req)
|
||||
if resp != nil {
|
||||
defer func() {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package historian
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/url"
|
||||
"testing"
|
||||
|
||||
@@ -21,7 +22,7 @@ func TestLokiHTTPClient(t *testing.T) {
|
||||
}, log.NewNopLogger())
|
||||
|
||||
// Unauthorized request should fail against Grafana Cloud.
|
||||
err = client.ping()
|
||||
err = client.ping(context.Background())
|
||||
require.Error(t, err)
|
||||
|
||||
client.cfg.BasicAuthUser = "<your_username>"
|
||||
@@ -32,7 +33,7 @@ func TestLokiHTTPClient(t *testing.T) {
|
||||
// client.cfg.TenantID = "<your_tenant_id>"
|
||||
|
||||
// Authorized request should fail against Grafana Cloud.
|
||||
err = client.ping()
|
||||
err = client.ping(context.Background())
|
||||
require.NoError(t, err)
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user