From 3e59939adc69b9ab3a71aa62ea3950dc93fc0514 Mon Sep 17 00:00:00 2001 From: Julia Bardi Date: Fri, 6 Mar 2026 12:32:54 +0100 Subject: [PATCH 1/4] add use_apm if dynamic_signal_types are enabled --- .../common/services/policy_template.test.ts | 25 +++++++++++++++++++ .../fleet/common/services/policy_template.ts | 4 ++- .../agent_policies/otel_collector.test.ts | 23 +++++++++++++++++ .../services/agent_policies/otel_collector.ts | 5 ++-- 4 files changed, 54 insertions(+), 3 deletions(-) diff --git a/x-pack/platform/plugins/shared/fleet/common/services/policy_template.test.ts b/x-pack/platform/plugins/shared/fleet/common/services/policy_template.test.ts index acf0e862c8399..0f1e434c2614e 100644 --- a/x-pack/platform/plugins/shared/fleet/common/services/policy_template.test.ts +++ b/x-pack/platform/plugins/shared/fleet/common/services/policy_template.test.ts @@ -379,6 +379,31 @@ describe('getNormalizedDataStreams', () => { expect(useApmVar?.default).toEqual(true); expect(useApmVar?.title).toEqual('Enable Elastic APM Enrichment'); }); + + it('should add use_apm var when otel input has dynamic_signal_types true', () => { + const result = getNormalizedDataStreams({ + ...integrationPkg, + type: 'input', + policy_templates: [ + { + input: 'otelcol', + type: 'logs', + name: 'otel-dynamic', + template_path: 'some/path.hbl', + title: 'OTel Dynamic', + description: 'OTel with dynamic signal types', + dynamic_signal_types: true, + vars: [], + }, + ], + }); + expect(result).toHaveLength(1); + expect(result[0].streams).toHaveLength(1); + const vars = result[0].streams![0].vars; + const useApmVar = vars?.find((v) => v.name === 'use_apm'); + expect(useApmVar).toBeDefined(); + expect(useApmVar?.default).toEqual(true); + }); }); describe('filterPolicyTemplatesTiles', () => { diff --git a/x-pack/platform/plugins/shared/fleet/common/services/policy_template.ts b/x-pack/platform/plugins/shared/fleet/common/services/policy_template.ts index 089b7cca49dc9..3d04b5df36c0b 100644 --- a/x-pack/platform/plugins/shared/fleet/common/services/policy_template.ts +++ b/x-pack/platform/plugins/shared/fleet/common/services/policy_template.ts @@ -113,9 +113,11 @@ export const getNormalizedDataStreams = ( const dataset = datasetName || createDefaultDatasetName(packageInfo, policyTemplate); let vars = addDatasetVarIfNotPresent(policyTemplate.vars, policyTemplate.name); + const isOtelTraces = (dataStreamType || policyTemplate.type) === dataTypes.Traces; + const isOtelDynamicSignalTypes = policyTemplate.dynamic_signal_types === true; if ( policyTemplate.input === OTEL_COLLECTOR_INPUT_TYPE && - (dataStreamType || policyTemplate.type) === dataTypes.Traces + (isOtelTraces || isOtelDynamicSignalTypes) ) { vars = addUseAPMVarIfNotPresent(vars); } diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts index 17ae3d2ac0b40..37bddca6f8782 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts @@ -787,6 +787,29 @@ describe('generateOtelcolConfig', () => { ], ]); + it('should add elasticapm connector and processor for dynamic_signal_types input and use_apm enabled', () => { + const inputWithUseApm: FullAgentPolicyInput = { + ...otelInputWithMultipleSignalTypes, + streams: otelInputWithMultipleSignalTypes.streams?.map((stream) => ({ + ...stream, + use_apm: true, + })), + }; + const inputs: FullAgentPolicyInput[] = [inputWithUseApm]; + const result = generateOtelcolConfig(inputs, defaultOutput, packageInfoCache); + + expect(result.connectors?.elasticapm).toEqual({}); + expect(result.processors?.elasticapm).toEqual({}); + expect(result.service?.pipelines?.['metrics/aggregated-otel-metrics']).toEqual({ + receivers: ['elasticapm'], + exporters: ['forward'], + }); + const tracesPipeline = + result.service?.pipelines?.['traces/otlp/test-multi-signal-stream-id-1']; + expect(tracesPipeline?.exporters).toContain('elasticapm'); + expect(tracesPipeline?.processors).toContain('elasticapm'); + }); + it('should generate transform with multiple signal type statements when dynamic_signal_types is true', () => { const inputs: FullAgentPolicyInput[] = [otelInputWithMultipleSignalTypes]; const result = generateOtelcolConfig(inputs, defaultOutput, packageInfoCache); diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts index c386dd8d25945..dd88926a8a5ce 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts @@ -62,7 +62,8 @@ export function generateOtelcolConfig( ); const shouldAddAPMConfig = - stream.data_stream.type === dataTypes.Traces && stream[USE_APM_VAR_NAME] === true; + (stream.data_stream.type === dataTypes.Traces || hasDynamicSignalTypes(packageInfo)) && + stream[USE_APM_VAR_NAME] === true; let otelConfig: OTelCollectorConfig = { ...addSuffixToOtelcolComponentsConfig('extensions', suffix, stream?.extensions), @@ -89,7 +90,7 @@ export function generateOtelcolConfig( otelConfig = appendOtelComponents(otelConfig, 'processors', [attributesTransform]); - if (stream.data_stream.type === dataTypes.Traces && stream[USE_APM_VAR_NAME] === true) { + if (shouldAddAPMConfig) { if (!otelConfig?.connectors) { otelConfig.connectors = {}; } From 647329232abced860adfe3396ac2234c39f5d049 Mon Sep 17 00:00:00 2001 From: Julia Bardi Date: Fri, 6 Mar 2026 12:43:02 +0100 Subject: [PATCH 2/4] revert condition change in otel_collector --- .../agent_policies/otel_collector.test.ts | 23 ------------------- .../services/agent_policies/otel_collector.ts | 3 +-- 2 files changed, 1 insertion(+), 25 deletions(-) diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts index 37bddca6f8782..17ae3d2ac0b40 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts @@ -787,29 +787,6 @@ describe('generateOtelcolConfig', () => { ], ]); - it('should add elasticapm connector and processor for dynamic_signal_types input and use_apm enabled', () => { - const inputWithUseApm: FullAgentPolicyInput = { - ...otelInputWithMultipleSignalTypes, - streams: otelInputWithMultipleSignalTypes.streams?.map((stream) => ({ - ...stream, - use_apm: true, - })), - }; - const inputs: FullAgentPolicyInput[] = [inputWithUseApm]; - const result = generateOtelcolConfig(inputs, defaultOutput, packageInfoCache); - - expect(result.connectors?.elasticapm).toEqual({}); - expect(result.processors?.elasticapm).toEqual({}); - expect(result.service?.pipelines?.['metrics/aggregated-otel-metrics']).toEqual({ - receivers: ['elasticapm'], - exporters: ['forward'], - }); - const tracesPipeline = - result.service?.pipelines?.['traces/otlp/test-multi-signal-stream-id-1']; - expect(tracesPipeline?.exporters).toContain('elasticapm'); - expect(tracesPipeline?.processors).toContain('elasticapm'); - }); - it('should generate transform with multiple signal type statements when dynamic_signal_types is true', () => { const inputs: FullAgentPolicyInput[] = [otelInputWithMultipleSignalTypes]; const result = generateOtelcolConfig(inputs, defaultOutput, packageInfoCache); diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts index dd88926a8a5ce..2765d0d4ba7fb 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts @@ -62,8 +62,7 @@ export function generateOtelcolConfig( ); const shouldAddAPMConfig = - (stream.data_stream.type === dataTypes.Traces || hasDynamicSignalTypes(packageInfo)) && - stream[USE_APM_VAR_NAME] === true; + stream.data_stream.type === dataTypes.Traces && stream[USE_APM_VAR_NAME] === true; let otelConfig: OTelCollectorConfig = { ...addSuffixToOtelcolComponentsConfig('extensions', suffix, stream?.extensions), From 1654130f82448d6d77f85a042ced2aa5d5070c88 Mon Sep 17 00:00:00 2001 From: Julia Bardi Date: Tue, 10 Mar 2026 11:12:20 +0100 Subject: [PATCH 3/4] add APM connector when there is traces pipeline --- .../agent_policies/otel_collector.test.ts | 25 +++++++++++++++++++ .../services/agent_policies/otel_collector.ts | 18 +++++++------ 2 files changed, 36 insertions(+), 7 deletions(-) diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts index 17ae3d2ac0b40..10c810a6d9cd6 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts @@ -787,6 +787,31 @@ describe('generateOtelcolConfig', () => { ], ]); + it('should add elasticapm connector and processor when stream has traces pipeline and use_apm enabled even if data_stream.type is not traces', () => { + const inputWithUseApm: FullAgentPolicyInput = { + ...otelInputWithMultipleSignalTypes, + streams: + otelInputWithMultipleSignalTypes.streams?.map((stream) => ({ + ...stream, + use_apm: true, + })) ?? [], + }; + const inputs: FullAgentPolicyInput[] = [inputWithUseApm]; + const result = generateOtelcolConfig(inputs, defaultOutput, packageInfoCache); + + expect(result.connectors?.elasticapm).toEqual({}); + expect(result.processors?.elasticapm).toEqual({}); + expect(result.service?.pipelines?.['metrics/aggregated-otel-metrics']).toEqual({ + receivers: ['elasticapm'], + exporters: ['forward'], + }); + const tracesPipelineKey = 'traces/otlp/test-multi-signal-stream-id-1'; + const tracesPipeline = result.service?.pipelines?.[tracesPipelineKey]; + expect(tracesPipeline).toBeDefined(); + expect(tracesPipeline?.exporters).toContain('elasticapm'); + expect(tracesPipeline?.processors).toContain('elasticapm'); + }); + it('should generate transform with multiple signal type statements when dynamic_signal_types is true', () => { const inputs: FullAgentPolicyInput[] = [otelInputWithMultipleSignalTypes]; const result = generateOtelcolConfig(inputs, defaultOutput, packageInfoCache); diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts index 2765d0d4ba7fb..c97ddbbb82e40 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts @@ -50,6 +50,10 @@ export function generateOtelcolConfig( const otelInputs: OTelCollectorConfig[] = (input?.streams ?? []).map((stream) => { // Avoid dots in keys, as they can create subobjects in agent config. const suffix = (input.id + '-' + stream.id).replaceAll('.', '-'); + // Extract signal types from pipeline IDs + const signalTypes = stream.service?.pipelines + ? extractSignalTypesFromPipelines(stream.service?.pipelines) + : []; const attributesTransform = generateOTelAttributesTransform( stream.data_stream.type ? stream.data_stream.type : 'logs', stream.data_stream.dataset, @@ -58,11 +62,13 @@ export function generateOtelcolConfig( : 'default', suffix, packageInfo, - stream.service?.pipelines + signalTypes ); + const hasTracesPipeline = signalTypes.includes('traces'); const shouldAddAPMConfig = - stream.data_stream.type === dataTypes.Traces && stream[USE_APM_VAR_NAME] === true; + (stream.data_stream.type === dataTypes.Traces || hasTracesPipeline) && + stream[USE_APM_VAR_NAME] === true; let otelConfig: OTelCollectorConfig = { ...addSuffixToOtelcolComponentsConfig('extensions', suffix, stream?.extensions), @@ -208,17 +214,15 @@ function generateOTelAttributesTransform( namespace: string, suffix: string, packageInfo?: PackageInfo, - streamPipelines?: Record + signalTypes?: string[] ): Record { const dynamicSignalTypes = hasDynamicSignalTypes(packageInfo); let transformStatements: Record = {}; - if (dynamicSignalTypes && streamPipelines) { - // When dynamic_signal_types is true, extract signal types from pipeline IDs - // and generate transforms for each. This allows the collector to route data + if (dynamicSignalTypes && signalTypes) { + // When dynamic_signal_types is true, generate transforms for each signal type. This allows the collector to route data // to the appropriate datastreams based on the pipelines configured in the policy. - const signalTypes = extractSignalTypesFromPipelines(streamPipelines); // Generate transforms for each signal type found in pipelines signalTypes.forEach((signalType) => { const typeTransforms = generateOtelTypeTransforms(signalType, dataset, namespace); From 5d95363628546fb4dfc9c547ab2eade4a2e0cf7c Mon Sep 17 00:00:00 2001 From: Julia Bardi Date: Tue, 10 Mar 2026 11:43:25 +0100 Subject: [PATCH 4/4] add elasticapm to only traces pipelines --- .../agent_policies/otel_collector.test.ts | 5 ++++ .../services/agent_policies/otel_collector.ts | 28 ++++++++----------- 2 files changed, 16 insertions(+), 17 deletions(-) diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts index 10c810a6d9cd6..2c79a51aab5c0 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.test.ts @@ -810,6 +810,11 @@ describe('generateOtelcolConfig', () => { expect(tracesPipeline).toBeDefined(); expect(tracesPipeline?.exporters).toContain('elasticapm'); expect(tracesPipeline?.processors).toContain('elasticapm'); + const metricsPipelineKey = 'metrics/otlp/test-multi-signal-stream-id-1'; + const metricsPipeline = result.service?.pipelines?.[metricsPipelineKey]; + expect(metricsPipeline).toBeDefined(); + expect(metricsPipeline?.exporters).not.toContain('elasticapm'); + expect(metricsPipeline?.processors).not.toContain('elasticapm'); }); it('should generate transform with multiple signal type statements when dynamic_signal_types is true', () => { diff --git a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts index c97ddbbb82e40..f0b6bda225182 100644 --- a/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts +++ b/x-pack/platform/plugins/shared/fleet/server/services/agent_policies/otel_collector.ts @@ -285,24 +285,18 @@ function conditionallyAddApmToPipelines( if (!shouldAddAPMConfig) { return pipelines; } - pipelines = addCompomentToPipelines(pipelines, 'elasticapm', 'exporters'); - pipelines = addCompomentToPipelines(pipelines, 'elasticapm', 'processors'); - return pipelines; -} - -function addCompomentToPipelines( - pipelines: any, - componentId: string, - type: string -): Record { - for (const pipelineId in pipelines) { - if (pipelines[pipelineId][type]) { - pipelines[pipelineId][type] = pipelines[pipelineId][type].concat([componentId]); - } else { - pipelines[pipelineId][type] = [componentId]; + const result: Record = {}; + Object.entries(pipelines as Record>).forEach( + ([pipelineID, pipeline]) => { + const signalType = getSignalType(pipelineID); + if (signalType === 'traces') { + pipeline.exporters = [...(pipeline.exporters || []), 'elasticapm']; + pipeline.processors = [...(pipeline.processors || []), 'elasticapm']; + } + result[pipelineID] = pipeline; } - } - return pipelines; + ); + return result; } function addSuffixToOtelcolPipelinesComponents(