fix(p6-stop): block delete-manifest on early stop; sync -m manifest access; overflow guard

This commit is contained in:
2026-09-10 17:39:34 +02:00
parent 37cff96537
commit ac7e9e3bc1
6 changed files with 210 additions and 39 deletions
+52 -26
View File
@@ -1403,6 +1403,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
if (stop_condition_reached(&context->stop_condition)) {
log_info_message(LOG_INFO_MISC,
"Stop deadline reached; stopping transfer at the next chunk boundary");
context->scan_stopped_early = true;
pipeline_cancel(context);
break;
}
@@ -1446,9 +1447,19 @@ static int send_chunks_multithreaded(void* pipeline_context) {
}
/* Completion tail: reached on natural exhaustion or an early stop deadline.
The keep-set manifest is committed exactly where it normally would be, so
--delete (delete-after timing) still commits its extras deletion. */
if (context->config->use_delete && !context->early_delete) {
A deadline that cut the scan short leaves an incomplete keep-set manifest;
transmitting it would make the receiver --delete the unscanned source
mirrors (data loss), so it is deliberately suppressed. Suppressing it also
means the manifest (which the scanner thread may still be appending) is
never read here on the early-stop path, so no scanner synchronization is
required to enter the tail. */
context->scan_stopped_early =
context->scan_stopped_early || stop_condition_reached(&context->stop_condition);
if (context->scan_stopped_early) {
log_message(LOG_LEVEL_WARNING,
"transfer stopped early (stop deadline); skipping --delete keep-set so "
"unscanned source mirrors are not deleted");
} else if (context->config->use_delete && !context->early_delete) {
/* Empty keep-set + scan I/O error must not delete the whole destination
(the source may not be genuinely empty -- see send_files). */
bool empty_io;
@@ -1826,15 +1837,18 @@ int send_files(Config* config) {
int total_files = 0;
time_t last_progress = 0;
time_t start = time(NULL);
/* True when the stop deadline cut the scan short so the keep-set manifest is
only a prefix of the source. */
bool scan_stopped_early = false;
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
/* Phase 6: stop-elegantly at the next chunk boundary once the deadline has
passed. The scanner may also have stopped early itself; either way the
completion tail below keeps everything already sent and still commits a
--delete keep-set manifest. */
completion tail below keeps everything already sent. */
if (stop_condition_reached(&stop)) {
chunk_destroy(current_chunk);
log_info_message(LOG_INFO_MISC,
"Stop deadline reached; stopping transfer at the next chunk boundary");
scan_stopped_early = true;
break;
}
unsigned long long chunk_bytes = 0;
@@ -1890,31 +1904,43 @@ int send_files(Config* config) {
goto send_fail;
if (directory_scanner_had_io_error(scanner))
had_scan_io = true;
if (had_scan_io && manifest && manifest->size == 0) {
/* A scan that hit an I/O error and produced no keep entries is ambiguous;
an empty keep-set would delete the whole destination. Refuse to delete
(see the early-timing comment above). */
log_message(LOG_LEVEL_ERROR,
"source scan hit an I/O error before finding any file; refusing to delete with "
"an empty keep-set (--delete)");
goto send_fail;
}
if ((manifest || config->delete_missing_args) && !delete_early) {
/* Late (commit) ordering: all file data is out; transmit the manifest so
the receiver commits the extras walk (--delete) and/or the
--delete-missing-args exact-path deletions only after the transfer
succeeds. In the early modes (--delete-before/--delete-during) the
manifest already went out up front, so nothing is re-sent here. */
if (send_delete_manifest(client->file_descriptor, manifest, excluded, missing_args) != 0) {
/* Phase 6: the scanner may have stopped early (returning NULL without a
failure) as soon as the deadline passed, so reflect that here too. A
deadline that cut the scan short leaves an incomplete keep-set; transmitting
it would make the receiver --delete the unscanned source mirrors (data
loss), so the late delete manifest is suppressed below. */
scan_stopped_early = scan_stopped_early || stop_condition_reached(&stop);
if (scan_stopped_early) {
log_message(LOG_LEVEL_WARNING,
"transfer stopped early (stop deadline); skipping --delete keep-set so "
"unscanned source mirrors are not deleted");
} else {
if (had_scan_io && manifest && manifest->size == 0) {
/* A scan that hit an I/O error and produced no keep entries is ambiguous;
an empty keep-set would delete the whole destination. Refuse to delete
(see the early-timing comment above). */
log_message(LOG_LEVEL_ERROR,
"source scan hit an I/O error before finding any file; refusing to delete with "
"an empty keep-set (--delete)");
goto send_fail;
}
if ((manifest || config->delete_missing_args) && !delete_early) {
/* Late (commit) ordering: all file data is out; transmit the manifest so
the receiver commits the extras walk (--delete) and/or the
--delete-missing-args exact-path deletions only after the transfer
succeeds. In the early modes (--delete-before/--delete-during) the
manifest already went out up front, so nothing is re-sent here. */
if (send_delete_manifest(client->file_descriptor, manifest, excluded, missing_args) != 0) {
if (manifest) {
array_list_delete(manifest);
manifest = NULL;
}
goto send_fail;
}
if (manifest) {
array_list_delete(manifest);
manifest = NULL;
}
goto send_fail;
}
if (manifest) {
array_list_delete(manifest);
manifest = NULL;
}
}
bool ok = finalize_transfer(client, config, remove_sources);
+3 -1
View File
@@ -188,7 +188,9 @@ void print_usage(void) {
printf(" integer); whatever was already transferred is kept\n");
printf(" --stop-at=TIME Stop at an absolute time: HH:MM, HH:MM:SS, or\n");
printf(" now+N[smhd] (a time already in the past stops the\n");
printf(" transfer immediately; client-only)\n");
printf(" transfer immediately; client-only). An early stop\n");
printf(" skips the late --delete keep-set so it cannot delete\n");
printf(" source mirrors that were not yet scanned\n");
printf(" --address <ip> Bind the outgoing client socket to this source address\n");
printf(" -4, --ipv4 Force IPv4 for destination resolution\n");
printf(" -6, --ipv6 Force IPv6 for destination resolution\n");
+7
View File
@@ -60,6 +60,13 @@ typedef struct {
/* Phase 6: client-only sender stop deadline, computed once before the worker
* threads start and shared read-only by the scanner and the sender thread. */
StopCondition stop_condition;
/* Phase 6: set when the scanner/sender reached the stop deadline before the
* scan (and thus the keep-set manifest) completed naturally. When true the
* completion tail must NOT transmit the partial manifest, or the receiver
* would delete unscanned source mirrors. Written by the sender thread
* before it reads the manifest, so no additional synchronization is needed
* to suppress the manifest. */
bool scan_stopped_early;
} PipelineContextSender;
typedef struct PipelineContextReceiver {
+21 -8
View File
@@ -4,16 +4,22 @@
#include <stdlib.h>
#include <string.h>
/* Parse a strictly positive decimal integer. Accepts leading '+' but not a
* leading '-', surrounding whitespace, fractional parts or trailing garbage. */
/* Parse a strictly positive decimal integer: only ASCII digits, no leading
* whitespace, sign or trailing garbage. */
static bool parse_positive_minutes(const char* value, long* out) {
if (!value || *value == '\0')
return false;
errno = 0;
char* end = NULL;
long v = strtol(value, &end, 10);
if (errno != 0 || end == value || *end != '\0')
if (*value < '0' || *value > '9')
return false;
long v = 0;
for (const char* p = value; *p != '\0'; p++) {
if (*p < '0' || *p > '9')
return false;
int digit = *p - '0';
if (v > (LONG_MAX - digit) / 10)
return false;
v = v * 10 + digit;
}
if (v <= 0 || v > INT_MAX)
return false;
*out = v;
@@ -45,10 +51,12 @@ bool stop_parse_at_time(const char* value, time_t now, time_t* out_deadline) {
/* now+N[smhd]: N whole units from the current wall clock. */
if (strncmp(value, "now+", 4) == 0) {
const char* p = value + 4;
if (*p == '\0')
/* The count must be a bare non-negative digit run: reject leading
whitespace ('now+ 5s') and a leading sign ('now++5s'). */
if (*p < '0' || *p > '9')
return false;
char* end = NULL;
errno = 0;
char* end = NULL;
long amount = strtol(p, &end, 10);
if (errno != 0 || end == p || amount < 0)
return false;
@@ -74,6 +82,11 @@ bool stop_parse_at_time(const char* value, time_t now, time_t* out_deadline) {
if (amount > LONG_MAX / unit_seconds)
return false;
long long delta = (long long)amount * unit_seconds;
/* Guard against signed overflow of now + delta. */
if ((long long)now > 0 && delta > (long long)LLONG_MAX - (long long)now)
return false;
if ((long long)now < 0 && delta < (long long)LLONG_MIN - (long long)now)
return false;
*out_deadline = now + (time_t)delta;
return true;
}
+112 -3
View File
@@ -7,6 +7,7 @@ elegantly at the next chunk/file boundary -- whatever was already transferred is
kept, the completion tail still runs, and the exit code is 0 (like rsync's
clean "stopped early" behavior). Malformed values are rejected up front.
"""
import filecmp
import os
import shutil
import time
@@ -58,6 +59,32 @@ def _seed_source(source):
fh.write(content)
def _seed_many(source, count=40, size=32 * 1024):
"""Create `count` same-size regular files (enough to span several chunks)."""
blob = os.urandom(size)
for i in range(count):
with open(os.path.join(source, f"f{i:04d}.dat"), "wb") as fh:
fh.write(blob)
def _seed_dest_by_transfer(source, dest, port, extra=None):
"""Do a plain full transfer source->dest so dest exactly mirrors source."""
run_client(source, dest, flags=(extra or []), port=port)
def _received_subset_matches(source, received):
"""Every file under `received` exists under `source` with identical bytes."""
if not os.path.isdir(received):
return not _received_files(received)
rels = _received_files(received)
for rel in rels:
src = os.path.join(source, rel)
dst = os.path.join(received, rel)
if not os.path.isfile(src) or not filecmp.cmp(src, dst, shallow=False):
return False
return True
class TestStopAfter:
@pytest.mark.ci
def test_stop_after_within_window(self, shared_server):
@@ -91,12 +118,15 @@ class TestStopAt:
cleanly (exit 0, nothing transferred)."""
source, dest = _make("past")
_seed_source(source)
now = time.localtime()
if now.tm_hour * 60 + now.tm_min >= 1:
# Use a same-day HH:MM two minutes in the past when that cannot roll
# over into the previous day (which would parse as a FUTURE time today);
# otherwise fall back to now+0s which is deterministically immediate.
lt = time.localtime()
if lt.tm_hour * 60 + lt.tm_min >= 3:
past = time.localtime(time.time() - 120)
stop_value = f"{past.tm_hour:02d}:{past.tm_min:02d}"
else:
stop_value = "now+0s" # first minute of the day: use "immediately now"
stop_value = "now+0s"
result, _ = run_client(source, dest, flags=[f"--stop-at={stop_value}"],
port=shared_server.port)
assert result.returncode == 0, \
@@ -130,3 +160,82 @@ class TestStopAt:
result, _ = run_client(source, dest, flags=[flag],
port=shared_server.port)
assert result.returncode != 0, f"{flag} should be rejected"
class TestStopPartial:
"""A genuine mid-transfer stop leaves a valid, strict non-empty prefix."""
@pytest.mark.ci
def test_stop_mid_transfer_leaves_valid_partial(self, shared_server):
"""With --bwlimit a real deadline cuts the transfer mid-way: what WAS
transferred is byte-identical, not everything is transferred, and the
run returns 0 without corrupting any file."""
source, dest = _make("partial")
_seed_many(source, count=60, size=32 * 1024)
flags = ["--chunk-size", "262144", "--bwlimit", "300", "--stop-at=now+3s"]
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, \
f"mid-transfer stop failed (rc {result.returncode}): " \
f"{(result.stderr or result.stdout)[:400]}"
received = get_dest_received_dir(dest, source)
got = _received_files(received)
assert len(got) > 0, "expected an early stop to still transfer a prefix"
assert len(got) < 60, \
f"expected a PARTIAL transfer (all 60 arrived): stopped too late"
assert _received_subset_matches(source, received), \
f"received files are not a byte-identical subset of the source"
class TestStopDelete:
"""--delete must never wipe the destination when the scan is cut short."""
def _seed(self, prefix, port, many=False):
source, dest = _make(prefix)
if many:
_seed_many(source, count=40, size=96 * 1024)
else:
_seed_source(source)
_seed_dest_by_transfer(source, dest, port)
return source, dest
@pytest.mark.ci
def test_stop_delete_immediate_preserves_source_mirrors(self, shared_server):
"""Immediate stop + --delete: the incomplete/empty keep-set must NOT
delete the seeded source mirrors (returncode 0, files survive)."""
source, dest = self._seed("del_imm", shared_server.port)
result, _ = run_client(source, dest, flags=["--delete", "--stop-at=now+0s"],
port=shared_server.port)
assert result.returncode == 0, \
f"--delete immediate stop failed: {(result.stderr or result.stdout)[:400]}"
received = get_dest_received_dir(dest, source)
mismatches, missing = verify_transfer(source, received)
assert not mismatches and not missing, \
f"--delete wiped source mirrors: missing={missing} mismatches={mismatches}"
@pytest.mark.ci
def test_stop_delete_midscan_preserves_source_mirrors(self, shared_server):
"""A mid-scan stop + --delete must suppress the partial keep-set so all
seeded source mirrors survive."""
source, dest = self._seed("del_mid", shared_server.port, many=True)
flags = ["--delete", "--chunk-size", "262144", "--bwlimit", "300", "--stop-at=now+3s"]
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, \
f"--delete mid-scan stop failed: {(result.stderr or result.stdout)[:400]}"
received = get_dest_received_dir(dest, source)
mismatches, missing = verify_transfer(source, received)
assert not mismatches and not missing, \
f"--delete mid-scan wiped source mirrors: missing={missing} mismatches={mismatches}"
@pytest.mark.ci
def test_stop_delete_multithreaded_preserves_source_mirrors(self, shared_server):
"""-m immediate stop + --delete: the completion tail must not read the
still-appendable manifest (no race) and must not delete the mirrors."""
source, dest = self._seed("del_mt", shared_server.port, many=True)
result, _ = run_client(source, dest, flags=["-m", "--delete", "--stop-at=now+0s"],
port=shared_server.port)
assert result.returncode == 0, \
f"-m --delete immediate stop failed: {(result.stderr or result.stdout)[:400]}"
received = get_dest_received_dir(dest, source)
mismatches, missing = verify_transfer(source, received)
assert not mismatches and not missing, \
f"-m --delete wiped source mirrors: missing={missing} mismatches={mismatches}"
+14
View File
@@ -24,7 +24,9 @@ static void test_stop_after_parse_invalid() {
EXPECT_FALSE(stop_parse_after_minutes("", &minutes));
EXPECT_FALSE(stop_parse_after_minutes("5x", &minutes));
EXPECT_FALSE(stop_parse_after_minutes("1.5", &minutes));
EXPECT_FALSE(stop_parse_after_minutes(" 5", &minutes));
EXPECT_FALSE(stop_parse_after_minutes(" 5 ", &minutes));
EXPECT_FALSE(stop_parse_after_minutes("+5", &minutes));
EXPECT_FALSE(stop_parse_after_minutes("2147483648", &minutes));
EXPECT_FALSE(stop_parse_after_minutes(NULL, &minutes));
}
@@ -87,6 +89,12 @@ static void test_stop_at_parse_invalid() {
EXPECT_FALSE(stop_parse_at_time("now+5x", now, &deadline));
EXPECT_FALSE(stop_parse_at_time("now-5m", now, &deadline));
EXPECT_FALSE(stop_parse_at_time("now+1w", now, &deadline));
EXPECT_FALSE(stop_parse_at_time("now+ 5s", now, &deadline));
EXPECT_FALSE(stop_parse_at_time("now++5s", now, &deadline));
/* Signed overflow of the destination deadline must be rejected, not wrap. */
EXPECT_FALSE(stop_parse_at_time("now+9223372036854775807s", now, &deadline));
/* 10^15 days is well beyond LONG_MAX/86400, so the amount itself is rejected. */
EXPECT_FALSE(stop_parse_at_time("now+1000000000000000d", now, &deadline));
EXPECT_FALSE(stop_parse_at_time("abc", now, &deadline));
EXPECT_FALSE(stop_parse_at_time("", now, &deadline));
EXPECT_FALSE(stop_parse_at_time(NULL, now, &deadline));
@@ -113,6 +121,12 @@ static void test_stop_deadline_latency() {
EXPECT_FALSE(no_after.has_wall);
EXPECT_FALSE(stop_condition_reached(&no_after));
/* An invalid (non-positive) after_minutes never arms the monotonic half. */
StopCondition zero_after = stop_condition_make(true, 0, false, 0, now);
EXPECT_FALSE(zero_after.has_monotonic);
StopCondition neg_after = stop_condition_make(true, -5, false, 0, now);
EXPECT_FALSE(neg_after.has_monotonic);
/* --stop-at: a wall-clock deadline in the past/now is reached; one in the
future is not, and it stays independent of the monotonic half. */
StopCondition wall_future = stop_condition_make(false, 0, true, time(NULL) + 3600, now);