Skip to content
This repository was archived by the owner on Jul 28, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
# Main (unreleased)

- [ENHANCEMENT] Strengthen readiness check for metrics instances. (@tpaschalis)

# v0.23.0 (2022-01-13)

- [ENHANCEMENT] Go 1.17 is now used for all builds of the Agent. (@tpaschalis)
Expand Down
6 changes: 6 additions & 0 deletions cmd/agent/entrypoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,12 @@ func (ep *Entrypoint) wire(mux *mux.Router, grpc *grpc.Server) {
})

mux.HandleFunc("/-/ready", func(w http.ResponseWriter, r *http.Request) {
if !ep.promMetrics.Ready() {
w.WriteHeader(http.StatusServiceUnavailable)
fmt.Fprint(w, "Metrics are not ready yet.\n")

return
}
w.WriteHeader(http.StatusOK)
fmt.Fprintf(w, "Agent is Ready.\n")
})
Expand Down
6 changes: 1 addition & 5 deletions pkg/integrations/v2/autoscrape/autoscrape_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ import (
"github.com/prometheus/prometheus/discovery"
"github.com/prometheus/prometheus/pkg/exemplar"
"github.com/prometheus/prometheus/pkg/labels"
"github.com/prometheus/prometheus/scrape"
"github.com/prometheus/prometheus/storage"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -98,11 +97,8 @@ func (ma *mockAppender) AppendExemplar(ref uint64, l labels.Labels, e exemplar.E
}

type mockInstance struct {
instance.NoOpInstance
app storage.Appender
}

func (mi *mockInstance) Run(ctx context.Context) error { return nil }
func (mi *mockInstance) Update(c instance.Config) error { return nil }
func (mi *mockInstance) TargetsActive() map[string][]*scrape.Target { return nil }
func (mi *mockInstance) StorageDirectory() string { return "" }
func (mi *mockInstance) Appender(ctx context.Context) storage.Appender { return mi.app }
22 changes: 22 additions & 0 deletions pkg/metrics/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"github.com/go-kit/log"
"github.com/go-kit/log/level"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/atomic"
"google.golang.org/grpc"

"github.com/grafana/agent/pkg/metrics/cluster"
Expand Down Expand Up @@ -144,6 +145,8 @@ type Agent struct {
stopped bool
stopOnce sync.Once
actor chan func()

initialBootDone atomic.Bool
}

// New creates and starts a new Agent.
Expand Down Expand Up @@ -268,6 +271,7 @@ func (a *Agent) ApplyConfig(cfg Config) error {

a.actor <- func() {
a.syncInstances(oldConfig, cfg)
a.initialBootDone.Store(true)
}

a.cfg = cfg
Expand Down Expand Up @@ -311,6 +315,24 @@ func (a *Agent) run() {
}
}

// Ready returns true if both the agent and all instances
// spawned by a Manager have completed startup.
func (a *Agent) Ready() bool {
// Wait for the initial load to complete so the instance manager has at least
// the base set of expected instances.
if !a.initialBootDone.Load() {
return false
}

for _, inst := range a.mm.ListInstances() {
if !inst.Ready() {
return false
}
}

return true
}

// WireGRPC wires gRPC services into the provided server.
func (a *Agent) WireGRPC(s *grpc.Server) {
a.cluster.WireGRPC(s)
Expand Down
4 changes: 4 additions & 0 deletions pkg/metrics/agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,10 @@ func (i *fakeInstance) Run(ctx context.Context) error {
}
}

func (i *fakeInstance) Ready() bool {
return true
}

func (i *fakeInstance) Update(_ instance.Config) error {
return instance.ErrInvalidUpdate{
Inner: fmt.Errorf("can't dynamically update fakeInstance"),
Expand Down
20 changes: 1 addition & 19 deletions pkg/metrics/http_test.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
package metrics

import (
"context"
"fmt"
"net/http"
"net/http/httptest"
Expand All @@ -15,7 +14,6 @@ import (
"github.com/prometheus/common/model"
"github.com/prometheus/prometheus/pkg/labels"
"github.com/prometheus/prometheus/scrape"
"github.com/prometheus/prometheus/storage"
"github.com/stretchr/testify/require"
)

Expand Down Expand Up @@ -135,26 +133,10 @@ func TestAgent_ListTargetsHandler(t *testing.T) {
}

type mockInstanceScrape struct {
instance.NoOpInstance
tgts map[string][]*scrape.Target
}

func (i *mockInstanceScrape) Run(ctx context.Context) error {
<-ctx.Done()
return nil
}

func (i *mockInstanceScrape) Update(_ instance.Config) error {
return nil
}

func (i *mockInstanceScrape) TargetsActive() map[string][]*scrape.Target {
return i.tgts
}

func (i *mockInstanceScrape) StorageDirectory() string {
return ""
}

func (i *mockInstanceScrape) Appender(ctx context.Context) storage.Appender {
return nil
}
13 changes: 12 additions & 1 deletion pkg/metrics/instance/instance.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import (
"github.com/prometheus/prometheus/scrape"
"github.com/prometheus/prometheus/storage"
"github.com/prometheus/prometheus/storage/remote"
"go.uber.org/atomic"
"gopkg.in/yaml.v2"
)

Expand Down Expand Up @@ -240,6 +241,9 @@ type Instance struct {
remoteStore *remote.Storage
storage storage.Storage

// ready is set to true after the initialization process finishes
ready atomic.Bool

hostFilter *HostFilter

logger log.Logger
Expand Down Expand Up @@ -318,7 +322,7 @@ func (i *Instance) Run(ctx context.Context) error {
// The actors defined here are defined in the order we want them to shut down.
// Primarily, we want to ensure that the following shutdown order is
// maintained:
// 1. The scrape manager stops
// 1. The scrape manager stops
// 2. WAL storage is closed
// 3. Remote write storage is closed
// This is done to allow the instance to write stale markers for all active
Expand Down Expand Up @@ -386,6 +390,7 @@ func (i *Instance) Run(ctx context.Context) error {
}

level.Debug(i.logger).Log("msg", "running instance", "name", cfg.Name)
i.ready.Store(true)
err := rg.Run()
if err != nil {
level.Error(i.logger).Log("msg", "agent instance stopped with error", "err", err)
Expand Down Expand Up @@ -446,6 +451,12 @@ func (i *Instance) initialize(ctx context.Context, reg prometheus.Registerer, cf
return nil
}

// Ready returns true if the Instance has been initialized and is ready
// to start scraping and delivering metrics.
func (i *Instance) Ready() bool {
return i.ready.Load()
}

// Update accepts a new Config for the Instance and will dynamically update any
// running Prometheus components with the new values from Config. Update will
// return an ErrInvalidUpdate if the Update could not be applied.
Expand Down
1 change: 1 addition & 0 deletions pkg/metrics/instance/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ type Manager interface {
// for the sake of testing from Manager implementations.
type ManagedInstance interface {
Run(ctx context.Context) error
Ready() bool
Update(c Config) error
TargetsActive() map[string][]*scrape.Target
StorageDirectory() string
Expand Down
8 changes: 8 additions & 0 deletions pkg/metrics/instance/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ func TestBasicManager_ApplyConfig(t *testing.T) {

type mockInstance struct {
RunFunc func(ctx context.Context) error
ReadyFunc func() bool
UpdateFunc func(c Config) error
TargetsActiveFunc func() map[string][]*scrape.Target
StorageDirectoryFunc func() string
Expand All @@ -112,6 +113,13 @@ func (m mockInstance) Run(ctx context.Context) error {
panic("RunFunc not provided")
}

func (m mockInstance) Ready() bool {
if m.ReadyFunc != nil {
return m.ReadyFunc()
}
panic("ReadyFunc not provided")
}

func (m mockInstance) Update(c Config) error {
if m.UpdateFunc != nil {
return m.UpdateFunc(c)
Expand Down
5 changes: 5 additions & 0 deletions pkg/metrics/instance/noop.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,11 @@ func (NoOpInstance) Run(ctx context.Context) error {
return nil
}

// Ready implements Instance.
func (NoOpInstance) Ready() bool {
return true
}

// Update implements Instance.
func (NoOpInstance) Update(_ Config) error {
return nil
Expand Down
10 changes: 1 addition & 9 deletions pkg/traces/remotewriteexporter/exporter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ import (
"github.com/grafana/agent/pkg/metrics/instance"
"github.com/prometheus/prometheus/pkg/exemplar"
"github.com/prometheus/prometheus/pkg/labels"
"github.com/prometheus/prometheus/scrape"
"github.com/prometheus/prometheus/storage"
"github.com/stretchr/testify/require"
"go.opentelemetry.io/collector/model/pdata"
Expand Down Expand Up @@ -103,17 +102,10 @@ func (m *mockManager) DeleteConfig(_ string) error { return nil }
func (m *mockManager) Stop() {}

type mockInstance struct {
instance.NoOpInstance
appender *mockAppender
}

func (m *mockInstance) Run(_ context.Context) error { return nil }

func (m *mockInstance) Update(_ instance.Config) error { return nil }

func (m *mockInstance) TargetsActive() map[string][]*scrape.Target { return nil }

func (m *mockInstance) StorageDirectory() string { return "" }

func (m *mockInstance) Appender(_ context.Context) storage.Appender {
if m.appender == nil {
m.appender = &mockAppender{}
Expand Down