codekingpro/portable-devtools
114k
1import io2import threading3from typing import Optional4 5from langsmith import utils as ls_utils6 7try:8 from zstandard import ZstdCompressor # type: ignore[import]9 10 ZSTD_AVAILABLE = True11except ImportError:12 ZSTD_AVAILABLE = False13 14compression_level = int(ls_utils.get_env_var("RUN_COMPRESSION_LEVEL") or 1)15compression_threads = int(ls_utils.get_env_var("RUN_COMPRESSION_THREADS") or -1)16 17DEFAULT_MAX_UNCOMPRESSED_QUEUE_BYTES = 1024 * 1024 * 1024 # 1GB18 19 20class CompressedTraces:21 def __init__(self, max_uncompressed_size_bytes: Optional[int] = None) -> None:22 if not ZSTD_AVAILABLE:23 raise ImportError(24 "zstandard is required for compressed trace ingestion. "25 "Install it with `pip install zstandard` or set the environment "26 "variable LANGSMITH_DISABLE_RUN_COMPRESSION=true to disable "27 "compression."28 )29 # Configure the maximum total uncompressed size for the in-memory queue.30 if max_uncompressed_size_bytes is None:31 max_bytes_str = ls_utils.get_env_var("MAX_INGEST_MEMORY_BYTES")32 if max_bytes_str is not None:33 max_uncompressed_size_bytes = int(max_bytes_str)34 else:35 max_uncompressed_size_bytes = DEFAULT_MAX_UNCOMPRESSED_QUEUE_BYTES36 37 self.max_uncompressed_size_bytes = max_uncompressed_size_bytes38 39 self.buffer: io.BytesIO = io.BytesIO()40 self.trace_count: int = 041 self.lock = threading.Lock()42 self.uncompressed_size: int = 043 self._context: list[str] = []44 45 self.compressor_writer = ZstdCompressor(46 level=compression_level, threads=compression_threads47 ).stream_writer(self.buffer, closefd=False)48 49 def reset(self) -> None:50 self.buffer = io.BytesIO()51 self.trace_count = 052 self.uncompressed_size = 053 self._context = []54 self.compressor_writer = ZstdCompressor(55 level=compression_level, threads=-156 ).stream_writer(self.buffer, closefd=False)57 