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