fix: correct throttle legacy resolution; add EXDEV temp-dir coverage
This commit is contained in:
+6
-1
@@ -1013,7 +1013,12 @@ integration tests unless it is explicitly listed as a limitation.
|
|||||||
|
|
||||||
- **`--temp-dir` is confined to the receive root on the receiver:** a relative
|
- **`--temp-dir` is confined to the receive root on the receiver:** a relative
|
||||||
dir resolves below it; an absolute path or one containing `..` is rejected.
|
dir resolves below it; an absolute path or one containing `..` is rejected.
|
||||||
An `EXDEV` install falls back to a non-atomic copy instead of aborting.
|
An `EXDEV` install falls back to a non-atomic copy instead of aborting. (The
|
||||||
|
confined receiver path cannot be mount-tested in the CI container — no
|
||||||
|
`CAP_SYS_ADMIN` and unprivileged user namespaces are disabled — so the
|
||||||
|
cross-filesystem fallback is exercised end-to-end through the unconfined local
|
||||||
|
`--read-batch` apply against a `/dev/shm` scratch dir, in
|
||||||
|
`tests/integration/test_temp_dir_exdev.py`.)
|
||||||
- **Deletion scoping:** the manifest carries the synchronized directories, so
|
- **Deletion scoping:** the manifest carries the synchronized directories, so
|
||||||
the extras walk only visits their subtrees; `--files-from` subsets no longer
|
the extras walk only visits their subtrees; `--files-from` subsets no longer
|
||||||
delete untransmitted paths outside the listed directories.
|
delete untransmitted paths outside the listed directories.
|
||||||
|
|||||||
@@ -185,7 +185,7 @@ bool file_send_sendfile_with_skip(File* file, int file_descriptor, bool use_meta
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
protocol_note_bytes_written((unsigned long long)sent);
|
protocol_note_bytes_written((unsigned long long)sent);
|
||||||
protocol_throttle_bytes((size_t)sent);
|
protocol_throttle_bytes(file_descriptor, (size_t)sent);
|
||||||
}
|
}
|
||||||
|
|
||||||
close(fd);
|
close(fd);
|
||||||
|
|||||||
@@ -302,11 +302,14 @@ static ProtocolSession* legacy_session(int read_fd, int write_fd) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/* Pace an out-of-band write that bypassed protocol_send_n_data (the plaintext
|
/* Pace an out-of-band write that bypassed protocol_send_n_data (the plaintext
|
||||||
* sendfile fast path). The bound/legacy session is resolved exactly as
|
* sendfile fast path). The bound/legacy session is resolved exactly as the
|
||||||
* send_n_data resolves it, so the same token-bucket state is throttled and the
|
* preceding send_n_data(fd, ...) resolved it, so the same token-bucket state is
|
||||||
* TLS and plaintext transports share identical --bwlimit semantics. */
|
* throttled and the TLS and plaintext transports share identical --bwlimit
|
||||||
void protocol_throttle_bytes(size_t bytes) {
|
* semantics. Passing the wire fd (rather than -1) is essential: the sendfile
|
||||||
bw_throttle_session(legacy_session(-1, -1), bytes);
|
* send left legacy_io_session.write_fd bound to it, so resolving with -1 would
|
||||||
|
* mismatch, re-initialize the session and hand out a second first-call burst. */
|
||||||
|
void protocol_throttle_bytes(int file_descriptor, size_t bytes) {
|
||||||
|
bw_throttle_session(legacy_session(-1, file_descriptor), bytes);
|
||||||
}
|
}
|
||||||
|
|
||||||
bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
|
bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
|
||||||
|
|||||||
@@ -225,11 +225,14 @@ unsigned long long protocol_bytes_written(void);
|
|||||||
unsigned long long protocol_bytes_read(void);
|
unsigned long long protocol_bytes_read(void);
|
||||||
void protocol_note_bytes_written(unsigned long long bytes);
|
void protocol_note_bytes_written(unsigned long long bytes);
|
||||||
/* Apply --bwlimit pacing to bytes written outside protocol_send_n_data (the
|
/* Apply --bwlimit pacing to bytes written outside protocol_send_n_data (the
|
||||||
* plaintext zero-copy sendfile fast path). Resolves the bound/legacy session
|
* plaintext zero-copy sendfile fast path). `file_descriptor` is the wire fd
|
||||||
* exactly as send_n_data does and runs the same token-bucket throttle, so the
|
* the bytes were written to, so the legacy session is resolved exactly as the
|
||||||
* sendfile transport is paced identically to the buffered/TLS paths. A no-op
|
* preceding send_n_data call resolved it (the bound TLS session still wins when
|
||||||
* when the effective session has no bandwidth limit. */
|
* set); resolving with the same fd avoids re-initializing the legacy session
|
||||||
void protocol_throttle_bytes(size_t bytes);
|
* and granting a second first-call burst. Runs the same token-bucket throttle,
|
||||||
|
* so the sendfile transport is paced identically to the buffered/TLS paths. A
|
||||||
|
* no-op when the effective session has no bandwidth limit. */
|
||||||
|
void protocol_throttle_bytes(int file_descriptor, size_t bytes);
|
||||||
|
|
||||||
void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd);
|
void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd);
|
||||||
/* Transitional bridge for helpers whose signatures still carry only an fd. */
|
/* Transitional bridge for helpers whose signatures still carry only an fd. */
|
||||||
|
|||||||
@@ -0,0 +1,80 @@
|
|||||||
|
"""End-to-end coverage for the `--temp-dir` EXDEV (cross-filesystem) fallback.
|
||||||
|
|
||||||
|
`file_to_disk_secure_impl` installs a completed temp file with `renameat(2)`;
|
||||||
|
when the scratch dir lives on a different filesystem the rename fails with
|
||||||
|
`EXDEV` and the engine retries with no scratch dir, writing the file directly in
|
||||||
|
the destination directory (a non-atomic copy), matching rsync.
|
||||||
|
|
||||||
|
The daemon receiver confines `--temp-dir` to the authorized receive root, so a
|
||||||
|
genuine cross-fs scratch there would require an in-root mount point. Bind/tmpfs
|
||||||
|
mounting is not permitted in the CI container (no `CAP_SYS_ADMIN`, and
|
||||||
|
unprivileged user namespaces are disabled), so this test reaches the exact same
|
||||||
|
code path through the local `--read-batch` apply instead: it has no
|
||||||
|
authorized-root confinement, so a relative `--temp-dir` that is a symlink to a
|
||||||
|
tmpfs (`/dev/shm`) is accepted and the final install then crosses filesystems.
|
||||||
|
"""
|
||||||
|
import os
|
||||||
|
import shutil
|
||||||
|
import subprocess
|
||||||
|
import sys
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
sys.path.insert(0, os.path.dirname(__file__))
|
||||||
|
from common import CLIENT_CMD, get_dest_received_dir
|
||||||
|
|
||||||
|
TMPFS = "/dev/shm"
|
||||||
|
|
||||||
|
|
||||||
|
def _run(args):
|
||||||
|
return subprocess.run(CLIENT_CMD + args, capture_output=True, text=True, timeout=180)
|
||||||
|
|
||||||
|
|
||||||
|
def _read(path):
|
||||||
|
with open(path, "rb") as fh:
|
||||||
|
return fh.read()
|
||||||
|
|
||||||
|
|
||||||
|
def test_read_batch_temp_dir_cross_filesystem_fallback(tmp_path):
|
||||||
|
if not os.path.isdir(TMPFS):
|
||||||
|
pytest.skip("no /dev/shm tmpfs available to force a cross-filesystem install")
|
||||||
|
|
||||||
|
source = tmp_path / "src"
|
||||||
|
dest = tmp_path / "dst"
|
||||||
|
source.mkdir()
|
||||||
|
dest.mkdir()
|
||||||
|
files = {
|
||||||
|
"payload.bin": bytes(range(256)) * 64,
|
||||||
|
"sub/nested.txt": b"nested exdev fallback\n" * 8,
|
||||||
|
}
|
||||||
|
for rel, data in files.items():
|
||||||
|
full = source / rel
|
||||||
|
full.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
full.write_bytes(data)
|
||||||
|
|
||||||
|
batch = tmp_path / "tree.batch"
|
||||||
|
r = _run(["--only-write-batch", str(batch), str(source)])
|
||||||
|
assert r.returncode == 0, (r.stdout, r.stderr)
|
||||||
|
|
||||||
|
# A cross-filesystem scratch dir, reached through a relative --temp-dir
|
||||||
|
# symlink (the local batch apply performs no authorized-root confinement).
|
||||||
|
scratch = os.path.join(TMPFS, "fastsync_exdev_%d" % os.getpid())
|
||||||
|
shutil.rmtree(scratch, ignore_errors=True)
|
||||||
|
os.makedirs(scratch)
|
||||||
|
os.symlink(scratch, dest / "scratch")
|
||||||
|
try:
|
||||||
|
assert os.stat(scratch).st_dev != os.stat(dest).st_dev, (
|
||||||
|
"scratch and destination share a filesystem; EXDEV cannot be exercised"
|
||||||
|
)
|
||||||
|
r = _run(["--read-batch", str(batch), str(dest), "--temp-dir=scratch"])
|
||||||
|
assert r.returncode == 0, (r.stdout, r.stderr)
|
||||||
|
# The engine must report the non-atomic cross-fs fallback rather than
|
||||||
|
# silently claiming an atomic install.
|
||||||
|
assert "different filesystem" in (r.stdout + r.stderr), (r.stdout, r.stderr)
|
||||||
|
# The tree is still byte-exact and the scratch dir is left clean.
|
||||||
|
received = get_dest_received_dir(str(dest), str(source))
|
||||||
|
for rel, data in files.items():
|
||||||
|
assert _read(os.path.join(received, rel)) == data, f"content mismatch for {rel}"
|
||||||
|
assert os.listdir(scratch) == [], "cross-fs temp file was not cleaned up"
|
||||||
|
finally:
|
||||||
|
shutil.rmtree(scratch, ignore_errors=True)
|
||||||
+43
-2
@@ -1,6 +1,8 @@
|
|||||||
#include "protocol.h"
|
#include "protocol.h"
|
||||||
#include "test_utils.h"
|
#include "test_utils.h"
|
||||||
|
#include <fcntl.h>
|
||||||
#include <limits.h>
|
#include <limits.h>
|
||||||
|
#include <stdlib.h>
|
||||||
#include <string.h>
|
#include <string.h>
|
||||||
#include <time.h>
|
#include <time.h>
|
||||||
#include <unistd.h>
|
#include <unistd.h>
|
||||||
@@ -715,7 +717,7 @@ static void test_protocol_throttle_bytes_paces() {
|
|||||||
|
|
||||||
struct timespec start;
|
struct timespec start;
|
||||||
clock_gettime(CLOCK_MONOTONIC, &start);
|
clock_gettime(CLOCK_MONOTONIC, &start);
|
||||||
protocol_throttle_bytes(150000);
|
protocol_throttle_bytes(-1, 150000);
|
||||||
struct timespec now;
|
struct timespec now;
|
||||||
clock_gettime(CLOCK_MONOTONIC, &now);
|
clock_gettime(CLOCK_MONOTONIC, &now);
|
||||||
long long elapsed_ms =
|
long long elapsed_ms =
|
||||||
@@ -736,7 +738,7 @@ static void test_protocol_throttle_bytes_unlimited() {
|
|||||||
|
|
||||||
struct timespec start;
|
struct timespec start;
|
||||||
clock_gettime(CLOCK_MONOTONIC, &start);
|
clock_gettime(CLOCK_MONOTONIC, &start);
|
||||||
protocol_throttle_bytes(100000000ULL);
|
protocol_throttle_bytes(-1, 100000000ULL);
|
||||||
struct timespec now;
|
struct timespec now;
|
||||||
clock_gettime(CLOCK_MONOTONIC, &now);
|
clock_gettime(CLOCK_MONOTONIC, &now);
|
||||||
long long elapsed_ms =
|
long long elapsed_ms =
|
||||||
@@ -746,6 +748,44 @@ static void test_protocol_throttle_bytes_unlimited() {
|
|||||||
protocol_session_unbind();
|
protocol_session_unbind();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* Regression for the plaintext sendfile path: it calls protocol_throttle_bytes()
|
||||||
|
* immediately after send_n_data(), which already bound legacy_io_session.write_fd
|
||||||
|
* to the wire fd. Resolving the throttle session with (read=-1, write=-1)
|
||||||
|
* mismatched that fd and re-initialized the legacy session, granting a *second*
|
||||||
|
* first-call burst and discarding the accumulated debt. This drives the same
|
||||||
|
* sequence and asserts the debt from send_n_data carries into the throttle. */
|
||||||
|
static void test_protocol_throttle_bytes_legacy_same_session() {
|
||||||
|
const size_t payload = 150000; /* 1.5x the 100 KB burst at --bwlimit=1 MB/s */
|
||||||
|
unsigned char* buffer = malloc(payload);
|
||||||
|
EXPECT_TRUE(buffer != NULL);
|
||||||
|
memset(buffer, 0, payload);
|
||||||
|
|
||||||
|
io_set_fds(-1, -1);
|
||||||
|
io_set_bwlimit(1000000ULL);
|
||||||
|
|
||||||
|
int fd = open("/dev/null", O_WRONLY);
|
||||||
|
EXPECT_TRUE(fd >= 0);
|
||||||
|
|
||||||
|
struct timespec start;
|
||||||
|
clock_gettime(CLOCK_MONOTONIC, &start);
|
||||||
|
/* send_n_data() consumes the whole 100 KB burst and sleeps ~50 ms. */
|
||||||
|
EXPECT_TRUE(send_n_data(fd, buffer, payload));
|
||||||
|
/* The throttle must share that session, so the 150 KB is all debt and sleeps
|
||||||
|
~150 ms (total ~200 ms). A re-initialized session would hand out a fresh
|
||||||
|
100 KB burst and sleep only ~50 ms (total ~100 ms). */
|
||||||
|
protocol_throttle_bytes(fd, payload);
|
||||||
|
struct timespec now;
|
||||||
|
clock_gettime(CLOCK_MONOTONIC, &now);
|
||||||
|
long long elapsed_ms =
|
||||||
|
(now.tv_sec - start.tv_sec) * 1000LL + (now.tv_nsec - start.tv_nsec) / 1000000LL;
|
||||||
|
EXPECT_TRUE(elapsed_ms >= 150);
|
||||||
|
|
||||||
|
close(fd);
|
||||||
|
free(buffer);
|
||||||
|
io_set_bwlimit(0);
|
||||||
|
io_set_fds(-1, -1);
|
||||||
|
}
|
||||||
|
|
||||||
void test_protocol() {
|
void test_protocol() {
|
||||||
test_send_receive_n_data();
|
test_send_receive_n_data();
|
||||||
test_send_receive_n_data_zero();
|
test_send_receive_n_data_zero();
|
||||||
@@ -778,4 +818,5 @@ void test_protocol() {
|
|||||||
test_data_create_starts_uncharged_and_unowned();
|
test_data_create_starts_uncharged_and_unowned();
|
||||||
test_protocol_throttle_bytes_paces();
|
test_protocol_throttle_bytes_paces();
|
||||||
test_protocol_throttle_bytes_unlimited();
|
test_protocol_throttle_bytes_unlimited();
|
||||||
|
test_protocol_throttle_bytes_legacy_same_session();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user