Skip to content

Commit 0fd006d

Browse files
committed
Address feedback
1 parent f9cfcc1 commit 0fd006d

10 files changed

Lines changed: 332 additions & 146 deletions

File tree

translator/tocwconfig/sampleConfig/opentelemetry/default_otel_config_windows.yaml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -684,16 +684,16 @@ service:
684684
exporters:
685685
- forward/opentelemetry
686686
processors:
687-
- transform/windows_events_scope
688687
- filter/windows_events_application
688+
- transform/windows_events_scope
689689
receivers:
690690
- windowseventlog/application
691691
logs/windows_events_system:
692692
exporters:
693693
- forward/opentelemetry
694694
processors:
695-
- transform/windows_events_scope
696695
- filter/windows_events_system
696+
- transform/windows_events_scope
697697
receivers:
698698
- windowseventlog/system
699699
metrics/opentelemetry:

translator/tocwconfig/sampleConfig/opentelemetry/windows_events_config.yaml

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,15 @@ processors:
105105
metrics: {}
106106
spans: {}
107107
traces: {}
108+
resource/windows_events_system:
109+
attributes:
110+
- action: upsert
111+
converted_type: ""
112+
from_attribute: ""
113+
from_context: ""
114+
key: aws.log.group.name
115+
pattern: ""
116+
value: /aws/cwagent/windows-events/System
108117
resourcedetection/opentelemetry:
109118
aks:
110119
resource_attributes:
@@ -482,7 +491,6 @@ receivers:
482491
poll_interval: 1s
483492
resource:
484493
aws.log.channel: System
485-
aws.log.group.name: /aws/cwagent/windows-events/System
486494
aws.log.source: windows_events
487495
retry_on_failure:
488496
enabled: false
@@ -516,16 +524,17 @@ service:
516524
exporters:
517525
- forward/opentelemetry
518526
processors:
519-
- transform/windows_events_scope
520527
- filter/windows_events_application
528+
- transform/windows_events_scope
521529
receivers:
522530
- windowseventlog/application
523531
logs/windows_events_system:
524532
exporters:
525533
- forward/opentelemetry
526534
processors:
527-
- transform/windows_events_scope
528535
- filter/windows_events_system
536+
- resource/windows_events_system
537+
- transform/windows_events_scope
529538
receivers:
530539
- windowseventlog/system
531540
telemetry:

translator/translate/otel/pipeline/opentelemetry/translator_logs.go

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,11 +45,14 @@ var otelLogsKeys = []string{
4545
common.OtelCollectLogsConfigKey,
4646
common.DatabaseInsightsConfigKey,
4747
common.ConfigKey(common.OpenTelemetryKey, common.CollectKey, common.OtlpKey),
48-
common.WindowsEventsConfigKey,
4948
}
5049

5150
func (t *baseLogsTranslator) Translate(conf *confmap.Conf) (*common.ComponentTranslators, error) {
52-
if err := common.ValidateAnySet(conf, t.ID(), otelLogsKeys); err != nil {
51+
keys := otelLogsKeys
52+
if runtime.GOOS == "windows" {
53+
keys = append(keys, common.WindowsEventsConfigKey)
54+
}
55+
if err := common.ValidateAnySet(conf, t.ID(), keys); err != nil {
5356
return nil, err
5457
}
5558

translator/translate/otel/pipeline/opentelemetry/translator_logs_test.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
package opentelemetry
55

66
import (
7+
"runtime"
78
"strings"
89
"testing"
910

@@ -19,7 +20,11 @@ func TestBaseLogsTranslator(t *testing.T) {
1920
tt := NewBaseLogsTranslator()
2021
assert.EqualValues(t, "logs/opentelemetry", tt.ID().String())
2122

22-
missingErr := &common.MissingKeyError{ID: tt.ID(), JsonKey: strings.Join(otelLogsKeys, " or ")}
23+
keys := otelLogsKeys
24+
if runtime.GOOS == "windows" {
25+
keys = append(keys, common.WindowsEventsConfigKey)
26+
}
27+
missingErr := &common.MissingKeyError{ID: tt.ID(), JsonKey: strings.Join(keys, " or ")}
2328
testCases := map[string]struct {
2429
input map[string]interface{}
2530
wantErr error

translator/translate/otel/pipeline/opentelemetry/windowsevents/translator.go

Lines changed: 32 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import (
1515
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/connector/forward"
1616
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/extension/filestorage"
1717
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/processor/filterprocessor"
18+
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/processor/resourceprocessor"
1819
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/processor/transformprocessor"
1920
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/receiver/windowseventlog"
2021
)
@@ -35,13 +36,26 @@ var severityNumbers = map[string]int{
3536
}
3637

3738
type eventEntry struct {
38-
name string
39-
receiverName string
40-
channel string
41-
raw bool
42-
resource map[string]string
43-
eventLevels []string
44-
eventIDs []int
39+
name string
40+
receiverName string
41+
channel string
42+
raw bool
43+
resource map[string]string
44+
logGroupName string
45+
logStreamName string
46+
eventLevels []string
47+
eventIDs []int
48+
}
49+
50+
func (e eventEntry) routingAttributes() map[string]string {
51+
attrs := make(map[string]string)
52+
if e.logGroupName != "" {
53+
attrs["aws.log.group.name"] = e.logGroupName
54+
}
55+
if e.logStreamName != "" {
56+
attrs["aws.log.stream.name"] = e.logStreamName
57+
}
58+
return attrs
4559
}
4660

4761
// filterCondition builds an OTTL drop condition. When both levels and IDs are present, they are ANDed.
@@ -99,17 +113,24 @@ func (t *windowsEventsPipelineTranslator) Translate(_ *confmap.Conf) (*common.Co
99113
receivers.Set(windowseventlog.NewTranslator(t.entry.receiverName, t.entry.channel, t.entry.raw, t.entry.resource))
100114

101115
processors := common.NewTranslatorMap[component.Config, component.ID]()
102-
processors.Set(transformprocessor.NewTranslatorWithName("windows_events_scope",
103-
transformprocessor.WithErrorMode(common.OTTLErrorModeIgnore),
104-
transformprocessor.WithLogScopeStatements(common.ScopeStatementsForSolution("otel-windows-events")),
105-
))
106116

107117
// TODO: Replace with upstream Query XML filtering when collector is bumped past v0.124.
108118
condition := t.entry.filterCondition()
109119
if condition != "" {
110120
processors.Set(filterprocessor.NewTranslatorWithLogCondition("windows_events_"+t.entry.name, condition, common.OTTLErrorModeIgnore))
111121
}
112122

123+
if attrs := t.entry.routingAttributes(); len(attrs) > 0 {
124+
processors.Set(resourceprocessor.NewTranslator(
125+
common.WithName("windows_events_"+t.entry.name),
126+
resourceprocessor.WithAttributes(attrs),
127+
))
128+
}
129+
processors.Set(transformprocessor.NewTranslatorWithName("windows_events_scope",
130+
transformprocessor.WithErrorMode(common.OTTLErrorModeIgnore),
131+
transformprocessor.WithLogScopeStatements(common.ScopeStatementsForSolution("otel-windows-events")),
132+
))
133+
113134
return &common.ComponentTranslators{
114135
Receivers: receivers,
115136
Processors: processors,

translator/translate/otel/pipeline/opentelemetry/windowsevents/translator_test.go

Lines changed: 42 additions & 106 deletions
Original file line numberDiff line numberDiff line change
@@ -8,58 +8,9 @@ import (
88

99
"github.com/stretchr/testify/assert"
1010
"github.com/stretchr/testify/require"
11-
"go.opentelemetry.io/collector/confmap"
1211
"go.opentelemetry.io/collector/pipeline"
13-
14-
translatorconfig "github.com/aws/amazon-cloudwatch-agent/translator/config"
15-
translatorcontext "github.com/aws/amazon-cloudwatch-agent/translator/context"
1612
)
1713

18-
func TestNewTranslators_Disabled(t *testing.T) {
19-
conf := confmap.NewFromStringMap(map[string]any{
20-
"opentelemetry": map[string]any{"collect": map[string]any{}},
21-
})
22-
translators := NewTranslators(conf)
23-
assert.Equal(t, 0, translators.Len())
24-
}
25-
26-
func TestNewTranslators_TwoEntries(t *testing.T) {
27-
translatorcontext.CurrentContext().SetOs(translatorconfig.OS_TYPE_WINDOWS)
28-
defer translatorcontext.CurrentContext().SetOs("")
29-
conf := confmap.NewFromStringMap(map[string]any{
30-
"opentelemetry": map[string]any{
31-
"collect": map[string]any{
32-
"windows_events": map[string]any{
33-
"collect_list": []any{
34-
map[string]any{"event_name": "System"},
35-
map[string]any{"event_name": "Application"},
36-
},
37-
},
38-
},
39-
},
40-
})
41-
translators := NewTranslators(conf)
42-
assert.Equal(t, 2, translators.Len())
43-
}
44-
45-
func TestNewTranslators_NonWindows(t *testing.T) {
46-
translatorcontext.CurrentContext().SetOs(translatorconfig.OS_TYPE_LINUX)
47-
defer translatorcontext.CurrentContext().SetOs("")
48-
conf := confmap.NewFromStringMap(map[string]any{
49-
"opentelemetry": map[string]any{
50-
"collect": map[string]any{
51-
"windows_events": map[string]any{
52-
"collect_list": []any{
53-
map[string]any{"event_name": "System"},
54-
},
55-
},
56-
},
57-
},
58-
})
59-
translators := NewTranslators(conf)
60-
assert.Equal(t, 0, translators.Len())
61-
}
62-
6314
func TestPipelineTranslator_ID(t *testing.T) {
6415
pt := &windowsEventsPipelineTranslator{entry: eventEntry{name: "system_0"}}
6516
assert.Equal(t, pipeline.NewIDWithName(pipeline.SignalLogs, "windows_events_system_0"), pt.ID())
@@ -101,37 +52,19 @@ func TestPipelineTranslator_Translate_WithFilter(t *testing.T) {
10152
assert.Equal(t, 2, result.Processors.Len())
10253
}
10354

104-
func TestParseEntries(t *testing.T) {
105-
conf := confmap.NewFromStringMap(map[string]any{
106-
"opentelemetry": map[string]any{
107-
"collect": map[string]any{
108-
"windows_events": map[string]any{
109-
"collect_list": []any{
110-
map[string]any{"event_name": "System", "event_levels": []any{"ERROR"}, "event_format": "xml", "log_group_name": "/custom/system"},
111-
map[string]any{"event_name": "Application", "event_ids": []any{float64(1001)}},
112-
map[string]any{"event_name": "Microsoft-Windows-PowerShell/Operational", "event_levels": []any{"WARNING"}},
113-
},
114-
},
115-
},
116-
},
117-
})
118-
entries := parseEntries(conf)
119-
require.Len(t, entries, 3)
120-
121-
assert.Equal(t, "system", entries[0].name)
122-
assert.Equal(t, "System", entries[0].channel)
123-
assert.True(t, entries[0].raw)
124-
assert.Equal(t, "/custom/system", entries[0].resource["aws.log.group.name"])
125-
assert.Equal(t, []string{"ERROR"}, entries[0].eventLevels)
126-
127-
assert.Equal(t, "application", entries[1].name)
128-
assert.Equal(t, "Application", entries[1].channel)
129-
assert.False(t, entries[1].raw)
130-
assert.Equal(t, []int{1001}, entries[1].eventIDs)
131-
132-
assert.Equal(t, "microsoft-windows-powershell_operational", entries[2].name)
133-
assert.Equal(t, "Microsoft-Windows-PowerShell/Operational", entries[2].channel)
134-
assert.Equal(t, []string{"WARNING"}, entries[2].eventLevels)
55+
func TestPipelineTranslator_Translate_WithRoutingAttrs(t *testing.T) {
56+
pt := &windowsEventsPipelineTranslator{entry: eventEntry{
57+
name: "system",
58+
receiverName: "system",
59+
channel: "System",
60+
resource: map[string]string{"aws.log.source": "windows_events", "aws.log.channel": "System"},
61+
logGroupName: "/custom/group",
62+
}}
63+
result, err := pt.Translate(nil)
64+
require.NoError(t, err)
65+
66+
// resource processor + scope transform = 2 processors
67+
assert.Equal(t, 2, result.Processors.Len())
13568
}
13669

13770
func TestPipelineTranslator_DuplicateChannels_SharedReceiver(t *testing.T) {
@@ -162,32 +95,6 @@ func TestPipelineTranslator_DuplicateChannels_SharedReceiver(t *testing.T) {
16295
assert.Equal(t, r1.Receivers.Keys(), r2.Receivers.Keys())
16396
}
16497

165-
func TestParseEntries_DuplicateChannels(t *testing.T) {
166-
conf := confmap.NewFromStringMap(map[string]any{
167-
"opentelemetry": map[string]any{
168-
"collect": map[string]any{
169-
"windows_events": map[string]any{
170-
"collect_list": []any{
171-
map[string]any{"event_name": "System", "event_levels": []any{"ERROR"}},
172-
map[string]any{"event_name": "System", "event_levels": []any{"WARNING"}},
173-
map[string]any{"event_name": "System", "event_ids": []any{float64(1001)}},
174-
},
175-
},
176-
},
177-
},
178-
})
179-
entries := parseEntries(conf)
180-
require.Len(t, entries, 3)
181-
assert.Equal(t, "system", entries[0].name)
182-
assert.Equal(t, "system_1", entries[1].name)
183-
assert.Equal(t, "system_2", entries[2].name)
184-
185-
// All entries share the same receiver (same channel = same checkpoint)
186-
assert.Equal(t, "system", entries[0].receiverName)
187-
assert.Equal(t, "system", entries[1].receiverName)
188-
assert.Equal(t, "system", entries[2].receiverName)
189-
}
190-
19198
func TestPipelineTranslator_Translate_XmlWithEventIDs_Error(t *testing.T) {
19299
pt := &windowsEventsPipelineTranslator{entry: eventEntry{
193100
name: "security",
@@ -237,3 +144,32 @@ func TestBuildFilterCondition(t *testing.T) {
237144
})
238145
}
239146
}
147+
148+
func TestRoutingAttributes(t *testing.T) {
149+
tests := []struct {
150+
name string
151+
entry eventEntry
152+
expected map[string]string
153+
}{
154+
{
155+
name: "no routing attrs",
156+
entry: eventEntry{name: "system"},
157+
expected: map[string]string{},
158+
},
159+
{
160+
name: "log group only",
161+
entry: eventEntry{name: "system", logGroupName: "/custom/group"},
162+
expected: map[string]string{"aws.log.group.name": "/custom/group"},
163+
},
164+
{
165+
name: "both",
166+
entry: eventEntry{name: "system", logGroupName: "/custom/group", logStreamName: "my-stream"},
167+
expected: map[string]string{"aws.log.group.name": "/custom/group", "aws.log.stream.name": "my-stream"},
168+
},
169+
}
170+
for _, tt := range tests {
171+
t.Run(tt.name, func(t *testing.T) {
172+
assert.Equal(t, tt.expected, tt.entry.routingAttributes())
173+
})
174+
}
175+
}

0 commit comments

Comments
 (0)