Predictive-Machine / database.py
justKevv's picture
fix the database.py indent
8ea561c
Raw
History Blame Contribute Delete
6.68 kB
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)}