Skip to content
Open
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
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ If you are still using the legacy [Access scopes][access-scopes], the `https://w
| ----------------------------------- | -------- |---------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
| `google.project-ids` | No | GCloud SDK auto-discovery | Repeatable flag of Google Project IDs |
| `google.projects.filter` | No | | GCloud projects filter expression. See more [here](https://cloud.google.com/sdk/gcloud/reference/projects/list). |
| `google.projects.filter-refresh-interval` | No | `0s` (disabled) | How often to re-evaluate `google.projects.filter` to pick up projects added to or removed from the matching org/folder without restarting. `0` keeps the pre-existing behavior of resolving the filter once at startup only. |
| `google.universe-domain` | No | `googleapis.com` | Target specific Google Cloud environments, such as public cloud, or specific sovereign clouds |
| `monitoring.metrics-ingest-delay` | No | | Offsets metric collection by a delay appropriate for each metric type, e.g. because bigquery metrics are slow to appear |
| `monitoring.drop-delegated-projects` | No | No | Drop metrics from attached projects and fetch `project_id` only. |
Expand Down
87 changes: 81 additions & 6 deletions collectors/runtime.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
"log/slog"
"slices"
"strings"
"sync/atomic"
"time"

"golang.org/x/oauth2/google"
Expand All @@ -37,13 +38,22 @@ type HistogramStoreFactory func(logger *slog.Logger, ttl time.Duration) DeltaHis

// Runtime holds the resolved state produced by NewRuntime.
type Runtime struct {
cfg *config.Config
projectIDs []string
cfg *config.Config
// projectIDs is a *atomic.Pointer[[]string] (rather than a plain []string)
// so that WithCache's shallow struct copy shares live state with the
// original: a background refresh (see StartProjectDiscoveryRefresh) must
// be visible through every Runtime value derived from the same NewRuntime
// call, and scrapes must be able to read it without lock contention.
projectIDs *atomic.Pointer[[]string]
service *monitoring.Service
logger *slog.Logger
counterStoreFactory CounterStoreFactory
histogramStoreFactory HistogramStoreFactory
cache *collectorCache
// discoverProjectIDs resolves cfg.ProjectsFilter to project IDs. It
// defaults to getProjectIDsFromFilter; tests override it to avoid
// depending on Application Default Credentials.
discoverProjectIDs func(ctx context.Context, filter string) ([]string, error)
}

// NewRuntime resolves project IDs and creates the monitoring service. The
Expand Down Expand Up @@ -84,16 +94,73 @@ func NewRuntime(ctx context.Context, logger *slog.Logger, cfg *config.Config, co
return nil, err
}

projectIDsPtr := &atomic.Pointer[[]string]{}
projectIDsPtr.Store(&projectIDs)

return &Runtime{
cfg: cfg,
projectIDs: projectIDs,
projectIDs: projectIDsPtr,
service: service,
logger: logger,
counterStoreFactory: counterFactory,
histogramStoreFactory: histogramFactory,
discoverProjectIDs: getProjectIDsFromFilter,
}, nil
}

// StartProjectDiscoveryRefresh periodically re-resolves cfg.ProjectsFilter in
// the background so that projects added to or removed from the matching
// org/folder are picked up without restarting the exporter. It is a no-op
// unless both cfg.ProjectsFilter and cfg.ProjectsRefreshInterval are set, so
// existing deployments see no behavior change unless they opt in.
//
// The refresh runs until ctx is canceled; callers must cancel ctx on shutdown
// to avoid leaking the goroutine.
func (r *Runtime) StartProjectDiscoveryRefresh(ctx context.Context) {
if r.cfg.ProjectsFilter == "" || r.cfg.ProjectsRefreshInterval <= 0 {
return
}
go r.refreshProjectIDsLoop(ctx)
}

func (r *Runtime) refreshProjectIDsLoop(ctx context.Context) {
ticker := time.NewTicker(r.cfg.ProjectsRefreshInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
r.refreshProjectIDs(ctx)
}
}
}

// refreshProjectIDs re-resolves cfg.ProjectsFilter and, on a successful
// non-empty result, atomically replaces the project list used by future
// scrapes. A transient API error or an unexpectedly empty result (e.g. IAM
// eventual consistency, a momentarily over-narrow filter) is logged and the
// previous known-good list is kept, so a background refresh hiccup never
// blanks out or fails a scrape.
func (r *Runtime) refreshProjectIDs(ctx context.Context) {
ids, err := r.discoverProjectIDs(ctx, r.cfg.ProjectsFilter)
if err != nil {
r.logger.Warn("failed to refresh project list from google.projects.filter; keeping previous list", "err", err)
return
}

ids = append(ids, r.cfg.ProjectIDs...)
ids = deduplicateProjectIDs(ids)

if len(ids) == 0 {
r.logger.Warn("google.projects.filter refresh returned zero projects; keeping previous list")
return
}

r.projectIDs.Store(&ids)
r.logger.Info("refreshed project list from google.projects.filter", "count", len(ids))
}

// WithCache returns a Runtime configured to cache its collectors per
// (project, prefix-filter). Subsequent calls to Collectors or
// CollectorsForPrefixes reuse cached entries until they expire, which lets
Expand Down Expand Up @@ -125,8 +192,9 @@ func (r *Runtime) CollectorsForPrefixes(prefixFilter []string) ([]*MonitoringCol
}

func (r *Runtime) buildCollectors(prefixFilter []string) ([]*MonitoringCollector, error) {
result := make([]*MonitoringCollector, 0, len(r.projectIDs))
for _, projectID := range r.projectIDs {
projectIDs := *r.projectIDs.Load()
result := make([]*MonitoringCollector, 0, len(projectIDs))
for _, projectID := range projectIDs {
c, err := r.collectorFor(projectID, prefixFilter)
if err != nil {
return nil, fmt.Errorf("collector for %q: %w", projectID, err)
Expand Down Expand Up @@ -222,9 +290,16 @@ func getProjectIDsFromFilter(ctx context.Context, filter string) ([]string, erro
if err != nil {
return nil, err
}
return listProjectIDs(ctx, service, filter)
}

// listProjectIDs returns the list of project IDs that match filter, using an
// already-constructed cloudresourcemanager service. Split out from
// getProjectIDsFromFilter so tests can inject a fake service pointed at an
// httptest.Server instead of relying on Application Default Credentials.
func listProjectIDs(ctx context.Context, service *cloudresourcemanager.Service, filter string) ([]string, error) {
var projectIDs []string
err = service.Projects.List().Filter(filter).Pages(ctx, func(page *cloudresourcemanager.ListProjectsResponse) error {
err := service.Projects.List().Filter(filter).Pages(ctx, func(page *cloudresourcemanager.ListProjectsResponse) error {
for _, project := range page.Projects {
projectIDs = append(projectIDs, project.ProjectId)
}
Expand Down
237 changes: 237 additions & 0 deletions collectors/runtime_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,21 @@ package collectors

import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"net/http"
"net/http/httptest"
"reflect"
"strings"
"sync/atomic"
"testing"
"time"

"google.golang.org/api/cloudresourcemanager/v1"
"google.golang.org/api/option"

"github.com/prometheus-community/stackdriver_exporter/config"
)

Expand Down Expand Up @@ -189,3 +198,231 @@ func TestRuntimeFilterMetricTypePrefixes(t *testing.T) {
})
}
}

// newTestRuntime builds a minimal Runtime for unit testing the project
// discovery refresh path, bypassing NewRuntime's ADC/GCP service setup.
func newTestRuntime(t *testing.T, cfg *config.Config, initialIDs []string, discover func(ctx context.Context, filter string) ([]string, error)) *Runtime {
t.Helper()
ptr := &atomic.Pointer[[]string]{}
ptr.Store(&initialIDs)
return &Runtime{
cfg: cfg,
projectIDs: ptr,
logger: slog.Default(),
discoverProjectIDs: discover,
}
}

// fakeCloudResourceManagerServer serves a canned ListProjectsResponse (or an
// error status) so listProjectIDs can be tested against a real HTTP client
// without depending on Application Default Credentials.
func fakeCloudResourceManagerServer(t *testing.T, statusCode int, resp *cloudresourcemanager.ListProjectsResponse) *httptest.Server {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(statusCode)
if resp != nil {
_ = json.NewEncoder(w).Encode(resp)
}
}))
t.Cleanup(server.Close)
return server
}

func fakeCloudResourceManagerService(t *testing.T, server *httptest.Server) *cloudresourcemanager.Service {
t.Helper()
service, err := cloudresourcemanager.NewService(context.Background(),
option.WithEndpoint(server.URL),
option.WithHTTPClient(server.Client()),
option.WithoutAuthentication(),
)
if err != nil {
t.Fatalf("cloudresourcemanager.NewService() error = %v", err)
}
return service
}

func TestListProjectIDsReturnsMatchingProjects(t *testing.T) {
t.Parallel()

server := fakeCloudResourceManagerServer(t, http.StatusOK, &cloudresourcemanager.ListProjectsResponse{
Projects: []*cloudresourcemanager.Project{
{ProjectId: "project-a"},
{ProjectId: "project-b"},
},
})
service := fakeCloudResourceManagerService(t, server)

got, err := listProjectIDs(context.Background(), service, "parent.id:12345")
if err != nil {
t.Fatalf("listProjectIDs() error = %v", err)
}
want := []string{"project-a", "project-b"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("listProjectIDs() = %#v, want %#v", got, want)
}
}

func TestListProjectIDsPropagatesAPIError(t *testing.T) {
t.Parallel()

server := fakeCloudResourceManagerServer(t, http.StatusInternalServerError, nil)
service := fakeCloudResourceManagerService(t, server)

_, err := listProjectIDs(context.Background(), service, "parent.id:12345")
if err == nil {
t.Fatal("listProjectIDs() expected error for 500 response, got nil")
}
}

func TestRefreshProjectIDsUpdatesListOnSuccess(t *testing.T) {
t.Parallel()

cfg := &config.Config{ProjectsFilter: "parent.id:12345"}
r := newTestRuntime(t, cfg, []string{"old-project"}, func(_ context.Context, _ string) ([]string, error) {
return []string{"new-project-b", "new-project-a"}, nil
})

r.refreshProjectIDs(context.Background())

got := *r.projectIDs.Load()
want := []string{"new-project-a", "new-project-b"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("projectIDs after refresh = %#v, want %#v", got, want)
}
}

func TestRefreshProjectIDsMergesStaticProjectIDs(t *testing.T) {
t.Parallel()

cfg := &config.Config{ProjectsFilter: "parent.id:12345", ProjectIDs: []string{"static-project", "new-project-a"}}
r := newTestRuntime(t, cfg, []string{"old-project"}, func(_ context.Context, _ string) ([]string, error) {
return []string{"new-project-a"}, nil
})

r.refreshProjectIDs(context.Background())

got := *r.projectIDs.Load()
want := []string{"new-project-a", "static-project"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("projectIDs after refresh = %#v, want %#v", got, want)
}
}

func TestRefreshProjectIDsKeepsPreviousListOnAPIError(t *testing.T) {
t.Parallel()

cfg := &config.Config{ProjectsFilter: "parent.id:12345"}
r := newTestRuntime(t, cfg, []string{"old-project"}, func(_ context.Context, _ string) ([]string, error) {
return nil, errors.New("transient GCP error")
})

r.refreshProjectIDs(context.Background())

got := *r.projectIDs.Load()
want := []string{"old-project"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("projectIDs after failed refresh = %#v, want unchanged %#v", got, want)
}
}

func TestRefreshProjectIDsKeepsPreviousListWhenFilterReturnsZero(t *testing.T) {
t.Parallel()

cfg := &config.Config{ProjectsFilter: "parent.id:12345"}
r := newTestRuntime(t, cfg, []string{"old-project"}, func(_ context.Context, _ string) ([]string, error) {
return nil, nil
})

r.refreshProjectIDs(context.Background())

got := *r.projectIDs.Load()
want := []string{"old-project"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("projectIDs after empty-result refresh = %#v, want unchanged %#v", got, want)
}
}

func TestStartProjectDiscoveryRefreshRunsOnlyWhenFilterAndIntervalSet(t *testing.T) {
tests := []struct {
name string
filter string
interval time.Duration
expectTriggered bool
}{
{name: "neither filter nor interval set", filter: "", interval: 0, expectTriggered: false},
{name: "filter set, interval zero", filter: "parent.id:1", interval: 0, expectTriggered: false},
{name: "interval set, filter empty", filter: "", interval: 10 * time.Millisecond, expectTriggered: false},
{name: "filter and interval both set", filter: "parent.id:1", interval: 10 * time.Millisecond, expectTriggered: true},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
var calls int32
cfg := &config.Config{ProjectsFilter: tt.filter, ProjectsRefreshInterval: tt.interval}
r := newTestRuntime(t, cfg, []string{"p"}, func(_ context.Context, _ string) ([]string, error) {
atomic.AddInt32(&calls, 1)
return []string{"p2"}, nil
})

ctx, cancel := context.WithCancel(context.Background())
r.StartProjectDiscoveryRefresh(ctx)

time.Sleep(50 * time.Millisecond)
cancel()
time.Sleep(20 * time.Millisecond)

triggered := atomic.LoadInt32(&calls) > 0
if triggered != tt.expectTriggered {
t.Fatalf("refresh triggered = %v, want %v (calls=%d)", triggered, tt.expectTriggered, calls)
}
})
}
}

func TestBuildCollectorsSafeUnderConcurrentRefresh(t *testing.T) {
cfg := &config.Config{MetricsPrefixes: []string{"compute.googleapis.com/"}}
r := newTestRuntime(t, cfg, []string{"project-a"}, nil)
r.counterStoreFactory = func(_ *slog.Logger, _ time.Duration) DeltaCounterStore { return nil }
r.histogramStoreFactory = func(_ *slog.Logger, _ time.Duration) DeltaHistogramStore { return nil }

done := make(chan struct{})
go func() {
defer close(done)
for i := 0; i < 200; i++ {
ids := []string{fmt.Sprintf("project-%d", i)}
r.projectIDs.Store(&ids)
}
}()

for i := 0; i < 200; i++ {
if _, err := r.buildCollectors(nil); err != nil {
t.Fatalf("buildCollectors() error = %v", err)
}
}
<-done
}

// TestWithCacheSharesProjectIDsPointer guards against a regression where
// WithCache's shallow struct copy (sibling := *r) would give the cached
// sibling a disconnected snapshot of the project list instead of sharing
// live state with the original Runtime that a background refresh updates.
func TestWithCacheSharesProjectIDsPointer(t *testing.T) {
cfg := &config.Config{MetricsPrefixes: []string{"compute.googleapis.com/"}}
r := newTestRuntime(t, cfg, []string{"project-a"}, nil)
r.counterStoreFactory = func(_ *slog.Logger, _ time.Duration) DeltaCounterStore { return nil }
r.histogramStoreFactory = func(_ *slog.Logger, _ time.Duration) DeltaHistogramStore { return nil }

sibling := r.WithCache()

updated := []string{"project-a", "project-b"}
r.projectIDs.Store(&updated)

cs, err := sibling.Collectors()
if err != nil {
t.Fatalf("sibling.Collectors() error = %v", err)
}
if len(cs) != 2 {
t.Fatalf("sibling.Collectors() returned %d collectors after original's projectIDs updated, want 2", len(cs))
}
}
Loading