214 lines
10 KiB
Python
214 lines
10 KiB
Python
# app/heartbeat_monitor.py
|
||
import threading
|
||
import time
|
||
import requests
|
||
import datetime
|
||
import logging
|
||
# import traceback
|
||
|
||
# Настройка логирования для HeartbeatMonitor
|
||
logger = logging.getLogger(__name__)
|
||
logger.setLevel(logging.INFO)
|
||
# Добавляем обработчик, если его еще нет (чтобы логи выводились в консоль)
|
||
if not logger.handlers:
|
||
handler = logging.StreamHandler()
|
||
formatter = logging.Formatter(
|
||
'%(asctime)s - %(name)s - %(levelname)s - %(message)s')
|
||
handler.setFormatter(formatter)
|
||
logger.addHandler(handler)
|
||
|
||
|
||
class HeartbeatMonitor:
|
||
"""
|
||
Мониторит доступность заданных внешних эндпоинтов путем периодических HTTP-запросов.
|
||
Работает в отдельном фоновом потоке, логируя результаты.
|
||
"""
|
||
|
||
def __init__(self, interval_minutes: int = 5):
|
||
self._interval_seconds = interval_minutes * 60
|
||
self._target_urls = [
|
||
"https://render-service-gsu7.onrender.com/heartbeat",
|
||
# Добавьте сюда другие URL для мониторинга, если они появятся
|
||
]
|
||
self._running = False
|
||
self._thread = None
|
||
self._last_check_results = {}
|
||
logger.info(
|
||
f"HeartbeatMonitor инициализирован с интервалом {interval_minutes} минут."
|
||
)
|
||
logger.info(f"Мониторинг URL: {', '.join(self._target_urls)}")
|
||
|
||
# Вывод call stack
|
||
# stack_trace = "".join(traceback.format_stack())
|
||
# logger.info(f"Call stack при инициализации HeartbeatMonitor:\n{stack_trace}")
|
||
|
||
def _run_heartbeat_loop(self):
|
||
"""
|
||
Основной цикл мониторинга, выполняющийся в отдельном потоке.
|
||
"""
|
||
while self._running:
|
||
logger.info("Начинается выполнение Heartbeat-проверок...")
|
||
current_results = {}
|
||
overall_healthy = True
|
||
|
||
for url in self._target_urls:
|
||
start_time = time.time()
|
||
try:
|
||
response = requests.get(url,
|
||
timeout=10) # Таймаут 10 секунд
|
||
if 200 <= response.status_code < 300:
|
||
status = {
|
||
"healthy": True,
|
||
"status_code": response.status_code,
|
||
"message": "OK"
|
||
}
|
||
logger.info(
|
||
f"✅ Heartbeat {url}: OK (HTTP {response.status_code})"
|
||
)
|
||
else:
|
||
status = {
|
||
"healthy": False,
|
||
"status_code": response.status_code,
|
||
"message": f"Non-2xx status code"
|
||
}
|
||
logger.warning(
|
||
f"⚠️ Heartbeat {url}: Ошибка (HTTP {response.status_code})"
|
||
)
|
||
overall_healthy = False
|
||
except requests.exceptions.RequestException as e:
|
||
status = {
|
||
"healthy": False,
|
||
"status_code": None,
|
||
"message": f"Request failed: {e}"
|
||
}
|
||
logger.error(f"❌ Heartbeat {url}: Ошибка запроса: {e}")
|
||
overall_healthy = False
|
||
except Exception as e:
|
||
status = {
|
||
"healthy": False,
|
||
"status_code": None,
|
||
"message": f"Unexpected error: {e}"
|
||
}
|
||
logger.critical(
|
||
f"🔥 Heartbeat {url}: Неожиданная ошибка: {e}")
|
||
overall_healthy = False
|
||
finally:
|
||
status["timestamp"] = datetime.datetime.now().isoformat()
|
||
status["response_time_ms"] = int(
|
||
(time.time() - start_time) * 1000)
|
||
current_results[url] = status
|
||
|
||
self._last_check_results = current_results
|
||
if overall_healthy:
|
||
logger.info("✔ Все Heartbeat-проверки успешно пройдены.")
|
||
else:
|
||
logger.warning("✖ Обнаружены проблемы в Heartbeat-проверках.")
|
||
|
||
time.sleep(self._interval_seconds)
|
||
|
||
def start(self):
|
||
"""
|
||
Запускает мониторинг в отдельном потоке.
|
||
"""
|
||
if not self._running:
|
||
self._running = True
|
||
# daemon=True позволяет приложению завершиться, даже если этот поток еще работает.
|
||
# Python автоматически завершит daemon-потоки при выходе из основной программы.
|
||
self._thread = threading.Thread(target=self._run_heartbeat_loop,
|
||
daemon=True)
|
||
self._thread.start()
|
||
logger.info("HeartbeatMonitor запущен.")
|
||
else:
|
||
logger.warning("HeartbeatMonitor уже запущен.")
|
||
|
||
def stop(self):
|
||
"""
|
||
Останавливает поток мониторинга.
|
||
"""
|
||
if self._running:
|
||
logger.info("Остановка HeartbeatMonitor...")
|
||
self._running = False
|
||
if self._thread and self._thread.is_alive():
|
||
self._thread.join(
|
||
timeout=5) # Даем потоку 5 секунд на завершение
|
||
if self._thread.is_alive():
|
||
logger.warning(
|
||
"HeartbeatMonitor поток не завершился в течение таймаута."
|
||
)
|
||
logger.info("HeartbeatMonitor остановлен.")
|
||
else:
|
||
logger.warning("HeartbeatMonitor не был запущен.")
|
||
|
||
def get_last_check_results(self):
|
||
"""
|
||
Возвращает результаты последней проверки (для возможного использования внутри приложения).
|
||
"""
|
||
return self._last_check_results
|
||
|
||
|
||
class VoiceHeartbeatMonitor(threading.Thread):
|
||
|
||
def __init__(self, voice_instance, generator, get_obsidian_settings_fn):
|
||
super().__init__(daemon=True)
|
||
self.vs = voice_instance # Тот самый экземпляр из api.py
|
||
self.generator = generator
|
||
self.get_obsidian_settings = get_obsidian_settings_fn # Функция, возвращающая текущий конфиг
|
||
self.last_regular_run = time.time()
|
||
self.last_command_time = ""
|
||
|
||
def run(self):
|
||
while True:
|
||
try:
|
||
if self.vs.is_running:
|
||
obsidian_settings = self.get_obsidian_settings()
|
||
# 1. Проверка команд <start>...<stop>
|
||
self._check_commands(obsidian_settings)
|
||
|
||
# 2. Регулярный запрос
|
||
if time.time() - self.last_regular_run > obsidian_settings['interval'] * 60:
|
||
self._run_regular_analysis(obsidian_settings)
|
||
self.last_regular_run = time.time()
|
||
time.sleep(3) # каждые 3 секунды проверяем
|
||
except Exception as e:
|
||
logger.critical(f"❌🎤 Voice Heartbeat: Неожиданная ошибка: {e}")
|
||
|
||
def _check_commands(self, obsidian_settings):
|
||
# Получаем данные за последние секунды
|
||
raw_text = self.vs.get_last_chars(obsidian_settings['context_limit'])
|
||
|
||
import re
|
||
# Чистим таймстемпы Vosk для LLM (удаляем паттерны [00:00:00])
|
||
clean_text = re.sub(r'\[.*?\]', '', raw_text)
|
||
|
||
pattern = f"{re.escape(obsidian_settings['marker_start'])}(.*?){re.escape(obsidian_settings['marker_stop'])}"
|
||
matches = re.findall(pattern, clean_text, re.DOTALL)
|
||
|
||
for match in matches:
|
||
cmd = match.strip()
|
||
# Проверяем, не обрабатывали ли мы ЭТУ конкретную строку только что
|
||
# Используем множество для хранения хешей обработанных команд
|
||
if not hasattr(self, '_processed_hashes'): self._processed_hashes = set()
|
||
|
||
cmd_hash = hash(cmd)
|
||
if cmd_hash not in self._processed_hashes:
|
||
print(f"🎤 New Voice Command detected: {cmd}")
|
||
|
||
# ОТПРАВЛЯЕМ В LLM через TitleGenerator
|
||
prompt = f"{(obsidian_settings.get('systemPrompt') or '')}\n\n{(obsidian_settings.get('prompt_commands') or '')}"
|
||
self.generator.add_voice_task_to_queue(
|
||
cmd,
|
||
prompt
|
||
)
|
||
|
||
self._processed_hashes.add(cmd_hash)
|
||
# Ограничиваем размер множества, чтобы не росло бесконечно
|
||
if len(self._processed_hashes) > 100: self._processed_hashes.clear() # Спорно, очень спорно
|
||
|
||
def _run_regular_analysis(self, obsidian_settings):
|
||
raw_text = self.vs.get_last_chars(obsidian_settings['voice_context_limit'])
|
||
if not raw_text.strip(): return
|
||
|
||
prompt = f"{(obsidian_settings.get('systemPrompt') or '')}\n\n{(obsidian_settings.get('prompt_regular') or '')}"
|
||
|
||
# Добавим новый тип 'voice_regular' вTitleGenerator аналогично 'voice_command'
|
||
self.generator.add_voice_regular_task(raw_text, prompt) |