From 78bad0b1c4186dd21dfc01da91ff60e0c892ebee Mon Sep 17 00:00:00 2001 From: Anthonin Bonnefoy Date: Tue, 28 Jul 2026 10:00:45 +0200 Subject: Add backend support for CompressedMessages This commit introduces a new CompressedMessages message which contains: - Algorithm used. Currently only zstd or lz4 are supported - Message types in the payload. For example, 'TDD' for 1 RowDescription followed by 2 DataRows - Compressed payload Compression can be enabled by setting protocol_backend_compression to a compression algorithm. Algorithm needs to be allowed using the protocol_backend_compression_allowed_algorithms list, and supported by the client. Client will report the supported compression algorithms through protocol negotiation using _pq_.supported_compressions. Once enabled, the global PQcommMethods is replaced by PqCommCompressMethods. - putmessage will buffer messages up to protocol_backend_compression_threshold bytes. Once this threshold is reached, messages will be stored in a compressed buffer. - flush will create a new CompressedMessages and send the content of the compressed buffer. The previous PQcommMethods is used to handle the transport. --- src/backend/libpq/Makefile | 3 + src/backend/libpq/meson.build | 3 + src/backend/libpq/pqcomm_compress.c | 734 ++++++++++++++++++ src/backend/libpq/pqcomm_compress_lz4.c | 258 ++++++ src/backend/libpq/pqcomm_compress_zstd.c | 257 ++++++ src/backend/utils/misc/guc_parameters.dat | 39 + src/backend/utils/misc/guc_tables.c | 2 + src/backend/utils/misc/postgresql.conf.sample | 7 + src/include/libpq/libpq-be.h | 6 + src/include/libpq/libpq.h | 45 ++ src/include/libpq/protocol.h | 5 +- src/include/utils/guc_hooks.h | 4 + 12 files changed, 1361 insertions(+), 2 deletions(-) create mode 100644 src/backend/libpq/pqcomm_compress.c create mode 100644 src/backend/libpq/pqcomm_compress_lz4.c create mode 100644 src/backend/libpq/pqcomm_compress_zstd.c diff --git a/src/backend/libpq/Makefile b/src/backend/libpq/Makefile index 98eb2a8242d..ce58536c51f 100644 --- a/src/backend/libpq/Makefile +++ b/src/backend/libpq/Makefile @@ -26,6 +26,9 @@ OBJS = \ hba.o \ ifaddr.o \ pqcomm.o \ + pqcomm_compress.o \ + pqcomm_compress_lz4.o \ + pqcomm_compress_zstd.o \ pqformat.o \ pqmq.o \ pqsignal.o diff --git a/src/backend/libpq/meson.build b/src/backend/libpq/meson.build index 8571f652844..a7ea28d45c6 100644 --- a/src/backend/libpq/meson.build +++ b/src/backend/libpq/meson.build @@ -12,6 +12,9 @@ backend_sources += files( 'hba.c', 'ifaddr.c', 'pqcomm.c', + 'pqcomm_compress.c', + 'pqcomm_compress_zstd.c', + 'pqcomm_compress_lz4.c', 'pqformat.c', 'pqmq.c', 'pqsignal.c', diff --git a/src/backend/libpq/pqcomm_compress.c b/src/backend/libpq/pqcomm_compress.c new file mode 100644 index 00000000000..24f60e5ee01 --- /dev/null +++ b/src/backend/libpq/pqcomm_compress.c @@ -0,0 +1,734 @@ +/*------------------------------------------------------------------------- + * + * pqcomm_compress.c + * Common routines for backend protocol compression + * + * Copyright (c) 2026, PostgreSQL Global Development Group + * + * src/backend/libpq/pqcomm_compress.c + * + *------------------------------------------------------------------------- + */ +#include "postgres.h" + +#include "access/xact.h" +#include "libpq/pqformat.h" +#include "port/pg_bswap.h" +#include "common/compression.h" +#include "libpq/libpq.h" +#include "nodes/pg_list.h" +#include "storage/ipc.h" +#include "utils/guc.h" +#include "utils/guc_hooks.h" +#include "utils/memutils.h" +#include "utils/varlena.h" +#include "miscadmin.h" + +/* Force flush after a specific number of messages */ +int protocol_backend_compression_number_messages; + +/* Minimum byte threshold before compressing messages */ +int protocol_backend_compression_threshold; + +/* Force an independent compression frame for each transaction */ +bool protocol_backend_compression_transaction_frame; + +/* Bitmask of pg_compress_algorithm values allowed */ +int protocol_backend_compression_allowed_algorithms; + +/* Internal functions */ +static void pq_compress_comm_reset(void); +static int pq_compress_flush(void); +static int pq_compress_flush_if_writable(void); +static int _pq_compress_flush(void); +static bool pq_compress_is_send_pending(void); +static int pq_compress_putmessage(char msgtype, const char *s, size_t len); +static void pq_compress_putmessage_noblock(char msgtype, const char *s, size_t len); +static int _pq_compress_putmessage(char msgtype, const char *s, size_t len, bool block); +static void pq_compress_close(int code, Datum arg); +static int pq_send_messages(bool block, bool partial); + +static bool PqCompressBusy; /* busy handling message */ + +static const PQcommMethods *PrevPQcommMethod = NULL; +static const PQcompressMethods *PqCompressMethods = NULL; + +static pqcomm_compress * cs = NULL; +static const PQcommMethods PqCommCompressMethods = { + .comm_reset = pq_compress_comm_reset, + .flush = pq_compress_flush, + .flush_if_writable = pq_compress_flush_if_writable, + .is_send_pending = pq_compress_is_send_pending, + .putmessage = pq_compress_putmessage, + .putmessage_noblock = pq_compress_putmessage_noblock +}; + +/* + * init_compress_state -- initialize compression context + */ +static void +init_compress_state(void) +{ + MemoryContext oldctx = MemoryContextSwitchTo(TopMemoryContext); + + Assert(cs == NULL); + + /* First time, create the context */ + cs = palloc0_object(pqcomm_compress); + + /* Initialize compressed message types */ + initStringInfo(&cs->msgTypes); + + /* Initialize uncompressed message buffer */ + initStringInfo(&cs->inBuf); + + /* + * If we init outBuf and have the compressors enlarge it, there will be + * some wasted memory due to enlargeStringInfo doubling the buffer size. + * Thus we let the compressor itself do the init using initStringInfoExt. + */ + + MemoryContextSwitchTo(oldctx); + + cs->algorithm = PG_COMPRESSION_NONE; + cs->pending_compressed_messages = false; + cs->opened_frame = false; + PqCompressMethods = NULL; +} + +/* + * pq_send_uncompressed_messages -- Send all uncompressed messages stored in + * compress state's inBuf. + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_send_uncompressed_messages(bool block) +{ + int cursor = 0; + + Assert(PqCompressBusy); + + /* + * We should only send uncompressed messages if the compress output buffer + * is empty + */ + Assert(cs->outBuf.len == 0); + + if (cs->inBuf.len <= 0) + /* Nothing to send, exit */ + return 0; + + while (cursor < cs->inBuf.len) + { + int length; + char msgtype = cs->inBuf.data[cursor++]; + + /* Read back the length from the stored message */ + memcpy(&length, cs->inBuf.data + cursor, 4); + /* Length includes the 4-byte length, subtract it */ + length = (int) pg_ntoh32(length) - 4; + cursor += 4; + + if (block) + { + if (PrevPQcommMethod->putmessage(msgtype, cs->inBuf.data + cursor, length)) + return EOF; + } + else + PrevPQcommMethod->putmessage_noblock(msgtype, cs->inBuf.data + cursor, length); + + cursor += length; + + /* Sanity check */ + Assert(cursor <= cs->inBuf.len); + }; + + /* + * Everything was sent, we can reset the uncompressed buffer and message + * types + */ + resetStringInfo(&cs->inBuf); + resetStringInfo(&cs->msgTypes); + return 0; +} + +/* + * pq_send_messages -- Send messages, both compressed and uncompressed messages + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_send_messages(bool block, bool partial) +{ + Assert(PqCompressBusy); + + /* Send compressed payload first if we have any */ + if (pq_send_compressed_message(block, partial)) + return EOF; + + /* then send any uncompressed payload */ + return pq_send_uncompressed_messages(block); +} + +/* + * pq_send_compressed_message -- Create a new CompressedMessages with the + * currently compressed messages and send it. + * + * Partial is true if the last compressed message is incomplete and the rest + * will need to be sent with additional CompressedMessages. This typically + * happens when the output buffer becomes full while compressing a message and + * the current content is sent to reset the compress buffer. + * + * returns 0 if OK, EOF if trouble + */ +int +pq_send_compressed_message(bool block, bool partial) +{ + StringInfoData buf; + + Assert(PqCompressBusy); + + if (cs->outBuf.len <= 0) + /* No compressed payload to send */ + return 0; + + pq_beginmessage(&buf, PqMsg_CompressedMessages); + + /* Send compression algorithm used. */ + pq_sendbyte(&buf, cs->algorithm); + + /* + * Send the message types, may be 0 len if first message is partial or if + * we just close the frame. + */ + cs->msgTypes.data[cs->msgTypes.len] = '\0'; + pq_send_ascii_string(&buf, cs->msgTypes.data); + resetStringInfo(&cs->msgTypes); + + /* And send the compressed payload itself. */ + pq_sendbytes(&buf, cs->outBuf.data, cs->outBuf.len); + + /* + * We can't rely on pq_endmessage since it would call + * pq_compress_putmessage. + */ + if (block) + { + if (PrevPQcommMethod->putmessage(buf.cursor, buf.data, buf.len)) + { + pfree(buf.data); + return EOF; + } + } + else + PrevPQcommMethod->putmessage_noblock(buf.cursor, buf.data, buf.len); + pfree(buf.data); + + /* With the compressed payload sent, we can reset the output buffer */ + resetStringInfo(&cs->outBuf); + + /* + * If the message was partial, we must keep pending_compressed_messages + * set as the rest of the message will need one or more + * PqMsg_CompressedMessages. + */ + if (!partial) + cs->pending_compressed_messages = false; + return 0; +} + +/* + * _pq_compress_putmessage -- Put message in compress state's inBuf. If compression + * threshold is exceeded, the message is compressed and staged in compress state's + * outBuf. + * + * During a flush, the compressed output buffer will be sent in a single + * CompressedMessages message. If the compression threshold wasn't reached, the + * uncompressed buffered messages are sent without compression. + * + * returns 0 if OK, EOF if trouble + */ +static int +_pq_compress_putmessage(char msgtype, const char *s, size_t len, bool block) +{ + /* No-op if reentrant call */ + if (PqCompressBusy) + return 0; + PqCompressBusy = true; + + pq_sendbyte(&cs->inBuf, msgtype); + pq_sendint32(&cs->inBuf, len + 4); + /* message may be empty */ + if (len > 0) + pq_sendbytes(&cs->inBuf, s, len); + + /* + * Transactional proxies need to inspect the content of some packets. To + * avoid having the proxy decompress the whole session for that purpose, + * we force those messages to always be sent uncompressed. For + * BackendKeyData, since the content is random bytes, there's no benefit + * compressing it. + */ + if (msgtype == PqMsg_ReadyForQuery + || msgtype == PqMsg_ErrorResponse + || msgtype == PqMsg_ParameterStatus + || msgtype == PqMsg_CommandComplete + || msgtype == PqMsg_BackendKeyData) + { + if (cs->opened_frame) + { + /* + * If transaction frame is enabled, we need to close the current + * frame if we just finished a transaction. + */ + bool end_frame = protocol_backend_compression_transaction_frame + && msgtype == PqMsg_ReadyForQuery + && cs->opened_frame + && !IsTransactionBlock(); + + /* Flush any compressed payload first */ + if (PqCompressMethods->flush(cs, block, end_frame)) + goto fail; + if (pq_send_compressed_message(block, false)) + goto fail; + if (end_frame) + cs->opened_frame = false; + } + + /* And send the uncompressed message */ + if (pq_send_uncompressed_messages(block)) + goto fail; + + PqCompressBusy = false; + return 0; + } + + if (cs->pending_compressed_messages || cs->inBuf.len > protocol_backend_compression_threshold) + { + /* + * We either crossed the compress threshold, or the threshold was + * already crossed with previous messages. Compress this message. + */ + bool start_frame = cs->opened_frame == false; + + cs->pending_compressed_messages = true; + cs->opened_frame = true; + if (PqCompressMethods->compress_message(cs, block, start_frame)) + goto fail; + resetStringInfo(&cs->inBuf); + } + + /* Keep track of the message type */ + appendStringInfoChar(&cs->msgTypes, msgtype); + + /* Flush if we've crossed the number of messages threshold */ + if (protocol_backend_compression_number_messages > 0 + && cs->msgTypes.len > protocol_backend_compression_number_messages) + if (_pq_compress_flush()) + goto fail; + + PqCompressBusy = false; + return 0; + +fail: + PqCompressBusy = false; + return EOF; +} + +/* + * pq_compress_comm_reset -- reset libpq during error recovery + */ +static void +pq_compress_comm_reset(void) +{ + /* Do not throw away pending data, but do reset the busy flag */ + PqCompressBusy = false; + PrevPQcommMethod->comm_reset(); +} + +/* + * _pq_compress_flush -- flush pending messages + * + * If we have a pending CompressedMessages, force the compressor to flush any + * buffered data and send it. + * + * returns 0 if OK, EOF if trouble + */ +static int +_pq_compress_flush(void) +{ + if (cs->pending_compressed_messages + && (PqCompressMethods->flush(cs, true, false))) + return EOF; + + /* + * We may have uncompressed messages currently buffered, so call + * pq_send_messages to send both uncompressed and compressed payload. + */ + if (pq_send_messages(true, false)) + return EOF; + + return PrevPQcommMethod->flush(); +} + +/* + * pq_compress_flush -- flush pending messages + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_flush(void) +{ + /* No-op if reentrant call */ + if (PqCompressBusy) + return 0; + + PqCompressBusy = true; + if (_pq_compress_flush()) + goto fail; + + PqCompressBusy = false; + return 0; + +fail: + PqCompressBusy = false; + return EOF; +} + +/* + * pq_compress_flush_if_writable -- flush pending messages without blocking + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_flush_if_writable(void) +{ + /* No-op if reentrant call */ + if (PqCompressBusy) + return 0; + + PqCompressBusy = true; + if (cs->pending_compressed_messages + && PqCompressMethods->flush(cs, false, false)) + goto fail; + + if (pq_send_messages(false, false)) + goto fail; + + if (PrevPQcommMethod->flush_if_writable()) + goto fail; + + PqCompressBusy = false; + return 0; + +fail: + PqCompressBusy = false; + return EOF; +} + +/* + * pq_compress_is_send_pending -- is there any pending data? + */ +static bool +pq_compress_is_send_pending(void) +{ + return cs->pending_compressed_messages + || cs->inBuf.len > 0 + || PrevPQcommMethod->is_send_pending(); +} + +/* + * pq_compress_putmessage -- add message to the compress buffer + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_putmessage(char msgtype, const char *s, size_t len) +{ + return _pq_compress_putmessage(msgtype, s, len, true); +} + +/* + * pq_compress_putmessage_noblock -- add message to the compress buffer without blocking + */ +static void +pq_compress_putmessage_noblock(char msgtype, const char *s, size_t len) +{ + _pq_compress_putmessage(msgtype, s, len, false); +} + +/* + * GUC assign hook for protocol_backend_compression: switch PqCommMethods to the + * implementation matching the newly selected compression method. + */ +void +assign_protocol_backend_compression(const char *newval, void *extra) +{ + char *algorithm_name = NULL; + char *detail = NULL; + void *new_private_data = NULL; + char *error_detail; + pg_compress_algorithm algorithm; + pg_compress_specification specification; + const PQcompressMethods *new_compress_methods = NULL; + + if (MyProcPort == NULL) + /* If there's no client connection, just ignore compression settings */ + return; + + /* Extract algorithm */ + parse_compress_options(newval, &algorithm_name, &detail); + if (!parse_compress_algorithm(algorithm_name, &algorithm)) + ereport(ERROR, + (errcode(ERRCODE_SYNTAX_ERROR), + errmsg("invalid value for parameter \"protocol_backend_compression\": \"%s\"", + newval))); + + /* Extract specification */ + parse_compress_specification(algorithm, detail, &specification); + error_detail = validate_compress_specification(&specification); + if (error_detail != NULL) + ereport(ERROR, + errcode(ERRCODE_SYNTAX_ERROR), + errmsg("invalid compression specification: %s", + error_detail)); + + if ((specification.options & PG_COMPRESSION_OPTION_WORKERS) && specification.workers >= 1) + ereport(ERROR, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("protocol compression does not support compression workers"))); + + if (!cs && algorithm == PG_COMPRESSION_NONE) + /* Compression was never enabled, nothing to do */ + return; + + if ((specification.options & PG_COMPRESSION_OPTION_LONG_DISTANCE) == 0) + { + /* If long distance parameter wasn't specified, default it to false. */ + specification.options |= PG_COMPRESSION_OPTION_LONG_DISTANCE; + specification.long_distance = false; + } + + if (cs && cs->algorithm == algorithm + && cs->specification.level == specification.level + && cs->specification.long_distance == specification.long_distance) + /* No compression change */ + return; + + if (algorithm != PG_COMPRESSION_NONE && + (protocol_backend_compression_allowed_algorithms & (1 << algorithm)) == 0) + ereport(ERROR, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("%s compression is not permitted by \"protocol_backend_compression_allowed_algorithms\"", + get_compress_algorithm_name(algorithm)))); + + if (!cs) + { + /* + * First time we have a non-null compression algorithm, initialize + * compression state. + */ + init_compress_state(); + on_proc_exit(pq_compress_close, 0); + } + + /* + * Before we create the new compressor, we need to close the ongoing + * frame. On a parameter change, the current compressor may be reset to + * handle the new parameters (like zstd's long distance). If creating the + * new compressor fails, this means the frame was closed unnecessarily, + * but this shouldn't be much of an issue. + */ + if (cs->opened_frame) + { + Assert(cs->algorithm != PG_COMPRESSION_NONE); + /* We have an opened frame, close it */ + if (PqCompressMethods->flush(cs, false, true)) + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not flush compressed data"))); + cs->opened_frame = false; + /* With the end frame inserted, a CompressedMessages needs to be sent */ + cs->pending_compressed_messages = true; + } + + /* + * We also need to send all uncompressed and compressed bytes before + * switching compressor. + */ + if (cs->inBuf.len > 0 || cs->pending_compressed_messages) + { + PqCompressBusy = true; + if (pq_send_messages(false, false)) + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not send compressed data"))); + PqCompressBusy = false; + } + + switch (algorithm) + { + case PG_COMPRESSION_ZSTD: + if (!MyProcPort->supported_compress_zstd) + { + ereport(ERROR, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("zstd compression is not supported by client"))); + } + new_compress_methods = pq_init_compressor_zstd(cs, &specification, &new_private_data); + break; + case PG_COMPRESSION_LZ4: + if (!MyProcPort->supported_compress_lz4) + { + ereport(ERROR, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("lz4 compression is not supported by client"))); + } + new_compress_methods = pq_init_compressor_lz4(cs, &specification, &new_private_data); + break; + case PG_COMPRESSION_GZIP: + ereport(ERROR, + (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), + errmsg("gzip compression is not supported"))); + break; + case PG_COMPRESSION_NONE: + break; + }; + + if (cs->algorithm != algorithm && PqCompressMethods != NULL) + { + /* + * That's an algorithm change, we need to drop the previous + * compression context. + */ + PqCompressMethods->free_compress_context(cs); + } + + if (PrevPQcommMethod == NULL && new_compress_methods != NULL) + { + /* Compression is enabled, replace the global PqCommMethods */ + PrevPQcommMethod = PqCommMethods; + PqCommMethods = &PqCommCompressMethods; + } + else if (new_compress_methods == NULL && PrevPQcommMethod != NULL) + { + /* + * Compression is disabled, restore PqCommMethods to the previous + * value + */ + PqCommMethods = PrevPQcommMethod; + PrevPQcommMethod = NULL; + } + + /* Update the global compress methods */ + PqCompressMethods = new_compress_methods; + + /* And update compress state */ + cs->private_data = new_private_data; + cs->algorithm = algorithm; + cs->specification = specification; +} + +/* + * GUC check hook for protocol_backend_compression_allowed_algorithms: parse the + * comma-separated list of algorithm names and store it as a bitmask of + * pg_compress_algorithm values. + * + * Returns an error if the compression algorithm is not supported. + */ +bool +check_protocol_backend_allowed_algorithms(char **newval, void **extra, GucSource source) +{ + char *rawstring; + List *elemlist; + ListCell *l; + int flags = 0; + bool result = true; + + /* Need a modifiable copy of string */ + rawstring = pstrdup(*newval); + + if (!SplitGUCList(rawstring, ',', &elemlist)) + { + GUC_check_errdetail("Invalid list syntax in parameter \"%s\".", + "protocol_backend_compression_allowed_algorithms"); + pfree(rawstring); + list_free(elemlist); + return false; + } + + foreach(l, elemlist) + { + char *item = (char *) lfirst(l); + pg_compress_algorithm algorithm; + + if (!parse_compress_algorithm(item, &algorithm)) + { + GUC_check_errdetail("Invalid compression algorithm \"%s\".", item); + result = false; + break; + } + if (algorithm == PG_COMPRESSION_GZIP) + { + GUC_check_errdetail("gzip compression is not supported"); + result = false; + break; + } +#ifndef USE_ZSTD + if (algorithm == PG_COMPRESSION_ZSTD) + { + GUC_check_errdetail("zstd compression is not supported by this build"); + result = false; + break; + } +#endif +#ifndef USE_LZ4 + if (algorithm == PG_COMPRESSION_LZ4) + { + GUC_check_errdetail("lz4 compression is not supported by this build"); + result = false; + break; + } +#endif + flags |= (1 << algorithm); + } + + pfree(rawstring); + list_free(elemlist); + + if (!result) + return result; + + *extra = guc_malloc(LOG, sizeof(int)); + if (!*extra) + return false; + *((int *) *extra) = flags; + + return result; +} + +/* + * GUC assign hook for protocol_backend_compression_allowed_algorithms. + */ +void +assign_protocol_backend_allowed_algorithms(const char *newval, void *extra) +{ + protocol_backend_compression_allowed_algorithms = *((int *) extra); +} + +/* + * pq_compress_close -- free compressor memory at backend exit + */ +static void +pq_compress_close(int code, Datum arg) +{ + Assert(cs); + if (PqCompressMethods != NULL) + PqCompressMethods->free_compress_context(cs); + if (cs->outBuf.data != NULL) + pfree(cs->outBuf.data); + pfree(cs->msgTypes.data); + pfree(cs->inBuf.data); + pfree(cs); +} diff --git a/src/backend/libpq/pqcomm_compress_lz4.c b/src/backend/libpq/pqcomm_compress_lz4.c new file mode 100644 index 00000000000..24a2d55fc3d --- /dev/null +++ b/src/backend/libpq/pqcomm_compress_lz4.c @@ -0,0 +1,258 @@ +/*------------------------------------------------------------------------- + * + * pqcomm_compress_lz4.c + * Compress backend messages with lz4 + * + * Portions Copyright (c) 2026, PostgreSQL Global Development Group + * + * src/backend/libpq/pqcomm_compress_lz4.c + * + *------------------------------------------------------------------------- + */ + +#include "postgres.h" + +#include "libpq/libpq.h" +#include "utils/memutils.h" + +#ifndef USE_LZ4 + +const PQcompressMethods * +pq_init_compressor_lz4(pqcomm_compress * cs, pg_compress_specification *specification, void **new_private_data) +{ + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("lz4 compression is not supported by this build"))); +} + +#else +#include + +#define CHUNK_SIZE (64 * 1024) + +typedef struct pqcomm_lz4 +{ + LZ4F_cctx *ctx; + LZ4F_preferences_t prefs; +} pqcomm_lz4; + +static int pq_compress_flush_lz4(pqcomm_compress * cs, bool block, bool end_frame); +static int pq_compress_message_lz4(pqcomm_compress * cs, bool block, bool start_frame); +static void pq_compress_free_lz4(pqcomm_compress * cs); + +static const PQcompressMethods PqCompressMethodsLz4 = { + .flush = pq_compress_flush_lz4, + .compress_message = pq_compress_message_lz4, + .free_compress_context = pq_compress_free_lz4, +}; + +/* + * pq_init_compressor_lz4 -- Initialize a lz4 compressor or change compress specification + * + * On success, the compressor state will be stored in new_private_data. + * + * returns PQcompressMethods of the lz4 compressor + */ +const PQcompressMethods * +pq_init_compressor_lz4(pqcomm_compress * cs, pg_compress_specification *specification, void **new_private_data) +{ + pqcomm_lz4 *lz4_state; + + Assert(!cs->opened_frame); + if (cs->algorithm != PG_COMPRESSION_LZ4) + { + /* We have an algorithm change, create a brand-new context */ + LZ4F_errorCode_t ctxError; + size_t outbuf_size; + + MemoryContext oldctx = MemoryContextSwitchTo(TopMemoryContext); + + /* Initialize state */ + lz4_state = palloc0_object(pqcomm_lz4); + lz4_state->prefs.compressionLevel = specification->level; + + /* + * LZ4F_compressBound provides the outbuf size needed in the worst + * case scenario, which is going to be 65544 for a 64KB srcSize. As we + * want to send chunks of 64KB of compressed data, we need to double + * that size so we can call LZ4F_compressUpdate until we have 64KB + * available to send. + */ + outbuf_size = 2 * LZ4F_compressBound(CHUNK_SIZE, &lz4_state->prefs); + + /* + * Initialize or enlarge outBuf to be able to store a full chunk + */ + if (cs->outBuf.data == NULL) + initStringInfoExt(&cs->outBuf, outbuf_size); + else + enlargeStringInfo(&cs->outBuf, outbuf_size); + + MemoryContextSwitchTo(oldctx); + + /* and create the lz4 ctx */ + ctxError = LZ4F_createCompressionContext(&lz4_state->ctx, LZ4F_VERSION); + if (LZ4F_isError(ctxError)) + { + pfree(lz4_state); + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not create lz4 compression context: %s", + LZ4F_getErrorName(ctxError)))); + } + } + else + { + /* This is just a parameter change, reuse existing context */ + lz4_state = (pqcomm_lz4 *) cs->private_data; + lz4_state->prefs.compressionLevel = specification->level; + } + + *new_private_data = lz4_state; + return &PqCompressMethodsLz4; +} + +/* + * pq_compress_start_lz4_frame -- Start a new lz4 frame + */ +static void +pq_compress_start_lz4_frame(pqcomm_compress * cs) +{ + size_t ret; + pqcomm_lz4 *lz4_state = (pqcomm_lz4 *) cs->private_data; + + Assert(cs->outBuf.len == 0); + ret = LZ4F_compressBegin(lz4_state->ctx, + cs->outBuf.data + cs->outBuf.len, + cs->outBuf.maxlen - cs->outBuf.len, + &lz4_state->prefs); + if (LZ4F_isError(ret)) + { + if (cs->algorithm != PG_COMPRESSION_LZ4) + { + /* Free the context and state if they were created */ + LZ4F_freeCompressionContext(lz4_state->ctx); + pfree(lz4_state); + } + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not write lz4 header: %s", LZ4F_getErrorName(ret)))); + } + cs->outBuf.len += ret; +} + +/* + * pq_compress_message_lz4 -- Compress a message using lz4 compressor + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_message_lz4(pqcomm_compress * cs, bool block, bool start_frame) +{ + pqcomm_lz4 *lz4_state = (pqcomm_lz4 *) cs->private_data; + size_t remaining = cs->inBuf.len; + + /* Start new lz4 frame if necessary */ + if (start_frame) + pq_compress_start_lz4_frame(cs); + + Assert(cs->inBuf.cursor == 0); + + while (remaining > 0) + { + size_t written; + + /* + * To keep memory usage constrained, we cap the amount of data to + * compress to CHUNK_SIZE. + */ + size_t chunk = Min(remaining, CHUNK_SIZE); + size_t bound = LZ4F_compressBound(chunk, &lz4_state->prefs); + + if (cs->outBuf.maxlen - cs->outBuf.len < bound) + { + /* We need to free space in the outBuf, send what we have */ + if (pq_send_compressed_message(block, true)) + return EOF; + } + Assert(cs->outBuf.maxlen - cs->outBuf.len >= bound); + + written = LZ4F_compressUpdate(lz4_state->ctx, + cs->outBuf.data + cs->outBuf.len, + cs->outBuf.maxlen - cs->outBuf.len, + cs->inBuf.data + cs->inBuf.cursor, + chunk, + NULL); + if (LZ4F_isError(written)) + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not compress data: %s", LZ4F_getErrorName(written)))); + + /* Update output buffer len */ + cs->outBuf.len += written; + + /* Update input buffer with consumed chunk */ + cs->inBuf.cursor += chunk; + remaining = cs->inBuf.len - cs->inBuf.cursor; + } + + return 0; +} + +/* + * pq_compress_flush_lz4 -- flush any pending data in the compress buffer + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_flush_lz4(pqcomm_compress * cs, bool block, bool end_frame) +{ + pqcomm_lz4 *lz4_state = (pqcomm_lz4 *) cs->private_data; + size_t ret; + size_t bound; + + /* Make sure the output buffer has enough room for the flush */ + bound = LZ4F_compressBound(0, &lz4_state->prefs); + if (cs->outBuf.maxlen - cs->outBuf.len < bound) + { + if (pq_send_compressed_message(block, false)) + return EOF; + } + Assert(cs->outBuf.maxlen - cs->outBuf.len >= bound); + + if (end_frame) + { + ret = LZ4F_compressEnd(lz4_state->ctx, + cs->outBuf.data + cs->outBuf.len, + cs->outBuf.maxlen - cs->outBuf.len, + NULL); + } + else + ret = LZ4F_flush(lz4_state->ctx, + cs->outBuf.data + cs->outBuf.len, + cs->outBuf.maxlen - cs->outBuf.len, + NULL); + if (LZ4F_isError(ret)) + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not flush lz4 compressed data: %s", LZ4F_getErrorName(ret)))); + + /* Update output buffer len */ + cs->outBuf.len += ret; + + return 0; +} + +/* + * pq_compress_free_lz4 -- free the compressor context + */ +static void +pq_compress_free_lz4(pqcomm_compress * cs) +{ + pqcomm_lz4 *lz4_state = (pqcomm_lz4 *) cs->private_data; + + LZ4F_freeCompressionContext(lz4_state->ctx); + pfree(lz4_state); +} + +#endif diff --git a/src/backend/libpq/pqcomm_compress_zstd.c b/src/backend/libpq/pqcomm_compress_zstd.c new file mode 100644 index 00000000000..b062399d1b2 --- /dev/null +++ b/src/backend/libpq/pqcomm_compress_zstd.c @@ -0,0 +1,257 @@ +/*------------------------------------------------------------------------- + * + * pqcomm_compress_zstd.c + * Compress backend messages with zstd + * + * Portions Copyright (c) 2026, PostgreSQL Global Development Group + * + * src/backend/libpq/pqcomm_compress_zstd.c + * + *------------------------------------------------------------------------- + */ + +#include "postgres.h" + +#include "libpq/libpq.h" +#include "utils/memutils.h" + +#ifndef USE_ZSTD + +const PQcompressMethods * +pq_init_compressor_zstd(pqcomm_compress * cs, pg_compress_specification *specification, void **new_private_data) +{ + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("zstd compression is not supported by this build"))); +} + +#else +#include + +typedef struct pqcomm_zstd +{ + ZSTD_CCtx *cctx; +} pqcomm_zstd; + +static int pq_compress_flush_zstd(pqcomm_compress * cs, bool block, bool end_frame); +static int pq_compress_message_zstd(pqcomm_compress * cs, bool block, bool start_frame); +static void pq_compress_free_zstd(pqcomm_compress * cs); + +static const PQcompressMethods PqCompressMethodsZstd = { + .flush = pq_compress_flush_zstd, + .compress_message = pq_compress_message_zstd, + .free_compress_context = pq_compress_free_zstd, +}; + +/* + * pq_init_compressor_zstd -- Initialize a zstd compressor or change compress + * specification + * + * On success, the compressor state will be stored in new_private_data. + * + * returns PQcompressMethods of the zstd compressor + */ +const PQcompressMethods * +pq_init_compressor_zstd(pqcomm_compress * cs, pg_compress_specification *specification, void **new_private_data) +{ + int ret; + pqcomm_zstd *zstate; + + Assert(!cs->opened_frame); + if (cs->algorithm != PG_COMPRESSION_ZSTD) + { + /* We have an algorithm change, create a brand-new context */ + MemoryContext oldctx = MemoryContextSwitchTo(TopMemoryContext); + size_t outbuf_size = ZSTD_CStreamOutSize(); + + /* Initialize compress buffer */ + if (cs->outBuf.data == NULL) + initStringInfoExt(&cs->outBuf, outbuf_size); + else + enlargeStringInfo(&cs->outBuf, outbuf_size); + + zstate = palloc0_object(pqcomm_zstd); + MemoryContextSwitchTo(oldctx); + + /* Create ctx */ + zstate->cctx = ZSTD_createCCtx(); + if (!zstate->cctx) + { + pfree(zstate); + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not create zstd compression context"))); + } + } + else + { + /* This is just a parameter change, reuse existing context */ + zstate = (pqcomm_zstd *) cs->private_data; + + /* + * Some parameters like long_distance can only be changed during + * zstd's init state, so reset cctx's session to allow such changes. + */ + ZSTD_CCtx_reset(zstate->cctx, ZSTD_reset_session_and_parameters); + } + + ret = ZSTD_CCtx_setParameter(zstate->cctx, ZSTD_c_compressionLevel, + specification->level); + if (ZSTD_isError(ret)) + { + if (cs->algorithm != PG_COMPRESSION_ZSTD) + { + /* Free the context and zstate if they were created */ + ZSTD_freeCCtx(zstate->cctx); + pfree(zstate); + } + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not set zstd compression level to %d: %s", + specification->level, ZSTD_getErrorName(ret)))); + } + + + if (specification->options & PG_COMPRESSION_OPTION_LONG_DISTANCE) + { + ret = ZSTD_CCtx_setParameter(zstate->cctx, + ZSTD_c_enableLongDistanceMatching, + specification->long_distance); + if (ZSTD_isError(ret)) + { + if (cs->algorithm != PG_COMPRESSION_ZSTD) + { + /* Free the context and zstate if they were created */ + ZSTD_freeCCtx(zstate->cctx); + pfree(zstate); + } + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not set zstd long distance to %i: %s", + specification->long_distance, ZSTD_getErrorName(ret)))); + } + } + + *new_private_data = zstate; + return &PqCompressMethodsZstd; +} + + +/* + * pq_compress_message_zstd -- Compress a message using zstd compressor + * + * If the compress buffer is full, it will be sent immediately, allowing to + * reset the buffer and compress the rest of the message. + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_message_zstd(pqcomm_compress * cs, bool block, bool start_frame) +{ + ZSTD_inBuffer inBuf; + size_t yet_to_flush = 0; + ZSTD_outBuffer outBuf; + pqcomm_zstd *zstate = (pqcomm_zstd *) cs->private_data; + + inBuf.src = cs->inBuf.data; + inBuf.size = cs->inBuf.len; + inBuf.pos = 0; + + outBuf.dst = cs->outBuf.data; + outBuf.size = cs->outBuf.maxlen; + outBuf.pos = cs->outBuf.len; + + do + { + yet_to_flush = + ZSTD_compressStream2(zstate->cctx, + &outBuf, + &inBuf, ZSTD_e_continue); + if (ZSTD_isError(yet_to_flush)) + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not compress data: %s", + ZSTD_getErrorName(yet_to_flush)))); + /* Keep compress buffer in sync */ + cs->outBuf.len = outBuf.pos; + + /* + * If the output buffer is left with not enough space, send the + * compressed bytes to the underlying pqcomm, which will empty the + * buffer. + */ + if (yet_to_flush > 0) + { + if (pq_send_compressed_message(block, true) == EOF) + return EOF; + /* sync position after pq_send_compressed_message call */ + outBuf.pos = cs->outBuf.len; + } + } while (yet_to_flush > 0); + + return 0; +} + +/* + * pq_compress_flush_zstd -- flush any pending data in the compress buffer + * + * If the compress buffer is full, it will be sent immediately, allowing to + * reset the buffer and compress the rest of the data. + * + * returns 0 if OK, EOF if trouble + */ +static int +pq_compress_flush_zstd(pqcomm_compress * cs, bool block, bool end_frame) +{ + size_t yet_to_flush; + pqcomm_zstd *zstate = (pqcomm_zstd *) cs->private_data; + ZSTD_outBuffer outBuf; + + outBuf.dst = cs->outBuf.data; + outBuf.size = cs->outBuf.maxlen; + outBuf.pos = cs->outBuf.len; + + do + { + ZSTD_inBuffer in = {NULL, 0, 0}; + size_t max_needed = ZSTD_compressBound(0); + + /* + * If the output buffer is left with not enough space, send the + * compressed bytes to the underlying pqcomm. + */ + if (outBuf.size - outBuf.pos < max_needed) + { + if (pq_send_compressed_message(block, false)) + return EOF; + outBuf.pos = cs->outBuf.len; + } + + yet_to_flush = ZSTD_compressStream2(zstate->cctx, + &outBuf, + &in, end_frame ? ZSTD_e_end : ZSTD_e_flush); + /* Keep compress buffer in sync */ + cs->outBuf.len = outBuf.pos; + + if (ZSTD_isError(yet_to_flush)) + ereport(ERROR, + (errcode(ERRCODE_INTERNAL_ERROR), + errmsg("could not compress data: %s", + ZSTD_getErrorName(yet_to_flush)))); + } while (yet_to_flush > 0); + return 0; +} + +/* + * pq_compress_free_zstd -- free the compressor context + */ +static void +pq_compress_free_zstd(pqcomm_compress * cs) +{ + pqcomm_zstd *zstate = (pqcomm_zstd *) cs->private_data; + + ZSTD_freeCCtx(zstate->cctx); + pfree(zstate); +} + +#endif diff --git a/src/backend/utils/misc/guc_parameters.dat b/src/backend/utils/misc/guc_parameters.dat index c57441f7d98..9063fb77890 100644 --- a/src/backend/utils/misc/guc_parameters.dat +++ b/src/backend/utils/misc/guc_parameters.dat @@ -2430,6 +2430,45 @@ check_hook => 'check_primary_slot_name', }, +{ name => 'protocol_backend_compression', type => 'string', context => 'PGC_USERSET', group => 'CLIENT_CONN_STATEMENT', + short_desc => 'Sets the compression algorithm to use for backend messages.', + variable => 'protocol_backend_compression', + boot_val => '"none"', + assign_hook => 'assign_protocol_backend_compression', +}, + +{ name => 'protocol_backend_compression_allowed_algorithms', type => 'string', context => 'PGC_SUSET', group => 'CLIENT_CONN_STATEMENT', + short_desc => 'Sets the list of compression algorithms that clients are allowed to request for backend messages.', + long_desc => 'Only algorithms supported by the server can be set.', + flags => 'GUC_SUPERUSER_ONLY | GUC_LIST_INPUT', + variable => 'protocol_backend_compression_allowed_algorithms_string', + boot_val => '"none"', + check_hook => 'check_protocol_backend_allowed_algorithms', + assign_hook => 'assign_protocol_backend_allowed_algorithms', +}, + +{ name => 'protocol_backend_compression_number_messages', type => 'int', context => 'PGC_USERSET', group => 'CLIENT_CONN_STATEMENT', + short_desc => 'Number of messages in a compressed payload before flushing.', + variable => 'protocol_backend_compression_number_messages', + boot_val => '0', + min => '0', + max => 'INT_MAX', +}, + +{ name => 'protocol_backend_compression_threshold', type => 'int', context => 'PGC_USERSET', group => 'CLIENT_CONN_STATEMENT', + short_desc => 'Use compression when the message bytes exceeds this threshold.', + variable => 'protocol_backend_compression_threshold', + boot_val => '100', + min => '0', + max => 'INT_MAX', +}, + +{ name => 'protocol_backend_compression_transaction_frame', type => 'bool', context => 'PGC_USERSET', group => 'CLIENT_CONN_STATEMENT', + short_desc => 'Forces an independent compression frame for each transaction.', + variable => 'protocol_backend_compression_transaction_frame', + boot_val => 'false', +}, + { name => 'quote_all_identifiers', type => 'bool', context => 'PGC_USERSET', group => 'COMPAT_OPTIONS_PREVIOUS', short_desc => 'When generating SQL fragments, quote all identifiers.', variable => 'quote_all_identifiers', diff --git a/src/backend/utils/misc/guc_tables.c b/src/backend/utils/misc/guc_tables.c index 342aaeef59a..e5575aa5201 100644 --- a/src/backend/utils/misc/guc_tables.c +++ b/src/backend/utils/misc/guc_tables.c @@ -619,6 +619,8 @@ int huge_pages_status = HUGE_PAGES_UNKNOWN; static char *syslog_ident_str; static double phony_random_seed; static char *client_encoding_string; +static char *protocol_backend_compression; +static char *protocol_backend_compression_allowed_algorithms_string; static char *datestyle_string; static char *server_encoding_string; static char *server_version_string; diff --git a/src/backend/utils/misc/postgresql.conf.sample b/src/backend/utils/misc/postgresql.conf.sample index e759f06b50f..43bbdd0b64c 100644 --- a/src/backend/utils/misc/postgresql.conf.sample +++ b/src/backend/utils/misc/postgresql.conf.sample @@ -826,6 +826,13 @@ #createrole_self_grant = '' # set and/or inherit #event_triggers = on +# - Protocol Compression - +#protocol_backend_compression = 'none'; # none, zstd or lz4 +#protocol_backend_compression_allowed_algorithms = 'none'; # command-separated list of algorithms +#protocol_backend_compression_threshold = 50; +#protocol_backend_compression_number_messages = 0; +#protocol_backend_compression_transaction_frame = false; + # - Locale and Formatting - #datestyle = 'iso, mdy' diff --git a/src/include/libpq/libpq-be.h b/src/include/libpq/libpq-be.h index 921b2daa4ff..8af452dbfe9 100644 --- a/src/include/libpq/libpq-be.h +++ b/src/include/libpq/libpq-be.h @@ -159,6 +159,12 @@ typedef struct Port */ char *application_name; + /* + * Supported compression algorithms + */ + bool supported_compress_zstd; + bool supported_compress_lz4; + /* * Information that needs to be held during the authentication cycle. */ diff --git a/src/include/libpq/libpq.h b/src/include/libpq/libpq.h index d15073a0a93..d0d70f2eba7 100644 --- a/src/include/libpq/libpq.h +++ b/src/include/libpq/libpq.h @@ -18,6 +18,7 @@ #include "lib/stringinfo.h" #include "libpq/libpq-be.h" +#include "common/compression.h" /* avoid including waiteventset.h */ @@ -175,4 +176,48 @@ extern bool check_ssl_key_file_permissions(const char *ssl_key_file, bool isServerStart); extern HostsFileLoadResult load_hosts(List **hosts, char **err_msg); +/* + * declarations for variable in pqcomm_compress.c + */ + +typedef struct pqcomm_compress +{ + pg_compress_algorithm algorithm; + pg_compress_specification specification; + + /* Is there a CompressedMessages to send? */ + bool pending_compressed_messages; + /* Is there an opened compression frame? */ + bool opened_frame; + + /* Buffer containing messages before compression */ + StringInfoData inBuf; + /* Track message types added to the current frame */ + StringInfoData msgTypes; + /* Output buffer for compression */ + StringInfoData outBuf; + + /* Private data to be used by the compressor. */ + void *private_data; +} pqcomm_compress; + +extern PGDLLIMPORT int protocol_backend_compression_number_messages; +extern PGDLLIMPORT int protocol_backend_compression_threshold; +extern PGDLLIMPORT bool protocol_backend_compression_transaction_frame; +extern PGDLLIMPORT int protocol_backend_compression_allowed_algorithms; + +/* + * prototypes for functions in pqcomm_compress.c + */ + +typedef struct +{ + int (*flush) (pqcomm_compress * cs, bool block, bool end_frame); + int (*compress_message) (pqcomm_compress * cs, bool block, bool start_frame); + void (*free_compress_context) (pqcomm_compress * cs); +} PQcompressMethods; +extern int pq_send_compressed_message(bool block, bool partial); +extern const PQcompressMethods *pq_init_compressor_zstd(pqcomm_compress * cs, pg_compress_specification *specification, void **new_private_data); +extern const PQcompressMethods *pq_init_compressor_lz4(pqcomm_compress * cs, pg_compress_specification *specification, void **new_private_data); + #endif /* LIBPQ_H */ diff --git a/src/include/libpq/protocol.h b/src/include/libpq/protocol.h index 503852ca988..16784d34f6e 100644 --- a/src/include/libpq/protocol.h +++ b/src/include/libpq/protocol.h @@ -61,8 +61,9 @@ /* These are the codes sent by both the frontend and backend. */ -#define PqMsg_CopyDone 'c' -#define PqMsg_CopyData 'd' +#define PqMsg_CopyDone 'c' +#define PqMsg_CopyData 'd' +#define PqMsg_CompressedMessages 'z' /* Additional codes sent by parallel workers to leader processes. */ diff --git a/src/include/utils/guc_hooks.h b/src/include/utils/guc_hooks.h index 06453a18c03..bc1c4812f0d 100644 --- a/src/include/utils/guc_hooks.h +++ b/src/include/utils/guc_hooks.h @@ -177,5 +177,9 @@ extern bool check_synchronized_standby_slots(char **newval, void **extra, extern void assign_synchronized_standby_slots(const char *newval, void *extra); extern bool check_log_min_messages(char **newval, void **extra, GucSource source); extern void assign_log_min_messages(const char *newval, void *extra); +extern void assign_protocol_backend_compression(const char *newval, void *extra); +extern bool check_protocol_backend_allowed_algorithms(char **newval, void **extra, + GucSource source); +extern void assign_protocol_backend_allowed_algorithms(const char *newval, void *extra); #endif /* GUC_HOOKS_H */ -- 2.50.1 (Apple Git-155)