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}')
class LangfuseCleanerAssistant(Assistants.assistant.Assistant):
 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, пакетно удаляет трейсы старше заданного периода хранения.

LangfuseCleanerAssistant(project_slug, batch_size: int = 100, period: int = 1)
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)
name
description
period
project_slug
batch_size
def should_delete_trace(self, trace):
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)

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

@observe(name=f'Цикл очистки {__qualname__}')
def run(self):
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. В случае ошибки логирует её и завершает выполнение.