refactor: structural splits, delete-engine consolidation, Config cleanup (no behavior change) #307
+260
-151
@@ -303,103 +303,121 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si
|
|||||||
return receiver_process_pending(config, file_descriptor, sink, NULL, NULL);
|
return receiver_process_pending(config, file_descriptor, sink, NULL, NULL);
|
||||||
}
|
}
|
||||||
|
|
||||||
/* Runs the whole receive loop. The delete manifest may legitimately arrive
|
/* Per-connection state threaded through the status handlers below. The parked
|
||||||
either FIRST (--delete-before / --delete-during: the sender transmits the
|
keep-set / per-directory session live here so one teardown helper can release
|
||||||
validated keep-set before any file data) or LAST (--delete-after /
|
them on every exit path. */
|
||||||
--delete-commit / --delete-delay: the manifest closes the data stream). In
|
typedef struct {
|
||||||
the early modes the receiver deletes as soon as the manifest has been read
|
Config* config;
|
||||||
and acknowledges with STATUS_OK so the sender only starts streaming once the
|
int fd;
|
||||||
deletion has committed (or failed); in the late modes the manifest is held
|
const ReceiverSink* sink;
|
||||||
and the deletion is committed only after the terminal STATUS_FINISHED proves
|
DeleteManifest** pending_manifest;
|
||||||
the whole transfer succeeded. A plain --delete defaults to the per-directory
|
DeletePlanSession** pending_plans;
|
||||||
delete-during plan mode (no manifest at all). See
|
/* Parked keep-set for the late/commit timing. Every exit path frees it
|
||||||
receiver_process_pending() for how the -m receiver defers that commit until
|
|
||||||
its disk writer has drained. */
|
|
||||||
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
|
|
||||||
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans) {
|
|
||||||
Status status;
|
|
||||||
if (!receive_status(file_descriptor, &status))
|
|
||||||
return -1;
|
|
||||||
/* Wall-clock (=CLOCK_MONOTONIC) anti-slowloris bookkeeping. session_start is
|
|
||||||
* fixed for the whole connection; last_progress is refreshed by every frame
|
|
||||||
* that is not a keepalive/abort. */
|
|
||||||
struct timespec session_start;
|
|
||||||
struct timespec last_progress;
|
|
||||||
clock_gettime(CLOCK_MONOTONIC, &session_start);
|
|
||||||
last_progress = session_start;
|
|
||||||
if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink))
|
|
||||||
return -1;
|
|
||||||
bool early_delete = config_delete_timing_early(config);
|
|
||||||
bool per_dir_delete = config_delete_timing_per_dir(config);
|
|
||||||
/* Parked keep-set for the late/commit timing. Every exit path below frees it
|
|
||||||
exactly once; the only exception is the successful FINISHED handoff, which
|
exactly once; the only exception is the successful FINISHED handoff, which
|
||||||
transfers ownership to *pending_manifest (used by the -m receiver). */
|
transfers ownership to *pending_manifest (used by the -m receiver). */
|
||||||
DeleteManifest* deferred_manifest = NULL;
|
DeleteManifest* deferred_manifest;
|
||||||
/* Per-directory delete session for --delete-during/--delete-delay. During the
|
/* Per-directory delete session for --delete-during/--delete-delay. During the
|
||||||
loop it applies plans inline (during) or snapshots their extras (delay); on
|
loop it applies plans inline (during) or snapshots their extras (delay); on
|
||||||
a successful FINISHED it is either committed here or handed to
|
a successful FINISHED it is either committed here or handed to
|
||||||
*pending_plans so the -m caller commits after its disk writer drained. */
|
*pending_plans so the -m caller commits after its disk writer drained. */
|
||||||
DeletePlanSession* plan_session = NULL;
|
DeletePlanSession* plan_session;
|
||||||
bool delete_limit_noted = false;
|
bool early_delete;
|
||||||
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
|
bool per_dir_delete;
|
||||||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH ||
|
bool delete_limit_noted;
|
||||||
status == STATUS_MKDIR || status == STATUS_MANIFEST || status == STATUS_HARDLINK ||
|
} ReceiverPendingState;
|
||||||
status == STATUS_SYMLINK || status == STATUS_SPECIAL || status == STATUS_DIR_TIMES ||
|
|
||||||
status == STATUS_DELETE_PLAN) {
|
/* Outcome of one frame handler. NEXT reads the following status frame; FAIL
|
||||||
if (status == STATUS_KEEPALIVE) {
|
tears the connection down without a peer STATUS_ERROR; ERROR tears it down
|
||||||
if (!send_status(file_descriptor, STATUS_KEEPALIVE))
|
and (when the sink owns error reporting) emits STATUS_ERROR. */
|
||||||
goto fail;
|
typedef enum {
|
||||||
goto next_status;
|
RECEIVER_STEP_NEXT,
|
||||||
}
|
RECEIVER_STEP_FAIL,
|
||||||
if (status == STATUS_ABORT) {
|
RECEIVER_STEP_ERROR,
|
||||||
|
} ReceiverStep;
|
||||||
|
|
||||||
|
static ReceiverStep receiver_handle_keepalive(ReceiverPendingState* state) {
|
||||||
|
if (!send_status(state->fd, STATUS_KEEPALIVE))
|
||||||
|
return RECEIVER_STEP_FAIL;
|
||||||
|
return RECEIVER_STEP_NEXT;
|
||||||
|
}
|
||||||
|
|
||||||
|
static ReceiverStep receiver_handle_abort(ReceiverPendingState* state) {
|
||||||
|
(void)state;
|
||||||
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
|
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
}
|
}
|
||||||
if (status == STATUS_CHECK) {
|
|
||||||
|
static ReceiverStep receiver_handle_check(ReceiverPendingState* state) {
|
||||||
bool skipped = false;
|
bool skipped = false;
|
||||||
bool would_transfer = false;
|
bool would_transfer = false;
|
||||||
File* file = receive_incremental_check_ex(file_descriptor, config, &skipped, &would_transfer);
|
File* file = receive_incremental_check_ex(state->fd, state->config, &skipped, &would_transfer);
|
||||||
if (config->dry_run) {
|
if (state->config->dry_run) {
|
||||||
/* Server-contacting --dry-run: the reply has already been sent
|
/* Server-contacting --dry-run: the reply has already been sent
|
||||||
(STATUS_OK = up to date, STATUS_DRY_RUN_TRANSFER = would transfer) and
|
(STATUS_OK = up to date, STATUS_DRY_RUN_TRANSFER = would transfer) and
|
||||||
nothing may be stored. Both flags false means a genuine protocol
|
nothing may be stored. Both flags false means a genuine protocol
|
||||||
error (STATUS_ERROR already sent or sent by receive_error below). */
|
error (STATUS_ERROR already sent or sent by receive_error below). */
|
||||||
if (!skipped && !would_transfer)
|
if (!skipped && !would_transfer)
|
||||||
goto receive_error;
|
return RECEIVER_STEP_ERROR;
|
||||||
} else if (!skipped && (!file || !sink->store_file(file, sink->context))) {
|
} else if (!skipped && (!file || !state->sink->store_file(file, state->sink->context))) {
|
||||||
goto receive_error;
|
return RECEIVER_STEP_ERROR;
|
||||||
}
|
}
|
||||||
} else if (status == STATUS_CHUNK) {
|
return RECEIVER_STEP_NEXT;
|
||||||
Chunk* chunk = receive_chunk_data(file_descriptor, config);
|
}
|
||||||
if (!chunk || !receiver_process_chunk(chunk, sink))
|
|
||||||
goto receive_error;
|
static ReceiverStep receiver_handle_chunk(ReceiverPendingState* state) {
|
||||||
} else if (status == STATUS_CHECK_BATCH) {
|
Chunk* chunk = receive_chunk_data(state->fd, state->config);
|
||||||
if (!receiver_process_batch(config, file_descriptor))
|
if (!chunk || !receiver_process_chunk(chunk, state->sink))
|
||||||
goto fail;
|
return RECEIVER_STEP_ERROR;
|
||||||
goto next_status;
|
return RECEIVER_STEP_NEXT;
|
||||||
} else if (status == STATUS_MKDIR) {
|
}
|
||||||
File* dir = file_receive_directory(file_descriptor, config);
|
|
||||||
if (!dir || !sink->store_file(dir, sink->context))
|
static ReceiverStep receiver_handle_check_batch(ReceiverPendingState* state) {
|
||||||
goto receive_error;
|
if (!receiver_process_batch(state->config, state->fd))
|
||||||
} else if (status == STATUS_DIR_TIMES) {
|
return RECEIVER_STEP_FAIL;
|
||||||
if (!receiver_process_dir_times(file_descriptor, config, sink))
|
return RECEIVER_STEP_NEXT;
|
||||||
goto receive_error;
|
}
|
||||||
} else if (status == STATUS_HARDLINK) {
|
|
||||||
File* file = file_receive_hardlink(file_descriptor);
|
static ReceiverStep receiver_handle_mkdir(ReceiverPendingState* state) {
|
||||||
if (!file || !sink->store_file(file, sink->context))
|
File* dir = file_receive_directory(state->fd, state->config);
|
||||||
goto receive_error;
|
if (!dir || !state->sink->store_file(dir, state->sink->context))
|
||||||
} else if (status == STATUS_SYMLINK) {
|
return RECEIVER_STEP_ERROR;
|
||||||
File* sym = file_receive_symlink(file_descriptor, config);
|
return RECEIVER_STEP_NEXT;
|
||||||
if (!sym || !sink->store_file(sym, sink->context))
|
}
|
||||||
goto receive_error;
|
|
||||||
} else if (status == STATUS_SPECIAL) {
|
static ReceiverStep receiver_handle_dir_times(ReceiverPendingState* state) {
|
||||||
File* file = file_receive_special(file_descriptor);
|
if (!receiver_process_dir_times(state->fd, state->config, state->sink))
|
||||||
if (!file || !sink->store_file(file, sink->context))
|
return RECEIVER_STEP_ERROR;
|
||||||
goto receive_error;
|
return RECEIVER_STEP_NEXT;
|
||||||
} else if (status == STATUS_MANIFEST) {
|
}
|
||||||
DeleteManifest* manifest = receive_manifest_entries(file_descriptor);
|
|
||||||
|
static ReceiverStep receiver_handle_hardlink(ReceiverPendingState* state) {
|
||||||
|
File* file = file_receive_hardlink(state->fd);
|
||||||
|
if (!file || !state->sink->store_file(file, state->sink->context))
|
||||||
|
return RECEIVER_STEP_ERROR;
|
||||||
|
return RECEIVER_STEP_NEXT;
|
||||||
|
}
|
||||||
|
|
||||||
|
static ReceiverStep receiver_handle_symlink(ReceiverPendingState* state) {
|
||||||
|
File* sym = file_receive_symlink(state->fd, state->config);
|
||||||
|
if (!sym || !state->sink->store_file(sym, state->sink->context))
|
||||||
|
return RECEIVER_STEP_ERROR;
|
||||||
|
return RECEIVER_STEP_NEXT;
|
||||||
|
}
|
||||||
|
|
||||||
|
static ReceiverStep receiver_handle_special(ReceiverPendingState* state) {
|
||||||
|
File* file = file_receive_special(state->fd);
|
||||||
|
if (!file || !state->sink->store_file(file, state->sink->context))
|
||||||
|
return RECEIVER_STEP_ERROR;
|
||||||
|
return RECEIVER_STEP_NEXT;
|
||||||
|
}
|
||||||
|
|
||||||
|
static ReceiverStep receiver_handle_manifest(ReceiverPendingState* state) {
|
||||||
|
Config* config = state->config;
|
||||||
|
int fd = state->fd;
|
||||||
|
const ReceiverSink* sink = state->sink;
|
||||||
|
DeleteManifest* manifest = receive_manifest_entries(fd);
|
||||||
if (!manifest)
|
if (!manifest)
|
||||||
goto fail; /* receive_manifest_entries already sent STATUS_ERROR */
|
return RECEIVER_STEP_FAIL; /* receive_manifest_entries already sent STATUS_ERROR */
|
||||||
if (config->dry_run) {
|
if (config->dry_run) {
|
||||||
/* Server-contacting --dry-run mutates nothing, so a keep-set manifest
|
/* Server-contacting --dry-run mutates nothing, so a keep-set manifest
|
||||||
is consumed and discarded. The early-delete mode still needs its ACK
|
is consumed and discarded. The early-delete mode still needs its ACK
|
||||||
@@ -412,11 +430,11 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
log_message(LOG_LEVEL_WARNING, "dry-run: could not enumerate would-delete paths");
|
log_message(LOG_LEVEL_WARNING, "dry-run: could not enumerate would-delete paths");
|
||||||
}
|
}
|
||||||
delete_manifest_free(manifest);
|
delete_manifest_free(manifest);
|
||||||
if (early_delete && !send_status(file_descriptor, STATUS_OK))
|
if (state->early_delete && !send_status(fd, STATUS_OK))
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
goto next_status;
|
return RECEIVER_STEP_NEXT;
|
||||||
}
|
}
|
||||||
if (early_delete) {
|
if (state->early_delete) {
|
||||||
/* --delete-before: the whole-tree manifest is authoritative the moment
|
/* --delete-before: the whole-tree manifest is authoritative the moment
|
||||||
it arrives, before any file data. Delete now and acknowledge so the
|
it arrives, before any file data. Delete now and acknowledge so the
|
||||||
sender only starts streaming once the deletion committed (or failed).
|
sender only starts streaming once the deletion committed (or failed).
|
||||||
@@ -434,61 +452,164 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
receiver_tally_deleted(sink, deleted);
|
receiver_tally_deleted(sink, deleted);
|
||||||
delete_manifest_free(manifest);
|
delete_manifest_free(manifest);
|
||||||
if (deletion == DELETE_COMMIT_ERROR) {
|
if (deletion == DELETE_COMMIT_ERROR) {
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(fd, STATUS_ERROR);
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
}
|
}
|
||||||
if (deletion == DELETE_COMMIT_LIMIT_REACHED && sink->note_delete_limit)
|
if (deletion == DELETE_COMMIT_LIMIT_REACHED && sink->note_delete_limit)
|
||||||
sink->note_delete_limit(sink->context);
|
sink->note_delete_limit(sink->context);
|
||||||
if (!send_status(file_descriptor, STATUS_OK))
|
if (!send_status(fd, STATUS_OK))
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
} else if (config->use_delete || config->delete_missing_args) {
|
} else if (config->use_delete || config->delete_missing_args) {
|
||||||
/* Plain --delete / --delete-after and the --delete-missing-args
|
/* Plain --delete / --delete-after and the --delete-missing-args
|
||||||
exact-path deletions: hold the manifest and commit it only after
|
exact-path deletions: hold the manifest and commit it only after
|
||||||
STATUS_FINISHED. The per-directory modes never send this frame. */
|
STATUS_FINISHED. The per-directory modes never send this frame. */
|
||||||
if (deferred_manifest) {
|
if (state->deferred_manifest) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
|
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
|
||||||
delete_manifest_free(deferred_manifest);
|
delete_manifest_free(state->deferred_manifest);
|
||||||
deferred_manifest = NULL;
|
state->deferred_manifest = NULL;
|
||||||
delete_manifest_free(manifest);
|
delete_manifest_free(manifest);
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(fd, STATUS_ERROR);
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
}
|
}
|
||||||
deferred_manifest = manifest;
|
state->deferred_manifest = manifest;
|
||||||
} else {
|
} else {
|
||||||
delete_manifest_free(manifest);
|
delete_manifest_free(manifest);
|
||||||
}
|
}
|
||||||
goto next_status;
|
return RECEIVER_STEP_NEXT;
|
||||||
} else if (status == STATUS_DELETE_PLAN) {
|
}
|
||||||
if (!per_dir_delete) {
|
|
||||||
|
static ReceiverStep receiver_handle_delete_plan(ReceiverPendingState* state) {
|
||||||
|
Config* config = state->config;
|
||||||
|
int fd = state->fd;
|
||||||
|
const ReceiverSink* sink = state->sink;
|
||||||
|
if (!state->per_dir_delete) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Received a per-directory delete plan without a per-dir "
|
log_message(LOG_LEVEL_ERROR, "Received a per-directory delete plan without a per-dir "
|
||||||
"delete timing");
|
"delete timing");
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(fd, STATUS_ERROR);
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
}
|
}
|
||||||
if (!plan_session) {
|
if (!state->plan_session) {
|
||||||
plan_session = delete_plan_session_create(config);
|
state->plan_session = delete_plan_session_create(config);
|
||||||
if (plan_session && config->report_deletes && sink->deleted_paths)
|
if (state->plan_session && config->report_deletes && sink->deleted_paths)
|
||||||
delete_plan_session_set_delete_observer(plan_session, receiver_record_deleted_path,
|
delete_plan_session_set_delete_observer(state->plan_session, receiver_record_deleted_path,
|
||||||
(void*)sink->deleted_paths);
|
(void*)sink->deleted_paths);
|
||||||
}
|
}
|
||||||
if (!plan_session || delete_plan_session_receive(plan_session, config, file_descriptor) != 0)
|
if (!state->plan_session || delete_plan_session_receive(state->plan_session, config, fd) != 0)
|
||||||
goto fail;
|
return RECEIVER_STEP_FAIL;
|
||||||
if (delete_plan_session_limit_reached(plan_session) && !delete_limit_noted &&
|
if (delete_plan_session_limit_reached(state->plan_session) && !state->delete_limit_noted &&
|
||||||
sink->note_delete_limit) {
|
sink->note_delete_limit) {
|
||||||
sink->note_delete_limit(sink->context);
|
sink->note_delete_limit(sink->context);
|
||||||
delete_limit_noted = true;
|
state->delete_limit_noted = true;
|
||||||
}
|
}
|
||||||
goto next_status;
|
return RECEIVER_STEP_NEXT;
|
||||||
} else {
|
}
|
||||||
File* file = file_receive(config, file_descriptor);
|
|
||||||
|
static ReceiverStep receiver_handle_file(ReceiverPendingState* state) {
|
||||||
|
File* file = file_receive(state->config, state->fd);
|
||||||
if (!file) {
|
if (!file) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
|
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
|
||||||
goto receive_error;
|
return RECEIVER_STEP_ERROR;
|
||||||
}
|
}
|
||||||
if (!sink->store_file(file, sink->context))
|
if (!state->sink->store_file(file, state->sink->context))
|
||||||
goto receive_error;
|
return RECEIVER_STEP_ERROR;
|
||||||
|
return RECEIVER_STEP_NEXT;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* One dispatch per admitted frame type; STATUS_NEXT (and any other
|
||||||
|
data-bearing status) falls through to the regular file receiver. */
|
||||||
|
static ReceiverStep receiver_dispatch_status(ReceiverPendingState* state, Status status) {
|
||||||
|
switch (status) {
|
||||||
|
case STATUS_KEEPALIVE:
|
||||||
|
return receiver_handle_keepalive(state);
|
||||||
|
case STATUS_ABORT:
|
||||||
|
return receiver_handle_abort(state);
|
||||||
|
case STATUS_CHECK:
|
||||||
|
return receiver_handle_check(state);
|
||||||
|
case STATUS_CHUNK:
|
||||||
|
return receiver_handle_chunk(state);
|
||||||
|
case STATUS_CHECK_BATCH:
|
||||||
|
return receiver_handle_check_batch(state);
|
||||||
|
case STATUS_MKDIR:
|
||||||
|
return receiver_handle_mkdir(state);
|
||||||
|
case STATUS_DIR_TIMES:
|
||||||
|
return receiver_handle_dir_times(state);
|
||||||
|
case STATUS_HARDLINK:
|
||||||
|
return receiver_handle_hardlink(state);
|
||||||
|
case STATUS_SYMLINK:
|
||||||
|
return receiver_handle_symlink(state);
|
||||||
|
case STATUS_SPECIAL:
|
||||||
|
return receiver_handle_special(state);
|
||||||
|
case STATUS_MANIFEST:
|
||||||
|
return receiver_handle_manifest(state);
|
||||||
|
case STATUS_DELETE_PLAN:
|
||||||
|
return receiver_handle_delete_plan(state);
|
||||||
|
default:
|
||||||
|
return receiver_handle_file(state);
|
||||||
}
|
}
|
||||||
next_status:
|
}
|
||||||
|
|
||||||
|
/* Release the parked keep-set / per-directory session exactly once on every
|
||||||
|
failure exit. Never commit a deletion for a failed stream. */
|
||||||
|
static void receiver_drop_pending(ReceiverPendingState* state) {
|
||||||
|
if (state->deferred_manifest) {
|
||||||
|
delete_manifest_free(state->deferred_manifest);
|
||||||
|
state->deferred_manifest = NULL;
|
||||||
|
}
|
||||||
|
if (state->plan_session) {
|
||||||
|
delete_plan_session_destroy(state->plan_session);
|
||||||
|
state->plan_session = NULL;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Runs the whole receive loop. The delete manifest may legitimately arrive
|
||||||
|
either FIRST (--delete-before / --delete-during: the sender transmits the
|
||||||
|
validated keep-set before any file data) or LAST (--delete-after /
|
||||||
|
--delete-commit / --delete-delay: the manifest closes the data stream). In
|
||||||
|
the early modes the receiver deletes as soon as the manifest has been read
|
||||||
|
and acknowledges with STATUS_OK so the sender only starts streaming once the
|
||||||
|
deletion has committed (or failed); in the late modes the manifest is held
|
||||||
|
and the deletion is committed only after the terminal STATUS_FINISHED proves
|
||||||
|
the whole transfer succeeded. A plain --delete defaults to the per-directory
|
||||||
|
delete-during plan mode (no manifest at all). See the per-frame handlers
|
||||||
|
above for how the -m receiver defers that commit until its disk writer has
|
||||||
|
drained. */
|
||||||
|
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
|
||||||
|
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans) {
|
||||||
|
Status status;
|
||||||
|
if (!receive_status(file_descriptor, &status))
|
||||||
|
return -1;
|
||||||
|
/* Wall-clock (=CLOCK_MONOTONIC) anti-slowloris bookkeeping. session_start is
|
||||||
|
* fixed for the whole connection; last_progress is refreshed by every frame
|
||||||
|
* that is not a keepalive/abort. */
|
||||||
|
struct timespec session_start;
|
||||||
|
struct timespec last_progress;
|
||||||
|
clock_gettime(CLOCK_MONOTONIC, &session_start);
|
||||||
|
last_progress = session_start;
|
||||||
|
if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink))
|
||||||
|
return -1;
|
||||||
|
ReceiverPendingState state = {
|
||||||
|
.config = config,
|
||||||
|
.fd = file_descriptor,
|
||||||
|
.sink = sink,
|
||||||
|
.pending_manifest = pending_manifest,
|
||||||
|
.pending_plans = pending_plans,
|
||||||
|
.deferred_manifest = NULL,
|
||||||
|
.plan_session = NULL,
|
||||||
|
.early_delete = config_delete_timing_early(config),
|
||||||
|
.per_dir_delete = config_delete_timing_per_dir(config),
|
||||||
|
.delete_limit_noted = false,
|
||||||
|
};
|
||||||
|
bool notify_peer = false;
|
||||||
|
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
|
||||||
|
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH ||
|
||||||
|
status == STATUS_MKDIR || status == STATUS_MANIFEST || status == STATUS_HARDLINK ||
|
||||||
|
status == STATUS_SYMLINK || status == STATUS_SPECIAL || status == STATUS_DIR_TIMES ||
|
||||||
|
status == STATUS_DELETE_PLAN) {
|
||||||
|
ReceiverStep step = receiver_dispatch_status(&state, status);
|
||||||
|
if (step == RECEIVER_STEP_FAIL)
|
||||||
|
goto fail;
|
||||||
|
if (step == RECEIVER_STEP_ERROR)
|
||||||
|
goto receive_error;
|
||||||
if (!receive_status(file_descriptor, &status))
|
if (!receive_status(file_descriptor, &status))
|
||||||
goto receive_error;
|
goto receive_error;
|
||||||
if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink))
|
if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink))
|
||||||
@@ -507,19 +628,19 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
disk writer may still be draining; the caller commits after the writer has
|
disk writer may still be draining; the caller commits after the writer has
|
||||||
joined so no extra file is removed unless the transfer is known to have
|
joined so no extra file is removed unless the transfer is known to have
|
||||||
succeeded. */
|
succeeded. */
|
||||||
if (deferred_manifest) {
|
if (state.deferred_manifest) {
|
||||||
if (pending_manifest) {
|
if (state.pending_manifest) {
|
||||||
*pending_manifest = deferred_manifest;
|
*state.pending_manifest = state.deferred_manifest;
|
||||||
deferred_manifest = NULL;
|
state.deferred_manifest = NULL;
|
||||||
} else {
|
} else {
|
||||||
size_t deleted = 0;
|
size_t deleted = 0;
|
||||||
DeletePathObserver observer =
|
DeletePathObserver observer =
|
||||||
(config->report_deletes && sink->deleted_paths) ? receiver_record_deleted_path : NULL;
|
(config->report_deletes && sink->deleted_paths) ? receiver_record_deleted_path : NULL;
|
||||||
DeleteCommitResult deletion = manifest_delete_all_observed(
|
DeleteCommitResult deletion = manifest_delete_all_observed(
|
||||||
config, deferred_manifest, &deleted, observer, (void*)sink->deleted_paths);
|
config, state.deferred_manifest, &deleted, observer, (void*)sink->deleted_paths);
|
||||||
receiver_tally_deleted(sink, deleted);
|
receiver_tally_deleted(sink, deleted);
|
||||||
delete_manifest_free(deferred_manifest);
|
delete_manifest_free(state.deferred_manifest);
|
||||||
deferred_manifest = NULL;
|
state.deferred_manifest = NULL;
|
||||||
if (deletion == DELETE_COMMIT_ERROR) {
|
if (deletion == DELETE_COMMIT_ERROR) {
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
goto fail;
|
goto fail;
|
||||||
@@ -533,28 +654,28 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
nothing yet and applies its decompressed snapshot here. The -m receiver
|
nothing yet and applies its decompressed snapshot here. The -m receiver
|
||||||
hands the session to its caller instead, which commits after the disk
|
hands the session to its caller instead, which commits after the disk
|
||||||
writer drained. */
|
writer drained. */
|
||||||
if (plan_session) {
|
if (state.plan_session) {
|
||||||
if (config->report_deletes && sink->deleted_paths)
|
if (config->report_deletes && sink->deleted_paths)
|
||||||
delete_plan_session_set_delete_observer(plan_session, receiver_record_deleted_path,
|
delete_plan_session_set_delete_observer(state.plan_session, receiver_record_deleted_path,
|
||||||
(void*)sink->deleted_paths);
|
(void*)sink->deleted_paths);
|
||||||
if (pending_plans) {
|
if (state.pending_plans) {
|
||||||
*pending_plans = plan_session;
|
*state.pending_plans = state.plan_session;
|
||||||
plan_session = NULL;
|
state.plan_session = NULL;
|
||||||
} else if (config->dry_run) {
|
} else if (config->dry_run) {
|
||||||
/* Central dry-run no-op: never commit a deletion for a -n run. */
|
/* Central dry-run no-op: never commit a deletion for a -n run. */
|
||||||
delete_plan_session_destroy(plan_session);
|
delete_plan_session_destroy(state.plan_session);
|
||||||
plan_session = NULL;
|
state.plan_session = NULL;
|
||||||
} else {
|
} else {
|
||||||
DeleteCommitResult deletion = delete_plan_session_commit(plan_session, config);
|
DeleteCommitResult deletion = delete_plan_session_commit(state.plan_session, config);
|
||||||
bool limit = delete_plan_session_limit_reached(plan_session);
|
bool limit = delete_plan_session_limit_reached(state.plan_session);
|
||||||
receiver_tally_deleted(sink, delete_plan_session_deleted(plan_session));
|
receiver_tally_deleted(sink, delete_plan_session_deleted(state.plan_session));
|
||||||
delete_plan_session_destroy(plan_session);
|
delete_plan_session_destroy(state.plan_session);
|
||||||
plan_session = NULL;
|
state.plan_session = NULL;
|
||||||
if (deletion == DELETE_COMMIT_ERROR) {
|
if (deletion == DELETE_COMMIT_ERROR) {
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
goto fail;
|
goto fail;
|
||||||
}
|
}
|
||||||
if (limit && !delete_limit_noted && sink->note_delete_limit)
|
if (limit && !state.delete_limit_noted && sink->note_delete_limit)
|
||||||
sink->note_delete_limit(sink->context);
|
sink->note_delete_limit(sink->context);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -568,26 +689,14 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
}
|
}
|
||||||
return 0;
|
return 0;
|
||||||
|
|
||||||
|
receive_error:
|
||||||
|
notify_peer = true;
|
||||||
fail:
|
fail:
|
||||||
/* Failure exits that must not (or already did) report a STATUS_ERROR. The
|
/* Failure exits that must not (or already did) report a STATUS_ERROR. The
|
||||||
parked keep-set/session is dropped: never commit a deletion for a failed
|
parked keep-set/session is dropped: never commit a deletion for a failed
|
||||||
stream. */
|
stream. */
|
||||||
if (deferred_manifest) {
|
receiver_drop_pending(&state);
|
||||||
delete_manifest_free(deferred_manifest);
|
if (notify_peer && sink->send_error)
|
||||||
deferred_manifest = NULL;
|
|
||||||
}
|
|
||||||
if (plan_session)
|
|
||||||
delete_plan_session_destroy(plan_session);
|
|
||||||
return -1;
|
|
||||||
|
|
||||||
receive_error:
|
|
||||||
if (deferred_manifest) {
|
|
||||||
delete_manifest_free(deferred_manifest);
|
|
||||||
deferred_manifest = NULL;
|
|
||||||
}
|
|
||||||
if (plan_session)
|
|
||||||
delete_plan_session_destroy(plan_session);
|
|
||||||
if (sink->send_error)
|
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
return -1;
|
return -1;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user