import json
import uuid
import paho.mqtt.client as mqtt
import ssl
import sys
import os
import io
import time
import logging
from datetime import datetime
from django.utils import timezone
from django.db import connection
from django.db.models import Q

# Импорт моделей Django
from .models import Ttnstate, Report, Product, Reportcurrentloop, Ttn, Logs

# Настройка логирования
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)

sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding='utf-8')

# Настройки MQTT
MQTT_BROKER = "localhost"
MQTT_PORT = 8883
MQTT_USERNAME = 'fomin_a'
MQTT_PASSWORD = 'Hs6#2vG#%8bxsKZf4'
MQTT_CLIENT_ID = str(uuid.uuid4())

# Используем маски (#), чтобы не перечислять сотни топиков вручную
MQTT_TOPICS_SUB = [
    ("momot_beton/report/1", 0), 
    ("momot_beton/report/2", 0), 
    ("momot_beton/report/4", 0),
    ("momot_beton/report/8", 0),
    ("momot_beton/report/16", 0),
    ("momot_beton/report/32", 0),
    ("momot_beton/report/64", 0),
    ("momot_beton/report/128", 0),
    ("momot_beton/report/256", 0),
    ("momot_beton/report/512", 0),
    ("momot_beton/report/1024", 0),
    ("momot_beton/production/1", 0),
    ("momot_beton/production/2", 0), 
    ("momot_beton/production/4", 0),
    ("momot_beton/production/8", 0),
    ("momot_beton/production/16", 0),
    ("momot_beton/production/32", 0),
    ("momot_beton/production/64", 0),
    ("momot_beton/production/128", 0),
    ("momot_beton/production/256", 0),
    ("momot_beton/production/512", 0),
    ("momot_beton/production/1024", 0),
    ("momot_beton/reference/1/confirm", 0),
    ("momot_beton/reference/2/confirm", 0), 
    ("momot_beton/reference/4/confirm", 0),
    ("momot_beton/reference/8/confirm", 0),
    ("momot_beton/reference/16/confirm", 0),
    ("momot_beton/reference/32/confirm", 0),
    ("momot_beton/reference/64/confirm", 0),
    ("momot_beton/reference/128/confirm", 0),
    ("momot_beton/reference/256/confirm", 0),
    ("momot_beton/reference/512/confirm", 0),
    ("momot_beton/reference/1024/confirm", 0),
]
MQTT_QOS = 0

# Dictionary to store objects by id
objects_by_id = {}

def ensure_connection():
    if connection.connection and not connection.is_usable():
        connection.close()

def save_log_event(event_type, message, topic=None, data=None):
    try:
        # Формируем детальное сообщение
        log_message = message
        if topic:
            log_message = f"[{topic}] {message}"
        if data:
            # Добавляем первые 200 символов данных для контекста
            data_str = json.dumps(data, ensure_ascii=False)[:200]
            log_message = f"{log_message} | Data: {data_str}"
        
        # Создаем запись в логах
        Logs.objects.create(
            event=event_type,
            date=timezone.now(),
            message=log_message
        )
        logger.info(f"Log saved: {event_type} - {log_message[:100]}")
    except Exception as e:
        logger.error(f"Failed to save log: {e}")

def update_status_for_other_topics(obj_id, topic_with_st_4, client):
    if obj_id in objects_by_id:
        for topic, status, payload in objects_by_id[obj_id]:
            if topic != topic_with_st_4 and status != 4:
                payload['st'] = 6
                client.publish(topic, json.dumps(payload), qos=0)
                logger.info(f"Updated id={obj_id} in topic {topic}, set st=6")

def get_type_code_from_component_code(code_component):
    if code_component is None:
        return "0000"
    
    try:
        code_str = str(code_component)
        first_digit = code_str[0]
        
        if first_digit == '1':
            return "1000"
        elif first_digit == '2':
            return "2000"
        elif first_digit == '3':
            return "3000"
        elif first_digit == '4':
            return "4000"
        elif first_digit == '5':
            return "5000"
        elif first_digit == '6':
            return "6000"
        elif first_digit == '7':
            return "7000"
        elif first_digit == '8':
            return "8000"
        elif first_digit == '9':
            return "9000"
        else:
            logger.warning(f"Unexpected first digit '{first_digit}' in code_component: {code_component}, using default '0000'")
            return "0000"
    except Exception as e:
        logger.error(f"Error determining typeCode for code_component {code_component}: {e}")
        return "0000"

def save_components_to_database(data, topic):
    try:
        component_list = data.get("component", [])
        if not component_list:
            logger.info(f"No components found in report from {topic}")
            return False
        
        saved_count = 0
        updated_count = 0
        with connection.cursor() as cursor:
            for component in component_list:
                code_component = component.get("code_component")
                name = component.get("name")
                
                if code_component is None or name is None:
                    logger.warning(f"Skipping component with missing data: {component}")
                    continue
                
                if name:
                    name = name.replace('\\/', '/')
                
                type_code = get_type_code_from_component_code(code_component)
                
                cursor.execute("""
                    SELECT id, name, typeCode FROM comp 
                    WHERE code = %s
                """, [code_component])
                
                existing = cursor.fetchone()
                
                if not existing:
                    cursor.execute("""
                        INSERT INTO comp (name, code, typeCode) 
                        VALUES (%s, %s, %s)
                    """, [name, code_component, type_code])
                    
                    saved_count += 1
                    logger.info(f"Saved new component: name='{name}', code={code_component}, typeCode={type_code}")
                else:
                    existing_id, existing_name, existing_type_code = existing
                    should_update = False
                    update_fields = []
                    
                    if existing_name != name:
                        update_fields.append(f"name = '{name}'")
                        should_update = True
                    
                    if existing_type_code != type_code:
                        update_fields.append(f"typeCode = '{type_code}'")
                        should_update = True
                    
                    if should_update:
                        update_query = f"UPDATE comp SET {', '.join(update_fields)} WHERE id = {existing_id}"
                        cursor.execute(update_query)
                        updated_count += 1
                        logger.info(f"Updated component: id={existing_id}, code={code_component}, "
                                  f"old_name='{existing_name}', new_name='{name}', "
                                  f"old_typeCode='{existing_type_code}', new_typeCode='{type_code}'")
        
        if saved_count > 0 or updated_count > 0:
            logger.info(f"Components processed from {topic}: {saved_count} new, {updated_count} updated")
        else:
            logger.info(f"No new components to save from {topic}")
        return True
        
    except Exception as e:
        error_msg = f"Error saving components to database: {e}"
        logger.error(error_msg)
        save_log_event('error', error_msg, topic, data)
        return False

def save_report_to_database(data, topic):
    try:
        if "current_loop" not in data:
            msg = "Skipping report without current_loop field"
            logger.info(f"{msg} from {topic}")
            save_log_event('skipped', msg, topic, data)
            return False
        
        if "product" not in data or not data["product"]:
            msg = "Skipping report without product data"
            logger.info(f"{msg} from {topic}")
            save_log_event('skipped', msg, topic, data)
            return False
        
        if "component" in data and data["component"]:
            save_components_to_database(data, topic)
        
        saved_count = 0
        skipped_count = 0
        
        # Получаем pwr_loop из данных
        pwr_loop_list = data.get("pwr_loop", [])
        
        for product in data["product"]:
            ind_product = product.get("ind_product")
            if ind_product is None:
                continue
            
            all_loops = [
                loop for loop in data.get("current_loop", [])
                if loop.get("ind_product") == ind_product
            ]
            
            all_weight_manual = [
                weight for weight in data.get("weight_manual", [])
                if weight.get("ind_product") == ind_product
            ]
            
            # Создаем копию данных для сохранения в JSON
            single_report_data = {
                "bsu": data.get("bsu"),
                "st": data.get("st"),
                "datetime": data.get("datetime"),
                "product": [product],
                "current_loop": all_loops,
                "weight_manual": all_weight_manual
            }
            
            # Добавляем pwr_loop, если он есть в исходных данных
            if "pwr_loop" in data:
                single_report_data["pwr_loop"] = data.get("pwr_loop")
            
            if "component" in data:
                single_report_data["component"] = data.get("component")
            
            if "update_product" in data:
                single_report_data["update_product"] = data.get("update_product")
            
            json_str = json.dumps(single_report_data, ensure_ascii=False)
            
            if Report.objects.filter(json=json_str).exists():
                msg = f"Report with same JSON already exists for ind_product={ind_product}, skipping"
                logger.info(msg)
                save_log_event('duplicate', msg, topic, {'ind_product': ind_product})
                skipped_count += 1
                continue
            
            Report.objects.create(json=json_str)
            saved_count += 1
            logger.info(f"Successfully saved Report for ind_product={ind_product} with {len(all_loops)} loops and {len(all_weight_manual)} weight_manual records")
        
        if saved_count > 0:
            logger.info(f"Successfully saved {saved_count} Report records from {topic}, skipped {skipped_count} duplicates")
        else:
            msg = f"No new Report records saved from {topic}, skipped {skipped_count} duplicates"
            logger.info(msg)
            save_log_event('info', msg, topic)
        
        return True
        
    except Exception as e:
        error_msg = f"Error saving report to database: {e}"
        logger.error(error_msg)
        save_log_event('error', error_msg, topic, data)
        return False

def process_reference_confirm(data, topic):
    try:
        topic_parts = topic.split('/')
        bsu_value = None
        
        for part in topic_parts:
            if part.isdigit():
                bsu_value = int(part)
                break
        
        if bsu_value is None:
            error_msg = f"Cannot extract BSU number from topic: {topic}"
            logger.error(error_msg)
            save_log_event('error', error_msg, topic, data)
            return False
        
        st_value = data.get("st")
        
        datetime_value = None
        if "recipe" in data and "datetime" in data["recipe"]:
            datetime_obj = data["recipe"]["datetime"]
            if isinstance(datetime_obj, dict):
                datetime_value = datetime_obj.get("datetime")
            else:
                datetime_value = datetime_obj
        
        if st_value is None:
            msg = f"Missing 'st' field in reference confirm message"
            logger.warning(msg)
            save_log_event('warning', msg, topic, data)
            return False
        
        if datetime_value is None:
            msg = f"Missing 'datetime' field in reference confirm message"
            logger.warning(msg)
            save_log_event('warning', msg, topic, data)
            return False
        
        logger.info(f"Processing reference confirm: topic={topic}, bsu={bsu_value}, st={st_value}, datetime={datetime_value}")
        
        with connection.cursor() as cursor:
            cursor.execute("""
                SELECT id, name, code, options, entityName, isPause 
                FROM mainState 
                WHERE code = %s
            """, [st_value])
            
            main_state_record = cursor.fetchone()
            
            if main_state_record:
                logger.info(f"Found mainState record with code={st_value}")
                
                cursor.execute("""
                    UPDATE recipeState 
                    SET state = %s, date = %s 
                    WHERE codeBsu = %s
                """, [st_value, datetime_value, bsu_value])
                
                if cursor.rowcount > 0:
                    logger.info(f"Successfully updated recipeState: codeBsu={bsu_value}, state={st_value}, date={datetime_value}")
                else:
                    msg = f"recipeState with codeBsu={bsu_value} does not exist - no rows updated"
                    logger.warning(msg)
                    save_log_event('warning', msg, topic, {'bsu': bsu_value, 'state': st_value})
                
                return True
            else:
                msg = f"No mainState record found with code={st_value}"
                logger.warning(msg)
                save_log_event('warning', msg, topic, {'state': st_value})
                return False
                
    except Exception as e:
        error_msg = f"Error processing reference confirm message: {e}"
        logger.error(error_msg)
        save_log_event('error', error_msg, topic, data)
        return False

def parse_iso_datetime(datetime_str):
    if not datetime_str:
        return None
    
    try:
        if datetime_str.endswith('Z'):
            datetime_str = datetime_str[:-1]
        
        if '.' in datetime_str:
            dt = datetime.strptime(datetime_str, "%Y-%m-%dT%H:%M:%S.%f")
        else:
            dt = datetime.strptime(datetime_str, "%Y-%m-%dT%H:%M:%S")
        
        return dt
    except Exception as e:
        logger.error(f"Error parsing datetime '{datetime_str}': {e}")
        return None

def check_product_duplicate(time_start, time_end, id_plant=None):
    try:
        if isinstance(time_start, str):
            start_dt = parse_iso_datetime(time_start)
        else:
            start_dt = time_start
        
        if isinstance(time_end, str):
            end_dt = parse_iso_datetime(time_end)
        else:
            end_dt = time_end
        
        if not start_dt or not end_dt:
            logger.warning(f"Cannot parse dates: start={time_start}, end={time_end}")
            return False
        
        start_dt_second = start_dt.replace(microsecond=0)
        end_dt_second = end_dt.replace(microsecond=0)
        
        query = Q(timeStart=start_dt_second, timeEnd=end_dt_second)
        
        if id_plant is not None:
            query &= Q(idPlant=id_plant)
        
        exists = Product.objects.filter(query).exists()
        
        if exists:
            logger.warning(f"Duplicate found: timeStart={start_dt_second}, timeEnd={end_dt_second}, idPlant={id_plant}")
        
        return exists
        
    except Exception as e:
        logger.error(f"Error checking product duplicate: {e}")
        return False

def save_message(client, topic, payload_str):
    ensure_connection()
    
    try:
        clean_payload = payload_str.strip()
        if not clean_payload:
            return None

        if "}{" in clean_payload:
            logger.warning(f"Extra data detected in {topic}. Taking first object.")
            clean_payload = clean_payload.split("}{")[0] + "}"
        
        data = json.loads(clean_payload)
    except Exception as e:
        error_msg = f"JSON Parse Error: {e}"
        logger.error(f"{error_msg} in {topic}")
        save_log_event('error', error_msg, topic, {'payload': payload_str[:200]})
        return None

    if "reference" in topic and "confirm" in topic:
        process_reference_confirm(data, topic)
        return data

    source_bsu = data.get("bsu")
    msg_type = None
    topic_bsu = topic.split('/')[-1] if topic.split('/')[-1].isdigit() else None
    
    if "ttn" in topic: msg_type = "ttn"
    elif "production" in topic: msg_type = "production"
    elif "report" in topic: msg_type = "report"

    code = data.get("st")
    logger.info(f"Msg: {msg_type} | Topic BSU: {topic_bsu} | Source BSU: {source_bsu} | st: {code}")

    if msg_type in ["ttn", "production"]:
        raw_date = data.get("datetime")
        processed_date = None
        if raw_date and isinstance(raw_date, str):
            try:
                dt = timezone.datetime.strptime(raw_date, "%Y-%m-%d %H:%M:%S")
                processed_date = timezone.make_aware(dt, timezone.get_current_timezone()).replace(tzinfo=None)
            except ValueError:
                error_msg = f"Date error: {raw_date}"
                logger.error(error_msg)
                save_log_event('error', error_msg, topic, {'date': raw_date})

        ttn_list = data.get("ttn", [])
        for ttn_item in ttn_list:
            idTtn_raw = ttn_item.get("ind_ttn")
            if idTtn_raw is None: 
                continue
            
            try:
                idTtn = int(str(idTtn_raw).strip())
                
                # Пропускаем idTtn = 0 (реальные данные с сервера)
                if idTtn == 0:
                    logger.info(f"Skipping Ttnstate for idTtn=0 (real data from server) from {topic}")
                    # Не логируем в logs, чтобы не засорять
                    continue
                
                # Проверка на дубликат
                if processed_date:
                    if not Ttnstate.objects.filter(idTtn=idTtn, date=processed_date).exists():
                        Ttnstate.objects.create(idTtn=idTtn, state=code, date=processed_date, json=json.dumps(data))
                        logger.info(f"Saved Ttnstate: {idTtn} (st {code})")
                    else:
                        # Логируем дубликат только для реальных idTtn
                        logger.info(f"Ttnstate already exists for idTtn={idTtn}")
                else:
                    if not Ttnstate.objects.filter(idTtn=idTtn).exists():
                        Ttnstate.objects.create(idTtn=idTtn, state=code, date=None, json=json.dumps(data))
                        logger.info(f"Saved Ttnstate: {idTtn} (st {code})")
                
                # Обновляем основной Ttn
                Ttn.objects.filter(id=idTtn).update(state=code)
            except Exception as e:
                logger.error(f"DB Error for Ttn {idTtn_raw}: {e}")

        if code == 14:
            bsu = data.get("bsu", None)
            if bsu:
                new_payload = data.copy()
                new_payload['st'] = 15
                new_topic = f"momot_beton/ttn/{bsu}"
                result = client.publish(new_topic, json.dumps(new_payload), qos=MQTT_QOS, retain=False)
                status = result.rc
        
                if status == 0:
                    logger.info(f"Message with st=15 published to {new_topic}")
                else:
                    error_msg = f"Failed to send message with st=15 to {new_topic}, result code: {status}"
                    logger.error(error_msg)
                    save_log_event('error', error_msg, topic, {'bsu': bsu, 'status': status})

        if code == 13:
            bsu = data.get("bsu", None)
            if bsu:
                new_payload = data.copy()
                new_payload['st'] = 13
                new_topic = f"momot_beton/ttn/{bsu}"
                result = client.publish(new_topic, json.dumps(new_payload), qos=MQTT_QOS, retain=False)
                status = result.rc
        
                if status == 0:
                    logger.info(f"Message with st=13 published to {new_topic}")
                else:
                    error_msg = f"Failed to send message with st=13 to {new_topic}, result code: {status}"
                    logger.error(error_msg)
                    save_log_event('error', error_msg, topic, {'bsu': bsu, 'status': status})
            
    elif msg_type == "report":
        # Сохраняем Report в БД
        save_report_to_database(data, topic)
        
        product_list = data.get("product", [])
        current_loop_list = data.get("current_loop", [])
        weight_manual_list = data.get("weight_manual", [])
        pwr_loop_list = data.get("pwr_loop", [])  # Получаем pwr_loop из данных
        idPlant = data.get("bsu")

        for prod in product_list:
            indProduct = prod.get("ind_product")
            idTtn = prod.get("ttn")
            time_start = prod.get("time_start")
            time_end = prod.get("time_done")
            
            if Product.objects.filter(indProduct=indProduct, idPlant=idPlant).exists():
                msg = f"Product with indProduct={indProduct} and idPlant={idPlant} already exists, skipping"
                logger.warning(msg)
                save_log_event('duplicate', msg, topic, {'indProduct': indProduct, 'idPlant': idPlant})
                continue
            
            if check_product_duplicate(time_start, time_end, idPlant):
                msg = f"SKIPPING product with duplicate time_start={time_start} and time_end={time_end} for idPlant={idPlant}"
                logger.warning(msg)
                save_log_event('duplicate', msg, topic, {'time_start': time_start, 'time_end': time_end, 'idPlant': idPlant})
                continue
            
            ds_raw = data.get("datetime")
            ds = None
            if ds_raw:
                try:
                    ds = timezone.make_aware(timezone.datetime.strptime(ds_raw, "%Y-%m-%d %H:%M:%S"))
                except Exception as e:
                    error_msg = f"Error parsing datetime: {e}"
                    logger.error(error_msg)
                    save_log_event('error', error_msg, topic, {'datetime': ds_raw})
            
            parsed_time_start = parse_iso_datetime(time_start)
            parsed_time_end = parse_iso_datetime(time_end)
            
            # Получаем значение v_loop как строку (может быть "2000 + 0" или "1000+2000")
            v_loop_value = prod.get("v_loop")
            if v_loop_value is None:
                v_loop_value = ""
            else:
                v_loop_value = str(v_loop_value)
            
            logger.info(f"Attempting to save Product: timeStart={parsed_time_start}, timeEnd={parsed_time_end}, idPlant={idPlant}, vLoop={v_loop_value}")
            
            try:
                product_obj = Product.objects.create(
                    dateStart=ds, 
                    timeEnd=parsed_time_end, 
                    vProduct=prod.get("v_product"),
                    loopNumber=prod.get("num_loop"), 
                    vLoop=v_loop_value,  # Теперь передаем строковое значение
                    driver=prod.get("driver"), 
                    car=prod.get("car"),
                    classRecipe=prod.get("class_recipe") or "", 
                    nameRecipe=prod.get("name_recipe"),
                    recipe=prod.get("recipe"), 
                    idTtn=idTtn, 
                    timeStart=parsed_time_start,
                    num_loop=prod.get("num_loop"), 
                    idPlant=idPlant, 
                    indProduct=indProduct
                )
                
                logger.info(f"Successfully saved Product with id={product_obj.id}, timeStart={parsed_time_start}, timeEnd={parsed_time_end}")
                
                # Сохраняем current_loop
                for loop in current_loop_list:
                    if loop.get("ind_product") == indProduct:
                        # Ищем соответствующее значение pwr для этого цикла
                        power_loop_value = None
                        for pwr in pwr_loop_list:
                            if pwr.get("ind_product") == indProduct and pwr.get("num_loop") == loop.get("num_loop"):
                                power_loop_value = pwr.get("pwr")
                                break
                        
                        # Создаем запись в Reportcurrentloop с полем powerLoop
                        Reportcurrentloop.objects.create(
                            vLoop=loop.get("v_loop"), 
                            loopNumber=loop.get("num_loop"),
                            code=loop.get("code"), 
                            dispencer=loop.get("dispenser"),
                            doisingError=loop.get("err_dosing"),
                            doisingErrorPersent=loop.get("err_dosing_persent"),
                            doisingKorr=loop.get("korr_dosing"),
                            humidityKorr=loop.get("korr_humidity"),
                            weightFactLoop=loop.get("weight_fact_loop"),
                            weightFactM3=loop.get("weight_fact_m3"),
                            weightRecipeLoop=loop.get("weight_recipe_loop"),
                            weightRecipeM3=loop.get("weight_recipe_m3"),
                            powerLoop=power_loop_value,  # Сохраняем значение pwr в поле powerLoop
                            idProduct=product_obj.id, 
                            indProduct=indProduct
                        )
                        logger.info(f"Saved loop: powerLoop={power_loop_value}, vLoop={loop.get('v_loop')}, code={loop.get('code')}")
                
                # Сохраняем weight_manual
                for weight_item in weight_manual_list:
                    if weight_item.get("ind_product") == indProduct:
                        try:
                            with connection.cursor() as cursor:
                                cursor.execute("""
                                    INSERT INTO weightManual 
                                    (indProduct, numLoop, code, weight, dispenser, idProduct) 
                                    VALUES (%s, %s, %s, %s, %s, %s)
                                """, [
                                    weight_item.get("ind_product"),
                                    weight_item.get("num_loop"),
                                    weight_item.get("code"),
                                    weight_item.get("weight"),
                                    weight_item.get("dispenser"),
                                    product_obj.id
                                ])
                            logger.info(f"Saved weight_manual: indProduct={weight_item.get('ind_product')}, "
                                      f"numLoop={weight_item.get('num_loop')}, code={weight_item.get('code')}, "
                                      f"idProduct={product_obj.id}")
                        except Exception as e:
                            error_msg = f"Error saving weight_manual: {e}, data: {weight_item}"
                            logger.error(error_msg)
                            save_log_event('error', error_msg, topic, weight_item)
                            
            except Exception as e:
                error_msg = f"Error creating Product: {e}"
                logger.error(error_msg)
                save_log_event('error', error_msg, topic, {'indProduct': indProduct, 'idPlant': idPlant})

        # Обработка дополнительных weight_manual, которые не попали в основной цикл
        if weight_manual_list:
            for weight_item in weight_manual_list:
                try:
                    product = Product.objects.filter(
                        indProduct=weight_item.get("ind_product"), 
                        idPlant=idPlant
                    ).order_by('-dateStart').first()
                    
                    if product:
                        with connection.cursor() as cursor:
                            cursor.execute("""
                                SELECT COUNT(*) FROM weightManual 
                                WHERE indProduct = %s AND numLoop = %s AND code = %s
                            """, [
                                weight_item.get("ind_product"),
                                weight_item.get("num_loop"),
                                weight_item.get("code")
                            ])
                            count = cursor.fetchone()[0]
                            
                            if count == 0:
                                cursor.execute("""
                                    INSERT INTO weightManual 
                                    (indProduct, numLoop, code, weight, dispenser, idProduct) 
                                    VALUES (%s, %s, %s, %s, %s, %s)
                                """, [
                                    weight_item.get("ind_product"),
                                    weight_item.get("num_loop"),
                                    weight_item.get("code"),
                                    weight_item.get("weight"),
                                    weight_item.get("dispenser"),
                                    product.id
                                ])
                                logger.info(f"Saved weight_manual (additional): indProduct={weight_item.get('ind_product')}, "
                                          f"idProduct={product.id}")
                            else:
                                msg = f"weight_manual duplicate found: indProduct={weight_item.get('ind_product')}, numLoop={weight_item.get('num_loop')}"
                                logger.info(msg)
                                save_log_event('duplicate', msg, topic, weight_item)
                    else:
                        msg = f"Product not found for weight_manual: indProduct={weight_item.get('ind_product')}"
                        logger.warning(msg)
                        save_log_event('warning', msg, topic, {'indProduct': weight_item.get('ind_product')})
                        
                except Exception as e:
                    error_msg = f"Error saving additional weight_manual: {e}, data: {weight_item}"
                    logger.error(error_msg)
                    save_log_event('error', error_msg, topic, weight_item)
    
    return data

def on_connect(client, userdata, flags, rc):
    logger.info(f"Connected to MQTT broker (code {rc})")
    if rc == 0:
        save_log_event('info', f"MQTT Client connected successfully", None, {'rc': rc})
    else:
        save_log_event('error', f"MQTT Client connection failed with code {rc}", None, {'rc': rc})
    client.subscribe(MQTT_TOPICS_SUB)

def on_message(client, userdata, msg):
    try:
        payload = msg.payload.decode('utf-8')
        data = save_message(client, msg.topic, payload)
        
        if data and data.get("id"):
            obj_id = data.get("id")
            obj_st = data.get("st")
            
            if obj_id not in objects_by_id:
                objects_by_id[obj_id] = []
            objects_by_id[obj_id].append((msg.topic, obj_st, data))
            
            if obj_st == 4:
                update_status_for_other_topics(obj_id, msg.topic, client)
    except Exception as e:
        error_msg = f"Error in on_message: {e}"
        logger.error(error_msg)
        save_log_event('error', error_msg, msg.topic, {'payload': msg.payload[:200]})

def start_mqtt_client():
    global client
    lock_file = '/tmp/mqtt_client.lock'
    
    if os.path.exists(lock_file):
        try:
            with open(lock_file, 'r') as f:
                old_pid = int(f.read())
            os.kill(old_pid, 0)
            logger.warning(f"MQTT client already running (PID {old_pid})")
            save_log_event('warning', f"MQTT client already running (PID {old_pid})", None)
            return None
        except (OSError, ValueError):
            os.remove(lock_file)

    with open(lock_file, 'w') as f:
        f.write(str(os.getpid()))

    try:
        client = mqtt.Client(client_id=MQTT_CLIENT_ID)
        client.on_connect = on_connect
        client.on_message = on_message
        client.username_pw_set(MQTT_USERNAME, MQTT_PASSWORD)
        client.tls_set(ca_certs=None, cert_reqs=ssl.CERT_NONE, tls_version=ssl.PROTOCOL_TLS_CLIENT)
        client.tls_insecure_set(True)
        
        logger.info(f"Connecting to {MQTT_BROKER}...")
        client.connect(MQTT_BROKER, MQTT_PORT, 60)
        client.loop_start()
        save_log_event('info', f"MQTT client started successfully", None)
        return client
    except Exception as e:
        error_msg = f"Failed to start: {e}"
        logger.error(error_msg)
        save_log_event('error', error_msg, None)
        if os.path.exists(lock_file): os.remove(lock_file)
        return None

def stop_mqtt_client():
    global client
    if client:
        client.loop_stop()
        client.disconnect()
        save_log_event('info', "MQTT client stopped", None)
    
    lock_file = '/tmp/mqtt_client.lock'
    if os.path.exists(lock_file):
        os.remove(lock_file)

if __name__ == "__main__":
    client = start_mqtt_client()
    if client:
        try:
            while True:
                time.sleep(5)
        except KeyboardInterrupt:
            logger.info("Stopping...")
            stop_mqtt_client()