Skip to content

Commit 9be8a28

Browse files
authored
Enable Config v2 in the OBI receiver (#2681)
1 parent 24042fc commit 9be8a28

11 files changed

Lines changed: 611 additions & 30 deletions

File tree

.github/workflows/pull_request.yml

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -241,3 +241,11 @@ jobs:
241241
- name: Build collector distribution
242242
working-directory: examples/otel-collector
243243
run: "$(go env GOPATH)/bin/builder --config ./builder-config.yaml"
244+
245+
- name: Validate Config v2 collector configuration
246+
working-directory: examples/otel-collector
247+
run: ./otelcol-dev/otelcol-dev validate --config ./config.yaml
248+
249+
- name: Smoke test Config v2 collector receiver
250+
working-directory: examples/otel-collector
251+
run: sudo ./smoke-test.sh

collector/config_linux.go

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
1+
// Copyright The OpenTelemetry Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
//go:build linux && (amd64 || arm64)
5+
6+
package collector // import "go.opentelemetry.io/obi/collector"
7+
8+
import (
9+
"errors"
10+
"fmt"
11+
12+
"go.yaml.in/yaml/v3"
13+
14+
"go.opentelemetry.io/collector/confmap"
15+
"go.opentelemetry.io/collector/consumer/consumertest"
16+
17+
"go.opentelemetry.io/obi/internal/config/convert"
18+
"go.opentelemetry.io/obi/internal/config/schema"
19+
"go.opentelemetry.io/obi/pkg/obi"
20+
)
21+
22+
type receiverConfig struct {
23+
runtime *obi.Config
24+
}
25+
26+
func (c *receiverConfig) Unmarshal(component *confmap.Conf) error {
27+
if component == nil {
28+
return nil
29+
}
30+
31+
data, err := yaml.Marshal(component.ToStringMap())
32+
if err != nil {
33+
return fmt.Errorf("marshal OBI receiver config: %w", err)
34+
}
35+
36+
extension, err := schema.ParseReceiverYAML(data)
37+
if err != nil {
38+
var notV2 *schema.NotV2Error
39+
if !errors.As(err, &notV2) {
40+
return fmt.Errorf("parse OBI receiver config v2: %w", err)
41+
}
42+
43+
cfg := defaultRuntimeConfig()
44+
if err := cfg.Unmarshal(component); err != nil {
45+
return fmt.Errorf("parse legacy OBI receiver config: %w", err)
46+
}
47+
c.runtime = cfg
48+
return nil
49+
}
50+
51+
cfg, err := convert.V2ToRuntime(extension)
52+
if err != nil {
53+
return fmt.Errorf("convert OBI receiver config v2: %w", err)
54+
}
55+
setReceiverConsumers(cfg)
56+
c.runtime = cfg
57+
return nil
58+
}
59+
60+
func (c *receiverConfig) Validate() error {
61+
if c == nil || c.runtime == nil {
62+
return errInvalidConfig
63+
}
64+
return c.runtime.Validate()
65+
}
66+
67+
func defaultRuntimeConfig() *obi.Config {
68+
cfg := obi.DefaultConfig
69+
setReceiverConsumers(&cfg)
70+
return &cfg
71+
}
72+
73+
func setReceiverConsumers(cfg *obi.Config) {
74+
cfg.Traces.TracesConsumer = consumertest.NewNop()
75+
cfg.OTELMetrics.MetricsConsumer = consumertest.NewNop()
76+
}

collector/config_linux_test.go

Lines changed: 214 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,214 @@
1+
// Copyright The OpenTelemetry Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
4+
//go:build linux && (amd64 || arm64)
5+
6+
package collector
7+
8+
import (
9+
"context"
10+
"testing"
11+
12+
"github.com/stretchr/testify/require"
13+
14+
"go.opentelemetry.io/collector/component"
15+
"go.opentelemetry.io/collector/component/componenttest"
16+
"go.opentelemetry.io/collector/confmap"
17+
"go.opentelemetry.io/collector/consumer/consumertest"
18+
"go.opentelemetry.io/collector/receiver/receivertest"
19+
20+
"go.opentelemetry.io/obi/internal/config/schema"
21+
"go.opentelemetry.io/obi/pkg/obi"
22+
)
23+
24+
func TestReceiverConfigStructure(t *testing.T) {
25+
require.NoError(t, componenttest.CheckConfigStruct(defaultConfig()))
26+
}
27+
28+
func TestReceiverConfigUnmarshalV2(t *testing.T) {
29+
cfg := newTestReceiverConfig(t)
30+
component := confmap.NewFromStringMap(map[string]any{
31+
"version": "2.0",
32+
"rules": []any{
33+
map[string]any{
34+
"action": "include",
35+
"match": map[string]any{
36+
"process": map[string]any{"open_ports": "8080"},
37+
},
38+
},
39+
},
40+
"channels": map[string]any{"buffer_len": 123},
41+
})
42+
43+
require.NoError(t, component.Unmarshal(cfg))
44+
require.NoError(t, cfg.Validate())
45+
require.Equal(t, 123, cfg.runtime.ChannelBufferLen)
46+
require.Len(t, cfg.runtime.Discovery.Instrument, 1)
47+
require.Equal(t, []int{8080}, cfg.runtime.Discovery.Instrument[0].OpenPorts.AllValues())
48+
require.NotNil(t, cfg.runtime.Traces.TracesConsumer)
49+
require.NotNil(t, cfg.runtime.OTELMetrics.MetricsConsumer)
50+
}
51+
52+
func TestReceiverConfigUnmarshalLegacy(t *testing.T) {
53+
cfg := newTestReceiverConfig(t)
54+
component := confmap.NewFromStringMap(map[string]any{
55+
"open_port": "8080",
56+
})
57+
58+
require.NoError(t, component.Unmarshal(cfg))
59+
require.NoError(t, cfg.Validate())
60+
require.Equal(t, []int{8080}, cfg.runtime.Port.AllValues())
61+
}
62+
63+
func TestReceiverConfigRejectsStandaloneSections(t *testing.T) {
64+
layouts := []struct {
65+
name string
66+
key string
67+
value any
68+
}{
69+
{name: "v2", key: "version", value: "2.0"},
70+
{name: "legacy selector", key: "open_port", value: "8080"},
71+
}
72+
for _, layout := range layouts {
73+
t.Run(layout.name, func(t *testing.T) {
74+
for _, section := range []string{"enrich", "correlation", "daemon"} {
75+
t.Run(section, func(t *testing.T) {
76+
cfg := newTestReceiverConfig(t)
77+
component := confmap.NewFromStringMap(map[string]any{
78+
layout.key: layout.value,
79+
section: map[string]any{},
80+
})
81+
82+
err := component.Unmarshal(cfg)
83+
84+
var notAllowed *schema.SectionNotAllowedError
85+
require.ErrorAs(t, err, &notAllowed)
86+
require.Equal(t, section, notAllowed.Section)
87+
require.Contains(t, err.Error(), "receiver config")
88+
require.Contains(t, err.Error(), "standalone mode")
89+
})
90+
}
91+
})
92+
}
93+
}
94+
95+
func TestReceiverConfigDoesNotFallbackFromInvalidV2(t *testing.T) {
96+
tests := []struct {
97+
name string
98+
component map[string]any
99+
check func(*testing.T, error)
100+
}{
101+
{
102+
name: "unsupported version",
103+
component: map[string]any{
104+
"version": "3.0",
105+
"channel_buffer_len": 123,
106+
},
107+
check: func(t *testing.T, err error) {
108+
var unsupported *schema.UnsupportedVersionError
109+
require.ErrorAs(t, err, &unsupported)
110+
require.Equal(t, "3.0", unsupported.Version)
111+
},
112+
},
113+
{
114+
name: "invalid capture value",
115+
component: map[string]any{
116+
"version": "2.0",
117+
"network": map[string]any{
118+
"capture": map[string]any{"source": "invalid"},
119+
},
120+
"channel_buffer_len": 123,
121+
},
122+
check: func(t *testing.T, err error) {
123+
require.Contains(t, err.Error(), "invalid source")
124+
var notV2 *schema.NotV2Error
125+
require.NotErrorAs(t, err, &notV2)
126+
},
127+
},
128+
{
129+
name: "standalone v2 layout with legacy selector",
130+
component: map[string]any{
131+
"file_format": "1.0",
132+
"extensions": map[string]any{
133+
"obi": map[string]any{
134+
"version": "2.0",
135+
"capture": map[string]any{},
136+
},
137+
},
138+
"open_port": "8080",
139+
"channel_buffer_len": 123,
140+
},
141+
check: func(t *testing.T, err error) {
142+
var wrongLayout *schema.ReceiverLayoutError
143+
require.ErrorAs(t, err, &wrongLayout)
144+
var notV2 *schema.NotV2Error
145+
require.NotErrorAs(t, err, &notV2)
146+
},
147+
},
148+
}
149+
150+
for _, test := range tests {
151+
t.Run(test.name, func(t *testing.T) {
152+
cfg := newTestReceiverConfig(t)
153+
154+
err := confmap.NewFromStringMap(test.component).Unmarshal(cfg)
155+
156+
require.Error(t, err)
157+
test.check(t, err)
158+
require.Equal(t, obi.DefaultConfig.ChannelBufferLen, cfg.runtime.ChannelBufferLen)
159+
})
160+
}
161+
}
162+
163+
func TestReceiverV2PipelineModes(t *testing.T) {
164+
tests := []struct {
165+
name string
166+
traces bool
167+
metrics bool
168+
}{
169+
{name: "traces", traces: true},
170+
{name: "metrics", metrics: true},
171+
{name: "traces and metrics", traces: true, metrics: true},
172+
}
173+
174+
for _, test := range tests {
175+
t.Run(test.name, func(t *testing.T) {
176+
cfg := newTestReceiverConfig(t)
177+
require.NoError(t, confmap.NewFromStringMap(map[string]any{
178+
"version": "2.0",
179+
}).Unmarshal(cfg))
180+
181+
factory := NewFactory()
182+
settings := receivertest.NewNopSettings(typeStr)
183+
settings.ID = component.MustNewIDWithName("obi", test.name)
184+
185+
if test.traces {
186+
consumer := consumertest.NewNop()
187+
receiver, err := factory.CreateTraces(t.Context(), settings, cfg, consumer)
188+
require.NoError(t, err)
189+
require.Same(t, consumer, cfg.runtime.Traces.TracesConsumer)
190+
t.Cleanup(func() {
191+
require.NoError(t, receiver.Shutdown(context.Background()))
192+
})
193+
}
194+
195+
if test.metrics {
196+
consumer := consumertest.NewNop()
197+
receiver, err := factory.CreateMetrics(t.Context(), settings, cfg, consumer)
198+
require.NoError(t, err)
199+
require.Same(t, consumer, cfg.runtime.OTELMetrics.MetricsConsumer)
200+
t.Cleanup(func() {
201+
require.NoError(t, receiver.Shutdown(context.Background()))
202+
})
203+
}
204+
})
205+
}
206+
}
207+
208+
func newTestReceiverConfig(t *testing.T) *receiverConfig {
209+
t.Helper()
210+
211+
cfg, ok := NewFactory().CreateDefaultConfig().(*receiverConfig)
212+
require.True(t, ok)
213+
return cfg
214+
}

collector/factory_linux.go

Lines changed: 14 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,6 @@ import (
1414

1515
"go.opentelemetry.io/collector/component"
1616
"go.opentelemetry.io/collector/consumer"
17-
"go.opentelemetry.io/collector/consumer/consumertest"
1817
"go.opentelemetry.io/collector/receiver"
1918

2019
"go.opentelemetry.io/obi/collector/internal"
@@ -37,7 +36,7 @@ func BuildTracesReceiver() receiver.CreateTracesFunc {
3736
) (receiver.Traces, error) {
3837
initLogger(rs)
3938

40-
cfg, ok := baseCfg.(*obi.Config)
39+
cfg, ok := runtimeConfig(baseCfg)
4140
if !ok {
4241
return nil, errInvalidConfig
4342
}
@@ -55,7 +54,7 @@ func BuildMetricsReceiver() receiver.CreateMetricsFunc {
5554
) (receiver.Metrics, error) {
5655
initLogger(rs)
5756

58-
cfg, ok := baseCfg.(*obi.Config)
57+
cfg, ok := runtimeConfig(baseCfg)
5958
if !ok {
6059
return nil, errInvalidConfig
6160
}
@@ -66,10 +65,16 @@ func BuildMetricsReceiver() receiver.CreateMetricsFunc {
6665
}
6766

6867
func defaultConfig() component.Config {
69-
cfg := obi.DefaultConfig
70-
// These are placeholders for the consumers; without these obi config will be invalid.
71-
// The actual consumers are set when the receiver is created.
72-
cfg.Traces.TracesConsumer = consumertest.NewNop()
73-
cfg.OTELMetrics.MetricsConsumer = consumertest.NewNop()
74-
return &cfg
68+
return &receiverConfig{runtime: defaultRuntimeConfig()}
69+
}
70+
71+
func runtimeConfig(baseCfg component.Config) (*obi.Config, bool) {
72+
switch cfg := baseCfg.(type) {
73+
case *receiverConfig:
74+
return cfg.runtime, cfg.runtime != nil
75+
case *obi.Config:
76+
return cfg, cfg != nil
77+
default:
78+
return nil, false
79+
}
7580
}

examples/otel-collector/README.md

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,17 @@ The collector requires `sudo` to attach eBPF probes to processes.
4141

4242
## Testing the Collector
4343

44+
After building the collector, run the same end-to-end smoke test used by the
45+
pull request workflow:
46+
47+
```bash
48+
sudo ./smoke-test.sh
49+
```
50+
51+
The smoke test starts an HTTP server and the Collector with `smoke-config.yaml`,
52+
then verifies that the Config v2 OBI receiver exports the server's trace through
53+
the debug exporter.
54+
4455
Once the collector is running, you can generate some test traces:
4556

4657
1. In a new terminal, start a simple HTTP server:
@@ -130,11 +141,12 @@ Once the collector is running, you can generate some test traces:
130141

131142
The `config.yaml` file defines:
132143

133-
- **OBI receiver**: Listens on port 8000 for HTTP traffic and automatically instruments services
144+
- **OBI receiver**: Uses Config v2 to capture services listening on port 8000
134145
- **OTLP receiver**: Accepts spans from manually instrumented applications
135146
- **Batch processor**: Groups spans for efficient export
136147
- **Debug exporter**: Prints spans to logs (useful for debugging)
137148
- **OTLP exporter**: Sends spans to a Jaeger backend (requires Jaeger to be running)
149+
- **Trace and metric pipelines**: Share one OBI receiver instance
138150

139151
You can modify `config.yaml` to:
140152

0 commit comments

Comments
 (0)