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
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.
