Fall back through the different transfer methods in a consistent way, so that after a copy_file_range failure, splice through intermediate pipes can still be tried (#196).

This commit is contained in:
Andrew Wood
2026-06-04 23:44:56 +01:00
parent 07c1798b5e
commit 14ee298bc3
4 changed files with 145 additions and 164 deletions
+8
View File
@@ -58,3 +58,11 @@ typedef bool _Bool;
#define __attribute__(x) /* GCC-only feature */
#endif
#endif
/* Support for the fallthrough attribute is not universal. */
/* From https://stackoverflow.com/questions/45349079/how-to-use-attribute-fallthrough-correctly-in-gcc */
#if defined(__GNUC__) && __GNUC__ >= 7
#define FALL_THROUGH __attribute__ ((fallthrough))
#else
#define FALL_THROUGH /*@-noeffect@*/ ((void)0) /*@+noeffect@*/
#endif /* __GNUC__ >= 7 */
+2 -1
View File
@@ -85,7 +85,8 @@ typedef enum {
PV_TRANSFERMETHOD_READWRITE, /* read() + write() */
PV_TRANSFERMETHOD_SPLICE, /* splice() */
PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE, /* splice() through an intermediate input pipe */
PV_TRANSFERMETHOD_COPY_FILE_RANGE /* copy_file_range() */
PV_TRANSFERMETHOD_COPY_FILE_RANGE, /* copy_file_range() */
PV_TRANSFERMETHOD__UNSET
} pvtransfermethod_t;
/*
+3
View File
@@ -42,6 +42,9 @@ pvdisplay_bytecount_t pv_formatter_buffer_percent(pvformatter_args_t args)
case PV_TRANSFERMETHOD_COPY_FILE_RANGE:
(void) pv_snprintf(content, sizeof(content), "{%s}", "--->");
break;
case PV_TRANSFERMETHOD__UNSET:
(void) pv_snprintf(content, sizeof(content), "{%s}", "????");
break;
}
return pv_formatter_segmentcontent(content, args);
+132 -163
View File
@@ -114,19 +114,13 @@ static int is_data_ready(int fd_in, /*@null@ */ bool *fd_in_ready, int fd_out, /
* TRANSFER_READ_TIMEOUT seconds have not yet elapsed, and more data is
* available to read according to is_data_ready().
*
* Sets state->transfer.method to PV_TRANSFERMETHOD_READWRITE as a side
* effect, so that if this is called as a fallback from another transfer
* method, the original caller can find out.
*
* Returns the total number of bytes read, or negative on error.
*/
static ssize_t pv__transfer__read_repeated(pvstate_t state, int fd, char *buf, size_t count)
static ssize_t pv__transfer__read_repeated(int fd, char *buf, size_t count)
{
struct timespec start_time;
ssize_t total_read;
state->transfer.method = PV_TRANSFERMETHOD_READWRITE;
memset(&start_time, 0, sizeof(start_time));
pv_elapsedtime_read(&start_time);
@@ -420,39 +414,20 @@ static ssize_t pv__transfer__splice_via_intermediate(pvstate_t state, int input_
*
* If splice() could not be used, sets state->transfer.splice_failed_fd to
* input_fd so splice() won't be tried again until the next input file, and
* then calls pv__transfer__read_repeated() to read up to
* "fallback_read_amount" bytes into the buffer, returning its result.
*
* Sets state->transfer.method to either PV_TRANSFERMETHOD_SPLICE or
* PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE as a side effect; if it falls back
* to pv__transfer__read_repeated(), that function then sets
* state->transfer.method to PV_TRANSFERMETHOD_READWRITE, so the caller can
* check which transfer method was actually used.
* sets *transfer_attempted to false, indicating that the caller should try
* another method, then returns 0.
*
* Returns the total number of bytes transferred, or negative on error.
*/
static ssize_t pv__transfer__splice_repeated(pvstate_t state, int input_fd, int output_fd, char *buf,
size_t fallback_read_amount, off_t max_to_read, off_t max_to_write)
static ssize_t pv__transfer__splice_repeated(pvstate_t state, int input_fd, int output_fd,
off_t max_to_read, off_t max_to_write, bool *transfer_attempted)
{
struct timespec start_time;
size_t bytes_to_splice;
ssize_t total_spliced;
bool use_intermediate_pipe;
debug("%d->%d, fallback_read_amount=%lu, max_to_read=%lld, max_to_write=%lld", input_fd, output_fd,
fallback_read_amount, max_to_read, max_to_write);
/*
* Early return via pv__transfer__read_repeated() if splice() is
* turned off, or if line mode is active, or if there's no output
* fd, or if splice() already failed on this input file descriptor,
* or if there's anything waiting in the transfer buffer.
*/
if (state->control.no_splice || state->control.linemode || (output_fd < 0)
|| (input_fd == state->transfer.splice_failed_fd)
|| (state->transfer.to_write > 0)) {
return pv__transfer__read_repeated(state, input_fd, buf, fallback_read_amount);
}
debug("%d->%d, max_to_read=%lld, max_to_write=%lld", input_fd, output_fd, max_to_read, max_to_write);
/*
* Cap the transfer amount if applicable.
@@ -484,16 +459,12 @@ static ssize_t pv__transfer__splice_repeated(pvstate_t state, int input_fd, int
/*
* Use an intermediate pipe if one is available and that transfer
* method is selected. Either way, update state->transfer.method to
* the method being used.
* method is selected.
*/
use_intermediate_pipe = false;
if ((PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE == state->transfer.method)
&& (-1 != state->transfer.intermediate_pipe[0]) && (-1 != state->transfer.intermediate_pipe[1])) {
use_intermediate_pipe = true;
state->transfer.method = PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE;
} else {
state->transfer.method = PV_TRANSFERMETHOD_SPLICE;
}
memset(&start_time, 0, sizeof(start_time));
@@ -537,8 +508,9 @@ static ssize_t pv__transfer__splice_repeated(pvstate_t state, int input_fd, int
/*
* If the splice failed, turn it off for this input file
* descriptor. Then, if nothing has been spliced so far,
* return the result of an ordinary read, otherwise return
* the amount spliced.
* return zero with *transfer_attempted set to false so
* another method is immediately attempted by the caller;
* otherwise return the amount spliced.
*/
if (nspliced < 0) {
debug("%s %d: %s: %s", "fd", input_fd, "disabling splice after failure", strerror(errno));
@@ -562,7 +534,8 @@ static ssize_t pv__transfer__splice_repeated(pvstate_t state, int input_fd, int
}
state->transfer.splice_failed_fd = input_fd;
if (0 == total_spliced) {
return pv__transfer__read_repeated(state, input_fd, buf, fallback_read_amount);
*transfer_attempted = false;
return 0;
} else {
return total_spliced;
}
@@ -630,39 +603,20 @@ static ssize_t pv__transfer__splice_repeated(pvstate_t state, int input_fd, int
*
* If copy_file_range() could not be used, sets
* state->transfer.copy_file_range_failed_fd to input_fd so it won't be
* tried again until the next input file, and then calls
* pv__transfer__read_repeated() to read up to "fallback_read_amount" bytes
* into the buffer, returning its result.
*
* Sets state->transfer.method to PV_TRANSFERMETHOD_COPY_FILE_RANGE as a
* side effect; if it falls back to pv__transfer__read_repeated(), that
* function then sets state->transfer.method to PV_TRANSFERMETHOD_READWRITE,
* so the caller can check which transfer method was actually used.
* tried again until the next input file, and sets *transfer_attempted to
* false, indicating that the caller should try another method, then returns
* 0.
*
* Returns the total number of bytes transferred, or negative on error.
*/
static ssize_t pv__transfer__copy_file_range_repeated(pvstate_t state, int input_fd, int output_fd, char *buf,
size_t fallback_read_amount, off_t max_to_read,
off_t max_to_write)
static ssize_t pv__transfer__copy_file_range_repeated(pvstate_t state, int input_fd, int output_fd,
off_t max_to_read, off_t max_to_write, bool *transfer_attempted)
{
struct timespec start_time;
size_t bytes_to_copy;
ssize_t total_copied;
debug("%d->%d, fallback_read_amount=%lu, max_to_read=%ld, max_to_write=%ld", input_fd, output_fd,
fallback_read_amount, max_to_read, max_to_write);
/*
* Early return via pv__transfer__read_repeated() if this feature is
* turned off, or if line mode is active, or if there's no output
* fd, or if it already failed on this input file descriptor, or if
* there's anything waiting in the transfer buffer.
*/
if (state->control.no_splice || state->control.linemode || (output_fd < 0)
|| (input_fd == state->transfer.copy_file_range_failed_fd)
|| (state->transfer.to_write > 0)) {
return pv__transfer__read_repeated(state, input_fd, buf, fallback_read_amount);
}
debug("%d->%d, max_to_read=%ld, max_to_write=%ld", input_fd, output_fd, max_to_read, max_to_write);
/*
* Cap the transfer amount if applicable.
@@ -692,9 +646,6 @@ static ssize_t pv__transfer__copy_file_range_repeated(pvstate_t state, int input
return 0;
}
/* Explicitly set the transfer method being used. */
state->transfer.method = PV_TRANSFERMETHOD_COPY_FILE_RANGE;
memset(&start_time, 0, sizeof(start_time));
pv_elapsedtime_read(&start_time);
@@ -728,16 +679,18 @@ static ssize_t pv__transfer__copy_file_range_repeated(pvstate_t state, int input
/*
* If the copy failed, turn it off for this input file
* descriptor. Then, if nothing has been copied so far,
* return the result of an ordinary read, otherwise return
* the amount copied.
* descriptor. Then, if nothing has been spliced so far,
* return zero with *transfer_attempted set to false so
* another method is immediately attempted by the caller;
* otherwise return the amount spliced.
*/
if (ncopied < 0) {
debug("%s %d: %s: %s", "fd", input_fd, "disabling copy_file_range after failure",
strerror(errno));
state->transfer.copy_file_range_failed_fd = input_fd;
if (0 == total_copied) {
return pv__transfer__read_repeated(state, input_fd, buf, fallback_read_amount);
*transfer_attempted = false;
return 0;
} else {
return total_copied;
}
@@ -826,6 +779,7 @@ static bool pv__transfer_read(pvstate_t state, int input_fd, bool *eof_in, bool
{
bool do_not_skip_errors;
bool zero_transfer_cap;
bool transfer_attempted;
size_t max_buffer_available;
off_t max_to_read, max_to_read_into_buffer;
off_t amount_to_skip, amount_skipped, orig_offset, skip_offset;
@@ -855,8 +809,6 @@ static bool pv__transfer_read(pvstate_t state, int input_fd, bool *eof_in, bool
max_to_read = state->control.size - state->transfer.total_bytes_read;
}
nread = 0;
/*
* Flag if the amount to read is going to be zero, so that a
* zero-sized read is not then mistakenly taken as EOF.
@@ -880,103 +832,120 @@ static bool pv__transfer_read(pvstate_t state, int input_fd, bool *eof_in, bool
if (max_to_read_into_buffer < 0)
max_to_read_into_buffer = 0;
/* Determine which transfer method to attempt. */
state->transfer.method = PV_TRANSFERMETHOD_READWRITE;
#ifdef HAVE_SPLICE
/*
* Attempt to use splice() only if splice() is not turned off, and
* line mode is inactive, and there's an output fd, and there's
* nothing waiting to be written from the transfer buffer, and
* splice() hasn't already failed on this input file descriptor.
*/
if (!(state->control.no_splice || state->control.linemode || (output_fd < 0)
|| (state->transfer.to_write > 0)
|| (input_fd == state->transfer.splice_failed_fd)
)) {
state->transfer.method = PV_TRANSFERMETHOD_SPLICE;
/*
* Splice through an intermediate input pipe if neither the
* input nor the output is a pipe, and an intermediate pipe
* is ready to use.
*/
if (!(state->status.output_is_pipe || state->status.current_input_is_pipe)
&& (-1 != state->transfer.intermediate_pipe[0]) && (-1 != state->transfer.intermediate_pipe[1])) {
state->transfer.method = PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE;
}
}
#endif /* HAVE_SPLICE */
#ifdef HAVE_COPY_FILE_RANGE
/*
* Attempt to use copy_file_range() only if both the input and the
* output are regular files, and splice() is not turned off (since
* it's one option for both this and splice), and line mode is
* inactive, and there's an output fd, and there's nothing waiting
* to be written from the transfer buffer, and copy_file_range()
* hasn't already failed on this input file descriptor.
*/
if (state->status.current_input_is_file && state->status.output_is_file
&& !(state->control.no_splice || state->control.linemode || (output_fd < 0)
|| (state->transfer.to_write > 0)
|| (input_fd == state->transfer.copy_file_range_failed_fd)
)) {
state->transfer.method = PV_TRANSFERMETHOD_COPY_FILE_RANGE;
}
#endif /* HAVE_COPY_FILE_RANGE */
debug
("%d->%d, max_to_write=%lld, max_to_read=%lld, max_to_read_into_buffer=%lld, max_buffer_available=%lld, method=%d",
input_fd, output_fd, max_to_write, max_to_read, max_to_read_into_buffer, max_buffer_available,
state->transfer.method);
("%d->%d, max_to_write=%lld, max_to_read=%lld, max_to_read_into_buffer=%lld, max_buffer_available=%lld",
input_fd, output_fd, max_to_write, max_to_read, max_to_read_into_buffer, max_buffer_available);
/* Determine which transfer method to attempt. */
/*
* Transfer data using the appropriate method.
*
* If one method falls back to another, it will update
* state->transfer.method, so that variable may have changed after
* this block.
* Use read/write alone, if in no-splice mode, or in line mode, or
* there's no output, or there's something in the buffer - in all
* these cases a read buffer is needed.
*/
switch (state->transfer.method) {
case PV_TRANSFERMETHOD_READWRITE:
nread =
pv__transfer__read_repeated(state, input_fd,
state->transfer.transfer_buffer + state->transfer.read_position,
(size_t) max_to_read_into_buffer);
break;
#ifdef HAVE_SPLICE
case PV_TRANSFERMETHOD_SPLICE:
case PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE:
nread =
pv__transfer__splice_repeated(state, input_fd, output_fd,
state->transfer.transfer_buffer + state->transfer.read_position,
(size_t) max_to_read_into_buffer, max_to_read, max_to_write);
break;
#else /* !HAVE_SPLICE */
case PV_TRANSFERMETHOD_SPLICE:
case PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE:
nread =
pv__transfer__read_repeated(state, input_fd,
state->transfer.transfer_buffer + state->transfer.read_position,
(size_t) max_to_read_into_buffer);
break;
#endif /* HAVE_SPLICE */
if (state->control.no_splice || state->control.linemode || output_fd < 0 || state->transfer.to_write > 0) {
state->transfer.method = PV_TRANSFERMETHOD_READWRITE;
} else {
state->transfer.method = PV_TRANSFERMETHOD__UNSET;
}
/*
* Choose, and then attempt, transfer methods until one gets a
* definitive result.
*/
nread = 0;
transfer_attempted = false;
while (!transfer_attempted) {
switch (state->transfer.method) {
case PV_TRANSFERMETHOD__UNSET:
/*
* Try copy_file_range first, if the input and
* output are both regular files, and the method
* hasn't already failed on this input.
*/
#ifdef HAVE_COPY_FILE_RANGE
case PV_TRANSFERMETHOD_COPY_FILE_RANGE:
nread =
pv__transfer__copy_file_range_repeated(state, input_fd, output_fd,
state->transfer.transfer_buffer +
state->transfer.read_position,
(size_t) max_to_read_into_buffer, max_to_read, max_to_write);
break;
if (state->status.current_input_is_file && state->status.output_is_file
&& input_fd != state->transfer.copy_file_range_failed_fd) {
state->transfer.method = PV_TRANSFERMETHOD_COPY_FILE_RANGE;
break;
}
#endif
FALL_THROUGH;
/*@fallthrough@ */
case PV_TRANSFERMETHOD_COPY_FILE_RANGE:
/*
* Try splice, if either the input or the output is
* a pipe, and splice hasn't already failed on this
* input.
*/
#ifdef HAVE_SPLICE
if ((state->status.current_input_is_pipe || state->status.output_is_pipe)
&& input_fd != state->transfer.splice_failed_fd) {
state->transfer.method = PV_TRANSFERMETHOD_SPLICE;
break;
}
#endif
FALL_THROUGH;
/*@fallthrough@ */
case PV_TRANSFERMETHOD_SPLICE:
/*
* Try splice via intermediate pipe, if an
* intermediate pipe is available, and neither the
* input nor the output is a pipe, and splice hasn't
* already failed on this input.
*/
#ifdef HAVE_SPLICE
if (-1 != state->transfer.intermediate_pipe[0] && -1 != state->transfer.intermediate_pipe[1]
&& input_fd != state->transfer.splice_failed_fd) {
state->transfer.method = PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE;
break;
}
#endif
FALL_THROUGH;
/*@fallthrough@ */
default:
/* Fall back to read/write. */
state->transfer.method = PV_TRANSFERMETHOD_READWRITE;
break;
}
debug("%s=%d", "transfer method", state->transfer.method);
/*
* Transfer data using the appropriate method.
*
* If the chosen method cannot be used, it will set
* transfer_attempted to false, so this loop will go around
* again.
*/
transfer_attempted = true;
switch (state->transfer.method) {
case PV_TRANSFERMETHOD_COPY_FILE_RANGE:
#ifdef HAVE_COPY_FILE_RANGE
nread =
pv__transfer__copy_file_range_repeated(state, input_fd, output_fd,
max_to_read, max_to_write, &transfer_attempted);
#else /* !HAVE_COPY_FILE_RANGE */
case PV_TRANSFERMETHOD_COPY_FILE_RANGE:
nread =
pv__transfer__read_repeated(state, input_fd,
state->transfer.transfer_buffer + state->transfer.read_position,
(size_t) max_to_read_into_buffer);
break;
transfer_attempted = false;
#endif /* HAVE_COPY_FILE_RANGE */
break;
case PV_TRANSFERMETHOD_SPLICE:
case PV_TRANSFERMETHOD_SPLICE_INTERMEDIATE:
#ifdef HAVE_SPLICE
nread = pv__transfer__splice_repeated(state, input_fd, output_fd,
max_to_read, max_to_write, &transfer_attempted);
#else /* !HAVE_SPLICE */
transfer_attempted = false;
#endif /* HAVE_SPLICE */
break;
case PV_TRANSFERMETHOD__UNSET:
case PV_TRANSFERMETHOD_READWRITE:
nread =
pv__transfer__read_repeated(input_fd,
state->transfer.transfer_buffer + state->transfer.read_position,
(size_t) max_to_read_into_buffer);
break;
}
}
/*