feat(p6-stop): --stop-after/--stop-at deadline transfer stop
This commit is contained in:
@@ -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 <errno.h>
|
||||
#include <limits.h>
|
||||
#include <time.h>
|
||||
#include <signal.h>
|
||||
#include <stdbool.h>
|
||||
#include <stddef.h>
|
||||
@@ -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,
|
||||
|
||||
+94
-49
@@ -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);
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
#include "hardlink.h"
|
||||
#include "protocol.h"
|
||||
#include "queue.h"
|
||||
#include "stop_condition.h"
|
||||
#include <dirent.h>
|
||||
#include <stdbool.h>
|
||||
#include <stdatomic.h>
|
||||
@@ -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 {
|
||||
|
||||
@@ -184,6 +184,11 @@ void print_usage(void) {
|
||||
printf(" --timeout <sec> I/O timeout in seconds (default: 30)\n");
|
||||
printf(" -T <sec> Alias for --timeout\n");
|
||||
printf(" --contimeout <sec> 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 <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");
|
||||
|
||||
Reference in New Issue
Block a user