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
10 changes: 8 additions & 2 deletions NEWS.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,13 @@
# Version (development version)

* ...

## Bug Fixes

* When a cluster worker that `plan()` had started was relaunched after
it was canceled, interrupted or had died, `plan()` never shut the
relaunched worker down, because the cluster registry kept the old
node. The worker process and its connection were left running until
garbage collection.


# Version 1.76.0 [2026-09-24]

Expand Down
15 changes: 15 additions & 0 deletions R/backend_api-11.ClusterFutureBackend-class.R
Original file line number Diff line number Diff line change
Expand Up @@ -1231,6 +1231,7 @@ requestNode <- local({

workers[[node_idx]] <- node2
backend[["workers"]] <- workers
clusterRegistry$replaceNode(workers, node_idx, node2)
node <- node2

if (debug) {
Expand Down Expand Up @@ -1609,6 +1610,7 @@ handleInterruptedFuture <- local({

## Update backend
backend[["workers"]] <- workers
clusterRegistry$replaceNode(workers, node_idx, node2)

node <- NULL

Expand Down Expand Up @@ -1743,6 +1745,18 @@ clusterRegistry <- local({
cluster
} ## startCluster()

## When a worker is relaunched, the registry must track the new node,
## otherwise stopCluster() closes the old, already closed, node and leaves
## the relaunched worker and its connection behind
replaceNode <- function(workers, idx, node) {
if (is.null(cluster)) return(invisible(FALSE))
if (!identical(attr(workers, "name", exact = TRUE), attr(cluster, "name", exact = TRUE))) {
return(invisible(FALSE))
}
cluster[[idx]] <<- node
invisible(TRUE)
} ## replaceNode()

stopCluster <- function(debug = FALSE) {
if (debug) {
mdebug_push("Stopping existing cluster ...")
Expand Down Expand Up @@ -1800,6 +1814,7 @@ clusterRegistry <- local({
list(
getCluster = getCluster,
startCluster = startCluster,
replaceNode = replaceNode,
stopCluster = stopCluster
)
}) ## clusterRegistry()
34 changes: 34 additions & 0 deletions inst/testme/test-cancel-relaunch-cluster.R
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
#' @tags cancel
#' @tags detritus-connections
#' @tags multisession

library(future)

message("Relaunched cluster workers are shut down with the plan ...")

n0 <- nrow(showConnections())

plan(multisession, workers = I(1))

## Cancel a running future, which terminates its worker
f <- future({ Sys.sleep(30); 42 })
Sys.sleep(1.0)
f <- cancel(f)

## The next future relaunches the terminated worker
f2 <- future(42)
stopifnot(value(f2) == 42)

## Keep a reference, so that garbage collection cannot hide the leak
backend <- plan("backend")

plan(sequential)

## The relaunched worker must be closed by plan(), not by garbage collection
n <- nrow(showConnections())
message("Connections: before = ", n0, ", after = ", n)
stopifnot(n <= n0)
rm(list = "backend")
gc()

message("Relaunched cluster workers are shut down with the plan ... done")
4 changes: 4 additions & 0 deletions tests/test-cancel-relaunch-cluster.R
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
#! /usr/bin/env Rscript
## This runs testme test script inst/testme/test-cancel-relaunch-cluster.R
## Don't edit - it was autogenerated by inst/testme/deploy.R
future:::testme("cancel-relaunch-cluster")