From 891db6cbe2894c8f55b39adbb751bb9d5b94c017 Mon Sep 17 00:00:00 2001 From: pablu Date: Thu, 16 Jul 2026 20:25:45 +0200 Subject: [PATCH 1/9] wip: generic adder dynamic, add dotted_template, WITH pylint errors --- logprep/abc/getter.py | 10 +- logprep/processor/generic_adder/processor.py | 8 +- logprep/processor/generic_adder/rule.py | 179 +++++++++++++++++-- 3 files changed, 184 insertions(+), 13 deletions(-) diff --git a/logprep/abc/getter.py b/logprep/abc/getter.py index 70095f30c..37ca827a1 100644 --- a/logprep/abc/getter.py +++ b/logprep/abc/getter.py @@ -114,10 +114,18 @@ def get_json(self) -> dict | list: content = self._resolve_content(raw) return self._parse_json(content) - def get_collection(self) -> dict | list: + def get_collection(self, content_field: str | None = None) -> dict | list: """Gets and parses the raw content to yaml or json""" content = self._resolve_content_by_content_type() + if isinstance(content, dict) and content_field is not None: + content = content[content_field] + elif content_field is not None: + raise ValueError( + "Expected mapping type, like json object when content_field is set, got %s.", + type(content), + ) + if isinstance(content, str): content = self._parse_yaml_or_json(content) diff --git a/logprep/processor/generic_adder/processor.py b/logprep/processor/generic_adder/processor.py index ae84b7ee1..e9889678d 100644 --- a/logprep/processor/generic_adder/processor.py +++ b/logprep/processor/generic_adder/processor.py @@ -37,8 +37,14 @@ class GenericAdder(Processor): rule_class = GenericAdderRule + def setup(self): + super().setup() + for rule in self.rules: + rule = typing.cast(GenericAdderRule, rule) + rule.set_job_tag(self._job_tag_for_cleanup) + def _apply_rules(self, event: dict, rule: Rule): rule = typing.cast(GenericAdderRule, rule) - items_to_add = rule.add + items_to_add = rule.add(event) if items_to_add: add_fields_to(event, items_to_add, rule, rule.merge_with_target, rule.overwrite_target) diff --git a/logprep/processor/generic_adder/rule.py b/logprep/processor/generic_adder/rule.py index e82e209b8..40298b92f 100644 --- a/logprep/processor/generic_adder/rule.py +++ b/logprep/processor/generic_adder/rule.py @@ -87,9 +87,47 @@ from attrs import define, field, validators +from logprep.filter.expression.filter_expression import FilterExpression from logprep.processor.base.rule import InvalidRuleDefinitionError from logprep.processor.field_manager.rule import FieldManagerRule -from logprep.util.getter import GetterFactory, RefreshableGetter +from logprep.util.converters import convert_from_dict +from logprep.util.dotted_template import DottedTemplate +from logprep.util.getter import GetterFactory, HttpGetter, RefreshableGetter +from logprep.util.helper import FieldValue, get_dotted_field_value + + +@define(kw_only=True, frozen=True) +class AddFromUrlConfig: + url: str = field( + validator=[validators.instance_of(str), validators.matches_re(r"^https?://.+")] + ) + + target_field: str | None = field( + default=None, validator=validators.optional(validators.instance_of(str)) + ) + + target_field_mapping: dict[str, str] | None = field( + validator=validators.optional( + validators.deep_mapping( + key_validator=validators.instance_of(str), + value_validator=validators.instance_of(str), + ) + ), + default=None, + ) + + def __attrs_post_init__(self) -> None: + if not self.target_field and not self.target_field_mapping: + raise ValueError("add_from_url requires target_field or target_field_mapping") + + +def _convert_add_from_url( + value: AddFromUrlConfig | dict | None, +) -> AddFromUrlConfig | None: + if value is None: + return None + + return convert_from_dict(AddFromUrlConfig, value) class GenericAdderRule(FieldManagerRule): @@ -113,13 +151,15 @@ class Config(FieldManagerRule.Config): """Contains a dictionary of field names and values that should be added. If dot notation is being used, then all fields on the path are being automatically created.""" - add_from_file: list = field( - validator=[ - validators.instance_of(list), - validators.deep_iterable(member_validator=validators.instance_of(str)), - ], - converter=lambda x: x if isinstance(x, list) else [x], - factory=list, + add_from_file: list | None = field( + validator=validators.optional( + validators.deep_iterable( + iterable_validator=validators.instance_of(list), + member_validator=validators.instance_of(str), + ) + ), + converter=lambda x: x if x is None or isinstance(x, list) else [x], + default=None, eq=False, ) """Contains the path or url to YML file that contains a dictionary of field names @@ -141,6 +181,19 @@ class Config(FieldManagerRule.Config): authenticity and integrity of the loaded values. """ + + add_from_url: AddFromUrlConfig | None = field( + default=None, + validator=validators.optional( + validators.instance_of(AddFromUrlConfig) + # validators.deep_iterable( + # iterable_validator=validators.instance_of(list), + # member_validator=validators.instance_of(AddFromUrlConfig), + # ) + ), + converter=_convert_add_from_url, + ) + only_first_existing_file: bool = field( validator=validators.instance_of(bool), default=False, eq=False ) @@ -160,6 +213,37 @@ def _refresh_add(self): def __attrs_post_init__(self): self._base_add = copy.deepcopy(self.add) + if ( + self.add_from_file is not None or self.add is not None + ) and self.add_from_url is not None: + raise ValueError( + "only one of add_from_file + add or add_from_url is allowed per GenericAdder rule" + ) + + if not self.add and self.add_from_file is None and self.add_from_url is None: + raise ValueError( + "one of add, add_from_file or add_from_url is neccessary per GenericAdder rule" + ) + + if ( + self.add_from_url is not None + and self.add_from_url.target_field is not None + and self.add_from_url.target_field_mapping is not None + ): + raise ValueError( + "only one of target_field or target_field_mapping is allowed per GenericAdder rule" + ) + + if ( + self.add_from_url is not None + and self.add_from_url.target_field is None + and self.add_from_url.target_field_mapping is None + ): + raise ValueError( + "one of target_field or target_field_mapping is neccessary per GenericAdder rule" + ) + + # Eagerly loaded from file if self.add_from_file: for add_file in self.add_from_file: # pylint: disable=not-an-iterable getter = GetterFactory.from_string(add_file) @@ -173,6 +257,7 @@ def __attrs_post_init__(self): def _add_from_path(self): """Reads add fields from file""" missing_files = [] + assert self.add_from_file is not None for add_file in self.add_from_file: # pylint: disable=not-an-iterable try: add_dict = GetterFactory.from_string(add_file).get_yaml() @@ -195,8 +280,80 @@ def _add_from_path(self): f"The following required files do not exist: '{missing_files}'" ) - @property - def add(self) -> dict: + def __init__(self, filter_rule: FilterExpression, config: Config, processor_name: str): + super().__init__(filter_rule, config, processor_name) + self._dynamic_content: dict[str, FieldValue] = {} + self._callback_tag = "" + + def set_job_tag(self, job_tag: str) -> None: + self._callback_tag = job_tag + + def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: + config = typing.cast(GenericAdderRule.Config, self._config) + + assert config.add_from_url + + dynamic_path_template = DottedTemplate(config.add_from_url.url) + key_val = { + identifier: get_dotted_field_value(event, identifier) + for identifier in dynamic_path_template.get_identifiers() + } + for identifier, val in key_val.items(): + if val is None: + raise ValueError( + f"missing event field {identifier!r} for dynamic list comparison path" + ) + if not isinstance(val, (str, int)): + raise ValueError( + f"value for list comparison field {identifier!r} is not a scalar value" + ) + pass + + dynamic_resolved = dynamic_path_template.substitute(key_val) + content: FieldValue | None = None + if dynamic_resolved not in self._dynamic_content: + http_getter = GetterFactory.from_string(dynamic_resolved) + assert isinstance(http_getter, HttpGetter) + + http_getter.keep_alive() + + content = http_getter.get_collection() + self._dynamic_content[dynamic_resolved] = content + + # Get tag here + # Add update callback + # Add cleanup callback + else: + RefreshableGetter.keep_alive_for_target(dynamic_resolved) + content = self._dynamic_content[dynamic_resolved] + + items_to_add: dict[str, FieldValue] = {} + + if config.add_from_url.target_field: + items_to_add[config.add_from_url.target_field] = content + else: + assert config.add_from_url.target_field_mapping is not None + + # TODO: Check this differently, what should this be? + assert isinstance(content, dict) + for ( + mapping_source_field, + mapping_target_field, + ) in config.add_from_url.target_field_mapping.items(): + # what should happen with missing fields? Right now None is skipped + # change this to get_dotted_field_value_with_missing + items_to_add[mapping_target_field] = get_dotted_field_value( + content, mapping_source_field + ) + + return items_to_add + + def add(self, event: dict) -> dict: """Returns the fields to add""" config = typing.cast(GenericAdderRule.Config, self._config) - return config.add + + if config.add_from_file is not None: + return config.add + + assert config.add_from_url is not None + return self._dynamic_add_from_url(event) From 1ab8eeb2b205327db04bc9736ea718df8cfdf0ba Mon Sep 17 00:00:00 2001 From: pablu Date: Tue, 21 Jul 2026 12:50:22 +0200 Subject: [PATCH 2/9] update to be closer to list_comparison, also update tests to pass correctly --- .../ng/processor/generic_adder/processor.py | 2 +- logprep/processor/generic_adder/processor.py | 2 +- logprep/processor/generic_adder/rule.py | 81 ++++++++++++------- .../generic_adder/test_generic_adder.py | 6 +- .../generic_adder/test_generic_adder_rule.py | 12 +-- 5 files changed, 65 insertions(+), 38 deletions(-) diff --git a/logprep/ng/processor/generic_adder/processor.py b/logprep/ng/processor/generic_adder/processor.py index 12f930595..6dc324d0a 100644 --- a/logprep/ng/processor/generic_adder/processor.py +++ b/logprep/ng/processor/generic_adder/processor.py @@ -39,6 +39,6 @@ class GenericAdder(Processor): def _apply_rules(self, event: dict, rule: Rule) -> None: rule = typing.cast(GenericAdderRule, rule) - items_to_add = rule.add + items_to_add = rule.add(event) if items_to_add: add_fields_to(event, items_to_add, rule, rule.merge_with_target, rule.overwrite_target) diff --git a/logprep/processor/generic_adder/processor.py b/logprep/processor/generic_adder/processor.py index e9889678d..ca70b7e50 100644 --- a/logprep/processor/generic_adder/processor.py +++ b/logprep/processor/generic_adder/processor.py @@ -41,7 +41,7 @@ def setup(self): super().setup() for rule in self.rules: rule = typing.cast(GenericAdderRule, rule) - rule.set_job_tag(self._job_tag_for_cleanup) + rule.init_generic_adder(self._job_tag_for_cleanup) def _apply_rules(self, event: dict, rule: Rule): rule = typing.cast(GenericAdderRule, rule) diff --git a/logprep/processor/generic_adder/rule.py b/logprep/processor/generic_adder/rule.py index 40298b92f..fed6ccc66 100644 --- a/logprep/processor/generic_adder/rule.py +++ b/logprep/processor/generic_adder/rule.py @@ -83,6 +83,7 @@ # pylint: enable=anomalous-backslash-in-string import copy +import os import typing from attrs import define, field, validators @@ -91,9 +92,8 @@ from logprep.processor.base.rule import InvalidRuleDefinitionError from logprep.processor.field_manager.rule import FieldManagerRule from logprep.util.converters import convert_from_dict -from logprep.util.dotted_template import DottedTemplate from logprep.util.getter import GetterFactory, HttpGetter, RefreshableGetter -from logprep.util.helper import FieldValue, get_dotted_field_value +from logprep.util.helper import DottedTemplate, FieldValue, get_dotted_field_value @define(kw_only=True, frozen=True) @@ -106,14 +106,12 @@ class AddFromUrlConfig: default=None, validator=validators.optional(validators.instance_of(str)) ) - target_field_mapping: dict[str, str] | None = field( - validator=validators.optional( - validators.deep_mapping( - key_validator=validators.instance_of(str), - value_validator=validators.instance_of(str), - ) + target_field_mapping: dict[str, str] = field( + validator=validators.deep_mapping( + key_validator=validators.instance_of(str), + value_validator=validators.instance_of(str), ), - default=None, + factory=dict, ) def __attrs_post_init__(self) -> None: @@ -184,13 +182,7 @@ class Config(FieldManagerRule.Config): add_from_url: AddFromUrlConfig | None = field( default=None, - validator=validators.optional( - validators.instance_of(AddFromUrlConfig) - # validators.deep_iterable( - # iterable_validator=validators.instance_of(list), - # member_validator=validators.instance_of(AddFromUrlConfig), - # ) - ), + validator=validators.optional(validators.instance_of(AddFromUrlConfig)), converter=_convert_add_from_url, ) @@ -214,7 +206,7 @@ def __attrs_post_init__(self): self._base_add = copy.deepcopy(self.add) if ( - self.add_from_file is not None or self.add is not None + self.add_from_file is not None or len(self.add) > 0 ) and self.add_from_url is not None: raise ValueError( "only one of add_from_file + add or add_from_url is allowed per GenericAdder rule" @@ -228,7 +220,7 @@ def __attrs_post_init__(self): if ( self.add_from_url is not None and self.add_from_url.target_field is not None - and self.add_from_url.target_field_mapping is not None + and len(self.add_from_url.target_field_mapping) > 0 ): raise ValueError( "only one of target_field or target_field_mapping is allowed per GenericAdder rule" @@ -237,7 +229,7 @@ def __attrs_post_init__(self): if ( self.add_from_url is not None and self.add_from_url.target_field is None - and self.add_from_url.target_field_mapping is None + and len(self.add_from_url.target_field_mapping) <= 0 ): raise ValueError( "one of target_field or target_field_mapping is neccessary per GenericAdder rule" @@ -284,19 +276,35 @@ def __init__(self, filter_rule: FilterExpression, config: Config, processor_name super().__init__(filter_rule, config, processor_name) self._dynamic_content: dict[str, FieldValue] = {} self._callback_tag = "" + self._is_dynamic: bool = False + self._dynamic_template: DottedTemplate + self._dynamic_identifiers: tuple[str, ...] = () - def set_job_tag(self, job_tag: str) -> None: + def init_generic_adder(self, job_tag: str) -> None: self._callback_tag = job_tag + config = typing.cast(GenericAdderRule.Config, self._config) + if config.add_from_file or config.add: + return + + assert config.add_from_url is not None + + base_template = DottedTemplate(config.add_from_url.url) + resolved_template = DottedTemplate(base_template.safe_substitute({**os.environ})) + self._dynamic_template = resolved_template + self._dynamic_identifiers = tuple(resolved_template.get_identifiers()) + + if len(self._dynamic_identifiers) > 0: + self._is_dynamic = True + def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: config = typing.cast(GenericAdderRule.Config, self._config) assert config.add_from_url - dynamic_path_template = DottedTemplate(config.add_from_url.url) key_val = { identifier: get_dotted_field_value(event, identifier) - for identifier in dynamic_path_template.get_identifiers() + for identifier in self._dynamic_identifiers } for identifier, val in key_val.items(): if val is None: @@ -309,7 +317,7 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: ) pass - dynamic_resolved = dynamic_path_template.substitute(key_val) + dynamic_resolved = self._dynamic_template.substitute(key_val) content: FieldValue | None = None if dynamic_resolved not in self._dynamic_content: http_getter = GetterFactory.from_string(dynamic_resolved) @@ -320,9 +328,21 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: content = http_getter.get_collection() self._dynamic_content[dynamic_resolved] = content - # Get tag here - # Add update callback - # Add cleanup callback + tag = self._callback_tag + + http_getter.add_callback( + tag, + self._update_dynamic_content, + deduplication_key=(tag, dynamic_resolved, id(self)), + fnc_args=[http_getter, dynamic_resolved], + ) + + http_getter.add_cleanup_callback( + tag, + self._cleanup, + deduplication_key=(tag, dynamic_resolved, id(self)), + fnc_args=[dynamic_resolved], + ) else: RefreshableGetter.keep_alive_for_target(dynamic_resolved) content = self._dynamic_content[dynamic_resolved] @@ -348,11 +368,18 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: return items_to_add + def _update_dynamic_content(self, http_getter: HttpGetter, resolved_uri: str): + content = http_getter.get_collection() + self._dynamic_content[resolved_uri] = content + + def _cleanup(self, resolved_uri: str): + self._dynamic_content.pop(resolved_uri, None) + def add(self, event: dict) -> dict: """Returns the fields to add""" config = typing.cast(GenericAdderRule.Config, self._config) - if config.add_from_file is not None: + if config.add_from_file is not None or len(config.add) > 0: return config.add assert config.add_from_url is not None diff --git a/tests/unit/processor/generic_adder/test_generic_adder.py b/tests/unit/processor/generic_adder/test_generic_adder.py index 6ef8fcd99..af5ace20b 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/processor/generic_adder/test_generic_adder.py @@ -303,8 +303,8 @@ "comp\\lex.field": "value", "comp\\lex.nested": {"field": 42}, "nested": {"comp\\lex.field": 1337}, - "\\u\\0\\1\\x\y": 1338, # pylint: disable=anomalous-backslash-in-string - "\\u\\0\\1\\x\z": "whatever", # pylint: disable=anomalous-backslash-in-string + "\\u\\0\\1\\x\\y": 1338, # pylint: disable=anomalous-backslash-in-string + "\\u\\0\\1\\x\\z": "whatever", # pylint: disable=anomalous-backslash-in-string }, id="Add from rule definition with escaping", ), @@ -454,7 +454,7 @@ def test_add_only_copies(self): event = {} instance.process(event) - rule_add = instance.rules[0].add + rule_add = instance.rules[0].add({}) assert event["some_list_field"] == ["some_value"] assert event["some_list_field"] is not rule_add["some_list_field"], "only copies in events" diff --git a/tests/unit/processor/generic_adder/test_generic_adder_rule.py b/tests/unit/processor/generic_adder/test_generic_adder_rule.py index 4052f3196..184a6a91a 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder_rule.py +++ b/tests/unit/processor/generic_adder/test_generic_adder_rule.py @@ -101,7 +101,7 @@ def test_rule_accepts_bool_type(self): "generic_adder": {"add": {"added_bool_field": True}}, } rule = GenericAdderRule.create_from_dict(rule_definition) - assert isinstance(rule.add.get("added_bool_field"), bool) + assert isinstance(rule.add({}).get("added_bool_field"), bool) @responses.activate def test_rule_callback_updates_additions_and_preserves_original_add(self, tmp_path): @@ -146,11 +146,11 @@ def test_rule_callback_updates_additions_and_preserves_original_add(self, tmp_pa with mock_env({ENV_NAME_LOGPREP_GETTER_CONFIG: str(http_getter_conf)}): scheduler = HttpGetter(protocol="http", target=url).scheduler rule = GenericAdderRule.create_from_dict(rule_definition) - assert rule.add == expected_1 + assert rule.add({}) == expected_1 HttpGetter.refresh() - assert rule.add == expected_1 + assert rule.add({}) == expected_1 scheduler.run_all() - assert rule.add == expected_2 - assert rule.add == expected_2 + assert rule.add({}) == expected_2 + assert rule.add({}) == expected_2 scheduler.run_all() - assert rule.add == expected_3 + assert rule.add({}) == expected_3 From 35228b49b93fa05b036f0413c5b00e7bd2443111 Mon Sep 17 00:00:00 2001 From: pablu Date: Thu, 23 Jul 2026 14:29:11 +0200 Subject: [PATCH 3/9] fix some review comments --- logprep/abc/getter.py | 26 +++++++------- logprep/processor/generic_adder/processor.py | 18 +++++++--- logprep/processor/generic_adder/rule.py | 36 +++++++++----------- 3 files changed, 43 insertions(+), 37 deletions(-) diff --git a/logprep/abc/getter.py b/logprep/abc/getter.py index 37ca827a1..64d2e7860 100644 --- a/logprep/abc/getter.py +++ b/logprep/abc/getter.py @@ -118,13 +118,7 @@ def get_collection(self, content_field: str | None = None) -> dict | list: """Gets and parses the raw content to yaml or json""" content = self._resolve_content_by_content_type() - if isinstance(content, dict) and content_field is not None: - content = content[content_field] - elif content_field is not None: - raise ValueError( - "Expected mapping type, like json object when content_field is set, got %s.", - type(content), - ) + content = Getter._apply_content_field(content, content_field) if isinstance(content, str): content = self._parse_yaml_or_json(content) @@ -151,11 +145,10 @@ def _parse_newline_separated_list(content: str) -> list: """Helper which tries to convert content to list""" return content.splitlines() - def get_list(self, content_field: str | None = None) -> list: - """Gets list and fails otherwise""" - - content = self._resolve_content_by_content_type() - + @staticmethod + def _apply_content_field( + content: dict | list | str, content_field: str | None = None + ) -> dict | list | str: if isinstance(content, dict) and content_field is not None: content = content[content_field] elif content_field is not None: @@ -163,6 +156,15 @@ def get_list(self, content_field: str | None = None) -> list: f"Expected mapping type when content_field is set, got {type(content)}" ) + return content + + def get_list(self, content_field: str | None = None) -> list: + """Gets list and fails otherwise""" + + content = self._resolve_content_by_content_type() + + content = Getter._apply_content_field(content, content_field) + if isinstance(content, str): content = self._parse_newline_separated_list(content) diff --git a/logprep/processor/generic_adder/processor.py b/logprep/processor/generic_adder/processor.py index ca70b7e50..5ae424bb7 100644 --- a/logprep/processor/generic_adder/processor.py +++ b/logprep/processor/generic_adder/processor.py @@ -25,10 +25,11 @@ """ import typing +from typing import Sequence from logprep.abc.processor import Processor -from logprep.processor.base.rule import Rule from logprep.processor.generic_adder.rule import GenericAdderRule +from logprep.util.getter import RefreshableGetter from logprep.util.helper import add_fields_to @@ -37,14 +38,21 @@ class GenericAdder(Processor): rule_class = GenericAdderRule + @property + def _rules(self) -> Sequence[GenericAdderRule]: + """Returns all rules""" + return typing.cast(Sequence[GenericAdderRule], self.rules) + def setup(self): super().setup() - for rule in self.rules: - rule = typing.cast(GenericAdderRule, rule) + for rule in self._rules: rule.init_generic_adder(self._job_tag_for_cleanup) - def _apply_rules(self, event: dict, rule: Rule): - rule = typing.cast(GenericAdderRule, rule) + def _apply_rules(self, event: dict, rule: GenericAdderRule): items_to_add = rule.add(event) if items_to_add: add_fields_to(event, items_to_add, rule, rule.merge_with_target, rule.overwrite_target) + + def _shut_down(self) -> None: + RefreshableGetter.remove_callbacks_for_tag(self._job_tag_for_cleanup) + return super()._shut_down() diff --git a/logprep/processor/generic_adder/rule.py b/logprep/processor/generic_adder/rule.py index fed6ccc66..209540775 100644 --- a/logprep/processor/generic_adder/rule.py +++ b/logprep/processor/generic_adder/rule.py @@ -116,7 +116,7 @@ class AddFromUrlConfig: def __attrs_post_init__(self) -> None: if not self.target_field and not self.target_field_mapping: - raise ValueError("add_from_url requires target_field or target_field_mapping") + raise ValueError("adding values from url requires target_field or target_field_mapping") def _convert_add_from_url( @@ -149,15 +149,13 @@ class Config(FieldManagerRule.Config): """Contains a dictionary of field names and values that should be added. If dot notation is being used, then all fields on the path are being automatically created.""" - add_from_file: list | None = field( - validator=validators.optional( - validators.deep_iterable( - iterable_validator=validators.instance_of(list), - member_validator=validators.instance_of(str), - ) + add_from_file: list[str] = field( + validator=validators.deep_iterable( + iterable_validator=validators.instance_of(list), + member_validator=validators.instance_of(str), ), - converter=lambda x: x if x is None or isinstance(x, list) else [x], - default=None, + converter=lambda x: x if isinstance(x, list) else [x], + factory=list, eq=False, ) """Contains the path or url to YML file that contains a dictionary of field names @@ -236,21 +234,18 @@ def __attrs_post_init__(self): ) # Eagerly loaded from file - if self.add_from_file: - for add_file in self.add_from_file: # pylint: disable=not-an-iterable - getter = GetterFactory.from_string(add_file) - if isinstance(getter, RefreshableGetter): - # TODO: This never gets cleaned up, Memory leak on a lot of new generic adders / generic resolvers - getter.add_callback( - f"generic_adder:{self.id}:{add_file}", self._refresh_add - ) - self._add_from_path() + for add_file in self.add_from_file: # pylint: disable=not-an-iterable + getter = GetterFactory.from_string(add_file) + if isinstance(getter, RefreshableGetter): + # TODO: This never gets cleaned up, Memory leak on a lot of new generic adders / generic resolvers + getter.add_callback(f"generic_adder:{self.id}:{add_file}", self._refresh_add) + self._add_from_path() def _add_from_path(self): """Reads add fields from file""" missing_files = [] - assert self.add_from_file is not None - for add_file in self.add_from_file: # pylint: disable=not-an-iterable + + for add_file in self.add_from_file: try: add_dict = GetterFactory.from_string(add_file).get_yaml() except FileNotFoundError: @@ -343,6 +338,7 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: deduplication_key=(tag, dynamic_resolved, id(self)), fnc_args=[dynamic_resolved], ) + else: RefreshableGetter.keep_alive_for_target(dynamic_resolved) content = self._dynamic_content[dynamic_resolved] From fd8faf3e44ba9ce76efcf7a9b0c4111491223e23 Mon Sep 17 00:00:00 2001 From: pablu Date: Thu, 23 Jul 2026 14:42:44 +0200 Subject: [PATCH 4/9] update stale tests, and remove unused pydoc comments --- tests/unit/ng/processor/generic_adder/test_generic_adder.py | 2 +- tests/unit/processor/generic_adder/test_generic_adder.py | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/tests/unit/ng/processor/generic_adder/test_generic_adder.py b/tests/unit/ng/processor/generic_adder/test_generic_adder.py index 8b1738f1b..e70a0f583 100644 --- a/tests/unit/ng/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/ng/processor/generic_adder/test_generic_adder.py @@ -95,7 +95,7 @@ async def test_add_only_copies(self): log_event = LogEvent(event, original=b"", input_meta=InputMeta()) await instance.process(log_event) - rule_add = instance.rules[0].add + rule_add = instance.rules[0].add({}) assert event["some_list_field"] == ["some_value"] assert event["some_list_field"] is not rule_add["some_list_field"], "only copies in events" diff --git a/tests/unit/processor/generic_adder/test_generic_adder.py b/tests/unit/processor/generic_adder/test_generic_adder.py index af5ace20b..e70a23bb2 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/processor/generic_adder/test_generic_adder.py @@ -295,7 +295,7 @@ { "add_generic_test": "Test", "event_id": 123, - "\\u\\0\\1\\x\\z": "whatever", # pylint: disable=anomalous-backslash-in-string + "\\u\\0\\1\\x\\z": "whatever", }, { "add_generic_test": "Test", @@ -303,8 +303,8 @@ "comp\\lex.field": "value", "comp\\lex.nested": {"field": 42}, "nested": {"comp\\lex.field": 1337}, - "\\u\\0\\1\\x\\y": 1338, # pylint: disable=anomalous-backslash-in-string - "\\u\\0\\1\\x\\z": "whatever", # pylint: disable=anomalous-backslash-in-string + "\\u\\0\\1\\x\\y": 1338, + "\\u\\0\\1\\x\\z": "whatever", }, id="Add from rule definition with escaping", ), From 94a01a4fc34c082a0206250cd6e0575181c5dc25 Mon Sep 17 00:00:00 2001 From: pablu Date: Fri, 24 Jul 2026 09:21:07 +0200 Subject: [PATCH 5/9] add integration / unit tests --- .../generic_adder/test_generic_adder.py | 33 +++++ .../generic_adder/test_generic_adder.py | 56 +++++++++ .../generic_adder/test_generic_adder_rule.py | 114 +++++++++++++++++- tests/unit/util/test_getter.py | 25 ++++ 4 files changed, 227 insertions(+), 1 deletion(-) diff --git a/tests/unit/ng/processor/generic_adder/test_generic_adder.py b/tests/unit/ng/processor/generic_adder/test_generic_adder.py index e70a0f583..60f33d2c1 100644 --- a/tests/unit/ng/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/ng/processor/generic_adder/test_generic_adder.py @@ -9,6 +9,7 @@ from copy import deepcopy import pytest +import responses from logprep.factory import Factory from logprep.ng.abc.event import InputMeta, LogEvent @@ -102,3 +103,35 @@ async def test_add_only_copies(self): assert event["some_dict_field"] == {"some_key": "some_value"} assert event["some_dict_field"] is not rule_add["some_dict_field"], "only copies in events" + + @responses.activate + async def test_adds_response_from_event_templated_url(self): + resolved_url = "https://values.example/acme" + response_content = {"user": {"name": "Alice"}, "risk": {"score": 7}} + responses.add(responses.GET, resolved_url, json=response_content) + processor = self._create_test_instance( + { + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant.id}", + "target_field": "enrichment", + } + }, + } + ] + } + ) + await processor.setup() + event = {"tenant": {"id": "acme"}} + + result = await processor.process(LogEvent(event, original=b"", input_meta=InputMeta())) + + assert result.errors == [] + assert event == { + "tenant": {"id": "acme"}, + "enrichment": response_content, + } + assert responses.calls[0].request.url == resolved_url diff --git a/tests/unit/processor/generic_adder/test_generic_adder.py b/tests/unit/processor/generic_adder/test_generic_adder.py index e70a23bb2..7a665f077 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/processor/generic_adder/test_generic_adder.py @@ -7,11 +7,14 @@ from copy import deepcopy import pytest +import responses from logprep.factory import Factory from logprep.processor.base.exceptions import InvalidRuleDefinitionError from logprep.processor.generic_adder.processor import GenericAdder +from logprep.util.getter import RefreshableGetter from tests.unit.processor.base import BaseProcessorTestCase +from tests.conftest import mock_env RULES_DIR_MISSING = "tests/testdata/unit/generic_adder/rules_missing" RULES_DIR_INVALID = "tests/testdata/unit/generic_adder/rules_invalid" @@ -461,3 +464,56 @@ def test_add_only_copies(self): assert event["some_dict_field"] == {"some_key": "some_value"} assert event["some_dict_field"] is not rule_add["some_dict_field"], "only copies in events" + + @responses.activate + def test_adds_mapped_response_fields_from_event_templated_url(self): + resolved_url = "https://values.example/acme" + responses.add( + responses.GET, + resolved_url, + json={ + "user": {"name": "Alice"}, + "risk": {"score": 7}, + }, + ) + configuration = { + "dynamic_generic_adder": { + "type": "generic_adder", + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://${GENERIC_ADDER_HOST}/${tenant.id}", + "target_field_mapping": { + "user.name": "enrichment.user", + "risk.score": "enrichment.score", + }, + } + }, + } + ], + } + } + + RefreshableGetter.reset() + with mock_env({"GENERIC_ADDER_HOST": "values.example"}): + processor = typing.cast(GenericAdder, self._create_test_instance(configuration)) + processor.setup() + + first_event = {"tenant": {"id": "acme"}} + second_event = {"tenant": {"id": "acme"}} + first_result = processor.process(first_event) + second_result = processor.process(second_event) + + processor.shut_down() + + assert first_result.errors == [] + assert second_result.errors == [] + assert first_event == { + "tenant": {"id": "acme"}, + "enrichment": {"user": "Alice", "score": 7}, + } + assert second_event == first_event + assert len(responses.calls) == 1 + assert responses.calls[0].request.url == resolved_url diff --git a/tests/unit/processor/generic_adder/test_generic_adder_rule.py b/tests/unit/processor/generic_adder/test_generic_adder_rule.py index 184a6a91a..217612eb7 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder_rule.py +++ b/tests/unit/processor/generic_adder/test_generic_adder_rule.py @@ -1,12 +1,13 @@ # pylint: disable=missing-docstring # pylint: disable=protected-access import json +import typing from pathlib import Path import pytest import responses -from logprep.processor.generic_adder.rule import GenericAdderRule +from logprep.processor.generic_adder.rule import AddFromUrlConfig, GenericAdderRule from logprep.util.defaults import ENV_NAME_LOGPREP_GETTER_CONFIG from logprep.util.getter import HttpGetter, RefreshableGetter from tests.conftest import mock_env @@ -28,6 +29,117 @@ def fixture_rule_definition(): class TestGenericAdderRule: + @pytest.mark.parametrize( + ("url_config", "expected_target_field", "expected_target_field_mapping"), + [ + pytest.param( + { + "url": "https://values.example/${tenant.id}", + "target_field": "enrichment", + }, + "enrichment", + {}, + id="whole-response", + ), + pytest.param( + { + "url": "https://values.example/${tenant.id}", + "target_field_mapping": { + "user.name": "enrichment.user", + "risk.score": "enrichment.score", + }, + }, + None, + { + "user.name": "enrichment.user", + "risk.score": "enrichment.score", + }, + id="field-mapping", + ), + ], + ) + def test_converts_add_from_url_configuration( + self, url_config, expected_target_field, expected_target_field_mapping + ): + rule = GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": {"add_from_url": url_config}, + } + ) + + config = typing.cast(GenericAdderRule.Config, rule._config) + + assert isinstance(config.add_from_url, AddFromUrlConfig) + assert config.add_from_url.target_field == expected_target_field + assert config.add_from_url.target_field_mapping == expected_target_field_mapping + + @pytest.mark.parametrize( + ("url_config", "error_message"), + [ + pytest.param( + {"url": "https://values.example/${tenant}"}, + "requires target_field or target_field_mapping", + id="missing-target", + ), + pytest.param( + { + "url": "https://values.example/${tenant}", + "target_field": "enrichment", + "target_field_mapping": {"risk": "enrichment.risk"}, + }, + "only one of target_field or target_field_mapping", + id="ambiguous-target", + ), + ], + ) + def test_rejects_invalid_add_from_url_target_configuration(self, url_config, error_message): + with pytest.raises(ValueError, match=error_message): + GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": {"add_from_url": url_config}, + } + ) + + def test_rejects_rule_without_addition_source(self): + with pytest.raises( + ValueError, + match="one of add, add_from_file or add_from_url", + ): + GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": {}, + } + ) + + @responses.activate + def test_resolves_dotted_event_field_and_adds_complete_response(self): + resolved_url = "https://values.example/acme" + response_content = { + "user": {"name": "Alice"}, + "risk": {"score": 7}, + } + responses.add(responses.GET, resolved_url, json=response_content) + rule = GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant.id}", + "target_field": "enrichment", + } + }, + } + ) + rule.init_generic_adder("generic-adder-test") + + additions = rule.add({"tenant": {"id": "acme"}}) + + assert additions == {"enrichment": response_content} + assert responses.calls[0].request.url == resolved_url + @pytest.mark.parametrize( "testcase, other_rule_definition, is_equal", [ diff --git a/tests/unit/util/test_getter.py b/tests/unit/util/test_getter.py index 014524d1a..a274311ee 100644 --- a/tests/unit/util/test_getter.py +++ b/tests/unit/util/test_getter.py @@ -1472,6 +1472,31 @@ def test_get_collection_parses_json_if_yaml_fails(self, mock_parse_yaml): http_getter.get_collection() mock_parse_yaml.assert_called_once() + @responses.activate + def test_get_collection_extracts_content_field(self): + responses.add( + responses.GET, + "http://something", + json={"payload": {"answer": 42}}, + ) + + http_getter = GetterFactory.from_string("http://something") + + assert http_getter.get_collection("payload") == {"answer": 42} + + @responses.activate + def test_get_collection_rejects_content_field_for_non_mapping_content(self): + responses.add( + responses.GET, + "http://something", + json=["one", "two"], + ) + + http_getter = GetterFactory.from_string("http://something") + + with pytest.raises(ValueError, match="Expected mapping type when content_field is set"): + http_getter.get_collection("payload") + @mock.patch("logprep.abc.getter.Getter.get_collection", return_value="not a dict") def test_get_dict_raises_exception_if_result_not_dict(self, _): http_getter = GetterFactory.from_string("http://something") From 635660b7604e3d27bbd94d1db74dd7965da1928c Mon Sep 17 00:00:00 2001 From: pablu Date: Fri, 24 Jul 2026 09:38:45 +0200 Subject: [PATCH 6/9] fix failing integration tests where config gets wrongly verified --- logprep/ng/processor/generic_adder/processor.py | 16 ++++++++++++++++ logprep/processor/generic_adder/rule.py | 16 +++++++--------- 2 files changed, 23 insertions(+), 9 deletions(-) diff --git a/logprep/ng/processor/generic_adder/processor.py b/logprep/ng/processor/generic_adder/processor.py index 6dc324d0a..1c8b84ea0 100644 --- a/logprep/ng/processor/generic_adder/processor.py +++ b/logprep/ng/processor/generic_adder/processor.py @@ -25,10 +25,12 @@ """ import typing +from typing import Sequence from logprep.ng.abc.processor import Processor from logprep.processor.base.rule import Rule from logprep.processor.generic_adder.rule import GenericAdderRule +from logprep.util.getter import RefreshableGetter from logprep.util.helper import add_fields_to @@ -37,8 +39,22 @@ class GenericAdder(Processor): rule_class = GenericAdderRule + @property + def _rules(self) -> Sequence[GenericAdderRule]: + """Returns all rules""" + return typing.cast(Sequence[GenericAdderRule], self.rules) + + async def setup(self): + await super().setup() + for rule in self._rules: + rule.init_generic_adder(self._job_tag_for_cleanup) + def _apply_rules(self, event: dict, rule: Rule) -> None: rule = typing.cast(GenericAdderRule, rule) items_to_add = rule.add(event) if items_to_add: add_fields_to(event, items_to_add, rule, rule.merge_with_target, rule.overwrite_target) + + def _shut_down(self) -> None: + RefreshableGetter.remove_callbacks_for_tag(self._job_tag_for_cleanup) + return super()._shut_down() diff --git a/logprep/processor/generic_adder/rule.py b/logprep/processor/generic_adder/rule.py index 209540775..ecada5375 100644 --- a/logprep/processor/generic_adder/rule.py +++ b/logprep/processor/generic_adder/rule.py @@ -203,22 +203,20 @@ def _refresh_add(self): def __attrs_post_init__(self): self._base_add = copy.deepcopy(self.add) - if ( - self.add_from_file is not None or len(self.add) > 0 - ) and self.add_from_url is not None: + if (self.add_from_file or self.add) and self.add_from_url is not None: raise ValueError( "only one of add_from_file + add or add_from_url is allowed per GenericAdder rule" ) - if not self.add and self.add_from_file is None and self.add_from_url is None: + if not self.add and not self.add_from_file and self.add_from_url is None: raise ValueError( "one of add, add_from_file or add_from_url is neccessary per GenericAdder rule" ) if ( self.add_from_url is not None - and self.add_from_url.target_field is not None - and len(self.add_from_url.target_field_mapping) > 0 + and self.add_from_url.target_field + and self.add_from_url.target_field_mapping ): raise ValueError( "only one of target_field or target_field_mapping is allowed per GenericAdder rule" @@ -226,8 +224,8 @@ def __attrs_post_init__(self): if ( self.add_from_url is not None - and self.add_from_url.target_field is None - and len(self.add_from_url.target_field_mapping) <= 0 + and not self.add_from_url.target_field + and not self.add_from_url.target_field_mapping ): raise ValueError( "one of target_field or target_field_mapping is neccessary per GenericAdder rule" @@ -375,7 +373,7 @@ def add(self, event: dict) -> dict: """Returns the fields to add""" config = typing.cast(GenericAdderRule.Config, self._config) - if config.add_from_file is not None or len(config.add) > 0: + if config.add_from_file or config.add: return config.add assert config.add_from_url is not None From 749c367b9a9fa67a2aab9ad4a2791f90091dde10 Mon Sep 17 00:00:00 2001 From: pablu Date: Fri, 24 Jul 2026 10:07:52 +0200 Subject: [PATCH 7/9] update changelog --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index b88676933..cbcacc28a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,7 @@ ### Breaking ### Features +* generic_adder: add support for templated http urls ### Improvements From 976248f941293152b2f61b33888c9f17d16abc4e Mon Sep 17 00:00:00 2001 From: pablu Date: Fri, 24 Jul 2026 10:24:36 +0200 Subject: [PATCH 8/9] use get_dotted_field_with_missing and update missing behaviour --- logprep/processor/generic_adder/rule.py | 26 +++++++++++++++++++------ 1 file changed, 20 insertions(+), 6 deletions(-) diff --git a/logprep/processor/generic_adder/rule.py b/logprep/processor/generic_adder/rule.py index ecada5375..225fb0e80 100644 --- a/logprep/processor/generic_adder/rule.py +++ b/logprep/processor/generic_adder/rule.py @@ -83,6 +83,7 @@ # pylint: enable=anomalous-backslash-in-string import copy +import logging import os import typing @@ -93,7 +94,15 @@ from logprep.processor.field_manager.rule import FieldManagerRule from logprep.util.converters import convert_from_dict from logprep.util.getter import GetterFactory, HttpGetter, RefreshableGetter -from logprep.util.helper import DottedTemplate, FieldValue, get_dotted_field_value +from logprep.util.helper import ( + MISSING, + DottedTemplate, + FieldValue, + get_dotted_field_value, + get_dotted_field_value_with_missing, +) + +logger = logging.getLogger("GenericAdder") @define(kw_only=True, frozen=True) @@ -354,11 +363,16 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: mapping_source_field, mapping_target_field, ) in config.add_from_url.target_field_mapping.items(): - # what should happen with missing fields? Right now None is skipped - # change this to get_dotted_field_value_with_missing - items_to_add[mapping_target_field] = get_dotted_field_value( - content, mapping_source_field - ) + item = get_dotted_field_value_with_missing(content, mapping_source_field) + if item is MISSING: + logger.warning( + "could not add source_field: %s for target_field: %s because was missing", + mapping_source_field, + mapping_target_field, + ) + continue + + items_to_add[mapping_target_field] = item return items_to_add From d0045f70a4a98cedc5d9cbf8893d153abf5d64b5 Mon Sep 17 00:00:00 2001 From: pablu Date: Fri, 24 Jul 2026 10:38:44 +0200 Subject: [PATCH 9/9] update rule / processor to behave similar to list_comparison --- .../ng/processor/generic_adder/processor.py | 7 +- logprep/processor/generic_adder/processor.py | 6 +- logprep/processor/generic_adder/rule.py | 51 ++++++- .../generic_adder/test_generic_adder.py | 119 ++++++++++++++- .../generic_adder/test_generic_adder.py | 138 +++++++++++++++++- .../generic_adder/test_generic_adder_rule.py | 85 ++++++++++- 6 files changed, 393 insertions(+), 13 deletions(-) diff --git a/logprep/ng/processor/generic_adder/processor.py b/logprep/ng/processor/generic_adder/processor.py index 1c8b84ea0..d722000e9 100644 --- a/logprep/ng/processor/generic_adder/processor.py +++ b/logprep/ng/processor/generic_adder/processor.py @@ -28,6 +28,7 @@ from typing import Sequence from logprep.ng.abc.processor import Processor +from logprep.processor.base.exceptions import ProcessingWarning from logprep.processor.base.rule import Rule from logprep.processor.generic_adder.rule import GenericAdderRule from logprep.util.getter import RefreshableGetter @@ -51,7 +52,11 @@ async def setup(self): def _apply_rules(self, event: dict, rule: Rule) -> None: rule = typing.cast(GenericAdderRule, rule) - items_to_add = rule.add(event) + + try: + items_to_add = rule.add(event) + except Exception as error: + raise ProcessingWarning(str(error), rule, event) from error if items_to_add: add_fields_to(event, items_to_add, rule, rule.merge_with_target, rule.overwrite_target) diff --git a/logprep/processor/generic_adder/processor.py b/logprep/processor/generic_adder/processor.py index 5ae424bb7..0d214502d 100644 --- a/logprep/processor/generic_adder/processor.py +++ b/logprep/processor/generic_adder/processor.py @@ -28,6 +28,7 @@ from typing import Sequence from logprep.abc.processor import Processor +from logprep.processor.base.exceptions import ProcessingWarning from logprep.processor.generic_adder.rule import GenericAdderRule from logprep.util.getter import RefreshableGetter from logprep.util.helper import add_fields_to @@ -49,7 +50,10 @@ def setup(self): rule.init_generic_adder(self._job_tag_for_cleanup) def _apply_rules(self, event: dict, rule: GenericAdderRule): - items_to_add = rule.add(event) + try: + items_to_add = rule.add(event) + except Exception as error: + raise ProcessingWarning(str(error), rule, event) from error if items_to_add: add_fields_to(event, items_to_add, rule, rule.merge_with_target, rule.overwrite_target) diff --git a/logprep/processor/generic_adder/rule.py b/logprep/processor/generic_adder/rule.py index 225fb0e80..ee1292775 100644 --- a/logprep/processor/generic_adder/rule.py +++ b/logprep/processor/generic_adder/rule.py @@ -93,6 +93,7 @@ from logprep.processor.base.rule import InvalidRuleDefinitionError from logprep.processor.field_manager.rule import FieldManagerRule from logprep.util.converters import convert_from_dict +from logprep.util.environ import ENV_VARS from logprep.util.getter import GetterFactory, HttpGetter, RefreshableGetter from logprep.util.helper import ( MISSING, @@ -281,6 +282,7 @@ def __init__(self, filter_rule: FilterExpression, config: Config, processor_name self._is_dynamic: bool = False self._dynamic_template: DottedTemplate self._dynamic_identifiers: tuple[str, ...] = () + self._static_uri: str | None = None def init_generic_adder(self, job_tag: str) -> None: self._callback_tag = job_tag @@ -292,13 +294,30 @@ def init_generic_adder(self, job_tag: str) -> None: assert config.add_from_url is not None base_template = DottedTemplate(config.add_from_url.url) - resolved_template = DottedTemplate(base_template.safe_substitute({**os.environ})) + resolved_template = DottedTemplate(base_template.safe_substitute({**ENV_VARS})) self._dynamic_template = resolved_template self._dynamic_identifiers = tuple(resolved_template.get_identifiers()) if len(self._dynamic_identifiers) > 0: self._is_dynamic = True + if not self._is_dynamic: + static_uri = resolved_template.substitute() + http_getter = GetterFactory.from_string(static_uri) + + assert isinstance(http_getter, HttpGetter) + + self._update_static_content(http_getter, static_uri) + + http_getter.add_callback( + self._callback_tag, + self._update_static_content, + deduplication_key=(self._callback_tag, static_uri, id(self)), + fnc_args=[http_getter, static_uri], + ) + + self._static_uri = static_uri + def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: config = typing.cast(GenericAdderRule.Config, self._config) @@ -311,16 +330,16 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: for identifier, val in key_val.items(): if val is None: raise ValueError( - f"missing event field {identifier!r} for dynamic list comparison path" + f"missing event field {identifier!r} for dynamic generic adder path" ) if not isinstance(val, (str, int)): raise ValueError( - f"value for list comparison field {identifier!r} is not a scalar value" + f"value for generic adder field {identifier!r} is not a scalar value" ) pass dynamic_resolved = self._dynamic_template.substitute(key_val) - content: FieldValue | None = None + content: FieldValue = None if dynamic_resolved not in self._dynamic_content: http_getter = GetterFactory.from_string(dynamic_resolved) assert isinstance(http_getter, HttpGetter) @@ -345,20 +364,26 @@ def _dynamic_add_from_url(self, event: dict) -> dict[str, FieldValue]: deduplication_key=(tag, dynamic_resolved, id(self)), fnc_args=[dynamic_resolved], ) - else: RefreshableGetter.keep_alive_for_target(dynamic_resolved) content = self._dynamic_content[dynamic_resolved] + return self._content_to_items_to_add(content) + + def _content_to_items_to_add(self, content: FieldValue): items_to_add: dict[str, FieldValue] = {} + config = typing.cast(GenericAdderRule.Config, self._config) + assert config.add_from_url + if config.add_from_url.target_field: items_to_add[config.add_from_url.target_field] = content else: assert config.add_from_url.target_field_mapping is not None - # TODO: Check this differently, what should this be? - assert isinstance(content, dict) + if not isinstance(content, dict): + raise ValueError("add_from_url.target_field_mapping requires a mapping response") + for ( mapping_source_field, mapping_target_field, @@ -380,6 +405,14 @@ def _update_dynamic_content(self, http_getter: HttpGetter, resolved_uri: str): content = http_getter.get_collection() self._dynamic_content[resolved_uri] = content + def _update_static_content(self, getter: HttpGetter, uri: str) -> None: + try: + self._update_dynamic_content(getter, uri) + except Exception as error: + self.mark_failed(error) + else: + self.clear_failed() + def _cleanup(self, resolved_uri: str): self._dynamic_content.pop(resolved_uri, None) @@ -390,5 +423,9 @@ def add(self, event: dict) -> dict: if config.add_from_file or config.add: return config.add + if not self._is_dynamic: + assert self._static_uri + return self._content_to_items_to_add(self._dynamic_content[self._static_uri]) + assert config.add_from_url is not None return self._dynamic_add_from_url(event) diff --git a/tests/unit/ng/processor/generic_adder/test_generic_adder.py b/tests/unit/ng/processor/generic_adder/test_generic_adder.py index 60f33d2c1..803acad83 100644 --- a/tests/unit/ng/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/ng/processor/generic_adder/test_generic_adder.py @@ -14,7 +14,8 @@ from logprep.factory import Factory from logprep.ng.abc.event import InputMeta, LogEvent from logprep.ng.processor.generic_adder.processor import GenericAdder -from logprep.processor.base.exceptions import InvalidRuleDefinitionError +from logprep.processor.base.exceptions import InvalidRuleDefinitionError, ProcessingWarning +from logprep.util.getter import HttpGetter, RefreshableGetter from tests.unit.ng.processor.base import BaseProcessorTestCase from tests.unit.processor.generic_adder.test_generic_adder import ( failure_test_cases as non_ng_failure_test_cases, @@ -135,3 +136,119 @@ async def test_adds_response_from_event_templated_url(self): "enrichment": response_content, } assert responses.calls[0].request.url == resolved_url + + @responses.activate + async def test_dynamic_url_failure_is_event_scoped(self): + failed_url = "https://values.example/acme" + successful_url = "https://values.example/beta" + responses.add(responses.GET, failed_url, status=500) + responses.add(responses.GET, successful_url, json={"risk": {"score": 7}}) + RefreshableGetter.reset() + processor = self._create_test_instance( + { + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant}", + "target_field": "enrichment", + } + }, + } + ] + } + ) + await processor.setup() + rule = processor.rules[0] + failed_event = {"tenant": "acme"} + successful_event = {"tenant": "beta"} + + failed_result = await processor.process( + LogEvent(failed_event, original=b"", input_meta=InputMeta()) + ) + successful_result = await processor.process( + LogEvent(successful_event, original=b"", input_meta=InputMeta()) + ) + + assert failed_result.errors == [] + assert len(failed_result.warnings) == 1 + assert isinstance(failed_result.warnings[0], ProcessingWarning) + assert failed_event == { + "tenant": "acme", + "tags": ["_generic_adder_failure"], + } + assert rule.data_error is None + assert len(HttpGetter._target_to_data_caches[failed_url].callbacks) == 0 + assert len(HttpGetter._target_to_data_caches[failed_url].cleanup_callbacks) == 0 + + assert successful_result.errors == [] + assert successful_result.warnings == [] + assert successful_event == { + "tenant": "beta", + "enrichment": {"risk": {"score": 7}}, + } + + async def test_missing_dynamic_url_field_adds_warning_without_clearing_event(self): + processor = self._create_test_instance( + { + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant.id}", + "target_field": "enrichment", + } + }, + } + ] + } + ) + await processor.setup() + event = {"message": "preserved"} + + result = await processor.process(LogEvent(event, original=b"", input_meta=InputMeta())) + + assert result.errors == [] + assert len(result.warnings) == 1 + assert "missing event field 'tenant.id'" in str(result.warnings[0]) + assert event == { + "message": "preserved", + "tags": ["_generic_adder_failure"], + } + + @responses.activate + async def test_mapping_response_type_error_adds_warning_without_clearing_event(self): + url = "https://values.example/acme" + responses.add(responses.GET, url, json=["not", "a", "mapping"]) + RefreshableGetter.reset() + processor = self._create_test_instance( + { + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant}", + "target_field_mapping": { + "risk.score": "enrichment.score", + }, + } + }, + } + ] + } + ) + await processor.setup() + event = {"tenant": "acme"} + + result = await processor.process(LogEvent(event, original=b"", input_meta=InputMeta())) + + assert result.errors == [] + assert len(result.warnings) == 1 + assert "target_field_mapping requires a mapping response" in str(result.warnings[0]) + assert event == { + "tenant": "acme", + "tags": ["_generic_adder_failure"], + } diff --git a/tests/unit/processor/generic_adder/test_generic_adder.py b/tests/unit/processor/generic_adder/test_generic_adder.py index 7a665f077..3fff06c8d 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder.py +++ b/tests/unit/processor/generic_adder/test_generic_adder.py @@ -10,9 +10,9 @@ import responses from logprep.factory import Factory -from logprep.processor.base.exceptions import InvalidRuleDefinitionError +from logprep.processor.base.exceptions import InvalidRuleDefinitionError, ProcessingWarning from logprep.processor.generic_adder.processor import GenericAdder -from logprep.util.getter import RefreshableGetter +from logprep.util.getter import HttpGetter, RefreshableGetter from tests.unit.processor.base import BaseProcessorTestCase from tests.conftest import mock_env @@ -517,3 +517,137 @@ def test_adds_mapped_response_fields_from_event_templated_url(self): assert second_event == first_event assert len(responses.calls) == 1 assert responses.calls[0].request.url == resolved_url + + @responses.activate + def test_dynamic_url_failure_is_event_scoped(self): + failed_url = "https://values.example/acme" + successful_url = "https://values.example/beta" + responses.add(responses.GET, failed_url, status=500) + responses.add(responses.GET, successful_url, json={"risk": {"score": 7}}) + RefreshableGetter.reset() + processor = typing.cast( + GenericAdder, + self._create_test_instance( + { + "dynamic_generic_adder": { + "type": "generic_adder", + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant}", + "target_field": "enrichment", + } + }, + } + ], + } + } + ), + ) + processor.setup() + rule = processor.rules[0] + failed_event = {"tenant": "acme"} + successful_event = {"tenant": "beta"} + + failed_result = processor.process(failed_event) + successful_result = processor.process(successful_event) + + assert failed_result.errors == [] + assert len(failed_result.warnings) == 1 + assert isinstance(failed_result.warnings[0], ProcessingWarning) + assert failed_event == { + "tenant": "acme", + "tags": ["_generic_adder_failure"], + } + assert rule.data_error is None + assert len(HttpGetter._target_to_data_caches[failed_url].callbacks) == 0 + assert len(HttpGetter._target_to_data_caches[failed_url].cleanup_callbacks) == 0 + + assert successful_result.errors == [] + assert successful_result.warnings == [] + assert successful_event == { + "tenant": "beta", + "enrichment": {"risk": {"score": 7}}, + } + + processor.shut_down() + + def test_missing_dynamic_url_field_adds_warning_without_clearing_event(self): + processor = typing.cast( + GenericAdder, + self._create_test_instance( + { + "dynamic_generic_adder": { + "type": "generic_adder", + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant.id}", + "target_field": "enrichment", + } + }, + } + ], + } + } + ), + ) + processor.setup() + event = {"message": "preserved"} + + result = processor.process(event) + + assert result.errors == [] + assert len(result.warnings) == 1 + assert "missing event field 'tenant.id'" in str(result.warnings[0]) + assert event == { + "message": "preserved", + "tags": ["_generic_adder_failure"], + } + + @responses.activate + def test_mapping_response_type_error_adds_warning_without_clearing_event(self): + url = "https://values.example/acme" + responses.add(responses.GET, url, json=["not", "a", "mapping"]) + RefreshableGetter.reset() + processor = typing.cast( + GenericAdder, + self._create_test_instance( + { + "dynamic_generic_adder": { + "type": "generic_adder", + "rules": [ + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant}", + "target_field_mapping": { + "risk.score": "enrichment.score", + }, + } + }, + } + ], + } + } + ), + ) + processor.setup() + event = {"tenant": "acme"} + + result = processor.process(event) + + processor.shut_down() + + assert result.errors == [] + assert len(result.warnings) == 1 + assert "target_field_mapping requires a mapping response" in str(result.warnings[0]) + assert event == { + "tenant": "acme", + "tags": ["_generic_adder_failure"], + } diff --git a/tests/unit/processor/generic_adder/test_generic_adder_rule.py b/tests/unit/processor/generic_adder/test_generic_adder_rule.py index 217612eb7..62c8aae4c 100644 --- a/tests/unit/processor/generic_adder/test_generic_adder_rule.py +++ b/tests/unit/processor/generic_adder/test_generic_adder_rule.py @@ -9,7 +9,7 @@ from logprep.processor.generic_adder.rule import AddFromUrlConfig, GenericAdderRule from logprep.util.defaults import ENV_NAME_LOGPREP_GETTER_CONFIG -from logprep.util.getter import HttpGetter, RefreshableGetter +from logprep.util.getter import GetterFactory, HttpGetter, RefreshableGetter from tests.conftest import mock_env @@ -140,6 +140,89 @@ def test_resolves_dotted_event_field_and_adds_complete_response(self): assert additions == {"enrichment": response_content} assert responses.calls[0].request.url == resolved_url + @responses.activate + def test_static_url_loads_during_setup_and_registers_only_refresh_callback(self): + url = "https://values.example/static" + response_content = {"risk": {"score": 7}} + responses.add(responses.GET, url, json=response_content) + RefreshableGetter.reset() + rule = GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": url, + "target_field": "enrichment", + } + }, + } + ) + + rule.init_generic_adder("generic-adder-test") + + assert rule.add({}) == {"enrichment": response_content} + assert rule.add({}) == {"enrichment": response_content} + assert len(responses.calls) == 1 + assert len(HttpGetter._target_to_data_caches[url].callbacks) == 1 + assert len(HttpGetter._target_to_data_caches[url].cleanup_callbacks) == 0 + + @responses.activate + def test_static_url_recovers_after_failed_initial_load(self, tmp_path): + url = "https://values.example/static" + response_content = {"risk": {"score": 7}} + responses.add(responses.GET, url, status=500) + RefreshableGetter.reset() + getter_config = tmp_path / "http_getter.json" + getter_config.write_text(json.dumps({url: {"refresh_interval": 1}})) + rule = GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": url, + "target_field": "enrichment", + } + }, + } + ) + + with mock_env({ENV_NAME_LOGPREP_GETTER_CONFIG: str(getter_config)}): + rule.init_generic_adder("generic-adder-test") + getter = GetterFactory.from_string(url) + assert isinstance(getter, HttpGetter) + assert getter.scheduler is not None + + assert rule.data_error is not None + assert len(getter.shared.callbacks) == 1 + assert len(getter.shared.cleanup_callbacks) == 0 + + responses.replace(responses.GET, url, json=response_content) + getter.scheduler.run_all() + + assert rule.data_error is None + assert rule.add({}) == {"enrichment": response_content} + + def test_target_field_mapping_skips_missing_values_but_preserves_none(self, caplog): + rule = GenericAdderRule.create_from_dict( + { + "filter": "*", + "generic_adder": { + "add_from_url": { + "url": "https://values.example/${tenant}", + "target_field_mapping": { + "present": "enrichment.present", + "missing": "enrichment.missing", + }, + } + }, + } + ) + + additions = rule._content_to_items_to_add({"present": None}) + + assert additions == {"enrichment.present": None} + assert "source_field: missing" in caplog.text + @pytest.mark.parametrize( "testcase, other_rule_definition, is_equal", [