From c0fd896df48a4cdaa7e4fb72b2ec1d2e6fe72fd6 Mon Sep 17 00:00:00 2001 From: Philipp Oppermann Date: Tue, 21 Feb 2023 16:09:28 +0100 Subject: [PATCH] Also close inputs when node finishes without sending `Stopped` message --- binaries/daemon/src/lib.rs | 66 +++++++++++++++++++++----------------- 1 file changed, 36 insertions(+), 30 deletions(-) diff --git a/binaries/daemon/src/lib.rs b/binaries/daemon/src/lib.rs index a60db6c0..27ac6b99 100644 --- a/binaries/daemon/src/lib.rs +++ b/binaries/daemon/src/lib.rs @@ -400,39 +400,44 @@ impl Daemon { let _ = reply_sender.send(DaemonReply::Result(Ok(()))); - // notify downstream nodes - let dataflow = self - .running - .get_mut(&dataflow_id) - .wrap_err_with(|| format!("failed to get downstream nodes: no running dataflow with ID `{dataflow_id}`"))?; - send_input_closed_events(dataflow, |(source_id, _)| source_id == &node_id).await; - - // TODO: notify remote nodes + self.handle_node_stop(dataflow_id, &node_id).await?; + } + } + Ok(()) + } - dataflow.running_nodes.remove(&node_id); - if dataflow.running_nodes.is_empty() { - tracing::info!( - "Dataflow `{dataflow_id}` finished on machine `{}`", - self.machine_id - ); - if let Some(addr) = self.coordinator_addr { - if coordinator::send_event( - addr, - self.machine_id.clone(), - DaemonEvent::AllNodesFinished { - dataflow_id, - result: Ok(()), - }, - ) - .await - .is_err() - { - tracing::warn!("failed to report dataflow finish to coordinator"); - } - } - self.running.remove(&dataflow_id); + #[tracing::instrument(skip(self))] + async fn handle_node_stop( + &mut self, + dataflow_id: Uuid, + node_id: &NodeId, + ) -> Result<(), eyre::ErrReport> { + let dataflow = self.running.get_mut(&dataflow_id).wrap_err_with(|| { + format!("failed to get downstream nodes: no running dataflow with ID `{dataflow_id}`") + })?; + send_input_closed_events(dataflow, |(source_id, _)| source_id == node_id).await; + dataflow.running_nodes.remove(node_id); + if dataflow.running_nodes.is_empty() { + tracing::info!( + "Dataflow `{dataflow_id}` finished on machine `{}`", + self.machine_id + ); + if let Some(addr) = self.coordinator_addr { + if coordinator::send_event( + addr, + self.machine_id.clone(), + DaemonEvent::AllNodesFinished { + dataflow_id, + result: Ok(()), + }, + ) + .await + .is_err() + { + tracing::warn!("failed to report dataflow finish to coordinator"); } } + self.running.remove(&dataflow_id); } Ok(()) } @@ -494,6 +499,7 @@ impl Daemon { tracing::warn!( "node `{dataflow_id}/{node_id}` finished without sending `Stopped` message" ); + self.handle_node_stop(dataflow_id, &node_id).await?; } match result { Ok(()) => {