From d880ef819392087485690f15a4cd6f7fa94a8feb Mon Sep 17 00:00:00 2001 From: kev365 <48181394+kev365@users.noreply.github.com> Date: Thu, 30 Jul 2026 18:57:35 -0500 Subject: [PATCH 1/2] Correction to resolve the event tag once per exported event #5185 The output and formatting engine resolved the event tag twice for every exported event, once in _ExportEvents for the event filter and again in _FlushExportBuffer for the output module. The event tag is now resolved once and carried through PsortEventHeap, or resolved when the export buffer is flushed when no event filter is used. Co-Authored-By: Claude Opus 5 (1M context) --- plaso/multi_process/output_engine.py | 65 +++++++++++++++++++++++----- 1 file changed, 54 insertions(+), 11 deletions(-) diff --git a/plaso/multi_process/output_engine.py b/plaso/multi_process/output_engine.py index 8a83b9b620..fc26dbf9ca 100644 --- a/plaso/multi_process/output_engine.py +++ b/plaso/multi_process/output_engine.py @@ -39,12 +39,18 @@ def PopEvent(self): EventObject: event. EventData: event data. EventDataStream: event data stream. + EventTag: event tag or None if not resolved. """ try: - event_values_hash, _, event, event_data, event_data_stream = heapq.heappop( - self._heap - ) - return event_values_hash, event, event_data, event_data_stream + ( + event_values_hash, + _, + event, + event_data, + event_data_stream, + event_tag, + ) = heapq.heappop(self._heap) + return event_values_hash, event, event_data, event_data_stream, event_tag except IndexError: return None @@ -67,13 +73,15 @@ def PopEvents(self): yield heap_values heap_values = self.PopEvent() - def PushEvent(self, event, event_data, event_data_stream): + def PushEvent(self, event, event_data, event_data_stream, event_tag=None): """Pushes an event onto the heap. Args: event (EventObject): event. event_data (EventData): event data. event_data_stream (EventDataStream): event data stream. + event_tag (Optional[EventTag]): event tag, or None if the event tag + should be resolved when the export buffer is flushed. """ event_values_hash = getattr(event_data, "_event_values_hash", None) if event_values_hash is None: @@ -90,7 +98,14 @@ def PushEvent(self, event, event_data, event_data_stream): # similar event values. heapq.heappush( self._heap, - (event_values_hash, timestamp_desc, event, event_data, event_data_stream), + ( + event_values_hash, + timestamp_desc, + event, + event_data, + event_data_stream, + event_tag, + ), ) @@ -117,6 +132,7 @@ def __init__(self): self._number_of_consumed_events = 0 self._output_mediator = None self._processing_configuration = None + self._resolve_event_tag_at_flush = True self._status = definitions.STATUS_INDICATOR_IDLE self._status_update_callback = None @@ -171,6 +187,7 @@ def _ExportEvent( event, event_data, event_data_stream, + event_tag=None, deduplicate_events=True, ): """Exports an event using an output module. @@ -181,6 +198,8 @@ def _ExportEvent( event (EventObject): event. event_data (EventData): event data. event_data_stream (EventDataStream): event data stream. + event_tag (Optional[EventTag]): event tag, or None if the event tag + should be resolved when the export buffer is flushed. deduplicate_events (Optional[bool]): True if events should be deduplicated. """ @@ -193,7 +212,9 @@ def _ExportEvent( ) self._export_event_timestamp = event.timestamp - self._export_event_heap.PushEvent(event, event_data, event_data_stream) + self._export_event_heap.PushEvent( + event, event_data, event_data_stream, event_tag=event_tag + ) def _ExportEvents( self, @@ -220,6 +241,12 @@ def _ExportEvents( """ self._status = definitions.STATUS_INDICATOR_EXPORTING + # An event filter can match on the labels of an event tag, in which case + # the event tag needs to be resolved before the filter is applied. + # Otherwise it is resolved when the export buffer is flushed, after + # duplicate events have been discarded. + self._resolve_event_tag_at_flush = event_filter is None + time_slice_buffer = None time_slice_range = None @@ -252,8 +279,10 @@ def _ExportEvents( else: event_data_stream = None - event_identifier = event.GetIdentifier() - event_tag = storage_reader.GetEventTagByEventIdentifer(event_identifier) + event_tag = None + if not self._resolve_event_tag_at_flush: + event_identifier = event.GetIdentifier() + event_tag = storage_reader.GetEventTagByEventIdentifer(event_identifier) if time_slice_range and event.timestamp != time_slice.event_timestamp: self._events_status.number_of_events_from_time_slice += 1 @@ -281,6 +310,7 @@ def _ExportEvents( event, event_data, event_data_stream, + event_tag=event_tag, deduplicate_events=deduplicate_events, ) self._number_of_consumed_events += 1 @@ -301,12 +331,22 @@ def _ExportEvents( event_in_buffer, event_data_in_buffer, ) in time_slice_buffer.Flush(): + event_tag_in_buffer = None + if not self._resolve_event_tag_at_flush: + event_identifier = event_in_buffer.GetIdentifier() + event_tag_in_buffer = ( + storage_reader.GetEventTagByEventIdentifer( + event_identifier + ) + ) + self._ExportEvent( storage_reader, output_module, event_in_buffer, event_data_in_buffer, event_data_stream, + event_tag=event_tag_in_buffer, deduplicate_events=deduplicate_events, ) self._number_of_consumed_events += 1 @@ -321,6 +361,7 @@ def _ExportEvents( event, event_data, event_data_stream, + event_tag=event_tag, deduplicate_events=deduplicate_events, ) self._number_of_consumed_events += 1 @@ -356,6 +397,7 @@ def _FlushExportBuffer( event, event_data, event_data_stream, + event_tag, ) in self._export_event_heap.PopEvents(): timestamp_desc = event.timestamp_desc @@ -367,8 +409,9 @@ def _FlushExportBuffer( self._events_status.number_of_duplicate_events += 1 continue - event_identifier = event.GetIdentifier() - event_tag = storage_reader.GetEventTagByEventIdentifer(event_identifier) + if self._resolve_event_tag_at_flush: + event_identifier = event.GetIdentifier() + event_tag = storage_reader.GetEventTagByEventIdentifer(event_identifier) if timestamp_desc in ( definitions.TIME_DESCRIPTION_LAST_ACCESS, From 2da553f1ef0d766b5f7907b51b7f7da415b26d5b Mon Sep 17 00:00:00 2001 From: kev365 <48181394+kev365@users.noreply.github.com> Date: Sun, 2 Aug 2026 15:02:24 -0500 Subject: [PATCH 2/2] Pass event tag resolution as an argument instead of object state #5185 Addresses review feedback: replaces the _resolve_event_tag_at_flush attribute with a resolve_event_tag argument passed to _ExportEvent and _FlushExportBuffer. Co-Authored-By: Claude Opus 5 (1M context) --- plaso/multi_process/output_engine.py | 33 +++++++++++++++++++++------- 1 file changed, 25 insertions(+), 8 deletions(-) diff --git a/plaso/multi_process/output_engine.py b/plaso/multi_process/output_engine.py index fc26dbf9ca..a25399f7fa 100644 --- a/plaso/multi_process/output_engine.py +++ b/plaso/multi_process/output_engine.py @@ -132,7 +132,6 @@ def __init__(self): self._number_of_consumed_events = 0 self._output_mediator = None self._processing_configuration = None - self._resolve_event_tag_at_flush = True self._status = definitions.STATUS_INDICATOR_IDLE self._status_update_callback = None @@ -189,6 +188,7 @@ def _ExportEvent( event_data_stream, event_tag=None, deduplicate_events=True, + resolve_event_tag=True, ): """Exports an event using an output module. @@ -202,13 +202,18 @@ def _ExportEvent( should be resolved when the export buffer is flushed. deduplicate_events (Optional[bool]): True if events should be deduplicated. + resolve_event_tag (Optional[bool]): True if the event tag should be + resolved when the export buffer is flushed. """ if ( event.timestamp != self._export_event_timestamp or self._export_event_heap.number_of_events > self._HEAP_MAXIMUM_EVENTS ): self._FlushExportBuffer( - storage_reader, output_module, deduplicate_events=deduplicate_events + storage_reader, + output_module, + deduplicate_events=deduplicate_events, + resolve_event_tag=resolve_event_tag, ) self._export_event_timestamp = event.timestamp @@ -245,7 +250,7 @@ def _ExportEvents( # the event tag needs to be resolved before the filter is applied. # Otherwise it is resolved when the export buffer is flushed, after # duplicate events have been discarded. - self._resolve_event_tag_at_flush = event_filter is None + resolve_event_tag = event_filter is None time_slice_buffer = None time_slice_range = None @@ -280,7 +285,7 @@ def _ExportEvents( event_data_stream = None event_tag = None - if not self._resolve_event_tag_at_flush: + if not resolve_event_tag: event_identifier = event.GetIdentifier() event_tag = storage_reader.GetEventTagByEventIdentifer(event_identifier) @@ -312,6 +317,7 @@ def _ExportEvents( event_data_stream, event_tag=event_tag, deduplicate_events=deduplicate_events, + resolve_event_tag=resolve_event_tag, ) self._number_of_consumed_events += 1 self._events_status.number_of_events_from_time_slice += 1 @@ -332,7 +338,7 @@ def _ExportEvents( event_data_in_buffer, ) in time_slice_buffer.Flush(): event_tag_in_buffer = None - if not self._resolve_event_tag_at_flush: + if not resolve_event_tag: event_identifier = event_in_buffer.GetIdentifier() event_tag_in_buffer = ( storage_reader.GetEventTagByEventIdentifer( @@ -348,6 +354,7 @@ def _ExportEvents( event_data_stream, event_tag=event_tag_in_buffer, deduplicate_events=deduplicate_events, + resolve_event_tag=resolve_event_tag, ) self._number_of_consumed_events += 1 self._events_status.number_of_filtered_events += 1 @@ -363,6 +370,7 @@ def _ExportEvents( event_data_stream, event_tag=event_tag, deduplicate_events=deduplicate_events, + resolve_event_tag=resolve_event_tag, ) self._number_of_consumed_events += 1 @@ -374,10 +382,16 @@ def _ExportEvents( ): break - self._FlushExportBuffer(storage_reader, output_module) + self._FlushExportBuffer( + storage_reader, output_module, resolve_event_tag=resolve_event_tag + ) def _FlushExportBuffer( - self, storage_reader, output_module, deduplicate_events=True + self, + storage_reader, + output_module, + deduplicate_events=True, + resolve_event_tag=True, ): """Flushes buffered events and writes them to the output module. @@ -386,6 +400,9 @@ def _FlushExportBuffer( output_module (OutputModule): output module. deduplicate_events (Optional[bool]): True if events should be deduplicated. + resolve_event_tag (Optional[bool]): True if the event tag should be + resolved, which is the case if it was not resolved before the + event filter was applied. """ last_event_values_hash = None last_macb_group_identifier = None @@ -409,7 +426,7 @@ def _FlushExportBuffer( self._events_status.number_of_duplicate_events += 1 continue - if self._resolve_event_tag_at_flush: + if resolve_event_tag: event_identifier = event.GetIdentifier() event_tag = storage_reader.GetEventTagByEventIdentifer(event_identifier)