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