Assistants.lf_cleaner
Ассистент очистки старых трейсов Langfuse.
Периодически удаляет трейсы старше заданного периода через HTTP API Langfuse, чтобы избежать превышения лимитов базы данных.
1""" 2Ассистент очистки старых трейсов Langfuse. 3 4Периодически удаляет трейсы старше заданного периода через 5HTTP API Langfuse, чтобы избежать превышения лимитов базы данных. 6""" 7 8import base64 9import os 10import requests 11 12from datetime import datetime, timedelta 13from loguru import logger 14from langfuse.decorators import observe 15 16from Assistants.assistant import Assistant 17 18 19class LangfuseCleanerAssistant(Assistant): 20 """Удаляет старые трейсы из Langfuse. 21 22 Использует Basic Auth и пагинированный обход трейсов через 23 публичное API Langfuse. Удаляет трейсы старше указанного 24 количества дней. 25 26 Attributes: 27 period: Период хранения трейсов в днях. 28 project_slug: Идентификатор проекта в Langfuse. 29 batch_size: Количество трейсов за один запрос. 30 """ 31 32 def __init__(self, project_slug, batch_size: int = 100, period: int = 1): 33 """Инициализирует ассистента очистки Langfuse. 34 35 Args: 36 project_slug: Идентификатор проекта в Langfuse. 37 batch_size: Количество трейсов за один запрос. 38 period: Период хранения трейсов в днях. 39 """ 40 super().__init__(name="Langfuse Cleaner") 41 self.description = "Cleans and optimizes Langfuse files for better performance." 42 self.period = period # in days 43 self.project_slug = project_slug 44 self.batch_size = batch_size 45 46 def should_delete_trace(self, trace): 47 """Определяет, нужно ли удалить трейс на основе его возраста. 48 49 Args: 50 trace: Объект трейса из Langfuse. 51 52 Returns: 53 bool: True если трейс старше указанного количества дней. 54 """ 55 # Получаем timestamp из трейса 56 trace_timestamp = trace.get("timestamp") # '2025-11-18T16:14:28.624Z' 57 58 if not trace_timestamp: 59 return False 60 61 # Преобразуем строку в datetime объект 62 try: 63 trace_date = datetime.fromisoformat(trace_timestamp.replace("Z", "+00:00")) 64 except ValueError as e: 65 logger.error(f"Ошибка преобразования даты {trace_timestamp}: {e}") 66 return False 67 68 # Вычисляем cutoff дату 69 cutoff_date = datetime.now().replace(tzinfo=trace_date.tzinfo) - timedelta(days=self.period) 70 71 # Сравниваем даты 72 return trace_date < cutoff_date 73 74 def delete_old_traces(self, days_old=30): 75 """Удаляет трейсы старше указанного количества дней через HTTP API Langfuse. 76 77 Проходит по всем страницам трейсов, отфильтровывает по дате 78 и удаляет каждый подходящий трейс. 79 80 Args: 81 days_old: Количество дней для фильтрации (устарел, используется self.period). 82 83 Returns: 84 int: Количество удалённых трейсов. 85 """ 86 87 logger.debug(f'Запуск цикла удаления старых трейсов "{self.name}"...') 88 89 base_url = os.getenv("LANGFUSE_HOST", "http://192.168.0.250:13000") 90 secret_key = os.getenv("LANGFUSE_SECRET_KEY") 91 public_key = os.getenv("LANGFUSE_PUBLIC_KEY") 92 93 credentials = base64.b64encode(f"{public_key}:{secret_key}".encode()).decode() 94 headers = {"Authorization": f"Basic {credentials}", "Content-Type": "application/json"} 95 96 cutoff = (datetime.utcnow() - timedelta(days=self.period)).isoformat() + "Z" 97 98 deleted_count = 0 99 failed_count = 0 100 page = 1 101 102 while True: 103 # Получаем список трейсов 104 url = f"{base_url}/api/public/traces" 105 params = { 106 "projectSlug": self.project_slug, 107 "limit": self.batch_size, 108 "page": page, 109 "toTimestamp": cutoff, 110 } 111 112 response = requests.get(url, headers=headers, params=params) 113 114 if response.status_code != 200: 115 logger.error(f"Ошибка получения трейсов: {response.text}") 116 break 117 118 data = response.json() 119 traces = data.get("data", []) 120 121 if not traces: 122 break 123 124 # Удаляем каждый трейс 125 for trace in traces: 126 if not self.should_delete_trace(trace): 127 continue 128 trace_id = trace["id"] 129 delete_url = f"{base_url}/api/public/traces/{trace_id}" 130 131 delete_response = requests.delete(delete_url, headers=headers) 132 133 if delete_response.status_code in [200, 204]: 134 deleted_count += 1 135 else: 136 failed_count += 1 137 138 page += 1 139 140 if failed_count > 0: 141 logger.warning(f"Не удалось удалить {failed_count} трейсов.") 142 if deleted_count > 0: 143 logger.debug(f"Всего удалено: {deleted_count} трейсов") 144 return deleted_count 145 146 @observe(name=f"Цикл очистки {__qualname__}") 147 def run(self): 148 """Выполняет один цикл очистки старых трейсов Langfuse.""" 149 # Main logic for cleaning and optimizing Langfuse files 150 deleted = self.delete_old_traces(days_old=self.period) 151 logger.debug(f"Цикл {self.name} завершен: удалено {deleted} старых трейсов.")
20class LangfuseCleanerAssistant(Assistant): 21 """Удаляет старые трейсы из Langfuse. 22 23 Использует Basic Auth и пагинированный обход трейсов через 24 публичное API Langfuse. Удаляет трейсы старше указанного 25 количества дней. 26 27 Attributes: 28 period: Период хранения трейсов в днях. 29 project_slug: Идентификатор проекта в Langfuse. 30 batch_size: Количество трейсов за один запрос. 31 """ 32 33 def __init__(self, project_slug, batch_size: int = 100, period: int = 1): 34 """Инициализирует ассистента очистки Langfuse. 35 36 Args: 37 project_slug: Идентификатор проекта в Langfuse. 38 batch_size: Количество трейсов за один запрос. 39 period: Период хранения трейсов в днях. 40 """ 41 super().__init__(name="Langfuse Cleaner") 42 self.description = "Cleans and optimizes Langfuse files for better performance." 43 self.period = period # in days 44 self.project_slug = project_slug 45 self.batch_size = batch_size 46 47 def should_delete_trace(self, trace): 48 """Определяет, нужно ли удалить трейс на основе его возраста. 49 50 Args: 51 trace: Объект трейса из Langfuse. 52 53 Returns: 54 bool: True если трейс старше указанного количества дней. 55 """ 56 # Получаем timestamp из трейса 57 trace_timestamp = trace.get("timestamp") # '2025-11-18T16:14:28.624Z' 58 59 if not trace_timestamp: 60 return False 61 62 # Преобразуем строку в datetime объект 63 try: 64 trace_date = datetime.fromisoformat(trace_timestamp.replace("Z", "+00:00")) 65 except ValueError as e: 66 logger.error(f"Ошибка преобразования даты {trace_timestamp}: {e}") 67 return False 68 69 # Вычисляем cutoff дату 70 cutoff_date = datetime.now().replace(tzinfo=trace_date.tzinfo) - timedelta(days=self.period) 71 72 # Сравниваем даты 73 return trace_date < cutoff_date 74 75 def delete_old_traces(self, days_old=30): 76 """Удаляет трейсы старше указанного количества дней через HTTP API Langfuse. 77 78 Проходит по всем страницам трейсов, отфильтровывает по дате 79 и удаляет каждый подходящий трейс. 80 81 Args: 82 days_old: Количество дней для фильтрации (устарел, используется self.period). 83 84 Returns: 85 int: Количество удалённых трейсов. 86 """ 87 88 logger.debug(f'Запуск цикла удаления старых трейсов "{self.name}"...') 89 90 base_url = os.getenv("LANGFUSE_HOST", "http://192.168.0.250:13000") 91 secret_key = os.getenv("LANGFUSE_SECRET_KEY") 92 public_key = os.getenv("LANGFUSE_PUBLIC_KEY") 93 94 credentials = base64.b64encode(f"{public_key}:{secret_key}".encode()).decode() 95 headers = {"Authorization": f"Basic {credentials}", "Content-Type": "application/json"} 96 97 cutoff = (datetime.utcnow() - timedelta(days=self.period)).isoformat() + "Z" 98 99 deleted_count = 0 100 failed_count = 0 101 page = 1 102 103 while True: 104 # Получаем список трейсов 105 url = f"{base_url}/api/public/traces" 106 params = { 107 "projectSlug": self.project_slug, 108 "limit": self.batch_size, 109 "page": page, 110 "toTimestamp": cutoff, 111 } 112 113 response = requests.get(url, headers=headers, params=params) 114 115 if response.status_code != 200: 116 logger.error(f"Ошибка получения трейсов: {response.text}") 117 break 118 119 data = response.json() 120 traces = data.get("data", []) 121 122 if not traces: 123 break 124 125 # Удаляем каждый трейс 126 for trace in traces: 127 if not self.should_delete_trace(trace): 128 continue 129 trace_id = trace["id"] 130 delete_url = f"{base_url}/api/public/traces/{trace_id}" 131 132 delete_response = requests.delete(delete_url, headers=headers) 133 134 if delete_response.status_code in [200, 204]: 135 deleted_count += 1 136 else: 137 failed_count += 1 138 139 page += 1 140 141 if failed_count > 0: 142 logger.warning(f"Не удалось удалить {failed_count} трейсов.") 143 if deleted_count > 0: 144 logger.debug(f"Всего удалено: {deleted_count} трейсов") 145 return deleted_count 146 147 @observe(name=f"Цикл очистки {__qualname__}") 148 def run(self): 149 """Выполняет один цикл очистки старых трейсов Langfuse.""" 150 # Main logic for cleaning and optimizing Langfuse files 151 deleted = self.delete_old_traces(days_old=self.period) 152 logger.debug(f"Цикл {self.name} завершен: удалено {deleted} старых трейсов.")
Удаляет старые трейсы из Langfuse.
Использует Basic Auth и пагинированный обход трейсов через публичное API Langfuse. Удаляет трейсы старше указанного количества дней.
Attributes: period: Период хранения трейсов в днях. project_slug: Идентификатор проекта в Langfuse. batch_size: Количество трейсов за один запрос.
33 def __init__(self, project_slug, batch_size: int = 100, period: int = 1): 34 """Инициализирует ассистента очистки Langfuse. 35 36 Args: 37 project_slug: Идентификатор проекта в Langfuse. 38 batch_size: Количество трейсов за один запрос. 39 period: Период хранения трейсов в днях. 40 """ 41 super().__init__(name="Langfuse Cleaner") 42 self.description = "Cleans and optimizes Langfuse files for better performance." 43 self.period = period # in days 44 self.project_slug = project_slug 45 self.batch_size = batch_size
Инициализирует ассистента очистки Langfuse.
Args: project_slug: Идентификатор проекта в Langfuse. batch_size: Количество трейсов за один запрос. period: Период хранения трейсов в днях.
47 def should_delete_trace(self, trace): 48 """Определяет, нужно ли удалить трейс на основе его возраста. 49 50 Args: 51 trace: Объект трейса из Langfuse. 52 53 Returns: 54 bool: True если трейс старше указанного количества дней. 55 """ 56 # Получаем timestamp из трейса 57 trace_timestamp = trace.get("timestamp") # '2025-11-18T16:14:28.624Z' 58 59 if not trace_timestamp: 60 return False 61 62 # Преобразуем строку в datetime объект 63 try: 64 trace_date = datetime.fromisoformat(trace_timestamp.replace("Z", "+00:00")) 65 except ValueError as e: 66 logger.error(f"Ошибка преобразования даты {trace_timestamp}: {e}") 67 return False 68 69 # Вычисляем cutoff дату 70 cutoff_date = datetime.now().replace(tzinfo=trace_date.tzinfo) - timedelta(days=self.period) 71 72 # Сравниваем даты 73 return trace_date < cutoff_date
Определяет, нужно ли удалить трейс на основе его возраста.
Args: trace: Объект трейса из Langfuse.
Returns: bool: True если трейс старше указанного количества дней.
75 def delete_old_traces(self, days_old=30): 76 """Удаляет трейсы старше указанного количества дней через HTTP API Langfuse. 77 78 Проходит по всем страницам трейсов, отфильтровывает по дате 79 и удаляет каждый подходящий трейс. 80 81 Args: 82 days_old: Количество дней для фильтрации (устарел, используется self.period). 83 84 Returns: 85 int: Количество удалённых трейсов. 86 """ 87 88 logger.debug(f'Запуск цикла удаления старых трейсов "{self.name}"...') 89 90 base_url = os.getenv("LANGFUSE_HOST", "http://192.168.0.250:13000") 91 secret_key = os.getenv("LANGFUSE_SECRET_KEY") 92 public_key = os.getenv("LANGFUSE_PUBLIC_KEY") 93 94 credentials = base64.b64encode(f"{public_key}:{secret_key}".encode()).decode() 95 headers = {"Authorization": f"Basic {credentials}", "Content-Type": "application/json"} 96 97 cutoff = (datetime.utcnow() - timedelta(days=self.period)).isoformat() + "Z" 98 99 deleted_count = 0 100 failed_count = 0 101 page = 1 102 103 while True: 104 # Получаем список трейсов 105 url = f"{base_url}/api/public/traces" 106 params = { 107 "projectSlug": self.project_slug, 108 "limit": self.batch_size, 109 "page": page, 110 "toTimestamp": cutoff, 111 } 112 113 response = requests.get(url, headers=headers, params=params) 114 115 if response.status_code != 200: 116 logger.error(f"Ошибка получения трейсов: {response.text}") 117 break 118 119 data = response.json() 120 traces = data.get("data", []) 121 122 if not traces: 123 break 124 125 # Удаляем каждый трейс 126 for trace in traces: 127 if not self.should_delete_trace(trace): 128 continue 129 trace_id = trace["id"] 130 delete_url = f"{base_url}/api/public/traces/{trace_id}" 131 132 delete_response = requests.delete(delete_url, headers=headers) 133 134 if delete_response.status_code in [200, 204]: 135 deleted_count += 1 136 else: 137 failed_count += 1 138 139 page += 1 140 141 if failed_count > 0: 142 logger.warning(f"Не удалось удалить {failed_count} трейсов.") 143 if deleted_count > 0: 144 logger.debug(f"Всего удалено: {deleted_count} трейсов") 145 return deleted_count
Удаляет трейсы старше указанного количества дней через HTTP API Langfuse.
Проходит по всем страницам трейсов, отфильтровывает по дате и удаляет каждый подходящий трейс.
Args: days_old: Количество дней для фильтрации (устарел, используется self.period).
Returns: int: Количество удалённых трейсов.
147 @observe(name=f"Цикл очистки {__qualname__}") 148 def run(self): 149 """Выполняет один цикл очистки старых трейсов Langfuse.""" 150 # Main logic for cleaning and optimizing Langfuse files 151 deleted = self.delete_old_traces(days_old=self.period) 152 logger.debug(f"Цикл {self.name} завершен: удалено {deleted} старых трейсов.")
Выполняет один цикл очистки старых трейсов Langfuse.