feat(protocol): client-message channel and rsync partial exit 23 (2.30.0)
This commit is contained in:
+15
-2
@@ -83,7 +83,7 @@ typedef struct {
|
||||
typedef enum SuperMode { SUPER_MODE_AUTO = 0, SUPER_MODE_ON = 1, SUPER_MODE_OFF = 2 } SuperMode;
|
||||
|
||||
/* ===========================================================================
|
||||
* Config wire-field table (single source of truth for protocol 2.29.0).
|
||||
* Config wire-field table (single source of truth for protocol 2.30.0).
|
||||
*
|
||||
* Every field below crosses the wire. The table is the ONLY place a
|
||||
* serialized field is named: config.h expands CONFIG_WIRE_FIELDS() to declare
|
||||
@@ -1084,7 +1084,20 @@ typedef struct Config {
|
||||
* version must bump; the strict same-version handshake (config_receive rejects a
|
||||
* mismatched version before parsing anything else) keeps a 2.29 client and a
|
||||
* 2.28 server from ever reaching that state. */
|
||||
#define PROTOCOL_VERSION "2.29.0"
|
||||
/* (11) Client-message channel + partial exit (protocol 2.30.0): the
|
||||
* config-frame LAYOUT is unchanged (no new config field), but the frame stream
|
||||
* gains two statuses. STATUS_CLIENT_MSG (client->server) carries a bounded,
|
||||
* length-prefixed diagnostic string so a client running with --stderr=client
|
||||
* (rsync's --no-msgs2stderr spelling) can forward its own diagnostics to the
|
||||
* server's stderr. STATUS_PARTIAL (receiver->client) is the terminal status
|
||||
* sent instead of STATUS_OK when a per-entry receiver failure (e.g. an
|
||||
* unprivileged --devices mknod) did not abort the stream; the sender exits 23
|
||||
* (rsync's partial transfer) and still removes successfully transferred
|
||||
* --remove-source-files sources. A 2.29 peer that does not know these status
|
||||
* values would reject them as an unknown status and tear the connection down,
|
||||
* so the protocol version must bump; the strict same-version handshake keeps a
|
||||
* 2.30 client and a 2.29 server from ever reaching that state. */
|
||||
#define PROTOCOL_VERSION "2.30.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
|
||||
|
||||
+80
-20
@@ -16,6 +16,7 @@ static bool info_flags_explicit = false;
|
||||
static FILE* log_fp = NULL;
|
||||
static _Thread_local bool eight_bit_output;
|
||||
static LogStderrMode stderr_mode = LOG_STDERR_ERRORS;
|
||||
static LogClientMsgSink client_msg_sink = NULL;
|
||||
|
||||
/* Serializes access to log_fp and makes each emitted line atomic: the
|
||||
* timestamp prefix, formatted body, and trailing newline are written as one
|
||||
@@ -77,6 +78,51 @@ LogStderrMode log_get_stderr_mode(void) {
|
||||
return stderr_mode;
|
||||
}
|
||||
|
||||
void log_set_client_msg_sink(LogClientMsgSink sink) {
|
||||
client_msg_sink = sink;
|
||||
}
|
||||
|
||||
LogClientMsgSink log_get_client_msg_sink(void) {
|
||||
return client_msg_sink;
|
||||
}
|
||||
|
||||
/* Format just the message body (no prefix/newline) into a freshly allocated
|
||||
* buffer. Shared by log_message (which may hand the body to a client-message
|
||||
* sink) and log_client_message. Returns NULL on allocation/format failure. */
|
||||
static char* format_log_body(const char* format, va_list args) {
|
||||
va_list copy;
|
||||
va_copy(copy, args);
|
||||
int body_len = vsnprintf(NULL, 0, format, copy);
|
||||
va_end(copy);
|
||||
if (body_len < 0)
|
||||
return NULL;
|
||||
char* body = malloc((size_t)body_len + 1);
|
||||
if (!body)
|
||||
return NULL;
|
||||
vsnprintf(body, (size_t)body_len + 1, format, args);
|
||||
return body;
|
||||
}
|
||||
|
||||
/* Assemble a complete log line (prefix + body + newline) from an already
|
||||
* formatted body. Returns NULL on allocation failure. */
|
||||
static char* format_log_line_from_body(LogLevel log_level, const struct tm* t, const char* body) {
|
||||
char prefix[64];
|
||||
int prefix_len = snprintf(
|
||||
prefix, sizeof(prefix), "%04d-%02d-%02d %02d:%02d:%02d [%s]: ", t->tm_year + 1900,
|
||||
t->tm_mon + 1, t->tm_mday, t->tm_hour, t->tm_min, t->tm_sec, log_level_strings[log_level]);
|
||||
if (prefix_len < 0 || prefix_len >= (int)sizeof(prefix))
|
||||
return NULL;
|
||||
size_t body_len = strlen(body);
|
||||
char* line = malloc((size_t)prefix_len + body_len + 2); /* body + '\n' + NUL */
|
||||
if (!line)
|
||||
return NULL;
|
||||
memcpy(line, prefix, (size_t)prefix_len);
|
||||
memcpy(line + prefix_len, body, body_len);
|
||||
line[(size_t)prefix_len + body_len] = '\n';
|
||||
line[(size_t)prefix_len + body_len + 1] = '\0';
|
||||
return line;
|
||||
}
|
||||
|
||||
/* Format one complete log line (timestamp prefix + body + newline) into a
|
||||
* freshly allocated buffer. This is pure CPU/malloc work and must happen
|
||||
* OUTSIDE the log mutex: the mutex only guards the log_fp pointer, so a
|
||||
@@ -84,26 +130,11 @@ LogStderrMode log_get_stderr_mode(void) {
|
||||
* on allocation/formatting failure. */
|
||||
static char* format_log_line(LogLevel log_level, const struct tm* t, const char* format,
|
||||
va_list args) {
|
||||
char prefix[64];
|
||||
int prefix_len = snprintf(
|
||||
prefix, sizeof(prefix), "%04d-%02d-%02d %02d:%02d:%02d [%s]: ", t->tm_year + 1900,
|
||||
t->tm_mon + 1, t->tm_mday, t->tm_hour, t->tm_min, t->tm_sec, log_level_strings[log_level]);
|
||||
if (prefix_len < 0 || prefix_len >= (int)sizeof(prefix))
|
||||
char* body = format_log_body(format, args);
|
||||
if (!body)
|
||||
return NULL;
|
||||
va_list copy;
|
||||
va_copy(copy, args);
|
||||
int body_len = vsnprintf(NULL, 0, format, copy);
|
||||
va_end(copy);
|
||||
if (body_len < 0)
|
||||
return NULL;
|
||||
size_t total = (size_t)prefix_len + (size_t)body_len;
|
||||
char* line = malloc(total + 2); /* body bytes + '\n' + NUL */
|
||||
if (!line)
|
||||
return NULL;
|
||||
memcpy(line, prefix, (size_t)prefix_len);
|
||||
vsnprintf(line + prefix_len, (size_t)body_len + 1, format, args);
|
||||
line[total] = '\n';
|
||||
line[total + 1] = '\0';
|
||||
char* line = format_log_line_from_body(log_level, t, body);
|
||||
free(body);
|
||||
return line;
|
||||
}
|
||||
|
||||
@@ -121,6 +152,20 @@ static void emit_log_line(FILE* console, const char* line) {
|
||||
mtx_unlock(&log_mutex);
|
||||
}
|
||||
|
||||
void log_client_message(const char* message) {
|
||||
if (!message)
|
||||
return;
|
||||
time_t now = time(NULL);
|
||||
struct tm t;
|
||||
if (!localtime_r(&now, &t))
|
||||
return;
|
||||
char* line = format_log_line_from_body(LOG_LEVEL_INFO, &t, message);
|
||||
if (!line)
|
||||
return;
|
||||
emit_log_line(stderr, line);
|
||||
free(line);
|
||||
}
|
||||
|
||||
void log_message(LogLevel log_level, const char* format, ...) {
|
||||
if (log_level < current_log_level)
|
||||
return;
|
||||
@@ -138,8 +183,23 @@ void log_message(LogLevel log_level, const char* format, ...) {
|
||||
|
||||
va_list args;
|
||||
va_start(args, format);
|
||||
char* line = format_log_line(log_level, &t, format, args);
|
||||
char* body = format_log_body(format, args);
|
||||
va_end(args);
|
||||
if (!body)
|
||||
return;
|
||||
/* LOG_STDERR_CLIENT: hand the diagnostic to the client-message channel. A
|
||||
sink that takes ownership suppresses the local write; otherwise (no sink
|
||||
yet, or the peer connection is not up) fall through to local output so the
|
||||
diagnostic is never lost. */
|
||||
if (stderr_mode == LOG_STDERR_CLIENT) {
|
||||
LogClientMsgSink sink = client_msg_sink;
|
||||
if (sink && sink(body)) {
|
||||
free(body);
|
||||
return;
|
||||
}
|
||||
}
|
||||
char* line = format_log_line_from_body(log_level, &t, body);
|
||||
free(body);
|
||||
if (!line)
|
||||
return;
|
||||
emit_log_line(dest_io, line);
|
||||
|
||||
+20
-1
@@ -15,7 +15,11 @@
|
||||
#endif
|
||||
|
||||
typedef enum { LOG_LEVEL_DEBUG, LOG_LEVEL_INFO, LOG_LEVEL_WARNING, LOG_LEVEL_ERROR } LogLevel;
|
||||
typedef enum { LOG_STDERR_ERRORS, LOG_STDERR_ALL } LogStderrMode;
|
||||
/* --stderr=MODE destinations. ERRORS keeps errors on stderr and everything
|
||||
* else on stdout; ALL sends every message to stderr; CLIENT routes the client's
|
||||
* own diagnostics over the protocol stream to the peer's stderr (rsync's
|
||||
* --stderr=client / --no-msgs2stderr). */
|
||||
typedef enum { LOG_STDERR_ERRORS, LOG_STDERR_ALL, LOG_STDERR_CLIENT } LogStderrMode;
|
||||
|
||||
typedef enum {
|
||||
LOG_DEBUG_IO = 1u << 0,
|
||||
@@ -86,5 +90,20 @@ void log_set_8_bit_output(bool enabled);
|
||||
bool log_get_8_bit_output(void);
|
||||
void log_set_stderr_mode(LogStderrMode mode);
|
||||
LogStderrMode log_get_stderr_mode(void);
|
||||
/* Write a message a peer forwarded over the client-message channel to this
|
||||
* process's stderr (and log file), with the standard log prefix. Used by the
|
||||
* server side of rsync's --stderr=client. */
|
||||
void log_client_message(const char* message);
|
||||
|
||||
/* Sink for LOG_STDERR_CLIENT. log_message() passes the un-prefixed message
|
||||
* body to the installed sink; a `true` return means the sink took ownership
|
||||
* (e.g. queued it for protocol transmission) and the message must NOT also be
|
||||
* written locally. A `false` return (or a NULL sink) makes log_message fall
|
||||
* back to the normal local destination, so a diagnostic emitted before the peer
|
||||
* connection exists is never lost (rsync's documented fallback). The sink may
|
||||
* be called from any thread and must be tolerant of that. */
|
||||
typedef bool (*LogClientMsgSink)(const char* message);
|
||||
void log_set_client_msg_sink(LogClientMsgSink sink);
|
||||
LogClientMsgSink log_get_client_msg_sink(void);
|
||||
|
||||
#endif
|
||||
|
||||
@@ -53,6 +53,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
|
||||
context->dir_entries_mutex_init = false;
|
||||
atomic_init(&context->dir_count, 0);
|
||||
context->delete_limit = false;
|
||||
context->partial = false;
|
||||
int init = 0;
|
||||
if (config->use_metadata) {
|
||||
context->dir_entries = array_list_create(file_destroy);
|
||||
|
||||
@@ -134,6 +134,11 @@ typedef struct {
|
||||
deletion (STATUS_DELETE_LIMIT): the transfer succeeded and the process must
|
||||
exit 25 like rsync. Read by the caller after the sender thread is joined. */
|
||||
bool delete_limit;
|
||||
/* Set by the sender thread when the receiver reported STATUS_PARTIAL (a
|
||||
per-entry receiver failure that did not abort the stream): the transfer
|
||||
otherwise succeeded, successfully stored --remove-source-files sources were
|
||||
removed, and the process must exit 23 like rsync. Read after join. */
|
||||
bool partial;
|
||||
} PipelineContextSender;
|
||||
|
||||
/* `config` is borrowed and must outlive the context: destroy does NOT free it,
|
||||
|
||||
+19
-2
@@ -664,6 +664,10 @@ static const char* status_to_string(Status status) {
|
||||
return "DELETE_LIMIT";
|
||||
case STATUS_DEST_INFO:
|
||||
return "DEST_INFO";
|
||||
case STATUS_CLIENT_MSG:
|
||||
return "CLIENT_MSG";
|
||||
case STATUS_PARTIAL:
|
||||
return "PARTIAL";
|
||||
default:
|
||||
return "UNKNOWN";
|
||||
}
|
||||
@@ -672,10 +676,10 @@ static const char* status_to_string(Status status) {
|
||||
/* Reject a raw wire status outside the known enum range before it is handed to
|
||||
* callers, so an unknown/corrupt frame fails as a protocol error instead of
|
||||
* being silently interpreted as an unexpected-but-valid verdict. STATUS_OK is
|
||||
* the first enumerator and STATUS_STATS the last, so the range check accepts
|
||||
* the first enumerator and STATUS_PARTIAL the last, so the range check accepts
|
||||
* every status the protocol defines. */
|
||||
static bool status_is_valid(Status status) {
|
||||
return status >= STATUS_OK && status <= STATUS_STATS;
|
||||
return status >= STATUS_OK && status <= STATUS_PARTIAL;
|
||||
}
|
||||
|
||||
/* Shared string send/receive implementation. `redact` selects whether the
|
||||
@@ -1144,6 +1148,19 @@ bool send_error_detail(int fd, const char* message) {
|
||||
return send_status(fd, STATUS_ERROR_DETAIL) && send_str(fd, message);
|
||||
}
|
||||
|
||||
bool send_client_message(int fd, const char* message) {
|
||||
if (!message)
|
||||
message = "";
|
||||
char bounded[MAX_CLIENT_MSG_BYTES + 1];
|
||||
size_t len = strlen(message);
|
||||
if (len > MAX_CLIENT_MSG_BYTES) {
|
||||
memcpy(bounded, message, MAX_CLIENT_MSG_BYTES);
|
||||
bounded[MAX_CLIENT_MSG_BYTES] = '\0';
|
||||
message = bounded;
|
||||
}
|
||||
return send_status(fd, STATUS_CLIENT_MSG) && send_str(fd, message);
|
||||
}
|
||||
|
||||
const char* protocol_last_error(void) {
|
||||
return io_error_detail;
|
||||
}
|
||||
|
||||
+30
-1
@@ -15,6 +15,12 @@
|
||||
* this for a rejection and the detail frame stays a small, fixed bound. */
|
||||
#define MAX_ERROR_DETAIL_BYTES 4096
|
||||
|
||||
/* Hard cap on a client diagnostic forwarded over the STATUS_CLIENT_MSG channel
|
||||
* (protocol 2.30.0, rsync's --stderr=client). The body is reused from the
|
||||
* bounded-string wire helper and sliced to this many bytes before it is sent,
|
||||
* so a peer can never be made to retain more than this per message. */
|
||||
#define MAX_CLIENT_MSG_BYTES 4096
|
||||
|
||||
/* Maximum uncompressed file payload accepted by the receiver's whole-file
|
||||
* paths. A single whole file is charged against the per-connection memory
|
||||
* reservation (MAX_CONNECTION_MEMORY) and against the server allocation
|
||||
@@ -240,7 +246,25 @@ enum NET_STATUS {
|
||||
* 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
|
||||
STATUS_STATS,
|
||||
/* Client diagnostic channel (protocol 2.30.0, rsync's --stderr=client /
|
||||
* --no-msgs2stderr). When the client's --stderr mode is `client`, the
|
||||
* client forwards its own diagnostics over this client->server frame
|
||||
* (STATUS_CLIENT_MSG followed by a bounded length-prefixed string, capped at
|
||||
* MAX_CLIENT_MSG_BYTES) instead of writing them to its local stderr. The
|
||||
* receiver reads the string and writes it to the server's stderr (respecting
|
||||
* the server log destination). Appended after STATUS_STATS so no existing
|
||||
* status is renumbered. */
|
||||
STATUS_CLIENT_MSG,
|
||||
/* Receiver-side partial transfer (protocol 2.30.0). Sent by the receiver as
|
||||
* the terminal status INSTEAD of STATUS_OK when one or more entries failed
|
||||
* per-entry without aborting the stream (currently a --devices mknod
|
||||
* EPERM/EACCES). The transfer otherwise succeeded and every successfully
|
||||
* stored file was acknowledged, so the sender may still remove
|
||||
* --remove-source-files sources; the sender maps this to rsync's exit code
|
||||
* 23 ("partial transfer due to error"), distinct from a fatal STATUS_ERROR.
|
||||
* Appended after STATUS_CLIENT_MSG so no existing status is renumbered. */
|
||||
STATUS_PARTIAL
|
||||
};
|
||||
|
||||
void io_set_fds(int read_fd, int write_fd);
|
||||
@@ -341,6 +365,11 @@ bool receive_status(int file_descriptor, Status* status);
|
||||
* length-prefixed string. Over-long messages are sliced and NULL is treated
|
||||
* as "". Returns false if the status or the string could not be sent. */
|
||||
bool send_error_detail(int file_descriptor, const char* message);
|
||||
/* Send STATUS_CLIENT_MSG followed by a bounded (<= MAX_CLIENT_MSG_BYTES)
|
||||
* length-prefixed string carrying a client diagnostic. Over-long messages are
|
||||
* sliced and NULL is treated as "". Returns false if the status or the string
|
||||
* could not be sent. */
|
||||
bool send_client_message(int file_descriptor, const char* message);
|
||||
/* Human-readable reason captured from the most recent STATUS_ERROR_DETAIL
|
||||
* received on this thread, or "" when the last status was a bare STATUS_ERROR
|
||||
* (or no detail was seen). Thread-local, and valid until the next non-keepalive
|
||||
|
||||
Reference in New Issue
Block a user