Merge branch 'feat/parity-wirestats' into feat/parity-completion
# Conflicts: # src/server/receiver_pipeline.c # src/shared/config.h # src/shared/protocol.h # tests/integration/test_fault_injection.py # tests/integration/test_preflight.py # tests/test_client_cli.c # tests/test_config.c
This commit is contained in:
+290
-120
@@ -66,9 +66,6 @@ static void log_server_rejection(const char* context) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Forward declaration for progress-reporting thread used in multithreaded send. */
|
||||
static int progress_thread_fn(void* arg);
|
||||
|
||||
static const char* display_bytes(unsigned long long bytes, bool human_readable, char* buffer,
|
||||
size_t buffer_size) {
|
||||
if (human_readable && format_human_size_decimal(bytes, buffer, buffer_size))
|
||||
@@ -86,20 +83,31 @@ static const char* stats_bytes(const Config* config, unsigned long long bytes, c
|
||||
return buffer;
|
||||
}
|
||||
|
||||
/* Print the rsync `--stats` block on stdout. FastSync is a push sender, so a
|
||||
few receiver-only counters (matched data, file-list bytes, deletion count)
|
||||
are not observable and are reported as 0; the labels and layout match rsync
|
||||
3.4.1. Shared by the single-threaded and multithreaded send paths. */
|
||||
/* Print the rsync `--stats` block on stdout. Byte totals use the process-wide
|
||||
wire counters and the receiver-only counters come from the STATUS_STATS frame;
|
||||
the labels, layout and rate/speedup formulas match rsync 3.4.1. Shared by the
|
||||
single-threaded and multithreaded send paths. */
|
||||
static void report_transfer_stats(const Config* config, int total_files,
|
||||
unsigned long long total_bytes, time_t start) {
|
||||
unsigned long long total_bytes, time_t start,
|
||||
const ReceiverStats* recv) {
|
||||
if (!config->stats || config->quiet)
|
||||
return;
|
||||
ReceiverStats none = {0};
|
||||
if (recv == NULL)
|
||||
recv = &none;
|
||||
unsigned long long sent = protocol_bytes_written();
|
||||
unsigned long long received = protocol_bytes_read();
|
||||
/* rsync: bytes_per_sec = (written + read) / (0.5 + (end - start)). */
|
||||
double elapsed = difftime(time(NULL), start);
|
||||
double rate = elapsed > 0.0 ? (double)total_bytes / elapsed : 0.0;
|
||||
double rate = (double)(sent + received) / (0.5 + elapsed);
|
||||
char total_buffer[32];
|
||||
char sent_buffer[32];
|
||||
char recv_buffer[32];
|
||||
char rate_buffer[32] = {0};
|
||||
char human_rate[32] = {0};
|
||||
const char* total = stats_bytes(config, total_bytes, total_buffer, sizeof(total_buffer));
|
||||
const char* sent_s = stats_bytes(config, sent, sent_buffer, sizeof(sent_buffer));
|
||||
const char* recv_s = stats_bytes(config, received, recv_buffer, sizeof(recv_buffer));
|
||||
const char* rate_str = rate_buffer;
|
||||
if (config->human_readable) {
|
||||
if (!format_human_size_decimal((unsigned long long)rate, human_rate, sizeof(human_rate)))
|
||||
@@ -108,23 +116,112 @@ static void report_transfer_stats(const Config* config, int total_files,
|
||||
} else {
|
||||
snprintf(rate_buffer, sizeof(rate_buffer), "%.2f", rate);
|
||||
}
|
||||
double speedup = (sent + received) > 0 ? (double)total_bytes / (double)(sent + received) : 0.0;
|
||||
printf("\n");
|
||||
printf("Number of files: %d\n", total_files);
|
||||
printf("Number of created files: %d\n", total_files);
|
||||
printf("Number of deleted files: 0\n");
|
||||
printf("Number of deleted files: %llu\n", recv->deleted_files);
|
||||
printf("Number of regular files transferred: %d\n", total_files);
|
||||
printf("Total file size: %s bytes\n", total);
|
||||
printf("Total transferred file size: %s bytes\n", total);
|
||||
printf("Literal data: %s bytes\n", total);
|
||||
printf("Matched data: 0 bytes\n");
|
||||
printf("Matched data: %llu bytes\n", recv->matched_data);
|
||||
printf("File list size: 0\n");
|
||||
printf("File list generation time: 0.000 seconds\n");
|
||||
printf("File list transfer time: 0.000 seconds\n");
|
||||
printf("Total bytes sent: %s\n", total);
|
||||
printf("Total bytes received: 0\n");
|
||||
printf("Total bytes sent: %s\n", sent_s);
|
||||
printf("Total bytes received: %s\n", recv_s);
|
||||
printf("\n");
|
||||
printf("sent %s bytes received 0 bytes %s bytes/sec\n", total, rate_str);
|
||||
printf("total size is %s speedup is %.2f\n", total, 1.0);
|
||||
printf("sent %s bytes received %s bytes %s bytes/sec\n", sent_s, recv_s, rate_str);
|
||||
printf("total size is %s speedup is %.2f%s\n", total, speedup,
|
||||
config->dry_run ? " (DRY RUN)" : "");
|
||||
fflush(stdout);
|
||||
}
|
||||
|
||||
/* ---- rsync-style per-file --progress ------------------------------------
|
||||
* rsync prints, for each transferred regular file, the file name followed by a
|
||||
* two-frame progress line: the first at the initial 32 KiB read window (always
|
||||
* 0.00 kB/s / 0:00:00 on a sub-second transfer) and a final 100% frame carrying
|
||||
* `(xfr#N, to-chk=X/Y)`. Rates are wall-clock dependent, so only the final
|
||||
* rate is measured here; the layout matches rsync 3.4.1's progress.c. */
|
||||
#define RSYNC_PROGRESS_IO_WINDOW (32ULL * 1024ULL)
|
||||
|
||||
static const char* delete_display_path(const Config* config, const char* path);
|
||||
|
||||
static bool g_progress_active;
|
||||
static unsigned long long g_progress_xferred;
|
||||
static unsigned long long g_progress_seen;
|
||||
static struct timespec g_progress_file_start;
|
||||
|
||||
static void progress_first_frame(unsigned long long size, char* out, size_t out_size) {
|
||||
char ofs_buf[32];
|
||||
unsigned long long ofs = size < RSYNC_PROGRESS_IO_WINDOW ? size : RSYNC_PROGRESS_IO_WINDOW;
|
||||
if (!format_big_num(ofs, false, ofs_buf, sizeof(ofs_buf)))
|
||||
snprintf(ofs_buf, sizeof(ofs_buf), "%llu", ofs);
|
||||
int pct = size == 0 ? 100 : (ofs == size ? 100 : (int)(100.0 * (double)ofs / (double)size));
|
||||
snprintf(out, out_size, "\r%15s %3d%% %7.2f%s %s%s", ofs_buf, pct, 0.0, "kB/s", " 0:00:00",
|
||||
" ");
|
||||
}
|
||||
|
||||
static void progress_final_frame(unsigned long long size, char* out, size_t out_size) {
|
||||
char ofs_buf[32];
|
||||
char rembuf[32];
|
||||
unsigned long long last_ofs = size < RSYNC_PROGRESS_IO_WINDOW ? size : RSYNC_PROGRESS_IO_WINDOW;
|
||||
if (!format_big_num(size, false, ofs_buf, sizeof(ofs_buf)))
|
||||
snprintf(ofs_buf, sizeof(ofs_buf), "%llu", size);
|
||||
struct timespec now;
|
||||
clock_gettime(CLOCK_MONOTONIC, &now);
|
||||
long long diff_ms = (long long)(now.tv_sec - g_progress_file_start.tv_sec) * 1000 +
|
||||
(now.tv_nsec - g_progress_file_start.tv_nsec) / 1000000;
|
||||
if (diff_ms <= 0)
|
||||
diff_ms = 1;
|
||||
double rate =
|
||||
size > last_ofs ? (double)(size - last_ofs) * 1000.0 / (double)diff_ms / 1024.0 : 0.0;
|
||||
const char* units = "kB/s";
|
||||
if (rate > 1024.0 * 1024.0) {
|
||||
rate /= 1024.0 * 1024.0;
|
||||
units = "GB/s";
|
||||
} else if (rate > 1024.0) {
|
||||
rate /= 1024.0;
|
||||
units = "MB/s";
|
||||
}
|
||||
unsigned long long remain = (unsigned long long)(diff_ms / 1000);
|
||||
snprintf(rembuf, sizeof(rembuf), "%4u:%02u:%02u", (unsigned)(remain / 3600),
|
||||
(unsigned)((remain / 60) % 60), (unsigned)(remain % 60));
|
||||
unsigned long long to_chk =
|
||||
g_progress_seen > g_progress_xferred ? g_progress_seen - g_progress_xferred : 0;
|
||||
snprintf(out, out_size, "\r%15s %3d%% %7.2f%s %s (xfr#%llu, to-chk=%llu/%llu)\n", ofs_buf, 100,
|
||||
rate, units, rembuf, g_progress_xferred, to_chk, g_progress_seen);
|
||||
}
|
||||
|
||||
static void client_progress_begin(const Config* config) {
|
||||
g_progress_active = config->show_progress && !config->quiet;
|
||||
g_progress_xferred = 0;
|
||||
g_progress_seen = 0;
|
||||
if (!g_progress_active)
|
||||
return;
|
||||
printf("sending incremental file list\n");
|
||||
fflush(stdout);
|
||||
}
|
||||
|
||||
/* Emit the name (unless itemize/out-format already did) and the two progress
|
||||
* frames for one transferred regular file. */
|
||||
static void client_progress_file(const Config* config, const File* file) {
|
||||
if (!g_progress_active || file == NULL || !file->data)
|
||||
return;
|
||||
g_progress_seen++;
|
||||
g_progress_xferred++;
|
||||
unsigned long long size = file->data->size;
|
||||
if (!config->itemize_changes && config->out_format == NULL) {
|
||||
const char* name = delete_display_path(config, file_wire_path(file));
|
||||
printf("%s\n", name ? name : "");
|
||||
}
|
||||
clock_gettime(CLOCK_MONOTONIC, &g_progress_file_start);
|
||||
char frame[160];
|
||||
progress_first_frame(size, frame, sizeof(frame));
|
||||
fputs(frame, stdout);
|
||||
progress_final_frame(size, frame, sizeof(frame));
|
||||
fputs(frame, stdout);
|
||||
fflush(stdout);
|
||||
}
|
||||
|
||||
@@ -802,37 +899,96 @@ static void mark_sender_done(PipelineContextSender* context) {
|
||||
mtx_unlock(&context->mutex_progress);
|
||||
}
|
||||
|
||||
/* Read the optional STATUS_STATS record (protocol 2.25.0) that the receiver
|
||||
* sends just before its terminal status when report_stats was negotiated.
|
||||
* Consumes the would-delete path list into `would_delete` (optional). */
|
||||
static bool receive_stats_record(int fd, ReceiverStats* stats, ArrayList* would_delete) {
|
||||
if (!format_stats_receive(fd, stats))
|
||||
return false;
|
||||
int count = 0;
|
||||
if (!receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES)
|
||||
return false;
|
||||
for (int i = 0; i < count; i++) {
|
||||
char* path = receive_wire_str(fd);
|
||||
if (!path)
|
||||
return false;
|
||||
if (would_delete) {
|
||||
char* copy = str_dup(path);
|
||||
free(path);
|
||||
if (!copy || !array_list_add(would_delete, copy)) {
|
||||
free(copy);
|
||||
return false;
|
||||
}
|
||||
} else {
|
||||
free(path);
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/* Strip the transfer-root prefix from a receiver-reported destination-relative
|
||||
* delete path so a `*deleting` line matches rsync's transfer-relative name
|
||||
* (FastSync's destination mirror includes the source's absolute path). */
|
||||
static const char* delete_display_path(const Config* config, const char* path) {
|
||||
if (!config || !path || !config->send_directory)
|
||||
return path;
|
||||
const char* root = config->send_directory;
|
||||
while (*root == '/')
|
||||
root++;
|
||||
const char* rel = path;
|
||||
while (*rel == '/')
|
||||
rel++;
|
||||
size_t root_len = strlen(root);
|
||||
while (root_len > 0 && root[root_len - 1] == '/')
|
||||
root_len--;
|
||||
if (root_len == 0)
|
||||
return rel;
|
||||
if (strncmp(rel, root, root_len) == 0 && (rel[root_len] == '/' || rel[root_len] == '\0'))
|
||||
return rel + root_len + (rel[root_len] == '/' ? 1 : 0);
|
||||
return rel;
|
||||
}
|
||||
|
||||
/* Send the final STATUS_FINISHED frame and await the receiver's verdict.
|
||||
When --remove-source-files is active the receiver acknowledges each data
|
||||
file it processed, in send order: STATUS_NEXT means the file was written,
|
||||
STATUS_OK means the file was skipped/unchanged. Skipped sources are marked
|
||||
so the later removal pass keeps them. */
|
||||
static bool finalize_transfer(Client* client, const Config* config, ArrayList* remove_sources,
|
||||
bool* delete_limit_out) {
|
||||
bool* delete_limit_out, ReceiverStats* stats_out) {
|
||||
if (delete_limit_out)
|
||||
*delete_limit_out = false;
|
||||
if (!send_status(client->file_descriptor, STATUS_FINISHED))
|
||||
return false;
|
||||
if (config->remove_source_files && remove_sources) {
|
||||
/* The receiver emits its optional wire-stats frame (protocol 2.25.0) FIRST,
|
||||
then any per-file --remove-source-files acks, then the terminal status. */
|
||||
Status status;
|
||||
if (!receive_status(client->file_descriptor, &status))
|
||||
return false;
|
||||
if (status == STATUS_STATS) {
|
||||
ReceiverStats scratch;
|
||||
if (!receive_stats_record(client->file_descriptor, stats_out ? stats_out : &scratch, NULL))
|
||||
return false;
|
||||
if (!receive_status(client->file_descriptor, &status))
|
||||
return false;
|
||||
}
|
||||
if (config->remove_source_files && remove_sources && remove_sources->size > 0) {
|
||||
for (int i = 0; i < remove_sources->size; i++) {
|
||||
Status per_file;
|
||||
if (!receive_status(client->file_descriptor, &per_file))
|
||||
if (i > 0 && !receive_status(client->file_descriptor, &status))
|
||||
return false;
|
||||
if (per_file == STATUS_ERROR) {
|
||||
if (status == STATUS_ERROR) {
|
||||
log_server_rejection("Receiver reported a per-file error");
|
||||
return false;
|
||||
}
|
||||
if (per_file == STATUS_OK) {
|
||||
if (status == STATUS_OK) {
|
||||
((SourceFile*)remove_sources->items[i])->skipped = true;
|
||||
} else if (per_file != STATUS_NEXT) {
|
||||
} else if (status != STATUS_NEXT) {
|
||||
log_message(LOG_LEVEL_ERROR, "Unexpected per-file status from receiver");
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if (!receive_status(client->file_descriptor, &status))
|
||||
return false;
|
||||
}
|
||||
Status status;
|
||||
if (!receive_status(client->file_descriptor, &status))
|
||||
return false;
|
||||
/* A capped --max-delete commit is a successful transfer that the client must
|
||||
report with rsync's exit code 25 (not an error). */
|
||||
if (status == STATUS_DELETE_LIMIT) {
|
||||
@@ -1584,13 +1740,6 @@ static int send_dry_run_remote(Config* config) {
|
||||
dry-run reports the same clear diagnostic instead of aborting mid-stream. */
|
||||
if (config_has_basis(config) && !basis_oversize_preflight(config))
|
||||
return 1;
|
||||
/* Would-delete reporting requires a receiver-side read-only extras walk that
|
||||
is not implemented yet; be explicit that --delete is a no-op in dry-run
|
||||
rather than silently ignoring it. */
|
||||
if ((config->use_delete || config->delete_missing_args) && !config->quiet)
|
||||
log_message(LOG_LEVEL_WARNING,
|
||||
"--dry-run: would-delete reporting is not available in this release; nothing is "
|
||||
"deleted");
|
||||
|
||||
/* A live session may follow, so arm graceful abort handling. */
|
||||
client_set_abort_armed(true);
|
||||
@@ -1609,9 +1758,14 @@ static int send_dry_run_remote(Config* config) {
|
||||
protocol_session_bind(&session);
|
||||
|
||||
int ret = 1;
|
||||
time_t dry_start = time(NULL);
|
||||
ReceiverStats dry_stats;
|
||||
memset(&dry_stats, 0, sizeof(dry_stats));
|
||||
PreparedScanner prepared;
|
||||
memset(&prepared, 0, sizeof(prepared));
|
||||
DirectoryScanner* scanner = NULL;
|
||||
ArrayList* dry_manifest = NULL;
|
||||
ArrayList* dry_dirs = NULL;
|
||||
if (!config_send(client->file_descriptor, config))
|
||||
goto dry_fail;
|
||||
receive_daemon_motd(client, config);
|
||||
@@ -1624,10 +1778,34 @@ static int send_dry_run_remote(Config* config) {
|
||||
int file_count = 0;
|
||||
unsigned long long total_bytes = 0;
|
||||
char size_buffer[32];
|
||||
/* -n --delete: build the same keep-set manifest a real run would send so the
|
||||
receiver can enumerate (read-only) the destination extras. Filter-excluded
|
||||
and size-pruned protections are not propagated here, so a filtered dry-run
|
||||
may over-report; the no-filter case is exact. */
|
||||
dry_manifest = config->use_delete ? array_list_create(free) : NULL;
|
||||
if (config->use_delete && !dry_manifest)
|
||||
goto dry_fail;
|
||||
/* Scope the receiver-side extras walk to the receive root (the "." sentinel),
|
||||
exactly as the recursive transfer path does. */
|
||||
if (config->use_delete) {
|
||||
dry_dirs = array_list_create(free);
|
||||
char* root_marker = dry_dirs ? str_dup(".") : NULL;
|
||||
if (!dry_dirs || !root_marker || !array_list_add(dry_dirs, root_marker)) {
|
||||
free(root_marker);
|
||||
if (dry_dirs)
|
||||
array_list_delete(dry_dirs);
|
||||
dry_dirs = NULL;
|
||||
goto dry_fail;
|
||||
}
|
||||
}
|
||||
if (!config->quiet)
|
||||
printf("Dry run: files to be transferred\n");
|
||||
Chunk* chunk;
|
||||
while ((chunk = directory_scanner_next(scanner)) != NULL) {
|
||||
if (dry_manifest && !add_chunk_to_manifest(dry_manifest, chunk)) {
|
||||
chunk_destroy(chunk);
|
||||
goto dry_fail;
|
||||
}
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
File* f = chunk->items[i];
|
||||
if (!f)
|
||||
@@ -1689,12 +1867,68 @@ static int send_dry_run_remote(Config* config) {
|
||||
goto dry_fail;
|
||||
if (io_error)
|
||||
log_message(LOG_LEVEL_WARNING, "source scan hit an unreadable directory");
|
||||
/* Terminate the stream so the receiver emits its success frame; no data
|
||||
frame and no delete manifest are ever sent in dry-run. */
|
||||
/* Send the keep-set manifest (no data frames) so the receiver can enumerate
|
||||
the destination extras; an early-timing delete ACKs before it will accept
|
||||
the terminal FINISHED. */
|
||||
bool early_delete = config->use_delete && config_delete_timing_early(config);
|
||||
if (dry_manifest) {
|
||||
if (send_delete_manifest(client->file_descriptor, dry_manifest, NULL, NULL, NULL, dry_dirs) !=
|
||||
0)
|
||||
goto dry_fail;
|
||||
if (early_delete) {
|
||||
Status ack;
|
||||
if (!receive_status_keepalive(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC,
|
||||
DELETE_ACK_KEEPALIVE_SEC, client_abort_pending) ||
|
||||
ack != STATUS_OK)
|
||||
goto dry_fail;
|
||||
}
|
||||
}
|
||||
/* Terminate the stream so the receiver emits its success frame; no data frame
|
||||
is ever sent in dry-run. */
|
||||
if (!send_status(client->file_descriptor, STATUS_FINISHED))
|
||||
goto dry_fail;
|
||||
Status status;
|
||||
if (!receive_status(client->file_descriptor, &status) || status != STATUS_OK)
|
||||
if (!receive_status(client->file_descriptor, &status))
|
||||
goto dry_fail;
|
||||
if (status == STATUS_STATS) {
|
||||
ArrayList* would_delete = array_list_create(free);
|
||||
if (!would_delete)
|
||||
goto dry_fail;
|
||||
if (!receive_stats_record(client->file_descriptor, &dry_stats, would_delete)) {
|
||||
array_list_delete(would_delete);
|
||||
goto dry_fail;
|
||||
}
|
||||
/* rsync prints `*deleting PATH` when itemizing (or `deleting PATH` with
|
||||
--out-format / -v); the plain-total output used here has no delete
|
||||
counterpart, so only the itemize/out-format cases are rendered. */
|
||||
if (!config->quiet && (config->itemize_changes || config->out_format != NULL)) {
|
||||
for (int i = 0; i < would_delete->size; i++) {
|
||||
const char* raw = (const char*)would_delete->items[i];
|
||||
const char* path = delete_display_path(config, raw);
|
||||
if (config->out_format != NULL) {
|
||||
ChangeEvent event;
|
||||
memset(&event, 0, sizeof(event));
|
||||
event.decision = CHANGE_SENT;
|
||||
event.deleted = true;
|
||||
event.name = path;
|
||||
event.path = path;
|
||||
char* line = change_render_format(config->out_format, config, &event);
|
||||
if (line) {
|
||||
printf("%s\n", line);
|
||||
free(line);
|
||||
}
|
||||
} else {
|
||||
char* escaped = output_escape(path, config->eight_bit_output);
|
||||
printf("*deleting %s\n", escaped ? escaped : path);
|
||||
free(escaped);
|
||||
}
|
||||
}
|
||||
}
|
||||
array_list_delete(would_delete);
|
||||
if (!receive_status(client->file_descriptor, &status))
|
||||
goto dry_fail;
|
||||
}
|
||||
if (status != STATUS_OK)
|
||||
goto dry_fail;
|
||||
if (!config->quiet) {
|
||||
if (config->human_readable)
|
||||
@@ -1703,9 +1937,14 @@ static int send_dry_run_remote(Config* config) {
|
||||
else
|
||||
printf("Total: %d files, %.1f MB\n", file_count, (double)total_bytes / (double)BYTES_PER_MIB);
|
||||
}
|
||||
report_transfer_stats(config, file_count, total_bytes, dry_start, &dry_stats);
|
||||
ret = io_error ? 1 : 0;
|
||||
|
||||
dry_fail:
|
||||
if (dry_manifest)
|
||||
array_list_delete(dry_manifest);
|
||||
if (dry_dirs)
|
||||
array_list_delete(dry_dirs);
|
||||
if (scanner)
|
||||
directory_scanner_destroy(scanner);
|
||||
prepared_scanner_destroy(&prepared);
|
||||
@@ -2029,6 +2268,8 @@ static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
|
||||
(stream && !config->use_compression)) &&
|
||||
source_is_regular_file(f);
|
||||
SourceFile* source = remove_sources ? source_file_create(f) : NULL;
|
||||
unsigned long long bytes_before = protocol_bytes_written();
|
||||
unsigned long long read_before = protocol_bytes_read();
|
||||
int rc = send_single_file(client, f, config, config->use_incremental, use_sendfile);
|
||||
if (rc == 1) {
|
||||
source_file_destroy(source);
|
||||
@@ -2038,7 +2279,9 @@ static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
|
||||
source_file_destroy(source);
|
||||
return -1;
|
||||
}
|
||||
change_emit_file_sent(config, f);
|
||||
change_emit_file_sent_bytes(config, f, protocol_bytes_written() - bytes_before,
|
||||
protocol_bytes_read() - read_before);
|
||||
client_progress_file(config, f);
|
||||
if (source && !array_list_add(remove_sources, source)) {
|
||||
source_file_destroy(source);
|
||||
return -1;
|
||||
@@ -2096,6 +2339,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
}
|
||||
}
|
||||
|
||||
client_progress_begin(context->config);
|
||||
while (true) {
|
||||
/* Graceful abort (Ctrl-C/SIGTERM): tell the receiver to clean up instead of
|
||||
dying abruptly. Best-effort: a failed send just means the peer is gone.
|
||||
@@ -2223,7 +2467,10 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
!send_dir_times(client, context->config, context->dir_entries))
|
||||
goto send_fail;
|
||||
bool delete_limit = false;
|
||||
bool ok = finalize_transfer(client, context->config, context->remove_source_files, &delete_limit);
|
||||
ReceiverStats recv_stats;
|
||||
memset(&recv_stats, 0, sizeof(recv_stats));
|
||||
bool ok = finalize_transfer(client, context->config, context->remove_source_files, &delete_limit,
|
||||
&recv_stats);
|
||||
context->delete_limit = delete_limit;
|
||||
if (!ok && context->config->use_delete)
|
||||
log_message(LOG_LEVEL_ERROR,
|
||||
@@ -2234,7 +2481,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
int total_files = context->total_files;
|
||||
unsigned long long total_bytes = context->total_bytes;
|
||||
mtx_unlock(&context->mutex_progress);
|
||||
report_transfer_stats(context->config, total_files, total_bytes, start);
|
||||
report_transfer_stats(context->config, total_files, total_bytes, start, &recv_stats);
|
||||
log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files,
|
||||
(double)total_bytes / (double)BYTES_PER_MIB);
|
||||
disconnect_transfer_client(client);
|
||||
@@ -2418,58 +2665,6 @@ static int load_files_multithreaded(void* pipeline_context) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Print a one-line transfer progress report to stderr. `suffix` ends the
|
||||
line (e.g. "Done.\n") or is "" for in-place refresh. Shared by the
|
||||
single-threaded loop and the multithreaded progress thread. */
|
||||
static void print_transfer_progress(unsigned long long total_bytes, time_t start,
|
||||
const char* suffix, bool human_readable) {
|
||||
double elapsed = difftime(time(NULL), start);
|
||||
double rate = elapsed > 0.0 ? (double)total_bytes / ((double)BYTES_PER_MIB * elapsed) : 0.0;
|
||||
if (human_readable) {
|
||||
char total_buffer[32];
|
||||
char rate_buffer[32];
|
||||
fprintf(stderr, "\rSent %s (%s/s) %s",
|
||||
display_bytes(total_bytes, true, total_buffer, sizeof(total_buffer)),
|
||||
display_bytes((unsigned long long)(rate * (double)BYTES_PER_MIB), true, rate_buffer,
|
||||
sizeof(rate_buffer)),
|
||||
suffix);
|
||||
} else {
|
||||
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) %s", (double)total_bytes / (double)BYTES_PER_MIB,
|
||||
rate, suffix);
|
||||
}
|
||||
fflush(stderr);
|
||||
}
|
||||
|
||||
/* Progress-reporting thread for multithreaded send. Runs in parallel with
|
||||
the scanner/loader/sender threads and prints periodic progress to stderr. */
|
||||
static int progress_thread_fn(void* arg) {
|
||||
PipelineContextSender* context = (PipelineContextSender*)arg;
|
||||
time_t last_progress = 0;
|
||||
time_t start = time(NULL);
|
||||
|
||||
while (true) {
|
||||
mtx_lock(&context->mutex_progress);
|
||||
bool done = context->sender_done;
|
||||
unsigned long long total = context->progress_bytes;
|
||||
mtx_unlock(&context->mutex_progress);
|
||||
|
||||
if (done) {
|
||||
print_transfer_progress(total, start, "Done.\n", context->config->human_readable);
|
||||
break;
|
||||
}
|
||||
|
||||
time_t now = time(NULL);
|
||||
if (now - last_progress >= 1) {
|
||||
last_progress = now;
|
||||
print_transfer_progress(total, start, "", context->config->human_readable);
|
||||
}
|
||||
|
||||
struct timespec ts = {0, 100 * 1000000L}; /* 100 ms */
|
||||
thrd_sleep(&ts, NULL);
|
||||
}
|
||||
return thrd_success;
|
||||
}
|
||||
|
||||
/* Phase 6 residual-batch (client-only). --write-batch=FILE / --only-write-batch
|
||||
* emit a self-contained single-file batch of a whole source tree from a
|
||||
* deterministic separate scan pass. Each chunk's file images are fully loaded
|
||||
@@ -2746,8 +2941,8 @@ int send_files(Config* config) {
|
||||
Chunk* current_chunk;
|
||||
unsigned long long total_bytes = 0;
|
||||
int total_files = 0;
|
||||
time_t last_progress = 0;
|
||||
time_t start = time(NULL);
|
||||
client_progress_begin(config);
|
||||
/* 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;
|
||||
@@ -2809,13 +3004,6 @@ int send_files(Config* config) {
|
||||
break;
|
||||
}
|
||||
total_bytes += chunk_bytes;
|
||||
if (config->show_progress && !config->quiet) {
|
||||
time_t now = time(NULL);
|
||||
if (now - last_progress >= 1) {
|
||||
last_progress = now;
|
||||
print_transfer_progress(total_bytes, start, "", config->human_readable);
|
||||
}
|
||||
}
|
||||
chunk_destroy(current_chunk);
|
||||
}
|
||||
if (send_failed) {
|
||||
@@ -2886,15 +3074,15 @@ int send_files(Config* config) {
|
||||
if (!send_dir_times(client, config, dir_entries))
|
||||
goto send_fail;
|
||||
bool delete_limit = false;
|
||||
bool ok = finalize_transfer(client, config, remove_sources, &delete_limit);
|
||||
ReceiverStats recv_stats;
|
||||
memset(&recv_stats, 0, sizeof(recv_stats));
|
||||
bool ok = finalize_transfer(client, config, remove_sources, &delete_limit, &recv_stats);
|
||||
if (!ok && config->use_delete)
|
||||
log_message(LOG_LEVEL_ERROR,
|
||||
"server reported a deletion failure (--delete); see the server log for the reason");
|
||||
if (ok)
|
||||
remove_transferred_sources(config, remove_sources);
|
||||
if (config->show_progress && !config->quiet)
|
||||
print_transfer_progress(total_bytes, start, "Done.\n", config->human_readable);
|
||||
report_transfer_stats(config, total_files, total_bytes, start);
|
||||
report_transfer_stats(config, total_files, total_bytes, start, &recv_stats);
|
||||
log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files,
|
||||
(double)total_bytes / (double)BYTES_PER_MIB);
|
||||
/* A skipped source entry (--ignore-errors past an unreadable directory, or a
|
||||
@@ -3138,29 +3326,11 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
return 1;
|
||||
}
|
||||
|
||||
thrd_t progress;
|
||||
bool progress_created = false;
|
||||
if (config->show_progress && !config->quiet) {
|
||||
progress_created = (thrd_create(&progress, progress_thread_fn, context) == thrd_success);
|
||||
if (!progress_created) {
|
||||
log_perror("Error creating progress thread");
|
||||
/* Non-fatal; continue without progress reporting */
|
||||
}
|
||||
}
|
||||
|
||||
int sender_result;
|
||||
thrd_join(scanner, NULL);
|
||||
thrd_join(loader, NULL);
|
||||
thrd_join(sender, &sender_result);
|
||||
|
||||
if (progress_created) {
|
||||
/* Signal progress thread to exit if it hasn't already */
|
||||
mtx_lock(&context->mutex_progress);
|
||||
context->sender_done = true;
|
||||
mtx_unlock(&context->mutex_progress);
|
||||
thrd_join(progress, NULL);
|
||||
}
|
||||
|
||||
bool scan_io;
|
||||
mtx_lock(&context->mutex_scanner);
|
||||
scan_io = context->scan_had_io_error;
|
||||
|
||||
Reference in New Issue
Block a user