From 24d144824694b1818d3345f5047f19d2b5421cef Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 10 Sep 2026 16:27:49 +0200 Subject: [PATCH 1/4] feat(p6-stop): --stop-after/--stop-at deadline transfer stop --- src/client/client_cli.c | 43 +++++++++++ src/client/client_send.c | 143 ++++++++++++++++++++++------------ src/client/scanner.c | 11 +++ src/client/scanner.h | 8 ++ src/client/usage.c | 5 ++ src/shared/config.c | 3 + src/shared/config.h | 13 ++++ src/shared/multiprocessing.h | 4 + src/shared/stop_condition.c | 144 +++++++++++++++++++++++++++++++++++ src/shared/stop_condition.h | 46 +++++++++++ 10 files changed, 371 insertions(+), 49 deletions(-) create mode 100644 src/shared/stop_condition.c create mode 100644 src/shared/stop_condition.h diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 4b14d88..87a1af2 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -11,12 +11,14 @@ #include "identity.h" #include "log.h" #include "protocol.h" +#include "stop_condition.h" #include "transport_tcp.h" #include "transport_tls.h" #include "usage.h" #include "utils.h" #include #include +#include #include #include #include @@ -847,6 +849,47 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, return -1; continue; } + /* --stop-after/--stop-at are client-only sender-side stop deadlines. They + * are parsed by stop_condition (so the unit tests exercise the same validate + * that production uses) and never serialized into the config frame. */ + if (strncmp(argv[i], "--stop-after=", 13) == 0) { + if (!stop_parse_after_minutes(argv[i] + 13, &config->stop_after_mins)) { + log_message(LOG_LEVEL_ERROR, "--stop-after must be a positive number of minutes"); + return -1; + } + continue; + } + if (strcmp(argv[i], "--stop-after") == 0) { + if (i + 1 >= argc) { + log_message(LOG_LEVEL_ERROR, "missing argument for --stop-after"); + return -1; + } + if (!stop_parse_after_minutes(argv[++i], &config->stop_after_mins)) { + log_message(LOG_LEVEL_ERROR, "--stop-after must be a positive number of minutes"); + return -1; + } + continue; + } + if (strncmp(argv[i], "--stop-at=", 10) == 0) { + if (!stop_parse_at_time(argv[i] + 10, time(NULL), &config->stop_at)) { + log_message(LOG_LEVEL_ERROR, "--stop-at must be HH:MM[:SS] or now+N[smhd]"); + return -1; + } + config->stop_at_set = true; + continue; + } + if (strcmp(argv[i], "--stop-at") == 0) { + if (i + 1 >= argc) { + log_message(LOG_LEVEL_ERROR, "missing argument for --stop-at"); + return -1; + } + if (!stop_parse_at_time(argv[++i], time(NULL), &config->stop_at)) { + log_message(LOG_LEVEL_ERROR, "--stop-at must be HH:MM[:SS] or now+N[smhd]"); + return -1; + } + config->stop_at_set = true; + continue; + } const char* threads_prefix = "--compress-threads="; if (strncmp(argv[i], threads_prefix, strlen(threads_prefix)) == 0) { if (set_compression_threads_option(&config->compression_threads, diff --git a/src/client/client_send.c b/src/client/client_send.c index e909810..edb1e3b 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -17,6 +17,7 @@ #include "protocol.h" #include "queue.h" #include "scanner.h" +#include "stop_condition.h" #include "transport_tcp.h" #include "transport_ssh.h" #include "transport_tls.h" @@ -1396,6 +1397,15 @@ static int send_chunks_multithreaded(void* pipeline_context) { } while (true) { + /* Phase 6: stop-elegantly at the next chunk boundary once the --stop-after + / --stop-at deadline has passed. Everything already sent is finalized by + the completion tail below; the run still returns success. */ + if (stop_condition_reached(&context->stop_condition)) { + log_info_message(LOG_INFO_MISC, + "Stop deadline reached; stopping transfer at the next chunk boundary"); + pipeline_cancel(context); + break; + } Chunk* current_chunk = queue_dequeue_multithreaded( context->queue_loader, &context->mutex_loader, &context->condition_not_empty_loader, &context->condition_not_full_loader, &context->loader_done); @@ -1407,55 +1417,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { protocol_session_unbind(); return thrd_error; } - 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; - mtx_lock(&context->mutex_scanner); - empty_io = context->scan_had_io_error && context->manifest && context->manifest->size == 0; - mtx_unlock(&context->mutex_scanner); - if (empty_io) { - 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 (send_delete_manifest(client->file_descriptor, context->manifest, - context->excluded_paths, context->missing_args) != 0) - goto send_fail; - } else if (context->config->delete_missing_args && !context->early_delete) { - /* --delete-missing-args without --delete: no keep-set is built, but the - exact-delete paths still ride the same manifest frame (commit once the - transfer succeeded). */ - if (send_delete_manifest(client->file_descriptor, NULL, NULL, context->missing_args) != 0) - goto send_fail; - } - bool ok = finalize_transfer(client, context->config, context->remove_source_files); - if (!ok && context->config->use_delete) - log_message(LOG_LEVEL_ERROR, - "server reported a deletion failure (--delete); see the server log for the " - "reason (a --max-delete limit that the run would exceed deletes nothing)"); - if (ok) - remove_transferred_sources(context->config, context->remove_source_files); - mtx_lock(&context->mutex_progress); - int total_files = context->total_files; - unsigned long long total_bytes = context->total_bytes; - mtx_unlock(&context->mutex_progress); - if (context->config->stats) - fprintf(stderr, "Stats: %d files, %.1f MB\n", total_files, total_bytes / 1048576.0); - log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files, - total_bytes / 1048576.0); - disconnect_transfer_client(client); - mark_sender_done(context); - protocol_session_unbind(); - return ok ? thrd_success : thrd_error; - - send_fail: - pipeline_cancel(context); - disconnect_transfer_client(client); - mark_sender_done(context); - protocol_session_unbind(); - return thrd_error; + break; } if (send_chunk_with_removal(client, current_chunk, context->config, context->remove_source_files) != 0) { @@ -1482,6 +1444,59 @@ static int send_chunks_multithreaded(void* pipeline_context) { mtx_unlock(&context->mutex_progress); chunk_destroy(current_chunk); } + + /* 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) { + /* 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; + mtx_lock(&context->mutex_scanner); + empty_io = context->scan_had_io_error && context->manifest && context->manifest->size == 0; + mtx_unlock(&context->mutex_scanner); + if (empty_io) { + 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 (send_delete_manifest(client->file_descriptor, context->manifest, context->excluded_paths, + context->missing_args) != 0) + goto send_fail; + } else if (context->config->delete_missing_args && !context->early_delete) { + /* --delete-missing-args without --delete: no keep-set is built, but the + exact-delete paths still ride the same manifest frame (commit once the + transfer succeeded). */ + if (send_delete_manifest(client->file_descriptor, NULL, NULL, context->missing_args) != 0) + goto send_fail; + } + bool ok = finalize_transfer(client, context->config, context->remove_source_files); + if (!ok && context->config->use_delete) + log_message(LOG_LEVEL_ERROR, + "server reported a deletion failure (--delete); see the server log for the " + "reason (a --max-delete limit that the run would exceed deletes nothing)"); + if (ok) + remove_transferred_sources(context->config, context->remove_source_files); + mtx_lock(&context->mutex_progress); + int total_files = context->total_files; + unsigned long long total_bytes = context->total_bytes; + mtx_unlock(&context->mutex_progress); + if (context->config->stats) + fprintf(stderr, "Stats: %d files, %.1f MB\n", total_files, total_bytes / 1048576.0); + log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files, + total_bytes / 1048576.0); + disconnect_transfer_client(client); + mark_sender_done(context); + protocol_session_unbind(); + return ok ? thrd_success : thrd_error; + +send_fail: + pipeline_cancel(context); + disconnect_transfer_client(client); + mark_sender_done(context); + protocol_session_unbind(); + return thrd_error; } /* Scan thread of the -m pipeline. --dirs disables recursive traversal (the @@ -1496,6 +1511,7 @@ static int scan_directory_multithreaded(void* pipeline_context) { protocol_session_unbind(); return thrd_error; } + prepared.options.stop_condition = &context->stop_condition; /* The keep-set manifest for the late modes is built from this data pass, so the parallel scanner records the protected excluded prefixes here. The early modes already transmitted the pre-scan keep-set and its protected @@ -1790,6 +1806,17 @@ int send_files(Config* config) { if (!manifest) goto send_fail; } + /* Phase 6: compute the client-only stop deadline once at transfer start. The + early-delete pre-scan above deliberately ignores it so the keep-set (and + its committed deletion) is always complete and correct. */ + struct timespec now_mono; + if (clock_gettime(CLOCK_MONOTONIC, &now_mono) != 0) { + now_mono.tv_sec = 0; + now_mono.tv_nsec = 0; + } + StopCondition stop = stop_condition_make(config->stop_after_mins > 0, config->stop_after_mins, + config->stop_at_set, config->stop_at, now_mono); + prepared.options.stop_condition = &stop; scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options); if (!scanner) goto send_fail; @@ -1800,6 +1827,16 @@ int send_files(Config* config) { time_t last_progress = 0; time_t start = time(NULL); 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. */ + 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"); + break; + } unsigned long long chunk_bytes = 0; for (int i = 0; i < current_chunk->element_count; i++) { chunk_bytes += current_chunk->items[i]->data->size; @@ -1987,6 +2024,14 @@ 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 */ + struct timespec now_mono; + if (clock_gettime(CLOCK_MONOTONIC, &now_mono) != 0) { + now_mono.tv_sec = 0; + now_mono.tv_nsec = 0; + } + context->stop_condition = + stop_condition_make(config->stop_after_mins > 0, config->stop_after_mins, config->stop_at_set, + config->stop_at, now_mono); bool collect_excluded = config->use_delete && !config->delete_excluded; if (config->use_delete) { context->manifest = array_list_create(free); diff --git a/src/client/scanner.c b/src/client/scanner.c index af0da24..f144fba 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -487,6 +487,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo scanner->relative_mode = options->relative && options->file_list != NULL; scanner->hardlinks = options->hardlinks; scanner->prune_empty_dirs = options->prune_empty_dirs; + scanner->stop_condition = options->stop_condition; scanner->dirs_root_emitted = false; scanner->list_index = 0; scanner->dirs_batch = NULL; @@ -843,6 +844,12 @@ static Chunk* dirs_flush_batch(DirectoryScanner* scanner) { static Chunk* directory_scanner_next_dirs(DirectoryScanner* scanner) { while (scanner->dirs_batch == NULL || scanner->dirs_batch_size <= scanner->chunk_size) { + if (scanner->stop_condition && stop_condition_reached(scanner->stop_condition)) { + Chunk* leftover = dirs_flush_batch(scanner); + if (leftover) + chunk_destroy(leftover); + return NULL; + } if (!scanner->dirs_batch) { scanner->dirs_batch = array_list_create(file_destroy); if (!scanner->dirs_batch) { @@ -886,6 +893,10 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { unsigned long long chunk_data_size = 0; while (1) { + if (scanner->stop_condition && stop_condition_reached(scanner->stop_condition)) { + array_list_delete(chunk_data); + return NULL; + } if (scanner->current_dir == NULL) { int ret = open_next_directory(scanner); if (ret == 0) diff --git a/src/client/scanner.h b/src/client/scanner.h index c306c5f..632571f 100644 --- a/src/client/scanner.h +++ b/src/client/scanner.h @@ -7,6 +7,7 @@ #include "hardlink.h" #include "protocol.h" #include "queue.h" +#include "stop_condition.h" #include #include #include @@ -92,6 +93,11 @@ typedef struct { * read-only here; the parallel scanner passes it unchanged to every worker so * one table detects every group across all subdirectories. */ HardLinkTable* hardlinks; + /* Phase 6: optional sender stop deadline. When non-NULL the scanner checks + * it at natural loop boundaries and stops emitting chunks once reached + * (without marking the scan as failed), so a busy scan itself stops early. + * Client-only, never serialized to the wire. */ + const StopCondition* stop_condition; } ScannerOptions; /* Internal per-scanner filter state. FilterNode chains represent the ordered @@ -165,6 +171,8 @@ typedef struct { /* --hard-links (-H): shared link-group detection table (see ScannerOptions). NULL when -H is off. */ HardLinkTable* hardlinks; + /* Phase 6: sender stop deadline (from ScannerOptions). */ + const StopCondition* stop_condition; } DirectoryScanner; typedef struct { diff --git a/src/client/usage.c b/src/client/usage.c index 24e9878..137b9d8 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -184,6 +184,11 @@ void print_usage(void) { printf(" --timeout I/O timeout in seconds (default: 30)\n"); printf(" -T Alias for --timeout\n"); printf(" --contimeout Connection timeout in seconds (default: 10)\n"); + printf(" --stop-after=MINS Stop the transfer after MINS minutes (a positive\n"); + 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(" --address 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"); diff --git a/src/shared/config.c b/src/shared/config.c index 0b8c2d3..bfb5e95 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -176,6 +176,9 @@ static void config_set_defaults(Config* config) { config->use_xattrs = false; config->fake_super = false; config->trust_sender = false; + config->stop_after_mins = 0; + config->stop_at = 0; + config->stop_at_set = false; } static bool valid_wire_bool(int value) { diff --git a/src/shared/config.h b/src/shared/config.h index 2be5acb..c3ce90e 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -6,6 +6,7 @@ #include #include #include +#include typedef enum { TRANSPORT_TCP, TRANSPORT_SSH } TransportType; @@ -446,6 +447,18 @@ typedef struct Config { * authorized root (see the phase-5 notes in RSYNC_COMPAT.md). Off by * default; only relaxes validation when explicitly requested. */ bool trust_sender; + + // Phase 6: --stop-after / --stop-at + /* Client-only sender-side transfer stop deadlines. --stop-after=MINS stops + * the transfer after a number of elapsed minutes (checked against + * CLOCK_MONOTONIC so clock changes do not skew it); --stop-at=TIME stops at + * an absolute wall-clock time (HH:MM, HH:MM:SS, or now+N[smhd]). At the + * deadline the run stops elegantly at the next chunk/file boundary and the + * completion tail still runs (exit 0). Both are LOCAL to the sending + * process and are NEVER serialized into the config frame. */ + int stop_after_mins; /* --stop-after=MINS minutes; 0 when unset */ + time_t stop_at; /* --stop-at=... absolute wall-clock deadline */ + bool stop_at_set; /* true when --stop-at was given */ } Config; /* Phase 5 (remote-option wave): 2.13.0 -> 2.14.0. diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 3d33a10..c1ff21c 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -10,6 +10,7 @@ #include "protocol.h" #include "queue.h" #include "receiver.h" +#include "stop_condition.h" #include typedef struct { @@ -56,6 +57,9 @@ typedef struct { bool sender_done; atomic_bool cancelled; ProtocolSession allocation_session; + /* 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; } PipelineContextSender; typedef struct PipelineContextReceiver { diff --git a/src/shared/stop_condition.c b/src/shared/stop_condition.c new file mode 100644 index 0000000..54e87f2 --- /dev/null +++ b/src/shared/stop_condition.c @@ -0,0 +1,144 @@ +#include "stop_condition.h" +#include +#include +#include +#include + +/* Parse a strictly positive decimal integer. Accepts leading '+' but not a + * leading '-', surrounding whitespace, fractional parts 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') + return false; + if (v <= 0 || v > INT_MAX) + return false; + *out = v; + return true; +} + +bool stop_parse_after_minutes(const char* value, int* out_minutes) { + if (!out_minutes) + return false; + long minutes = 0; + if (!parse_positive_minutes(value, &minutes)) + return false; + *out_minutes = (int)minutes; + return true; +} + +/* Two consecutive ASCII digits -> 0..99. */ +static bool parse_two_digits(const char* s, int* out) { + if (s[0] < '0' || s[0] > '9' || s[1] < '0' || s[1] > '9') + return false; + *out = (s[0] - '0') * 10 + (s[1] - '0'); + return true; +} + +bool stop_parse_at_time(const char* value, time_t now, time_t* out_deadline) { + if (!value || !out_deadline) + return false; + + /* 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') + return false; + char* end = NULL; + errno = 0; + long amount = strtol(p, &end, 10); + if (errno != 0 || end == p || amount < 0) + return false; + long unit_seconds; + switch (*end) { + case 's': + unit_seconds = 1; + break; + case 'm': + unit_seconds = 60; + break; + case 'h': + unit_seconds = 3600; + break; + case 'd': + unit_seconds = 86400; + break; + default: + return false; + } + if (end[1] != '\0') + return false; + if (amount > LONG_MAX / unit_seconds) + return false; + long long delta = (long long)amount * unit_seconds; + *out_deadline = now + (time_t)delta; + return true; + } + + /* HH:MM or HH:MM:SS on the current local day. */ + size_t len = strlen(value); + if (len != 5 && len != 8) + return false; + if (value[2] != ':' || (len == 8 && value[5] != ':')) + return false; + int hh, mm, ss = 0; + if (!parse_two_digits(value, &hh) || !parse_two_digits(value + 3, &mm)) + return false; + if (len == 8 && !parse_two_digits(value + 6, &ss)) + return false; + if (hh > 23 || mm > 59 || ss > 59) + return false; + + struct tm today; + if (!localtime_r(&now, &today)) + return false; + today.tm_hour = hh; + today.tm_min = mm; + today.tm_sec = ss; + today.tm_isdst = -1; + time_t deadline = mktime(&today); + if (deadline == (time_t)-1) + return false; + *out_deadline = deadline; + return true; +} + +StopCondition stop_condition_make(bool has_after, int after_minutes, bool has_at, time_t at_time, + struct timespec now_mono) { + StopCondition condition; + condition.has_monotonic = false; + condition.monotonic_deadline.tv_sec = 0; + condition.monotonic_deadline.tv_nsec = 0; + condition.has_wall = false; + condition.wall_deadline = 0; + if (has_after && after_minutes > 0) { + condition.has_monotonic = true; + condition.monotonic_deadline.tv_sec = now_mono.tv_sec + (time_t)after_minutes * 60; + condition.monotonic_deadline.tv_nsec = now_mono.tv_nsec; + } + if (has_at) { + condition.has_wall = true; + condition.wall_deadline = at_time; + } + return condition; +} + +bool stop_condition_reached(const StopCondition* condition) { + if (!condition) + return false; + if (condition->has_wall && time(NULL) >= condition->wall_deadline) + return true; + if (condition->has_monotonic) { + struct timespec now; + if (clock_gettime(CLOCK_MONOTONIC, &now) != 0) + return false; + if (now.tv_sec > condition->monotonic_deadline.tv_sec || + (now.tv_sec == condition->monotonic_deadline.tv_sec && + now.tv_nsec >= condition->monotonic_deadline.tv_nsec)) + return true; + } + return false; +} \ No newline at end of file diff --git a/src/shared/stop_condition.h b/src/shared/stop_condition.h new file mode 100644 index 0000000..f1ceb7a --- /dev/null +++ b/src/shared/stop_condition.h @@ -0,0 +1,46 @@ +#ifndef STOP_CONDITION_H +#define STOP_CONDITION_H + +#include +#include + +/* Client-only transfer stop conditions (--stop-after=MINS / --stop-at=TIME). + * Both are local sender-side deadlines: they are never serialized into the + * config frame and never bump PROTOCOL_VERSION. A transfer checks the + * condition at natural chunk/file boundaries and, once reached, stops + * elegantly (everything already sent is finalized normally, exit 0). + * + * A condition combines an optional CLOCK_MONOTONIC instant (the relative + * --stop-after duration, immune to wall-clock changes) with an optional + * wall-clock instant (the absolute --stop-at form). Either one being reached + * ends the transfer. */ +typedef struct StopCondition { + bool has_monotonic; + struct timespec monotonic_deadline; + bool has_wall; + time_t wall_deadline; +} StopCondition; + +/* Parse --stop-after=MINS: a positive integer count of minutes. Zero, + * negative, empty and non-numeric values are rejected. Returns true when + * accepted and stores the value in *out_minutes. */ +bool stop_parse_after_minutes(const char* value, int* out_minutes); + +/* Parse --stop-at=TIME. Accepted forms are HH:MM, HH:MM:SS and + * now+N[smhd] (seconds/minutes/hours/days from now). The absolute forms are + * resolved against `now` (local wall clock) and written to *out_deadline; a + * time already in the past yields a deadline <= now ("stop immediately"). + * Returns false on any malformed value. */ +bool stop_parse_at_time(const char* value, time_t now, time_t* out_deadline); + +/* Build the runtime condition at transfer start. after_minutes is the + * relative --stop-after duration (<= 0 disables it); at_time is the absolute + * --stop-at deadline (only consulted when has_at is true); now_mono is the + * CLOCK_MONOTONIC reading at start. */ +StopCondition stop_condition_make(bool has_after, int after_minutes, bool has_at, time_t at_time, + struct timespec now_mono); + +/* True once either deadline has passed (wall clock first, then monotonic). */ +bool stop_condition_reached(const StopCondition* condition); + +#endif From 37cff965372da283c17d456292e30189d784e680 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 10 Sep 2026 16:27:49 +0200 Subject: [PATCH 2/4] test(p6-stop): unit and integration tests for stop deadlines --- tests/integration/test_stop.py | 132 ++++++++++++++++++++++++++++++++ tests/runner.c | 2 + tests/test_stop.c | 136 +++++++++++++++++++++++++++++++++ tests/test_stop.h | 6 ++ 4 files changed, 276 insertions(+) create mode 100644 tests/integration/test_stop.py create mode 100644 tests/test_stop.c create mode 100644 tests/test_stop.h diff --git a/tests/integration/test_stop.py b/tests/integration/test_stop.py new file mode 100644 index 0000000..d10f03c --- /dev/null +++ b/tests/integration/test_stop.py @@ -0,0 +1,132 @@ +"""--stop-after / --stop-at deadline-stop integration tests. + +These cover the client-only sender stop conditions: --stop-after=MINS stops +after N elapsed minutes, --stop-at=HH:MM[:SS] or now+N[smhd] stops at an +absolute (or relative) wall-clock time. A reached deadline ends the transfer +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 os +import shutil +import time + +import pytest + +from common import ( + TEST_DATA_DIR, + run_client, + clean_dir, + get_dest_received_dir, + verify_transfer, +) + + +def _make(self_prefix): + source = os.path.join(TEST_DATA_DIR, f"stop_{self_prefix}_src") + dest = os.path.join(TEST_DATA_DIR, f"stop_{self_prefix}_dst") + clean_dir(source) + shutil.rmtree(dest, ignore_errors=True) + os.makedirs(dest) + return source, dest + + +def _received_files(root): + """All files under `root`, relative paths.""" + if not os.path.isdir(root): + return [] + return [ + os.path.relpath(os.path.join(dirpath, name), root) + for dirpath, _, names in os.walk(root) + for name in names + ] + + +def _seed_source(source): + """Create a handful of regular and nested files.""" + files = { + "small.txt": b"hello world\n", + "medium.txt": b"the quick brown fox jumps over the lazy dog\n" * 400, + "binary.bin": bytes(range(256)) * 100, + "nested/deep.txt": b"deeply nested file\n", + "nested/another.txt": b"another nested file\n" * 40, + } + for rel, content in files.items(): + path = os.path.join(source, rel) + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, "wb") as fh: + fh.write(content) + + +class TestStopAfter: + @pytest.mark.ci + def test_stop_after_within_window(self, shared_server): + """A --stop-after set well past the run's duration lets it finish fully.""" + source, dest = _make("within") + _seed_source(source) + result, _ = run_client(source, dest, flags=["--stop-after=60"], + port=shared_server.port) + assert result.returncode == 0, \ + f"--stop-after full run 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"full transfer mismatch: missing={missing} mismatches={mismatches}" + + @pytest.mark.ci + def test_stop_after_rejects_nonpositive(self, shared_server): + """0 and negative minutes are invalid (must be a positive integer).""" + source, dest = _make("reject") + _seed_source(source) + for bad in ("0", "-1"): + result, _ = run_client(source, dest, flags=[f"--stop-after={bad}"], + port=shared_server.port) + assert result.returncode != 0, f"--stop-after={bad} should be rejected" + + +class TestStopAt: + @pytest.mark.ci + def test_stop_at_past(self, shared_server): + """A --stop-at already in the past stops the transfer immediately but + cleanly (exit 0, nothing transferred).""" + source, dest = _make("past") + _seed_source(source) + now = time.localtime() + if now.tm_hour * 60 + now.tm_min >= 1: + 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" + result, _ = run_client(source, dest, flags=[f"--stop-at={stop_value}"], + port=shared_server.port) + assert result.returncode == 0, \ + f"--stop-at past run failed (rc {result.returncode}): " \ + f"{(result.stderr or result.stdout)[:400]}" + received = get_dest_received_dir(dest, source) + assert _received_files(received) == [], \ + f"expected nothing transferred, got {_received_files(received)}" + + @pytest.mark.ci + def test_stop_at_now_plus_stops_immediately(self, shared_server): + """now+0s resolves to the current instant, so the transfer stops at once.""" + source, dest = _make("nowplus") + _seed_source(source) + result, _ = run_client(source, dest, flags=["--stop-at=now+0s"], + port=shared_server.port) + assert result.returncode == 0, \ + f"--stop-at=now+0s should stop cleanly: " \ + f"{(result.stderr or result.stdout)[:400]}" + received = get_dest_received_dir(dest, source) + assert _received_files(received) == [], \ + f"expected nothing transferred, got {_received_files(received)}" + + @pytest.mark.ci + def test_stop_rejects_garbage(self, shared_server): + """Malformed --stop-at/--stop-after values are rejected up front.""" + source, dest = _make("garbage") + _seed_source(source) + for flag in ("--stop-after=abc", "--stop-at=12:99", "--stop-at=12", + "--stop-at=now+5x", "--stop-at=now-5s"): + result, _ = run_client(source, dest, flags=[flag], + port=shared_server.port) + assert result.returncode != 0, f"{flag} should be rejected" \ No newline at end of file diff --git a/tests/runner.c b/tests/runner.c index 5cf2d54..f376aa4 100644 --- a/tests/runner.c +++ b/tests/runner.c @@ -27,6 +27,7 @@ #include "test_server_cli.h" #include "test_shared_utils.h" #include "test_stress.h" +#include "test_stop.h" #include "test_transport_tcp.h" #include "test_transport_ssh.h" #include "test_transport_tls.h" @@ -67,6 +68,7 @@ int main() { RUN_TEST(test_log); RUN_TEST(test_robustness); RUN_TEST(test_stress); + RUN_TEST(test_stop); RUN_TEST(test_property); RUN_TEST(test_transport_tcp); RUN_TEST(test_transport_ssh); diff --git a/tests/test_stop.c b/tests/test_stop.c new file mode 100644 index 0000000..1e20239 --- /dev/null +++ b/tests/test_stop.c @@ -0,0 +1,136 @@ +#include "test_stop.h" +#include "stop_condition.h" +#include "test_utils.h" +#include +#include + +static void test_stop_after_parse_valid() { + int minutes = 0; + EXPECT_TRUE(stop_parse_after_minutes("5", &minutes)); + EXPECT_EQ_INT(minutes, 5); + EXPECT_TRUE(stop_parse_after_minutes("1", &minutes)); + EXPECT_EQ_INT(minutes, 1); + EXPECT_TRUE(stop_parse_after_minutes("1440", &minutes)); + EXPECT_EQ_INT(minutes, 1440); + EXPECT_TRUE(stop_parse_after_minutes("2147483647", &minutes)); + EXPECT_EQ_INT(minutes, INT_MAX); +} + +static void test_stop_after_parse_invalid() { + int minutes = 0; + EXPECT_FALSE(stop_parse_after_minutes("0", &minutes)); + EXPECT_FALSE(stop_parse_after_minutes("-1", &minutes)); + EXPECT_FALSE(stop_parse_after_minutes("abc", &minutes)); + 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("2147483648", &minutes)); + EXPECT_FALSE(stop_parse_after_minutes(NULL, &minutes)); +} + +static void test_stop_at_parse_hhmm() { + time_t now = 1700000000; + time_t deadline = 0; + + EXPECT_TRUE(stop_parse_at_time("12:30", now, &deadline)); + struct tm t; + EXPECT_NOT_NULL(localtime_r(&deadline, &t)); + EXPECT_EQ_INT(t.tm_hour, 12); + EXPECT_EQ_INT(t.tm_min, 30); + EXPECT_EQ_INT(t.tm_sec, 0); + + EXPECT_TRUE(stop_parse_at_time("12:30:59", now, &deadline)); + EXPECT_NOT_NULL(localtime_r(&deadline, &t)); + EXPECT_EQ_INT(t.tm_hour, 12); + EXPECT_EQ_INT(t.tm_min, 30); + EXPECT_EQ_INT(t.tm_sec, 59); + + EXPECT_TRUE(stop_parse_at_time("00:00", now, &deadline)); + EXPECT_NOT_NULL(localtime_r(&deadline, &t)); + EXPECT_EQ_INT(t.tm_hour, 0); + EXPECT_EQ_INT(t.tm_min, 0); + EXPECT_EQ_INT(t.tm_sec, 0); +} + +static void test_stop_at_parse_now_plus() { + time_t now = 1700000000; + time_t deadline = 0; + + EXPECT_TRUE(stop_parse_at_time("now+90s", now, &deadline)); + EXPECT_EQ_INT(deadline, now + 90); + EXPECT_TRUE(stop_parse_at_time("now+5m", now, &deadline)); + EXPECT_EQ_INT(deadline, now + 300); + EXPECT_TRUE(stop_parse_at_time("now+2h", now, &deadline)); + EXPECT_EQ_INT(deadline, now + 7200); + EXPECT_TRUE(stop_parse_at_time("now+1d", now, &deadline)); + EXPECT_EQ_INT(deadline, now + 86400); + EXPECT_TRUE(stop_parse_at_time("now+0s", now, &deadline)); + EXPECT_EQ_INT(deadline, now); +} + +static void test_stop_at_parse_invalid() { + time_t now = 1700000000; + time_t deadline = 0; + EXPECT_FALSE(stop_parse_at_time("12", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("12:3", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("1234", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("12:30:5", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("12:30:5x", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("24:00", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("12:60", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("12:30:61", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("12;00", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("now", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("now+", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("now+5", now, &deadline)); + 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("abc", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time("", now, &deadline)); + EXPECT_FALSE(stop_parse_at_time(NULL, now, &deadline)); +} + +static void test_stop_deadline_latency() { + struct timespec now; + EXPECT_EQ_INT(clock_gettime(CLOCK_MONOTONIC, &now), 0); + + StopCondition future = stop_condition_make(true, 60, false, 0, now); + EXPECT_TRUE(future.has_monotonic); + EXPECT_EQ_INT(future.monotonic_deadline.tv_sec, now.tv_sec + 3600); + EXPECT_EQ_INT(future.monotonic_deadline.tv_nsec, now.tv_nsec); + EXPECT_FALSE(future.has_wall); + EXPECT_FALSE(stop_condition_reached(&future)); + + /* Move the 60-minute deadline into the past: the check now reports reached. */ + StopCondition past = stop_condition_make(true, 60, false, 0, now); + past.monotonic_deadline.tv_sec -= 7200; + EXPECT_TRUE(stop_condition_reached(&past)); + + StopCondition no_after = stop_condition_make(false, 0, false, 0, now); + EXPECT_FALSE(no_after.has_monotonic); + EXPECT_FALSE(no_after.has_wall); + EXPECT_FALSE(stop_condition_reached(&no_after)); + + /* --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); + EXPECT_TRUE(wall_future.has_wall); + EXPECT_FALSE(wall_future.has_monotonic); + EXPECT_FALSE(stop_condition_reached(&wall_future)); + + StopCondition wall_past = stop_condition_make(false, 0, true, time(NULL) - 1, now); + EXPECT_TRUE(stop_condition_reached(&wall_past)); + + EXPECT_FALSE(stop_condition_reached(NULL)); +} + +void test_stop(void) { + test_stop_after_parse_valid(); + test_stop_after_parse_invalid(); + test_stop_at_parse_hhmm(); + test_stop_at_parse_now_plus(); + test_stop_at_parse_invalid(); + test_stop_deadline_latency(); +} \ No newline at end of file diff --git a/tests/test_stop.h b/tests/test_stop.h new file mode 100644 index 0000000..3c5b85e --- /dev/null +++ b/tests/test_stop.h @@ -0,0 +1,6 @@ +#ifndef TEST_STOP_H +#define TEST_STOP_H + +void test_stop(void); + +#endif \ No newline at end of file From ac7e9e3bc1c170015c749e0ea1ec824a3736f852 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 10 Sep 2026 17:39:34 +0200 Subject: [PATCH 3/4] fix(p6-stop): block delete-manifest on early stop; sync -m manifest access; overflow guard --- src/client/client_send.c | 78 ++++++++++++++-------- src/client/usage.c | 4 +- src/shared/multiprocessing.h | 7 ++ src/shared/stop_condition.c | 29 +++++--- tests/integration/test_stop.py | 117 +++++++++++++++++++++++++++++++-- tests/test_stop.c | 14 ++++ 6 files changed, 210 insertions(+), 39 deletions(-) diff --git a/src/client/client_send.c b/src/client/client_send.c index edb1e3b..fe6dcdf 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -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); diff --git a/src/client/usage.c b/src/client/usage.c index 137b9d8..0b6e34a 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -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 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"); diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index c1ff21c..31d545b 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -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 { diff --git a/src/shared/stop_condition.c b/src/shared/stop_condition.c index 54e87f2..b49053c 100644 --- a/src/shared/stop_condition.c +++ b/src/shared/stop_condition.c @@ -4,16 +4,22 @@ #include #include -/* 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; } diff --git a/tests/integration/test_stop.py b/tests/integration/test_stop.py index d10f03c..148fe87 100644 --- a/tests/integration/test_stop.py +++ b/tests/integration/test_stop.py @@ -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, \ @@ -129,4 +159,83 @@ class TestStopAt: "--stop-at=now+5x", "--stop-at=now-5s"): result, _ = run_client(source, dest, flags=[flag], port=shared_server.port) - assert result.returncode != 0, f"{flag} should be rejected" \ No newline at end of file + 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}" \ No newline at end of file diff --git a/tests/test_stop.c b/tests/test_stop.c index 1e20239..4fc71bf 100644 --- a/tests/test_stop.c +++ b/tests/test_stop.c @@ -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); From 09a07179c92c36f7a1ab48beb69e9aa806d2ff76 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 10 Sep 2026 18:16:35 +0200 Subject: [PATCH 4/4] fix(p6-stop): init -m scan_stopped_early; gate delete warning; stabilize partial-stop test --- src/client/client_send.c | 18 ++++++++++++------ src/shared/multiprocessing.c | 1 + tests/integration/test_stop.py | 2 +- 3 files changed, 14 insertions(+), 7 deletions(-) diff --git a/src/client/client_send.c b/src/client/client_send.c index fe6dcdf..1ab04dd 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -1456,9 +1456,12 @@ static int send_chunks_multithreaded(void* pipeline_context) { 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"); + if (context->config->use_delete || context->config->delete_missing_args) + log_message(LOG_LEVEL_WARNING, + "transfer stopped early (stop deadline); skipping --delete keep-set so " + "unscanned source mirrors are not deleted"); + else + log_message(LOG_LEVEL_WARNING, "transfer stopped early (stop deadline)"); } 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). */ @@ -1911,9 +1914,12 @@ int send_files(Config* config) { 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"); + if (config->use_delete || config->delete_missing_args) + log_message(LOG_LEVEL_WARNING, + "transfer stopped early (stop deadline); skipping --delete keep-set so " + "unscanned source mirrors are not deleted"); + else + log_message(LOG_LEVEL_WARNING, "transfer stopped early (stop deadline)"); } else { if (had_scan_io && manifest && manifest->size == 0) { /* A scan that hit an I/O error and produced no keep entries is ambiguous; diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index a77ffb6..d45b569 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -31,6 +31,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->scan_had_io_error = false; context->remove_source_files = NULL; context->early_delete = false; + context->scan_stopped_early = false; context->total_files = 0; context->progress_bytes = 0; context->total_bytes = 0; diff --git a/tests/integration/test_stop.py b/tests/integration/test_stop.py index 148fe87..f2bdc35 100644 --- a/tests/integration/test_stop.py +++ b/tests/integration/test_stop.py @@ -172,7 +172,7 @@ class TestStopPartial: 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"] + flags = ["--chunk-size", "262144", "--bwlimit", "100", "--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}): " \