File size: 6,676 Bytes
abf024a c7cc1b5 abf024a c7cc1b5 abf024a c7cc1b5 abf024a c7cc1b5 c341bcb c7cc1b5 77e6989 c341bcb 77e6989 c7cc1b5 c341bcb c7cc1b5 8c90b73 8ea561c c7cc1b5 13ba4b4 c7cc1b5 13ba4b4 c7cc1b5 13ba4b4 c7cc1b5 13ba4b4 c7cc1b5 c341bcb c7cc1b5 c341bcb c7cc1b5 c341bcb c7cc1b5 8ea561c 13ba4b4 c7cc1b5 8c90b73 c7cc1b5 8c90b73 c7cc1b5 8c90b73 c7cc1b5 8c90b73 13ba4b4 c7cc1b5 13ba4b4 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 | 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)} |