Team Ai
Datasetpublic

codekingpro/portable-devtools

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

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

codekingpro/portable-devtools · Team Ai