Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import com.google.common.collect.Sets;
import com.google.common.net.InetAddresses;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
import diskCacheV111.namespace.EventNotifier;
import diskCacheV111.util.CacheException;
Expand Down Expand Up @@ -638,13 +639,8 @@ public void destroy() throws IOException {
public void messageArrived(PoolPassiveIoFileMessage<?> message) {

String poolName = message.getPoolName();
long verifier = message.getVerifier();
InetSocketAddress[] poolAddresses = message.socketAddresses();

_log.debug("NFS mover ready: {}", poolName);

PoolDS device = _poolDeviceMap.getOrCreateDS(poolName, verifier, poolAddresses);


// REVISIT 11.0: remove drop legacy support. Old polls will send legacy stateid.
// stateid4 stateid = message.challange();
Expand All @@ -659,6 +655,31 @@ public void messageArrived(PoolPassiveIoFileMessage<?> message) {
* Door reboot.
*/
if (transfer != null) {

long verifier = message.getVerifier();
InetSocketAddress[] poolAddresses = message.socketAddresses();

PoolDS device = _poolDeviceMap.getOrCreateDS(poolName, verifier, poolAddresses);

if (transfer.getMoverId() == null) {
_log.warn("NFS mover ready for transfer without mover: {}", stateid);
// we have not got a reply from pool manager yet, try to start mover again.
// as mover start requests is idempotent, we can safely re-issue it.
// Only redirect once the (re-)started mover is actually known; otherwise the
// client would be sent to a pool without a running mover. Also guard against
// re-starting a mover for a transfer that is already being torn down.
if (!transfer.hasMover()) {
transfer.startMoverAsync(_poolStub.getTimeoutInMillis()).addListener(() -> {
if (transfer.getMoverId() != null) {
transfer.redirect(device);
} else {
_log.error("Failed to re-start mover for transfer without mover: {}", stateid);
}
}, MoreExecutors.directExecutor());
}
return;
}

transfer.redirect(device);
}
}
Expand Down
1 change: 1 addition & 0 deletions modules/dcache/src/main/java/org/dcache/util/Transfer.java
Original file line number Diff line number Diff line change
Expand Up @@ -1073,6 +1073,7 @@ public ListenableFuture<Void> selectPoolAsync(long timeout) {

/**
* Creates a mover for the transfer.
* @param timeout timeout in milliseconds
*/
public ListenableFuture<Void> startMoverAsync(long timeout) {
FileAttributes fileAttributes = getFileAttributes();
Expand Down
Loading