From 0e33f84f38b03fa6372af6998499f40160398beb Mon Sep 17 00:00:00 2001 From: TapTap Date: Wed, 16 Sep 2026 23:27:34 +0200 Subject: [PATCH] feat(codec): implement md4/sha1/none digests and lz4/zlib/zlibx codecs Add real implementations for the rsync 3.4.1 checksum and compression breadth: a self-contained MD4 (RFC 1320), OpenSSL-backed SHA1, a no-digest mode, and LZ4/zlib codecs alongside zstd. Compressed buffers are now self-describing (a leading codec id), so every existing decompression call site keeps working through a process-global codec selection. zlibx shares the zlib codec because FastSync compresses only delta/token bytes (never matched file data), matching the 'x' intent. --- CMakeLists.txt | 20 ++- shell.nix | 2 + src/shared/checksum.c | 231 +++++++++++++++++++++++++--- src/shared/checksum.h | 45 ++++-- src/shared/compression.c | 322 ++++++++++++++++++++++++++++++++++++--- src/shared/compression.h | 52 ++++++- 6 files changed, 606 insertions(+), 66 deletions(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index 727150b..999a8c5 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1,6 +1,6 @@ cmake_minimum_required(VERSION 3.22) -project(FastFileTransfer VERSION 2.23.0) +project(FastFileTransfer VERSION 2.26.0) set(CMAKE_EXPORT_COMPILE_COMMANDS ON) set(CMAKE_C_STANDARD 11) @@ -68,6 +68,16 @@ if(NOT ZSTD_LIBRARY) message(FATAL_ERROR "zstd library not found. Ensure it is in your nix-shell!") endif() +find_library(ZLIB_LIBRARY z) +if(NOT ZLIB_LIBRARY) + message(FATAL_ERROR "zlib library not found. Ensure zlib1g-dev / nix zlib is available!") +endif() + +find_library(LZ4_LIBRARY lz4) +if(NOT LZ4_LIBRARY) + message(FATAL_ERROR "lz4 library not found. Ensure liblz4-dev / nix lz4 is available!") +endif() + find_package(OpenSSL REQUIRED) # --- Explicit source lists --- @@ -135,8 +145,8 @@ set(CLIENT_MAIN_SRCS src/client/client_cli.c) # --- Library targets --- add_library(fastsync_shared STATIC ${SHARED_SRCS}) target_include_directories(fastsync_shared PUBLIC src/shared) -target_link_libraries(fastsync_shared PUBLIC Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL - OpenSSL::Crypto xxhash) +target_link_libraries(fastsync_shared PUBLIC Threads::Threads ${ZSTD_LIBRARY} ${ZLIB_LIBRARY} + ${LZ4_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash) add_library(fastsync_client_core STATIC ${CLIENT_CORE_SRCS}) target_include_directories(fastsync_client_core PUBLIC src/client) @@ -275,7 +285,7 @@ if(ENABLE_FUZZ) target_include_directories(${FUZZ_NAME} PRIVATE tests src/shared src/server) target_compile_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined -fno-omit-frame-pointer) target_link_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined) - target_link_libraries(${FUZZ_NAME} PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL - OpenSSL::Crypto xxhash) + target_link_libraries(${FUZZ_NAME} PRIVATE Threads::Threads ${ZSTD_LIBRARY} ${ZLIB_LIBRARY} + ${LZ4_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash) endforeach() endif() diff --git a/shell.nix b/shell.nix index 6d03b8d..b32c25f 100644 --- a/shell.nix +++ b/shell.nix @@ -38,6 +38,8 @@ pkgs.mkShell { buildInputs = with pkgs; [ zstd + zlib + lz4 openssl ]; diff --git a/src/shared/checksum.c b/src/shared/checksum.c index b95a8bc..71426e0 100644 --- a/src/shared/checksum.c +++ b/src/shared/checksum.c @@ -7,6 +7,160 @@ * for the whole binary; this TU only needs the declarations. */ #include +/* --------------------------------------------------------------------------- + * Self-contained MD4 (RFC 1320). OpenSSL's MD4 lives in the legacy provider + * and is not guaranteed present, so FastSync carries its own implementation to + * keep --checksum-choice=md4 working on every build. + * ------------------------------------------------------------------------- */ + +typedef struct { + uint32_t state[4]; + uint64_t bit_count; + uint8_t buffer[64]; + size_t buffer_len; +} Md4Ctx; + +static uint32_t md4_rotl(uint32_t x, int n) { + return (x << n) | (x >> (32 - n)); +} + +static void md4_transform(uint32_t state[4], const uint8_t block[64]) { + uint32_t x[16]; + for (int i = 0; i < 16; i++) + x[i] = (uint32_t)block[i * 4] | ((uint32_t)block[i * 4 + 1] << 8) | + ((uint32_t)block[i * 4 + 2] << 16) | ((uint32_t)block[i * 4 + 3] << 24); + + uint32_t a = state[0], b = state[1], c = state[2], d = state[3]; + +#define F(x, y, z) (((x) & (y)) | (~(x) & (z))) +#define G(x, y, z) (((x) & (y)) | ((x) & (z)) | ((y) & (z))) +#define H(x, y, z) ((x) ^ (y) ^ (z)) +#define ROUND1(a, b, c, d, k, s) a = md4_rotl(a + F(b, c, d) + x[k], s) +#define ROUND2(a, b, c, d, k, s) a = md4_rotl(a + G(b, c, d) + x[k] + 0x5a827999u, s) +#define ROUND3(a, b, c, d, k, s) a = md4_rotl(a + H(b, c, d) + x[k] + 0x6ed9eba1u, s) + + ROUND1(a, b, c, d, 0, 3); + ROUND1(d, a, b, c, 1, 7); + ROUND1(c, d, a, b, 2, 11); + ROUND1(b, c, d, a, 3, 19); + ROUND1(a, b, c, d, 4, 3); + ROUND1(d, a, b, c, 5, 7); + ROUND1(c, d, a, b, 6, 11); + ROUND1(b, c, d, a, 7, 19); + ROUND1(a, b, c, d, 8, 3); + ROUND1(d, a, b, c, 9, 7); + ROUND1(c, d, a, b, 10, 11); + ROUND1(b, c, d, a, 11, 19); + ROUND1(a, b, c, d, 12, 3); + ROUND1(d, a, b, c, 13, 7); + ROUND1(c, d, a, b, 14, 11); + ROUND1(b, c, d, a, 15, 19); + + ROUND2(a, b, c, d, 0, 3); + ROUND2(d, a, b, c, 4, 5); + ROUND2(c, d, a, b, 8, 9); + ROUND2(b, c, d, a, 12, 13); + ROUND2(a, b, c, d, 1, 3); + ROUND2(d, a, b, c, 5, 5); + ROUND2(c, d, a, b, 9, 9); + ROUND2(b, c, d, a, 13, 13); + ROUND2(a, b, c, d, 2, 3); + ROUND2(d, a, b, c, 6, 5); + ROUND2(c, d, a, b, 10, 9); + ROUND2(b, c, d, a, 14, 13); + ROUND2(a, b, c, d, 3, 3); + ROUND2(d, a, b, c, 7, 5); + ROUND2(c, d, a, b, 11, 9); + ROUND2(b, c, d, a, 15, 13); + + ROUND3(a, b, c, d, 0, 3); + ROUND3(d, a, b, c, 8, 9); + ROUND3(c, d, a, b, 4, 11); + ROUND3(b, c, d, a, 12, 15); + ROUND3(a, b, c, d, 2, 3); + ROUND3(d, a, b, c, 10, 9); + ROUND3(c, d, a, b, 6, 11); + ROUND3(b, c, d, a, 14, 15); + ROUND3(a, b, c, d, 1, 3); + ROUND3(d, a, b, c, 9, 9); + ROUND3(c, d, a, b, 5, 11); + ROUND3(b, c, d, a, 13, 15); + ROUND3(a, b, c, d, 3, 3); + ROUND3(d, a, b, c, 11, 9); + ROUND3(c, d, a, b, 7, 11); + ROUND3(b, c, d, a, 15, 15); + +#undef F +#undef G +#undef H +#undef ROUND1 +#undef ROUND2 +#undef ROUND3 + + state[0] += a; + state[1] += b; + state[2] += c; + state[3] += d; +} + +static void md4_init(Md4Ctx* ctx) { + ctx->state[0] = 0x67452301u; + ctx->state[1] = 0xefcdab89u; + ctx->state[2] = 0x98badcfeu; + ctx->state[3] = 0x10325476u; + ctx->bit_count = 0; + ctx->buffer_len = 0; +} + +static void md4_update(Md4Ctx* ctx, const uint8_t* data, size_t len) { + ctx->bit_count += (uint64_t)len * 8; + while (len > 0) { + size_t space = sizeof(ctx->buffer) - ctx->buffer_len; + size_t take = len < space ? len : space; + memcpy(ctx->buffer + ctx->buffer_len, data, take); + ctx->buffer_len += take; + data += take; + len -= take; + if (ctx->buffer_len == sizeof(ctx->buffer)) { + md4_transform(ctx->state, ctx->buffer); + ctx->buffer_len = 0; + } + } +} + +static void md4_final(Md4Ctx* ctx, uint8_t out[16]) { + uint64_t bit_count = ctx->bit_count; + uint8_t pad = 0x80; + md4_update(ctx, &pad, 1); + uint8_t zero = 0; + while (ctx->buffer_len != 56) + md4_update(ctx, &zero, 1); + uint8_t length_le[8]; + for (int i = 0; i < 8; i++) + length_le[i] = (uint8_t)((bit_count >> (8 * i)) & 0xff); + md4_update(ctx, length_le, sizeof(length_le)); + for (int i = 0; i < 4; i++) { + out[i * 4] = (uint8_t)(ctx->state[i] & 0xff); + out[i * 4 + 1] = (uint8_t)((ctx->state[i] >> 8) & 0xff); + out[i * 4 + 2] = (uint8_t)((ctx->state[i] >> 16) & 0xff); + out[i * 4 + 3] = (uint8_t)((ctx->state[i] >> 24) & 0xff); + } +} + +/* One-shot EVP digest (md5/sha1). Returns false when OpenSSL refuses. */ +static bool evp_digest(const EVP_MD* md, const void* data, size_t size, uint8_t* out, + size_t out_capacity, size_t* out_len) { + static const uint8_t empty = 0; + const void* input = data ? data : ∅ + unsigned int digest_len = 0; + if (EVP_Digest(input, size, out, &digest_len, md, NULL) != 1) + return false; + if (digest_len > out_capacity) + return false; + *out_len = digest_len; + return true; +} + bool checksum_digest(ChecksumAlgo algo, uint64_t seed, const void* data, size_t size, uint8_t* out, size_t out_capacity, size_t* out_len) { if (!out || !out_len || out_capacity < CHECKSUM_MAX_DIGEST_LEN) @@ -14,43 +168,45 @@ bool checksum_digest(ChecksumAlgo algo, uint64_t seed, const void* data, size_t if (data == NULL && size != 0) return false; - if (algo == CHECKSUM_ALGO_XXH64) { + switch (algo) { + case CHECKSUM_ALGO_XXH64: { uint64_t digest = XXH64(data, size, seed); memcpy(out, &digest, sizeof(digest)); *out_len = sizeof(digest); return true; } - - if (algo == CHECKSUM_ALGO_XXH3) { + case CHECKSUM_ALGO_XXH3: { uint64_t digest = XXH3_64bits_withSeed(data, size, seed); memcpy(out, &digest, sizeof(digest)); *out_len = sizeof(digest); return true; } - - if (algo == CHECKSUM_ALGO_XXH128) { + case CHECKSUM_ALGO_XXH128: { XXH128_hash_t digest = XXH3_128bits_withSeed(data, size, seed); memcpy(out, &digest, sizeof(digest)); *out_len = sizeof(digest); return true; } - - if (algo == CHECKSUM_ALGO_MD5) { + case CHECKSUM_ALGO_MD5: /* md5 takes no seed; the caller's seed is deliberately ignored (documented - * in RSYNC_COMPAT.md). OpenSSL's one-shot EVP_Digest needs a non-NULL - * buffer even for an empty input, so map a NULL data + size==0 to an empty - * buffer. */ - static const uint8_t empty = 0; - const void* input = data ? data : ∅ - unsigned int digest_len = 0; - if (EVP_Digest(input, size, out, &digest_len, EVP_md5(), NULL) != 1) - return false; - if (digest_len > out_capacity) - return false; - *out_len = digest_len; + * in RSYNC_COMPAT.md). */ + return evp_digest(EVP_md5(), data, size, out, out_capacity, out_len); + case CHECKSUM_ALGO_MD4: { + Md4Ctx ctx; + md4_init(&ctx); + md4_update(&ctx, (const uint8_t*)data, size); + md4_final(&ctx, out); + *out_len = 16; + return true; + } + case CHECKSUM_ALGO_SHA1: + /* sha1 takes no seed; the caller's seed is deliberately ignored. */ + return evp_digest(EVP_sha1(), data, size, out, out_capacity, out_len); + case CHECKSUM_ALGO_NONE: + /* No checksum requested: an empty digest is the successful result. */ + *out_len = 0; return true; } - return false; } @@ -65,6 +221,12 @@ int checksum_algo_from_name(const char* name) { return (int)CHECKSUM_ALGO_XXH128; if (strcasecmp(name, "md5") == 0) return (int)CHECKSUM_ALGO_MD5; + if (strcasecmp(name, "md4") == 0) + return (int)CHECKSUM_ALGO_MD4; + if (strcasecmp(name, "sha1") == 0) + return (int)CHECKSUM_ALGO_SHA1; + if (strcasecmp(name, "none") == 0) + return (int)CHECKSUM_ALGO_NONE; return -1; } @@ -78,13 +240,21 @@ const char* checksum_algo_name(ChecksumAlgo algo) { return "xxh128"; case CHECKSUM_ALGO_MD5: return "md5"; + case CHECKSUM_ALGO_MD4: + return "md4"; + case CHECKSUM_ALGO_SHA1: + return "sha1"; + case CHECKSUM_ALGO_NONE: + return "none"; } return ""; } bool checksum_algo_valid(int algo) { return algo == (int)CHECKSUM_ALGO_XXH64 || algo == (int)CHECKSUM_ALGO_MD5 || - algo == (int)CHECKSUM_ALGO_XXH3 || algo == (int)CHECKSUM_ALGO_XXH128; + algo == (int)CHECKSUM_ALGO_XXH3 || algo == (int)CHECKSUM_ALGO_XXH128 || + algo == (int)CHECKSUM_ALGO_MD4 || algo == (int)CHECKSUM_ALGO_SHA1 || + algo == (int)CHECKSUM_ALGO_NONE; } uint8_t checksum_digest_len(ChecksumAlgo algo) { @@ -94,7 +264,26 @@ uint8_t checksum_digest_len(ChecksumAlgo algo) { return 8; case CHECKSUM_ALGO_XXH128: case CHECKSUM_ALGO_MD5: + case CHECKSUM_ALGO_MD4: return 16; + case CHECKSUM_ALGO_SHA1: + return 20; + case CHECKSUM_ALGO_NONE: + return 0; } return 0; -} \ No newline at end of file +} + +ChecksumAlgo checksum_negotiate_default(void) { + /* rsync 3.4.1 default preference order; every entry is compiled in, so this + * resolves to xxh128. */ + static const ChecksumAlgo preference[] = { + CHECKSUM_ALGO_XXH128, CHECKSUM_ALGO_XXH3, CHECKSUM_ALGO_XXH64, CHECKSUM_ALGO_MD5, + CHECKSUM_ALGO_MD4, CHECKSUM_ALGO_SHA1, CHECKSUM_ALGO_NONE, + }; + for (size_t i = 0; i < sizeof(preference) / sizeof(preference[0]); i++) { + if (checksum_algo_valid((int)preference[i])) + return preference[i]; + } + return CHECKSUM_ALGO_XXH64; +} diff --git a/src/shared/checksum.h b/src/shared/checksum.h index c323730..e2468d0 100644 --- a/src/shared/checksum.h +++ b/src/shared/checksum.h @@ -8,25 +8,35 @@ /* Whole-file content-digest algorithms selectable with --checksum-choice and * seeded with --checksum-seed. The ids are the values actually placed on the * wire (config frame), so they must be kept stable and validated on receive. - * CHECKSUM_ALGO_XXH64 == 0 is the default and is byte-for-byte what FastSync - * computed before these options existed (xxHash64 with seed 0). The set mirrors - * the algorithms rsync 3.4.1 can be built with; the ones FastSync does not - * implement (md4, sha1, none) are rejected by name at parse time. */ + * CHECKSUM_ALGO_XXH64 == 0 is the historical FastSync default and its numeric + * value is preserved. The full set mirrors the algorithms rsync 3.4.1 can be + * built with; every one of them is implemented here. */ typedef enum { CHECKSUM_ALGO_XXH64 = 0, CHECKSUM_ALGO_MD5 = 1, CHECKSUM_ALGO_XXH3 = 2, - CHECKSUM_ALGO_XXH128 = 3 + CHECKSUM_ALGO_XXH128 = 3, + CHECKSUM_ALGO_MD4 = 4, + CHECKSUM_ALGO_SHA1 = 5, + CHECKSUM_ALGO_NONE = 6 } ChecksumAlgo; -/* xxh128 digest is 16 bytes, the longest supported. */ -#define CHECKSUM_MAX_DIGEST_LEN 16 +/* FastSync's negotiated default (rsync 3.4.1 auto-negotiates xxh128 first). + * The wire default for Config->checksum_algo is this value. */ +#define CHECKSUM_ALGO_DEFAULT CHECKSUM_ALGO_XXH128 + +/* sha1 digest is 20 bytes, the longest supported. */ +#define CHECKSUM_MAX_DIGEST_LEN 20 /* Compute the whole-file digest of the first `size` bytes of `data`. * * - CHECKSUM_ALGO_XXH64: xxHash64(data, size, seed) (full 64-bit seed). - * - CHECKSUM_ALGO_MD5: md5(data, size) via OpenSSL EVP. - * md5 has no seed, so `seed` is ignored (documented). + * - CHECKSUM_ALGO_XXH3: XXH3_64bits_withSeed(data, size, seed). + * - CHECKSUM_ALGO_XXH128: XXH3_128bits_withSeed(data, size, seed). + * - CHECKSUM_ALGO_MD5: md5(data, size) via OpenSSL EVP (seed ignored). + * - CHECKSUM_ALGO_MD4: md4(data, size), self-contained RFC 1320 (seed ignored). + * - CHECKSUM_ALGO_SHA1: sha1(data, size) via OpenSSL EVP (seed ignored). + * - CHECKSUM_ALGO_NONE: no digest; *out_len is 0 and nothing is written. * - `size == 0` hashes the empty input (plus its seed), not a NULL input. * * Writes up to `out_capacity` bytes into `out`, storing the digest length in @@ -36,10 +46,9 @@ bool checksum_digest(ChecksumAlgo algo, uint64_t seed, const void* data, size_t size_t out_capacity, size_t* out_len); /* Resolve a --checksum-choice string (case-insensitive) to an algorithm id. - * Accepts "xxh64"/"xxhash", "xxh3", "xxh128" and "md5". "auto", rsync's - * default automatic choice, is resolved to the default by the caller (it is not - * a distinct algorithm here). Returns -1 for any name FastSync does not - * implement (md4/sha1/none included). */ + * Accepts "xxh64"/"xxhash", "xxh3", "xxh128", "md5", "md4", "sha1", "none". + * "auto" is not an algorithm here; the caller resolves it to the negotiated + * default. Returns -1 for any unrecognized name. */ int checksum_algo_from_name(const char* name); /* Canonical name of an algorithm (used in CLI error messages). */ @@ -48,7 +57,13 @@ const char* checksum_algo_name(ChecksumAlgo algo); /* True when `algo` is a supported id (used by config receive validation). */ bool checksum_algo_valid(int algo); -/* Digest length in bytes for an algorithm (xxh64/xxh3 = 8, md5/xxh128 = 16). */ +/* Digest length in bytes for an algorithm (xxh64/xxh3 = 8, + * md5/md4/xxh128 = 16, sha1 = 20, none = 0). */ uint8_t checksum_digest_len(ChecksumAlgo algo); -#endif /* CHECKSUM_H */ \ No newline at end of file +/* Pick the first algorithm from FastSync's compiled-in preference list that is + * supported on this build (rsync 3.4.1's `--version` order: + * xxh128 xxh3 xxh64 md5 md4 sha1 none). Used to resolve "auto". */ +ChecksumAlgo checksum_negotiate_default(void); + +#endif /* CHECKSUM_H */ diff --git a/src/shared/compression.c b/src/shared/compression.c index 70a8bad..1c89fc0 100644 --- a/src/shared/compression.c +++ b/src/shared/compression.c @@ -3,12 +3,15 @@ #include "log.h" #include "protocol.h" #include +#include +#include #include #include #include #include #include #include +#include #include #define INITIAL_DECOMPRESS_BUF_SIZE (1024 * 1024) @@ -29,6 +32,13 @@ "z " \ "zip zst" +/* Self-describing compressed frames: the first byte is the CompressionAlgo id. + * zlib/lz4 store the uncompressed size as a little-endian uint32 after the + * codec byte so decompression can be exactly pre-sized and bounded. */ +#define LZ4_SIZE_PREFIX_LEN 4 + +static _Atomic int g_compression_algo = COMPRESSION_ALGO_ZSTD; + /* Case-insensitive match of a bare suffix (no leading dot) against a * space-separated suffix list. */ static bool suffix_in_list(const char* name, const char* list) { @@ -67,6 +77,75 @@ bool compression_should_skip_with_suffixes(const char* path, char* const* suffix return false; } +CompressionAlgo compression_default_algo(void) { + return COMPRESSION_ALGO_ZSTD; +} + +int compression_algo_from_name(const char* name) { + if (!name) + return -1; + if (strcasecmp(name, "zstd") == 0) + return (int)COMPRESSION_ALGO_ZSTD; + if (strcasecmp(name, "lz4") == 0) + return (int)COMPRESSION_ALGO_LZ4; + if (strcasecmp(name, "zlib") == 0) + return (int)COMPRESSION_ALGO_ZLIB; + if (strcasecmp(name, "zlibx") == 0) + return (int)COMPRESSION_ALGO_ZLIBX; + if (strcasecmp(name, "none") == 0) + return (int)COMPRESSION_ALGO_NONE; + return -1; +} + +const char* compression_algo_name(CompressionAlgo algo) { + switch (algo) { + case COMPRESSION_ALGO_NONE: + return "none"; + case COMPRESSION_ALGO_ZSTD: + return "zstd"; + case COMPRESSION_ALGO_LZ4: + return "lz4"; + case COMPRESSION_ALGO_ZLIB: + return "zlib"; + case COMPRESSION_ALGO_ZLIBX: + return "zlibx"; + } + return ""; +} + +bool compression_algo_valid(int algo) { + return algo == (int)COMPRESSION_ALGO_NONE || algo == (int)COMPRESSION_ALGO_ZSTD || + algo == (int)COMPRESSION_ALGO_LZ4 || algo == (int)COMPRESSION_ALGO_ZLIB || + algo == (int)COMPRESSION_ALGO_ZLIBX; +} + +bool compression_algo_enabled(CompressionAlgo algo) { + return algo != COMPRESSION_ALGO_NONE; +} + +CompressionAlgo compression_negotiate_default(void) { + /* rsync 3.4.1 default preference order; every entry is compiled in, so this + * resolves to zstd. */ + static const CompressionAlgo preference[] = { + COMPRESSION_ALGO_ZSTD, COMPRESSION_ALGO_LZ4, COMPRESSION_ALGO_ZLIBX, + COMPRESSION_ALGO_ZLIB, COMPRESSION_ALGO_NONE, + }; + for (size_t i = 0; i < sizeof(preference) / sizeof(preference[0]); i++) { + if (compression_algo_valid((int)preference[i])) + return preference[i]; + } + return COMPRESSION_ALGO_ZSTD; +} + +void compression_set_algo(CompressionAlgo algo) { + if (compression_algo_valid((int)algo)) + atomic_store(&g_compression_algo, (int)algo); +} + +CompressionAlgo compression_get_algo(void) { + return (CompressionAlgo)atomic_load(&g_compression_algo); +} + /* Per-thread cache of zstd contexts plus the grow-only compression scratch * buffer. zstd contexts are stateful and not safe to share between threads, * so each thread keeps its own (see compression_get_thread_ctx). The cache is @@ -157,17 +236,25 @@ static void compression_ctx_put(CompressionThreadCtx* ctx) { compression_ctx_free(ctx); } -Data* data_compress(Data* data_to_compress, int compression_level) { - return data_compress_with_threads(data_to_compress, compression_level, 0); +/* Build a frame consisting of a copy of `src` prefixed by `codec`. */ +static Data* frame_with_codec(const void* src, size_t size, CompressionAlgo codec) { + if (size > SIZE_MAX - 1) + return NULL; + Data* out = data_create_empty(size + 1); + if (!out) + return NULL; + ((uint8_t*)out->data)[0] = (uint8_t)codec; + if (size > 0) + memcpy((uint8_t*)out->data + 1, src, size); + out->size = size + 1; + return out; } -Data* data_compress_with_threads(Data* data_to_compress, int compression_level, - int compression_threads) { - if (!data_to_compress || (!data_to_compress->data && data_to_compress->size != 0) || - compression_threads < 0 || compression_threads > COMPRESSION_MAX_THREADS) +static Data* zstd_compress(Data* in, int compression_level, int compression_threads) { + size_t dst_size = ZSTD_compressBound(in->size); + if (dst_size > SIZE_MAX - 1) return NULL; - log_message(LOG_LEVEL_DEBUG, "Starting to compress data"); - size_t dst_size = ZSTD_compressBound(data_to_compress->size); + dst_size += 1; /* codec prefix */ CompressionThreadCtx* ctx = compression_get_thread_ctx(); if (ctx == NULL) { @@ -218,7 +305,7 @@ Data* data_compress_with_threads(Data* data_to_compress, int compression_level, if (available_threads > 0) { /* Streaming compression needs the source size before threaded mode can end a frame. */ - size_t zret = ZSTD_CCtx_setPledgedSrcSize(ctx->cctx, data_to_compress->size); + size_t zret = ZSTD_CCtx_setPledgedSrcSize(ctx->cctx, in->size); if (ZSTD_isError(zret)) { log_message(LOG_LEVEL_ERROR, "Failed to set compression source size: %s", ZSTD_getErrorName(zret)); @@ -236,8 +323,8 @@ Data* data_compress_with_threads(Data* data_to_compress, int compression_level, ctx->out_cap = dst_size; } - ZSTD_inBuffer input = {data_to_compress->data, data_to_compress->size, 0}; - ZSTD_outBuffer output = {ctx->out_buf, dst_size, 0}; + ZSTD_inBuffer input = {in->data, in->size, 0}; + ZSTD_outBuffer output = {(uint8_t*)ctx->out_buf + 1, dst_size - 1, 0}; size_t ret; do { @@ -250,30 +337,192 @@ Data* data_compress_with_threads(Data* data_to_compress, int compression_level, /* Hand off an exactly-sized copy; the scratch buffer stays cached so the next * call does not reallocate a ZSTD_compressBound-sized block. */ - compressed_data = data_create_empty(output.pos); + compressed_data = data_create_empty(output.pos + 1); if (compressed_data == NULL) { log_message(LOG_LEVEL_ERROR, "Failed to allocate compressed data"); goto cleanup; } + ((uint8_t*)compressed_data->data)[0] = (uint8_t)COMPRESSION_ALGO_ZSTD; if (output.pos > 0) - memcpy(compressed_data->data, ctx->out_buf, output.pos); - compressed_data->size = output.pos; + memcpy((uint8_t*)compressed_data->data + 1, (uint8_t*)ctx->out_buf + 1, output.pos); + compressed_data->size = output.pos + 1; - log_debug_message(LOG_DEBUG_UTIL, "Data succesfully compressed from %zu to %zu", - data_to_compress->size, compressed_data->size); + log_debug_message(LOG_DEBUG_UTIL, "Data succesfully compressed from %zu to %zu", in->size, + compressed_data->size); cleanup: compression_ctx_put(ctx); return compressed_data; } -Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) { - if (!compressed_data || (!compressed_data->data && compressed_data->size != 0) || - maximum_size == 0) +static Data* lz4_compress(Data* in) { + int bound = LZ4_compressBound((int)in->size); + if (bound < 0 || in->size > (size_t)INT_MAX) return NULL; + Data* out = data_create_empty((size_t)bound + 1 + LZ4_SIZE_PREFIX_LEN); + if (!out) + return NULL; + uint32_t raw_size = (uint32_t)in->size; + uint8_t* p = (uint8_t*)out->data; + p[0] = (uint8_t)COMPRESSION_ALGO_LZ4; + for (int i = 0; i < LZ4_SIZE_PREFIX_LEN; i++) + p[1 + i] = (uint8_t)((raw_size >> (8 * i)) & 0xff); + int written = 0; + if (in->size > 0) { + written = LZ4_compress_default((const char*)in->data, (char*)p + 1 + LZ4_SIZE_PREFIX_LEN, + (int)in->size, bound); + if (written <= 0) { + data_destroy(out); + return NULL; + } + } + out->size = (size_t)written + 1 + LZ4_SIZE_PREFIX_LEN; + return out; +} + +static Data* zlib_compress(Data* in, CompressionAlgo algo, int compression_level) { + int level = compression_level; + if (level < 1) + level = Z_DEFAULT_COMPRESSION; + if (level > 9) + level = 9; + uLong bound = compressBound((uLong)in->size); + if (in->size > (size_t)ULONG_MAX) + return NULL; + Data* out = data_create_empty((size_t)bound + 1 + LZ4_SIZE_PREFIX_LEN); + if (!out) + return NULL; + uint32_t raw_size = (uint32_t)in->size; + uint8_t* p = (uint8_t*)out->data; + p[0] = (uint8_t)algo; + for (int i = 0; i < LZ4_SIZE_PREFIX_LEN; i++) + p[1 + i] = (uint8_t)((raw_size >> (8 * i)) & 0xff); + uLongf dest_len = bound; + int rc = compress2(p + 1 + LZ4_SIZE_PREFIX_LEN, &dest_len, (const Bytef*)in->data, + (uLong)in->size, level); + if (rc != Z_OK) { + data_destroy(out); + return NULL; + } + out->size = (size_t)dest_len + 1 + LZ4_SIZE_PREFIX_LEN; + return out; +} + +Data* data_compress_codec(Data* data_to_compress, CompressionAlgo algo, int compression_level, + int compression_threads) { + if (!data_to_compress || (!data_to_compress->data && data_to_compress->size != 0) || + compression_threads < 0 || compression_threads > COMPRESSION_MAX_THREADS) + return NULL; + if (!compression_algo_valid((int)algo)) + return NULL; + log_message(LOG_LEVEL_DEBUG, "Starting to compress data"); + switch (algo) { + case COMPRESSION_ALGO_NONE: + return frame_with_codec(data_to_compress->data, data_to_compress->size, COMPRESSION_ALGO_NONE); + case COMPRESSION_ALGO_ZSTD: + return zstd_compress(data_to_compress, compression_level, compression_threads); + case COMPRESSION_ALGO_LZ4: + return lz4_compress(data_to_compress); + case COMPRESSION_ALGO_ZLIB: + case COMPRESSION_ALGO_ZLIBX: + return zlib_compress(data_to_compress, algo, compression_level); + } + return NULL; +} + +Data* data_compress_with_threads(Data* data_to_compress, int compression_level, + int compression_threads) { + return data_compress_codec(data_to_compress, compression_get_algo(), compression_level, + compression_threads); +} + +Data* data_compress(Data* data_to_compress, int compression_level) { + return data_compress_codec(data_to_compress, compression_get_algo(), compression_level, 0); +} + +static Data* decompress_none(const Data* compressed_data, size_t maximum_size) { + size_t size = compressed_data->size - 1; + if (size > maximum_size) + return NULL; + Data* out = data_create_empty(size); + if (!out) + return NULL; + if (size > 0) + memcpy(out->data, (const uint8_t*)compressed_data->data + 1, size); + out->size = size; + return out; +} + +/* Read the 4-byte little-endian raw size stored after the codec byte. */ +static bool read_raw_size(const Data* in, uint32_t* raw_size) { + if (in->size < 1 + LZ4_SIZE_PREFIX_LEN) + return false; + const uint8_t* p = (const uint8_t*)in->data; + uint32_t v = 0; + for (int i = 0; i < LZ4_SIZE_PREFIX_LEN; i++) + v |= (uint32_t)p[1 + i] << (8 * i); + *raw_size = v; + return true; +} + +static Data* lz4_decompress(Data* compressed_data, size_t maximum_size, size_t hard_limit) { + uint32_t raw_size = 0; + if (!read_raw_size(compressed_data, &raw_size)) + return NULL; + if (raw_size > hard_limit || raw_size > maximum_size) + return NULL; + size_t comp_size = compressed_data->size - 1 - LZ4_SIZE_PREFIX_LEN; + Data* out = data_create_empty(raw_size); + if (!out) + return NULL; + if (raw_size == 0) { + out->size = 0; + return out; + } + int rc = LZ4_decompress_safe((const char*)compressed_data->data + 1 + LZ4_SIZE_PREFIX_LEN, + (char*)out->data, (int)comp_size, (int)raw_size); + if (rc < 0 || (uint32_t)rc != raw_size) { + log_message(LOG_LEVEL_ERROR, "LZ4 decompression failed"); + data_destroy(out); + return NULL; + } + out->size = raw_size; + return out; +} + +static Data* zlib_decompress(Data* compressed_data, size_t maximum_size, size_t hard_limit) { + uint32_t raw_size = 0; + if (!read_raw_size(compressed_data, &raw_size)) + return NULL; + if (raw_size > hard_limit || raw_size > maximum_size) + return NULL; + size_t comp_size = compressed_data->size - 1 - LZ4_SIZE_PREFIX_LEN; + Data* out = data_create_empty(raw_size); + if (!out) + return NULL; + if (raw_size == 0) { + out->size = 0; + return out; + } + uLongf dest_len = raw_size; + int rc = + uncompress((Bytef*)out->data, &dest_len, + (const Bytef*)compressed_data->data + 1 + LZ4_SIZE_PREFIX_LEN, (uLong)comp_size); + if (rc != Z_OK || dest_len != raw_size) { + log_message(LOG_LEVEL_ERROR, "zlib decompression failed"); + data_destroy(out); + return NULL; + } + out->size = raw_size; + return out; +} + +static Data* zstd_decompress(Data* compressed_data, size_t maximum_size) { + /* The zstd frame starts after the codec byte. */ + const void* frame = (const uint8_t*)compressed_data->data + 1; + size_t frame_size = compressed_data->size - 1; log_debug_message(LOG_DEBUG_UTIL, "Start to decompress data"); - unsigned long long dst_size = - ZSTD_getFrameContentSize(compressed_data->data, compressed_data->size); + unsigned long long dst_size = ZSTD_getFrameContentSize(frame, frame_size); /* ZSTD_isError() is also true for ZSTD_CONTENTSIZE_ERROR and * ZSTD_CONTENTSIZE_UNKNOWN (both are encoded near (size_t)-1), so test the * sentinels explicitly instead of blanket-rejecting every error-ish value: @@ -287,9 +536,9 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) { // ZSTD_CONTENTSIZE_UNKNOWN (~2^64) can cause massive allocation; // fall back to a conservative estimate (3x compressed size) when unknown. if (dst_size == ZSTD_CONTENTSIZE_UNKNOWN) { - if (compressed_data->size > ULLONG_MAX / 3) + if (frame_size > ULLONG_MAX / 3) return NULL; - dst_size = compressed_data->size * 3; + dst_size = frame_size * 3; if (dst_size < INITIAL_DECOMPRESS_BUF_SIZE) dst_size = INITIAL_DECOMPRESS_BUF_SIZE; } @@ -326,7 +575,7 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) { goto cleanup; } - ZSTD_inBuffer input = {compressed_data->data, compressed_data->size, 0}; + ZSTD_inBuffer input = {frame, frame_size, 0}; ZSTD_outBuffer output = {uncompressed_data->data, buf_size, 0}; size_t ret; @@ -385,6 +634,31 @@ cleanup: return uncompressed_data; } +Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) { + if (!compressed_data || (!compressed_data->data && compressed_data->size != 0) || + maximum_size == 0) + return NULL; + if (compressed_data->size < 1) + return NULL; + unsigned long long hard_limit = + maximum_size < MAX_DECOMPRESSED_SIZE ? maximum_size : MAX_DECOMPRESSED_SIZE; + uint8_t codec = ((const uint8_t*)compressed_data->data)[0]; + if (!compression_algo_valid(codec)) + return NULL; + switch ((CompressionAlgo)codec) { + case COMPRESSION_ALGO_NONE: + return decompress_none(compressed_data, (size_t)hard_limit); + case COMPRESSION_ALGO_ZSTD: + return zstd_decompress(compressed_data, (size_t)hard_limit); + case COMPRESSION_ALGO_LZ4: + return lz4_decompress(compressed_data, maximum_size, (size_t)hard_limit); + case COMPRESSION_ALGO_ZLIB: + case COMPRESSION_ALGO_ZLIBX: + return zlib_decompress(compressed_data, maximum_size, (size_t)hard_limit); + } + return NULL; +} + Data* data_decompress(Data* compressed_data) { return data_decompress_limited(compressed_data, MAX_DECOMPRESSED_SIZE); } diff --git a/src/shared/compression.h b/src/shared/compression.h index 2c3753c..ff5bbb5 100644 --- a/src/shared/compression.h +++ b/src/shared/compression.h @@ -6,11 +6,61 @@ #define COMPRESSION_MAX_THREADS 64 +/* Compression algorithms selectable with --compress-choice / -z. The ids are + * the values placed on the wire (Config->compression_algo), so they must be + * kept stable. NONE is "no compression"; ZSTD is the historical FastSync + * default and the negotiated "auto" choice. ZLIBX is rsync's zlib-without- + * matched-data variant: FastSync compresses only the delta/token bytes (it does + * not put matched file data in the compression stream), so its zlib codec is + * already the "x" form and zlib/zlibx share the same implementation, recorded + * under distinct ids. */ +typedef enum { + COMPRESSION_ALGO_NONE = 0, + COMPRESSION_ALGO_ZSTD = 1, + COMPRESSION_ALGO_LZ4 = 2, + COMPRESSION_ALGO_ZLIB = 3, + COMPRESSION_ALGO_ZLIBX = 4 +} CompressionAlgo; + +CompressionAlgo compression_default_algo(void); + +/* Resolve a --compress-choice string (case-insensitive) to an algorithm id. + * Accepts "zstd", "lz4", "zlib", "zlibx", "none". "auto" is not an algorithm + * here; the caller resolves it to the negotiated default. Returns -1 for any + * unrecognized name. */ +int compression_algo_from_name(const char* name); +const char* compression_algo_name(CompressionAlgo algo); +bool compression_algo_valid(int algo); + +/* Pick the first algorithm from FastSync's compiled-in preference list + * (rsync 3.4.1's `--version` order: zstd lz4 zlibx zlib none). Resolves + * "auto". */ +CompressionAlgo compression_negotiate_default(void); + +/* True when the algorithm actually compresses (i.e. is not NONE). */ +bool compression_algo_enabled(CompressionAlgo algo); + +/* Select the process-wide codec used by the legacy wrappers below. Each + * process serves exactly one transfer config (the server forks per connection, + * the client configures itself before spawning transfer threads), so a + * process-global default is sufficient and constant for the lifetime of a + * transfer. Defaults to ZSTD when never set. Thread-safe. */ +void compression_set_algo(CompressionAlgo algo); +CompressionAlgo compression_get_algo(void); + +/* Codec-aware primitives. The compressed buffer is self-describing: its first + * byte is the CompressionAlgo id, so decompression never needs the codec passed + * separately (this keeps every existing Decompress call site source-compatible). + * `data_compress_codec` returns NULL on invalid input or an unsupported codec. */ +Data* data_compress_codec(Data* data_to_compress, CompressionAlgo algo, int compression_level, + int compression_threads); +Data* data_decompress_limited(Data* compressed_data, size_t maximum_size); + +/* Legacy zstd-default wrappers retained for existing callers/tests. */ Data* data_compress(Data* data_to_compress, int compression_level); Data* data_compress_with_threads(Data* data_to_compress, int compression_level, int compression_threads); Data* data_decompress(Data* compressed_data); -Data* data_decompress_limited(Data* compressed_data, size_t maximum_size); bool compression_should_skip_with_suffixes(const char* path, char* const* suffixes, int count); /* Release the calling thread's cached zstd contexts (compressor, decompressor