Skip to content

Commit 3aa92a2

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 f9d2c94 commit 3aa92a2

13 files changed

Lines changed: 891 additions & 9 deletions

File tree

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
# Use this changelog template to create an entry for release notes.
2+
3+
# One of 'breaking', 'deprecation', 'new_component', 'enhancement', 'bug_fix'
4+
change_type: enhancement
5+
6+
# The name of the component, or a single word describing the area of concern, (e.g. receiver/otlp)
7+
component: pkg/scraperhelper
8+
9+
# A brief description of the change. Surround your text with quotes ("") if it needs to start with a backtick (`).
10+
note: Add ControllerExtension interface for extensible scraper controllers
11+
12+
# One or more tracking issues or pull requests related to the change
13+
issues: [12449]
14+
15+
# (Optional) One or more lines of additional information to render under the primary note.
16+
# These lines will be padded with 2 spaces and then inserted directly into the document.
17+
# Use pipe (|) for multiline entries.
18+
subtext:
19+
20+
# Optional: The change log or logs in which this entry should be included.
21+
# e.g. '[user]' or '[user, api]'
22+
# Include 'user' if the change is relevant to end users.
23+
# Include 'api' if there is a change to a library API.
24+
# Default: '[user]'
25+
change_logs: [api]

cmd/mdatagen/go.mod

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,8 +88,10 @@ require (
8888
go.opentelemetry.io/collector/config/configtls v1.56.0 // indirect
8989
go.opentelemetry.io/collector/confmap/xconfmap v0.150.0 // indirect
9090
go.opentelemetry.io/collector/consumer/consumererror v0.150.0 // indirect
91+
go.opentelemetry.io/collector/extension v1.56.0 // indirect
9192
go.opentelemetry.io/collector/extension/extensionauth v1.56.0 // indirect
9293
go.opentelemetry.io/collector/extension/extensionmiddleware v0.150.0 // indirect
94+
go.opentelemetry.io/collector/extension/xextension v0.150.0 // indirect
9395
go.opentelemetry.io/collector/internal/componentalias v0.150.0 // indirect
9496
go.opentelemetry.io/collector/internal/fanoutconsumer v0.150.0 // indirect
9597
go.opentelemetry.io/collector/pdata/testdata v0.150.0 // indirect

docs/rfcs/scraper-controller-extension.md

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -40,12 +40,19 @@ 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+
//
44+
// Implementations may call scrapeFunc concurrently. After Deregister on the
45+
// returned handle returns, implementations must not start new invocations of
46+
// scrapeFunc; already in-flight invocations may continue to run.
47+
RegisterScraper(ctx context.Context, scrapeFunc func(context.Context) error) (RegistrationHandle, error)
4448
}
4549

4650
// RegistrationHandle provides a way to deregister a scraper from a controller extension
4751
type RegistrationHandle interface {
48-
// Deregister removes the scraper from the controller extension
52+
// Deregister removes the scraper from the controller extension.
53+
// After Deregister returns, the extension must not start any new
54+
// invocations of the associated scrapeFunc. Deregister need not wait
55+
// for in-flight invocations to complete.
4956
Deregister(ctx context.Context) error
5057
}
5158
```
@@ -76,8 +83,8 @@ Modify `ControllerConfig` to support:
7683
```golang
7784
type ControllerConfig struct {
7885
// CollectionInterval sets how frequently the scraper should be called.
79-
// If zero or negative, the timer-based scraping is disabled.
80-
// At least one controller extension must be configured if timer is disabled.
86+
// Must be positive, or zero to disable timer-based scraping. If zero,
87+
// at least one controller extension must be configured.
8188
CollectionInterval time.Duration `mapstructure:"collection_interval"`
8289

8390
// InitialDelay sets the initial start delay for the scraper timer.
@@ -96,7 +103,7 @@ type ControllerConfig struct {
96103
### 3. Controller Implementation Changes
97104

98105
Modify `controller.Controller` to:
99-
- Support disabling the timer when `CollectionInterval <= 0`
106+
- Support disabling the timer when `CollectionInterval == 0` (negative values remain invalid)
100107
- Register scrapers with configured controller extensions during `Start()`
101108
- Deregister scrapers during `Shutdown()`
102109
- Validate that at least one controller extension is configured if timer is disabled
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: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
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+
//
22+
// Implementations may call scrapeFunc concurrently. After Deregister on
23+
// the returned handle returns, implementations must not start new
24+
// invocations of scrapeFunc; already in-flight invocations may continue
25+
// to run.
26+
RegisterScraper(ctx context.Context, scrapeFunc func(context.Context) error) (RegistrationHandle, error)
27+
}
28+
29+
// RegistrationHandle is returned by ControllerExtension.RegisterScraper and
30+
// is used to deregister the scraper during shutdown.
31+
type RegistrationHandle interface {
32+
// Deregister removes the scraper registration from the extension.
33+
// After Deregister returns, the extension must not start any new
34+
// invocations of the associated scrapeFunc. Deregister need not wait
35+
// for in-flight invocations to complete.
36+
Deregister(ctx context.Context) error
37+
}
38+
39+
// DeregisterFunc is a function that implements RegistrationHandle.
40+
// A nil DeregisterFunc is valid and returns nil on Deregister.
41+
type DeregisterFunc func(ctx context.Context) error
42+
43+
var _ RegistrationHandle = DeregisterFunc(nil)
44+
45+
// Deregister calls the underlying function. If the receiver is nil, it returns nil.
46+
func (f DeregisterFunc) Deregister(ctx context.Context) error {
47+
if f == nil {
48+
return nil
49+
}
50+
return f(ctx)
51+
}
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+
}

0 commit comments

Comments
 (0)