@@ -301,69 +301,99 @@ delta_delta_compressor_alloc(void)
301301 return compressor ;
302302}
303303
304- static DeltaDeltaCompressed *
305- delta_delta_from_parts (uint64 last_value , uint64 last_delta , Simple8bRleSerialized * deltas ,
306- Simple8bRleSerialized * nulls )
304+ static void *
305+ delta_delta_set_header_and_advance (uint64 last_value , uint64 last_delta , bool has_nulls ,
306+ size_t compressed_size , void * dest )
307307{
308- uint32 nulls_size = 0 ;
309- Size compressed_size ;
310- char * compressed_data ;
311- DeltaDeltaCompressed * compressed ;
312-
313- if (nulls != NULL )
314- nulls_size = simple8brle_serialized_total_size (nulls );
315-
316- compressed_size =
317- sizeof (DeltaDeltaCompressed ) + simple8brle_serialized_total_size (deltas ) + nulls_size ;
318-
319- if (!AllocSizeIsValid (compressed_size ))
320- ereport (ERROR ,
321- (errcode (ERRCODE_PROGRAM_LIMIT_EXCEEDED ),
322- errmsg ("compressed size exceeds the maximum allowed (%d)" , (int ) MaxAllocSize )));
323-
324- compressed_data = palloc (compressed_size );
325- compressed = (DeltaDeltaCompressed * ) compressed_data ;
308+ DeltaDeltaCompressed * compressed = (DeltaDeltaCompressed * ) dest ;
326309 SET_VARSIZE (& compressed -> vl_len_ , compressed_size );
327-
328310 compressed -> compression_algorithm = COMPRESSION_ALGORITHM_DELTADELTA ;
329311 compressed -> last_value = last_value ;
330312 compressed -> last_delta = last_delta ;
331- compressed -> has_nulls = nulls_size != 0 ? 1 : 0 ;
313+ compressed -> has_nulls = has_nulls ? 1 : 0 ;
314+ compressed -> padding [0 ] = 0 ;
315+ compressed -> padding [1 ] = 0 ;
316+ return (char * ) compressed + sizeof (* compressed );
317+ }
332318
333- compressed_data += sizeof ( * compressed );
334- compressed_data =
335- bytes_serialize_simple8b_and_advance ( compressed_data ,
336- simple8brle_serialized_total_size ( deltas ),
337- deltas ) ;
319+ size_t
320+ delta_delta_compressor_compressed_size ( DeltaDeltaCompressor * compressor , size_t * nulls_size_out )
321+ {
322+ size_t compressed_size = sizeof ( DeltaDeltaCompressed );
323+ size_t nulls_size_actual ;
338324
339- if (compressed -> has_nulls == 1 && nulls != NULL )
325+ /* If there are no elements, the compressed size is 0 even if there are nulls */
326+ if ((compressor -> delta_delta .num_elements +
327+ compressor -> delta_delta .num_uncompressed_elements ) == 0 )
340328 {
341- CheckCompressedData (nulls -> num_elements > deltas -> num_elements );
342- bytes_serialize_simple8b_and_advance (compressed_data , nulls_size , nulls );
329+ if (nulls_size_out != NULL )
330+ * nulls_size_out = 0 ;
331+ return 0 ;
343332 }
344333
345- return compressed ;
334+ compressed_size += simple8brle_compressor_compressed_const_size (& compressor -> delta_delta );
335+
336+ if (compressor -> has_nulls )
337+ {
338+ nulls_size_actual = simple8brle_compressor_compressed_const_size (& compressor -> nulls );
339+ compressed_size += nulls_size_actual ;
340+ if (nulls_size_out != NULL )
341+ * nulls_size_out = nulls_size_actual ;
342+ }
343+ else if (nulls_size_out != NULL )
344+ * nulls_size_out = 0 ;
345+
346+ return compressed_size ;
346347}
347348
348349void *
349350delta_delta_compressor_finish (DeltaDeltaCompressor * compressor )
350351{
351- Simple8bRleSerialized * deltas = simple8brle_compressor_finish (& compressor -> delta_delta );
352- Simple8bRleSerialized * nulls = simple8brle_compressor_finish (& compressor -> nulls );
353- DeltaDeltaCompressed * compressed ;
352+ size_t total_size = delta_delta_compressor_compressed_size (compressor , NULL );
353+ char * compressed = NULL ;
354354
355- if (deltas == NULL )
355+ if (total_size == 0 )
356356 return NULL ;
357357
358- compressed = delta_delta_from_parts (compressor -> prev_val ,
359- compressor -> prev_delta ,
360- deltas ,
361- compressor -> has_nulls ? nulls : NULL );
362-
363- Assert (compressed -> compression_algorithm == COMPRESSION_ALGORITHM_DELTADELTA );
358+ compressed = palloc (total_size );
359+ delta_delta_compressor_finish_into (compressor , compressed );
364360 return compressed ;
365361}
366362
363+ void *
364+ delta_delta_compressor_finish_into (DeltaDeltaCompressor * compressor , void * dest )
365+ {
366+ size_t data_size ;
367+ size_t nulls_size ;
368+ size_t compressed_size ;
369+ char * result = (char * ) dest ;
370+
371+ /* The compressed size includes the header and the nulls */
372+ compressed_size = delta_delta_compressor_compressed_size (compressor , & nulls_size );
373+ if (compressed_size == 0 )
374+ return dest ;
375+
376+ /* Check if the data size is valid */
377+ data_size = compressed_size - sizeof (DeltaDeltaCompressed ) - nulls_size ;
378+ Assert (compressed_size > (sizeof (DeltaDeltaCompressed ) + nulls_size ));
379+
380+ result = delta_delta_set_header_and_advance (compressor -> prev_val ,
381+ compressor -> prev_delta ,
382+ compressor -> has_nulls ,
383+ compressed_size ,
384+ dest );
385+
386+ result = simple8brle_compressor_finish_into (& compressor -> delta_delta , result , data_size );
387+
388+ if (compressor -> has_nulls )
389+ {
390+ Assert (nulls_size > 0 );
391+ result = simple8brle_compressor_finish_into (& compressor -> nulls , result , nulls_size );
392+ }
393+
394+ return result ;
395+ }
396+
367397Datum
368398tsl_deltadelta_compressor_finish (PG_FUNCTION_ARGS )
369399{
@@ -423,7 +453,11 @@ int64_decompression_iterator_init_forward(DeltaDeltaDecompressionIterator *iter,
423453 Oid element_type )
424454{
425455 StringInfoData si = { .data = compressed , .len = VARSIZE (compressed ) };
456+ CheckCompressedData (VARSIZE (compressed ) >= sizeof (DeltaDeltaCompressed ));
457+ CheckCompressedData (((DeltaDeltaCompressed * ) compressed )-> compression_algorithm ==
458+ COMPRESSION_ALGORITHM_DELTADELTA );
426459 DeltaDeltaCompressed * header = consumeCompressedData (& si , sizeof (DeltaDeltaCompressed ));
460+
427461 Simple8bRleSerialized * deltas = bytes_deserialize_simple8b_and_advance (& si );
428462
429463 const bool has_nulls = header -> has_nulls == 1 ;
@@ -456,6 +490,9 @@ int64_decompression_iterator_init_reverse(DeltaDeltaDecompressionIterator *iter,
456490 Oid element_type )
457491{
458492 StringInfoData si = { .data = compressed , .len = VARSIZE (compressed ) };
493+ CheckCompressedData (VARSIZE (compressed ) >= sizeof (DeltaDeltaCompressed ));
494+ CheckCompressedData (((DeltaDeltaCompressed * ) compressed )-> compression_algorithm ==
495+ COMPRESSION_ALGORITHM_DELTADELTA );
459496 DeltaDeltaCompressed * header = consumeCompressedData (& si , sizeof (DeltaDeltaCompressed ));
460497 Simple8bRleSerialized * deltas = bytes_deserialize_simple8b_and_advance (& si );
461498
@@ -723,20 +760,54 @@ deltadelta_compressed_recv(StringInfo buffer)
723760 uint8 has_nulls ;
724761 uint64 last_value ;
725762 uint64 last_delta ;
726- Simple8bRleSerialized * delta_deltas ;
763+ Simple8bRleSerialized * delta_delta_values ;
727764 Simple8bRleSerialized * nulls = NULL ;
728765 DeltaDeltaCompressed * compressed ;
766+ void * buf_ptr = NULL ;
767+ int delta_size = 0 ;
768+ int nulls_size = 0 ;
769+ size_t compressed_size = 0 ;
770+
771+ size_t allocated_size = 2 * sizeof (Simple8bRleSerialized ) + sizeof (DeltaDeltaCompressed ) +
772+ (buffer -> len - buffer -> cursor );
773+ compressed = palloc (allocated_size );
729774
730775 has_nulls = pq_getmsgbyte (buffer );
731776 CheckCompressedData (has_nulls == 0 || has_nulls == 1 );
732777
733778 last_value = pq_getmsgint64 (buffer );
734779 last_delta = pq_getmsgint64 (buffer );
735- delta_deltas = simple8brle_serialized_recv (buffer );
780+
781+ /* Leave space for the header, but we don't yet know the size of the compressed data */
782+ buf_ptr = (char * ) compressed + sizeof (DeltaDeltaCompressed );
783+
784+ /* Calculate the size of the delta delta values based on the number of bytes read */
785+ delta_size = buffer -> cursor ;
786+ buf_ptr = simple8brle_serialized_recv_into (buffer , buf_ptr , & delta_delta_values );
787+ delta_size = buffer -> cursor - delta_size ;
788+
736789 if (has_nulls )
737- nulls = simple8brle_serialized_recv (buffer );
790+ {
791+ /* Calculate the size of the nulls based on the number of bytes read */
792+ nulls_size = buffer -> cursor ;
793+ buf_ptr = simple8brle_serialized_recv_into (buffer , buf_ptr , & nulls );
794+ nulls_size = buffer -> cursor - nulls_size ;
795+ CheckCompressedData (delta_delta_values -> num_elements < nulls -> num_elements );
796+ }
797+
798+ compressed_size = sizeof (DeltaDeltaCompressed ) + delta_size + nulls_size ;
799+
800+ if (!AllocSizeIsValid (compressed_size ))
801+ ereport (ERROR ,
802+ (errcode (ERRCODE_PROGRAM_LIMIT_EXCEEDED ),
803+ errmsg ("compressed size exceeds the maximum allowed (%d)" , (int ) MaxAllocSize )));
738804
739- compressed = delta_delta_from_parts (last_value , last_delta , delta_deltas , nulls );
805+ /* Set header but don't change the buffer pointer */
806+ delta_delta_set_header_and_advance (last_value ,
807+ last_delta ,
808+ has_nulls ,
809+ compressed_size ,
810+ compressed );
740811
741812 PG_RETURN_POINTER (compressed );
742813}
0 commit comments