Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
backend_cffi.py4479 linesDownload Raw Back to zstandard
1# Copyright (c) 2016-present, Gregory Szorc
2# All rights reserved.
3#
4# This software may be modified and distributed under the terms
5# of the BSD license. See the LICENSE file for details.
6
7"""Python interface to the Zstandard (zstd) compression library."""
8
9from __future__ import absolute_import, unicode_literals
10
11# This should match what the C extension exports.
12__all__ = [
13    "BufferSegment",
14    "BufferSegments",
15    "BufferWithSegments",
16    "BufferWithSegmentsCollection",
17    "ZstdCompressionChunker",
18    "ZstdCompressionDict",
19    "ZstdCompressionObj",
20    "ZstdCompressionParameters",
21    "ZstdCompressionReader",
22    "ZstdCompressionWriter",
23    "ZstdCompressor",
24    "ZstdDecompressionObj",
25    "ZstdDecompressionReader",
26    "ZstdDecompressionWriter",
27    "ZstdDecompressor",
28    "ZstdError",
29    "FrameParameters",
30    "backend_features",
31    "estimate_decompression_context_size",
32    "frame_content_size",
33    "frame_header_size",
34    "get_frame_parameters",
35    "train_dictionary",
36    # Constants.
37    "FLUSH_BLOCK",
38    "FLUSH_FRAME",
39    "COMPRESSOBJ_FLUSH_FINISH",
40    "COMPRESSOBJ_FLUSH_BLOCK",
41    "ZSTD_VERSION",
42    "FRAME_HEADER",
43    "CONTENTSIZE_UNKNOWN",
44    "CONTENTSIZE_ERROR",
45    "MAX_COMPRESSION_LEVEL",
46    "COMPRESSION_RECOMMENDED_INPUT_SIZE",
47    "COMPRESSION_RECOMMENDED_OUTPUT_SIZE",
48    "DECOMPRESSION_RECOMMENDED_INPUT_SIZE",
49    "DECOMPRESSION_RECOMMENDED_OUTPUT_SIZE",
50    "MAGIC_NUMBER",
51    "BLOCKSIZELOG_MAX",
52    "BLOCKSIZE_MAX",
53    "WINDOWLOG_MIN",
54    "WINDOWLOG_MAX",
55    "CHAINLOG_MIN",
56    "CHAINLOG_MAX",
57    "HASHLOG_MIN",
58    "HASHLOG_MAX",
59    "MINMATCH_MIN",
60    "MINMATCH_MAX",
61    "SEARCHLOG_MIN",
62    "SEARCHLOG_MAX",
63    "SEARCHLENGTH_MIN",
64    "SEARCHLENGTH_MAX",
65    "TARGETLENGTH_MIN",
66    "TARGETLENGTH_MAX",
67    "LDM_MINMATCH_MIN",
68    "LDM_MINMATCH_MAX",
69    "LDM_BUCKETSIZELOG_MAX",
70    "STRATEGY_FAST",
71    "STRATEGY_DFAST",
72    "STRATEGY_GREEDY",
73    "STRATEGY_LAZY",
74    "STRATEGY_LAZY2",
75    "STRATEGY_BTLAZY2",
76    "STRATEGY_BTOPT",
77    "STRATEGY_BTULTRA",
78    "STRATEGY_BTULTRA2",
79    "DICT_TYPE_AUTO",
80    "DICT_TYPE_RAWCONTENT",
81    "DICT_TYPE_FULLDICT",
82    "FORMAT_ZSTD1",
83    "FORMAT_ZSTD1_MAGICLESS",
84]
85
86import io
87import os
88
89from ._cffi import (  # type: ignore
90    ffi,
91    lib,
92)
93
94backend_features = set()  # type: ignore
95
96COMPRESSION_RECOMMENDED_INPUT_SIZE = lib.ZSTD_CStreamInSize()
97COMPRESSION_RECOMMENDED_OUTPUT_SIZE = lib.ZSTD_CStreamOutSize()
98DECOMPRESSION_RECOMMENDED_INPUT_SIZE = lib.ZSTD_DStreamInSize()
99DECOMPRESSION_RECOMMENDED_OUTPUT_SIZE = lib.ZSTD_DStreamOutSize()
100
101new_nonzero = ffi.new_allocator(should_clear_after_alloc=False)
102
103MAX_COMPRESSION_LEVEL = lib.ZSTD_maxCLevel()
104MAGIC_NUMBER = lib.ZSTD_MAGICNUMBER
105FRAME_HEADER = b"\x28\xb5\x2f\xfd"
106CONTENTSIZE_UNKNOWN = lib.ZSTD_CONTENTSIZE_UNKNOWN
107CONTENTSIZE_ERROR = lib.ZSTD_CONTENTSIZE_ERROR
108ZSTD_VERSION = (
109    lib.ZSTD_VERSION_MAJOR,
110    lib.ZSTD_VERSION_MINOR,
111    lib.ZSTD_VERSION_RELEASE,
112)
113
114BLOCKSIZELOG_MAX = lib.ZSTD_BLOCKSIZELOG_MAX
115BLOCKSIZE_MAX = lib.ZSTD_BLOCKSIZE_MAX
116WINDOWLOG_MIN = lib.ZSTD_WINDOWLOG_MIN
117WINDOWLOG_MAX = lib.ZSTD_WINDOWLOG_MAX
118CHAINLOG_MIN = lib.ZSTD_CHAINLOG_MIN
119CHAINLOG_MAX = lib.ZSTD_CHAINLOG_MAX
120HASHLOG_MIN = lib.ZSTD_HASHLOG_MIN
121HASHLOG_MAX = lib.ZSTD_HASHLOG_MAX
122MINMATCH_MIN = lib.ZSTD_MINMATCH_MIN
123MINMATCH_MAX = lib.ZSTD_MINMATCH_MAX
124SEARCHLOG_MIN = lib.ZSTD_SEARCHLOG_MIN
125SEARCHLOG_MAX = lib.ZSTD_SEARCHLOG_MAX
126SEARCHLENGTH_MIN = lib.ZSTD_MINMATCH_MIN
127SEARCHLENGTH_MAX = lib.ZSTD_MINMATCH_MAX
128TARGETLENGTH_MIN = lib.ZSTD_TARGETLENGTH_MIN
129TARGETLENGTH_MAX = lib.ZSTD_TARGETLENGTH_MAX
130LDM_MINMATCH_MIN = lib.ZSTD_LDM_MINMATCH_MIN
131LDM_MINMATCH_MAX = lib.ZSTD_LDM_MINMATCH_MAX
132LDM_BUCKETSIZELOG_MAX = lib.ZSTD_LDM_BUCKETSIZELOG_MAX
133
134STRATEGY_FAST = lib.ZSTD_fast
135STRATEGY_DFAST = lib.ZSTD_dfast
136STRATEGY_GREEDY = lib.ZSTD_greedy
137STRATEGY_LAZY = lib.ZSTD_lazy
138STRATEGY_LAZY2 = lib.ZSTD_lazy2
139STRATEGY_BTLAZY2 = lib.ZSTD_btlazy2
140STRATEGY_BTOPT = lib.ZSTD_btopt
141STRATEGY_BTULTRA = lib.ZSTD_btultra
142STRATEGY_BTULTRA2 = lib.ZSTD_btultra2
143
144DICT_TYPE_AUTO = lib.ZSTD_dct_auto
145DICT_TYPE_RAWCONTENT = lib.ZSTD_dct_rawContent
146DICT_TYPE_FULLDICT = lib.ZSTD_dct_fullDict
147
148FORMAT_ZSTD1 = lib.ZSTD_f_zstd1
149FORMAT_ZSTD1_MAGICLESS = lib.ZSTD_f_zstd1_magicless
150
151FLUSH_BLOCK = 0
152FLUSH_FRAME = 1
153
154COMPRESSOBJ_FLUSH_FINISH = 0
155COMPRESSOBJ_FLUSH_BLOCK = 1
156
157
158def _cpu_count():
159    # os.cpu_count() was introducd in Python 3.4.
160    try:
161        return os.cpu_count() or 0
162    except AttributeError:
163        pass
164
165    # Linux.
166    try:
167        return os.sysconf("SC_NPROCESSORS_ONLN")
168    except (AttributeError, ValueError):
169        pass
170
171    # TODO implement on other platforms.
172    return 0
173
174
175class BufferSegment:
176    """Represents a segment within a ``BufferWithSegments``.
177
178    This type is essentially a reference to N bytes within a
179    ``BufferWithSegments``.
180
181    The object conforms to the buffer protocol.
182    """
183
184    @property
185    def offset(self):
186        """The byte offset of this segment within its parent buffer."""
187        raise NotImplementedError()
188
189    def __len__(self):
190        """Obtain the length of the segment, in bytes."""
191        raise NotImplementedError()
192
193    def tobytes(self):
194        """Obtain bytes copy of this segment."""
195        raise NotImplementedError()
196
197
198class BufferSegments:
199    """Represents an array of ``(offset, length)`` integers.
200
201    This type is effectively an index used by :py:class:`BufferWithSegments`.
202
203    The array members are 64-bit unsigned integers using host/native bit order.
204
205    Instances conform to the buffer protocol.
206    """
207
208
209class BufferWithSegments:
210    """A memory buffer containing N discrete items of known lengths.
211
212    This type is essentially a fixed size memory address and an array
213    of 2-tuples of ``(offset, length)`` 64-bit unsigned native-endian
214    integers defining the byte offset and length of each segment within
215    the buffer.
216
217    Instances behave like containers.
218
219    Instances also conform to the buffer protocol. So a reference to the
220    backing bytes can be obtained via ``memoryview(o)``. A *copy* of the
221    backing bytes can be obtained via ``.tobytes()``.
222
223    This type exists to facilitate operations against N>1 items without
224    the overhead of Python object creation and management. Used with
225    APIs like :py:meth:`ZstdDecompressor.multi_decompress_to_buffer`, it
226    is possible to decompress many objects in parallel without the GIL
227    held, leading to even better performance.
228    """
229
230    @property
231    def size(self):
232        """Total sizein bytes of the backing buffer."""
233        raise NotImplementedError()
234
235    def __len__(self):
236        raise NotImplementedError()
237
238    def __getitem__(self, i):
239        """Obtains a segment within the buffer.
240
241        The returned object references memory within this buffer.
242
243        :param i:
244           Integer index of segment to retrieve.
245        :return:
246           :py:class:`BufferSegment`
247        """
248        raise NotImplementedError()
249
250    def segments(self):
251        """Obtain the array of ``(offset, length)`` segments in the buffer.
252
253        :return:
254           :py:class:`BufferSegments`
255        """
256        raise NotImplementedError()
257
258    def tobytes(self):
259        """Obtain bytes copy of this instance."""
260        raise NotImplementedError()
261
262
263class BufferWithSegmentsCollection:
264    """A virtual spanning view over multiple BufferWithSegments.
265
266    Instances are constructed from 1 or more :py:class:`BufferWithSegments`
267    instances. The resulting object behaves like an ordered sequence whose
268    members are the segments within each ``BufferWithSegments``.
269
270    If the object is composed of 2 ``BufferWithSegments`` instances with the
271    first having 2 segments and the second have 3 segments, then ``b[0]``
272    and ``b[1]`` access segments in the first object and ``b[2]``, ``b[3]``,
273    and ``b[4]`` access segments from the second.
274    """
275
276    def __len__(self):
277        """The number of segments within all ``BufferWithSegments``."""
278        raise NotImplementedError()
279
280    def __getitem__(self, i):
281        """Obtain the ``BufferSegment`` at an offset."""
282        raise NotImplementedError()
283
284
285class ZstdError(Exception):
286    pass
287
288
289def _zstd_error(zresult):
290    # Resolves to bytes on Python 2 and 3. We use the string for formatting
291    # into error messages, which will be literal unicode. So convert it to
292    # unicode.
293    return ffi.string(lib.ZSTD_getErrorName(zresult)).decode("utf-8")
294
295
296def _make_cctx_params(params):
297    res = lib.ZSTD_createCCtxParams()
298    if res == ffi.NULL:
299        raise MemoryError()
300
301    res = ffi.gc(res, lib.ZSTD_freeCCtxParams)
302
303    attrs = [
304        (lib.ZSTD_c_format, params.format),
305        (lib.ZSTD_c_compressionLevel, params.compression_level),
306        (lib.ZSTD_c_windowLog, params.window_log),
307        (lib.ZSTD_c_hashLog, params.hash_log),
308        (lib.ZSTD_c_chainLog, params.chain_log),
309        (lib.ZSTD_c_searchLog, params.search_log),
310        (lib.ZSTD_c_minMatch, params.min_match),
311        (lib.ZSTD_c_targetLength, params.target_length),
312        (lib.ZSTD_c_strategy, params.strategy),
313        (lib.ZSTD_c_contentSizeFlag, params.write_content_size),
314        (lib.ZSTD_c_checksumFlag, params.write_checksum),
315        (lib.ZSTD_c_dictIDFlag, params.write_dict_id),
316        (lib.ZSTD_c_nbWorkers, params.threads),
317        (lib.ZSTD_c_jobSize, params.job_size),
318        (lib.ZSTD_c_overlapLog, params.overlap_log),
319        (lib.ZSTD_c_forceMaxWindow, params.force_max_window),
320        (lib.ZSTD_c_enableLongDistanceMatching, params.enable_ldm),
321        (lib.ZSTD_c_ldmHashLog, params.ldm_hash_log),
322        (lib.ZSTD_c_ldmMinMatch, params.ldm_min_match),
323        (lib.ZSTD_c_ldmBucketSizeLog, params.ldm_bucket_size_log),
324        (lib.ZSTD_c_ldmHashRateLog, params.ldm_hash_rate_log),
325    ]
326
327    for param, value in attrs:
328        _set_compression_parameter(res, param, value)
329
330    return res
331
332
333class ZstdCompressionParameters(object):
334    """Low-level zstd compression parameters.
335
336    This type represents a collection of parameters to control how zstd
337    compression is performed.
338
339    Instances can be constructed from raw parameters or derived from a
340    base set of defaults specified from a compression level (recommended)
341    via :py:meth:`ZstdCompressionParameters.from_level`.
342
343    >>> # Derive compression settings for compression level 7.
344    >>> params = zstandard.ZstdCompressionParameters.from_level(7)
345
346    >>> # With an input size of 1MB
347    >>> params = zstandard.ZstdCompressionParameters.from_level(7, source_size=1048576)
348
349    Using ``from_level()``, it is also possible to override individual compression
350    parameters or to define additional settings that aren't automatically derived.
351    e.g.:
352
353    >>> params = zstandard.ZstdCompressionParameters.from_level(4, window_log=10)
354    >>> params = zstandard.ZstdCompressionParameters.from_level(5, threads=4)
355
356    Or you can define low-level compression settings directly:
357
358    >>> params = zstandard.ZstdCompressionParameters(window_log=12, enable_ldm=True)
359
360    Once a ``ZstdCompressionParameters`` instance is obtained, it can be used to
361    configure a compressor:
362
363    >>> cctx = zstandard.ZstdCompressor(compression_params=params)
364
365    Some of these are very low-level settings. It may help to consult the official
366    zstandard documentation for their behavior. Look for the ``ZSTD_p_*`` constants
367    in ``zstd.h`` (https://github.com/facebook/zstd/blob/dev/lib/zstd.h).
368    """
369
370    @staticmethod
371    def from_level(level, source_size=0, dict_size=0, **kwargs):
372        """Create compression parameters from a compression level.
373
374        :param level:
375           Integer compression level.
376        :param source_size:
377           Integer size in bytes of source to be compressed.
378        :param dict_size:
379           Integer size in bytes of compression dictionary to use.
380        :return:
381           :py:class:`ZstdCompressionParameters`
382        """
383        params = lib.ZSTD_getCParams(level, source_size, dict_size)
384
385        args = {
386            "window_log": "windowLog",
387            "chain_log": "chainLog",
388            "hash_log": "hashLog",
389            "search_log": "searchLog",
390            "min_match": "minMatch",
391            "target_length": "targetLength",
392            "strategy": "strategy",
393        }
394
395        for arg, attr in args.items():
396            if arg not in kwargs:
397                kwargs[arg] = getattr(params, attr)
398
399        return ZstdCompressionParameters(**kwargs)
400
401    def __init__(
402        self,
403        format=0,
404        compression_level=0,
405        window_log=0,
406        hash_log=0,
407        chain_log=0,
408        search_log=0,
409        min_match=0,
410        target_length=0,
411        strategy=-1,
412        write_content_size=1,
413        write_checksum=0,
414        write_dict_id=0,
415        job_size=0,
416        overlap_log=-1,
417        force_max_window=0,
418        enable_ldm=0,
419        ldm_hash_log=0,
420        ldm_min_match=0,
421        ldm_bucket_size_log=0,
422        ldm_hash_rate_log=-1,
423        threads=0,
424    ):
425        params = lib.ZSTD_createCCtxParams()
426        if params == ffi.NULL:
427            raise MemoryError()
428
429        params = ffi.gc(params, lib.ZSTD_freeCCtxParams)
430
431        self._params = params
432
433        if threads < 0:
434            threads = _cpu_count()
435
436        # We need to set ZSTD_c_nbWorkers before ZSTD_c_jobSize and ZSTD_c_overlapLog
437        # because setting ZSTD_c_nbWorkers resets the other parameters.
438        _set_compression_parameter(params, lib.ZSTD_c_nbWorkers, threads)
439
440        _set_compression_parameter(params, lib.ZSTD_c_format, format)
441        _set_compression_parameter(
442            params, lib.ZSTD_c_compressionLevel, compression_level
443        )
444        _set_compression_parameter(params, lib.ZSTD_c_windowLog, window_log)
445        _set_compression_parameter(params, lib.ZSTD_c_hashLog, hash_log)
446        _set_compression_parameter(params, lib.ZSTD_c_chainLog, chain_log)
447        _set_compression_parameter(params, lib.ZSTD_c_searchLog, search_log)
448        _set_compression_parameter(params, lib.ZSTD_c_minMatch, min_match)
449        _set_compression_parameter(
450            params, lib.ZSTD_c_targetLength, target_length
451        )
452
453        if strategy == -1:
454            strategy = 0
455
456        _set_compression_parameter(params, lib.ZSTD_c_strategy, strategy)
457        _set_compression_parameter(
458            params, lib.ZSTD_c_contentSizeFlag, write_content_size
459        )
460        _set_compression_parameter(
461            params, lib.ZSTD_c_checksumFlag, write_checksum
462        )
463        _set_compression_parameter(params, lib.ZSTD_c_dictIDFlag, write_dict_id)
464        _set_compression_parameter(params, lib.ZSTD_c_jobSize, job_size)
465
466        if overlap_log == -1:
467            overlap_log = 0
468
469        _set_compression_parameter(params, lib.ZSTD_c_overlapLog, overlap_log)
470        _set_compression_parameter(
471            params, lib.ZSTD_c_forceMaxWindow, force_max_window
472        )
473        _set_compression_parameter(
474            params, lib.ZSTD_c_enableLongDistanceMatching, enable_ldm
475        )
476        _set_compression_parameter(params, lib.ZSTD_c_ldmHashLog, ldm_hash_log)
477        _set_compression_parameter(
478            params, lib.ZSTD_c_ldmMinMatch, ldm_min_match
479        )
480        _set_compression_parameter(
481            params, lib.ZSTD_c_ldmBucketSizeLog, ldm_bucket_size_log
482        )
483
484        if ldm_hash_rate_log == -1:
485            ldm_hash_rate_log = 0
486
487        _set_compression_parameter(
488            params, lib.ZSTD_c_ldmHashRateLog, ldm_hash_rate_log
489        )
490
491    @property
492    def format(self):
493        return _get_compression_parameter(self._params, lib.ZSTD_c_format)
494
495    @property
496    def compression_level(self):
497        return _get_compression_parameter(
498            self._params, lib.ZSTD_c_compressionLevel
499        )
500
501    @property
502    def window_log(self):
503        return _get_compression_parameter(self._params, lib.ZSTD_c_windowLog)
504
505    @property
506    def hash_log(self):
507        return _get_compression_parameter(self._params, lib.ZSTD_c_hashLog)
508
509    @property
510    def chain_log(self):
511        return _get_compression_parameter(self._params, lib.ZSTD_c_chainLog)
512
513    @property
514    def search_log(self):
515        return _get_compression_parameter(self._params, lib.ZSTD_c_searchLog)
516
517    @property
518    def min_match(self):
519        return _get_compression_parameter(self._params, lib.ZSTD_c_minMatch)
520
521    @property
522    def target_length(self):
523        return _get_compression_parameter(self._params, lib.ZSTD_c_targetLength)
524
525    @property
526    def strategy(self):
527        return _get_compression_parameter(self._params, lib.ZSTD_c_strategy)
528
529    @property
530    def write_content_size(self):
531        return _get_compression_parameter(
532            self._params, lib.ZSTD_c_contentSizeFlag
533        )
534
535    @property
536    def write_checksum(self):
537        return _get_compression_parameter(self._params, lib.ZSTD_c_checksumFlag)
538
539    @property
540    def write_dict_id(self):
541        return _get_compression_parameter(self._params, lib.ZSTD_c_dictIDFlag)
542
543    @property
544    def job_size(self):
545        return _get_compression_parameter(self._params, lib.ZSTD_c_jobSize)
546
547    @property
548    def overlap_log(self):
549        return _get_compression_parameter(self._params, lib.ZSTD_c_overlapLog)
550
551    @property
552    def force_max_window(self):
553        return _get_compression_parameter(
554            self._params, lib.ZSTD_c_forceMaxWindow
555        )
556
557    @property
558    def enable_ldm(self):
559        return _get_compression_parameter(
560            self._params, lib.ZSTD_c_enableLongDistanceMatching
561        )
562
563    @property
564    def ldm_hash_log(self):
565        return _get_compression_parameter(self._params, lib.ZSTD_c_ldmHashLog)
566
567    @property
568    def ldm_min_match(self):
569        return _get_compression_parameter(self._params, lib.ZSTD_c_ldmMinMatch)
570
571    @property
572    def ldm_bucket_size_log(self):
573        return _get_compression_parameter(
574            self._params, lib.ZSTD_c_ldmBucketSizeLog
575        )
576
577    @property
578    def ldm_hash_rate_log(self):
579        return _get_compression_parameter(
580            self._params, lib.ZSTD_c_ldmHashRateLog
581        )
582
583    @property
584    def threads(self):
585        return _get_compression_parameter(self._params, lib.ZSTD_c_nbWorkers)
586
587    def estimated_compression_context_size(self):
588        """Estimated size in bytes needed to compress with these parameters."""
589        return lib.ZSTD_estimateCCtxSize_usingCCtxParams(self._params)
590
591
592def estimate_decompression_context_size():
593    """Estimate the memory size requirements for a decompressor instance.
594
595    :return:
596       Integer number of bytes.
597    """
598    return lib.ZSTD_estimateDCtxSize()
599
600
601def _set_compression_parameter(params, param, value):
602    zresult = lib.ZSTD_CCtxParams_setParameter(params, param, value)
603    if lib.ZSTD_isError(zresult):
604        raise ZstdError(
605            "unable to set compression context parameter: %s"
606            % _zstd_error(zresult)
607        )
608
609
610def _get_compression_parameter(params, param):
611    result = ffi.new("int *")
612
613    zresult = lib.ZSTD_CCtxParams_getParameter(params, param, result)
614    if lib.ZSTD_isError(zresult):
615        raise ZstdError(
616            "unable to get compression context parameter: %s"
617            % _zstd_error(zresult)
618        )
619
620    return result[0]
621
622
623class ZstdCompressionWriter(object):
624    """Writable compressing stream wrapper.
625
626    ``ZstdCompressionWriter`` is a write-only stream interface for writing
627    compressed data to another stream.
628
629    This type conforms to the ``io.RawIOBase`` interface and should be usable
630    by any type that operates against a *file-object* (``typing.BinaryIO``
631    in Python type hinting speak). Only methods that involve writing will do
632    useful things.
633
634    As data is written to this stream (e.g. via ``write()``), that data
635    is sent to the compressor. As compressed data becomes available from
636    the compressor, it is sent to the underlying stream by calling its
637    ``write()`` method.
638
639    Both ``write()`` and ``flush()`` return the number of bytes written to the
640    object's ``write()``. In many cases, small inputs do not accumulate enough
641    data to cause a write and ``write()`` will return ``0``.
642
643    Calling ``close()`` will mark the stream as closed and subsequent I/O
644    operations will raise ``ValueError`` (per the documented behavior of
645    ``io.RawIOBase``). ``close()`` will also call ``close()`` on the underlying
646    stream if such a method exists and the instance was constructed with
647    ``closefd=True``
648
649    Instances are obtained by calling :py:meth:`ZstdCompressor.stream_writer`.
650
651    Typically usage is as follows:
652
653    >>> cctx = zstandard.ZstdCompressor(level=10)
654    >>> compressor = cctx.stream_writer(fh)
655    >>> compressor.write(b"chunk 0\\n")
656    >>> compressor.write(b"chunk 1\\n")
657    >>> compressor.flush()
658    >>> # Receiver will be able to decode ``chunk 0\\nchunk 1\\n`` at this point.
659    >>> # Receiver is also expecting more data in the zstd *frame*.
660    >>>
661    >>> compressor.write(b"chunk 2\\n")
662    >>> compressor.flush(zstandard.FLUSH_FRAME)
663    >>> # Receiver will be able to decode ``chunk 0\\nchunk 1\\nchunk 2``.
664    >>> # Receiver is expecting no more data, as the zstd frame is closed.
665    >>> # Any future calls to ``write()`` at this point will construct a new
666    >>> # zstd frame.
667
668    Instances can be used as context managers. Exiting the context manager is
669    the equivalent of calling ``close()``, which is equivalent to calling
670    ``flush(zstandard.FLUSH_FRAME)``:
671
672    >>> cctx = zstandard.ZstdCompressor(level=10)
673    >>> with cctx.stream_writer(fh) as compressor:
674    ...     compressor.write(b'chunk 0')
675    ...     compressor.write(b'chunk 1')
676    ...     ...
677
678    .. important::
679
680       If ``flush(FLUSH_FRAME)`` is not called, emitted data doesn't
681       constitute a full zstd *frame* and consumers of this data may complain
682       about malformed input. It is recommended to use instances as a context
683       manager to ensure *frames* are properly finished.
684
685    If the size of the data being fed to this streaming compressor is known,
686    you can declare it before compression begins:
687
688    >>> cctx = zstandard.ZstdCompressor()
689    >>> with cctx.stream_writer(fh, size=data_len) as compressor:
690    ...     compressor.write(chunk0)
691    ...     compressor.write(chunk1)
692    ...     ...
693
694    Declaring the size of the source data allows compression parameters to
695    be tuned. And if ``write_content_size`` is used, it also results in the
696    content size being written into the frame header of the output data.
697
698    The size of chunks being ``write()`` to the destination can be specified:
699
700    >>> cctx = zstandard.ZstdCompressor()
701    >>> with cctx.stream_writer(fh, write_size=32768) as compressor:
702    ...     ...
703
704    To see how much memory is being used by the streaming compressor:
705
706    >>> cctx = zstandard.ZstdCompressor()
707    >>> with cctx.stream_writer(fh) as compressor:
708    ...     ...
709    ...     byte_size = compressor.memory_size()
710
711    Thte total number of bytes written so far are exposed via ``tell()``:
712
713    >>> cctx = zstandard.ZstdCompressor()
714    >>> with cctx.stream_writer(fh) as compressor:
715    ...     ...
716    ...     total_written = compressor.tell()
717
718    ``stream_writer()`` accepts a ``write_return_read`` boolean argument to
719    control the return value of ``write()``. When ``False`` (the default),
720    ``write()`` returns the number of bytes that were ``write()``'en to the
721    underlying object. When ``True``, ``write()`` returns the number of bytes
722    read from the input that were subsequently written to the compressor.
723    ``True`` is the *proper* behavior for ``write()`` as specified by the
724    ``io.RawIOBase`` interface and will become the default value in a future
725    release.
726    """
727
728    def __init__(
729        self,
730        compressor,
731        writer,
732        source_size,
733        write_size,
734        write_return_read,
735        closefd=True,
736    ):
737        self._compressor = compressor
738        self._writer = writer
739        self._write_size = write_size
740        self._write_return_read = bool(write_return_read)
741        self._closefd = bool(closefd)
742        self._entered = False
743        self._closing = False
744        self._closed = False
745        self._bytes_compressed = 0
746
747        self._dst_buffer = ffi.new("char[]", write_size)
748        self._out_buffer = ffi.new("ZSTD_outBuffer *")
749        self._out_buffer.dst = self._dst_buffer
750        self._out_buffer.size = len(self._dst_buffer)
751        self._out_buffer.pos = 0
752
753        zresult = lib.ZSTD_CCtx_setPledgedSrcSize(compressor._cctx, source_size)
754        if lib.ZSTD_isError(zresult):
755            raise ZstdError(
756                "error setting source size: %s" % _zstd_error(zresult)
757            )
758
759    def __enter__(self):
760        if self._closed:
761            raise ValueError("stream is closed")
762
763        if self._entered:
764            raise ZstdError("cannot __enter__ multiple times")
765
766        self._entered = True
767        return self
768
769    def __exit__(self, exc_type, exc_value, exc_tb):
770        self._entered = False
771        self.close()
772        self._compressor = None
773
774        return False
775
776    def __iter__(self):
777        raise io.UnsupportedOperation()
778
779    def __next__(self):
780        raise io.UnsupportedOperation()
781
782    def memory_size(self):
783        return lib.ZSTD_sizeof_CCtx(self._compressor._cctx)
784
785    def fileno(self):
786        f = getattr(self._writer, "fileno", None)
787        if f:
788            return f()
789        else:
790            raise OSError("fileno not available on underlying writer")
791
792    def close(self):
793        if self._closed:
794            return
795
796        try:
797            self._closing = True
798            self.flush(FLUSH_FRAME)
799        finally:
800            self._closing = False
801            self._closed = True
802
803        # Call close() on underlying stream as well.
804        f = getattr(self._writer, "close", None)
805        if self._closefd and f:
806            f()
807
808    @property
809    def closed(self):
810        return self._closed
811
812    def isatty(self):
813        return False
814
815    def readable(self):
816        return False
817
818    def readline(self, size=-1):
819        raise io.UnsupportedOperation()
820
821    def readlines(self, hint=-1):
822        raise io.UnsupportedOperation()
823
824    def seek(self, offset, whence=None):
825        raise io.UnsupportedOperation()
826
827    def seekable(self):
828        return False
829
830    def truncate(self, size=None):
831        raise io.UnsupportedOperation()
832
833    def writable(self):
834        return True
835
836    def writelines(self, lines):
837        raise NotImplementedError("writelines() is not yet implemented")
838
839    def read(self, size=-1):
840        raise io.UnsupportedOperation()
841
842    def readall(self):
843        raise io.UnsupportedOperation()
844
845    def readinto(self, b):
846        raise io.UnsupportedOperation()
847
848    def write(self, data):
849        """Send data to the compressor and possibly to the inner stream."""
850        if self._closed:
851            raise ValueError("stream is closed")
852
853        total_write = 0
854
855        data_buffer = ffi.from_buffer(data)
856
857        in_buffer = ffi.new("ZSTD_inBuffer *")
858        in_buffer.src = data_buffer
859        in_buffer.size = len(data_buffer)
860        in_buffer.pos = 0
861
862        out_buffer = self._out_buffer
863        out_buffer.pos = 0
864
865        while in_buffer.pos < in_buffer.size:
866            zresult = lib.ZSTD_compressStream2(
867                self._compressor._cctx,
868                out_buffer,
869                in_buffer,
870                lib.ZSTD_e_continue,
871            )
872            if lib.ZSTD_isError(zresult):
873                raise ZstdError(
874                    "zstd compress error: %s" % _zstd_error(zresult)
875                )
876
877            if out_buffer.pos:
878                self._writer.write(
879                    ffi.buffer(out_buffer.dst, out_buffer.pos)[:]
880                )
881                total_write += out_buffer.pos
882                self._bytes_compressed += out_buffer.pos
883                out_buffer.pos = 0
884
885        if self._write_return_read:
886            return in_buffer.pos
887        else:
888            return total_write
889
890    def flush(self, flush_mode=FLUSH_BLOCK):
891        """Evict data from compressor's internal state and write it to inner stream.
892
893        Calling this method may result in 0 or more ``write()`` calls to the
894        inner stream.
895
896        This method will also call ``flush()`` on the inner stream, if such a
897        method exists.
898
899        :param flush_mode:
900           How to flush the zstd compressor.
901
902           ``zstandard.FLUSH_BLOCK`` will flush data already sent to the
903           compressor but not emitted to the inner stream. The stream is still
904           writable after calling this. This is the default behavior.
905
906           See documentation for other ``zstandard.FLUSH_*`` constants for more
907           flushing options.
908        :return:
909           Integer number of bytes written to the inner stream.
910        """
911
912        if flush_mode == FLUSH_BLOCK:
913            flush = lib.ZSTD_e_flush
914        elif flush_mode == FLUSH_FRAME:
915            flush = lib.ZSTD_e_end
916        else:
917            raise ValueError("unknown flush_mode: %r" % flush_mode)
918
919        if self._closed:
920            raise ValueError("stream is closed")
921
922        total_write = 0
923
924        out_buffer = self._out_buffer
925        out_buffer.pos = 0
926
927        in_buffer = ffi.new("ZSTD_inBuffer *")
928        in_buffer.src = ffi.NULL
929        in_buffer.size = 0
930        in_buffer.pos = 0
931
932        while True:
933            zresult = lib.ZSTD_compressStream2(
934                self._compressor._cctx, out_buffer, in_buffer, flush
935            )
936            if lib.ZSTD_isError(zresult):
937                raise ZstdError(
938                    "zstd compress error: %s" % _zstd_error(zresult)
939                )
940
941            if out_buffer.pos:
942                self._writer.write(
943                    ffi.buffer(out_buffer.dst, out_buffer.pos)[:]
944                )
945                total_write += out_buffer.pos
946                self._bytes_compressed += out_buffer.pos
947                out_buffer.pos = 0
948
949            if not zresult:
950                break
951
952        f = getattr(self._writer, "flush", None)
953        if f and not self._closing:
954            f()
955
956        return total_write
957
958    def tell(self):
959        return self._bytes_compressed
960
961
962class ZstdCompressionObj(object):
963    """A compressor conforming to the API in Python's standard library.
964
965    This type implements an API similar to compression types in Python's
966    standard library such as ``zlib.compressobj`` and ``bz2.BZ2Compressor``.
967    This enables existing code targeting the standard library API to swap
968    in this type to achieve zstd compression.
969
970    .. important::
971
972       The design of this API is not ideal for optimal performance.
973
974       The reason performance is not optimal is because the API is limited to
975       returning a single buffer holding compressed data. When compressing
976       data, we don't know how much data will be emitted. So in order to
977       capture all this data in a single buffer, we need to perform buffer
978       reallocations and/or extra memory copies. This can add significant
979       overhead depending on the size or nature of the compressed data how
980       much your application calls this type.
981
982       If performance is critical, consider an API like
983       :py:meth:`ZstdCompressor.stream_reader`,
984       :py:meth:`ZstdCompressor.stream_writer`,
985       :py:meth:`ZstdCompressor.chunker`, or
986       :py:meth:`ZstdCompressor.read_to_iter`, which result in less overhead
987       managing buffers.
988
989    Instances are obtained by calling :py:meth:`ZstdCompressor.compressobj`.
990
991    Here is how this API should be used:
992
993    >>> cctx = zstandard.ZstdCompressor()
994    >>> cobj = cctx.compressobj()
995    >>> data = cobj.compress(b"raw input 0")
996    >>> data = cobj.compress(b"raw input 1")
997    >>> data = cobj.flush()
998
999    Or to flush blocks:
1000
1001    >>> cctx.zstandard.ZstdCompressor()
1002    >>> cobj = cctx.compressobj()
1003    >>> data = cobj.compress(b"chunk in first block")
1004    >>> data = cobj.flush(zstandard.COMPRESSOBJ_FLUSH_BLOCK)
1005    >>> data = cobj.compress(b"chunk in second block")
1006    >>> data = cobj.flush()
1007
1008    For best performance results, keep input chunks under 256KB. This avoids
1009    extra allocations for a large output object.
1010
1011    It is possible to declare the input size of the data that will be fed
1012    into the compressor:
1013
1014    >>> cctx = zstandard.ZstdCompressor()
1015    >>> cobj = cctx.compressobj(size=6)
1016    >>> data = cobj.compress(b"foobar")
1017    >>> data = cobj.flush()
1018    """
1019
1020    def __init__(
1021        self, compressor, write_size=COMPRESSION_RECOMMENDED_OUTPUT_SIZE
1022    ):
1023        self._compressor = compressor
1024        self._out = ffi.new("ZSTD_outBuffer *")
1025        self._dst_buffer = ffi.new("char[]", write_size)
1026        self._out.dst = self._dst_buffer
1027        self._out.size = write_size
1028        self._out.pos = 0
1029        self._finished = False
1030
1031    def compress(self, data):
1032        """Send data to the compressor.
1033
1034        This method receives bytes to feed to the compressor and returns
1035        bytes constituting zstd compressed data.
1036
1037        The zstd compressor accumulates bytes and the returned bytes may be
1038        substantially smaller or larger than the size of the input data on
1039        any given call. The returned value may be the empty byte string
1040        (``b""``).
1041
1042        :param data:
1043           Data to write to the compressor.
1044        :return:
1045           Compressed data.
1046        """
1047        if self._finished:
1048            raise ZstdError("cannot call compress() after compressor finished")
1049
1050        data_buffer = ffi.from_buffer(data)
1051        source = ffi.new("ZSTD_inBuffer *")
1052        source.src = data_buffer
1053        source.size = len(data_buffer)
1054        source.pos = 0
1055
1056        chunks = []
1057
1058        while source.pos < len(data):
1059            zresult = lib.ZSTD_compressStream2(
1060                self._compressor._cctx, self._out, source, lib.ZSTD_e_continue
1061            )
1062            if lib.ZSTD_isError(zresult):
1063                raise ZstdError(
1064                    "zstd compress error: %s" % _zstd_error(zresult)
1065                )
1066
1067            if self._out.pos:
1068                chunks.append(ffi.buffer(self._out.dst, self._out.pos)[:])
1069                self._out.pos = 0
1070
1071        return b"".join(chunks)
1072
1073    def flush(self, flush_mode=COMPRESSOBJ_FLUSH_FINISH):
1074        """Emit data accumulated in the compressor that hasn't been outputted yet.
1075
1076        The ``flush_mode`` argument controls how to end the stream.
1077
1078        ``zstandard.COMPRESSOBJ_FLUSH_FINISH`` (the default) ends the
1079        compression stream and finishes a zstd frame. Once this type of flush
1080        is performed, ``compress()`` and ``flush()`` can no longer be called.
1081        This type of flush **must** be called to end the compression context. If
1082        not called, the emitted data may be incomplete and may not be readable
1083        by a decompressor.
1084
1085        ``zstandard.COMPRESSOBJ_FLUSH_BLOCK`` will flush a zstd block. This
1086        ensures that all data fed to this instance will have been omitted and
1087        can be decoded by a decompressor. Flushes of this type can be performed
1088        multiple times. The next call to ``compress()`` will begin a new zstd
1089        block.
1090
1091        :param flush_mode:
1092           How to flush the zstd compressor.
1093        :return:
1094           Compressed data.
1095        """
1096        if flush_mode not in (
1097            COMPRESSOBJ_FLUSH_FINISH,
1098            COMPRESSOBJ_FLUSH_BLOCK,
1099        ):
1100            raise ValueError("flush mode not recognized")
1101
1102        if self._finished:
1103            raise ZstdError("compressor object already finished")
1104
1105        if flush_mode == COMPRESSOBJ_FLUSH_BLOCK:
1106            z_flush_mode = lib.ZSTD_e_flush
1107        elif flush_mode == COMPRESSOBJ_FLUSH_FINISH:
1108            z_flush_mode = lib.ZSTD_e_end
1109            self._finished = True
1110        else:
1111            raise ZstdError("unhandled flush mode")
1112
1113        assert self._out.pos == 0
1114
1115        in_buffer = ffi.new("ZSTD_inBuffer *")
1116        in_buffer.src = ffi.NULL
1117        in_buffer.size = 0
1118        in_buffer.pos = 0
1119
1120        chunks = []
1121
1122        while True:
1123            zresult = lib.ZSTD_compressStream2(
1124                self._compressor._cctx, self._out, in_buffer, z_flush_mode
1125            )
1126            if lib.ZSTD_isError(zresult):
1127                raise ZstdError(
1128                    "error ending compression stream: %s" % _zstd_error(zresult)
1129                )
1130
1131            if self._out.pos:
1132                chunks.append(ffi.buffer(self._out.dst, self._out.pos)[:])
1133                self._out.pos = 0
1134
1135            if not zresult:
1136                break
1137
1138        return b"".join(chunks)
1139
1140
1141class ZstdCompressionChunker(object):
1142    """Compress data to uniformly sized chunks.
1143
1144    This type allows you to iteratively feed chunks of data into a compressor
1145    and produce output chunks of uniform size.
1146
1147    ``compress()``, ``flush()``, and ``finish()`` all return an iterator of
1148    ``bytes`` instances holding compressed data. The iterator may be empty.
1149    Callers MUST iterate through all elements of the returned iterator before
1150    performing another operation on the object or else the compressor's
1151    internal state may become confused. This can result in an exception being
1152    raised or malformed data being emitted.
1153
1154    All chunks emitted by ``compress()`` will have a length of the configured
1155    chunk size.
1156
1157    ``flush()`` and ``finish()`` may return a final chunk smaller than
1158    the configured chunk size.
1159
1160    Instances are obtained by calling :py:meth:`ZstdCompressor.chunker`.
1161
1162    Here is how the API should be used:
1163
1164    >>> cctx = zstandard.ZstdCompressor()
1165    >>> chunker = cctx.chunker(chunk_size=32768)
1166    >>>
1167    >>> with open(path, 'rb') as fh:
1168    ...     while True:
1169    ...         in_chunk = fh.read(32768)
1170    ...         if not in_chunk:
1171    ...             break
1172    ...
1173    ...         for out_chunk in chunker.compress(in_chunk):
1174    ...             # Do something with output chunk of size 32768.
1175    ...
1176    ...     for out_chunk in chunker.finish():
1177    ...         # Do something with output chunks that finalize the zstd frame.
1178
1179    This compressor type is often a better alternative to
1180    :py:class:`ZstdCompressor.compressobj` because it has better performance
1181    properties.
1182
1183    ``compressobj()`` will emit output data as it is available. This results
1184    in a *stream* of output chunks of varying sizes. The consistency of the
1185    output chunk size with ``chunker()`` is more appropriate for many usages,
1186    such as sending compressed data to a socket.
1187
1188    ``compressobj()`` may also perform extra memory reallocations in order
1189    to dynamically adjust the sizes of the output chunks. Since ``chunker()``
1190    output chunks are all the same size (except for flushed or final chunks),
1191    there is less memory allocation/copying overhead.
1192    """
1193
1194    def __init__(self, compressor, chunk_size):
1195        self._compressor = compressor
1196        self._out = ffi.new("ZSTD_outBuffer *")
1197        self._dst_buffer = ffi.new("char[]", chunk_size)
1198        self._out.dst = self._dst_buffer
1199        self._out.size = chunk_size
1200        self._out.pos = 0

Showing the first 1,200 of 4479 lines. Download the file for the rest.

codekingpro/portable-devtools · Team Ai