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: # Multiple readings mode - no duplicate check, used for time-series 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)}