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)}