codekingpro/portable-devtools
114k
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
