Skip to content

Commit 242b323

Browse files
committed
Add OTLP logs pipeline for application signals
1 parent 0e07319 commit 242b323

7 files changed

Lines changed: 234 additions & 1 deletion

File tree

translator/cmdutil/translatorutil.go

Lines changed: 58 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -251,7 +251,64 @@ func TranslateJsonMapToYamlConfig(jsonConfigValue interface{}) (interface{}, err
251251
if err != nil {
252252
return nil, err
253253
}
254-
return mapstructure.Marshal(cfg)
254+
result, err := mapstructure.Marshal(cfg)
255+
if err != nil {
256+
return nil, err
257+
}
258+
// configopaque.String values are nil'd during marshal to prevent secret leakage.
259+
// Restore non-secret headers for otlphttp/application_signals exporter (log group/stream).
260+
injectAppSignalsLogsHeaders(result, jsonConfigValue)
261+
return result, nil
262+
}
263+
264+
// injectAppSignalsLogsHeaders restores x-aws-log-group and x-aws-log-stream headers
265+
// on the otlphttp/application_signals exporter. These are nil'd by NilHookFunc during
266+
// mapstructure marshal because configopaque.String is treated as a secret, but log
267+
// group/stream names are not sensitive and must survive the YAML round-trip.
268+
func injectAppSignalsLogsHeaders(yamlMap map[string]any, jsonConfig interface{}) {
269+
jsonMap, ok := jsonConfig.(map[string]interface{})
270+
if !ok {
271+
return
272+
}
273+
exporters, ok := yamlMap["exporters"].(map[string]any)
274+
if !ok {
275+
return
276+
}
277+
exporterCfg, ok := exporters["otlphttp/application_signals"].(map[string]any)
278+
if !ok {
279+
return
280+
}
281+
282+
// Read log_group_name and log_stream_name from JSON config
283+
logGroup := "/aws/application-signals/data"
284+
logStream := "default"
285+
if logs, ok := jsonMap["logs"].(map[string]interface{}); ok {
286+
if lc, ok := logs["logs_collected"].(map[string]interface{}); ok {
287+
if as, ok := lc["application_signals"].(map[string]interface{}); ok {
288+
if v, ok := as["log_group_name"].(string); ok && v != "" {
289+
logGroup = v
290+
}
291+
if v, ok := as["log_stream_name"].(string); ok && v != "" {
292+
logStream = v
293+
}
294+
} else if as, ok := lc["app_signals"].(map[string]interface{}); ok {
295+
if v, ok := as["log_group_name"].(string); ok && v != "" {
296+
logGroup = v
297+
}
298+
if v, ok := as["log_stream_name"].(string); ok && v != "" {
299+
logStream = v
300+
}
301+
}
302+
}
303+
}
304+
305+
headers, ok := exporterCfg["headers"].(map[string]any)
306+
if !ok {
307+
headers = make(map[string]any)
308+
exporterCfg["headers"] = headers
309+
}
310+
headers["x-aws-log-group"] = logGroup
311+
headers["x-aws-log-stream"] = logStream
255312
}
256313

257314
func ConfigToTomlFile(config interface{}, tomlConfigFilePath string) error {

translator/config/schema.json

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -633,6 +633,34 @@
633633
},
634634
"windows_events": {
635635
"$ref": "#/definitions/logsDefinition/definitions/logsWindowsEventsDefinition"
636+
},
637+
"application_signals": {
638+
"type": "object",
639+
"properties": {
640+
"log_group_name": {
641+
"type": "string",
642+
"minLength": 1
643+
},
644+
"log_stream_name": {
645+
"type": "string",
646+
"minLength": 1
647+
}
648+
},
649+
"additionalProperties": true
650+
},
651+
"app_signals": {
652+
"type": "object",
653+
"properties": {
654+
"log_group_name": {
655+
"type": "string",
656+
"minLength": 1
657+
},
658+
"log_stream_name": {
659+
"type": "string",
660+
"minLength": 1
661+
}
662+
},
663+
"additionalProperties": true
636664
}
637665
},
638666
"minProperties": 1,

translator/translate/otel/common/common.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -136,12 +136,15 @@ const (
136136
var (
137137
AppSignalsTraces = ConfigKey(TracesKey, TracesCollectedKey, AppSignals)
138138
AppSignalsMetrics = ConfigKey(LogsKey, MetricsCollectedKey, AppSignals)
139+
AppSignalsLogs = ConfigKey(LogsKey, LogsCollectedKey, AppSignals)
139140
AppSignalsTracesFallback = ConfigKey(TracesKey, TracesCollectedKey, AppSignalsFallback)
140141
AppSignalsMetricsFallback = ConfigKey(LogsKey, MetricsCollectedKey, AppSignalsFallback)
142+
AppSignalsLogsFallback = ConfigKey(LogsKey, LogsCollectedKey, AppSignalsFallback)
141143

142144
AppSignalsConfigKeys = map[pipeline.Signal][]string{
143145
pipeline.SignalTraces: {AppSignalsTraces, AppSignalsTracesFallback},
144146
pipeline.SignalMetrics: {AppSignalsMetrics, AppSignalsMetricsFallback},
147+
pipeline.SignalLogs: {AppSignalsLogs, AppSignalsLogsFallback},
145148
}
146149
SystemMetricsEnabledConfigKey = ConfigKey(AgentKey, SystemMetricsEnabledKey)
147150
JmxConfigKey = ConfigKey(MetricsKey, MetricsCollectedKey, JmxKey)
Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
2+
// SPDX-License-Identifier: MIT
3+
4+
package otlphttp
5+
6+
import (
7+
"fmt"
8+
9+
"go.opentelemetry.io/collector/component"
10+
"go.opentelemetry.io/collector/config/configauth"
11+
"go.opentelemetry.io/collector/confmap"
12+
"go.opentelemetry.io/collector/exporter"
13+
"go.opentelemetry.io/collector/exporter/otlphttpexporter"
14+
15+
"github.com/aws/amazon-cloudwatch-agent/translator/translate/agent"
16+
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/common"
17+
)
18+
19+
type translator struct {
20+
name string
21+
factory exporter.Factory
22+
}
23+
24+
var _ common.ComponentTranslator = (*translator)(nil)
25+
26+
func NewTranslatorWithName(name string) common.ComponentTranslator {
27+
return &translator{name, otlphttpexporter.NewFactory()}
28+
}
29+
30+
func (t *translator) ID() component.ID {
31+
return component.NewIDWithName(t.factory.Type(), t.name)
32+
}
33+
34+
// Translate creates an otlphttp exporter config that sends OTLP logs to the
35+
// CloudWatch OTLP endpoint with SigV4 authentication.
36+
func (t *translator) Translate(_ *confmap.Conf) (component.Config, error) {
37+
cfg := t.factory.CreateDefaultConfig().(*otlphttpexporter.Config)
38+
39+
region := agent.Global_Config.Region
40+
if region == "" {
41+
return nil, fmt.Errorf("region is required for otlphttp exporter")
42+
}
43+
44+
cfg.ClientConfig.Endpoint = fmt.Sprintf("https://logs.%s.amazonaws.com", region)
45+
cfg.ClientConfig.Auth = &configauth.Authentication{
46+
AuthenticatorID: component.NewID(component.MustNewType(common.SigV4Auth)),
47+
}
48+
49+
// Note: x-aws-log-group and x-aws-log-stream headers are injected as raw strings
50+
// in translatorutil.go:injectAppSignalsLogsHeaders() because configopaque.String
51+
// values are nil'd during mapstructure marshal (NilHookFunc) to prevent secret leakage.
52+
53+
return cfg, nil
54+
}

translator/translate/otel/pipeline/applicationsignals/translator.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,9 +15,11 @@ import (
1515
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/exporter/awsemf"
1616
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/exporter/awsxray"
1717
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/exporter/debug"
18+
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/exporter/otlphttp"
1819
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/extension/agenthealth"
1920
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/extension/awsproxy"
2021
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/extension/k8smetadata"
22+
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/extension/sigv4auth"
2123
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/processor/awsapplicationsignals"
2224
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/processor/awsentity"
2325
"github.com/aws/amazon-cloudwatch-agent/translator/translate/otel/processor/metricstransformprocessor"
@@ -58,6 +60,20 @@ func (t *translator) Translate(conf *confmap.Conf) (*common.ComponentTranslators
5860
Extensions: common.NewTranslatorMap[component.Config, component.ID](),
5961
}
6062

63+
if t.signal == pipeline.SignalLogs {
64+
// OTLP Logs pipeline: receive OTLP logs from instrumentation, forward to
65+
// CloudWatch OTLP endpoint via otlphttp exporter with SigV4 auth.
66+
// No processors needed — logs are forwarded as-is to preserve full OTLP structure.
67+
if enabled, _ := common.GetBool(conf, common.AgentDebugConfigKey); enabled {
68+
translators.Exporters.Set(debug.NewTranslator(common.WithName(common.AppSignals)))
69+
}
70+
translators.Exporters.Set(otlphttp.NewTranslatorWithName(common.AppSignals))
71+
translators.Extensions.Set(sigv4auth.NewTranslator())
72+
translators.Extensions.Set(agenthealth.NewTranslator(agenthealth.LogsName, []string{agenthealth.OperationPutLogEvents}))
73+
translators.Extensions.Set(agenthealth.NewTranslatorWithStatusCode(agenthealth.StatusCodeName, nil, true))
74+
return translators, nil
75+
}
76+
6177
if t.signal == pipeline.SignalMetrics {
6278
translators.Processors.Set(metricstransformprocessor.NewTranslatorWithName(common.AppSignals))
6379
}

translator/translate/otel/pipeline/applicationsignals/translator_test.go

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -339,3 +339,77 @@ func TestTranslatorMetricsForECS(t *testing.T) {
339339
})
340340
}
341341
}
342+
343+
func TestTranslatorLogs(t *testing.T) {
344+
type want struct {
345+
receivers []string
346+
processors []string
347+
exporters []string
348+
extensions []string
349+
}
350+
tt := NewTranslator(pipeline.SignalLogs)
351+
assert.EqualValues(t, "logs/application_signals", tt.ID().String())
352+
testCases := map[string]struct {
353+
input map[string]interface{}
354+
want *want
355+
wantErr error
356+
isEKSCache func() eksdetector.IsEKSCache
357+
}{
358+
"WithoutLogsCollectedKey": {
359+
input: map[string]interface{}{},
360+
wantErr: &common.MissingKeyError{ID: tt.ID(), JsonKey: fmt.Sprint(common.AppSignalsLogs)},
361+
},
362+
"WithAppSignalsEnabledLogs": {
363+
input: map[string]interface{}{
364+
"logs": map[string]interface{}{
365+
"logs_collected": map[string]interface{}{
366+
"application_signals": map[string]interface{}{},
367+
},
368+
},
369+
},
370+
want: &want{
371+
receivers: []string{"otlp/grpc_0_0_0_0_4315", "otlp/http_0_0_0_0_4316"},
372+
processors: []string{},
373+
exporters: []string{"otlphttp/application_signals"},
374+
extensions: []string{"sigv4auth", "agenthealth/logs", "agenthealth/statuscode"},
375+
},
376+
isEKSCache: eksdetector.TestIsEKSCacheEKS,
377+
},
378+
"WithAppSignalsLogsAndDebug": {
379+
input: map[string]interface{}{
380+
"agent": map[string]interface{}{
381+
"debug": true,
382+
},
383+
"logs": map[string]interface{}{
384+
"logs_collected": map[string]interface{}{
385+
"application_signals": map[string]interface{}{},
386+
},
387+
},
388+
},
389+
want: &want{
390+
receivers: []string{"otlp/grpc_0_0_0_0_4315", "otlp/http_0_0_0_0_4316"},
391+
processors: []string{},
392+
exporters: []string{"debug/application_signals", "otlphttp/application_signals"},
393+
extensions: []string{"sigv4auth", "agenthealth/logs", "agenthealth/statuscode"},
394+
},
395+
isEKSCache: eksdetector.TestIsEKSCacheEKS,
396+
},
397+
}
398+
for name, testCase := range testCases {
399+
t.Run(name, func(t *testing.T) {
400+
eksdetector.IsEKS = testCase.isEKSCache
401+
conf := confmap.NewFromStringMap(testCase.input)
402+
got, err := tt.Translate(conf)
403+
assert.Equal(t, testCase.wantErr, err)
404+
if testCase.want == nil {
405+
assert.Nil(t, got)
406+
} else {
407+
require.NotNil(t, got)
408+
assert.Equal(t, testCase.want.receivers, collections.MapSlice(got.Receivers.Keys(), component.ID.String))
409+
assert.Equal(t, testCase.want.processors, collections.MapSlice(got.Processors.Keys(), component.ID.String))
410+
assert.Equal(t, testCase.want.exporters, collections.MapSlice(got.Exporters.Keys(), component.ID.String))
411+
assert.Equal(t, testCase.want.extensions, collections.MapSlice(got.Extensions.Keys(), component.ID.String))
412+
}
413+
})
414+
}
415+
}

translator/translate/otel/translate_otel.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ func Translate(jsonConfig interface{}, os string) (*otelcol.Config, error) {
7373
translators.Merge(containerInsightsTranslators)
7474
translators.Set(applicationsignals.NewTranslator(pipeline.SignalTraces))
7575
translators.Set(applicationsignals.NewTranslator(pipeline.SignalMetrics))
76+
translators.Set(applicationsignals.NewTranslator(pipeline.SignalLogs))
7677
translators.Merge(prometheus.NewTranslators(conf))
7778
translators.Set(emf_logs.NewTranslator())
7879
translators.Set(xray.NewTranslator())

0 commit comments

Comments
 (0)