Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions dispatcher/src/queue.rs
Original file line number Diff line number Diff line change
Expand Up @@ -500,6 +500,11 @@ impl WorkQueue {
self.remote_nodes.lock().unwrap().insert(node_id, channel);
}

/// Removes the offload channel for a remote node, e.g. when it has disconnected.
pub fn remove_remote_channel(&self, node_id: u64) {
self.remote_nodes.lock().unwrap().remove(&node_id);
}

/// Put work back into queue after trying to offload without success.
pub async fn reenqueue(&self, work: WorkToDo, debt: Debt) {
self.push(work, debt, false);
Expand Down
17 changes: 15 additions & 2 deletions dispatcher/src/queue/policy/data_locality.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,8 +54,21 @@ pub fn prepare_io_element(
if max_size * 2 > total_input_size {
let maybe_sender = remote_nodes.lock().unwrap().get(node_id).cloned();
if let Some(node_sender) = maybe_sender {
node_sender.send((work, debt)).unwrap();
return None;
// If the remote node has disconnected its receiver is gone, in which case
// we recover the work and fall back to executing it locally.
match node_sender.send((work, debt)) {
Ok(()) => return None,
Err(mpsc::error::SendError((work, debt))) => {
return Some((
work,
debt,
IOElementData {
remote_data,
total_input_size,
},
));
}
}
}
}
}
Expand Down
Loading
Loading