From 097b5b4ca40ecfd3c3dba9c6462bcb863d0f93e3 Mon Sep 17 00:00:00 2001 From: Alexis Sellier Date: Wed, 12 Apr 2023 11:27:57 +0200 Subject: [PATCH] node: Improve worker code logging and errors --- radicle-node/src/worker.rs | 55 +++++++++++++++++++++----------------- 1 file changed, 30 insertions(+), 25 deletions(-) diff --git a/radicle-node/src/worker.rs b/radicle-node/src/worker.rs index c79a89d0..9b407580 100644 --- a/radicle-node/src/worker.rs +++ b/radicle-node/src/worker.rs @@ -55,6 +55,8 @@ pub enum FetchError { Identity(#[from] IdentityError), #[error("upload failed: {0}")] Upload(#[from] UploadError), + #[error("worker channel error: {0}")] + Channel(#[from] chan::SendError), #[error(transparent)] StagingInit(#[from] fetch::error::Init), #[error(transparent)] @@ -203,11 +205,7 @@ impl Worker { } => { log::debug!(target: "worker", "Worker processing outgoing fetch for {}", rid); - let result = self.fetch(*rid, *remote, stream, namespaces, channels); - if let Err(err) = &result { - log::error!(target: "worker", "Fetch error: {err}"); - } - result + self.fetch(*rid, *remote, stream, namespaces, channels) } Fetch::Responder { .. } => { log::debug!(target: "worker", "Worker processing incoming fetch.."); @@ -237,32 +235,31 @@ impl Worker { mut channels: Channels, ) -> Result, FetchError> { let staging = fetch::StagingPhaseInitial::new(&self.storage, rid, namespaces.clone())?; - - self._fetch( + match self._fetch( &staging.repo, remote, staging.refspecs(), stream, &mut channels, - )?; - if let Err(e) = self.handle.flush(remote, stream) { - log::error!(target: "worker", "Error flushing worker stream: {e}"); + ) { + Ok(()) => log::debug!(target: "worker", "Initial fetch for {rid} exited successfully"), + Err(e) => log::error!(target: "worker", "Initial fetch for {rid} failed: {e}"), } + Self::eof(remote, stream, &mut channels.sender, &mut self.handle)?; - let staging = match staging.into_final().map_err(FetchError::from) { - Ok(staging) => staging, - Err(e) => return Err(e), - }; - self._fetch( + let staging = staging.into_final()?; + match self._fetch( &staging.repo, remote, staging.refspecs(), stream, &mut channels, - )?; - if let Err(e) = self.handle.flush(remote, stream) { - log::error!(target: "worker", "Error flushing worker stream: {e}"); + ) { + Ok(()) => log::debug!(target: "worker", "Final fetch for {rid} exited successfully"), + Err(e) => log::error!(target: "worker", "Final fetch for {rid} failed: {e}"), } + Self::eof(remote, stream, &mut channels.sender, &mut self.handle)?; + staging.transfer().map_err(FetchError::from) } @@ -368,7 +365,6 @@ impl Worker { S: fetch::AsRefspecs, { let mut tunnel = Tunnel::with(channels, stream, remote, self.handle.clone())?; - let rid = repo.id; let tunnel_addr = tunnel.local_addr(); let mut cmd = process::Command::new("git"); cmd.current_dir(repo.path()) @@ -410,22 +406,31 @@ impl Worker { tunnel.run(self.timeout)?; let result = child.wait()?; - let result = if result.success() { - log::debug!(target: "worker", "Fetch for {} exited successfully", rid); + if result.success() { Ok(()) } else { - log::error!(target: "worker", "Fetch for {} failed", rid); Err(FetchError::CommandFailed { code: result.code().unwrap_or(1), }) - }; + } + } + fn eof( + remote: NodeId, + stream: StreamId, + sender: &mut ChannelWriter, + handle: &mut Handle, + ) -> Result<(), FetchError> { log::debug!(target: "worker", "Sending `EOF` to remote.."); - if let Err(e) = channels.sender.eof() { + if let Err(e) = sender.eof() { log::error!(target: "worker", "Fetch error: error sending `EOF` message: {e}"); + return Err(e.into()); } - result + if let Err(e) = handle.flush(remote, stream) { + log::error!(target: "worker", "Error flushing worker stream: {e}"); + } + Ok(()) } }