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} старых трейсов.")
class LangfuseCleanerAssistant(Assistants.assistant.Assistant):
 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: Количество трейсов за один запрос.

LangfuseCleanerAssistant(project_slug, batch_size: int = 100, period: int = 1)
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: Период хранения трейсов в днях.

description
period
project_slug
batch_size
def should_delete_trace(self, trace):
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 если трейс старше указанного количества дней.

def delete_old_traces(self, days_old=30):
 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: Количество удалённых трейсов.

@observe(name=f'Цикл очистки {__qualname__}')
def run(self):
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.