# 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 LOCAL_TZ = ZoneInfo("America/Toronto") # ------------------------------------------------- # Conversion d'un timestamp Decimal en datetime UTC # ------------------------------------------------- def ts_to_utc(ts): ts = Decimal(str(ts)) seconds = int(ts) micro = int((ts - Decimal(seconds)) * Decimal(1_000_000)) return datetime.fromtimestamp(seconds, tz=timezone.utc).replace(microsecond=micro) # ------------------------------------------------- # 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: continue 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" # ------------------------------------------------- # 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 ) if status == 0: 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 ) continue # ------------------------------------------------- # 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 last_dt = get_last_valid_lecture_time(cur, detail_table, tp_id, dossard) if last_dt and interval_sec > 0: delta = (dt_utc - last_dt).total_seconds() if delta < interval_sec: status = 0 reason_code = f"SKIP_INTERVAL {delta:.1f}s < {interval_sec}s" reason_ref_id = det_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 ) continue # ------------------------------------------------- # Calcul tour (placeholder) # ------------------------------------------------- tour = 1 tempstour = Decimal("0") tempstour_str = "0" 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=tour, tempstour=tempstour, tempstour_str=tempstour_str ) if status == 1: log_info(f"[OK] id={rcourse_id} tag={tag} dossard={dossard} tp={tp_name} prio={prio}") # ------------------------------------------------- # Mise à jour du pointeur # ------------------------------------------------- update_last_processed_read_id(cur, cnx, engine_table, max_id) # ------------------------------------------------- # Dernière lecture valide # ------------------------------------------------- def get_last_valid_lecture_time(cur, detail_table, tp_id, dossard): cur.execute(f""" SELECT timestamp_read FROM `{detail_table}` WHERE tp_id = %s AND cla_info_dossard = %s AND status = 1 ORDER BY detail_id DESC LIMIT 1 """, (tp_id, dossard)) row = cur.fetchone() if not row: return None ts_decimal = Decimal(str(row["timestamp_read"])) return ts_to_utc(ts_decimal)