Skip to content
Merged
Changes from all 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
52 changes: 42 additions & 10 deletions iotfunctions/bif.py
Original file line number Diff line number Diff line change
Expand Up @@ -627,7 +627,7 @@ def execute(self, df):
except KeyError:
last_ts_had_alert = False

if is_first_cycle and first_alert_from_previous_run is not None and first_alert_in_this_run is not None:
if is_first_cycle and first_alert_from_previous_run is not None and pd.notna(first_alert_from_previous_run) and first_alert_in_this_run is not None and pd.notna(first_alert_in_this_run):
first_alert_in_this_run = min(first_alert_from_previous_run, first_alert_in_this_run)

cache_df.loc[entity_id] = {
Expand Down Expand Up @@ -1396,7 +1396,9 @@ def execute(self, df, start_ts=None, end_ts=None, entities=None):
has_current_data = device_id in devices_in_current_df
device_registration_time = None
first_alert_from_previous_run = None
if cache_df.empty or self.dms.running_with_backtrack or not is_cached_device:
if self.dms.entity_type_type != 'DEVICE_TYPE' and (cache_df.empty or self.dms.running_with_backtrack):
device_registration_time = self._get_resource_registration_time(device_id)
elif (cache_df.empty or self.dms.running_with_backtrack or not is_cached_device) and self.dms.entity_type_type == 'DEVICE_TYPE':
device_registration_time = self._get_device_registration_time(device_id)

# Load cache for this device
Expand All @@ -1414,9 +1416,14 @@ def execute(self, df, start_ts=None, end_ts=None, entities=None):

first_alert_in_this_run = None
# During backtrack first cycle with previous alert, apply cooldown
if self.dms.running_with_backtrack and is_first_cycle and first_alert_from_previous_run is not None and self.cooldown:
cooldown_until = pd.Timestamp(first_alert_from_previous_run) + self.cooldown_timedelta
logger.info(f'Backtrack: Applying cooldown from previous run first alert {first_alert_from_previous_run} until {cooldown_until}')
if self.dms.running_with_backtrack and is_first_cycle and first_alert_from_previous_run is not None and pd.notna(first_alert_from_previous_run) and self.cooldown:
if last_event_timestamp is not None:
no_data_condition = ( df.empty or ( self.input_item is not None and ( self.input_item not in df.columns or df[self.input_item].isna().all())))
if no_data_condition:
cooldown_until = pd.Timestamp(last_event_timestamp) + self.duration_timedelta + self.cooldown_timedelta
else:
cooldown_until = pd.Timestamp(first_alert_from_previous_run) + self.cooldown_timedelta
logger.info(f'Backtrack: Applying cooldown from previous run first alert {first_alert_from_previous_run} until {cooldown_until}')

if has_current_data:
df, last_event_timestamp, cooldown_until, first_alert_in_this_run = self._process_device_with_data(df, device_id, last_event_timestamp, cooldown_until,
Expand All @@ -1426,11 +1433,16 @@ def execute(self, df, start_ts=None, end_ts=None, entities=None):
metrics_to_monitor, device_registration_time, start_ts, end_ts, is_first_cycle, first_alert_in_this_run, first_alert_from_previous_run)

# Update cache for device
if is_first_cycle and first_alert_from_previous_run is not None and first_alert_in_this_run is not None:
if is_first_cycle and first_alert_from_previous_run is not None and pd.notna(first_alert_from_previous_run) and first_alert_in_this_run is not None and pd.notna(first_alert_in_this_run):
first_alert_in_this_run = min(first_alert_from_previous_run, first_alert_in_this_run)
cache_df.at[device_id, 'last_event_timestamp'] = last_event_timestamp
cache_df.at[device_id, 'cooldown_until'] = cooldown_until
cache_df.at[device_id, 'first_alert_time'] = first_alert_in_this_run if is_first_cycle else cache_df.loc[device_id]['first_alert_time']
if is_first_cycle:
cache_df.at[device_id, 'first_alert_time'] = first_alert_in_this_run
elif is_cached_device and device_id in cache_df.index:
cache_df.at[device_id, 'first_alert_time'] = cache_df.loc[device_id, 'first_alert_time']
else:
cache_df.at[device_id, 'first_alert_time'] = None

# Store updated cache
self.cache.store_alert_cache(kpi_function_id, cache_df, self.dms.running_with_backtrack)
Expand Down Expand Up @@ -1459,7 +1471,7 @@ def _process_device_without_data(self, df, device_id, last_event_timestamp,coold
logger.info( f"Device : {device_id} has no data in current batch")
# Get last event timestamp if not in cache
now = self.dms.launch_date
if last_event_timestamp is None or pd.isna(last_event_timestamp):
if (last_event_timestamp is None or pd.isna(last_event_timestamp)) and self.dms.entity_type_type == 'DEVICE_TYPE':
last_event_timestamp = self._get_last_event_from_database(device_id, metrics_to_monitor)

if last_event_timestamp is None:
Expand Down Expand Up @@ -1717,13 +1729,33 @@ def _get_device_registration_time(self, device_id):
logger.error(f"Error querying DEVICE_LIST for {device_id}: {str(e)}")
return None

def _get_resource_registration_time(self, device_id):
"""Get resource registration time Entity_type"""
try:
logger.info(f"Quering DB to get resource registration time, resource_id = {self.dms.entity_type_id}")
query, table = self.dms.db.query( 'ENTITY_TYPE', 'IOTANALYTICS', column_names=['CREATED_TS'] )

query = query.filter( self.dms.db.get_column_object(table, 'ENTITY_TYPE_ID') == self.dms.entity_type_id)

result = self.dms.db.read_sql_query(query.statement)

if result is not None and not result.empty:
created_ts = result['created_ts'].iloc[0]
if pd.notna(created_ts):
return pd.Timestamp(created_ts)
return None
except Exception as e:
logger.error(f"Error querying ENTITY_TYPE for {device_id}: {str(e)}")
return None

def _get_gap_measurement_start_time(self, backtrack_start_ts, last_event_timestamp, device_registration_time):
"""Determine starting point for gap measurement"""
if last_event_timestamp is not None and pd.notna(last_event_timestamp) and backtrack_start_ts is None:
# Prioritize last_event_timestamp if available (most recent data point)
if last_event_timestamp is not None and pd.notna(last_event_timestamp):
logger.info(f"nodata calculation started from last_event_timestamp : {last_event_timestamp}")
return last_event_timestamp
if backtrack_start_ts is not None:
return max(backtrack_start_ts,device_registration_time) if device_registration_time is not None else backtrack_start_ts
return max(backtrack_start_ts, device_registration_time) if device_registration_time is not None else backtrack_start_ts
return device_registration_time

def _create_synthetic_alert_row(self, device_id, alert_timestamp):
Expand Down
Loading