# process_course.py from engine.engine_components.tp_config import load_tp_config_for_con, resolve_tp_for_lecture from engine.engine_components.course_config import load_course_config from engine.engine_components.engine_state import get_last_processed_read_id, update_last_processed_read_id from engine.engine_components.lectures import get_new_lectures from engine.engine_components.classement import get_classement_by_tag from engine.engine_components.detail_writer import insert_detail_record from engine.logger import log_info, log_debug from datetime import datetime, timezone from zoneinfo import ZoneInfo from decimal import Decimal from engine.utils import format_temps_decimal from engine.engine_components.regles import resolve_priority from engine.engine_components.regles import check_interval_rule from engine.engine_components.regles import get_last_valid_lecture_time from engine.engine_components.regles import ts_to_utc from engine.engine_components.regles import update_detail_status_zero LOCAL_TZ = ZoneInfo("America/Toronto") # ------------------------------------------------- # Normalisation du temps de course (start/end) # ------------------------------------------------- def _normalize_course_time(dt): if dt is None: return None, None dt_utc = dt.replace(tzinfo=timezone.utc) dt_local = dt_utc return dt_local, dt_utc # ------------------------------------------------- # Moteur principal : traitement d'une course # ------------------------------------------------- def process_course(cur, cnx, con_id: int): tp_config = load_tp_config_for_con(cur, con_id) course_cfg = load_course_config(cur, con_id) raw_start = course_cfg["start_time"] raw_end = course_cfg["end_time"] start_local, start_utc = _normalize_course_time(raw_start) end_local, end_utc = _normalize_course_time(raw_end) engine_table = f"resultv2_{con_id}_lecture" classement_table = f"resultv2_{con_id}_classement" # ------------------------------------------------- # Pause # ------------------------------------------------- cur.execute(f"SELECT process_status FROM `{engine_table}` LIMIT 1") row = cur.fetchone() status = (row["process_status"] or "").upper() if status == "PAUSE_REQUEST": cur.execute(f"UPDATE `{engine_table}` SET process_status = 'PAUSE'") cnx.commit() log_info(f"[PAUSE] con_id={con_id} — pause appliquée avant traitement.") return if status == "PAUSE": log_info(f"[PAUSE] con_id={con_id} — course déjà en pause.") return # ------------------------------------------------- # Lire dernier ID traité # ------------------------------------------------- last_processed = get_last_processed_read_id(cur, engine_table) if last_processed is None: last_processed = 0 log_info(f"[STATE] con_id={con_id} last_processed={last_processed}") # ------------------------------------------------- # Récupère les lectures non traitées # ------------------------------------------------- lectures = get_new_lectures(cur, last_processed) log_info(f"[NEW LECTURES] con_id={con_id} count={len(lectures)}") if not lectures: return max_id = last_processed # ------------------------------------------------- # Boucle sur chaque lecture brute # ------------------------------------------------- for lec in lectures: rcourse_id = int(lec["rcourse_id"]) chr_loc = lec["chrinfo_location"] timestamp = Decimal(lec["chrinfo_time"]) max_id = max(max_id, rcourse_id) dt_utc = ts_to_utc(timestamp) dt_local = dt_utc.astimezone(LOCAL_TZ) log_debug(f"[READ] id={rcourse_id} utc={dt_utc} loc={chr_loc}") # ------------------------------------------------- # Validation start/end # ------------------------------------------------- if start_utc and dt_utc < start_utc: log_debug(f"[REJECT START] id={rcourse_id} dt_utc={dt_utc} < start_utc={start_utc}") continue if end_utc and dt_utc > end_utc: log_debug(f"[REJECT END] id={rcourse_id} dt_utc={dt_utc} > end_utc={end_utc}") continue status = 1 reason_code = None reason_ref_id = None # ------------------------------------------------- # Vérifier BIB dans classement # ------------------------------------------------- tag = lec["chrinfo_tag"] coureur = get_classement_by_tag(cur, classement_table, tag) if not coureur: dossard = "???" status = 0 reason_code = "UNKNOWN_BIB" reason_ref_id = None else: dossard = coureur["cla_info_dossard"] # ------------------------------------------------- # Résoudre timing point # ------------------------------------------------- tp_info = resolve_tp_for_lecture(tp_config, con_id, chr_loc) if not tp_info: continue tp_id = tp_info["tp_id"] tp_name = tp_info["tp_name"] prio = tp_info["priority"] controller = tp_info["controller_name"] det_id = tp_info["det_id"] detail_table = f"resultv2_{con_id}_detail" # ------------------------------------------------- # 🔥 INSÉRER IMMÉDIATEMENT DANS DETAILS 🔥 # ------------------------------------------------- detail_id=insert_detail_record( cur, cnx, con_id, dossard, lec, tp_id, tp_name, chr_loc, prio, controller, dt_utc, status=status, reason_code=reason_code, reason_ref_id=reason_ref_id, tour=None, tempstour=None, tempstour_str=None ) # 👉 Après ceci, TU CONTINUES TA LOGIQUE (priorité, intervalle, tours) # mais la donnée brute est déjà sauvée. # on traite pas les dossard ??? if dossard != "???": # ------------------------------------------------- # Règle de priorité # ------------------------------------------------- #status, reason_code, reason_ref_id, forced_tour = resolve_priority( # cur, cnx, detail_table, tp_id, dossard, timestamp, prio, det_id, tp_config #) # ------------------------------------------------- # Règle d'intervalle # ------------------------------------------------- interval_cfg = ( tp_config[tp_id]["config"] .get("config_tp_interval", {}) .get("clef_temp_minimum_read_entre", []) ) interval_sec = 0 if interval_cfg: try: interval_sec = int(interval_cfg[0]["value1"]) except: interval_sec = 0 # Appel de la règle d’intervalle status_interval, rc_interval, rr_interval = check_interval_rule( cur, detail_table, tp_id, dossard, timestamp, interval_sec, det_id ) # Si la règle d’intervalle rejette if status_interval == 0: update_detail_status_zero( cur, cnx, detail_table, detail_id, # <-- l'ID que tu as obtenu juste après insert_detail_record rc_interval, # reason_code retourné par la règle rr_interval # reason_ref_id retourné par la règle ) # Et on arrête le traitement de cette lecture continue # ------------------------------------------------- # Calcul du tour + temps de tour # ------------------------------------------------- # Appel de la règle tour status_tour = check_tour_rule( cur, detail_table, tp_id, dossard, timestamp, det_id ) # Si la règle d’tour rejette if status_interval == 0: update_detail_status_zero( cur, cnx, detail_table, detail_id, # <-- l'ID que tu as obtenu juste après insert_detail_record rc_interval, # reason_code retourné par la règle rr_interval # reason_ref_id retourné par la règle ) # Et on arrête le traitement de cette lecture continue #tour, prev_ts = compute_tour(cur, detail_table, tp_id, dossard) #if forced_tour is not None: # tour = forced_tour #if prev_ts is None: # tempstour = Decimal("0") #else: # tempstour = timestamp - prev_ts #tempstour_str = format_temps_decimal(tempstour) if status == 1: log_info(f"[OK] id={rcourse_id} tag={tag} dossard={dossard} tp={tp_name} prio={prio}") else: log_info(f"[dossard ???] id={rcourse_id} ") # ------------------------------------------------- # Mise à jour du pointeur # ------------------------------------------------- update_last_processed_read_id(cur, cnx, engine_table, max_id)