| import os |
| import logging |
| from dotenv import load_dotenv |
| from supabase import create_client, Client |
|
|
| load_dotenv() |
| logger = logging.getLogger(__name__) |
|
|
| class Database(): |
| def __init__(self): |
| self.__url = os.environ.get("SUPABASE_URL") or "" |
| self.__key: str = os.environ.get("SUPABASE_KEY") or "" |
| |
| if not self.__url or not self.__key: |
| logger.critical("SUPABASE_URL or SUPABASE_KEY is missing in environment") |
| raise ValueError("SUPABASE_URL or SUPABASE_KEY is missing in environment") |
| |
| try: |
| self.__supabase: Client = create_client(self.__url, self.__key) |
| self.__last_processed_ids = {} |
| logger.info("Database connection established successfully") |
| except Exception as e: |
| logger.critical(f"Failed to create Supabase client: {str(e)}") |
| raise |
| |
| def mark_as_processed(self, machine_id, udi): |
| if machine_id not in self.__last_processed_ids: |
| self.__last_processed_ids[machine_id] = [] |
| self.__last_processed_ids[machine_id].append(udi) |
| logger.info(f"Marked UDI {udi} as processed") |
| |
| def reset_last_processed_id(self, machine_id): |
| if machine_id in self.__last_processed_ids: |
| self.__last_processed_ids[machine_id] = [] |
| logger.info("Reset last_processed_id - data can be reprocessed") |
| |
| def get_all_machine_id(self): |
| try: |
| response = ( |
| self.__supabase.table("machines") |
| .select("id", "name") |
| .execute() |
| ) |
| return response |
| except Exception as e: |
| return e |
| |
| def get_sensor_readings(self, limit, machine_id): |
| try: |
| query = self.__supabase.table("sensor_readings").select("*") |
| |
| if machine_id: |
| query = query.eq("machine_id", machine_id) |
| |
| response = ( |
| query |
| .order("created_at", desc=True) |
| .limit(limit) |
| .execute() |
| ) |
| |
| if not response.data: |
| logger.warning("No sensor readings found in database") |
| return None |
| |
| if limit == 1: |
| current_reading = response.data[0] |
| current_udi = current_reading.get("udi") |
| |
| machine_processed_udis = self.__last_processed_ids.get(machine_id, []) |
| if current_udi in machine_processed_udis: |
| logger.warning(f"Duplicate data detected for Machine {machine_id} - UDI {current_udi} already processed") |
| return { |
| "success": False, |
| "message": "Data already predicted", |
| } |
| |
| logger.debug(f"Retrieved single sensor reading - UDI: {current_udi}, Machine: {machine_id}") |
| return current_reading |
| |
| else: |
| |
| logger.debug(f"Retrieved {len(response.data)} sensor readings for time-series prediction") |
| history_data = response.data[::-1] |
| |
| return history_data |
| |
| except Exception as e: |
| logger.error(f"Failed to retrieve sensor readings: {str(e)}") |
| return None |
| |
| def update_or_create_new_predictions(self, machine_id, timestamp, risk_score, failure_predicted, failure_type, predicted_failure_time, confidence): |
| try: |
| check_database = ( |
| self.__supabase.table("prediction_results") |
| .select('machine_id', 'failure_predicted') |
| .eq("machine_id", machine_id) |
| .execute() |
| ) |
| |
| if check_database.data and len(check_database.data) > 0: |
| if check_database.data[0].get("failure_predicted") == False and failure_predicted == True: |
| response = ( |
| self.__supabase.table("prediction_results") |
| .update({ |
| "machine_id": machine_id, |
| "timestamp": timestamp, |
| "risk_score": risk_score, |
| "failure_predicted": failure_predicted, |
| "failure_type": failure_type, |
| "predicted_failure_time": predicted_failure_time, |
| "confidence": confidence |
| }) |
| .eq("machine_id", machine_id) |
| .execute() |
| ) |
| elif check_database.data[0].get("failure_predicted") == True and failure_predicted == False: |
| return {"success": False, "message": "Cannot updated already predicted data until it is resolved"} |
| else: |
| response = ( |
| self.__supabase.table("prediction_results") |
| .update({ |
| "machine_id": machine_id, |
| "timestamp": timestamp, |
| "risk_score": risk_score, |
| "failure_predicted": failure_predicted, |
| "failure_type": failure_type, |
| "predicted_failure_time": predicted_failure_time, |
| "confidence": confidence |
| }) |
| .eq("machine_id", machine_id) |
| .execute() |
| ) |
| elif not check_database.data: |
| response = ( |
| self.__supabase.table("prediction_results") |
| .insert({ |
| "machine_id": machine_id, |
| "timestamp": timestamp, |
| "risk_score": risk_score, |
| "failure_predicted": failure_predicted, |
| "failure_type": failure_type, |
| "predicted_failure_time": predicted_failure_time, |
| "confidence": confidence |
| }) |
| .execute() |
| ) |
| logger.info(f"Prediction saved - Machine: {machine_id}, Failure: {failure_predicted}, Risk: {risk_score:.2f}") |
| return {"success": True, "data": response} |
| except Exception as e: |
| logger.error(f"Failed to save prediction to database: {str(e)}") |
| return {"success": False, "error": str(e)} |