259 lines
12 KiB
Python
259 lines
12 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_metrics_run = time.time()
|
||
with self.vs.history_lock:
|
||
self.last_gm_index = len(self.vs.history) - 1
|
||
self.gm_buffer = ""
|
||
|
||
def run(self):
|
||
if self.last_gm_index == -1:
|
||
with self.vs.history_lock:
|
||
self.last_gm_index = len(self.vs.history) - 1
|
||
|
||
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[
|
||
'voiceInterval'] * 60:
|
||
self._run_regular_analysis(obsidian_settings)
|
||
self.last_regular_run = time.time()
|
||
|
||
# 3. НОВОЕ: Интервал метрик (Графики + Вкладка 4 LLM)
|
||
if time.time() - self.last_metrics_run > obsidian_settings.get('metricsInterval', 15) * 60:
|
||
# Запускаем питоновские графики
|
||
threading.Thread(target=self.vs.generate_metrics_graphs, daemon=True).start()
|
||
# Запускаем LLM анализ качественных метрик
|
||
self._run_metrics_analysis(obsidian_settings)
|
||
self.last_metrics_run = time.time()
|
||
|
||
time.sleep(3) # каждые 3 секунды проверяем
|
||
except Exception as e:
|
||
logger.critical(f"❌🎤 Voice Heartbeat: Неожиданная ошибка: {e}")
|
||
|
||
def _check_commands(self, obsidian_settings):
|
||
# 1. Забираем новые реплики
|
||
new_text, new_index = self.vs.get_new_gm_text(self.last_gm_index)
|
||
|
||
# Добавляем в буфер (через пробел, чтобы не склеились слова)
|
||
if new_text:
|
||
self.gm_buffer = (self.gm_buffer + " " + new_text).strip()
|
||
self.last_gm_index = new_index
|
||
|
||
if not self.gm_buffer:
|
||
return
|
||
|
||
import re
|
||
start_marker = obsidian_settings['markerStart']
|
||
stop_marker = obsidian_settings['markerStop']
|
||
|
||
# Регулярка для поиска самой короткой подходящей пары (нежадный поиск)
|
||
pattern = f"{re.escape(start_marker)}(.*?){re.escape(stop_marker)}"
|
||
|
||
last_stop_pos = 0
|
||
# finditer позволяет получить позиции вхождений
|
||
matches = list(re.finditer(pattern, self.gm_buffer, re.DOTALL))
|
||
|
||
if matches:
|
||
for m in matches:
|
||
cmd = m.group(1).strip()
|
||
if cmd:
|
||
print(f"🎤 New GM Voice Command detected: {cmd}")
|
||
prompt = f"{(obsidian_settings.get('systemPrompt') or '')}\n\n{(obsidian_settings.get('promptCommands') or '')}"
|
||
self.generator.add_voice_task_to_queue(cmd, prompt)
|
||
|
||
# Запоминаем позицию конца последнего найденного стоп-маркера
|
||
last_stop_pos = m.end()
|
||
|
||
# ОБРЕЗАЕМ БУФЕР: оставляем только то, что идет за последним стоп-маркером
|
||
# (там может быть начало следующей команды "старт команда 3...")
|
||
self.gm_buffer = self.gm_buffer[last_stop_pos:].strip()
|
||
else:
|
||
# Если матчей нет, но буфер стал слишком огромным (забыли стоп-маркер),
|
||
# стоит его ограничить, чтобы не переполнять память, например, последними 2000 симв.
|
||
if len(self.gm_buffer) > 2000:
|
||
# Ищем последний старт-маркер, чтобы не отрезать начало потенциальной команды
|
||
last_start = self.gm_buffer.rfind(start_marker)
|
||
if last_start != -1:
|
||
self.gm_buffer = self.gm_buffer[last_start:]
|
||
else:
|
||
self.gm_buffer = "" # Маркеров старта тоже нет - чистим
|
||
|
||
def _run_regular_analysis(self, obsidian_settings):
|
||
# ТЕПЕРЬ БЕРЕМ ИЗ ПАМЯТИ: хронологический список реплик всех ролей
|
||
context_text = self.vs.get_context_by_limit(obsidian_settings['voiceContextLimit'])
|
||
|
||
if not context_text.strip():
|
||
return
|
||
|
||
print(f"🎤 Запуск регулярного анализа контекста ({len(context_text)} симв.)")
|
||
prompt = f"{(obsidian_settings.get('systemPrompt') or '')}\n\n{(obsidian_settings.get('promptRegular') or '')}"
|
||
self.generator.add_voice_regular_task(context_text, prompt)
|
||
|
||
# Метод для вызова LLM для аналитики
|
||
def _run_metrics_analysis(self, obsidian_settings):
|
||
context_text = self.vs.get_context_by_limit(obsidian_settings.get('voiceContextLimit', 30000))
|
||
if not context_text.strip(): return
|
||
|
||
prompt = obsidian_settings.get('promptMetrics')
|
||
if prompt:
|
||
print("🎤 Запуск LLM анализа качественных метрик сессии...")
|
||
self.generator.add_voice_metrics_task(context_text, prompt) |