Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -481,7 +481,7 @@ export const validatePackagePolicy = (
if (input.streams.length) {
input.streams.forEach((stream) => {
const streamValidationResults: PackagePolicyConfigValidationResults = {};
const streamKey = `${stream.data_stream.dataset}-${input.type}`;
const streamKey = `${stream.data_stream.dataset}-${getInputEffectiveName(input)}`;

const streamVarDefs = streamVarDefsByDatasetAndInput[streamKey];
const streamVarGroups = streamVarGroupsByDatasetAndInput[streamKey];
Expand Down
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import {

import {
getStreamsForInputType,
getInputEffectiveName,
packageToPackagePolicy,
} from '../../../../common/services/package_to_package_policy';
import { _compilePackagePolicyInputs } from '../../package_policy';
Expand Down Expand Up @@ -292,10 +293,11 @@ function buildIndexedPackage(packageInfo: PackageInfo): PackageWithInputAndStrea
const inputs = getNormalizedInputs(policyTemplate);

inputs.forEach((packageInput) => {
const inputId = `${policyTemplate.name}-${packageInput.type}`;
const inputEffectiveName = getInputEffectiveName(packageInput);
const inputId = `${policyTemplate.name}-${inputEffectiveName}`;

const streams = getStreamsForInputType(
packageInput.type,
inputEffectiveName,
packageInfo,
isIntegrationPolicyTemplate(policyTemplate) && policyTemplate.data_streams
? policyTemplate.data_streams
Expand All @@ -308,7 +310,7 @@ function buildIndexedPackage(packageInfo: PackageInfo): PackageWithInputAndStrea
}
>
>((acc, stream) => {
const streamId = `${packageInput.type}-${stream.data_stream.dataset}`;
const streamId = `${inputEffectiveName}-${stream.data_stream.dataset}`;
acc[streamId] = {
...stream,
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2143,6 +2143,55 @@ describe('Package policy service', () => {
]);
});

it('should compile stream templates when stream.input references a named input (name ?? type)', async () => {
const inputs = await _compilePackagePolicyInputs(
{
name: 'test',
version: '1.0.0',
data_streams: [
{
type: 'logs',
dataset: 'package.dataset1',
streams: [{ input: 'named_logfile', template_path: 'some_template_path.yml' }],
path: 'dataset1',
},
],
policy_templates: [
{
inputs: [{ type: 'log', name: 'named_logfile' }],
},
],
} as unknown as PackageInfo,
{},
[
{
type: 'log',
name: 'named_logfile',
enabled: true,
streams: [
{
id: 'datastream01',
data_stream: { dataset: 'package.dataset1', type: 'logs' },
enabled: true,
vars: {
paths: {
value: ['/var/log/set.log'],
},
},
},
],
},
],
ASSETS_MAP_FIXTURES
);

expect(inputs[0].streams[0].compiled_stream).toEqual({
metricset: ['dataset1'],
paths: ['/var/log/set.log'],
type: 'log',
});
});

it('should compile integration stream data_stream.dataset from package stream var default value', async () => {
const pkgInfo = {
name: 'dsvarpkg',
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3808,12 +3808,13 @@ function _compilePackageStream(

stream = _applyIndexPrivileges(packageDataStream, streamIn);

const effectiveInputName = getInputEffectiveName(input);
const streamFromPkg = (packageDataStream.streams || []).find(
(pkgStream) => pkgStream.input === input.type
(pkgStream) => pkgStream.input === effectiveInputName
);
if (!streamFromPkg) {
throw new StreamNotFoundError(
`Stream template not found, unable to find stream for input ${input.type}`
`Stream template not found, unable to find stream for input ${effectiveInputName} (type: ${input.type})`
);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import type { NewPackagePolicy, RegistryStream, UpdatePackagePolicy } from '../.
import { SO_SEARCH_LIMIT } from '../../../common';
import {
doesPackageHaveIntegrations,
getInputEffectiveName,
getNormalizedDataStreams,
getNormalizedInputs,
} from '../../../common/services';
Expand Down Expand Up @@ -391,7 +392,9 @@ function _getInputSecretPaths(
if (input.streams.length) {
input.streams.forEach((stream, streamIndex) => {
const streamVarDefs =
streamSecretVarDefsByDatasetAndInput[`${stream.data_stream.dataset}-${input.type}`];
streamSecretVarDefsByDatasetAndInput[
`${stream.data_stream.dataset}-${getInputEffectiveName(input)}`
];
if (streamVarDefs && Object.keys(streamVarDefs).length) {
Object.entries(stream.vars || {}).forEach(([name, configEntry]) => {
if (streamVarDefs[name]) {
Expand Down
Loading