""" Data Processing Module - Data Transformation and Manipulation Handles JSON parsing, data cleaning, and export functionality """ import json from typing import Any import pandas as pd def process_audit_data(df: pd.DataFrame) -> pd.DataFrame: """ Process raw audit data: parse JSON fields, clean data, add computed columns Args: df: Raw DataFrame from database query Returns: Processed DataFrame """ if df is None or df.empty: return df processed_df = df.copy() # Convert event_time to datetime if not already (for safety, though SQL should provide this) if 'event_time' in processed_df.columns: processed_df['event_time'] = pd.to_datetime(processed_df['event_time']) # NOTE: event_date, event_hour, event_day_of_week are now computed in SQL # NOTE: tenant_type and tenant_id are now computed in SQL # NOTE: Most null handling should be done in SQL via COALESCE/IFNULL for better performance # Only handle nulls for columns not yet processed in SQL # Minimal null handling - only for fields not yet handled in SQL # Most columns should have COALESCE in SQL queries (per refactor.md phase 3.3) if 'changed_by' in processed_df.columns: processed_df['changed_by'] = processed_df['changed_by'].fillna('Unknown') # Vectorized profile_id formatting - convert to int where possible # Keep this minimal pandas transform for final display if 'profile_id' in processed_df.columns: processed_df['profile_id'] = pd.to_numeric( processed_df['profile_id'], errors='coerce' ).fillna('') return processed_df def parse_json_field(value: Any) -> Any: """ Safely parse a JSON field Args: value: Field value (string, dict, or other) Returns: Parsed JSON object or original value if parsing fails """ if pd.isna(value) or value == '' or value is None: return None if isinstance(value, str): try: return json.loads(value) except (json.JSONDecodeError, ValueError): return value return value def parse_target_labels(value: Any) -> list[str]: """ Parse target labels array Args: value: Target labels value Returns: List of label strings """ if pd.isna(value) or value == '' or value is None: return [] parsed = parse_json_field(value) if isinstance(parsed, list): return parsed elif isinstance(parsed, str): return [parsed] else: return [] def get_property_changes(before: Any, after: Any) -> dict[str, dict[str, Any]]: """ Compare before and after properties to identify changes Args: before: Properties before the change after: Properties after the change Returns: Dictionary of changed properties with before/after values """ changes = {} before_dict = parse_json_field(before) if before else {} after_dict = parse_json_field(after) if after else {} if not isinstance(before_dict, dict): before_dict = {} if not isinstance(after_dict, dict): after_dict = {} # Find all unique keys all_keys = set(before_dict.keys()) | set(after_dict.keys()) for key in all_keys: before_val = before_dict.get(key) after_val = after_dict.get(key) if before_val != after_val: changes[key] = {'before': before_val, 'after': after_val} return changes def export_to_csv(df: pd.DataFrame) -> str: """ Export DataFrame to CSV string Args: df: DataFrame to export Returns: CSV string """ # Select columns to export (exclude complex JSON fields and raw tenant columns) export_columns = [ 'event_time', 'operation', 'changed_by', 'profile_uuid', 'profile_type', 'profile_id', 'tenant_type', 'tenant_id', 'relationship_id', 'transaction_id', 'kafka_offset', 'kafka_partition', ] # Filter to existing columns available_columns = [col for col in export_columns if col in df.columns] # Column slicing already creates a copy, no need for explicit .copy() export_df = df[available_columns] # Use pre-formatted time from SQL if available, otherwise format if 'event_time_formatted' in export_df.columns: export_df = export_df.rename(columns={'event_time_formatted': 'event_time'}) elif 'event_time' in export_df.columns: export_df['event_time'] = export_df['event_time'].dt.strftime('%Y-%m-%d %H:%M:%S') # Convert to CSV csv_string = export_df.to_csv(index=False) return csv_string def filter_by_operation(df: pd.DataFrame, operation: str) -> pd.DataFrame: """ Filter DataFrame by operation type Args: df: Input DataFrame operation: Operation type to filter (created/updated/deleted) Returns: Filtered DataFrame """ if operation.lower() == 'all': return df return df[df['operation'].str.lower() == operation.lower()] def filter_by_date_range(df: pd.DataFrame, start_date: str, end_date: str) -> pd.DataFrame: """ Filter DataFrame by date range Args: df: Input DataFrame start_date: Start date (YYYY-MM-DD) end_date: End date (YYYY-MM-DD) Returns: Filtered DataFrame """ start = pd.to_datetime(start_date) end = pd.to_datetime(end_date) + pd.Timedelta(days=1) - pd.Timedelta(seconds=1) return df[(df['event_time'] >= start) & (df['event_time'] <= end)] def filter_by_profile_type(df: pd.DataFrame, profile_type: str) -> pd.DataFrame: """ Filter DataFrame by profile type Args: df: Input DataFrame profile_type: Profile type to filter Returns: Filtered DataFrame """ if profile_type.lower() == 'all': return df return df[df['profile_type'] == profile_type] def filter_by_target_type(df: pd.DataFrame, target_type: str) -> pd.DataFrame: """ Filter DataFrame by target type Args: df: Input DataFrame target_type: Target type to filter Returns: Filtered DataFrame """ if target_type.lower() == 'all': return df return df[df['primary_target_type'] == target_type] def get_summary_statistics(df: pd.DataFrame) -> dict[str, Any]: """ Calculate summary statistics for the dataset Args: df: Input DataFrame Returns: Dictionary of summary statistics """ if df.empty: return {'total_events': 0, 'unique_profiles': 0, 'date_range': 'N/A', 'operations': {}} min_date = df['event_time'].min().strftime('%Y-%m-%d') max_date = df['event_time'].max().strftime('%Y-%m-%d') stats = { 'total_events': len(df), 'unique_profiles': df['profile_uuid'].nunique() if 'profile_uuid' in df.columns else 0, 'date_range': f'{min_date} to {max_date}', 'operations': df['operation'].value_counts().to_dict() if 'operation' in df.columns else {}, 'profile_types': df['profile_type'].value_counts().to_dict() if 'profile_type' in df.columns else {}, 'resource_types': df['tenant_type'].value_counts().to_dict() if 'tenant_type' in df.columns else {}, } return stats