summaryrefslogtreecommitdiff
path: root/mqtt_log.py
blob: cac4a0176a4ed83b635c425c118c2c6870723554 (plain)
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
import os
import json
import time
from pathlib import Path
import requests
import paho.mqtt.client as mqtt
from dotenv import load_dotenv

load_dotenv()

def file_print(data_string, data_data):
    """ Обработчик данных от сервера MQTT. Передаем: топик, данные """
    try:
        # Извлекаем имя файла из топика
        if "/" in data_string:
            file_name = data_string[1:data_string.rfind("/")]
            # Получаем расширение или вторую часть
            value_part = data_string[data_string.rfind("/") + 1:]
        else:
            file_name = data_string
            value_part = ""
        
        # Формируем пути
        csv_path = f"{file_name}.csv"
        timestamp = time.strftime('%Y/%m/%d - %H:%M:%S', time.localtime())
        formatted_data = f"/{timestamp}_{data_data}"
        
        # Запись в файл
        fileWrite(os.getenv('FOLDER_ARHIVE_NAME'), csv_path, formatted_data)
        
        # Обновление кэша JSON
        update_json_cache(file_name, f"{data_data}_{time.strftime('%H:%M:%S_%d/%m/%Y', time.localtime())}")
        
    except Exception as e:
        print(f"Ошибка в file_print: {e}")

def update_json_cache(file_name, value):
    """Обновление JSON кэша"""
    path = os.getenv('FILE_CACHE_JSON_NAME')
    data_json = {file_name: value}
    
    try:
        if Path(path).exists():
            with open(path, 'r', encoding='utf-8') as json_file:
                data = json.load(json_file)
                data.update(data_json)
        else:
            data = data_json
            
        with open(path, 'w', encoding='utf-8') as f:
            json.dump(data, f, ensure_ascii=False, indent=2)
    except Exception as e:
        print(f"Ошибка при обновлении JSON: {e}")

def connect_mqtt():
    """ Подключение к серверу MQTT """
    try:
        # Попытка для новых версий
        client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)
    except AttributeError:
        # Для старых версий
        client = mqtt.Client()
    
    client.username_pw_set(os.getenv('USERNAME'), os.getenv('PASSWORD'))
    client.connect(
        os.getenv('SERVER_IP'), 
        int(os.getenv('PORT_MQTT'))
    )
    return client

def subscribe(client: mqtt):
    """ Подписываемся на топики """
    def on_message(client, userdata, msg):
        try:
            file_print(msg.topic, msg.payload.decode())
        except Exception as error:
            print(f"Ошибка в on_message: {error}")

    client.subscribe(os.getenv('TOPIC'))
    client.on_message = on_message

def fileWrite(folderArhiveName, path, myFile):
    """ Запись в файл. Передаем: название папки архива, название файла (топика), данные """
    try:
        # Извлекаем имя папки
        space_pos = myFile.find(' ')
        if space_pos == -1:
            folder_name = folderArhiveName + myFile
        else:
            folder_name = folderArhiveName + myFile[:space_pos]
        
        # Создаем путь
        path_temp = Path.cwd() / folder_name
        path_temp.mkdir(parents=True, exist_ok=True)
        
        # Записываем данные
        file_path = path_temp / path
        with open(file_path, "a", encoding='utf-8') as f:
            # Извлекаем данные после пробела
            data_part = myFile[myFile.rfind(' ') + 1:] if ' ' in myFile else myFile
            f.write(data_part + '\n')
            
    except Exception as e:
        print(f"Ошибка записи файла: {e}")

def run():
    try:
        print("Подключение к MQTT...")
        client = connect_mqtt()
        subscribe(client)
        print("Ожидание сообщений...")
        client.loop_forever()
    except KeyboardInterrupt:
        print("\nПрограмма остановлена")
    except Exception as e:
        print(f"Критическая ошибка: {e}")

if __name__ == "__main__":
    run()