Merge Wave 4: performance (packed metadata 2.20.0, indexed lookups, zstd reuse, byte-bounded queues)
CI / lint (push) Successful in 1m30s
CI / sanitizers (undefined) (push) Successful in 1m1s
CI / sanitizers (address) (push) Successful in 1m7s
CI / fuzz-build (push) Successful in 29s
CI / coverage (push) Successful in 50s
CI / valgrind (push) Successful in 3m12s
CI / build-and-test (push) Successful in 5m22s
CI / lint (push) Successful in 1m30s
CI / sanitizers (undefined) (push) Successful in 1m1s
CI / sanitizers (address) (push) Successful in 1m7s
CI / fuzz-build (push) Successful in 29s
CI / coverage (push) Successful in 50s
CI / valgrind (push) Successful in 3m12s
CI / build-and-test (push) Successful in 5m22s
This commit is contained in:
@@ -4,6 +4,35 @@ All notable changes to FastSync are documented here. Versions match
|
||||
`PROTOCOL_VERSION` (printed by `fastsync --version`); the client and server must
|
||||
run the same version because the handshake is strict.
|
||||
|
||||
## [2.20.0] - 2026-09-13
|
||||
|
||||
### Security
|
||||
|
||||
- Cap cumulative `DirTimeList` growth and bound pre-auth config-string memory
|
||||
(remote memory-exhaustion DoS).
|
||||
- Daemon host access control (`hosts allow`/`hosts deny`, IPv4/IPv6/CIDR),
|
||||
configurable global `max connections`, connection audit logging, and a
|
||||
bounded `auth failure delay` throttle. IPv4-mapped peers are normalized and
|
||||
invalid patterns are rejected at parse time (no silent fail-open).
|
||||
- Honor `--timeout` for protocol I/O and bound idle/session time to defeat
|
||||
keepalive slowloris; child-safe signal handling in the forked daemon.
|
||||
- Compiler/linker hardening (`_FORTIFY_SOURCE`, stack protector, PIE, RELRO)
|
||||
and pinned build dependencies.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Use-after-free in the basis-dir oversize preflight.
|
||||
- Placeholder `Data` leaks, `missing_args` leak, scanner chunk leak.
|
||||
- Thread-safe logging; single fd owner and cleanup epilogue in the server
|
||||
handler.
|
||||
|
||||
### Performance
|
||||
|
||||
- Metadata now crosses the wire as one packed frame (protocol 2.20.0).
|
||||
- Delete keep-set and `--files-from` lookups indexed (O(n*m) → O(n)).
|
||||
- Reused per-thread zstd contexts; `TCP_NODELAY` by default.
|
||||
- Byte-bounded sender queues; removed a redundant scanner `stat()`.
|
||||
|
||||
## [2.19.0] - 2026-09-12
|
||||
|
||||
### Security
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
cmake_minimum_required(VERSION 3.22)
|
||||
|
||||
project(FastFileTransfer VERSION 2.19.0)
|
||||
project(FastFileTransfer VERSION 2.20.0)
|
||||
|
||||
set(CMAKE_EXPORT_COMPILE_COMMANDS ON)
|
||||
set(CMAKE_C_STANDARD 11)
|
||||
|
||||
@@ -551,7 +551,7 @@ before the module list, before authentication, and the connecting peer address
|
||||
|
||||
## Protocol and Security
|
||||
|
||||
FastSync protocol version `2.19.0` is shared by the client and server. The
|
||||
FastSync protocol version `2.20.0` is shared by the client and server. The
|
||||
current protocol is sender-driven and includes configuration negotiation,
|
||||
including the maximum allocation limit, incremental checks, checksums,
|
||||
manifests, keep-alives, abort handling, per-file remove-source results, and
|
||||
|
||||
+19
-2
@@ -679,7 +679,7 @@ now transmits targets (the prior behavior was broken/partial); its status moved
|
||||
| `--stop-after=MINS` | Stop after N minutes | ✅ Implemented | Client-only sender stop deadline (Phase 6): computing `--stop-after=MINS` (a positive minute count; 0/negative/garbage rejected) and `--stop-at=TIME` (`HH:MM`, `HH:MM:SS`, or `now+N[smhd]`; a past time stops immediately). The transfer stops ELEGANTLY at the next chunk boundary: everything already fully sent is kept and applied, the run returns 0, and --delete (late/delete-after timing) does NOT wipe the destination — when the scan is cut short the partial keep-set manifest is suppressed with a warning (the delete walk is skipped rather than acting on an incomplete keep-set, so unscanned source mirrors survive). `--delete-before`/`--delete-during` still run their complete pre-scan (which ignores the deadline). Local client-only fields: never serialized into the wire config frame, so no PROTOCOL_VERSION bump. `--stop-after` uses CLOCK_MONOTONIC; `--stop-at` uses the wall clock. Works single-threaded and under `-j`/`--threads` (multithreaded). Divergence: rsync computes `--stop-after` from the run start; FastSync likewise. When both are given, the earlier of the two deadlines wins (checked per iteration). See the Phase-6 stop notes below |
|
||||
| `--stop-at=TIME` | Stop at specified time | ✅ Implemented | Same feature as `--stop-after` (deadline transfer stop), absolute wall-clock form (`HH:MM[:SS]` or `now+N[smhd]`). See the row above and the Phase-6 stop notes |
|
||||
| `--fsync` | Fsync every written file before publication | ✅ Implemented | |
|
||||
| `--protocol=NUM` | Force older protocol version | ✅ Implemented | Forces the wire protocol version for this transfer. FastSync has exactly ONE wire format (`PROTOCOL_VERSION`, currently 2.19.0) with no downgrade/backward-compat code paths, so `--protocol=2.19.0` is accepted (it sets the version claim the client sends, which the server already requires to match exactly) and **every other value is rejected up front** with a clear error before any connection — it does not and cannot speak an older or virtual wire format. Divergence from rsync (which negotiates a range and downgrades to an integer 0..31): FastSync's honest contract is force-to-the-one-supported-value; a genuine downgrade would require a per-version compatibility layer that does not exist. Client-only; the server-side exact-match check is unchanged. `--protocol=2.18.0`/`2.18`/`2.17.0`/`2.16.0`/`2.15.0`/`216`/`31`/garbage are all rejected. See the Phase-6 protocol note below |
|
||||
| `--protocol=NUM` | Force older protocol version | ✅ Implemented | Forces the wire protocol version for this transfer. FastSync has exactly ONE wire format (`PROTOCOL_VERSION`, currently 2.20.0) with no downgrade/backward-compat code paths, so `--protocol=2.20.0` is accepted (it sets the version claim the client sends, which the server already requires to match exactly) and **every other value is rejected up front** with a clear error before any connection — it does not and cannot speak an older or virtual wire format. Divergence from rsync (which negotiates a range and downgrades to an integer 0..31): FastSync's honest contract is force-to-the-one-supported-value; a genuine downgrade would require a per-version compatibility layer that does not exist. Client-only; the server-side exact-match check is unchanged. `--protocol=2.19.0`/`2.18.0`/`2.18`/`2.17.0`/`2.16.0`/`2.15.0`/`216`/`31`/garbage are all rejected. See the Phase-6 protocol note below |
|
||||
| `--iconv=CONVERT_SPEC` | Charset conversion | ✅ Implemented | Charset conversion of FILE NAMES (not content) at the protocol boundary via iconv(3): `--iconv=LOCAL[,REMOTE]` — the sender converts each local filename LOCAL→REMOTE before transmitting, and the receiver converts each wire filename REMOTE→LOCAL before creating/writing. The full CONVERT_SPEC is serialized into the config frame as a new trailing string field so the peer knows the wire charset; **PROTOCOL_VERSION bumped 2.15.0 → 2.16.0**. `LOCAL[,REMOTE]` parse: single charset ⇒ LOCAL==REMOTE (identity both ways); garbage rejected up front. Validation probes BOTH directions (a spec that only opens one way is refused, as is a NUL-emitting target charset like utf-16/utf-32/ucs-2, since filenames cannot contain NUL). An unrepresentable name (EILSEQ/EINVAL) fails that path cleanly with a logged `--iconv: cannot convert file name ...` and is never written mangled/truncated. Conversion is applied at EVERY wire-path site (regular/MKDIR/hardlink path+target/symlink path+target/SPECIAL, the delete manifest, the incremental-check path, and the `-s`/`chunk_serialize` embedded blob path), on both client and server (`--iconv` is also a server/daemon option). Zero overhead when unset. See the Phase-6 iconv notes below |
|
||||
| `--checksum-seed=NUM` | Set checksum seed | ✅ Implemented | Sets the seed for FastSync's whole-file xxHash64 digest (full 64-bit seed) and for the delta path's per-block xxHash32 strong checksum (low 32 bits of the seed). An explicit seed deterministically changes every computed digest on BOTH endpoints (sender and receiver share the seed via the config frame, protocol 2.10.0), so identical runs with the same seed skip the same files and a changed seed changes the digests — the explicit-seed path that makes xxHash comparisons deterministic. `--checksum-choice=md5` has no seed and ignores it (documented). The value is a strict decimal 0..2⁶⁴-1 (blank, signed, or non-numeric values are rejected). Like rsync, a seed only matters where a digest is actually computed (`--checksum` or a basis-dir run, or a delta transfer); it does not by itself enable `--checksum`/`--delta`. Divergence from rsync: the default is seed 0, and FastSync never randomizes the seed (rsync uses a random per-transfer seed when `--checksum-seed` is unset); FastSync's unset default therefore reproduces its historical byte-for-byte behavior |
|
||||
| `--secluded-args`, `-s` | Use protocol to send args | ⛔ Impossible/Divergence | Accepted for CLI compatibility (including the rsync short `-s`, Phase 7 Wave A) but a documented **no-op / divergence**. rsync's `-s` protects arguments from shell expansion by shipping them over the protocol; FastSync never passes remote arguments through a shell expansion boundary in the first place — its SSH transport builds the remote argv as **single-quote-escaped shell words** (`ssh_build_remote_command`), so the injection/leak that `-s` guards against does not exist and there is nothing to "seclude". Implementing a true arg-send protocol would mean replacing the argv-based SSH launch with an in-band argument channel, a large redesign of the transport that buys no security here. Chunk serialization remains the long-only `--chunk-serialization`. |
|
||||
@@ -795,7 +795,7 @@ These are the hardest compatibility items because they require durable formats o
|
||||
|
||||
**Phase 6, Wave B (iconv) shipping note (PROTOCOL 2.15.0 → 2.16.0):** `--iconv=LOCAL[,REMOTE]` converts file NAMES at the wire boundary (never content). The full CONVERT_SPEC is serialized into the config frame as a new trailing string field (empty→NULL canonicalized), so both ends share the same wire charset interpretation; this required the PROTOCOL bump because the frame is a strict ordered sequence and a peer that does not parse the new trailing field would desynchronize. Each end derives LOCAL (its own charset) and REMOTE (the wire charset): the sender opens LOCAL→REMOTE and converts every transmitted filename; the receiver opens REMOTE→LOCAL and converts every received filename before creating/writing. Conversion is applied at every wire-path site (regular/MKDIR/hardlink path+target/symlink path+target/SPECIAL, the delete manifest keep/protected/missing entries, the incremental-check path, and the embedded `-s`/chunk-blob path). A name it cannot convert (EILSEQ/EINVAL) is failed cleanly with a logged `--iconv: cannot convert file name ...` and is never written truncated/mangled. Validation probes both directions up front (both the sender local→remote and the receiver remote→local, and, for a server/daemon with its own `--iconv`, the client-REMOTE→server-LOCAL pair) so an unusable spec is rejected before the connection rather than mid-transfer, and NUL-emitting target charsets (utf-16/utf-32/ucs-2) are refused because filenames cannot contain NUL. Divergence documented upstream: the receiver does NOT half-swap; the wire charset always comes from the sender's REMOTE half, so a server whose local charset differs from the client's LOCAL must declare it with its own `--iconv`. Conversion is process-global and runs on a single thread per process (sender thread / receiver-loop thread), initialized before worker threads start and freed after they join.
|
||||
|
||||
**Phase 6, Wave C (protocol-version) shipping note (no PROTOCOL_VERSION change):** `--protocol=NUM` lets the client force the wire protocol version for a transfer. FastSync's protocol is a single lockstep format: the config frame is a strict ordered sequence and the server requires the client's version string to equal `PROTOCOL_VERSION` exactly (`config_receive_with_validate`, src/shared/config.c) — there are no older-format code paths and no downgrade/negotiation machinery, so a lower/higher/virtual version can never be spoken. The honest contract is therefore: `--protocol=2.19.0` (the current `PROTOCOL_VERSION`, as of the A7 auth redesign) is accepted and stored into the client's `version` claim (which `config_send` already transmits), and every other value — `2.18.0`, `2.18`, `2.17.0`, `2.16.0`, `2.15.0`, `3.0.0`, rsync-integer spellings like `216`/`31`, garbage, empty — is rejected up front in `validate_config()` before any connection, with a clear error that FastSync supports only its current wire protocol and cannot speak an older or virtual one. Implementation is client-only: a server-side `--protocol` is intentionally not added because the server has no negotiation (it only enforces exact match), and it could only ever be the current version. This preserves (and slightly tightens) existing validation: the client now also refuses to launch with a version it cannot actually speak, rather than only the server rejecting it later. A genuine downgrade would require a per-version compatibility layer for every frame/feature added since (append 2.10, preallocate 2.11, hardlinks 2.12, devices/specials/symlink-trust/xattr 2.13, remote-option 2.14, daemon module/auth 2.15, iconv 2.16, dir/symlink times 2.17, privilege flags --super/--copy-as 2.18, SCRAM daemon auth 2.19) and is intentionally out of scope — documented divergences from rsync's integer-negotiated downgrade remain.
|
||||
**Phase 6, Wave C (protocol-version) shipping note (no PROTOCOL_VERSION change):** `--protocol=NUM` lets the client force the wire protocol version for a transfer. FastSync's protocol is a single lockstep format: the config frame is a strict ordered sequence and the server requires the client's version string to equal `PROTOCOL_VERSION` exactly (`config_receive_with_validate`, src/shared/config.c) — there are no older-format code paths and no downgrade/negotiation machinery, so a lower/higher/virtual version can never be spoken. The honest contract is therefore: `--protocol=2.20.0` (the current `PROTOCOL_VERSION`, as of the packed-metadata wave) is accepted and stored into the client's `version` claim (which `config_send` already transmits), and every other value — `2.19.0`, `2.18.0`, `2.18`, `2.17.0`, `2.16.0`, `2.15.0`, `3.0.0`, rsync-integer spellings like `216`/`31`, garbage, empty — is rejected up front in `validate_config()` before any connection, with a clear error that FastSync supports only its current wire protocol and cannot speak an older or virtual one. Implementation is client-only: a server-side `--protocol` is intentionally not added because the server has no negotiation (it only enforces exact match), and it could only ever be the current version. This preserves (and slightly tightens) existing validation: the client now also refuses to launch with a version it cannot actually speak, rather than only the server rejecting it later. A genuine downgrade would require a per-version compatibility layer for every frame/feature added since (append 2.10, preallocate 2.11, hardlinks 2.12, devices/specials/symlink-trust/xattr 2.13, remote-option 2.14, daemon module/auth 2.15, iconv 2.16, dir/symlink times 2.17, privilege flags --super/--copy-as 2.18, SCRAM daemon auth 2.19, packed metadata 2.20) and is intentionally out of scope — documented divergences from rsync's integer-negotiated downgrade remain.
|
||||
|
||||
**Phase-1/2 selection-and-update status correction (docs):** `-I/--ignore-times`, `--size-only`, `-@/--modify-window`, `--existing`, `--ignore-existing`, `-u/--update`, `-W/--whole-file`, and `--compress-threads` were previously listed as not-implemented in this document but are in fact fully implemented and tested on `dev`. This pass corrects the matrix to match the code. The realistic model of these is that FastSync is a *sender-driven* whole-tree copy, so the size+mtime quick-check and all three receiver-policy skips (`--existing`, `--ignore-existing`, `-u`) are evaluated against the **destination** on the receiver side, and their booleans cross the wire in the config frame. `-I`/`--size-only`/`--modify-window` modify the `--incremental` per-file `STATUS_CHECK` handshake's match predicate (`-I` disables the mtime leg and forces transfer; `--size-only` drops only the mtime leg; `--modify-window` adds tolerance to `metadata_mtime_matches`); they require `--incremental` (or a basis dir) to have a handshake to affect, mirroring how they only matter where a quick-check exists in rsync. `--existing`/`--ignore-existing`/`-u` are receiver write-time policies (skipping the write / newer-destination guard) applied across the regular-file, `--delay-updates`-staged, hardlink-sibling, and special/device paths; `-u` implies `-M` metadata and uses a second-then-nanosecond strict `>` newer check; both correctly influence `--remove-source-files` (a skipped source is not removed). `-W/--whole-file` disables block-level delta (opt-in via `--delta`), folded into the wire `use_delta` so no protocol bump was needed, and makes `--fuzzy` inert; `--append`/`--append-verify` are rejected with `-W`. `--compress-threads=NUM` (1..64, client-only, never crosses the wire) sizes the zstd compression worker pool. No code was changed by this correction; the implementation had landed in earlier merge waves (feat/ignore-times, feat/ignore-existing via the newer `file_to_disk_secure_no_replace`/`linkat EEXIST` path, feat/size-only, feat/modify-window, feat/whole-file, feat/update, compression-threads).
|
||||
|
||||
@@ -840,6 +840,23 @@ These are the last compatibility items and the closing phase toward rsync flag p
|
||||
|
||||
**Post-Phase-7 Summary (after Waves A–E).** ✅143 / 🔀0 / ⛔4 / ⚠️0 / 🔄0 / ❌0 = 147. The 3 `🔀 Alt Arg` rows (`-a`, `-p`, `-z`) are ✅ (Wave A). All 10 prior `⚠️ Partial` rows are resolved to ✅ (`-S`, `-P`, `--block-size`, `--fake-super`, `--devices`, `--copy-devices`, `--write-devices`) or ⛔ (`--stderr=client`, `-N/--crtimes`, `--specials` for the impossible socket case). The 3 `🔄 Compatibility No-op` rows are resolved: `-O`/`-J` are now real ✅ (Wave D), `--secluded-args` is ⛔. The **Impossible/Divergence** bucket holds the 4 physically-impossible/divergent flags: `--stderr=client`, `-N/--crtimes`, `--specials` (sockets), `--secluded-args`. The last two `❌ Not Implemented` rows — `--super` and `--copy-as=USER[:GROUP]` — are now ✅ (Wave E). **No `❌ Not Implemented` rows remain.**
|
||||
|
||||
## Packed Metadata Frame (protocol 2.20.0)
|
||||
|
||||
A file's metadata used to cross the wire as up to 12 separate per-field framed
|
||||
messages (a present flag followed by mode/uid/gid/mtime/atime/crtime writes),
|
||||
which cost ~11 extra protocol frames per file on many-small-file trees. FastSync
|
||||
now sends the metadata as ONE packed frame: a single `int32` present flag
|
||||
(`0` = absent) followed, when present, by the fixed
|
||||
`FILE_METADATA_WIRE_SIZE`-byte (68-byte) field record already emitted by the
|
||||
shared `metadata_to_buf()`/`metadata_from_buf()` chunk codec. Absent metadata is
|
||||
a lone `int32` zero. The encoded field layout is unchanged (only the framing
|
||||
collapses), so chunk-serialized blobs remain byte-identical. Protocol data is an
|
||||
unframed byte stream, so the packed encoding is byte-for-byte identical to the
|
||||
old field-by-field writes; `PROTOCOL_VERSION` was bumped `2.19.0 → 2.20.0` as a
|
||||
deliberate lockstep-release marker rather than because of a
|
||||
desynchronization. The strict same-version handshake rejects any mismatch before
|
||||
a byte of the frame is parsed.
|
||||
|
||||
### Recommended Delivery Order
|
||||
|
||||
1. Resolve short-option conflicts (`-m`, `-M`, `-T`, `-f`, `-s`) and define the compatibility contract.
|
||||
|
||||
@@ -37,6 +37,13 @@
|
||||
|
||||
#define STREAM_THRESHOLD (64ULL * 1024 * 1024)
|
||||
|
||||
/* Aggregate loaded payload bytes the sender may buffer across the loader queue
|
||||
and the chunk in flight. Sending one chunk adds up to ~2 * MAX_CHUNK_SIZE of
|
||||
transient serialize/compress buffers on top of the queued payloads, so this
|
||||
ceiling keeps total pipeline memory within MAX_CONNECTION_MEMORY (mirrors the
|
||||
receiver's RECEIVER_QUEUE_MAX_BYTES). */
|
||||
#define SENDER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE)
|
||||
|
||||
/* Forward declaration for progress-reporting thread used in multithreaded send. */
|
||||
static int progress_thread_fn(void* arg);
|
||||
|
||||
@@ -1493,6 +1500,10 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
protocol_session_unbind();
|
||||
return thrd_error;
|
||||
}
|
||||
/* Payload bytes this chunk was charged for on the loader's byte budget.
|
||||
Computed before destruction and released after the memory is actually
|
||||
freed, so a loader blocked on the budget wakes only once room exists. */
|
||||
size_t queued_payload = pipeline_context_sender_chunk_bytes(current_chunk);
|
||||
unsigned long long chunk_bytes = 0;
|
||||
int chunk_files = 0;
|
||||
for (int i = 0; i < current_chunk->element_count; i++) {
|
||||
@@ -1507,6 +1518,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
context->progress_bytes = context->total_bytes;
|
||||
mtx_unlock(&context->mutex_progress);
|
||||
chunk_destroy(current_chunk);
|
||||
pipeline_context_sender_note_bytes_released(context, queued_payload);
|
||||
}
|
||||
|
||||
/* Completion tail: reached on natural exhaustion or an early stop deadline.
|
||||
@@ -1728,11 +1740,7 @@ static int load_files_multithreaded(void* pipeline_context) {
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!queue_enqueue_multithreaded_cancel(context->queue_loader, chunk, &context->mutex_loader,
|
||||
&context->condition_not_empty_loader,
|
||||
&context->condition_not_full_loader,
|
||||
&context->cancelled)) {
|
||||
chunk_destroy(chunk);
|
||||
if (!pipeline_context_sender_enqueue_chunk(context, chunk)) {
|
||||
pipeline_cancel(context);
|
||||
protocol_session_unbind();
|
||||
return thrd_error;
|
||||
@@ -2197,8 +2205,13 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
unsigned long long available_memory =
|
||||
pages > 0 && page_size > 0 ? (unsigned long long)pages * (unsigned long long)page_size
|
||||
: 512ULL * 1024 * 1024;
|
||||
unsigned long long avg_file_size = 1024 * 1024;
|
||||
int qsize = (int)(available_memory / avg_file_size);
|
||||
/* Size the chunk queues from the actual chunk size rather than a fixed 1 MiB
|
||||
average: a chunk holds roughly `chunk_size` bytes of file data, so counting
|
||||
chunks at 1 MiB over-estimated the queue capacity by up to 10x. The byte
|
||||
budget below is the authoritative bound; this count keeps the unloaded
|
||||
chunks waiting in the scanner queue bounded too. */
|
||||
unsigned long long chunk_size = config->chunk_size > 0 ? config->chunk_size : DEFAULT_CHUNK_SIZE;
|
||||
int qsize = (int)(available_memory / chunk_size);
|
||||
if (qsize < 10)
|
||||
qsize = 10;
|
||||
if (qsize > 1000)
|
||||
@@ -2223,7 +2236,8 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
}
|
||||
context->missing_args = missing_args;
|
||||
missing_args = NULL; /* owned by the context from here on */
|
||||
*config_ptr = NULL; /* context now owns config through all remaining paths */
|
||||
pipeline_context_sender_set_queue_byte_limit(context, SENDER_QUEUE_MAX_BYTES);
|
||||
*config_ptr = NULL; /* context now owns config through all remaining paths */
|
||||
struct timespec now_mono;
|
||||
if (clock_gettime(CLOCK_MONOTONIC, &now_mono) != 0) {
|
||||
now_mono.tv_sec = 0;
|
||||
|
||||
@@ -388,9 +388,13 @@ static int scanner_inspect_entry(const ScannerOptions* options, const char* sour
|
||||
goto apply_filters;
|
||||
|
||||
regular:
|
||||
if (stat(entry->path, &entry->stats) != 0)
|
||||
goto skip;
|
||||
entry->is_directory = S_ISDIR(entry->stats.st_mode);
|
||||
/* Not a symlink: the lstat() above already described this entry, and lstat
|
||||
and stat are identical for every non-symlink, so reuse that result instead
|
||||
of issuing a redundant stat() on the scanner hot path. stat() is still
|
||||
used on the dereference paths above/below for actual symlinks (copy-links,
|
||||
safe/copy-unsafe links, and -k symlinks-to-directories). */
|
||||
entry->stats = link_stats;
|
||||
entry->is_directory = S_ISDIR(link_stats.st_mode);
|
||||
if (entry->is_directory)
|
||||
return 1;
|
||||
|
||||
|
||||
+182
-45
@@ -2,11 +2,12 @@
|
||||
#include "data.h"
|
||||
#include "log.h"
|
||||
#include "protocol.h"
|
||||
#include <stdlib.h>
|
||||
#include <limits.h>
|
||||
#include <stdint.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <strings.h>
|
||||
#include <threads.h>
|
||||
#include <unistd.h>
|
||||
#include <zstd.h>
|
||||
|
||||
@@ -39,6 +40,96 @@ bool compression_should_skip_with_suffixes(const char* path, char* const* suffix
|
||||
return false;
|
||||
}
|
||||
|
||||
/* 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
|
||||
* stored in a C11 thread-specific storage slot whose destructor releases the
|
||||
* contexts when the thread exits; this keeps LeakSanitizer clean for the
|
||||
* short-lived sender/receiver/scanner worker threads without every worker
|
||||
* entry point having to remember to call compression_free_thread_contexts().
|
||||
* The main thread's slot is not torn down by tss at process exit, so an atexit
|
||||
* hook releases it (and compression_free_thread_contexts allows eager
|
||||
* release). */
|
||||
typedef struct {
|
||||
ZSTD_CCtx* cctx;
|
||||
ZSTD_DCtx* dctx;
|
||||
void* out_buf; /* reusable ZSTD_compressBound-sized output scratch */
|
||||
size_t out_cap; /* bytes currently allocated for out_buf */
|
||||
int level; /* compression level currently applied to cctx */
|
||||
int workers; /* nbWorkers currently applied to cctx */
|
||||
bool params_set;
|
||||
bool cached; /* false when the TSS slot could not be used: caller owns */
|
||||
} CompressionThreadCtx;
|
||||
|
||||
static once_flag compression_tls_once = ONCE_FLAG_INIT;
|
||||
static tss_t compression_tls_key;
|
||||
static bool compression_tls_ready;
|
||||
|
||||
static void compression_tls_make_key(void);
|
||||
|
||||
static void compression_ctx_free(CompressionThreadCtx* ctx) {
|
||||
if (!ctx)
|
||||
return;
|
||||
if (ctx->cctx)
|
||||
ZSTD_freeCCtx(ctx->cctx);
|
||||
if (ctx->dctx)
|
||||
ZSTD_freeDCtx(ctx->dctx);
|
||||
free(ctx->out_buf);
|
||||
free(ctx);
|
||||
}
|
||||
|
||||
static void compression_tls_destructor(void* value) {
|
||||
compression_ctx_free((CompressionThreadCtx*)value);
|
||||
}
|
||||
|
||||
void compression_free_thread_contexts(void) {
|
||||
call_once(&compression_tls_once, compression_tls_make_key);
|
||||
if (!compression_tls_ready)
|
||||
return;
|
||||
CompressionThreadCtx* ctx = (CompressionThreadCtx*)tss_get(compression_tls_key);
|
||||
if (!ctx)
|
||||
return;
|
||||
/* Clear the slot first so the thread-exit destructor cannot free it twice. */
|
||||
tss_set(compression_tls_key, NULL);
|
||||
compression_ctx_free(ctx);
|
||||
}
|
||||
|
||||
static void compression_atexit_cleanup(void) {
|
||||
compression_free_thread_contexts();
|
||||
}
|
||||
|
||||
static void compression_tls_make_key(void) {
|
||||
if (tss_create(&compression_tls_key, compression_tls_destructor) == thrd_success) {
|
||||
compression_tls_ready = true;
|
||||
atexit(compression_atexit_cleanup);
|
||||
}
|
||||
}
|
||||
|
||||
static CompressionThreadCtx* compression_get_thread_ctx(void) {
|
||||
call_once(&compression_tls_once, compression_tls_make_key);
|
||||
if (!compression_tls_ready) {
|
||||
/* Extremely unlikely: fall back to an uncached context the caller frees. */
|
||||
return (CompressionThreadCtx*)calloc(1, sizeof(CompressionThreadCtx));
|
||||
}
|
||||
CompressionThreadCtx* ctx = (CompressionThreadCtx*)tss_get(compression_tls_key);
|
||||
if (ctx)
|
||||
return ctx;
|
||||
ctx = (CompressionThreadCtx*)calloc(1, sizeof(CompressionThreadCtx));
|
||||
if (!ctx)
|
||||
return NULL;
|
||||
ctx->cached = true;
|
||||
if (tss_set(compression_tls_key, ctx) != thrd_success)
|
||||
ctx->cached = false;
|
||||
return ctx;
|
||||
}
|
||||
|
||||
/* Release an uncached context immediately; cached contexts are owned by the
|
||||
* thread's TSS slot and freed on thread exit / compression_free_thread_contexts. */
|
||||
static void compression_ctx_put(CompressionThreadCtx* ctx) {
|
||||
if (ctx && !ctx->cached)
|
||||
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);
|
||||
}
|
||||
@@ -50,68 +141,102 @@ Data* data_compress_with_threads(Data* data_to_compress, int compression_level,
|
||||
return NULL;
|
||||
log_message(LOG_LEVEL_DEBUG, "Starting to compress data");
|
||||
size_t dst_size = ZSTD_compressBound(data_to_compress->size);
|
||||
Data* compressed_data = data_create_empty(dst_size);
|
||||
if (compressed_data == NULL)
|
||||
return NULL;
|
||||
|
||||
ZSTD_CCtx* cctx = ZSTD_createCCtx();
|
||||
if (!cctx) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to create ZSTD compression context");
|
||||
data_destroy(compressed_data);
|
||||
CompressionThreadCtx* ctx = compression_get_thread_ctx();
|
||||
if (ctx == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to allocate ZSTD compression context");
|
||||
return NULL;
|
||||
}
|
||||
Data* compressed_data = NULL;
|
||||
|
||||
size_t zret = ZSTD_CCtx_setParameter(cctx, ZSTD_c_compressionLevel, compression_level);
|
||||
if (ZSTD_isError(zret)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to set compression level: %s", ZSTD_getErrorName(zret));
|
||||
ZSTD_freeCCtx(cctx);
|
||||
data_destroy(compressed_data);
|
||||
return NULL;
|
||||
if (!ctx->cctx) {
|
||||
ctx->cctx = ZSTD_createCCtx();
|
||||
if (!ctx->cctx) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to create ZSTD compression context");
|
||||
goto cleanup;
|
||||
}
|
||||
ctx->params_set = false;
|
||||
}
|
||||
|
||||
/* Reset only the session: parameters (and any already-allocated zstd worker
|
||||
* pool) stay attached to the context, so compressing the next file does not
|
||||
* rebuild the pool. */
|
||||
ZSTD_CCtx_reset(ctx->cctx, ZSTD_reset_session_only);
|
||||
|
||||
if (!ctx->params_set || ctx->level != compression_level) {
|
||||
size_t zret = ZSTD_CCtx_setParameter(ctx->cctx, ZSTD_c_compressionLevel, compression_level);
|
||||
if (ZSTD_isError(zret)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to set compression level: %s", ZSTD_getErrorName(zret));
|
||||
goto cleanup;
|
||||
}
|
||||
ctx->level = compression_level;
|
||||
}
|
||||
|
||||
int available_threads = 0;
|
||||
if (compression_threads > 0) {
|
||||
long online_cpus = sysconf(_SC_NPROCESSORS_ONLN);
|
||||
int available_threads = online_cpus > 0 && online_cpus < compression_threads
|
||||
? (int)online_cpus
|
||||
: compression_threads;
|
||||
zret = ZSTD_CCtx_setParameter(cctx, ZSTD_c_nbWorkers, available_threads);
|
||||
available_threads = online_cpus > 0 && online_cpus < compression_threads ? (int)online_cpus
|
||||
: compression_threads;
|
||||
}
|
||||
if (!ctx->params_set || ctx->workers != available_threads) {
|
||||
size_t zret = ZSTD_CCtx_setParameter(ctx->cctx, ZSTD_c_nbWorkers, available_threads);
|
||||
if (ZSTD_isError(zret)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to set compression threads: %s",
|
||||
ZSTD_getErrorName(zret));
|
||||
ZSTD_freeCCtx(cctx);
|
||||
data_destroy(compressed_data);
|
||||
return NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
ctx->workers = available_threads;
|
||||
}
|
||||
ctx->params_set = true;
|
||||
|
||||
if (available_threads > 0) {
|
||||
/* Streaming compression needs the source size before threaded mode can end a frame. */
|
||||
zret = ZSTD_CCtx_setPledgedSrcSize(cctx, data_to_compress->size);
|
||||
size_t zret = ZSTD_CCtx_setPledgedSrcSize(ctx->cctx, data_to_compress->size);
|
||||
if (ZSTD_isError(zret)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to set compression source size: %s",
|
||||
ZSTD_getErrorName(zret));
|
||||
ZSTD_freeCCtx(cctx);
|
||||
data_destroy(compressed_data);
|
||||
return NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
}
|
||||
|
||||
if (ctx->out_cap < dst_size) {
|
||||
void* grown = protocol_realloc(ctx->out_buf, dst_size);
|
||||
if (grown == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to allocate compression buffer");
|
||||
goto cleanup;
|
||||
}
|
||||
ctx->out_buf = grown;
|
||||
ctx->out_cap = dst_size;
|
||||
}
|
||||
|
||||
ZSTD_inBuffer input = {data_to_compress->data, data_to_compress->size, 0};
|
||||
ZSTD_outBuffer output = {compressed_data->data, dst_size, 0};
|
||||
ZSTD_outBuffer output = {ctx->out_buf, dst_size, 0};
|
||||
|
||||
size_t ret;
|
||||
do {
|
||||
ret = ZSTD_compressStream2(cctx, &output, &input, ZSTD_e_end);
|
||||
ret = ZSTD_compressStream2(ctx->cctx, &output, &input, ZSTD_e_end);
|
||||
if (ZSTD_isError(ret)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Compression failed: %s", ZSTD_getErrorName(ret));
|
||||
ZSTD_freeCCtx(cctx);
|
||||
data_destroy(compressed_data);
|
||||
return NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
} while (ret > 0);
|
||||
|
||||
/* 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);
|
||||
if (compressed_data == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to allocate compressed data");
|
||||
goto cleanup;
|
||||
}
|
||||
if (output.pos > 0)
|
||||
memcpy(compressed_data->data, ctx->out_buf, output.pos);
|
||||
compressed_data->size = output.pos;
|
||||
ZSTD_freeCCtx(cctx);
|
||||
|
||||
log_debug_message(LOG_DEBUG_UTIL, "Data succesfully compressed from %zu to %zu",
|
||||
data_to_compress->size, compressed_data->size);
|
||||
|
||||
cleanup:
|
||||
compression_ctx_put(ctx);
|
||||
return compressed_data;
|
||||
}
|
||||
|
||||
@@ -144,20 +269,30 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
ZSTD_DCtx* dctx = ZSTD_createDCtx();
|
||||
if (!dctx) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to create ZSTD decompression context");
|
||||
CompressionThreadCtx* ctx = compression_get_thread_ctx();
|
||||
if (ctx == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to allocate ZSTD decompression context");
|
||||
return NULL;
|
||||
}
|
||||
Data* uncompressed_data = NULL;
|
||||
|
||||
if (!ctx->dctx) {
|
||||
ctx->dctx = ZSTD_createDCtx();
|
||||
if (!ctx->dctx) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to create ZSTD decompression context");
|
||||
goto cleanup;
|
||||
}
|
||||
}
|
||||
/* Reset only the session; decompression parameters are sticky. */
|
||||
ZSTD_DCtx_reset(ctx->dctx, ZSTD_reset_session_only);
|
||||
|
||||
size_t buf_size = (dst_size > 0) ? (size_t)dst_size : INITIAL_DECOMPRESS_BUF_SIZE;
|
||||
if (buf_size > maximum_size)
|
||||
buf_size = maximum_size;
|
||||
Data* uncompressed_data = data_create_empty(buf_size);
|
||||
uncompressed_data = data_create_empty(buf_size);
|
||||
if (!uncompressed_data) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to allocate decompression buffer");
|
||||
ZSTD_freeDCtx(dctx);
|
||||
return NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
|
||||
ZSTD_inBuffer input = {compressed_data->data, compressed_data->size, 0};
|
||||
@@ -165,20 +300,20 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) {
|
||||
|
||||
size_t ret;
|
||||
do {
|
||||
ret = ZSTD_decompressStream(dctx, &output, &input);
|
||||
ret = ZSTD_decompressStream(ctx->dctx, &output, &input);
|
||||
if (ZSTD_isError(ret)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Decompression failed: %s", ZSTD_getErrorName(ret));
|
||||
ZSTD_freeDCtx(dctx);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
uncompressed_data = NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
if (ret > 0 && output.pos == output.size) {
|
||||
if (buf_size >= hard_limit || buf_size > SIZE_MAX / 2) {
|
||||
log_message(LOG_LEVEL_ERROR, "Decompressed data exceeds %llu bytes",
|
||||
(unsigned long long)MAX_DECOMPRESSED_SIZE);
|
||||
ZSTD_freeDCtx(dctx);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
uncompressed_data = NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
buf_size *= 2;
|
||||
if (buf_size > hard_limit)
|
||||
@@ -186,9 +321,9 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) {
|
||||
void* new_data = protocol_realloc(uncompressed_data->data, buf_size);
|
||||
if (!new_data) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to grow decompression buffer");
|
||||
ZSTD_freeDCtx(dctx);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
uncompressed_data = NULL;
|
||||
goto cleanup;
|
||||
}
|
||||
uncompressed_data->data = new_data;
|
||||
output.dst = new_data;
|
||||
@@ -197,9 +332,11 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) {
|
||||
} while (ret > 0);
|
||||
|
||||
uncompressed_data->size = output.pos;
|
||||
ZSTD_freeDCtx(dctx);
|
||||
|
||||
log_debug_message(LOG_DEBUG_UTIL, "Decompressed data successfully");
|
||||
|
||||
cleanup:
|
||||
compression_ctx_put(ctx);
|
||||
return uncompressed_data;
|
||||
}
|
||||
|
||||
|
||||
@@ -14,4 +14,12 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size);
|
||||
bool compression_should_skip(const char* path);
|
||||
bool compression_should_skip_with_suffixes(const char* path, char* const* suffixes, int count);
|
||||
|
||||
/* Release the calling thread's cached zstd contexts (compressor, decompressor
|
||||
* and scratch buffer). The cache is thread-local and is also released
|
||||
* automatically when a worker thread exits (via a C11 tss destructor) and for
|
||||
* the main thread at process exit; this explicit entry point exists so tests
|
||||
* and long-lived callers can drop the cache deterministically. Safe to call
|
||||
* when no context has been created, and idempotent. */
|
||||
void compression_free_thread_contexts(void);
|
||||
|
||||
#endif
|
||||
|
||||
+17
-2
@@ -642,8 +642,23 @@ typedef struct Config {
|
||||
* anything else) is what keeps a 2.19 client and a 2.18 server from ever
|
||||
* reaching that state. SECURITY: a 2.19 store holds a salted PBKDF2 verifier
|
||||
* and cannot verify (and refuses to load) a legacy unsalted-SHA-256 store line,
|
||||
* so an old bearer digest can never be replayed against a 2.19 daemon. */
|
||||
#define PROTOCOL_VERSION "2.19.0"
|
||||
* so an old bearer digest can never be replayed against a 2.19 daemon.
|
||||
*
|
||||
* Packed Metadata Wave: 2.19.0 -> 2.20.0.
|
||||
*
|
||||
* WHY the bump: metadata_send()/metadata_receive() no longer emit/consume the
|
||||
* metadata as up to 12 separate per-field writes. A file's metadata now
|
||||
* crosses the wire as ONE packed frame: a single int32 present flag (0 =
|
||||
* absent, 1 = present) followed, when present, by the fixed
|
||||
* FILE_METADATA_WIRE_SIZE-byte (68-byte) field record produced by
|
||||
* metadata_to_buf(). Protocol data is an unframed byte stream, so the packed
|
||||
* encoding is byte-for-byte identical to the old field-by-field writes (same
|
||||
* fields, same order, same widths); the change only removes per-field syscalls.
|
||||
* The bump is therefore a deliberate lockstep-release marker, not a
|
||||
* desynchronization fix — the strict same-version handshake still rejects a
|
||||
* mixed 2.19/2.20 deployment. The chunk codec, which already used the packed
|
||||
* metadata_to_buf()/metadata_from_buf() form, is unchanged. */
|
||||
#define PROTOCOL_VERSION "2.20.0"
|
||||
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
|
||||
/* Upper bound on total basis-dir entries (rsync caps --link-dest at 20). */
|
||||
#define MAX_BASIS_DIRS 64
|
||||
|
||||
+50
-19
@@ -104,8 +104,27 @@ static int normalize_entry(const char* raw, size_t len, bool strip_line_endings,
|
||||
return result;
|
||||
}
|
||||
|
||||
/* Build the membership index over the exact entries only. `file_list_affects`
|
||||
combines the exact/descendant lookups with a walk of the query's own ancestor
|
||||
prefixes, so no ancestor prefix is ever materialized as a copy and the index
|
||||
stays O(entry count) memory regardless of path depth. An empty entry (the
|
||||
source root) sets whole_tree and short-circuits every query. */
|
||||
static bool file_list_index_build(FileListSet* set, char* err, size_t err_size) {
|
||||
if (!path_index_build(&set->index, (const char* const*)set->entries, (size_t)set->count)) {
|
||||
snprintf(err, err_size, "memory allocation failed");
|
||||
return false;
|
||||
}
|
||||
for (int i = 0; i < set->count; i++) {
|
||||
if (set->entries[i][0] == '\0') {
|
||||
set->whole_tree = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
static FileListSet* string_list_to_set(StringList* raw, char* err, size_t err_size) {
|
||||
FileListSet* set = malloc(sizeof(FileListSet));
|
||||
FileListSet* set = calloc(1, sizeof(FileListSet));
|
||||
if (!set) {
|
||||
snprintf(err, err_size, "memory allocation failed");
|
||||
return NULL;
|
||||
@@ -114,6 +133,10 @@ static FileListSet* string_list_to_set(StringList* raw, char* err, size_t err_si
|
||||
set->entries = raw->items;
|
||||
raw->items = NULL;
|
||||
raw->count = 0;
|
||||
if (!file_list_index_build(set, err, err_size)) {
|
||||
file_list_destroy(set);
|
||||
return NULL;
|
||||
}
|
||||
return set;
|
||||
}
|
||||
|
||||
@@ -160,34 +183,42 @@ FileListSet* file_list_load(const char* path, bool null_separated, char* err, si
|
||||
void file_list_destroy(FileListSet* set) {
|
||||
if (!set)
|
||||
return;
|
||||
path_index_free(&set->index);
|
||||
for (int i = 0; i < set->count; i++)
|
||||
free(set->entries[i]);
|
||||
free(set->entries);
|
||||
free(set);
|
||||
}
|
||||
|
||||
static bool path_has_prefix(const char* path, const char* prefix) {
|
||||
size_t plen = strlen(prefix);
|
||||
if (strncmp(path, prefix, plen) != 0)
|
||||
return false;
|
||||
return path[plen] == '/' || path[plen] == '\0';
|
||||
}
|
||||
|
||||
bool file_list_affects(const FileListSet* set, const char* rel) {
|
||||
if (!set)
|
||||
return true;
|
||||
if (!rel)
|
||||
return false;
|
||||
for (int i = 0; i < set->count; i++) {
|
||||
const char* entry = set->entries[i];
|
||||
if (entry[0] == '\0')
|
||||
return true; /* whole tree listed */
|
||||
if (strcmp(rel, entry) == 0)
|
||||
return true; /* the entry itself is listed */
|
||||
if (path_has_prefix(rel, entry))
|
||||
return true; /* rel lives under a listed directory */
|
||||
if (path_has_prefix(entry, rel))
|
||||
return true; /* rel is an ancestor directory of a listed entry */
|
||||
if (set->whole_tree)
|
||||
return true; /* whole tree listed */
|
||||
/* An exact entry match means `rel` itself is listed. */
|
||||
if (path_index_contains(&set->index, rel))
|
||||
return true;
|
||||
/* Otherwise `rel` is affected when a listed entry is an ancestor directory of
|
||||
it; walk rel's own directory prefixes (which preserve path-boundary
|
||||
semantics) and test each for an exact entry. No prefixes are stored. */
|
||||
size_t len = strlen(rel);
|
||||
while (len > 0) {
|
||||
const char* slash = NULL;
|
||||
for (size_t i = len; i-- > 0;) {
|
||||
if (rel[i] == '/') {
|
||||
slash = rel + i;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (!slash)
|
||||
break;
|
||||
len = (size_t)(slash - rel);
|
||||
if (path_index_contains_n(&set->index, rel, len))
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
/* Finally `rel` is affected when it is an ancestor directory of a listed
|
||||
entry (binary search for the first entry at or after `rel` + '/'). */
|
||||
return path_index_has_descendant(&set->index, rel);
|
||||
}
|
||||
|
||||
+10
-2
@@ -1,6 +1,7 @@
|
||||
#ifndef FILE_LIST_H
|
||||
#define FILE_LIST_H
|
||||
|
||||
#include "utils.h"
|
||||
#include <stdbool.h>
|
||||
#include <stddef.h>
|
||||
|
||||
@@ -12,11 +13,18 @@
|
||||
* of "." means the whole tree, absolute entries and ".." traversal are
|
||||
* rejected at parse time. The set is immutable and shared read-only across
|
||||
* scanner worker threads.
|
||||
*/
|
||||
|
||||
*
|
||||
* Membership is answered from `index`, built once at load time over the exact
|
||||
* entries only: `index.exact` matches a listed path, the sorted view detects an
|
||||
* ancestor directory of a listed entry, and `rel`'s own directory prefixes are
|
||||
* matched against the exact set while descending. No ancestor prefix is stored
|
||||
* as a separate string, so the index is O(entry count) memory however deep the
|
||||
* paths are, and each query is O(path length) comparisons. */
|
||||
typedef struct {
|
||||
char** entries; /* normalized rel paths; "" means the whole tree */
|
||||
int count;
|
||||
PathIndex index;
|
||||
bool whole_tree; /* an entry of "" lists the source root */
|
||||
} FileListSet;
|
||||
|
||||
/* Load and validate a --files-from file. When `null_separated` (-0/--from0)
|
||||
|
||||
+21
-124
@@ -157,33 +157,17 @@ FileMetadata* metadata_from_buf(char** buf) {
|
||||
|
||||
bool metadata_send(int file_descriptor, const FileMetadata* m) {
|
||||
if (m == NULL) {
|
||||
int32_t zero = 0;
|
||||
return send_n_data(file_descriptor, &zero, sizeof(zero));
|
||||
int32_t absent = 0;
|
||||
return send_n_data(file_descriptor, &absent, sizeof(absent));
|
||||
}
|
||||
int32_t present = 1;
|
||||
int32_t mode = (int32_t)m->mode;
|
||||
int32_t uid = (int32_t)m->uid;
|
||||
int32_t gid = (int32_t)m->gid;
|
||||
int64_t mtime_sec = (int64_t)m->mtime_sec;
|
||||
int64_t mtime_nsec = (int64_t)m->mtime_nsec;
|
||||
int32_t atime_valid = m->atime_valid ? 1 : 0;
|
||||
int64_t atime_sec = (int64_t)m->atime_sec;
|
||||
int64_t atime_nsec = (int64_t)m->atime_nsec;
|
||||
int32_t crtime_valid = m->crtime_valid ? 1 : 0;
|
||||
int64_t crtime_sec = (int64_t)m->crtime_sec;
|
||||
int64_t crtime_nsec = (int64_t)m->crtime_nsec;
|
||||
return send_n_data(file_descriptor, &present, sizeof(present)) &&
|
||||
send_n_data(file_descriptor, &mode, sizeof(mode)) &&
|
||||
send_n_data(file_descriptor, &uid, sizeof(uid)) &&
|
||||
send_n_data(file_descriptor, &gid, sizeof(gid)) &&
|
||||
send_n_data(file_descriptor, &mtime_sec, sizeof(mtime_sec)) &&
|
||||
send_n_data(file_descriptor, &mtime_nsec, sizeof(mtime_nsec)) &&
|
||||
send_n_data(file_descriptor, &atime_valid, sizeof(atime_valid)) &&
|
||||
send_n_data(file_descriptor, &atime_sec, sizeof(atime_sec)) &&
|
||||
send_n_data(file_descriptor, &atime_nsec, sizeof(atime_nsec)) &&
|
||||
send_n_data(file_descriptor, &crtime_valid, sizeof(crtime_valid)) &&
|
||||
send_n_data(file_descriptor, &crtime_sec, sizeof(crtime_sec)) &&
|
||||
send_n_data(file_descriptor, &crtime_nsec, sizeof(crtime_nsec));
|
||||
/* One packed frame (protocol 2.20.0): the int32 present flag followed by the
|
||||
fixed FILE_METADATA_WIRE_SIZE-byte field record. metadata_to_buf() emits
|
||||
exactly that layout (present + fields), so build it once and write the
|
||||
whole record in a single call instead of one frame per field. */
|
||||
char packed[sizeof(int32_t) + FILE_METADATA_WIRE_SIZE];
|
||||
char* cursor = packed;
|
||||
metadata_to_buf(&cursor, m);
|
||||
return send_n_data(file_descriptor, packed, sizeof(packed));
|
||||
}
|
||||
|
||||
FileMetadata* metadata_receive(int file_descriptor, int* ok) {
|
||||
@@ -203,109 +187,22 @@ FileMetadata* metadata_receive(int file_descriptor, int* ok) {
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
FileMetadata* m = protocol_alloc(sizeof(FileMetadata));
|
||||
/* Rebuild the packed record metadata_from_buf() expects: the present flag we
|
||||
just read, followed by exactly FILE_METADATA_WIRE_SIZE field bytes. */
|
||||
char packed[sizeof(int32_t) + FILE_METADATA_WIRE_SIZE];
|
||||
memcpy(packed, &present, sizeof(present));
|
||||
if (!receive_n_data(file_descriptor, packed + sizeof(present), FILE_METADATA_WIRE_SIZE)) {
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
char* cursor = packed;
|
||||
FileMetadata* m = metadata_from_buf(&cursor);
|
||||
if (m == NULL) {
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
int32_t mode;
|
||||
if (!receive_n_data(file_descriptor, &mode, sizeof(mode))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
m->mode = (mode_t)mode;
|
||||
int32_t uid;
|
||||
if (!receive_n_data(file_descriptor, &uid, sizeof(uid))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
m->uid = (uid_t)uid;
|
||||
int32_t gid;
|
||||
if (!receive_n_data(file_descriptor, &gid, sizeof(gid))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
m->gid = (gid_t)gid;
|
||||
int64_t mtime_sec;
|
||||
if (!receive_n_data(file_descriptor, &mtime_sec, sizeof(mtime_sec))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
m->mtime_sec = (time_t)mtime_sec;
|
||||
int64_t mtime_nsec;
|
||||
if (!receive_n_data(file_descriptor, &mtime_nsec, sizeof(mtime_nsec))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
m->mtime_nsec = (long)mtime_nsec;
|
||||
int32_t atime_valid;
|
||||
if (!receive_n_data(file_descriptor, &atime_valid, sizeof(atime_valid))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
int64_t atime_sec;
|
||||
if (!receive_n_data(file_descriptor, &atime_sec, sizeof(atime_sec))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
int64_t atime_nsec;
|
||||
if (!receive_n_data(file_descriptor, &atime_nsec, sizeof(atime_nsec))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
int32_t crtime_valid;
|
||||
if (!receive_n_data(file_descriptor, &crtime_valid, sizeof(crtime_valid))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
int64_t crtime_sec;
|
||||
if (!receive_n_data(file_descriptor, &crtime_sec, sizeof(crtime_sec))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
int64_t crtime_nsec;
|
||||
if (!receive_n_data(file_descriptor, &crtime_nsec, sizeof(crtime_nsec))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
m->atime_valid = atime_valid != 0;
|
||||
m->atime_sec = (time_t)atime_sec;
|
||||
m->atime_nsec = (long)atime_nsec;
|
||||
m->crtime_valid = crtime_valid != 0;
|
||||
m->crtime_sec = (time_t)crtime_sec;
|
||||
m->crtime_nsec = (long)crtime_nsec;
|
||||
if (mtime_nsec < 0 || mtime_nsec >= 1000000000LL || mode < 0 || uid < 0 || gid < 0 ||
|
||||
atime_valid < 0 || atime_valid > 1 || crtime_valid < 0 || crtime_valid > 1 ||
|
||||
(atime_valid && (atime_nsec < 0 || atime_nsec >= 1000000000LL)) ||
|
||||
(crtime_valid && (crtime_nsec < 0 || crtime_nsec >= 1000000000LL))) {
|
||||
free(m);
|
||||
if (ok)
|
||||
*ok = 0;
|
||||
return NULL;
|
||||
}
|
||||
if (ok)
|
||||
*ok = 1;
|
||||
return m;
|
||||
|
||||
@@ -29,7 +29,13 @@
|
||||
|
||||
/* Size of metadata fields on wire, excluding the int32_t `present` field that
|
||||
* is always sent first. The total wire size for present metadata is
|
||||
* sizeof(int32_t) + FILE_METADATA_WIRE_SIZE (68 bytes on most platforms). */
|
||||
* sizeof(int32_t) + FILE_METADATA_WIRE_SIZE (68 bytes on most platforms).
|
||||
*
|
||||
* metadata_send()/metadata_receive() (protocol 2.20.0) frame the metadata as a
|
||||
* single packed record: one int32 present flag (0 = absent) followed, when
|
||||
* present, by exactly FILE_METADATA_WIRE_SIZE bytes of field data. This is the
|
||||
* same present+fields byte layout metadata_to_buf()/metadata_from_buf() use, so
|
||||
* the wire metadata is now one frame instead of one frame per field. */
|
||||
#define FILE_METADATA_WIRE_SIZE (sizeof(int32_t) * 5 + sizeof(int64_t) * 6)
|
||||
|
||||
void metadata_to_buf(char** buf, const FileMetadata* m);
|
||||
|
||||
@@ -11,6 +11,7 @@
|
||||
#include "protocol.h"
|
||||
#include "queue.h"
|
||||
#include "utils.h"
|
||||
#include <stdint.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
@@ -26,6 +27,8 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
|
||||
context->queue_loader = queue_loader;
|
||||
context->scanner_done = false;
|
||||
context->loader_done = false;
|
||||
context->queued_bytes = 0;
|
||||
context->max_queue_bytes = 0;
|
||||
context->manifest = NULL;
|
||||
context->excluded_paths = NULL;
|
||||
context->missing_args = NULL;
|
||||
@@ -97,6 +100,84 @@ fail:
|
||||
return NULL;
|
||||
}
|
||||
|
||||
void pipeline_context_sender_set_queue_byte_limit(PipelineContextSender* context,
|
||||
size_t max_bytes) {
|
||||
if (context == NULL)
|
||||
return;
|
||||
mtx_lock(&context->mutex_loader);
|
||||
context->max_queue_bytes = max_bytes;
|
||||
context->queued_bytes = 0;
|
||||
cnd_broadcast(&context->condition_not_full_loader);
|
||||
mtx_unlock(&context->mutex_loader);
|
||||
}
|
||||
|
||||
size_t pipeline_context_sender_chunk_bytes(const Chunk* chunk) {
|
||||
if (chunk == NULL || chunk->items == NULL)
|
||||
return 0;
|
||||
size_t total = 0;
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
const File* file = chunk->items[i];
|
||||
if (file == NULL || file->data == NULL || file->data->data == NULL)
|
||||
continue;
|
||||
if (file->data->size > SIZE_MAX - total)
|
||||
return SIZE_MAX;
|
||||
total += file->data->size;
|
||||
}
|
||||
return total;
|
||||
}
|
||||
|
||||
void pipeline_context_sender_note_bytes_released(PipelineContextSender* context,
|
||||
size_t released_bytes) {
|
||||
if (context == NULL || context->max_queue_bytes == 0 || released_bytes == 0)
|
||||
return;
|
||||
mtx_lock(&context->mutex_loader);
|
||||
if (released_bytes >= context->queued_bytes)
|
||||
context->queued_bytes = 0;
|
||||
else
|
||||
context->queued_bytes -= released_bytes;
|
||||
cnd_signal(&context->condition_not_full_loader);
|
||||
mtx_unlock(&context->mutex_loader);
|
||||
}
|
||||
|
||||
bool pipeline_context_sender_enqueue_chunk(PipelineContextSender* context, Chunk* chunk) {
|
||||
if (context == NULL || chunk == NULL)
|
||||
return false;
|
||||
size_t chunk_bytes = pipeline_context_sender_chunk_bytes(chunk);
|
||||
mtx_lock(&context->mutex_loader);
|
||||
while (!atomic_load(&context->cancelled)) {
|
||||
bool blocked_by_count = queue_is_full(context->queue_loader);
|
||||
bool blocked_by_budget = false;
|
||||
if (context->max_queue_bytes > 0) {
|
||||
size_t budget = context->max_queue_bytes;
|
||||
size_t used = context->queued_bytes;
|
||||
if (used >= budget) {
|
||||
blocked_by_budget = true;
|
||||
} else if (chunk_bytes > budget - used) {
|
||||
/* A single payload larger than the whole budget is only admitted to an
|
||||
empty pipeline so the wait can never deadlock. */
|
||||
blocked_by_budget = used != 0;
|
||||
}
|
||||
}
|
||||
if (!blocked_by_count && !blocked_by_budget)
|
||||
break;
|
||||
cnd_wait(&context->condition_not_full_loader, &context->mutex_loader);
|
||||
}
|
||||
if (atomic_load(&context->cancelled)) {
|
||||
mtx_unlock(&context->mutex_loader);
|
||||
chunk_destroy(chunk);
|
||||
return false;
|
||||
}
|
||||
if (!queue_enqueue(context->queue_loader, chunk)) {
|
||||
mtx_unlock(&context->mutex_loader);
|
||||
chunk_destroy(chunk);
|
||||
return false;
|
||||
}
|
||||
context->queued_bytes += chunk_bytes;
|
||||
cnd_signal(&context->condition_not_empty_loader);
|
||||
mtx_unlock(&context->mutex_loader);
|
||||
return true;
|
||||
}
|
||||
|
||||
void pipeline_context_sender_destroy(PipelineContextSender* context) {
|
||||
if (context->manifest) {
|
||||
array_list_delete(context->manifest);
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
#include <stdatomic.h>
|
||||
|
||||
#include "array_list.h"
|
||||
#include "chunk.h"
|
||||
#include "config.h"
|
||||
#include "file.h"
|
||||
#include "protocol.h"
|
||||
@@ -25,6 +26,15 @@ typedef struct {
|
||||
cnd_t condition_not_full_loader;
|
||||
cnd_t condition_not_empty_loader;
|
||||
bool loader_done;
|
||||
/* Aggregate loaded payload bytes queued on queue_loader but not yet released
|
||||
by the sender. Guarded by `mutex_loader`. When `max_queue_bytes` is
|
||||
non-zero the loader blocks before enqueueing a chunk that would push this
|
||||
total over it, so the sender buffers a bounded number of bytes rather than
|
||||
an unbounded count of chunks that may each be up to chunk_size (or a single
|
||||
file) in size. Files streamed straight from disk by sendfile hold no
|
||||
payload, so only in-memory (`data->data`) payloads are counted. */
|
||||
size_t queued_bytes;
|
||||
size_t max_queue_bytes;
|
||||
ArrayList* manifest;
|
||||
/* Protected prefixes (paths the source scan excluded by user rules) sent
|
||||
with the keep-set manifest so --delete leaves them alone unless
|
||||
@@ -113,6 +123,22 @@ typedef struct PipelineContextReceiver {
|
||||
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
|
||||
Queue* queue_loader);
|
||||
void pipeline_context_sender_destroy(PipelineContextSender* context);
|
||||
/* Bound the loaded payload bytes the sender may buffer ahead of the network
|
||||
writer (see max_queue_bytes). */
|
||||
void pipeline_context_sender_set_queue_byte_limit(PipelineContextSender* context, size_t max_bytes);
|
||||
/* Total payload bytes a chunk currently holds in memory (loaded file data
|
||||
only; zero for entries with no payload or data streamed from disk). */
|
||||
size_t pipeline_context_sender_chunk_bytes(const Chunk* chunk);
|
||||
/* Blocking enqueue used by the sender's loader stage. Blocks while
|
||||
queue_loader is full by element count or when adding `chunk` would push the
|
||||
queued payload bytes over the configured byte limit; waits until the sender
|
||||
releases bytes. Takes ownership of `chunk` on success and destroys it on
|
||||
failure/cancel. */
|
||||
bool pipeline_context_sender_enqueue_chunk(PipelineContextSender* context, Chunk* chunk);
|
||||
/* Account for `released_bytes` of payload memory that the sender freed after
|
||||
destroying a chunk, unblocking a loader waiting on the byte limit. */
|
||||
void pipeline_context_sender_note_bytes_released(PipelineContextSender* context,
|
||||
size_t released_bytes);
|
||||
PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver,
|
||||
int file_descriptor, SSL* ssl);
|
||||
void pipeline_context_receiver_destroy(PipelineContextReceiver* context);
|
||||
|
||||
@@ -19,6 +19,7 @@
|
||||
static volatile sig_atomic_t g_active_connections = 0;
|
||||
|
||||
static void tcp_apply_socket_timeout(int fd);
|
||||
static void tcp_enable_nodelay_default(int fd, int family);
|
||||
|
||||
static void sigchld_handler(int sig) {
|
||||
(void)sig;
|
||||
@@ -148,6 +149,7 @@ static void accept_loop(Server* server, void (*child_fn)(int, void*), void* chil
|
||||
continue;
|
||||
}
|
||||
tcp_apply_socket_timeout(fd);
|
||||
tcp_enable_nodelay_default(fd, client_addr.ss_family);
|
||||
char peer[128];
|
||||
if (!utils_sockaddr_to_string((const struct sockaddr*)&client_addr, peer, sizeof(peer)))
|
||||
snprintf(peer, sizeof(peer), "unknown");
|
||||
@@ -231,6 +233,19 @@ static void tcp_apply_socket_timeout(int fd) {
|
||||
setsockopt(fd, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
|
||||
}
|
||||
|
||||
/* Enable TCP_NODELAY by default on a transfer socket: the protocol emits many
|
||||
* small messages and Nagle's algorithm would otherwise coalesce/delay them.
|
||||
* Best-effort only: the family guard keeps this to IP/TCP sockets, and a
|
||||
* setsockopt failure is ignored. A caller-provided --sockopts TCP_NODELAY=0
|
||||
* is applied afterwards on the connect path, so an explicit user choice still
|
||||
* wins. */
|
||||
static void tcp_enable_nodelay_default(int fd, int family) {
|
||||
if (family != AF_INET && family != AF_INET6)
|
||||
return;
|
||||
int value = 1;
|
||||
setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &value, sizeof(value));
|
||||
}
|
||||
|
||||
Client* client_create() {
|
||||
Client* client = (Client*)malloc(sizeof(Client));
|
||||
if (client == NULL) {
|
||||
@@ -376,6 +391,9 @@ bool tcp_connect_socket_ex(Client* client, const char* host, int port,
|
||||
if (client->file_descriptor < 0)
|
||||
continue;
|
||||
|
||||
/* Default first; a user --sockopts TCP_NODELAY=0 applied below overrides. */
|
||||
tcp_enable_nodelay_default(client->file_descriptor, rp->ai_family);
|
||||
|
||||
if (opts && opts->sockopt_count > 0 &&
|
||||
!tcp_apply_sockopts(client->file_descriptor, opts->sockopts, opts->sockopt_count)) {
|
||||
close(client->file_descriptor);
|
||||
|
||||
+278
-36
@@ -13,6 +13,7 @@
|
||||
#include <sys/socket.h>
|
||||
#include <sys/stat.h>
|
||||
#include <unistd.h>
|
||||
#include <xxhash.h>
|
||||
|
||||
static int authorized_root_fd = -1;
|
||||
static char* authorized_root_path;
|
||||
@@ -96,6 +97,241 @@ char* str_dup(const char* string) {
|
||||
return new_string;
|
||||
}
|
||||
|
||||
#define STR_HASH_SET_MIN_CAPACITY 16
|
||||
|
||||
static size_t str_hash_set_hash(const char* key, size_t len) {
|
||||
return (size_t)XXH64(key, len, 0);
|
||||
}
|
||||
|
||||
/* Store a borrowed key. Returns 1 when a new slot was filled and 0 for a
|
||||
* duplicate. */
|
||||
static int str_hash_set_put(StrHashSet* set, const char* key, size_t len) {
|
||||
size_t mask = set->capacity - 1;
|
||||
size_t index = str_hash_set_hash(key, len) & mask;
|
||||
while (true) {
|
||||
StrHashSetSlot* slot = &set->slots[index];
|
||||
if (!slot->key) {
|
||||
slot->key = key;
|
||||
set->size++;
|
||||
return 1;
|
||||
}
|
||||
if (strlen(slot->key) == len && memcmp(slot->key, key, len) == 0)
|
||||
return 0;
|
||||
index = (index + 1) & mask;
|
||||
}
|
||||
}
|
||||
|
||||
static bool str_hash_set_resize(StrHashSet* set, size_t new_capacity) {
|
||||
StrHashSetSlot* old_slots = set->slots;
|
||||
size_t old_capacity = set->capacity;
|
||||
StrHashSetSlot* slots = calloc(new_capacity, sizeof(StrHashSetSlot));
|
||||
if (!slots)
|
||||
return false;
|
||||
set->slots = slots;
|
||||
set->capacity = new_capacity;
|
||||
set->size = 0;
|
||||
for (size_t i = 0; i < old_capacity; i++) {
|
||||
if (old_slots[i].key)
|
||||
(void)str_hash_set_put(set, old_slots[i].key, strlen(old_slots[i].key));
|
||||
}
|
||||
free(old_slots);
|
||||
return true;
|
||||
}
|
||||
|
||||
static bool str_hash_set_grow(StrHashSet* set) {
|
||||
if (set->capacity != 0 && (set->size + 1) * 4 <= set->capacity * 3)
|
||||
return true;
|
||||
size_t new_capacity = set->capacity ? set->capacity * 2 : STR_HASH_SET_MIN_CAPACITY;
|
||||
return str_hash_set_resize(set, new_capacity);
|
||||
}
|
||||
|
||||
bool str_hash_set_init(StrHashSet* set, size_t hint) {
|
||||
if (!set)
|
||||
return false;
|
||||
set->slots = NULL;
|
||||
set->capacity = 0;
|
||||
set->size = 0;
|
||||
size_t capacity = STR_HASH_SET_MIN_CAPACITY;
|
||||
while (capacity < (hint + 1) * 2 && capacity <= SIZE_MAX / 2)
|
||||
capacity *= 2;
|
||||
set->slots = calloc(capacity, sizeof(StrHashSetSlot));
|
||||
if (!set->slots)
|
||||
return false;
|
||||
set->capacity = capacity;
|
||||
return true;
|
||||
}
|
||||
|
||||
void str_hash_set_free(StrHashSet* set) {
|
||||
if (!set)
|
||||
return;
|
||||
free(set->slots);
|
||||
set->slots = NULL;
|
||||
set->capacity = 0;
|
||||
set->size = 0;
|
||||
}
|
||||
|
||||
bool str_hash_set_insert_ref(StrHashSet* set, const char* key) {
|
||||
if (!set || !key)
|
||||
return false;
|
||||
if (!str_hash_set_grow(set))
|
||||
return false;
|
||||
/* put() returns 1 for a new slot and 0 for a duplicate; both are success. */
|
||||
(void)str_hash_set_put(set, key, strlen(key));
|
||||
return true;
|
||||
}
|
||||
|
||||
static const StrHashSetSlot* str_hash_set_find_n(const StrHashSet* set, const char* key,
|
||||
size_t len) {
|
||||
if (!set || set->capacity == 0 || !key)
|
||||
return NULL;
|
||||
size_t mask = set->capacity - 1;
|
||||
size_t index = str_hash_set_hash(key, len) & mask;
|
||||
while (true) {
|
||||
const StrHashSetSlot* slot = &set->slots[index];
|
||||
if (!slot->key)
|
||||
return NULL;
|
||||
if (strlen(slot->key) == len && memcmp(slot->key, key, len) == 0)
|
||||
return slot;
|
||||
index = (index + 1) & mask;
|
||||
}
|
||||
}
|
||||
|
||||
bool str_hash_set_lookup_n(const StrHashSet* set, const char* key, size_t len) {
|
||||
return str_hash_set_find_n(set, key, len) != NULL;
|
||||
}
|
||||
|
||||
bool str_hash_set_lookup(const StrHashSet* set, const char* key) {
|
||||
if (!key)
|
||||
return false;
|
||||
return str_hash_set_lookup_n(set, key, strlen(key));
|
||||
}
|
||||
|
||||
static int str_sorted_array_compare(const void* left, const void* right) {
|
||||
const char* const* left_key = left;
|
||||
const char* const* right_key = right;
|
||||
return strcmp(*left_key, *right_key);
|
||||
}
|
||||
|
||||
bool str_sorted_array_build(StrSortedArray* array, const char* const* items, size_t count) {
|
||||
if (!array)
|
||||
return false;
|
||||
array->items = NULL;
|
||||
array->count = 0;
|
||||
if (count == 0)
|
||||
return true;
|
||||
if (!items || count > SIZE_MAX / sizeof(const char*))
|
||||
return false;
|
||||
const char** sorted = malloc(count * sizeof(*sorted));
|
||||
if (!sorted)
|
||||
return false;
|
||||
for (size_t i = 0; i < count; i++)
|
||||
sorted[i] = items[i];
|
||||
qsort(sorted, count, sizeof(*sorted), str_sorted_array_compare);
|
||||
array->items = sorted;
|
||||
array->count = count;
|
||||
return true;
|
||||
}
|
||||
|
||||
void str_sorted_array_free(StrSortedArray* array) {
|
||||
if (!array)
|
||||
return;
|
||||
free(array->items);
|
||||
array->items = NULL;
|
||||
array->count = 0;
|
||||
}
|
||||
|
||||
bool str_sorted_array_contains(const StrSortedArray* array, const char* key) {
|
||||
if (!array || !key || array->count == 0)
|
||||
return false;
|
||||
size_t lo = 0;
|
||||
size_t hi = array->count;
|
||||
while (lo < hi) {
|
||||
size_t mid = lo + (hi - lo) / 2;
|
||||
int cmp = strcmp(array->items[mid], key);
|
||||
if (cmp < 0)
|
||||
lo = mid + 1;
|
||||
else if (cmp > 0)
|
||||
hi = mid;
|
||||
else
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/* Compare `entry` against the virtual key `key` + '/' without allocating the
|
||||
* concatenation. Returns <0, 0 or >0 as `entry` sorts before, equal to, or
|
||||
* after that virtual key. */
|
||||
static int str_sorted_array_compare_prefix(const char* entry, const char* key, size_t key_len) {
|
||||
int cmp = strncmp(entry, key, key_len);
|
||||
if (cmp != 0)
|
||||
return cmp;
|
||||
unsigned char next = (unsigned char)entry[key_len];
|
||||
if (next == '\0')
|
||||
return -1; /* entry == key sorts before key + '/' */
|
||||
return (int)next - (int)'/';
|
||||
}
|
||||
|
||||
bool str_sorted_array_has_child_prefix(const StrSortedArray* array, const char* key) {
|
||||
if (!array || !key || array->count == 0 || key[0] == '\0')
|
||||
return false;
|
||||
size_t key_len = strlen(key);
|
||||
size_t lo = 0;
|
||||
size_t hi = array->count;
|
||||
while (lo < hi) {
|
||||
size_t mid = lo + (hi - lo) / 2;
|
||||
if (str_sorted_array_compare_prefix(array->items[mid], key, key_len) < 0)
|
||||
lo = mid + 1;
|
||||
else
|
||||
hi = mid;
|
||||
}
|
||||
if (lo >= array->count)
|
||||
return false;
|
||||
const char* entry = array->items[lo];
|
||||
return strncmp(entry, key, key_len) == 0 && entry[key_len] == '/';
|
||||
}
|
||||
|
||||
bool path_index_build(PathIndex* index, const char* const* entries, size_t count) {
|
||||
if (!index)
|
||||
return false;
|
||||
index->exact.slots = NULL;
|
||||
index->exact.capacity = 0;
|
||||
index->exact.size = 0;
|
||||
index->sorted.items = NULL;
|
||||
index->sorted.count = 0;
|
||||
if (!str_hash_set_init(&index->exact, count))
|
||||
return false;
|
||||
if (!str_sorted_array_build(&index->sorted, entries, count)) {
|
||||
str_hash_set_free(&index->exact);
|
||||
return false;
|
||||
}
|
||||
for (size_t i = 0; i < count; i++) {
|
||||
if (!str_hash_set_insert_ref(&index->exact, entries[i])) {
|
||||
path_index_free(index);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
void path_index_free(PathIndex* index) {
|
||||
if (!index)
|
||||
return;
|
||||
str_hash_set_free(&index->exact);
|
||||
str_sorted_array_free(&index->sorted);
|
||||
}
|
||||
|
||||
bool path_index_contains(const PathIndex* index, const char* path) {
|
||||
return index && str_hash_set_lookup(&index->exact, path);
|
||||
}
|
||||
|
||||
bool path_index_contains_n(const PathIndex* index, const char* path, size_t len) {
|
||||
return index && str_hash_set_lookup_n(&index->exact, path, len);
|
||||
}
|
||||
|
||||
bool path_index_has_descendant(const PathIndex* index, const char* path) {
|
||||
return index && str_sorted_array_has_child_prefix(&index->sorted, path);
|
||||
}
|
||||
|
||||
char* output_escape(const char* string, bool eight_bit_output) {
|
||||
if (!string)
|
||||
return NULL;
|
||||
@@ -196,15 +432,23 @@ bool format_human_bytes(unsigned long long bytes, char* buffer, size_t buffer_si
|
||||
return written >= 0 && (size_t)written < buffer_size;
|
||||
}
|
||||
|
||||
static bool is_dir_in_manifest(const char* rel_path, ArrayList* manifest) {
|
||||
size_t len = strlen(rel_path);
|
||||
for (int i = 0; i < manifest->size; i++) {
|
||||
const char* entry = (const char*)manifest->items[i];
|
||||
// Check if entry starts with rel_path + '/' or matches exactly
|
||||
if (strncmp(entry, rel_path, len) == 0 && (entry[len] == '/' || entry[len] == '\0'))
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
/* Build the keep-set index from the exact manifest entries only. A lookup of
|
||||
`rel` succeeds iff `rel` is a kept entry, a kept directory, or an ancestor
|
||||
directory of kept content (the old is_dir_in_manifest predicate); the sorted
|
||||
view answers "is an ancestor of kept content" without materializing any
|
||||
per-component prefix copy, so the index is O(manifest size) memory. */
|
||||
static bool build_keep_index(const ArrayList* manifest, PathIndex* index) {
|
||||
if (!manifest || manifest->size <= 0)
|
||||
return path_index_build(index, NULL, 0);
|
||||
return path_index_build(index, (const char* const*)manifest->items, (size_t)manifest->size);
|
||||
}
|
||||
|
||||
static bool keep_is_dir(const PathIndex* index, const char* rel_path) {
|
||||
return path_index_contains(index, rel_path) || path_index_has_descendant(index, rel_path);
|
||||
}
|
||||
|
||||
static bool keep_is_file(const PathIndex* index, const char* rel_path) {
|
||||
return path_index_contains(index, rel_path);
|
||||
}
|
||||
|
||||
/* True when child_rel is, or lies below, a protected entry. A prefix "a"
|
||||
@@ -234,7 +478,7 @@ bool path_under_skip_prefix(const char* child_rel, bool at_root, const DeleteSki
|
||||
prefixes) mark the enclosing directory as surviving, exactly as they would
|
||||
make a real rmdir fail with ENOTEMPTY. Stops early once *count reaches the
|
||||
cap (sets *exceeds). Returns false on a traversal error. */
|
||||
static bool count_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest, size_t cap,
|
||||
static bool count_extras_fd(int dirfd, const char* rel_path, const PathIndex* keep, size_t cap,
|
||||
size_t* count, bool* exceeds, const DeleteSkipEntry* skips,
|
||||
int skip_count, bool* survives) {
|
||||
/* openat(dirfd, ".") opens an independent file description: a dup() would
|
||||
@@ -284,15 +528,15 @@ static bool count_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest
|
||||
bool child_ok = true;
|
||||
bool child_survives = true;
|
||||
if (childfd >= 0) {
|
||||
child_ok = count_extras_fd(childfd, child_rel, manifest, cap, count, exceeds, skips,
|
||||
skip_count, &child_survives);
|
||||
child_ok = count_extras_fd(childfd, child_rel, keep, cap, count, exceeds, skips, skip_count,
|
||||
&child_survives);
|
||||
close(childfd);
|
||||
} else if (errno != ENOENT) {
|
||||
operation_ok = false;
|
||||
}
|
||||
if (!child_ok)
|
||||
operation_ok = false;
|
||||
if (is_dir_in_manifest(child_rel, manifest)) {
|
||||
if (keep_is_dir(keep, child_rel)) {
|
||||
/* A directory with kept content below it is never removed. */
|
||||
local_survives = true;
|
||||
} else if (child_survives) {
|
||||
@@ -308,13 +552,7 @@ static bool count_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest
|
||||
}
|
||||
}
|
||||
} else {
|
||||
bool found = false;
|
||||
for (int i = 0; i < manifest->size; i++) {
|
||||
if (strcmp((char*)manifest->items[i], child_rel) == 0) {
|
||||
found = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
bool found = keep_is_file(keep, child_rel);
|
||||
if (!found) {
|
||||
if (*count >= cap) {
|
||||
*exceeds = true;
|
||||
@@ -330,7 +568,7 @@ static bool count_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest
|
||||
return operation_ok;
|
||||
}
|
||||
|
||||
static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest,
|
||||
static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* keep,
|
||||
size_t max_delete, size_t* deleted_count, const DeleteSkipEntry* skips,
|
||||
int skip_count) {
|
||||
/* Independent file description (see count_extras_fd). */
|
||||
@@ -379,15 +617,15 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
|
||||
int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
|
||||
bool child_removed = false;
|
||||
if (childfd >= 0) {
|
||||
child_removed = delete_extras_fd(childfd, child_rel, manifest, max_delete, deleted_count,
|
||||
skips, skip_count);
|
||||
child_removed = delete_extras_fd(childfd, child_rel, keep, max_delete, deleted_count, skips,
|
||||
skip_count);
|
||||
if (!child_removed)
|
||||
operation_ok = false;
|
||||
close(childfd);
|
||||
} else if (errno != ENOENT) {
|
||||
operation_ok = false;
|
||||
}
|
||||
if (child_removed && !is_dir_in_manifest(child_rel, manifest)) {
|
||||
if (child_removed && !keep_is_dir(keep, child_rel)) {
|
||||
if (*deleted_count >= max_delete) {
|
||||
operation_ok = false;
|
||||
} else {
|
||||
@@ -406,13 +644,7 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
|
||||
}
|
||||
} else {
|
||||
// Check if relative path is in manifest
|
||||
bool found = false;
|
||||
for (int i = 0; i < manifest->size; i++) {
|
||||
if (strcmp((char*)manifest->items[i], child_rel) == 0) {
|
||||
found = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
bool found = keep_is_file(keep, child_rel);
|
||||
if (!found) {
|
||||
if (*deleted_count >= max_delete) {
|
||||
operation_ok = false;
|
||||
@@ -436,13 +668,18 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
|
||||
return operation_ok;
|
||||
}
|
||||
|
||||
DeleteWalkResult delete_extras_limited(const char* dest_root, ArrayList* manifest,
|
||||
DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest,
|
||||
size_t max_delete, const DeleteSkipEntry* skips,
|
||||
int skip_count, size_t* deleted_out) {
|
||||
if (deleted_out)
|
||||
*deleted_out = 0;
|
||||
if (!manifest)
|
||||
return DELETE_WALK_ERROR;
|
||||
/* Index the keep-set once so both passes answer membership in O(path length)
|
||||
instead of scanning every manifest entry for every destination entry. */
|
||||
PathIndex keep;
|
||||
if (!build_keep_index(manifest, &keep))
|
||||
return DELETE_WALK_ERROR;
|
||||
int rootfd;
|
||||
if (authorized_root_fd >= 0) {
|
||||
if (authorized_root_path)
|
||||
@@ -454,35 +691,40 @@ DeleteWalkResult delete_extras_limited(const char* dest_root, ArrayList* manifes
|
||||
} else {
|
||||
rootfd = open(dest_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
|
||||
}
|
||||
if (rootfd < 0)
|
||||
if (rootfd < 0) {
|
||||
path_index_free(&keep);
|
||||
return DELETE_WALK_ERROR;
|
||||
}
|
||||
if (max_delete != SIZE_MAX) {
|
||||
/* Rehearse the deletion first so a run that would exceed the cap removes
|
||||
nothing (rsync's all-or-nothing --max-delete contract). */
|
||||
size_t count = 0;
|
||||
bool exceeds = false;
|
||||
bool survives = false;
|
||||
bool counted_ok = count_extras_fd(rootfd, "", manifest, max_delete, &count, &exceeds, skips,
|
||||
bool counted_ok = count_extras_fd(rootfd, "", &keep, max_delete, &count, &exceeds, skips,
|
||||
skip_count, &survives);
|
||||
if (!counted_ok) {
|
||||
close(rootfd);
|
||||
path_index_free(&keep);
|
||||
return DELETE_WALK_ERROR;
|
||||
}
|
||||
if (exceeds) {
|
||||
close(rootfd);
|
||||
path_index_free(&keep);
|
||||
return DELETE_WALK_LIMIT_EXCEEDED;
|
||||
}
|
||||
}
|
||||
size_t deleted_count = 0;
|
||||
bool ok = delete_extras_fd(rootfd, "", manifest, max_delete, &deleted_count, skips, skip_count);
|
||||
bool ok = delete_extras_fd(rootfd, "", &keep, max_delete, &deleted_count, skips, skip_count);
|
||||
if (close(rootfd) != 0)
|
||||
ok = false;
|
||||
path_index_free(&keep);
|
||||
if (deleted_out)
|
||||
*deleted_out = deleted_count;
|
||||
return ok ? DELETE_WALK_OK : DELETE_WALK_ERROR;
|
||||
}
|
||||
|
||||
bool delete_extras(const char* dest_root, ArrayList* manifest) {
|
||||
bool delete_extras(const char* dest_root, const ArrayList* manifest) {
|
||||
return delete_extras_limited(dest_root, manifest, SIZE_MAX, NULL, 0, NULL) == DELETE_WALK_OK;
|
||||
}
|
||||
|
||||
|
||||
+73
-2
@@ -6,6 +6,77 @@
|
||||
#include <stdbool.h>
|
||||
#include <sys/socket.h>
|
||||
|
||||
/* Small open-addressing string hash set used to turn quadratic membership
|
||||
* scans into O(path length) exact-match lookups (the --delete keep-set and the
|
||||
* --files-from allow-set). Keys are hashed with xxHash64 (seed 0); collisions
|
||||
* are resolved by linear probing over a power-of-two table that grows at 75%
|
||||
* load. Keys are always borrowed from the caller and must outlive the set; the
|
||||
* set never copies or owns keys, so indexing M entries costs O(M) memory. The
|
||||
* set is not thread-safe for mutation, but a fully built set supports
|
||||
* concurrent read-only lookups. */
|
||||
typedef struct {
|
||||
const char* key; /* NULL marks an empty slot */
|
||||
} StrHashSetSlot;
|
||||
|
||||
typedef struct {
|
||||
StrHashSetSlot* slots;
|
||||
size_t capacity; /* power of two, zero before init */
|
||||
size_t size;
|
||||
} StrHashSet;
|
||||
|
||||
/* Initialize an empty set sized for roughly `hint` entries. Returns false on
|
||||
* allocation failure. */
|
||||
bool str_hash_set_init(StrHashSet* set, size_t hint);
|
||||
void str_hash_set_free(StrHashSet* set);
|
||||
/* Insert a borrowed key (must outlive the set). A duplicate is ignored.
|
||||
* Returns false on allocation failure. */
|
||||
bool str_hash_set_insert_ref(StrHashSet* set, const char* key);
|
||||
/* Look up a NUL-terminated key / a key of `len` bytes. */
|
||||
bool str_hash_set_lookup(const StrHashSet* set, const char* key);
|
||||
bool str_hash_set_lookup_n(const StrHashSet* set, const char* key, size_t len);
|
||||
|
||||
/* Sorted, non-owning view of NUL-terminated strings. Built from borrowed
|
||||
* pointers (qsort), so indexing M entries costs O(M) memory and O(M log M)
|
||||
* time; exact membership and ancestor-prefix existence are binary searches
|
||||
* that never materialize a prefix copy. */
|
||||
typedef struct {
|
||||
const char** items; /* sorted with strcmp; borrowed, never freed */
|
||||
size_t count;
|
||||
} StrSortedArray;
|
||||
|
||||
/* Build `array` over the borrowed `items`. Only the pointer array is copied,
|
||||
* never the strings. Returns false on allocation failure. */
|
||||
bool str_sorted_array_build(StrSortedArray* array, const char* const* items, size_t count);
|
||||
void str_sorted_array_free(StrSortedArray* array);
|
||||
/* True when some item equals `key`. */
|
||||
bool str_sorted_array_contains(const StrSortedArray* array, const char* key);
|
||||
/* True when some item starts with `key` followed by '/' (i.e. `key` is a proper
|
||||
* ancestor directory of an item). Allocates nothing. */
|
||||
bool str_sorted_array_has_child_prefix(const StrSortedArray* array, const char* key);
|
||||
|
||||
/* Read-only membership index over exact relative paths. `exact` answers
|
||||
* O(path length) equality; `sorted` answers whether any indexed path lies
|
||||
* strictly below a query directory. Both borrow their keys from the caller and
|
||||
* no ancestor prefix is stored as a separate string, so an index over M entries
|
||||
* is O(M) memory regardless of path depth. Not thread-safe to build, but safe
|
||||
* for concurrent read-only queries once built. */
|
||||
typedef struct {
|
||||
StrHashSet exact;
|
||||
StrSortedArray sorted;
|
||||
} PathIndex;
|
||||
|
||||
/* Build an index borrowing `entries` (which must outlive the index). Returns
|
||||
* false on allocation failure, freeing any partial state. */
|
||||
bool path_index_build(PathIndex* index, const char* const* entries, size_t count);
|
||||
void path_index_free(PathIndex* index);
|
||||
/* True when `path` is an indexed entry. */
|
||||
bool path_index_contains(const PathIndex* index, const char* path);
|
||||
/* Length-bounded form of path_index_contains (`path` need not be terminated). */
|
||||
bool path_index_contains_n(const PathIndex* index, const char* path, size_t len);
|
||||
/* True when some indexed entry lies strictly below `path` (starts with
|
||||
* `path` + '/'). */
|
||||
bool path_index_has_descendant(const PathIndex* index, const char* path);
|
||||
|
||||
char* str_dup(const char* string);
|
||||
char* output_escape(const char* string, bool eight_bit_output);
|
||||
char* path_cat(const char* path1, const char* path2);
|
||||
@@ -48,10 +119,10 @@ bool path_under_skip_prefix(const char* child_rel, bool at_root, const DeleteSki
|
||||
and the delete pass are two separate walks, so a concurrent change between
|
||||
them (another process adding/removing entries) can make the second pass
|
||||
delete a different set than the first one counted. */
|
||||
DeleteWalkResult delete_extras_limited(const char* dest_root, ArrayList* manifest,
|
||||
DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest,
|
||||
size_t max_delete, const DeleteSkipEntry* skips,
|
||||
int skip_count, size_t* deleted_out);
|
||||
bool delete_extras(const char* dest_root, ArrayList* manifest);
|
||||
bool delete_extras(const char* dest_root, const ArrayList* manifest);
|
||||
bool utils_set_authorized_root(int fd, const char* canonical_path);
|
||||
/* The fd-only compatibility form is fail-closed for path-based operations;
|
||||
* callers should use utils_set_authorized_root with the canonical identity. */
|
||||
|
||||
@@ -94,14 +94,14 @@ def _seed_protocol_source(source):
|
||||
class TestProtocol:
|
||||
@pytest.mark.ci
|
||||
def test_protocol_current_version_accepted(self, shared_server):
|
||||
"""--protocol=2.19.0 (the current PROTOCOL_VERSION) is accepted and the
|
||||
"""--protocol=2.20.0 (the current PROTOCOL_VERSION) is accepted and the
|
||||
transfer completes normally."""
|
||||
source = os.path.join(TEST_DATA_DIR, "proto_ok_src")
|
||||
dest = os.path.join(TEST_DATA_DIR, "proto_ok_dst")
|
||||
shutil.rmtree(dest, ignore_errors=True)
|
||||
os.makedirs(dest)
|
||||
_seed_protocol_source(source)
|
||||
result, _ = run_client(source, dest, flags=["--protocol=2.19.0"],
|
||||
result, _ = run_client(source, dest, flags=["--protocol=2.20.0"],
|
||||
port=shared_server.port)
|
||||
assert result.returncode == 0, \
|
||||
f"--protocol current run failed: {(result.stderr or result.stdout)[:400]}"
|
||||
@@ -118,7 +118,7 @@ class TestProtocol:
|
||||
shutil.rmtree(dest, ignore_errors=True)
|
||||
os.makedirs(dest)
|
||||
_seed_protocol_source(source)
|
||||
for bad in ("2.18.0", "2.17.0", "2.15.0", "2.16.0", "216", "31"):
|
||||
for bad in ("2.19.0", "2.18.0", "2.17.0", "2.15.0", "2.16.0", "216", "31"):
|
||||
result, _ = run_client(source, dest, flags=[f"--protocol={bad}"],
|
||||
port=shared_server.port)
|
||||
assert result.returncode != 0, f"--protocol={bad} should be rejected"
|
||||
|
||||
@@ -12,6 +12,7 @@
|
||||
#include "test_delay_updates.h"
|
||||
#include "test_delta.h"
|
||||
#include "test_file.h"
|
||||
#include "test_file_list.h"
|
||||
#include "test_file_sendfile.h"
|
||||
#include "test_fuzz_smoke.h"
|
||||
#include "test_glob.h"
|
||||
@@ -67,6 +68,7 @@ int main() {
|
||||
RUN_TEST(test_glob);
|
||||
RUN_TEST(test_iconv);
|
||||
RUN_TEST(test_file);
|
||||
RUN_TEST(test_file_list);
|
||||
RUN_TEST(test_trust_sender);
|
||||
RUN_TEST(test_delay_updates);
|
||||
RUN_TEST(test_file_sendfile);
|
||||
|
||||
@@ -259,7 +259,7 @@ static void test_parse_args_protocol_accept_current() {
|
||||
Config* cfg = valid_client_config();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
char* argv_equals[] = {"fastsync", "--source-dir", "/src",
|
||||
"--dest-dir", "/dst", "--protocol=2.19.0"};
|
||||
"--dest-dir", "/dst", "--protocol=2.20.0"};
|
||||
int positional_args[2];
|
||||
int positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 6, argv_equals, positional_args, &positional_count), 0);
|
||||
@@ -269,7 +269,7 @@ static void test_parse_args_protocol_accept_current() {
|
||||
cfg = valid_client_config();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
char* argv_space[] = {"fastsync", "--source-dir", "/src", "--dest-dir",
|
||||
"/dst", "--protocol", "2.19.0"};
|
||||
"/dst", "--protocol", "2.20.0"};
|
||||
positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 7, argv_space, positional_args, &positional_count), 0);
|
||||
EXPECT_EQ_STR(cfg->version, PROTOCOL_VERSION);
|
||||
@@ -279,8 +279,8 @@ static void test_parse_args_protocol_accept_current() {
|
||||
/* Any --protocol value other than the current PROTOCOL_VERSION must end in
|
||||
* failure (parse_args simply stores it; validate_config rejects it up front). */
|
||||
static void test_parse_args_protocol_rejects_other_versions() {
|
||||
static const char* const bad_versions[] = {"2.17", "2.16", "2.15.0", "2.16.0", "2.17.0",
|
||||
"2.18.0", "216", "31", "abc", ""};
|
||||
static const char* const bad_versions[] = {
|
||||
"2.17", "2.16", "2.15.0", "2.16.0", "2.17.0", "2.18.0", "2.19.0", "216", "31", "abc", ""};
|
||||
for (size_t i = 0; i < sizeof(bad_versions) / sizeof(bad_versions[0]); i++) {
|
||||
Config* cfg = valid_client_config();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
#include "utils.h"
|
||||
#include <string.h>
|
||||
#include <sys/stat.h>
|
||||
#include <threads.h>
|
||||
#include <unistd.h>
|
||||
|
||||
static void test_data_compress_decompress_roundtrip() {
|
||||
@@ -137,10 +138,84 @@ static void test_chunk_compress_decompress_roundtrip() {
|
||||
unlink(path2);
|
||||
}
|
||||
|
||||
typedef struct {
|
||||
int id;
|
||||
int iterations;
|
||||
bool ok;
|
||||
} CompressionThreadArg;
|
||||
|
||||
/* Each worker exercises the per-thread cached zstd contexts: several
|
||||
* compress/decompress round-trips with varying payload sizes, levels and
|
||||
* worker counts so the context is reused (and its parameters re-applied)
|
||||
* across calls, concurrently with other workers. */
|
||||
static int compression_reuse_worker(void* arg) {
|
||||
CompressionThreadArg* a = (CompressionThreadArg*)arg;
|
||||
a->ok = true;
|
||||
for (int it = 0; it < a->iterations; it++) {
|
||||
size_t size = 512 + (size_t)((a->id * 7919 + it * 104729) % (48 * 1024));
|
||||
char* original = malloc(size);
|
||||
if (!original) {
|
||||
a->ok = false;
|
||||
break;
|
||||
}
|
||||
for (size_t i = 0; i < size; i++)
|
||||
original[i] = (char)((i * 31 + (size_t)a->id + (size_t)it * 7) % 251);
|
||||
Data* input = data_create(original, size);
|
||||
if (!input) { /* data_create takes ownership of original, even on failure */
|
||||
a->ok = false;
|
||||
break;
|
||||
}
|
||||
int level = 1 + ((it / 2) % 5);
|
||||
int threads = ((it / 2) % 2 == 0) ? 2 : 0;
|
||||
Data* compressed = data_compress_with_threads(input, level, threads);
|
||||
if (!compressed) {
|
||||
data_destroy(input);
|
||||
a->ok = false;
|
||||
break;
|
||||
}
|
||||
Data* decompressed = data_decompress(compressed);
|
||||
bool roundtrip_ok = decompressed != NULL && decompressed->size == size &&
|
||||
memcmp(decompressed->data, original, size) == 0;
|
||||
data_destroy(decompressed);
|
||||
data_destroy(compressed);
|
||||
data_destroy(input);
|
||||
if (!roundtrip_ok) {
|
||||
a->ok = false;
|
||||
break;
|
||||
}
|
||||
}
|
||||
/* Deliberately do NOT free the thread context here: the C11 tss destructor
|
||||
* must release it when this thread exits (validated by LeakSanitizer). */
|
||||
return thrd_success;
|
||||
}
|
||||
|
||||
static void test_data_compress_reused_contexts_multithreaded() {
|
||||
enum { NTHREADS = 8, ITERATIONS = 6 };
|
||||
thrd_t threads[NTHREADS];
|
||||
CompressionThreadArg args[NTHREADS];
|
||||
bool all_created = true;
|
||||
for (int i = 0; i < NTHREADS; i++) {
|
||||
args[i].id = i;
|
||||
args[i].iterations = ITERATIONS;
|
||||
args[i].ok = false;
|
||||
if (thrd_create(&threads[i], compression_reuse_worker, &args[i]) != thrd_success) {
|
||||
all_created = false;
|
||||
break;
|
||||
}
|
||||
}
|
||||
EXPECT_TRUE(all_created);
|
||||
for (int i = 0; i < NTHREADS; i++)
|
||||
EXPECT_EQ_INT(thrd_join(threads[i], NULL), thrd_success);
|
||||
for (int i = 0; i < NTHREADS; i++)
|
||||
EXPECT_TRUE(args[i].ok);
|
||||
compression_free_thread_contexts();
|
||||
}
|
||||
|
||||
void test_compression() {
|
||||
test_data_compress_decompress_roundtrip();
|
||||
test_data_compress_decompress_large();
|
||||
test_skip_compress_suffix_matching();
|
||||
test_data_compress_with_threads_roundtrip();
|
||||
test_data_compress_reused_contexts_multithreaded();
|
||||
test_chunk_compress_decompress_roundtrip();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,196 @@
|
||||
#include "test_file_list.h"
|
||||
#include "file_list.h"
|
||||
#include "test_utils.h"
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
/* Reference implementation of the ORIGINAL file_list_affects linear scan. The
|
||||
indexed implementation must agree with it on every query; this pins the
|
||||
subtle semantics: empty entry == whole tree, exact match, rel under a listed
|
||||
directory, and rel an ancestor directory of a listed entry. */
|
||||
static bool reference_affects(const FileListSet* set, const char* rel) {
|
||||
if (!set)
|
||||
return true;
|
||||
if (!rel)
|
||||
return false;
|
||||
for (int i = 0; i < set->count; i++) {
|
||||
const char* entry = set->entries[i];
|
||||
if (entry[0] == '\0')
|
||||
return true;
|
||||
if (strcmp(rel, entry) == 0)
|
||||
return true;
|
||||
size_t entry_len = strlen(entry);
|
||||
if (strncmp(rel, entry, entry_len) == 0 && (rel[entry_len] == '/' || rel[entry_len] == '\0'))
|
||||
return true;
|
||||
size_t rel_len = strlen(rel);
|
||||
if (strncmp(entry, rel, rel_len) == 0 && (entry[rel_len] == '/' || entry[rel_len] == '\0'))
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
static void write_list(const char* path, const char* bytes) {
|
||||
FILE* fp = fopen(path, "wb");
|
||||
EXPECT_NOT_NULL(fp);
|
||||
size_t len = strlen(bytes);
|
||||
EXPECT_EQ_INT((int)fwrite(bytes, 1, len, fp), (int)len);
|
||||
fclose(fp);
|
||||
}
|
||||
|
||||
static void check_queries(const FileListSet* set, const char* const* queries, int query_count) {
|
||||
for (int i = 0; i < query_count; i++) {
|
||||
bool expected = reference_affects(set, queries[i]);
|
||||
bool actual = file_list_affects(set, queries[i]);
|
||||
if (expected != actual) {
|
||||
printf(" [FAIL] affects(\"%s\"): expected %d, got %d\n", queries[i], expected, actual);
|
||||
current_test_failed = true;
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
static void test_membership_matches_reference() {
|
||||
const char* path = "test_file_list_case.txt";
|
||||
char err[160];
|
||||
|
||||
/* Nested directories, an ancestor of a listed entry, an exact file, a
|
||||
non-matching neighbor with the same prefix, and a literal '*'. */
|
||||
write_list(path, "a\na/b\na/b/c\nab\nc.txt\nsub/b.bin\n*\n");
|
||||
FileListSet* set = file_list_load(path, false, err, sizeof(err));
|
||||
EXPECT_NOT_NULL(set);
|
||||
const char* queries[] = {
|
||||
"a", "a/b", "a/b/c", "a/b/c/d", "a/bx", "a/x", "ab",
|
||||
"abc", "c.txt", "c.txt/x", "c", "sub", "sub/b.bin", "sub/b.bin/z",
|
||||
"sub2", "*", "x", "", "a/b/cd", "/", "a/b/",
|
||||
};
|
||||
check_queries(set, queries, (int)(sizeof(queries) / sizeof(queries[0])));
|
||||
EXPECT_TRUE(file_list_affects(set, "a/b/c/d"));
|
||||
EXPECT_TRUE(file_list_affects(set, "a/bx")); /* under listed directory "a" */
|
||||
EXPECT_TRUE(file_list_affects(set, "a/b/c"));
|
||||
EXPECT_TRUE(file_list_affects(set, "a/x")); /* under listed directory "a" */
|
||||
EXPECT_FALSE(file_list_affects(set, "abc")); /* component boundary: not "a" */
|
||||
EXPECT_TRUE(file_list_affects(set, "sub"));
|
||||
EXPECT_FALSE(file_list_affects(set, "sub2"));
|
||||
file_list_destroy(set);
|
||||
remove(path);
|
||||
|
||||
/* A single "." entry means the whole tree: every non-NULL query is true. */
|
||||
write_list(path, ".\n");
|
||||
set = file_list_load(path, false, err, sizeof(err));
|
||||
EXPECT_NOT_NULL(set);
|
||||
const char* root_queries[] = {"", "a", "a/b/c", "unrelated", "*", "/"};
|
||||
for (int i = 0; i < (int)(sizeof(root_queries) / sizeof(root_queries[0])); i++)
|
||||
EXPECT_TRUE(file_list_affects(set, root_queries[i]));
|
||||
EXPECT_FALSE(file_list_affects(set, NULL));
|
||||
check_queries(set, root_queries, (int)(sizeof(root_queries) / sizeof(root_queries[0])));
|
||||
file_list_destroy(set);
|
||||
remove(path);
|
||||
|
||||
/* Trailing slashes and "./" prefixes normalize away, so the query matches the
|
||||
clean path (and not the raw spelling). */
|
||||
write_list(path, "./dir/\ndir2/./x\n");
|
||||
set = file_list_load(path, false, err, sizeof(err));
|
||||
EXPECT_NOT_NULL(set);
|
||||
EXPECT_TRUE(file_list_affects(set, "dir"));
|
||||
EXPECT_TRUE(file_list_affects(set, "dir/x"));
|
||||
EXPECT_TRUE(file_list_affects(set, "dir2/x"));
|
||||
EXPECT_TRUE(file_list_affects(set, "dir2"));
|
||||
EXPECT_TRUE(file_list_affects(set, "dir/")); /* boundary prefix of listed "dir" */
|
||||
check_queries(set, (const char*[]){"dir", "dir/", "dir/x", "dir2", "dir2/x", "dir3"}, 6);
|
||||
file_list_destroy(set);
|
||||
remove(path);
|
||||
|
||||
/* An empty file yields an empty set: nothing is affected, and NULL set still
|
||||
means "everything". */
|
||||
write_list(path, "");
|
||||
set = file_list_load(path, false, err, sizeof(err));
|
||||
EXPECT_NOT_NULL(set);
|
||||
EXPECT_EQ_INT(set->count, 0);
|
||||
EXPECT_FALSE(file_list_affects(set, "a"));
|
||||
EXPECT_FALSE(file_list_affects(set, ""));
|
||||
check_queries(set, (const char*[]){"a", "a/b", ""}, 3);
|
||||
file_list_destroy(set);
|
||||
remove(path);
|
||||
|
||||
/* NULL set is the unrestricted case. */
|
||||
EXPECT_TRUE(file_list_affects(NULL, "anything"));
|
||||
EXPECT_TRUE(file_list_affects(NULL, NULL));
|
||||
}
|
||||
|
||||
/* Explicit ancestor/descendant coverage: a query that is a proper ancestor of
|
||||
a listed entry is affected, and a query below a listed entry is affected,
|
||||
while a component-boundary neighbor is not. */
|
||||
static void test_ancestor_and_descendant_queries() {
|
||||
const char* path = "test_file_list_ancestor.txt";
|
||||
char err[160];
|
||||
write_list(path, "top/mid/leaf.txt\nsingle.txt\n");
|
||||
FileListSet* set = file_list_load(path, false, err, sizeof(err));
|
||||
EXPECT_NOT_NULL(set);
|
||||
|
||||
/* q is an ancestor of a listed entry. */
|
||||
EXPECT_TRUE(file_list_affects(set, "top"));
|
||||
EXPECT_TRUE(file_list_affects(set, "top/mid"));
|
||||
EXPECT_FALSE(file_list_affects(set, "top/other")); /* neither direction */
|
||||
EXPECT_FALSE(file_list_affects(set, "to")); /* component boundary */
|
||||
|
||||
/* A listed entry is an ancestor of q. */
|
||||
EXPECT_TRUE(file_list_affects(set, "single.txt"));
|
||||
EXPECT_TRUE(file_list_affects(set, "single.txt/deeper"));
|
||||
EXPECT_FALSE(file_list_affects(set, "single.txtx")); /* boundary */
|
||||
|
||||
check_queries(set,
|
||||
(const char*[]){"top", "top/mid", "top/mid/leaf.txt", "top/other", "single.txt",
|
||||
"single.txt/deeper", "single.txtx", "to"},
|
||||
8);
|
||||
file_list_destroy(set);
|
||||
remove(path);
|
||||
}
|
||||
|
||||
/* Regression for the remote OOM: an adversarial --files-from entry made of a
|
||||
very deep chain of repeated components must be indexed with memory
|
||||
proportional to the entry count. The old implementation stored one copied
|
||||
ancestor prefix per component (O(L^2) bytes for a single entry); the sorted
|
||||
index stores the exact entries only. */
|
||||
static void test_deep_paths_are_bounded() {
|
||||
const char* path = "test_file_list_deep.txt";
|
||||
enum { COMPONENTS = 20000 };
|
||||
size_t entry_len = (size_t)COMPONENTS * 2; /* "a/" per component */
|
||||
char* entry = malloc(entry_len + 1);
|
||||
EXPECT_NOT_NULL(entry);
|
||||
for (size_t i = 0; i < entry_len; i += 2) {
|
||||
entry[i] = 'a';
|
||||
entry[i + 1] = '/';
|
||||
}
|
||||
entry[entry_len - 1] = 'z'; /* .../a/z: a deep leaf name */
|
||||
entry[entry_len] = '\0';
|
||||
|
||||
FILE* fp = fopen(path, "wb");
|
||||
EXPECT_NOT_NULL(fp);
|
||||
EXPECT_EQ_INT((int)fwrite(entry, 1, entry_len, fp), (int)entry_len);
|
||||
EXPECT_EQ_INT(fputc('\n', fp), '\n');
|
||||
fclose(fp);
|
||||
|
||||
char err[160];
|
||||
FileListSet* set = file_list_load(path, false, err, sizeof(err));
|
||||
EXPECT_NOT_NULL(set);
|
||||
EXPECT_EQ_INT(set->count, 1);
|
||||
/* One exact entry stored, not one node per path component. */
|
||||
EXPECT_EQ_INT((int)set->index.sorted.count, 1);
|
||||
EXPECT_EQ_INT((int)set->index.exact.size, 1);
|
||||
EXPECT_TRUE(file_list_affects(set, entry)); /* exact */
|
||||
EXPECT_TRUE(file_list_affects(set, "a")); /* ancestor of the entry */
|
||||
EXPECT_TRUE(file_list_affects(set, "a/a")); /* deeper ancestor */
|
||||
EXPECT_FALSE(file_list_affects(set, "b")); /* unrelated */
|
||||
EXPECT_FALSE(file_list_affects(set, "aa")); /* component boundary */
|
||||
|
||||
file_list_destroy(set);
|
||||
remove(path);
|
||||
free(entry);
|
||||
}
|
||||
|
||||
void test_file_list() {
|
||||
test_membership_matches_reference();
|
||||
test_ancestor_and_descendant_queries();
|
||||
test_deep_paths_are_bounded();
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
#ifndef TEST_FILE_LIST_H
|
||||
#define TEST_FILE_LIST_H
|
||||
|
||||
void test_file_list();
|
||||
|
||||
#endif
|
||||
+139
-31
@@ -5,6 +5,8 @@
|
||||
#include "test_utils.h"
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/ioctl.h>
|
||||
#include <sys/socket.h>
|
||||
#include <sys/stat.h>
|
||||
#include <unistd.h>
|
||||
|
||||
@@ -134,6 +136,115 @@ static void test_metadata_send_null() {
|
||||
close(p[1]);
|
||||
}
|
||||
|
||||
/* protocol 2.20.0: metadata is one packed frame. With metadata present the
|
||||
* wire record is exactly sizeof(int32_t) + FILE_METADATA_WIRE_SIZE bytes (the
|
||||
* present flag followed by the fixed field record); absent metadata is a lone
|
||||
* int32 zero. */
|
||||
static void test_metadata_wire_is_one_packed_frame() {
|
||||
io_set_bwlimit(0);
|
||||
int p[2];
|
||||
EXPECT_EQ_INT(pipe(p), 0);
|
||||
io_set_fds(p[0], p[1]);
|
||||
|
||||
FileMetadata original = {.mode = 0640,
|
||||
.uid = 42,
|
||||
.gid = 43,
|
||||
.mtime_sec = 111,
|
||||
.mtime_nsec = 222,
|
||||
.atime_valid = true,
|
||||
.atime_sec = 333,
|
||||
.atime_nsec = 444,
|
||||
.crtime_valid = false,
|
||||
.crtime_sec = 0,
|
||||
.crtime_nsec = 0};
|
||||
EXPECT_TRUE(metadata_send(p[1], &original));
|
||||
|
||||
unsigned char wire[sizeof(int32_t) + FILE_METADATA_WIRE_SIZE];
|
||||
EXPECT_EQ_INT((int)read(p[0], wire, sizeof(wire)), (int)sizeof(wire));
|
||||
int32_t flag;
|
||||
memcpy(&flag, wire, sizeof(flag));
|
||||
EXPECT_EQ_INT(flag, 1);
|
||||
int avail = -1;
|
||||
EXPECT_EQ_INT(ioctl(p[0], FIONREAD, &avail), 0);
|
||||
EXPECT_EQ_INT(avail, 0);
|
||||
|
||||
/* The present frame decodes in one shot with the shared codec. */
|
||||
char* cursor = (char*)wire;
|
||||
FileMetadata* decoded = metadata_from_buf(&cursor);
|
||||
EXPECT_NOT_NULL(decoded);
|
||||
EXPECT_EQ_INT((int)(cursor - (char*)wire), (int)sizeof(wire));
|
||||
EXPECT_EQ_INT(decoded->mode, 0640);
|
||||
EXPECT_EQ_INT(decoded->uid, 42);
|
||||
EXPECT_EQ_INT(decoded->gid, 43);
|
||||
EXPECT_EQ_INT(decoded->mtime_sec, 111);
|
||||
EXPECT_EQ_INT(decoded->mtime_nsec, 222);
|
||||
EXPECT_TRUE(decoded->atime_valid);
|
||||
EXPECT_EQ_INT(decoded->atime_sec, 333);
|
||||
EXPECT_EQ_INT(decoded->atime_nsec, 444);
|
||||
EXPECT_FALSE(decoded->crtime_valid);
|
||||
free(decoded);
|
||||
|
||||
/* Absent metadata is a lone int32 zero (4 bytes). */
|
||||
EXPECT_TRUE(metadata_send(p[1], NULL));
|
||||
EXPECT_EQ_INT((int)read(p[0], wire, sizeof(int32_t)), (int)sizeof(int32_t));
|
||||
memcpy(&flag, wire, sizeof(flag));
|
||||
EXPECT_EQ_INT(flag, 0);
|
||||
avail = -1;
|
||||
EXPECT_EQ_INT(ioctl(p[0], FIONREAD, &avail), 0);
|
||||
EXPECT_EQ_INT(avail, 0);
|
||||
|
||||
close(p[0]);
|
||||
close(p[1]);
|
||||
}
|
||||
|
||||
/* Round-trip over a socketpair (not just a pipe): present metadata compares
|
||||
* equal field-by-field and absent metadata yields NULL with ok == 1. */
|
||||
static void test_metadata_send_receive_socketpair() {
|
||||
io_set_bwlimit(0);
|
||||
int sv[2];
|
||||
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, sv), 0);
|
||||
io_set_fds(sv[0], sv[1]);
|
||||
|
||||
FileMetadata original = {.mode = 0600,
|
||||
.uid = 7,
|
||||
.gid = 8,
|
||||
.mtime_sec = 1000,
|
||||
.mtime_nsec = 1,
|
||||
.atime_valid = true,
|
||||
.atime_sec = 2000,
|
||||
.atime_nsec = 2,
|
||||
.crtime_valid = true,
|
||||
.crtime_sec = 3000,
|
||||
.crtime_nsec = 3};
|
||||
EXPECT_TRUE(metadata_send(sv[1], &original));
|
||||
|
||||
int ok = 0;
|
||||
FileMetadata* received = metadata_receive(sv[0], &ok);
|
||||
EXPECT_NOT_NULL(received);
|
||||
EXPECT_EQ_INT(ok, 1);
|
||||
EXPECT_EQ_INT(received->mode, 0600);
|
||||
EXPECT_EQ_INT(received->uid, 7);
|
||||
EXPECT_EQ_INT(received->gid, 8);
|
||||
EXPECT_EQ_INT(received->mtime_sec, 1000);
|
||||
EXPECT_EQ_INT(received->mtime_nsec, 1);
|
||||
EXPECT_TRUE(received->atime_valid);
|
||||
EXPECT_EQ_INT(received->atime_sec, 2000);
|
||||
EXPECT_EQ_INT(received->atime_nsec, 2);
|
||||
EXPECT_TRUE(received->crtime_valid);
|
||||
EXPECT_EQ_INT(received->crtime_sec, 3000);
|
||||
EXPECT_EQ_INT(received->crtime_nsec, 3);
|
||||
free(received);
|
||||
|
||||
EXPECT_TRUE(metadata_send(sv[1], NULL));
|
||||
ok = 0;
|
||||
received = metadata_receive(sv[0], &ok);
|
||||
EXPECT_NULL(received);
|
||||
EXPECT_EQ_INT(ok, 1);
|
||||
|
||||
close(sv[0]);
|
||||
close(sv[1]);
|
||||
}
|
||||
|
||||
static void test_metadata_rejects_invalid_values() {
|
||||
int p[2];
|
||||
EXPECT_EQ_INT(pipe(p), 0);
|
||||
@@ -148,36 +259,30 @@ static void test_metadata_rejects_invalid_values() {
|
||||
}
|
||||
|
||||
/* metadata_receive must reject an out-of-range atime/crtime nsec even when the
|
||||
* flag would otherwise be valid (defense-in-depth on the -U/-N wire fields). */
|
||||
* flag would otherwise be valid (defense-in-depth on the -U/-N wire fields).
|
||||
* The packed record is built by the shared codec so the out-of-range value
|
||||
* actually reaches the wire. */
|
||||
static void test_metadata_receive_rejects_bad_optional_times() {
|
||||
int p[2];
|
||||
EXPECT_EQ_INT(pipe(p), 0);
|
||||
io_set_fds(p[0], p[1]);
|
||||
|
||||
int32_t present = 1;
|
||||
int32_t mode = 0644;
|
||||
int32_t uid = 1000;
|
||||
int32_t gid = 1000;
|
||||
int64_t mtime_sec = 1;
|
||||
int64_t mtime_nsec = 0;
|
||||
int32_t atime_valid = 1;
|
||||
int64_t atime_sec = 1;
|
||||
int64_t atime_nsec = 2000000000; /* invalid: >= 1e9 */
|
||||
EXPECT_TRUE(send_n_data(p[1], &present, sizeof(present)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &mode, sizeof(mode)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &uid, sizeof(uid)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &gid, sizeof(gid)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &mtime_sec, sizeof(mtime_sec)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &atime_valid, sizeof(atime_valid)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &atime_sec, sizeof(atime_sec)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &atime_nsec, sizeof(atime_nsec)));
|
||||
int32_t crtime_valid = 0;
|
||||
int64_t crtime_sec = 0;
|
||||
int64_t crtime_nsec = 0;
|
||||
EXPECT_TRUE(send_n_data(p[1], &crtime_valid, sizeof(crtime_valid)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &crtime_sec, sizeof(crtime_sec)));
|
||||
EXPECT_TRUE(send_n_data(p[1], &crtime_nsec, sizeof(crtime_nsec)));
|
||||
FileMetadata bad = {.mode = 0644,
|
||||
.uid = 1000,
|
||||
.gid = 1000,
|
||||
.mtime_sec = 1,
|
||||
.mtime_nsec = 0,
|
||||
.atime_valid = true,
|
||||
.atime_sec = 1,
|
||||
.atime_nsec = 2000000000, /* invalid: >= 1e9 */
|
||||
.crtime_valid = false,
|
||||
.crtime_sec = 0,
|
||||
.crtime_nsec = 0};
|
||||
char packed[sizeof(int32_t) + FILE_METADATA_WIRE_SIZE];
|
||||
char* write_ptr = packed;
|
||||
metadata_to_buf(&write_ptr, &bad);
|
||||
EXPECT_TRUE(send_n_data(p[1], packed, sizeof(packed)));
|
||||
|
||||
int ok = 1;
|
||||
EXPECT_NULL(metadata_receive(p[0], &ok));
|
||||
EXPECT_EQ_INT(ok, 0);
|
||||
@@ -235,12 +340,13 @@ static void test_file_restore_metadata() {
|
||||
const char* content = "test content";
|
||||
EXPECT_TRUE(file_write_to_disk(path, content, strlen(content), false, false));
|
||||
|
||||
FileMetadata m;
|
||||
m.mode = 0644;
|
||||
m.uid = getuid();
|
||||
m.gid = getgid();
|
||||
m.mtime_sec = 1234567890;
|
||||
m.mtime_nsec = 0;
|
||||
FileMetadata m = {.mode = 0644,
|
||||
.uid = getuid(),
|
||||
.gid = getgid(),
|
||||
.mtime_sec = 1234567890,
|
||||
.mtime_nsec = 0,
|
||||
.atime_valid = false,
|
||||
.crtime_valid = false};
|
||||
|
||||
file_restore_metadata(path, &m, false);
|
||||
|
||||
@@ -347,6 +453,8 @@ void test_metadata() {
|
||||
test_metadata_from_buf_null();
|
||||
test_metadata_send_receive_roundtrip();
|
||||
test_metadata_send_null();
|
||||
test_metadata_wire_is_one_packed_frame();
|
||||
test_metadata_send_receive_socketpair();
|
||||
test_metadata_rejects_invalid_values();
|
||||
test_metadata_receive_rejects_bad_optional_times();
|
||||
test_metadata_mtime_window();
|
||||
|
||||
@@ -332,6 +332,123 @@ static void test_receiver_enqueue_byte_budget() {
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
/* Build a one-file chunk that appears to hold `bytes` of loaded payload by
|
||||
handing it a real buffer of that size (the byte accounting counts only
|
||||
in-memory `data->data`, mirroring sendfile's streamed chunks). */
|
||||
static Chunk* make_loaded_chunk(const char* name, size_t bytes) {
|
||||
File* file = file_create(name);
|
||||
if (!file)
|
||||
return NULL;
|
||||
void* buffer = malloc(bytes > 0 ? bytes : 1);
|
||||
if (!buffer) {
|
||||
file_destroy(file);
|
||||
return NULL;
|
||||
}
|
||||
file->data->data = buffer;
|
||||
file->data->size = bytes;
|
||||
File* items[1] = {file};
|
||||
Chunk* chunk = chunk_create(items, 1);
|
||||
if (!chunk)
|
||||
file_destroy(file);
|
||||
return chunk;
|
||||
}
|
||||
|
||||
/* Only loaded (in-memory) payload is charged: a chunk whose files have no
|
||||
buffer (e.g. sendfile streams the bytes from disk) accounts for zero. */
|
||||
static void test_sender_chunk_bytes_accounting() {
|
||||
EXPECT_EQ_INT((int)pipeline_context_sender_chunk_bytes(NULL), 0);
|
||||
|
||||
File* streamed = file_create("sender_account_streamed");
|
||||
EXPECT_NOT_NULL(streamed);
|
||||
streamed->data->size = 4096; /* declared size, but no in-memory buffer */
|
||||
File* streamed_items[1] = {streamed};
|
||||
Chunk* streamed_chunk = chunk_create(streamed_items, 1);
|
||||
EXPECT_NOT_NULL(streamed_chunk);
|
||||
EXPECT_EQ_INT((int)pipeline_context_sender_chunk_bytes(streamed_chunk), 0);
|
||||
chunk_destroy(streamed_chunk);
|
||||
|
||||
Chunk* loaded = make_loaded_chunk("sender_account_loaded", 2000);
|
||||
EXPECT_NOT_NULL(loaded);
|
||||
EXPECT_EQ_INT((int)pipeline_context_sender_chunk_bytes(loaded), 2000);
|
||||
chunk_destroy(loaded);
|
||||
}
|
||||
|
||||
typedef struct {
|
||||
PipelineContextSender* context;
|
||||
Chunk* chunk;
|
||||
atomic_bool* done;
|
||||
atomic_bool* result;
|
||||
} SenderByteBudgetArg;
|
||||
|
||||
static int sender_byte_budget_worker(void* arg) {
|
||||
SenderByteBudgetArg* worker = arg;
|
||||
bool ok = pipeline_context_sender_enqueue_chunk(worker->context, worker->chunk);
|
||||
atomic_store(worker->result, ok);
|
||||
atomic_store(worker->done, true);
|
||||
return thrd_success;
|
||||
}
|
||||
|
||||
/* The sender's loader stage must not buffer more loaded payload bytes ahead of
|
||||
the network writer than the configured byte budget: an enqueue that would
|
||||
exceed the budget blocks until the sender releases bytes. */
|
||||
static void test_sender_enqueue_byte_budget() {
|
||||
Config* cfg = config_create();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
free(cfg->version);
|
||||
cfg->version = str_dup(PROTOCOL_VERSION);
|
||||
cfg->send_directory = str_dup("/src");
|
||||
cfg->receive_root_directory = str_dup("/dst");
|
||||
|
||||
Queue* q_scanner = queue_create(16, chunk_destroy);
|
||||
Queue* q_loader = queue_create(16, chunk_destroy);
|
||||
EXPECT_NOT_NULL(q_scanner);
|
||||
EXPECT_NOT_NULL(q_loader);
|
||||
PipelineContextSender* ctx = pipeline_context_sender_create(cfg, q_scanner, q_loader);
|
||||
EXPECT_NOT_NULL(ctx);
|
||||
pipeline_context_sender_set_queue_byte_limit(ctx, 3000);
|
||||
|
||||
Chunk* first = make_loaded_chunk("sender_budget_1", 2000);
|
||||
EXPECT_NOT_NULL(first);
|
||||
EXPECT_TRUE(pipeline_context_sender_enqueue_chunk(ctx, first));
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000);
|
||||
|
||||
/* A second 2000-byte chunk would push the pipeline to 4000 > 3000 budget, so
|
||||
its enqueue must block until the first chunk's bytes are released. */
|
||||
Chunk* second = make_loaded_chunk("sender_budget_2", 2000);
|
||||
EXPECT_NOT_NULL(second);
|
||||
atomic_bool done;
|
||||
atomic_bool result;
|
||||
atomic_init(&done, false);
|
||||
atomic_init(&result, false);
|
||||
SenderByteBudgetArg arg = {ctx, second, &done, &result};
|
||||
thrd_t enqueuer;
|
||||
EXPECT_EQ_INT(thrd_create(&enqueuer, sender_byte_budget_worker, &arg), thrd_success);
|
||||
|
||||
/* Give a broken (unbounded) implementation every chance to enqueue. */
|
||||
struct timespec wait = {0, 200 * 1000000L};
|
||||
thrd_sleep(&wait, NULL);
|
||||
EXPECT_FALSE(atomic_load(&done));
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* budget still honored */
|
||||
|
||||
/* Simulate the sender: dequeue + destroy the first chunk, then release its
|
||||
bytes. Only the post-join state (below) is deterministic. */
|
||||
Chunk* drained =
|
||||
queue_dequeue_multithreaded(q_loader, &ctx->mutex_loader, &ctx->condition_not_empty_loader,
|
||||
&ctx->condition_not_full_loader, &ctx->loader_done);
|
||||
EXPECT_NOT_NULL(drained);
|
||||
chunk_destroy(drained);
|
||||
pipeline_context_sender_note_bytes_released(ctx, 2000);
|
||||
|
||||
EXPECT_EQ_INT(thrd_join(enqueuer, NULL), thrd_success);
|
||||
EXPECT_TRUE(atomic_load(&done));
|
||||
EXPECT_TRUE(atomic_load(&result));
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* second payload now in flight */
|
||||
|
||||
/* pipeline_context_sender_destroy frees the still-queued second chunk and
|
||||
owns cfg/q_scanner/q_loader from here on. */
|
||||
pipeline_context_sender_destroy(ctx);
|
||||
}
|
||||
|
||||
void test_multiprocessing() {
|
||||
test_sender_create_destroy();
|
||||
test_receiver_create_destroy();
|
||||
@@ -344,4 +461,6 @@ void test_multiprocessing() {
|
||||
}
|
||||
test_write_thread_done();
|
||||
test_receiver_enqueue_byte_budget();
|
||||
test_sender_chunk_bytes_accounting();
|
||||
test_sender_enqueue_byte_budget();
|
||||
}
|
||||
|
||||
@@ -1369,6 +1369,139 @@ static void test_scanner_chunk_ownership() {
|
||||
rmdir(dir);
|
||||
}
|
||||
|
||||
/* The scanner derives entry type from a single lstat() for non-symlinks
|
||||
* (regular files and directories) and only calls stat() to dereference real
|
||||
* symlinks. Guard the regular-file/directory/symlink distinction across the
|
||||
* default (symlinks skipped), --copy-links (dereferenced) and -l (carried)
|
||||
* modes so the lstat/stat reuse cannot misclassify entries. */
|
||||
static void test_scanner_entry_classification() {
|
||||
const char* root = "test_scan_classify";
|
||||
const char* sub = "test_scan_classify/sub";
|
||||
const char* file = "test_scan_classify/file.txt";
|
||||
const char* nested = "test_scan_classify/sub/nested.txt";
|
||||
const char* link_file = "test_scan_classify/link_file";
|
||||
const char* link_dir = "test_scan_classify/link_dir";
|
||||
|
||||
EXPECT_EQ_INT(mkdir(root, 0755), 0);
|
||||
EXPECT_EQ_INT(mkdir(sub, 0755), 0);
|
||||
create_test_file(file, "hello"); /* 5 bytes */
|
||||
create_test_file(nested, "nested"); /* 6 bytes */
|
||||
EXPECT_EQ_INT(symlink("file.txt", link_file), 0);
|
||||
EXPECT_EQ_INT(symlink("sub", link_dir), 0);
|
||||
|
||||
/* Default: no link option -> symlinks are skipped entirely; regular files and
|
||||
directories (descended, not emitted) are classified as before. */
|
||||
{
|
||||
ScannerOptions options = {0};
|
||||
DirectoryScanner* scanner = directory_scanner_create_with_options(root, &options);
|
||||
EXPECT_NOT_NULL(scanner);
|
||||
size_t root_len = strlen(root);
|
||||
bool file_ok = false, nested_ok = false, link_seen = false;
|
||||
Chunk* chunk;
|
||||
while ((chunk = directory_scanner_next(scanner)) != NULL) {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
const File* f = chunk->items[i];
|
||||
const char* rel = f->path + root_len;
|
||||
if (*rel == '/')
|
||||
rel++;
|
||||
if (strcmp(rel, "file.txt") == 0) {
|
||||
file_ok = !f->is_dir && !f->is_symlink && f->data->size == 5;
|
||||
} else if (strcmp(rel, "sub/nested.txt") == 0) {
|
||||
nested_ok = !f->is_dir && !f->is_symlink && f->data->size == 6;
|
||||
} else {
|
||||
link_seen = true;
|
||||
}
|
||||
}
|
||||
chunk_destroy(chunk);
|
||||
}
|
||||
EXPECT_FALSE(directory_scanner_failed(scanner));
|
||||
EXPECT_TRUE(file_ok);
|
||||
EXPECT_TRUE(nested_ok);
|
||||
EXPECT_FALSE(link_seen);
|
||||
directory_scanner_destroy(scanner);
|
||||
}
|
||||
|
||||
/* --copy-links: symlinks are dereferenced. A link to a file becomes a
|
||||
regular file with the referent's size; a link to a directory is traversed. */
|
||||
{
|
||||
ScannerOptions options = {0};
|
||||
options.copy_links = true;
|
||||
DirectoryScanner* scanner = directory_scanner_create_with_options(root, &options);
|
||||
EXPECT_NOT_NULL(scanner);
|
||||
size_t root_len = strlen(root);
|
||||
bool file_ok = false, nested_ok = false, link_file_ok = false;
|
||||
bool link_dir_nested_ok = false, symlink_leaked = false;
|
||||
Chunk* chunk;
|
||||
while ((chunk = directory_scanner_next(scanner)) != NULL) {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
const File* f = chunk->items[i];
|
||||
const char* rel = f->path + root_len;
|
||||
if (*rel == '/')
|
||||
rel++;
|
||||
if (f->is_symlink)
|
||||
symlink_leaked = true;
|
||||
if (strcmp(rel, "file.txt") == 0)
|
||||
file_ok = !f->is_dir && f->data->size == 5;
|
||||
else if (strcmp(rel, "sub/nested.txt") == 0)
|
||||
nested_ok = !f->is_dir && f->data->size == 6;
|
||||
else if (strcmp(rel, "link_file") == 0)
|
||||
link_file_ok = !f->is_dir && f->data->size == 5;
|
||||
else if (strcmp(rel, "link_dir/nested.txt") == 0)
|
||||
link_dir_nested_ok = !f->is_dir && f->data->size == 6;
|
||||
}
|
||||
chunk_destroy(chunk);
|
||||
}
|
||||
EXPECT_FALSE(directory_scanner_failed(scanner));
|
||||
EXPECT_TRUE(file_ok);
|
||||
EXPECT_TRUE(nested_ok);
|
||||
EXPECT_TRUE(link_file_ok);
|
||||
EXPECT_TRUE(link_dir_nested_ok);
|
||||
EXPECT_FALSE(symlink_leaked);
|
||||
directory_scanner_destroy(scanner);
|
||||
}
|
||||
|
||||
/* -l (--links): symlinks are carried through as symlinks, not dereferenced. */
|
||||
{
|
||||
ScannerOptions options = {0};
|
||||
options.follow_symlinks = true;
|
||||
DirectoryScanner* scanner = directory_scanner_create_with_options(root, &options);
|
||||
EXPECT_NOT_NULL(scanner);
|
||||
size_t root_len = strlen(root);
|
||||
bool file_ok = false, link_file_ok = false, link_dir_ok = false, leaked_dir = false;
|
||||
Chunk* chunk;
|
||||
while ((chunk = directory_scanner_next(scanner)) != NULL) {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
const File* f = chunk->items[i];
|
||||
const char* rel = f->path + root_len;
|
||||
if (*rel == '/')
|
||||
rel++;
|
||||
if (strcmp(rel, "file.txt") == 0)
|
||||
file_ok = !f->is_dir && !f->is_symlink && f->data->size == 5;
|
||||
else if (strcmp(rel, "link_file") == 0)
|
||||
link_file_ok = f->is_symlink && !f->is_dir;
|
||||
else if (strcmp(rel, "link_dir") == 0)
|
||||
link_dir_ok = f->is_symlink && !f->is_dir;
|
||||
else if (strcmp(rel, "link_dir/nested.txt") == 0)
|
||||
leaked_dir = true;
|
||||
}
|
||||
chunk_destroy(chunk);
|
||||
}
|
||||
EXPECT_FALSE(directory_scanner_failed(scanner));
|
||||
EXPECT_TRUE(file_ok);
|
||||
EXPECT_TRUE(link_file_ok);
|
||||
EXPECT_TRUE(link_dir_ok);
|
||||
EXPECT_FALSE(leaked_dir);
|
||||
directory_scanner_destroy(scanner);
|
||||
}
|
||||
|
||||
unlink(link_file);
|
||||
unlink(link_dir);
|
||||
unlink(nested);
|
||||
unlink(file);
|
||||
rmdir(sub);
|
||||
rmdir(root);
|
||||
}
|
||||
|
||||
void test_scanner() {
|
||||
test_scanner_single_file();
|
||||
test_scanner_multiple_files();
|
||||
@@ -1406,4 +1539,5 @@ void test_scanner() {
|
||||
test_files_from_relative_send_path();
|
||||
test_scanner_captures_directory_times();
|
||||
test_scanner_chunk_ownership();
|
||||
test_scanner_entry_classification();
|
||||
}
|
||||
|
||||
@@ -150,6 +150,43 @@ static void test_walker_removes_extras_keeps_manifest_and_protected() {
|
||||
free(root);
|
||||
}
|
||||
|
||||
static void test_walker_keeps_nested_manifest_dirs() {
|
||||
/* The keep-set index must preserve deep content: a directory is protected
|
||||
when its own name is a keep entry OR when kept content lives below it, and
|
||||
an exact kept file survives while its siblings are removed. */
|
||||
char* root = make_walk_root("nestedkeep");
|
||||
EXPECT_NOT_NULL(root);
|
||||
EXPECT_TRUE(write_file_at(root, "extra.txt", "extra"));
|
||||
EXPECT_EQ_INT(make_subdir(root, "keepdir"), 0);
|
||||
EXPECT_EQ_INT(make_subdir(root, "keepdir/deep"), 0);
|
||||
EXPECT_TRUE(write_file_at(root, "keepdir/deep/keep.txt", "kept"));
|
||||
EXPECT_TRUE(write_file_at(root, "keepdir/extra2.txt", "extra"));
|
||||
EXPECT_EQ_INT(make_subdir(root, "dropdir"), 0);
|
||||
EXPECT_EQ_INT(make_subdir(root, "keep2"), 0);
|
||||
EXPECT_TRUE(write_file_at(root, "keep2/inner.txt", "kept"));
|
||||
EXPECT_EQ_INT(make_subdir(root, "keep3"), 0);
|
||||
|
||||
const char* keeps[] = {"keepdir/deep/keep.txt", "keep2/inner.txt", "keep3"};
|
||||
ArrayList* manifest = make_manifest_strings(keeps, 3);
|
||||
EXPECT_NOT_NULL(manifest);
|
||||
size_t deleted = 0;
|
||||
DeleteWalkResult result = delete_extras_limited(root, manifest, 100000, NULL, 0, &deleted);
|
||||
EXPECT_EQ_INT((int)result, (int)DELETE_WALK_OK);
|
||||
EXPECT_FALSE(file_exists(root, "extra.txt"));
|
||||
EXPECT_TRUE(file_exists(root, "keepdir/deep/keep.txt"));
|
||||
EXPECT_FALSE(file_exists(root, "keepdir/extra2.txt"));
|
||||
EXPECT_TRUE(dir_exists(root, "keepdir"));
|
||||
EXPECT_TRUE(dir_exists(root, "keepdir/deep"));
|
||||
EXPECT_FALSE(dir_exists(root, "dropdir"));
|
||||
EXPECT_TRUE(dir_exists(root, "keep2"));
|
||||
EXPECT_TRUE(file_exists(root, "keep2/inner.txt"));
|
||||
EXPECT_TRUE(dir_exists(root, "keep3")); /* an exact directory keep entry survives */
|
||||
EXPECT_EQ_INT((int)deleted, 3);
|
||||
array_list_delete(manifest);
|
||||
remove_walk_tree(root);
|
||||
free(root);
|
||||
}
|
||||
|
||||
static void test_walker_max_delete_exceeded_deletes_nothing() {
|
||||
char* root = make_walk_root("maxdel");
|
||||
EXPECT_NOT_NULL(root);
|
||||
@@ -419,8 +456,74 @@ static void test_fd_peer_ip() {
|
||||
EXPECT_EQ_STR(peer_string, "");
|
||||
}
|
||||
|
||||
/* The keep/files-from indexes must store exactly the input entries (one node
|
||||
each), never a copied ancestor prefix per component. This builds a PathIndex
|
||||
over paths thousands of components deep and checks the structural bound plus
|
||||
the exact / descendant query semantics. */
|
||||
static void test_path_index_bounded() {
|
||||
enum { COUNT = 8, COMPONENTS = 5000 };
|
||||
size_t entry_len = (size_t)COMPONENTS * 2 + 2; /* trailing "xN" */
|
||||
char* storage = malloc((size_t)COUNT * (entry_len + 1));
|
||||
EXPECT_NOT_NULL(storage);
|
||||
const char** entries = calloc(COUNT, sizeof(char*));
|
||||
EXPECT_NOT_NULL(entries);
|
||||
for (int i = 0; i < COUNT; i++) {
|
||||
char* entry = storage + (size_t)i * (entry_len + 1);
|
||||
size_t pos = 0;
|
||||
for (int c = 0; c < COMPONENTS; c++) {
|
||||
entry[pos++] = 'a';
|
||||
entry[pos++] = '/';
|
||||
}
|
||||
entry[pos++] = 'x';
|
||||
entry[pos++] = (char)('0' + i);
|
||||
entry[pos] = '\0';
|
||||
entries[i] = entry;
|
||||
}
|
||||
|
||||
PathIndex index;
|
||||
EXPECT_TRUE(path_index_build(&index, entries, COUNT));
|
||||
EXPECT_EQ_INT((int)index.sorted.count, COUNT);
|
||||
EXPECT_EQ_INT((int)index.exact.size, COUNT);
|
||||
EXPECT_TRUE(path_index_contains(&index, entries[0]));
|
||||
EXPECT_FALSE(path_index_contains(&index, "a"));
|
||||
EXPECT_TRUE(path_index_has_descendant(&index, "a"));
|
||||
EXPECT_TRUE(path_index_has_descendant(&index, "a/a"));
|
||||
EXPECT_FALSE(path_index_has_descendant(&index, "aa"));
|
||||
path_index_free(&index);
|
||||
|
||||
free((void*)entries);
|
||||
free(storage);
|
||||
}
|
||||
|
||||
static void test_path_index_semantics() {
|
||||
const char* entries[] = {"a/b/c.txt", "a/b/d.txt", "x.txt", "deep/deeper/deepest"};
|
||||
PathIndex index;
|
||||
EXPECT_TRUE(path_index_build(&index, entries, 4));
|
||||
EXPECT_TRUE(path_index_contains(&index, "a/b/c.txt"));
|
||||
EXPECT_FALSE(path_index_contains(&index, "a/b"));
|
||||
EXPECT_TRUE(path_index_contains_n(&index, "a/b/c.txt/ignored", 9));
|
||||
EXPECT_FALSE(path_index_contains_n(&index, "a/b/c.txt/ignored", 10));
|
||||
EXPECT_TRUE(path_index_has_descendant(&index, "a"));
|
||||
EXPECT_TRUE(path_index_has_descendant(&index, "a/b"));
|
||||
EXPECT_FALSE(path_index_has_descendant(&index, "a/b/c.txt"));
|
||||
EXPECT_FALSE(path_index_has_descendant(&index, "ab"));
|
||||
EXPECT_FALSE(path_index_has_descendant(&index, ""));
|
||||
path_index_free(&index);
|
||||
|
||||
/* A zero-entry index answers no queries. */
|
||||
PathIndex empty;
|
||||
EXPECT_TRUE(path_index_build(&empty, NULL, 0));
|
||||
EXPECT_EQ_INT((int)empty.sorted.count, 0);
|
||||
EXPECT_FALSE(path_index_contains(&empty, "a"));
|
||||
EXPECT_FALSE(path_index_has_descendant(&empty, "a"));
|
||||
path_index_free(&empty);
|
||||
}
|
||||
|
||||
void test_shared_utils() {
|
||||
test_path_index_bounded();
|
||||
test_path_index_semantics();
|
||||
test_walker_removes_extras_keeps_manifest_and_protected();
|
||||
test_walker_keeps_nested_manifest_dirs();
|
||||
test_walker_max_delete_exceeded_deletes_nothing();
|
||||
test_walker_max_delete_exact_bound_deletes();
|
||||
test_walker_unlimited_deletes_all();
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
#include "protocol.h"
|
||||
#include "test_utils.h"
|
||||
#include "transport_tcp.h"
|
||||
#include <arpa/inet.h>
|
||||
#include <netinet/in.h>
|
||||
#include <netinet/tcp.h>
|
||||
#include <string.h>
|
||||
@@ -195,6 +196,49 @@ static void test_client_disconnect_delete() {
|
||||
client_delete(c);
|
||||
}
|
||||
|
||||
/* TCP_NODELAY is enabled by default on a connected transfer socket, and an
|
||||
* explicit --sockopts TCP_NODELAY=0 still overrides it. */
|
||||
static void test_tcp_nodelay_default_and_override() {
|
||||
Server* s = server_create(0);
|
||||
EXPECT_NOT_NULL(s);
|
||||
EXPECT_EQ_INT(listen(s->file_descriptor, 1), 0);
|
||||
struct sockaddr_in bound;
|
||||
socklen_t bound_len = sizeof(bound);
|
||||
EXPECT_EQ_INT(getsockname(s->file_descriptor, (struct sockaddr*)&bound, &bound_len), 0);
|
||||
int port = (int)ntohs(bound.sin_port);
|
||||
EXPECT_TRUE(port > 0);
|
||||
|
||||
Client* c = client_create();
|
||||
EXPECT_NOT_NULL(c);
|
||||
EXPECT_TRUE(client_connect(c, "127.0.0.1", port));
|
||||
int got = 0;
|
||||
socklen_t len = sizeof(got);
|
||||
EXPECT_EQ_INT(getsockopt(c->file_descriptor, IPPROTO_TCP, TCP_NODELAY, &got, &len), 0);
|
||||
EXPECT_EQ_INT(got, 1);
|
||||
client_disconnect(c);
|
||||
client_delete(c);
|
||||
|
||||
SockOptEntry* entries = NULL;
|
||||
int count = 0;
|
||||
EXPECT_EQ_INT(config_sockopts_parse("TCP_NODELAY=0", &entries, &count), 0);
|
||||
TcpConnectOptions opts;
|
||||
memset(&opts, 0, sizeof(opts));
|
||||
opts.sockopts = entries;
|
||||
opts.sockopt_count = count;
|
||||
|
||||
Client* c2 = client_create();
|
||||
EXPECT_NOT_NULL(c2);
|
||||
EXPECT_TRUE(client_connect_ex(c2, "127.0.0.1", port, &opts));
|
||||
got = 0;
|
||||
len = sizeof(got);
|
||||
EXPECT_EQ_INT(getsockopt(c2->file_descriptor, IPPROTO_TCP, TCP_NODELAY, &got, &len), 0);
|
||||
EXPECT_EQ_INT(got, 0);
|
||||
client_disconnect(c2);
|
||||
client_delete(c2);
|
||||
free(entries);
|
||||
server_delete(&s);
|
||||
}
|
||||
|
||||
void test_transport_tcp() {
|
||||
test_server_create_ephemeral();
|
||||
test_server_delete_null();
|
||||
@@ -211,4 +255,5 @@ void test_transport_tcp() {
|
||||
test_sockopts_apply_sets_option();
|
||||
test_server_create_bind_address();
|
||||
test_server_create_bind_ipv6();
|
||||
test_tcp_nodelay_default_and_override();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user