""" Production Bio-Signal Telemetry Processing Engine. Parses multi-lead ECG waveforms and processes diagnostic statements into clean targets. """ import os import ast import numpy as np import pandas as pd import wfdb from scipy.signal import resample class ECGPreprocessingPipeline: def __init__(self, ecg_dir, target_sampling_rate=100, signal_length=1000): self.ecg_dir = ecg_dir self.target_sampling_rate = target_sampling_rate self.signal_length = signal_length def parse_diagnostic_statements(self, statement_string): """Converts raw string dictionaries into clean evaluation arrays.""" try: return ast.literal_eval(statement_string) except (ValueError, SyntaxError): return {} def extract_signal_waveform(self, relative_record_path): """Extracts and scales 12-lead raw data points directly from binary assets.""" absolute_path = os.path.join(self.ecg_dir, relative_record_path) if not os.path.exists(absolute_path + ".dat"): return None # Load the raw record using WFDB bindings signal_array, metadata = wfdb.rdsamp(absolute_path) # Ensure the signal fits our target runtime input dimensions exactly if len(signal_array) != self.signal_length: signal_array = resample(signal_array, self.signal_length) # Fill any missing values with zeros to maintain tensor consistency if np.isnan(signal_array).any(): signal_array = np.nan_to_num(signal_array, nan=0.0) return signal_array def map_diagnostic_to_binary_target(self, dict_statements): """ Maps multi-label classifications into clean clinical risk scores. 1: High Cardiovascular Risk, 0: Normal Baseline. """ # Clinical risk markers defined within PTB-XL schema guidelines high_risk_diagnostic_anchors = {'AMI', 'IMI', 'ALV', 'ILV', 'LVH', 'LAO/LAE', 'RVH'} for diagnostic_key in dict_statements.keys(): if diagnostic_key in high_risk_diagnostic_anchors: return 1 return 0 def process_dataset(self, limits=1000): metadata_csv = os.path.join(self.ecg_dir, "ptbxl_database.csv") if not os.path.exists(metadata_csv): print(f"[WARNING] Database tracking sheet missing at {metadata_csv}") return None, None df = pd.read_csv(metadata_csv, index_col='ecg_id') df['scp_codes'] = df['scp_codes'].apply(self.parse_diagnostic_statements) df['Target'] = df['scp_codes'].apply(self.map_diagnostic_to_binary_target) # Subsample rows to fit local memory constraints during development df_subsample = df.head(limits) signal_accumulator = [] target_accumulator = [] print(f"[INFO] Ingesting {len(df_subsample)} raw biosignal channels from disk...") for ecg_id, row in df_subsample.iterrows(): waveform = self.extract_signal_waveform(row['filename_lr']) if waveform is not None: signal_accumulator.append(waveform) target_accumulator.append(row['Target']) return np.array(signal_accumulator), np.array(target_accumulator) if __name__ == "__main__": project_root = os.path.abspath(os.path.join(os.path.dirname(__file__), "..")) ecg_data_path = os.path.join(project_root, "data", "ecg") if os.path.exists(os.path.join(ecg_data_path, "ptbxl_database.csv")): processor = ECGPreprocessingPipeline(ecg_data_path) signals, targets = processor.process_dataset(limits=100) print(f"[SUCCESS] Signal Tensor Layout: {signals.shape}, Targets Layout: {targets.shape}")