diff --git a/src/client/change_list.c b/src/client/change_list.c index dbe0ca1..3cd0c12 100644 --- a/src/client/change_list.c +++ b/src/client/change_list.c @@ -1,5 +1,7 @@ #include "change_list.h" +#include "checksum.h" #include "utils.h" +#include #include #include #include @@ -200,6 +202,82 @@ char* change_render_itemize(const Config* config, const ChangeEvent* event) { /* ---- --out-format / --log-file-format ---- */ +/* rsync 3.4.1's `%C` uses the negotiated transfer checksum; with the default + * "auto" choice on both ends that is xxh128. FastSync's internal XXH64 default + * is not an rsync algorithm, so map it to xxh128 for parity. */ +static ChecksumAlgo out_format_checksum_algo(const Config* config) { + switch ((ChecksumAlgo)config->checksum_algo) { + case CHECKSUM_ALGO_MD5: + return CHECKSUM_ALGO_MD5; + case CHECKSUM_ALGO_XXH3: + return CHECKSUM_ALGO_XXH3; + case CHECKSUM_ALGO_XXH128: + return CHECKSUM_ALGO_XXH128; + case CHECKSUM_ALGO_XXH64: + default: + return CHECKSUM_ALGO_XXH128; + } +} + +/* Render a digest as rsync's sum_as_hex: for xxh128 the HIGH 64-bit half is + * printed before the low half; every other algorithm prints its bytes in order. */ +static void digest_to_hex(ChecksumAlgo algo, const uint8_t* digest, size_t len, char* out) { + if (algo == CHECKSUM_ALGO_XXH128 && len == 16) { + uint64_t low = 0; + uint64_t high = 0; + memcpy(&low, digest, sizeof(low)); + memcpy(&high, digest + 8, sizeof(high)); + snprintf(out, len * 2 + 1, "%016llx%016llx", (unsigned long long)high, (unsigned long long)low); + return; + } + static const char hex[] = "0123456789abcdef"; + for (size_t i = 0; i < len; i++) { + out[i * 2] = hex[(digest[i] >> 4) & 0xf]; + out[i * 2 + 1] = hex[digest[i] & 0xf]; + } + out[len * 2] = '\0'; +} + +static bool format_uses_checksum(const char* format) { + if (format == NULL) + return false; + for (const char* p = format; *p != '\0';) { + if (*p != '%') { + p++; + continue; + } + char token = p[1]; + if (token == '\0') + break; + if (token == 'C') + return true; + p += 2; + } + return false; +} + +/* Fill event->checksum/checksum_known for a transferred regular file. A + * non-regular entry (or a hard-link sibling) leaves checksum_known false, which + * renders as spaces like rsync. */ +static void fill_event_checksum(const Config* config, const File* file, ChangeEvent* event) { + if (file == NULL || file->is_dir || file->is_symlink || file->is_special || + (file->link_group != 0 && !file->link_first)) + return; + if (!format_uses_checksum(config->out_format) && !format_uses_checksum(config->log_file_format)) + return; + if (file->path == NULL) + return; + ChecksumAlgo algo = out_format_checksum_algo(config); + uint8_t digest[CHECKSUM_MAX_DIGEST_LEN]; + size_t len = 0; + /* rsync's %C is the transfer checksum, which is always seeded with 0 (it is + * independent of --checksum-seed, as rsync 3.4.1 demonstrates). */ + if (!checksum_digest_file(algo, 0, file->path, digest, sizeof(digest), &len)) + return; + digest_to_hex(algo, digest, len, event->checksum); + event->checksum_known = true; +} + char* change_render_format(const char* format, const Config* config, const ChangeEvent* event) { if (format == NULL || event == NULL) return NULL; @@ -221,6 +299,11 @@ char* change_render_format(const char* format, const Config* config, const Chang ok = strbuf_append_char(&line, '%'); break; case 'i': { + if (event->deleted) { + /* rsync's ITEM_DELETED itemize code: `*deleting ` (11 chars). */ + ok = strbuf_append(&line, "*deleting "); + break; + } char code[12]; itemize_code(config, event, code); ok = strbuf_append(&line, code); @@ -245,6 +328,22 @@ char* change_render_format(const char* format, const Config* config, const Chang int written = snprintf(digits, sizeof(digits), "%llu", event->bytes_sent); ok = written >= 0 && (size_t)written < sizeof(digits) && strbuf_append(&line, digits); } break; + case 'c': { + char digits[32]; + int written = snprintf(digits, sizeof(digits), "%llu", event->bytes_read); + ok = written >= 0 && (size_t)written < sizeof(digits) && strbuf_append(&line, digits); + } break; + case 'C': { + if (event->checksum_known) { + ok = strbuf_append(&line, event->checksum); + } else { + /* rsync pads a non-regular / untransferred entry with spaces. */ + ChecksumAlgo algo = out_format_checksum_algo(config); + int width = checksum_digest_len(algo) * 2; + for (int i = 0; i < width && ok; i++) + ok = strbuf_append_char(&line, ' '); + } + } break; case 'M': { char when[32]; if (format_rsync_datetime(event->mtime_sec, true, when, sizeof(when))) @@ -464,7 +563,8 @@ static void fill_event_from_file(const Config* config, const File* file, ChangeE } } -void change_emit_file_sent(const Config* config, const File* file) { +void change_emit_file_sent_bytes(const Config* config, const File* file, + unsigned long long bytes_sent, unsigned long long bytes_read) { if (file == NULL || !change_list_enabled(config)) return; ChangeEvent event; @@ -489,19 +589,27 @@ void change_emit_file_sent(const Config* config, const File* file) { event.hardlink_target = file->hardlink_target; event.bytes_sent = 0; } else { - /* Literal payload bytes delivered; compressed/delta wire bytes are not - * separately counted. */ - event.bytes_sent = event.size; + event.bytes_sent = bytes_sent; + event.bytes_read = bytes_read; } char* name = NULL; char* path = NULL; fill_event_from_file(config, file, &event, &name, &path); - if (name != NULL && path != NULL) + if (name != NULL && path != NULL) { + fill_event_checksum(config, file, &event); change_emit(config, &event); + } free(name); free(path); } +void change_emit_file_sent(const Config* config, const File* file) { + if (file == NULL) + return; + unsigned long long payload = file->data != NULL ? file->data->size : 0; + change_emit_file_sent_bytes(config, file, payload, 0); +} + void change_emit_dir_sent(const Config* config, const File* file) { if (file == NULL || !change_list_enabled(config)) return; diff --git a/src/client/change_list.h b/src/client/change_list.h index 874b9fe..6771d63 100644 --- a/src/client/change_list.h +++ b/src/client/change_list.h @@ -2,6 +2,7 @@ #define CHANGE_LIST_H #include "config.h" +#include "checksum.h" #include "file_types.h" #include "format.h" #include @@ -34,10 +35,17 @@ typedef struct { bool is_symlink; bool is_special; bool is_hardlink; /* a hard-link sibling (linked, no data sent) */ + bool deleted; /* a would-delete report (-n --delete); no source file */ const char* symlink_target; const char* hardlink_target; unsigned long long size; /* source file length in bytes */ - unsigned long long bytes_sent; /* literal data bytes actually transferred */ + unsigned long long bytes_sent; /* wire bytes actually transferred (rsync %b) */ + unsigned long long bytes_read; /* wire bytes read back for this file (rsync %c) */ + /* rsync %C: whole-file checksum hex for a transferred regular file. Only + * filled when the active format uses %C (checksum_known == false otherwise, + * which renders as spaces like rsync for non-regular entries). */ + bool checksum_known; + char checksum[CHECKSUM_MAX_DIGEST_LEN * 2 + 1]; time_t mtime_sec; long mtime_nsec; mode_t mode; @@ -61,7 +69,9 @@ char* change_render_itemize_code(const Config* config, const ChangeEvent* event) /* Expand an --out-format/--log-file-format template. Supported tokens: * %i itemize code %n transfer-relative name (dir: trailing /) * %f long display path %l file length in bytes - * %b bytes actually sent %M mtime (YYYY/MM/DD-HH:MM:SS) + * %b wire bytes transferred %c wire bytes read back for the file + * %C whole-file checksum hex (xxh128 by default; spaces for non-regular) + * %M mtime (YYYY/MM/DD-HH:MM:SS) * %t current time %o operation ("send"/"del.") * %p pid %B permission bits without the type char * %U uid %G gid @@ -80,7 +90,14 @@ char* change_render_list_line(const Config* config, const ChangeEvent* event); * CHANGE_UP_TO_DATE events produce no output. */ void change_emit(const Config* config, const ChangeEvent* event); -/* Build and emit a CHANGE_SENT event for a file the client just sent. */ +/* Build and emit a CHANGE_SENT event for a file the client just sent. `bytes_sent` + * / `bytes_read` are the process-wide wire-byte deltas for this file (rsync's + * %b / %c); pass 0 when unknown. */ +void change_emit_file_sent_bytes(const Config* config, const File* file, + unsigned long long bytes_sent, unsigned long long bytes_read); + +/* Build and emit a CHANGE_SENT event for a file the client just sent, deriving + * the wire byte counts from the source payload length. */ void change_emit_file_sent(const Config* config, const File* file); /* Build and emit a CHANGE_SENT event for an explicit directory entry (-d). */ diff --git a/src/client/client_cli.c b/src/client/client_cli.c index ab7a1d9..eb5b41c 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -2438,6 +2438,21 @@ static int cli_finalize_config(Config* config, bool verbose, bool no_delta, bool * check. This is a wire field. */ config->report_dest_info = config->itemize_changes || config->out_format != NULL || (config->log_file != NULL && config->log_file_format != NULL); + /* Wire-stats parity: --stats, --progress/-P, an --out-format token that needs + * a wire counter (%b/%c), or a dry-run --delete need the receiver's + * end-of-transfer STATUS_STATS report. This is a wire field (protocol + * 2.25.0). */ + bool format_needs_wire = false; + if (config->out_format != NULL) { + for (const char* p = config->out_format; *p != '\0'; p++) { + if (p[0] == '%' && (p[1] == 'b' || p[1] == 'c')) { + format_needs_wire = true; + break; + } + } + } + config->report_stats = config->stats || config->show_progress || format_needs_wire || + (config->dry_run && config->use_delete); return 0; } diff --git a/src/client/client_send.c b/src/client/client_send.c index 4d900c7..49dd445 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -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; diff --git a/src/client/usage.c b/src/client/usage.c index d0d8a75..e654e5f 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -270,7 +270,7 @@ void print_usage(void) { printf(" --suffix Backup suffix (default: ~)\n"); printf(" --stats Print transfer statistics at end\n"); printf(" -i, --itemize-changes Print an rsync-style per-file change line\n"); - printf(" --out-format=FORMAT Output format for changed files (%%f %%n %%l %%b %%M %%%%)\n"); + printf(" --out-format=FORMAT Output format (%%f %%n %%l %%b %%c %%C %%i %%M %%%%)\n"); printf(" --list-only List source files instead of transferring\n"); printf(" --log-file-format=FORMAT Per-file log line format (needs --log-file)\n"); printf(" -h, --human-readable Print byte sizes in human-readable form\n"); diff --git a/src/server/receiver.c b/src/server/receiver.c index fc2346e..bd88edf 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -12,6 +12,7 @@ #include "protocol.h" #include "utils.h" #include +#include #include #include @@ -59,6 +60,29 @@ bool receiver_send_final_success(int fd, const Config* config, const ReceiverOut return send_status(fd, final_status); } +bool receiver_send_stats_frame(int fd, const Config* config, const ReceiverStats* stats, + const struct ArrayList* would_delete) { + if (!config->report_stats) + return true; + ReceiverStats local; + memset(&local, 0, sizeof(local)); + const ReceiverStats* out = stats ? stats : &local; + size_t count = would_delete ? (size_t)would_delete->size : 0; + if (count > (size_t)MAX_MANIFEST_ENTRIES) + count = MAX_MANIFEST_ENTRIES; + ReceiverStats record = *out; + record.would_delete_count = count; + if (!send_status(fd, STATUS_STATS) || !format_stats_send(fd, &record) || + !send_int(fd, (int)count)) + return false; + for (size_t i = 0; i < count; i++) { + const char* path = (const char*)would_delete->items[i]; + if (!send_wire_str(fd, path ? path : "")) + return false; + } + return true; +} + static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) { if (!chunk || !sink || !sink->store_file) return false; @@ -346,7 +370,14 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver if (config->dry_run) { /* Server-contacting --dry-run mutates nothing, so a keep-set manifest is consumed and discarded. The early-delete mode still needs its ACK - so a sender blocked on the delete handshake is not left hanging. */ + so a sender blocked on the delete handshake is not left hanging. + When would-delete reporting is armed, enumerate (read-only) the + destination extras so the terminal STATUS_STATS frame can list them. */ + if (config->use_delete && sink->would_delete) { + size_t count = 0; + if (!manifest_would_delete_list(config, manifest, sink->would_delete, &count)) + log_message(LOG_LEVEL_WARNING, "dry-run: could not enumerate would-delete paths"); + } delete_manifest_free(manifest); if (early_delete && !send_status(file_descriptor, STATUS_OK)) goto fail; @@ -517,6 +548,10 @@ typedef struct { /* Set when a --max-delete commit was capped; the terminal frame then carries STATUS_DELETE_LIMIT so the sender exits 25 like rsync. */ bool delete_limit_reached; + /* End-of-transfer wire counters (protocol 2.25.0) and the -n/--dry-run + --delete would-delete path list collected while processing the manifest. */ + ReceiverStats stats; + ArrayList* would_delete; } ReceiverSaveContext; static bool receiver_save_file(File* file, void* context_pointer) { @@ -565,6 +600,8 @@ static void receiver_note_delete_limit(void* context_pointer) { static bool receiver_send_success_frame(int fd, void* context_pointer) { ReceiverSaveContext* context = context_pointer; Status final_status = context->delete_limit_reached ? STATUS_DELETE_LIMIT : STATUS_OK; + if (!receiver_send_stats_frame(fd, context->config, &context->stats, context->would_delete)) + return false; /* Server-contacting --dry-run: nothing was staged or written, so there is nothing to publish and no directory times to stamp. */ if (context->config->dry_run) @@ -592,12 +629,22 @@ static bool receiver_send_success_frame(int fd, void* context_pointer) { int receiver_receive_files(Config* config, int file_descriptor) { ReceiverSaveContext context = {.config = config, .outcomes = {0}}; dir_time_list_init(&context.dir_times); - ReceiverSink sink = {receiver_save_file, &context, true, true, receiver_send_success_frame, - receiver_note_delete_limit}; + context.would_delete = array_list_create(free); + if (!context.would_delete) + return -1; + ReceiverSink sink = {receiver_save_file, + &context, + true, + true, + receiver_send_success_frame, + receiver_note_delete_limit, + &context.stats, + context.would_delete}; int ret = receiver_process(config, file_descriptor, &sink); if (ret != 0 && config->delay_updates && config->delay_context) delay_updates_cleanup(config->delay_context); receiver_outcomes_destroy(&context.outcomes); dir_time_list_free(&context.dir_times); + array_list_delete(context.would_delete); return ret; } diff --git a/src/server/receiver.h b/src/server/receiver.h index e2248d5..0dc2277 100644 --- a/src/server/receiver.h +++ b/src/server/receiver.h @@ -40,15 +40,28 @@ typedef struct { ReceiverSuccessFrame send_success_frame; /* Optional; may be NULL when the sink has no --max-delete handling. */ ReceiverNoteDeleteLimit note_delete_limit; + /* Optional end-of-transfer wire counters (protocol 2.25.0). When non-NULL + and the wire config carries report_stats, the success frame is preceded by + a STATUS_STATS record; `would_delete` (optional, receiver-owned strings) + carries the -n/--dry-run --delete path list. */ + ReceiverStats* stats; + struct ArrayList* would_delete; } ReceiverSink; bool receiver_outcomes_append(ReceiverOutcomes* outcomes, unsigned char code); void receiver_outcomes_destroy(ReceiverOutcomes* outcomes); + /* Send the terminal success frame. `final_status` is usually STATUS_OK, or STATUS_DELETE_LIMIT when a --max-delete commit was capped. */ bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes, Status final_status); +/* Emit STATUS_STATS (a fixed ReceiverStats record plus, when `would_delete` is + non-NULL, a count and that many wire strings) when the wire config requested + report_stats. A no-op otherwise. */ +bool receiver_send_stats_frame(int fd, const Config* config, const ReceiverStats* stats, + const struct ArrayList* would_delete); + int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink); /* receiver_process with an escape hatch for the commit-style (late) deletion: when `pending_manifest` is non-NULL the receiver does NOT delete at diff --git a/src/server/receiver_pipeline.c b/src/server/receiver_pipeline.c index 0f0bc3a..770bb3e 100644 --- a/src/server/receiver_pipeline.c +++ b/src/server/receiver_pipeline.c @@ -165,8 +165,14 @@ int receive_thread(void* pipeline_context) { const Config* config = context->config; mtx_unlock(&context->mutex); - ReceiverSink sink = { - receiver_enqueue_file, context, false, false, NULL, receiver_pipeline_note_delete_limit}; + ReceiverSink sink = {receiver_enqueue_file, + context, + false, + false, + NULL, + receiver_pipeline_note_delete_limit, + NULL, + NULL}; if (receiver_process_pending((Config*)config, file_descriptor, &sink, &context->deferred_manifest, &context->deferred_plans) != 0) { receiver_thread_fail(context); diff --git a/src/shared/checksum.c b/src/shared/checksum.c index b95a8bc..f4706b8 100644 --- a/src/shared/checksum.c +++ b/src/shared/checksum.c @@ -1,10 +1,14 @@ #include "checksum.h" +#include #include #include #include +#include /* delta.c owns the single XXH_IMPLEMENTATION that provides the xxHash symbols - * for the whole binary; this TU only needs the declarations. */ + * for the whole binary; this TU only needs the declarations. The streaming + * state structs and XXH3_update are exposed only with XXH_STATIC_LINKING_ONLY. */ +#define XXH_STATIC_LINKING_ONLY #include bool checksum_digest(ChecksumAlgo algo, uint64_t seed, const void* data, size_t size, uint8_t* out, @@ -54,6 +58,99 @@ bool checksum_digest(ChecksumAlgo algo, uint64_t seed, const void* data, size_t return false; } +bool checksum_digest_file(ChecksumAlgo algo, uint64_t seed, const char* path, uint8_t* out, + size_t out_capacity, size_t* out_len) { + if (!path || !out || !out_len || out_capacity < CHECKSUM_MAX_DIGEST_LEN) + return false; + + int fd = open(path, O_RDONLY | O_CLOEXEC); + if (fd < 0) + return false; + + uint8_t buffer[64 * 1024]; + bool ok = false; + + if (algo == CHECKSUM_ALGO_MD5) { + EVP_MD_CTX* ctx = EVP_MD_CTX_new(); + if (!ctx) { + close(fd); + return false; + } + unsigned int digest_len = 0; + if (EVP_DigestInit_ex(ctx, EVP_md5(), NULL) == 1) { + ok = true; + ssize_t got; + while ((got = read(fd, buffer, sizeof(buffer))) > 0) { + if (EVP_DigestUpdate(ctx, buffer, (size_t)got) != 1) { + ok = false; + break; + } + } + if (got < 0) + ok = false; + if (ok && EVP_DigestFinal_ex(ctx, out, &digest_len) == 1 && digest_len <= out_capacity) + *out_len = digest_len; + else + ok = false; + } + EVP_MD_CTX_free(ctx); + close(fd); + return ok; + } + + XXH64_state_t xxh64; + XXH3_state_t* xxh3 = NULL; + if (algo == CHECKSUM_ALGO_XXH64) { + XXH64_reset(&xxh64, seed); + } else if (algo == CHECKSUM_ALGO_XXH3 || algo == CHECKSUM_ALGO_XXH128) { + xxh3 = XXH3_createState(); + if (!xxh3) { + close(fd); + return false; + } + if (algo == CHECKSUM_ALGO_XXH3) + XXH3_64bits_reset_withSeed(xxh3, seed); + else + XXH3_128bits_reset_withSeed(xxh3, seed); + } else { + close(fd); + return false; + } + + ok = true; + ssize_t got; + while ((got = read(fd, buffer, sizeof(buffer))) > 0) { + if (algo == CHECKSUM_ALGO_XXH64) + XXH64_update(&xxh64, buffer, (size_t)got); + else if (XXH3_64bits_update(xxh3, buffer, (size_t)got) == XXH_ERROR) { + ok = false; + break; + } + } + if (got < 0) + ok = false; + + if (ok) { + if (algo == CHECKSUM_ALGO_XXH64) { + uint64_t digest = XXH64_digest(&xxh64); + memcpy(out, &digest, sizeof(digest)); + *out_len = sizeof(digest); + } else if (algo == CHECKSUM_ALGO_XXH3) { + uint64_t digest = XXH3_64bits_digest(xxh3); + memcpy(out, &digest, sizeof(digest)); + *out_len = sizeof(digest); + } else { + XXH128_hash_t digest = XXH3_128bits_digest(xxh3); + memcpy(out, &digest, sizeof(digest)); + *out_len = sizeof(digest); + } + } + if (xxh3) + XXH3_freeState(xxh3); + close(fd); + return ok; +} + int checksum_algo_from_name(const char* name) { if (!name) return -1; diff --git a/src/shared/checksum.h b/src/shared/checksum.h index c323730..fe32218 100644 --- a/src/shared/checksum.h +++ b/src/shared/checksum.h @@ -35,11 +35,16 @@ typedef enum { bool checksum_digest(ChecksumAlgo algo, uint64_t seed, const void* data, size_t size, uint8_t* out, size_t out_capacity, size_t* out_len); -/* Resolve a --checksum-choice string (case-insensitive) to an algorithm id. - * Accepts "xxh64"/"xxhash", "xxh3", "xxh128" and "md5". "auto", rsync's - * default automatic choice, is resolved to the default by the caller (it is not - * a distinct algorithm here). Returns -1 for any name FastSync does not - * implement (md4/sha1/none included). */ +/* Streaming whole-file digest: hash the contents of `path` without holding the + * whole file in memory. Same digest/capacity contract as checksum_digest. + * Returns false on open/read failure or an undersized buffer. */ +bool checksum_digest_file(ChecksumAlgo algo, uint64_t seed, const char* path, uint8_t* out, + size_t out_capacity, size_t* out_len); + +/* Resolve a --checksum-choice string (case-insensitive) to an algorithm id. * Accepts + * "xxh64"/"xxhash", "xxh3", "xxh128" and "md5". "auto", rsync's default automatic choice, is + * resolved to the default by the caller (it is not a distinct algorithm here). Returns -1 for any + * name FastSync does not implement (md4/sha1/none included). */ int checksum_algo_from_name(const char* name); /* Canonical name of an algorithm (used in CLI error messages). */ diff --git a/src/shared/config.h b/src/shared/config.h index c8654e3..8e7c0fe 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -250,8 +250,17 @@ typedef enum SuperMode { SUPER_MODE_AUTO = 0, SUPER_MODE_ON = 1, SUPER_MODE_OFF * answer every per-file STATUS_CHECK with a STATUS_DEST_INFO snapshot of the * pre-transfer destination entry (see protocol.h). It is set by the client * only when -i/--itemize-changes or --out-format asks for per-file change - * output; the transfer decision itself is unchanged. */ -#define CONFIG_WIRE_OUTPUT_FIELDS(X) X(report_dest_info, bool, false, BOOL) + * output; the transfer decision itself is unchanged. + * + * Wire-stats wave (protocol 2.25.0). report_stats tells the receiver to send a + * STATUS_STATS frame immediately before its terminal success status carrying + * the receiver-only counters (matched data, deleted-file count) and, + * for -n/--dry-run --delete, the destination-relative paths it WOULD have + * deleted. It is set by the client only when --stats, --progress/-P, an + * --out-format token needs a wire counter (%b/%c), or a dry-run carries + * --delete; the transfer decision itself is unchanged. */ +#define CONFIG_WIRE_OUTPUT_FIELDS(X) \ + X(report_dest_info, bool, false, BOOL) X(report_stats, bool, false, BOOL) /* All serialized fields, in exact wire order. Concatenating the per-segment * lists here is what keeps the declaration order = the wire order. */ @@ -919,8 +928,22 @@ typedef struct Config { * snapshot of the old entry) before its ordinary verdict when the config frame * carries the new report_dest_info bool appended after the --copy-as block. * This is both a config-frame layout change (one trailing bool) and a frame - * sequence change (the new status). */ -#define PROTOCOL_VERSION "2.24.0" + * sequence change (the new status). + * + * (4) Delete timing (protocol 2.24.0): the sender streams one delete plan per + * source directory so --delete-during/--delete-delay reproduce rsync's deletion + * timing (the plan fields and STATUS_DELETE_PLAN are documented at the keep-set + * / delete-plan definitions below). + * + * (5) Wire-stats parity (protocol 2.25.0): --stats, --progress/-P and the + * --out-format %b/%c tokens need receiver-only and wire counters that the push + * sender cannot observe, and -n/--dry-run --delete must report the extras it + * would have removed without deleting anything. The config frame gains one + * trailing report_stats bool and the receiver emits a new STATUS_STATS frame + * (carrying matched data, the deleted-file count and the would-delete path + * list) immediately before its terminal success status. Both a config-frame + * layout change and a frame-sequence change, hence the bump. */ +#define PROTOCOL_VERSION "2.25.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) /* Upper bound on total basis-dir entries (rsync caps --link-dest at 20). */ #define MAX_BASIS_DIRS 64 diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index 07aef4f..4fd73e5 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -3460,6 +3460,42 @@ static bool delete_missing_args_budgeted(const Config* config, DeleteManifest* m /* Public wrappers used outside the commit path (and by unit tests): no --max-delete budget. */ +bool manifest_would_delete_list(const Config* config, DeleteManifest* manifest, ArrayList* out, + size_t* count_out) { + if (count_out) + *count_out = 0; + if (!config || !manifest || !manifest->keeps || !out) + return false; + int skip_count = (config->delay_updates ? 1 : 0) + config->basis_count + + (manifest->protected ? manifest->protected->size : 0); + DeleteSkipEntry* skips = NULL; + if (skip_count > 0) { + skips = calloc((size_t)skip_count, sizeof(DeleteSkipEntry)); + if (!skips) + return false; + int idx = 0; + if (config->delay_updates) { + skips[idx].prefix = DELAY_UPDATES_STAGING_DIR; + skips[idx].top_level_only = true; + idx++; + } + for (int i = 0; i < config->basis_count; i++) { + skips[idx].prefix = config->basis_dirs[i].path; + skips[idx].top_level_only = false; + idx++; + } + for (int i = 0; i < manifest->protected->size; i++) { + skips[idx].prefix = (const char*)manifest->protected->items[i]; + skips[idx].top_level_only = false; + idx++; + } + } + bool ok = delete_extras_list(config->receive_root_directory, manifest->keeps, manifest->dirs, + skips, skip_count, out, count_out); + free(skips); + return ok; +} + bool manifest_delete_extras(const Config* config, DeleteManifest* manifest) { DeleteBudgetState budget = { .max_delete = SIZE_MAX, .deleted = 0, .skipped = 0, .limit_hit = false}; diff --git a/src/shared/file_receive.h b/src/shared/file_receive.h index 5cddd99..7419c0a 100644 --- a/src/shared/file_receive.h +++ b/src/shared/file_receive.h @@ -143,6 +143,14 @@ typedef enum { stopped part of the work, or DELETE_COMMIT_ERROR on a genuine failure. */ DeleteCommitResult manifest_delete_all(const Config* config, DeleteManifest* manifest); +/* -n/--dry-run --delete would-delete reporting: walk the destination exactly as + the delete pass would and append (strdup'd) destination-relative paths that + WOULD be removed to `out`, without touching disk. Uses the same staging-dir, + basis-dir and protected-prefix skips as the real commit. Returns true on a + clean walk; `*count_out` receives the number of paths appended. */ +bool manifest_would_delete_list(const Config* config, DeleteManifest* manifest, ArrayList* out, + size_t* count_out); + /* Outcome of a single file_save_to_disk operation. The receiver needs to distinguish "written" from "skipped" so --remove-source-files can be told which sources were actually stored. */ diff --git a/src/shared/file_send.c b/src/shared/file_send.c index dfdeb5c..f064f93 100644 --- a/src/shared/file_send.c +++ b/src/shared/file_send.c @@ -182,6 +182,7 @@ bool file_send_sendfile_with_skip(File* file, int file_descriptor, bool use_meta close(fd); return false; } + protocol_note_bytes_written((unsigned long long)sent); } close(fd); diff --git a/src/shared/format.c b/src/shared/format.c index d690e1b..147d002 100644 --- a/src/shared/format.c +++ b/src/shared/format.c @@ -101,3 +101,29 @@ bool format_dest_state_receive(int fd, OutputDestState* state) { state->gid = gid; return true; } + +bool format_stats_send(int fd, const ReceiverStats* stats) { + if (!stats) + return false; + unsigned long long matched = stats->matched_data; + unsigned long long deleted = stats->deleted_files; + unsigned long long would = stats->would_delete_count; + return send_n_data(fd, &matched, sizeof(matched)) && send_n_data(fd, &deleted, sizeof(deleted)) && + send_n_data(fd, &would, sizeof(would)); +} + +bool format_stats_receive(int fd, ReceiverStats* stats) { + if (!stats) + return false; + unsigned long long matched = 0; + unsigned long long deleted = 0; + unsigned long long would = 0; + if (!receive_n_data(fd, &matched, sizeof(matched)) || + !receive_n_data(fd, &deleted, sizeof(deleted)) || !receive_n_data(fd, &would, sizeof(would))) + return false; + memset(stats, 0, sizeof(*stats)); + stats->matched_data = matched; + stats->deleted_files = deleted; + stats->would_delete_count = would; + return true; +} diff --git a/src/shared/format.h b/src/shared/format.h index 7f4fc13..fc904d4 100644 --- a/src/shared/format.h +++ b/src/shared/format.h @@ -56,4 +56,21 @@ bool format_rsync_datetime(time_t when, bool dash, char* buffer, size_t buffer_s bool format_dest_state_send(int fd, const OutputDestState* state); bool format_dest_state_receive(int fd, OutputDestState* state); +/* End-of-transfer receiver counters reported through STATUS_STATS (protocol + * 2.25.0) when the wire config carries report_stats. `would_delete_count` is + * the number of destination-relative paths the receiver would have deleted in a + * -n/--dry-run --delete run; that many wire strings immediately follow the + * fixed record (sent/read by the caller). */ +typedef struct { + unsigned long long matched_data; + unsigned long long deleted_files; + unsigned long long would_delete_count; +} ReceiverStats; + +/* Fixed-width STATUS_STATS counter record. The status frame and the optional + * would-delete path list are sent/received by the caller. Returns false on I/O + * failure. */ +bool format_stats_send(int fd, const ReceiverStats* stats); +bool format_stats_receive(int fd, ReceiverStats* stats); + #endif diff --git a/src/shared/protocol.c b/src/shared/protocol.c index f3402a4..70c0d8a 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -29,6 +29,13 @@ static unsigned long long io_bwlimit = 0; static mtx_t bw_mutex; static once_flag bw_mutex_once = ONCE_FLAG_INIT; +/* Process-wide wire byte counters, used by the client to render rsync's + * --stats/--progress totals and the --out-format %b/%c tokens. The zero-copy + * sendfile path bypasses protocol_send_n_data, so it reports its bytes through + * protocol_note_bytes_written. */ +static atomic_ullong io_bytes_written = 0; +static atomic_ullong io_bytes_read = 0; + static unsigned long long global_bwlimit(void); static bool protocol_reserve_memory(ProtocolSession* session, size_t charge) { @@ -242,6 +249,18 @@ SSL* io_get_ssl(void) { return io_ssl; } +unsigned long long protocol_bytes_written(void) { + return atomic_load(&io_bytes_written); +} + +unsigned long long protocol_bytes_read(void) { + return atomic_load(&io_bytes_read); +} + +void protocol_note_bytes_written(unsigned long long bytes) { + atomic_fetch_add(&io_bytes_written, bytes); +} + static ProtocolSession* legacy_session(int read_fd, int write_fd) { if (bound_session) return bound_session; @@ -343,6 +362,7 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat wait_events = POLLOUT; } log_debug_message(LOG_DEBUG_IO, " Send n Data: %zu", total_bytes_send); + atomic_fetch_add(&io_bytes_written, (unsigned long long)total_bytes_send); return true; } @@ -430,6 +450,7 @@ static bool protocol_receive_n_data_until(ProtocolSession* session, void* data, wait_events = POLLIN; } log_debug_message(LOG_DEBUG_IO, " Received n Data: %zu", total_bytes_received); + atomic_fetch_add(&io_bytes_read, (unsigned long long)total_bytes_received); return true; } diff --git a/src/shared/protocol.h b/src/shared/protocol.h index 4dc91a0..30bc5af 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -195,7 +195,14 @@ enum NET_STATUS { * ("." for the receive root); then the child-directory count + names and the * child-file count + names that must be kept. Appended after * STATUS_DEST_INFO so no existing status is renumbered. */ - STATUS_DELETE_PLAN + STATUS_DELETE_PLAN, + /* End-of-transfer receiver counter report (protocol 2.25.0). When the wire + * config carries report_stats=true, the receiver sends this status once, + * immediately before its terminal success status, followed by a fixed stats + * record (see format_stats_send/receive in format.h) and, when the run is a + * --dry-run with --delete, the would-delete path list. Appended after + * STATUS_DELETE_PLAN so no existing status is renumbered. */ + STATUS_STATS }; void io_set_fds(int read_fd, int write_fd); @@ -203,6 +210,14 @@ void io_set_bwlimit(unsigned long long bytes_per_sec); void io_set_ssl(SSL* ssl); SSL* io_get_ssl(void); +/* Process-wide wire byte counters. protocol_send_n_data/protocol_receive_n_data + * update them; the zero-copy sendfile path reports through + * protocol_note_bytes_written. Used by the client to render rsync's + * --stats/--progress totals and the --out-format %b/%c tokens. */ +unsigned long long protocol_bytes_written(void); +unsigned long long protocol_bytes_read(void); +void protocol_note_bytes_written(unsigned long long bytes); + void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd); /* Transitional bridge for helpers whose signatures still carry only an fd. */ void protocol_session_bind(ProtocolSession* session); diff --git a/src/shared/utils.c b/src/shared/utils.c index 8c1a20d..14bdd7d 100644 --- a/src/shared/utils.c +++ b/src/shared/utils.c @@ -725,6 +725,148 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* k return operation_ok; } +/* Read-only mirror of delete_extras_fd: records the paths that WOULD be removed + without unlinking anything. A child directory is reported after its own + reportable children (depth-first), matching the delete pass's ordering. */ +static bool list_extras_fd(int dirfd, const char* rel_path, const PathIndex* keep, + const PathIndex* dirs, ArrayList* out, size_t* recorded, + const DeleteSkipEntry* skips, int skip_count, bool parent_deletable, + bool* all_removed) { + int scanfd = openat(dirfd, ".", O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); + if (scanfd < 0) + return false; + DIR* dir = fdopendir(scanfd); + if (!dir) { + close(scanfd); + return false; + } + bool operation_ok = true; + bool local_survives = false; + bool deletable = parent_deletable || is_synced_dir(dirs, rel_path); + const struct dirent* entry; + while ((entry = readdir(dir)) != NULL) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char* child_rel = path_cat((char*)rel_path, entry->d_name); + if (!child_rel) { + operation_ok = false; + continue; + } + if (path_under_skip_prefix(child_rel, rel_path[0] == '\0', skips, skip_count)) { + local_survives = true; + free(child_rel); + continue; + } + struct stat st; + if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) { + if (errno != ENOENT) + operation_ok = false; + free(child_rel); + continue; + } + if (S_ISDIR(st.st_mode)) { + int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); + bool child_all_removed = false; + if (childfd >= 0) { + if (!list_extras_fd(childfd, child_rel, keep, dirs, out, recorded, skips, skip_count, + deletable, &child_all_removed)) + operation_ok = false; + close(childfd); + } else if (errno != ENOENT) { + operation_ok = false; + } + bool child_synced = dirs && path_index_contains(dirs, child_rel); + if (child_synced || keep_is_dir(keep, child_rel)) { + local_survives = true; + } else if (child_all_removed && deletable) { + size_t len = strlen(child_rel); + char* copy = malloc(len + 2); + if (!copy) { + operation_ok = false; + } else { + memcpy(copy, child_rel, len); + copy[len] = '/'; + copy[len + 1] = '\0'; + if (!array_list_add(out, copy)) { + free(copy); + operation_ok = false; + } else { + (*recorded)++; + } + } + } else { + local_survives = true; + } + } else { + bool found = keep_is_file(keep, child_rel); + if (found || !deletable) { + local_survives = true; + } else { + char* copy = str_dup(child_rel); + if (!copy || !array_list_add(out, copy)) { + free(copy); + operation_ok = false; + } else { + (*recorded)++; + } + } + } + free(child_rel); + } + closedir(dir); + *all_removed = !local_survives; + return operation_ok; +} + +bool delete_extras_list(const char* dest_root, const ArrayList* manifest, + const ArrayList* synced_dirs, const DeleteSkipEntry* skips, int skip_count, + ArrayList* out, size_t* count_out) { + if (count_out) + *count_out = 0; + if (!manifest || !out) + return false; + PathIndex keep; + if (!build_keep_index(manifest, &keep)) + return false; + PathIndex dirs; + bool have_dirs = synced_dirs != NULL; + if (have_dirs && + !path_index_build(&dirs, (const char* const*)synced_dirs->items, (size_t)synced_dirs->size)) { + path_index_free(&keep); + return false; + } + int rootfd; + int root_fd = utils_get_authorized_root_fd(); + if (root_fd >= 0) { + if (utils_get_authorized_root_path()) + rootfd = utils_open_authorized_destination(dest_root); + else if (dest_root == NULL) + rootfd = dup(root_fd); + else + rootfd = -1; + } else { + rootfd = open(dest_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); + } + if (rootfd < 0) { + path_index_free(&keep); + if (have_dirs) + path_index_free(&dirs); + return false; + } + bool all_removed = false; + size_t recorded = 0; + bool ok = list_extras_fd(rootfd, "", &keep, have_dirs ? &dirs : NULL, out, &recorded, skips, + skip_count, false, &all_removed); + if (close(rootfd) != 0) + ok = false; + path_index_free(&keep); + if (have_dirs) + path_index_free(&dirs); + if (count_out) + *count_out = recorded; + return ok; +} + DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest, const ArrayList* synced_dirs, size_t max_delete, const DeleteSkipEntry* skips, int skip_count, diff --git a/src/shared/utils.h b/src/shared/utils.h index 8321fa3..578706f 100644 --- a/src/shared/utils.h +++ b/src/shared/utils.h @@ -138,6 +138,14 @@ DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* m const ArrayList* synced_dirs, size_t max_delete, const DeleteSkipEntry* skips, int skip_count, size_t* deleted_out, size_t* skipped_out); +/* Read-only companion to delete_extras_limited: walk the destination exactly as + the delete pass would and APPEND (strdup'd) destination-relative paths that + WOULD be removed, without touching disk. Used for -n/--dry-run --delete + would-delete reporting. Returns true on a clean walk; the caller owns the + strings appended to `out` and receives their count in *count_out. */ +bool delete_extras_list(const char* dest_root, const ArrayList* manifest, + const ArrayList* synced_dirs, const DeleteSkipEntry* skips, int skip_count, + ArrayList* out, size_t* count_out); bool delete_extras(const char* dest_root, const ArrayList* manifest); /* Open the existing destination directory at `dest_root`, confined to the authorized root with an O_NOFOLLOW component walk (the same confinement the diff --git a/tests/integration/test_fault_injection.py b/tests/integration/test_fault_injection.py index fd93651..1c2fd22 100644 --- a/tests/integration/test_fault_injection.py +++ b/tests/integration/test_fault_injection.py @@ -36,7 +36,7 @@ from common import ( # noqa: E402 verify_transfer, ) -PROTOCOL_VERSION = b"2.24.0" +PROTOCOL_VERSION = b"2.25.0" STATUS_MANIFEST = 5 STATUS_OK = 0 diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index de0c4f0..2a382d0 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -1855,8 +1855,8 @@ class TestDelete: ) assert result.returncode == 0, f"Exit {result.returncode}: {result.stderr[:100]}" output = result.stdout + result.stderr - assert "Sent " in output and "MB" in output, "--progress produced no stable byte marker" - assert "Done." in output, "--progress did not report completion" + assert "sending incremental file list" in output, "--progress produced no rsync header" + assert "(xfr#" in output, "--progress produced no per-file xfr block" def test_human_readable_stats(self, shared_server): clean_dir(DEST_DIR) @@ -1891,8 +1891,8 @@ class TestDelete: ) assert result.returncode == 0, f"Exit {result.returncode}: {result.stderr[:100]}" output = result.stdout + result.stderr - assert "Sent " in output - assert "Done." in output + assert "sending incremental file list" in output + assert "(xfr#" in output class TestInfo: diff --git a/tests/integration/test_output_parity.py b/tests/integration/test_output_parity.py index 0843ca2..0d49117 100644 --- a/tests/integration/test_output_parity.py +++ b/tests/integration/test_output_parity.py @@ -12,7 +12,7 @@ import sys import pytest sys.path.insert(0, os.path.dirname(__file__)) -from common import TEST_DATA_DIR, run_client, clean_dir, get_dest_received_dir +from common import TEST_DATA_DIR, run_client, clean_dir, get_dest_received_dir, ServerManager RSYNC = shutil.which("rsync") requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed") @@ -280,3 +280,174 @@ class TestListOnlyParity: assert fast_lines == rsync_lines, ( f"rsync={rsync_lines}\nfastsync={fast_lines}" ) + + +def _make_one_file(root, name="f.bin", size=100): + clean_dir(root) + with open(os.path.join(root, name), "wb") as fh: + fh.write(bytes((i * 7 + 3) & 0xFF for i in range(size))) + + +class TestWireStatsParity: + """Wire-counter output parity: --out-format %b/%c/%C, --progress and + --stats versus real rsync 3.4.1.""" + + @requires_rsync + @pytest.mark.ci + def test_out_format_checksum_matches_rsync(self, shared_server): + """%C (whole-file xxh128, seed 0) is protocol-independent, so the full + `%C %l %n` line must be byte-identical to rsync.""" + source = os.path.join(TEST_DATA_DIR, "wire_ck_src") + dest = os.path.join(TEST_DATA_DIR, "wire_ck_dst") + rdst = os.path.join(TEST_DATA_DIR, "wire_ck_rdst") + _make_one_file(source, "f.bin", 200000) + clean_dir(dest) + clean_dir(rdst) + fmt = "%C %l %n" + rsync_result = _rsync(["-a", "--out-format=" + fmt, source + "/", rdst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + result, _ = run_client(source, dest, flags=["-a", "--out-format=" + fmt], + port=shared_server.port) + assert result.returncode == 0, result.stderr[:300] + + def file_lines(text): + # Ignore the root directory entry: fastsync does not transfer the + # source-root dir itself (a separate pre-existing divergence). + return [ + line for line in text.splitlines() if not line.rsplit(" ", 1)[-1].endswith("/") + ] + + assert file_lines(result.stdout) == file_lines(rsync_result.stdout), ( + f"rsync={rsync_result.stdout!r} fastsync={result.stdout!r}" + ) + + @requires_rsync + @pytest.mark.ci + def test_out_format_b_is_wire_bytes(self, shared_server): + """%b is true transferred (wire) bytes, not the source length: it must + differ from %l (the source length) and exceed it for a framed transfer.""" + source = os.path.join(TEST_DATA_DIR, "wire_b_src") + dest = os.path.join(TEST_DATA_DIR, "wire_b_dst") + _make_one_file(source, "f.bin", 5000) + clean_dir(dest) + result, _ = run_client(source, dest, flags=["-a", "--out-format=%b %l %c"], + port=shared_server.port) + assert result.returncode == 0, result.stderr[:300] + line = result.stdout.strip() + parts = line.split() + assert len(parts) == 3 and all(p.isdigit() for p in parts), line + wire_b, src_l, wire_c = (int(p) for p in parts) + assert src_l == 5000, line + assert wire_b > src_l, f"%b must include wire framing: {line}" + + @requires_rsync + @pytest.mark.ci + def test_progress_first_frame_matches_rsync(self, shared_server): + """For a sub-32 KiB file the first --progress frame is deterministic + (0.00 kB/s, 0:00:00) and must be byte-identical to rsync's.""" + source = os.path.join(TEST_DATA_DIR, "wire_pg_src") + dest = os.path.join(TEST_DATA_DIR, "wire_pg_dst") + rdst = os.path.join(TEST_DATA_DIR, "wire_pg_rdst") + _make_one_file(source, "f.bin", 100) + clean_dir(dest) + clean_dir(rdst) + rsync_result = _rsync(["-a", "--progress", source + "/", rdst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + result, _ = run_client(source, dest, flags=["-a", "--progress"], + port=shared_server.port) + assert result.returncode == 0, result.stderr[:300] + + def frames(text): + # subprocess text mode normalizes \r to \n (universal newlines). + return [p for p in text.split("\n") if "%" in p] + + rsync_frames = frames(rsync_result.stdout) + fast_frames = frames(result.stdout) + assert rsync_frames and fast_frames, (rsync_result.stdout, result.stdout) + assert fast_frames[0] == rsync_frames[0], (rsync_frames[0], fast_frames[0]) + assert "(xfr#1," in fast_frames[-1], fast_frames[-1] + + @requires_rsync + @pytest.mark.ci + def test_stats_selected_lines_match_rsync(self, shared_server): + """The protocol-independent --stats lines must match rsync exactly.""" + source = os.path.join(TEST_DATA_DIR, "wire_st_src") + dest = os.path.join(TEST_DATA_DIR, "wire_st_dst") + rdst = os.path.join(TEST_DATA_DIR, "wire_st_rdst") + _make_one_file(source, "f.bin", 6000) + clean_dir(dest) + clean_dir(rdst) + rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + result, _ = run_client(source, dest, flags=["-a", "--stats"], + port=shared_server.port) + assert result.returncode == 0, result.stderr[:300] + keys = ( + "Number of regular files transferred", + "Total file size", + "Total transferred file size", + "Literal data", + "Matched data", + "Number of deleted files", + ) + + def pick(text): + out = {} + for line in text.splitlines(): + for key in keys: + if line.startswith(key + ":"): + out[key] = line + return out + + assert pick(result.stdout) == pick(rsync_result.stdout), ( + f"rsync={pick(rsync_result.stdout)} fastsync={pick(result.stdout)}" + ) + + @requires_rsync + @pytest.mark.ci + def test_dry_run_delete_lines_match_rsync(self): + """-n --delete emits transfer-relative `*deleting` lines like rsync.""" + source = os.path.join(TEST_DATA_DIR, "wire_del_src") + dest = os.path.join(TEST_DATA_DIR, "wire_del_dst") + rdst = os.path.join(TEST_DATA_DIR, "wire_del_rdst") + clean_dir(source) + clean_dir(dest) + clean_dir(rdst) + with open(os.path.join(source, "a.txt"), "wb") as fh: + fh.write(b"a\n") + for root, entries in ( + (rdst, {"extra.txt": b"x\n"}), + (rdst, {"sub/y.txt": b"y\n", "extradir/z.txt": b"z\n"}), + ): + for rel, data in entries.items(): + full = os.path.join(root, rel) + os.makedirs(os.path.dirname(full), exist_ok=True) + with open(full, "wb") as fh: + fh.write(data) + # FastSync mirrors the source's absolute path under dest. + received = get_dest_received_dir(dest, source) + for rel, data in ( + ("extra.txt", b"x\n"), + ("sub/y.txt", b"y\n"), + ("extradir/z.txt", b"z\n"), + ): + full = os.path.join(received, rel) + os.makedirs(os.path.dirname(full), exist_ok=True) + with open(full, "wb") as fh: + fh.write(data) + + rsync_result = _rsync(["-a", "-n", "--delete", "-i", source + "/", rdst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + rsync_del = sorted( + line for line in rsync_result.stdout.splitlines() if line.startswith("*deleting") + ) + # The shared session server refuses deletion; start one that allows it. + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, flags=["-a", "-n", "--delete", "-i"], + port=server.port) + assert result.returncode == 0, result.stderr[:300] + fast_del = sorted( + line for line in result.stdout.splitlines() if line.startswith("*deleting") + ) + assert fast_del == rsync_del, f"rsync={rsync_del}\nfastsync={fast_del}" diff --git a/tests/integration/test_preflight.py b/tests/integration/test_preflight.py index 547d8f0..01aeaaf 100644 --- a/tests/integration/test_preflight.py +++ b/tests/integration/test_preflight.py @@ -94,14 +94,14 @@ def _seed_protocol_source(source): class TestProtocol: @pytest.mark.ci def test_protocol_current_version_accepted(self, shared_server): - """--protocol=2.24.0 (the current PROTOCOL_VERSION) is accepted and the + """--protocol=2.25.0 (the current PROTOCOL_VERSION) is accepted and the transfer completes normally.""" source = os.path.join(TEST_DATA_DIR, "proto_ok_src") dest = os.path.join(TEST_DATA_DIR, "proto_ok_dst") shutil.rmtree(dest, ignore_errors=True) os.makedirs(dest) _seed_protocol_source(source) - result, _ = run_client(source, dest, flags=["--protocol=2.24.0"], + result, _ = run_client(source, dest, flags=["--protocol=2.25.0"], port=shared_server.port) assert result.returncode == 0, \ f"--protocol current run failed: {(result.stderr or result.stdout)[:400]}" diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index 01bc272..ace49fb 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -317,7 +317,7 @@ static void test_parse_args_protocol_accept_current() { Config* cfg = valid_client_config(); EXPECT_NOT_NULL(cfg); char* argv_equals[] = {"fastsync", "--source-dir", "/src", - "--dest-dir", "/dst", "--protocol=2.24.0"}; + "--dest-dir", "/dst", "--protocol=2.25.0"}; int positional_args[2]; int positional_count = 0; EXPECT_EQ_INT(parse_args(cfg, 6, argv_equals, positional_args, &positional_count), 0); @@ -327,7 +327,7 @@ static void test_parse_args_protocol_accept_current() { cfg = valid_client_config(); EXPECT_NOT_NULL(cfg); char* argv_space[] = {"fastsync", "--source-dir", "/src", "--dest-dir", - "/dst", "--protocol", "2.24.0"}; + "/dst", "--protocol", "2.25.0"}; positional_count = 0; EXPECT_EQ_INT(parse_args(cfg, 7, argv_space, positional_args, &positional_count), 0); EXPECT_EQ_STR(cfg->version, PROTOCOL_VERSION); diff --git a/tests/test_config.c b/tests/test_config.c index b506bad..668b7e5 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -2775,14 +2775,14 @@ static void golden_config_populate(Config* c) { c->copy_as_gid = 222; } -/* The pinned golden frame (protocol 2.24.0). The values below are the only +/* The pinned golden frame (protocol 2.25.0). The values below are the only * thing that ties the generated table to the historical wire format; update * them ONLY with a PROTOCOL_VERSION bump and a documented reason. The 2.24.0 - * per-directory delete-plan wave changes only the version string in the config - * frame (the frame layout itself is unchanged from 2.23.0); the byte-exact hash - * is recomputed for the new version bytes. */ -#define GOLDEN_WIRE_LEN 697 -#define GOLDEN_WIRE_HASH 13736055061412501670ULL + * per-directory delete-plan wave changed only the version string in the config + * frame; the 2.25.0 wire-stats wave appends one report_stats bool. The + * byte-exact values are recomputed for the merged layout. */ +#define GOLDEN_WIRE_LEN 701 +#define GOLDEN_WIRE_HASH 16170466870400670271ULL static unsigned long long fnv1a_64(const unsigned char* buf, size_t len) { unsigned long long h = 1469598103934665603ULL; diff --git a/tests/test_fuzz_smoke.c b/tests/test_fuzz_smoke.c index d1965ac..279ff76 100644 --- a/tests/test_fuzz_smoke.c +++ b/tests/test_fuzz_smoke.c @@ -18,9 +18,10 @@ /* P8 config-frame tail: super_mode (4) + copy-as presence (4) + uid (4) + gid (4). */ #define P8_TAIL_BYTES 16 -/* Protocol 2.23.0 appends one trailing bool (report_dest_info) AFTER the P8 - * tail, so the P8 fields sit this many bytes before the end of the frame. */ -#define OUTPUT_TAIL_BYTES 4 +/* Protocol 2.25.0 appends two trailing bools (report_dest_info, report_stats) + * AFTER the P8 tail, so the P8 fields sit this many bytes before the end of the + * frame. */ +#define OUTPUT_TAIL_BYTES 8 /* Smoke test for chunk_deserialize fuzz target */ static void test_fuzz_chunk_deserialize() {