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
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
package org.forgerock.openicf.common.rpc;

import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.ReentrantLock;

import org.forgerock.util.Function;
Expand All @@ -49,15 +50,55 @@ public abstract class RemoteRequest<V, E extends Exception, G extends RemoteConn
private final long requestId;
private final RemoteRequestFactory.CompletionCallback<V, E, G, H, P> completionCallback;

private Long requestTime = null;
private PromiseImpl<V, E> promise = null;
// The request is registered in RemoteConnectionGroup#remoteRequests
// before it is sent, so it can be cancelled from another thread while
// the send is still in progress. The promise therefore exists from
// construction on, and requestTime records whether the message has left
// (read by tryCancel on the cancelling thread).
private volatile Long requestTime = null;
private final PromiseImpl<V, E> promise;
private final ReentrantLock lock = new ReentrantLock();

// tryCancel(true) sets this before it reads requestTime; the send
// function sets requestTime before it reads this. Whichever runs second
// sees the other's write, so a cancel racing with the send never leaves
// the remote side uninformed - and remoteCancelSent keeps it to one
// cancel message when both do.
private volatile boolean remoteCancelRequested = false;
private final AtomicBoolean remoteCancelSent = new AtomicBoolean(false);

public RemoteRequest(P context, long requestId,
RemoteRequestFactory.CompletionCallback<V, E, G, H, P> completionCallback) {
this.context = context;
this.requestId = requestId;
this.completionCallback = completionCallback;
this.promise = new PromiseImpl<V, E>() {

@Override
protected E tryCancel(boolean mayInterruptIfRunning) {
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
if (mayInterruptIfRunning) {
remoteCancelRequested = true;
// Nothing to cancel remotely while the message has not
// been delivered: the send function does not send a
// cancelled request.
if (isSent()) {
try {
notifyRemoteCancelOnce();
} catch (final Throwable t) {
return createCancellationException(t);
}
}
}
return createCancellationException(null);
}

};
this.promise.thenOnResultOrException(new Runnable() {
@Override
public void run() {
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
RemoteRequest.this.completionCallback.complete(RemoteRequest.this);
}
});
}

/**
Expand Down Expand Up @@ -89,6 +130,13 @@ public Long getRequestTime() {
return requestTime;
}

/**
* Returns the promise of this request. It exists from construction on,
* so a request that is registered but not yet sent can be cancelled; the
* result arrives only once the message has been sent and answered.
*
* @return the promise, never {@code null}.
*/
public Promise<V, E> getPromise() {
return promise;
}
Expand All @@ -109,79 +157,77 @@ protected boolean cancel() {
return promise.cancel(false);
}

private boolean isSent() {
return null != requestTime;
}

private void notifyRemoteCancelOnce() {
if (remoteCancelSent.compareAndSet(false, true)) {
tryCancelRemote(context, requestId);
}
}

public Function<H, Promise<V, E>, Exception> getSendFunction() {
final Promise<V, E> resultPromise = promise;
if (null == resultPromise) {
final MessageElement message = createMessageElement(context, requestId);
if (message == null || !(message.isString() || message.isByte())) {
throw new IllegalStateException("RemoteRequest has empty message");
}
if (isSent()) {
return new Function<H, Promise<V, E>, Exception>() {

public Promise<V, E> apply(H remoteConnectionHolder) throws Exception {
if (null == promise) {
// Single thread should process it so it should not
// return false
if (lock.tryLock(1, TimeUnit.MINUTES)) {
try {
if (null == promise) {

promise = new PromiseImpl<V, E>() {

protected E tryCancel(boolean mayInterruptIfRunning) {
if (mayInterruptIfRunning) {
try {
tryCancelRemote(context, requestId);
} catch (final Throwable t) {
return createCancellationException(t);
}
}
return createCancellationException(null);
}

};

promise.thenOnResultOrException(new Runnable() {
public void run() {
completionCallback.complete(RemoteRequest.this);
}
});

try {
if (message.isByte()) {
remoteConnectionHolder.sendBytes(message.byteMessage)
.get();
} else if (message.isString()) {
remoteConnectionHolder
.sendString(message.stringMessage).get();
}
} catch (final Exception e) {
promise = null;
throw e;
} catch (final Throwable t) {
promise = null;
throw new Exception(t);
}
// Message has been delivered - Report
// success
requestTime = System.currentTimeMillis();
}
} finally {
lock.unlock();
}
}
}
@Override
public Promise<V, E> apply(H value) throws Exception {
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
return promise;
}
};
} else {
return new Function<H, Promise<V, E>, Exception>() {
}
final MessageElement message = createMessageElement(context, requestId);
if (message == null || !(message.isString() || message.isByte())) {
throw new IllegalStateException("RemoteRequest has empty message");
}
return new Function<H, Promise<V, E>, Exception>() {

public Promise<V, E> apply(H value) throws Exception {
return resultPromise;
@Override
public Promise<V, E> apply(H remoteConnectionHolder) throws Exception {
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
if (isSent()) {
return promise;
}
};
}
// Single thread should process it so it should not
// return false
if (!lock.tryLock(1, TimeUnit.MINUTES)) {
throw new IllegalStateException("RemoteRequest " + requestId
+ " is still being sent by another thread");
}
try {
// A request cancelled before it was sent stays unsent:
// the caller gets the cancelled promise back instead of
// waiting for an answer that can never arrive.
if (!isSent() && !promise.isDone()) {
// A failed send propagates to the group, which
// retries on its next connection with this same
// promise.
if (message.isByte()) {
remoteConnectionHolder.sendBytes(message.byteMessage).get();
} else if (message.isString()) {
remoteConnectionHolder.sendString(message.stringMessage).get();
}
// Message has been delivered - Report success
requestTime = System.currentTimeMillis();
if (remoteCancelRequested) {
// Cancelled while the message was on its way:
// tryCancel saw it unsent. The promise is
// already cancelled; a cancel that cannot be
// sent must not fail the delivered request, or
// the group would report it as never sent.
try {
notifyRemoteCancelOnce();
} catch (final Throwable ignored) {
// The transport is gone - nothing left to tell.
}
}
}
} finally {
lock.unlock();
}
return promise;
}
};
}

// --- inner Classes
Expand Down
Loading
Loading