node: Improve worker code logging and errors
This commit is contained in:
parent
eeec2008e8
commit
097b5b4ca4
|
|
@ -55,6 +55,8 @@ pub enum FetchError {
|
||||||
Identity(#[from] IdentityError),
|
Identity(#[from] IdentityError),
|
||||||
#[error("upload failed: {0}")]
|
#[error("upload failed: {0}")]
|
||||||
Upload(#[from] UploadError),
|
Upload(#[from] UploadError),
|
||||||
|
#[error("worker channel error: {0}")]
|
||||||
|
Channel(#[from] chan::SendError<ChannelEvent>),
|
||||||
#[error(transparent)]
|
#[error(transparent)]
|
||||||
StagingInit(#[from] fetch::error::Init),
|
StagingInit(#[from] fetch::error::Init),
|
||||||
#[error(transparent)]
|
#[error(transparent)]
|
||||||
|
|
@ -203,11 +205,7 @@ impl Worker {
|
||||||
} => {
|
} => {
|
||||||
log::debug!(target: "worker", "Worker processing outgoing fetch for {}", rid);
|
log::debug!(target: "worker", "Worker processing outgoing fetch for {}", rid);
|
||||||
|
|
||||||
let result = self.fetch(*rid, *remote, stream, namespaces, channels);
|
self.fetch(*rid, *remote, stream, namespaces, channels)
|
||||||
if let Err(err) = &result {
|
|
||||||
log::error!(target: "worker", "Fetch error: {err}");
|
|
||||||
}
|
|
||||||
result
|
|
||||||
}
|
}
|
||||||
Fetch::Responder { .. } => {
|
Fetch::Responder { .. } => {
|
||||||
log::debug!(target: "worker", "Worker processing incoming fetch..");
|
log::debug!(target: "worker", "Worker processing incoming fetch..");
|
||||||
|
|
@ -237,32 +235,31 @@ impl Worker {
|
||||||
mut channels: Channels,
|
mut channels: Channels,
|
||||||
) -> Result<Vec<RefUpdate>, FetchError> {
|
) -> Result<Vec<RefUpdate>, FetchError> {
|
||||||
let staging = fetch::StagingPhaseInitial::new(&self.storage, rid, namespaces.clone())?;
|
let staging = fetch::StagingPhaseInitial::new(&self.storage, rid, namespaces.clone())?;
|
||||||
|
match self._fetch(
|
||||||
self._fetch(
|
|
||||||
&staging.repo,
|
&staging.repo,
|
||||||
remote,
|
remote,
|
||||||
staging.refspecs(),
|
staging.refspecs(),
|
||||||
stream,
|
stream,
|
||||||
&mut channels,
|
&mut channels,
|
||||||
)?;
|
) {
|
||||||
if let Err(e) = self.handle.flush(remote, stream) {
|
Ok(()) => log::debug!(target: "worker", "Initial fetch for {rid} exited successfully"),
|
||||||
log::error!(target: "worker", "Error flushing worker stream: {e}");
|
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) {
|
let staging = staging.into_final()?;
|
||||||
Ok(staging) => staging,
|
match self._fetch(
|
||||||
Err(e) => return Err(e),
|
|
||||||
};
|
|
||||||
self._fetch(
|
|
||||||
&staging.repo,
|
&staging.repo,
|
||||||
remote,
|
remote,
|
||||||
staging.refspecs(),
|
staging.refspecs(),
|
||||||
stream,
|
stream,
|
||||||
&mut channels,
|
&mut channels,
|
||||||
)?;
|
) {
|
||||||
if let Err(e) = self.handle.flush(remote, stream) {
|
Ok(()) => log::debug!(target: "worker", "Final fetch for {rid} exited successfully"),
|
||||||
log::error!(target: "worker", "Error flushing worker stream: {e}");
|
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)
|
staging.transfer().map_err(FetchError::from)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -368,7 +365,6 @@ impl Worker {
|
||||||
S: fetch::AsRefspecs,
|
S: fetch::AsRefspecs,
|
||||||
{
|
{
|
||||||
let mut tunnel = Tunnel::with(channels, stream, remote, self.handle.clone())?;
|
let mut tunnel = Tunnel::with(channels, stream, remote, self.handle.clone())?;
|
||||||
let rid = repo.id;
|
|
||||||
let tunnel_addr = tunnel.local_addr();
|
let tunnel_addr = tunnel.local_addr();
|
||||||
let mut cmd = process::Command::new("git");
|
let mut cmd = process::Command::new("git");
|
||||||
cmd.current_dir(repo.path())
|
cmd.current_dir(repo.path())
|
||||||
|
|
@ -410,22 +406,31 @@ impl Worker {
|
||||||
tunnel.run(self.timeout)?;
|
tunnel.run(self.timeout)?;
|
||||||
|
|
||||||
let result = child.wait()?;
|
let result = child.wait()?;
|
||||||
let result = if result.success() {
|
if result.success() {
|
||||||
log::debug!(target: "worker", "Fetch for {} exited successfully", rid);
|
|
||||||
Ok(())
|
Ok(())
|
||||||
} else {
|
} else {
|
||||||
log::error!(target: "worker", "Fetch for {} failed", rid);
|
|
||||||
Err(FetchError::CommandFailed {
|
Err(FetchError::CommandFailed {
|
||||||
code: result.code().unwrap_or(1),
|
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..");
|
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}");
|
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(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue