From ba7fcbbd85d36f79d0b102016c0d2c86fff6f747 Mon Sep 17 00:00:00 2001 From: "Chao Li (Evan)" Date: Fri, 7 Aug 2026 17:09:55 +0800 Subject: [PATCH v7] Fix detection of truncated compressed backups pg_verifybackup failed to detect truncated compressed tar backups for zstd, lz4 and gzip compression. astreamer_zstd_decompressor checked whether ZSTD_decompressStream() returned an error. However, a positive return value at end-of-input means that the final zstd frame has not completed. Record the state and verify at finalization that zstd completed the frame. Also flush any pending internally buffered output before performing that check. astreamer_gzip_decompressor did not check whether inflate() had reached Z_STREAM_END before finalizing. Record whether inflate() completed the gzip stream, and reject finalization if it did not. astreamer_lz4_decompressor did not check whether LZ4F_decompress() completed the final frame before finalizing. Record the status of decompression and reject finalization unless the frame completed. This prevents pg_verifybackup from accepting a truncated compressed tar backup. Author: Chao Li Co-authored-by: Daniel Gustafsson Reviewed-by: Osama Abdul Qader Reviewed-by: Zsolt Parragi Reviewed-by: Japin Li Discussion: https://postgr.es/m/5962B878-C43D-4EBC-9E95-1F945CE5E586@gmail.com --- src/bin/pg_verifybackup/t/010_client_untar.pl | 13 ++++++ src/fe_utils/astreamer_gzip.c | 13 +++++- src/fe_utils/astreamer_lz4.c | 14 ++++++ src/fe_utils/astreamer_zstd.c | 46 +++++++++++++++++++ src/include/fe_utils/astreamer.h | 8 ++++ src/tools/pgindent/typedefs.list | 1 + 6 files changed, 94 insertions(+), 1 deletion(-) diff --git a/src/bin/pg_verifybackup/t/010_client_untar.pl b/src/bin/pg_verifybackup/t/010_client_untar.pl index db8b96ceb43..438658b52cb 100644 --- a/src/bin/pg_verifybackup/t/010_client_untar.pl +++ b/src/bin/pg_verifybackup/t/010_client_untar.pl @@ -138,6 +138,19 @@ for my $tc (@test_configuration) [ 'pg_verifybackup', '--exit-on-error', $backup_path, ], "verify backup, compression $method"); + if ($method ne 'none') + { + my $flen = -s $backup_path . "/" . $tc->{'backup_archive'}; + ok( truncate( + $backup_path . "/" . $tc->{'backup_archive'}, + $flen - 1), + "file truncated"); + $primary->command_fails_like( + [ 'pg_verifybackup', '--exit-on-error', $backup_path, ], + qr/could not decompress data/, + "backup is corrupted, compression $method"); + } + # Cleanup. rmtree($backup_path); } diff --git a/src/fe_utils/astreamer_gzip.c b/src/fe_utils/astreamer_gzip.c index bc3d53076e1..53934a9b2dc 100644 --- a/src/fe_utils/astreamer_gzip.c +++ b/src/fe_utils/astreamer_gzip.c @@ -48,6 +48,7 @@ typedef struct astreamer_gzip_decompressor astreamer base; z_stream zstream; size_t bytes_written; + astreamer_decompression_state state; } astreamer_gzip_decompressor; static void astreamer_gzip_writer_content(astreamer *streamer, @@ -270,6 +271,7 @@ astreamer_gzip_decompressor_new(astreamer *next) if (inflateInit2(zs, 15 + 16) != Z_OK) pg_fatal("could not initialize compression library"); + streamer->state = ASTREAMER_STREAM_NEW; return &streamer->base; #else pg_fatal("this build does not support compression with %s", "gzip"); @@ -318,7 +320,11 @@ astreamer_gzip_decompressor_content(astreamer *streamer, */ res = inflate(zs, Z_NO_FLUSH); - if (res != Z_OK && res != Z_STREAM_END && res != Z_BUF_ERROR) + if (res == Z_STREAM_END) + mystreamer->state = ASTREAMER_FRAME_COMPLETE; + else if (res == Z_OK || res == Z_BUF_ERROR) + mystreamer->state = ASTREAMER_FRAME_INCOMPLETE; + else pg_fatal("could not decompress data: %s", zs->msg ? zs->msg : "unknown error"); @@ -346,6 +352,11 @@ astreamer_gzip_decompressor_finalize(astreamer *streamer) mystreamer = (astreamer_gzip_decompressor *) streamer; + if (unlikely(mystreamer->state == ASTREAMER_STREAM_NEW)) + pg_fatal("could not decompress data: compressed stream is empty"); + else if (mystreamer->state != ASTREAMER_FRAME_COMPLETE) + pg_fatal("could not decompress data: compressed stream is incomplete"); + /* * End of the stream, if there is some pending data in output buffers then * we must forward it to next streamer. diff --git a/src/fe_utils/astreamer_lz4.c b/src/fe_utils/astreamer_lz4.c index 12dfde2c837..949e0cc1232 100644 --- a/src/fe_utils/astreamer_lz4.c +++ b/src/fe_utils/astreamer_lz4.c @@ -25,6 +25,7 @@ #include "fe_utils/astreamer.h" #ifdef USE_LZ4 + typedef struct astreamer_lz4_frame { astreamer base; @@ -35,6 +36,7 @@ typedef struct astreamer_lz4_frame size_t bytes_written; bool header_written; + astreamer_decompression_state state; } astreamer_lz4_frame; static void astreamer_lz4_compressor_content(astreamer *streamer, @@ -297,6 +299,7 @@ astreamer_lz4_decompressor_new(astreamer *next) pg_fatal("could not initialize compression library: %s", LZ4F_getErrorName(ctxError)); + streamer->state = ASTREAMER_STREAM_NEW; return &streamer->base; #else pg_fatal("this build does not support compression with %s", "LZ4"); @@ -358,6 +361,12 @@ astreamer_lz4_decompressor_content(astreamer *streamer, pg_fatal("could not decompress data: %s", LZ4F_getErrorName(ret)); + /* The frame is only done when LZ4F_decompress returns 0 */ + if (ret) + mystreamer->state = ASTREAMER_FRAME_INCOMPLETE; + else + mystreamer->state = ASTREAMER_FRAME_COMPLETE; + /* Update input buffer based on number of bytes consumed */ avail_in -= read_size; next_in += read_size; @@ -395,6 +404,11 @@ astreamer_lz4_decompressor_finalize(astreamer *streamer) mystreamer = (astreamer_lz4_frame *) streamer; + if (unlikely(mystreamer->state == ASTREAMER_STREAM_NEW)) + pg_fatal("could not decompress data: compressed stream is empty"); + else if (mystreamer->state != ASTREAMER_FRAME_COMPLETE) + pg_fatal("could not decompress data: compressed stream is incomplete"); + /* * End of the stream, if there is some pending data in output buffers then * we must forward it to next streamer. diff --git a/src/fe_utils/astreamer_zstd.c b/src/fe_utils/astreamer_zstd.c index 98e8a700efe..d7b0b2778c8 100644 --- a/src/fe_utils/astreamer_zstd.c +++ b/src/fe_utils/astreamer_zstd.c @@ -33,6 +33,7 @@ typedef struct astreamer_zstd_frame ZSTD_CCtx *cctx; ZSTD_DCtx *dctx; ZSTD_outBuffer zstd_outBuf; + astreamer_decompression_state state; } astreamer_zstd_frame; static void astreamer_zstd_compressor_content(astreamer *streamer, @@ -280,6 +281,7 @@ astreamer_zstd_decompressor_new(astreamer *next) streamer->zstd_outBuf.size = streamer->base.bbs_buffer.maxlen; streamer->zstd_outBuf.pos = 0; + streamer->state = ASTREAMER_STREAM_NEW; return &streamer->base; #else pg_fatal("this build does not support compression with %s", "ZSTD"); @@ -329,6 +331,12 @@ astreamer_zstd_decompressor_content(astreamer *streamer, if (ZSTD_isError(ret)) pg_fatal("could not decompress data: %s", ZSTD_getErrorName(ret)); + + /* The frame is only done when ZSTD_decompressStream returns 0 */ + if (ret) + mystreamer->state = ASTREAMER_FRAME_INCOMPLETE; + else + mystreamer->state = ASTREAMER_FRAME_COMPLETE; } } @@ -340,6 +348,44 @@ astreamer_zstd_decompressor_finalize(astreamer *streamer) { astreamer_zstd_frame *mystreamer = (astreamer_zstd_frame *) streamer; + /* + * A full output buffer with a positive return value might leave data in + * zstd's internal buffers. Call the decompressor with empty input until + * it has flushed that data. + */ + while (mystreamer->state == ASTREAMER_FRAME_INCOMPLETE && + mystreamer->zstd_outBuf.pos == mystreamer->zstd_outBuf.size) + { + ZSTD_inBuffer empty = {NULL, 0, 0}; + size_t ret; + + astreamer_content(mystreamer->base.bbs_next, NULL, + mystreamer->zstd_outBuf.dst, + mystreamer->zstd_outBuf.pos, ASTREAMER_UNKNOWN); + + mystreamer->zstd_outBuf.dst = mystreamer->base.bbs_buffer.data; + mystreamer->zstd_outBuf.size = mystreamer->base.bbs_buffer.maxlen; + mystreamer->zstd_outBuf.pos = 0; + + ret = ZSTD_decompressStream(mystreamer->dctx, + &mystreamer->zstd_outBuf, &empty); + + if (ZSTD_isError(ret)) + pg_fatal("could not decompress data: %s", + ZSTD_getErrorName(ret)); + + /* The frame is only done when ZSTD_decompressStream returns 0 */ + if (ret) + mystreamer->state = ASTREAMER_FRAME_INCOMPLETE; + else + mystreamer->state = ASTREAMER_FRAME_COMPLETE; + } + + if (unlikely(mystreamer->state == ASTREAMER_STREAM_NEW)) + pg_fatal("could not decompress data: compressed stream is empty"); + else if (mystreamer->state != ASTREAMER_FRAME_COMPLETE) + pg_fatal("could not decompress data: compressed stream is incomplete"); + /* * End of the stream, if there is some pending data in output buffers then * we must forward it to next streamer. diff --git a/src/include/fe_utils/astreamer.h b/src/include/fe_utils/astreamer.h index 8329e4efbc5..e7f1bf2cd45 100644 --- a/src/include/fe_utils/astreamer.h +++ b/src/include/fe_utils/astreamer.h @@ -68,6 +68,14 @@ typedef enum ASTREAMER_ARCHIVE_TRAILER, } astreamer_archive_context; +/* State of the most recently processed compressed frame. */ +typedef enum +{ + ASTREAMER_STREAM_NEW, + ASTREAMER_FRAME_INCOMPLETE, + ASTREAMER_FRAME_COMPLETE, +} astreamer_decompression_state; + /* * Each chunk of data that is classified as ASTREAMER_MEMBER_HEADER, * ASTREAMER_MEMBER_CONTENTS, or ASTREAMER_MEMBER_TRAILER should also diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list index 8e3213954d8..463a9f26645 100644 --- a/src/tools/pgindent/typedefs.list +++ b/src/tools/pgindent/typedefs.list @@ -3615,6 +3615,7 @@ array_unnest_fctx assign_collations_context astreamer astreamer_archive_context +astreamer_decompression_state astreamer_extractor astreamer_gzip_decompressor astreamer_gzip_writer -- 2.50.1 (Apple Git-155)