Skip to content

Commit 43fa995

Browse files
committed
[scraperhelper] Add extension-based scraper controller support
Implement the extensionscrapercontroller package and integrate it with the scraper controller, allowing extensions to trigger scrapes based on external events instead of (or in addition to) a fixed timer interval. - Add extensionscrapercontroller.ControllerExtension interface and DeregisterFunc helper in extension/xextension/extensionscrapercontroller - Add Controllers field to ControllerConfig for specifying extension IDs - Allow CollectionInterval=0 when Controllers are configured - Register scrapers with controller extensions on Start, deregister on Shutdown - Update RFC to remove scraperID parameter from RegisterScraper Assisted-by: Claude Opus 4.5
1 parent 28265fc commit 43fa995

10 files changed

Lines changed: 367 additions & 3 deletions

File tree

docs/rfcs/scraper-controller-extension.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ type ControllerExtension interface {
4040
// RegisterScraper registers a scraper with this controller extension.
4141
// The extension will invoke the provided scrape function according to its implementation.
4242
// Returns a registration handle that can be used to deregister the scraper.
43-
RegisterScraper(ctx context.Context, scraperID component.ID, scrapeFunc func(context.Context) error) (RegistrationHandle, error)
43+
RegisterScraper(ctx context.Context, scrapeFunc func(context.Context) error) (RegistrationHandle, error)
4444
}
4545

4646
// RegistrationHandle provides a way to deregister a scraper from a controller extension
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
// Copyright The OpenTelemetry Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
// Package extensionscrapercontroller defines the interface for extensions
5+
// that control when scraper-based receivers invoke their scrapers.
6+
package extensionscrapercontroller // import "go.opentelemetry.io/collector/extension/xextension/extensionscrapercontroller"
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
// Copyright The OpenTelemetry Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
package extensionscrapercontroller // import "go.opentelemetry.io/collector/extension/xextension/extensionscrapercontroller"
5+
6+
import (
7+
"context"
8+
9+
"go.opentelemetry.io/collector/extension"
10+
)
11+
12+
// ControllerExtension is an extension that controls when scraper-based
13+
// receivers invoke their scrapers.
14+
type ControllerExtension interface {
15+
extension.Extension
16+
17+
// RegisterScraper registers a scraper with the extension. The extension
18+
// will call scrapeFunc when it determines a scrape should occur.
19+
// The returned RegistrationHandle must be used to deregister the scraper
20+
// during shutdown.
21+
RegisterScraper(ctx context.Context, scrapeFunc func(context.Context) error) (RegistrationHandle, error)
22+
}
23+
24+
// RegistrationHandle is returned by ControllerExtension.RegisterScraper and
25+
// is used to deregister the scraper during shutdown.
26+
type RegistrationHandle interface {
27+
Deregister(ctx context.Context) error
28+
}
29+
30+
// DeregisterFunc is a function that implements RegistrationHandle.
31+
// A nil DeregisterFunc is valid and returns nil on Deregister.
32+
type DeregisterFunc func(ctx context.Context) error
33+
34+
var _ RegistrationHandle = (DeregisterFunc)(nil)
35+
36+
// Deregister calls the underlying function. If the receiver is nil, it returns nil.
37+
func (f DeregisterFunc) Deregister(ctx context.Context) error {
38+
if f == nil {
39+
return nil
40+
}
41+
return f(ctx)
42+
}
Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
// Copyright The OpenTelemetry Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
package extensionscrapercontroller
5+
6+
import (
7+
"context"
8+
"errors"
9+
"testing"
10+
11+
"github.com/stretchr/testify/assert"
12+
"github.com/stretchr/testify/require"
13+
)
14+
15+
func TestDeregisterFuncNil(t *testing.T) {
16+
var f DeregisterFunc
17+
assert.NoError(t, f.Deregister(context.Background()))
18+
}
19+
20+
func TestDeregisterFuncDelegates(t *testing.T) {
21+
called := false
22+
expectedErr := errors.New("deregister error")
23+
f := DeregisterFunc(func(context.Context) error {
24+
called = true
25+
return expectedErr
26+
})
27+
err := f.Deregister(context.Background())
28+
require.True(t, called)
29+
assert.Equal(t, expectedErr, err)
30+
}

scraper/scraperhelper/controller_test.go

Lines changed: 212 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
"go.opentelemetry.io/collector/component"
2323
"go.opentelemetry.io/collector/component/componenttest"
2424
"go.opentelemetry.io/collector/consumer/consumertest"
25+
"go.opentelemetry.io/collector/extension/xextension/extensionscrapercontroller"
2526
"go.opentelemetry.io/collector/pdata/plog"
2627
"go.opentelemetry.io/collector/pdata/pmetric"
2728
"go.opentelemetry.io/collector/receiver"
@@ -832,3 +833,214 @@ func TestNewMetricsController_ScraperIDInErrorLogs(t *testing.T) {
832833
assert.Equal(t, "test log from receiver", receiverLog.Message)
833834
assert.NotContains(t, receiverLog.ContextMap(), "scraper")
834835
}
836+
837+
// mockHost implements component.Host with configurable extensions.
838+
type mockHost struct {
839+
ext map[component.ID]component.Component
840+
}
841+
842+
func (h *mockHost) GetExtensions() map[component.ID]component.Component {
843+
return h.ext
844+
}
845+
846+
// mockControllerExtension implements extensionscrapercontroller.ControllerExtension.
847+
type mockControllerExtension struct {
848+
component.StartFunc
849+
component.ShutdownFunc
850+
scrapeFunc func(context.Context) error
851+
deregistered bool
852+
}
853+
854+
func (m *mockControllerExtension) RegisterScraper(_ context.Context, scrapeFunc func(context.Context) error) (extensionscrapercontroller.RegistrationHandle, error) {
855+
m.scrapeFunc = scrapeFunc
856+
return extensionscrapercontroller.DeregisterFunc(func(context.Context) error {
857+
m.deregistered = true
858+
return nil
859+
}), nil
860+
}
861+
862+
func TestExtensionTriggersMetricsScrape(t *testing.T) {
863+
t.Parallel()
864+
865+
scrapeCh := make(chan int, 10)
866+
ts := &testScrape{ch: scrapeCh}
867+
868+
scp, err := scraper.NewMetrics(ts.scrapeMetrics)
869+
require.NoError(t, err)
870+
871+
extID := component.MustNewID("myext")
872+
mockExt := &mockControllerExtension{}
873+
874+
cfg := &ControllerConfig{
875+
CollectionInterval: 0,
876+
InitialDelay: 0,
877+
Controllers: []component.ID{extID},
878+
}
879+
880+
recv, err := NewMetricsController(
881+
cfg,
882+
receivertest.NewNopSettings(receivertest.NopType),
883+
new(consumertest.MetricsSink),
884+
AddMetricsScraper(component.MustNewType("scraper"), scp),
885+
)
886+
require.NoError(t, err)
887+
888+
host := &mockHost{ext: map[component.ID]component.Component{extID: mockExt}}
889+
require.NoError(t, recv.Start(context.Background(), host))
890+
891+
// Extension triggers a scrape
892+
require.NotNil(t, mockExt.scrapeFunc, "scrapeFunc should have been registered")
893+
require.NoError(t, mockExt.scrapeFunc(context.Background()))
894+
895+
// Verify scrape was called
896+
select {
897+
case <-scrapeCh:
898+
case <-time.After(time.Second):
899+
t.Fatal("expected scrape to be triggered by extension")
900+
}
901+
902+
require.NoError(t, recv.Shutdown(context.Background()))
903+
assert.True(t, mockExt.deregistered, "expected deregister to be called on shutdown")
904+
}
905+
906+
func TestExtensionAndTimerBothTriggerScrapes(t *testing.T) {
907+
t.Parallel()
908+
909+
scrapeCh := make(chan int, 10)
910+
ts := &testScrape{ch: scrapeCh}
911+
912+
scp, err := scraper.NewMetrics(ts.scrapeMetrics)
913+
require.NoError(t, err)
914+
915+
extID := component.MustNewID("myext")
916+
mockExt := &mockControllerExtension{}
917+
918+
tickerCh := make(chan time.Time)
919+
920+
cfg := &ControllerConfig{
921+
CollectionInterval: time.Second,
922+
InitialDelay: 0,
923+
Controllers: []component.ID{extID},
924+
}
925+
926+
recv, err := NewMetricsController(
927+
cfg,
928+
receivertest.NewNopSettings(receivertest.NopType),
929+
new(consumertest.MetricsSink),
930+
AddMetricsScraper(component.MustNewType("scraper"), scp),
931+
WithTickerChannel(tickerCh),
932+
)
933+
require.NoError(t, err)
934+
935+
host := &mockHost{ext: map[component.ID]component.Component{extID: mockExt}}
936+
require.NoError(t, recv.Start(context.Background(), host))
937+
938+
// Initial scrape on start (from ticker goroutine)
939+
<-scrapeCh
940+
941+
// Extension triggers a scrape
942+
require.NoError(t, mockExt.scrapeFunc(context.Background()))
943+
<-scrapeCh
944+
945+
// Ticker triggers a scrape
946+
tickerCh <- time.Now()
947+
<-scrapeCh
948+
949+
require.NoError(t, recv.Shutdown(context.Background()))
950+
assert.True(t, mockExt.deregistered)
951+
}
952+
953+
func TestExtensionNotFound(t *testing.T) {
954+
t.Parallel()
955+
956+
scp, err := scraper.NewMetrics(func(context.Context) (pmetric.Metrics, error) {
957+
return pmetric.NewMetrics(), nil
958+
})
959+
require.NoError(t, err)
960+
961+
cfg := &ControllerConfig{
962+
CollectionInterval: 0,
963+
InitialDelay: 0,
964+
Controllers: []component.ID{component.MustNewID("missing")},
965+
}
966+
967+
recv, err := NewMetricsController(
968+
cfg,
969+
receivertest.NewNopSettings(receivertest.NopType),
970+
new(consumertest.MetricsSink),
971+
AddMetricsScraper(component.MustNewType("scraper"), scp),
972+
)
973+
require.NoError(t, err)
974+
975+
host := &mockHost{ext: map[component.ID]component.Component{}}
976+
err = recv.Start(context.Background(), host)
977+
require.Error(t, err)
978+
assert.Contains(t, err.Error(), `extension "missing" not found`)
979+
}
980+
981+
func TestExtensionWrongType(t *testing.T) {
982+
t.Parallel()
983+
984+
scp, err := scraper.NewMetrics(func(context.Context) (pmetric.Metrics, error) {
985+
return pmetric.NewMetrics(), nil
986+
})
987+
require.NoError(t, err)
988+
989+
extID := component.MustNewID("wrongtype")
990+
991+
cfg := &ControllerConfig{
992+
CollectionInterval: 0,
993+
InitialDelay: 0,
994+
Controllers: []component.ID{extID},
995+
}
996+
997+
recv, err := NewMetricsController(
998+
cfg,
999+
receivertest.NewNopSettings(receivertest.NopType),
1000+
new(consumertest.MetricsSink),
1001+
AddMetricsScraper(component.MustNewType("scraper"), scp),
1002+
)
1003+
require.NoError(t, err)
1004+
1005+
// Use a plain component that does not implement ControllerExtension
1006+
host := &mockHost{ext: map[component.ID]component.Component{extID: &struct {
1007+
component.StartFunc
1008+
component.ShutdownFunc
1009+
}{}}}
1010+
err = recv.Start(context.Background(), host)
1011+
require.Error(t, err)
1012+
assert.Contains(t, err.Error(), `extension "wrongtype" is not a scraper controller extension`)
1013+
}
1014+
1015+
func TestDeregisterOnShutdown(t *testing.T) {
1016+
t.Parallel()
1017+
1018+
scp, err := scraper.NewMetrics(func(context.Context) (pmetric.Metrics, error) {
1019+
return pmetric.NewMetrics(), nil
1020+
})
1021+
require.NoError(t, err)
1022+
1023+
extID := component.MustNewID("myext")
1024+
mockExt := &mockControllerExtension{}
1025+
1026+
cfg := &ControllerConfig{
1027+
CollectionInterval: 0,
1028+
InitialDelay: 0,
1029+
Controllers: []component.ID{extID},
1030+
}
1031+
1032+
recv, err := NewMetricsController(
1033+
cfg,
1034+
receivertest.NewNopSettings(receivertest.NopType),
1035+
new(consumertest.MetricsSink),
1036+
AddMetricsScraper(component.MustNewType("scraper"), scp),
1037+
)
1038+
require.NoError(t, err)
1039+
1040+
host := &mockHost{ext: map[component.ID]component.Component{extID: mockExt}}
1041+
require.NoError(t, recv.Start(context.Background(), host))
1042+
assert.False(t, mockExt.deregistered)
1043+
1044+
require.NoError(t, recv.Shutdown(context.Background()))
1045+
assert.True(t, mockExt.deregistered)
1046+
}

scraper/scraperhelper/go.mod

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ require (
88
go.opentelemetry.io/collector/component/componenttest v0.150.0
99
go.opentelemetry.io/collector/consumer v1.56.0
1010
go.opentelemetry.io/collector/consumer/consumertest v0.150.0
11+
go.opentelemetry.io/collector/extension/xextension v0.150.0
1112
go.opentelemetry.io/collector/pdata v1.56.0
1213
go.opentelemetry.io/collector/pdata/testdata v0.150.0
1314
go.opentelemetry.io/collector/pipeline v1.56.0
@@ -39,6 +40,7 @@ require (
3940
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
4041
go.opentelemetry.io/collector/consumer/consumererror v0.150.0 // indirect
4142
go.opentelemetry.io/collector/consumer/xconsumer v0.150.0 // indirect
43+
go.opentelemetry.io/collector/extension v1.56.0 // indirect
4244
go.opentelemetry.io/collector/featuregate v1.56.0 // indirect
4345
go.opentelemetry.io/collector/internal/componentalias v0.150.0 // indirect
4446
go.opentelemetry.io/collector/pdata/pprofile v0.150.0 // indirect
@@ -87,4 +89,8 @@ replace go.opentelemetry.io/collector/internal/testutil => ../../internal/testut
8789

8890
replace go.opentelemetry.io/collector/pipeline/xpipeline => ../../pipeline/xpipeline
8991

92+
replace go.opentelemetry.io/collector/extension => ../../extension
93+
94+
replace go.opentelemetry.io/collector/extension/xextension => ../../extension/xextension
95+
9096
replace go.opentelemetry.io/collector/internal/componentalias => ../../internal/componentalias

scraper/scraperhelper/internal/controller/config.go

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ import (
99
"time"
1010

1111
"go.uber.org/multierr"
12+
13+
"go.opentelemetry.io/collector/component"
1214
)
1315

1416
var errNonPositiveInterval = errors.New("requires positive value")
@@ -26,6 +28,11 @@ type ControllerConfig struct {
2628
InitialDelay time.Duration `mapstructure:"initial_delay"`
2729
// Timeout is an optional value used to set scraper's context deadline.
2830
Timeout time.Duration `mapstructure:"timeout"`
31+
// Controllers is a list of extension IDs that control when scrapes occur.
32+
// When specified, extensions can trigger scrapes based on external events.
33+
// If Controllers is non-empty, CollectionInterval may be zero to disable
34+
// time-based scraping entirely.
35+
Controllers []component.ID `mapstructure:"controllers"`
2936
// prevent unkeyed literal initialization
3037
_ struct{}
3138
}
@@ -41,7 +48,7 @@ func NewDefaultControllerConfig() ControllerConfig {
4148
}
4249

4350
func (set *ControllerConfig) Validate() (errs error) {
44-
if set.CollectionInterval <= 0 {
51+
if set.CollectionInterval <= 0 && len(set.Controllers) == 0 {
4552
errs = multierr.Append(errs, fmt.Errorf(`"collection_interval": %w`, errNonPositiveInterval))
4653
}
4754
if set.Timeout < 0 {

0 commit comments

Comments
 (0)