From 5ad047ab6b916c46a3566781977e746f77453ca7 Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Mon, 20 Jul 2026 11:04:15 +0300 Subject: [PATCH 1/2] Fix lost pre-handshake responses causing testAsynchronousBatch to hang Messages are dispatched to a thread pool in OpenICFServerAdapter, so an operation response arriving on a freshly opened socket could be processed before that socket's handshake and was silently discarded. For a batch operation this lost the batch-token response, and CompletionListener then timed out without completing the promise, leaving executeBatch() blocked until the connection group checker failed it 9 minutes later with "Operation finished on remote server with unknown result". Wait briefly for the in-flight handshake before dispatching non-handshake messages instead of dropping them, and make CompletionListener fail the promise when no batch token arrives so a lost token surfaces as an error instead of a hang. Fixes #107 --- .../framework/async/impl/BatchApiOpImpl.java | 6 +++- .../remote/OpenICFServerAdapter.java | 36 +++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java index 9da19228..056587bb 100644 --- a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java +++ b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java @@ -397,10 +397,14 @@ && new Date().getTime() < completionTimeout.get()) { if (returnToken != null) { getResultHandler().handleResult(returnToken); logger.ok("Token returned."); + observer.onCompleted(); } else { + // Without failing the promise here the caller blocked on it would + // wait forever for a token message that will never arrive. logger.ok("Batch CompletionListener timed out. Unable to return batch token."); + getExceptionHandler().handleException(new ConnectorException( + "Batch operation failed to return a batch token")); } - observer.onCompleted(); logger.ok("CompletionListener finished."); } } diff --git a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java index b701fbe3..0e70b080 100644 --- a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java +++ b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java @@ -12,6 +12,7 @@ * information: "Portions copyright [year] [name of copyright owner]". * * Copyright 2015-2016 ForgeRock AS. + * Portions Copyrighted 2026 3A Systems, LLC */ package org.forgerock.openicf.framework.remote; @@ -20,6 +21,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Locale; +import java.util.concurrent.TimeUnit; import org.forgerock.openicf.common.protobuf.CommonObjectMessages; import org.forgerock.openicf.common.protobuf.OperationMessages; @@ -141,6 +143,9 @@ public void processMessage(final WebSocketConnectionHolder socket, byte[] bytes) if (logger.isOk()) { logger.ok("{0} onMessage({1})", loggerName(), message.toString()); } + if (!isHandshakeMessage(message)) { + awaitHandshake(socket, message.getMessageId()); + } if (message.hasRequest()) { if (message.getRequest().hasHandshakeMessage()) { if (isClient()) { @@ -200,6 +205,37 @@ public void processMessage(final WebSocketConnectionHolder socket, byte[] bytes) } } + private static boolean isHandshakeMessage(final RemoteMessage message) { + return (message.hasRequest() && message.getRequest().hasHandshakeMessage()) + || (message.hasResponse() && message.getResponse().hasHandshakeMessage()); + } + + /** + * Messages are processed on a thread pool, so a message received on a + * freshly opened connection can overtake the handshake that arrived just + * before it and would be dropped by the pre-handshake branches below. The + * handshake is being processed concurrently, so briefly wait for it + * instead of discarding the message. + */ + private void awaitHandshake(final WebSocketConnectionHolder socket, long messageId) { + if (socket.isHandHooked()) { + return; + } + final long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10); + while (!socket.isHandHooked() && System.currentTimeMillis() < deadline) { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + } + if (!socket.isHandHooked()) { + logger.info("{0} handshake still pending after wait, message('{1}') via socket:{2} ", + loggerName(), messageId, socket.hashCode()); + } + } + protected void handleRemoteMessage(final WebSocketConnectionHolder socket, final RemoteMessage message) { logger.warn("{0} received unknown message('{1}') via socket:{2} ", loggerName(), message.getMessageId(), From 2efeda513440193d666e1923a8763998d70abf8f Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Mon, 20 Jul 2026 13:02:13 +0300 Subject: [PATCH 2/2] Address review feedback on pre-handshake message fix - Cancel the remote batch operation when CompletionListener times out without a token, mirroring the sendResponse failure path; guard with promise.isDone() and tolerate an already-dead transport. - Make the RemoteOperationContext fields volatile in the client, grizzly and jetty connection holders: they are written on the handshake- processing thread and read from other pool threads via isHandHooked(). - Shorten the awaitHandshake wait from 10s to 500ms, skip it for error responses (handled before the isHandHooked guard) and log the give-up at warn instead of info. - Add BatchApiOpImplTest covering the lost-token case: executeBatch now fails with ConnectorException in ~5s, observer.onError fires and onCompleted does not. --- .../framework/async/impl/BatchApiOpImpl.java | 12 +- .../ClientRemoteConnectorInfoManager.java | 4 +- .../remote/OpenICFServerAdapter.java | 21 ++- .../async/impl/BatchApiOpImplTest.java | 172 ++++++++++++++++++ .../grizzly/OpenICFWebSocketApplication.java | 6 +- .../server/jetty/SinglePrincipal.java | 6 +- 6 files changed, 209 insertions(+), 12 deletions(-) create mode 100644 OpenICF-java-framework/connector-framework-server/src/test/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImplTest.java diff --git a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java index 056587bb..7003b810 100644 --- a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java +++ b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImpl.java @@ -402,8 +402,16 @@ && new Date().getTime() < completionTimeout.get()) { // Without failing the promise here the caller blocked on it would // wait forever for a token message that will never arrive. logger.ok("Batch CompletionListener timed out. Unable to return batch token."); - getExceptionHandler().handleException(new ConnectorException( - "Batch operation failed to return a batch token")); + if (!getPromise().isDone()) { + getExceptionHandler().handleException(new ConnectorException( + "Batch operation failed to return a batch token")); + try { + tryCancelRemote(getConnectionContext(), getRequestId()); + } catch (Exception e) { + // The transport may already be down when the token was lost. + logger.ok(e, "Failed to cancel remote batch operation."); + } + } } logger.ok("CompletionListener finished."); } diff --git a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/client/ClientRemoteConnectorInfoManager.java b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/client/ClientRemoteConnectorInfoManager.java index 598afff8..8121ecdc 100644 --- a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/client/ClientRemoteConnectorInfoManager.java +++ b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/client/ClientRemoteConnectorInfoManager.java @@ -533,7 +533,9 @@ private class ICFWebSocket extends SimpleWebSocket { protected final Queue listeners = new ConcurrentLinkedQueue(); - private RemoteOperationContext context = null; + // Written on the handshake-processing pool thread, read by other message + // threads via getRemoteConnectionContext()/isHandHooked(). + private volatile RemoteOperationContext context = null; private final WebSocketConnectionHolder adapter = new WebSocketConnectionHolder() { diff --git a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java index 0e70b080..9c3d9521 100644 --- a/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java +++ b/OpenICF-java-framework/connector-framework-server/src/main/java/org/forgerock/openicf/framework/remote/OpenICFServerAdapter.java @@ -21,7 +21,6 @@ import java.util.ArrayList; import java.util.List; import java.util.Locale; -import java.util.concurrent.TimeUnit; import org.forgerock.openicf.common.protobuf.CommonObjectMessages; import org.forgerock.openicf.common.protobuf.OperationMessages; @@ -143,7 +142,7 @@ public void processMessage(final WebSocketConnectionHolder socket, byte[] bytes) if (logger.isOk()) { logger.ok("{0} onMessage({1})", loggerName(), message.toString()); } - if (!isHandshakeMessage(message)) { + if (!isHandshakeMessage(message) && !isErrorResponse(message)) { awaitHandshake(socket, message.getMessageId()); } if (message.hasRequest()) { @@ -210,18 +209,28 @@ private static boolean isHandshakeMessage(final RemoteMessage message) { || (message.hasResponse() && message.getResponse().hasHandshakeMessage()); } + /** + * Error responses are handled before the isHandHooked() guards below, so + * they never need to wait for the handshake. + */ + private static boolean isErrorResponse(final RemoteMessage message) { + return message.hasResponse() && message.getResponse().hasError(); + } + /** * Messages are processed on a thread pool, so a message received on a * freshly opened connection can overtake the handshake that arrived just * before it and would be dropped by the pre-handshake branches below. The - * handshake is being processed concurrently, so briefly wait for it - * instead of discarding the message. + * handshake is being processed concurrently and normally lands within + * milliseconds, so briefly wait for it instead of discarding the message. */ + private static final long HANDSHAKE_WAIT_MS = 500; + private void awaitHandshake(final WebSocketConnectionHolder socket, long messageId) { if (socket.isHandHooked()) { return; } - final long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10); + final long deadline = System.currentTimeMillis() + HANDSHAKE_WAIT_MS; while (!socket.isHandHooked() && System.currentTimeMillis() < deadline) { try { Thread.sleep(10); @@ -231,7 +240,7 @@ private void awaitHandshake(final WebSocketConnectionHolder socket, long message } } if (!socket.isHandHooked()) { - logger.info("{0} handshake still pending after wait, message('{1}') via socket:{2} ", + logger.warn("{0} handshake still pending after wait, message('{1}') via socket:{2} ", loggerName(), messageId, socket.hashCode()); } } diff --git a/OpenICF-java-framework/connector-framework-server/src/test/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImplTest.java b/OpenICF-java-framework/connector-framework-server/src/test/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImplTest.java new file mode 100644 index 00000000..51c606d1 --- /dev/null +++ b/OpenICF-java-framework/connector-framework-server/src/test/java/org/forgerock/openicf/framework/async/impl/BatchApiOpImplTest.java @@ -0,0 +1,172 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2026 3A Systems, LLC. + */ +package org.forgerock.openicf.framework.async.impl; + +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; + +import org.forgerock.openicf.common.protobuf.OperationMessages.BatchOpResult; +import org.forgerock.openicf.common.protobuf.OperationMessages.OperationResponse; +import org.forgerock.openicf.common.protobuf.RPCMessages.HandshakeMessage; +import org.forgerock.openicf.common.rpc.RemoteRequest; +import org.forgerock.openicf.common.rpc.RemoteRequestFactory; +import org.forgerock.openicf.common.rpc.RequestDistributor; +import org.forgerock.openicf.framework.remote.rpc.RemoteOperationContext; +import org.forgerock.openicf.framework.remote.rpc.WebSocketConnectionGroup; +import org.forgerock.openicf.framework.remote.rpc.WebSocketConnectionHolder; +import org.forgerock.util.Function; +import org.identityconnectors.framework.api.ConnectorKey; +import org.identityconnectors.framework.api.Observer; +import org.identityconnectors.framework.api.operations.APIOperation; +import org.identityconnectors.framework.api.operations.batch.BatchTask; +import org.identityconnectors.framework.api.operations.batch.DeleteBatchTask; +import org.identityconnectors.framework.common.exceptions.ConnectorException; +import org.identityconnectors.framework.common.objects.BatchResult; +import org.identityconnectors.framework.common.objects.ObjectClass; +import org.identityconnectors.framework.common.objects.OperationOptionsBuilder; +import org.identityconnectors.framework.common.objects.Uid; +import org.testng.Assert; +import org.testng.annotations.Test; + +import com.google.protobuf.ByteString; + +/** + * Tests that a batch operation whose batch-token response is never delivered + * fails the caller with a {@link ConnectorException} when the + * CompletionListener times out, instead of leaving the caller blocked forever. + */ +public class BatchApiOpImplTest { + + private static final ConnectorKey CONNECTOR_KEY = + new ConnectorKey("testbundle", "1.0", "test.Connector"); + + private static class RecordingObserver implements Observer { + final AtomicBoolean completed = new AtomicBoolean(false); + final AtomicReference error = new AtomicReference(); + + public void onCompleted() { + completed.set(true); + } + + public void onError(Throwable e) { + error.set(e); + } + + public void onNext(BatchResult batchResult) { + } + } + + private static final WebSocketConnectionHolder FAKE_SOCKET = new WebSocketConnectionHolder() { + + protected void handshake(HandshakeMessage message) { + } + + protected void tryClose() { + } + + public boolean isOperational() { + return true; + } + + public RemoteOperationContext getRemoteConnectionContext() { + return null; + } + + public Future sendBytes(byte[] data) { + return CompletableFuture.completedFuture(null); + } + + public Future sendString(String data) { + return CompletableFuture.completedFuture(null); + } + + public void sendPing(byte[] applicationData) throws Exception { + } + + public void sendPong(byte[] applicationData) throws Exception { + } + }; + + /** + * Submits the request over the fake socket and immediately delivers a + * "complete" BatchOpResult carrying no batch token, simulating a lost + * token response. + */ + private static class NoTokenRequestDistributor + implements + RequestDistributor { + + public , V, E extends Exception> R trySubmitRequest( + RemoteRequestFactory requestFactory) { + R request = + requestFactory.createRemoteRequest(null, 1L, + new RemoteRequestFactory.CompletionCallback() { + public void complete( + RemoteRequest request) { + } + }); + try { + request.getSendFunction().apply(FAKE_SOCKET); + } catch (Exception e) { + throw ConnectorException.wrap(e); + } + request.handleIncomingMessage(FAKE_SOCKET, OperationResponse.newBuilder() + .setBatchOpResult(BatchOpResult.newBuilder().setComplete(true)).build()); + return request; + } + + public boolean isOperational() { + return true; + } + } + + @Test(timeOut = 30000) + public void testExecuteBatchFailsWhenNoTokenArrives() { + BatchApiOpImpl operation = + new BatchApiOpImpl(new NoTokenRequestDistributor(), CONNECTOR_KEY, + new Function() { + public ByteString apply(RemoteOperationContext context) { + return ByteString.EMPTY; + } + }, APIOperation.NO_TIMEOUT); + + List tasks = + Collections. singletonList(new DeleteBatchTask(ObjectClass.ACCOUNT, + new Uid("1"), new OperationOptionsBuilder().build())); + RecordingObserver observer = new RecordingObserver(); + + long start = System.currentTimeMillis(); + try { + operation.executeBatch(tasks, observer, new OperationOptionsBuilder().build()); + Assert.fail("executeBatch must fail when no batch token arrives"); + } catch (ConnectorException expected) { + Assert.assertTrue(expected.getMessage().contains("batch token"), + "Unexpected failure: " + expected.getMessage()); + } + long elapsed = System.currentTimeMillis() - start; + + // CompletionListener gives up ~5s after the request was created. + Assert.assertTrue(elapsed < 20000, "Failure took too long: " + elapsed + "ms"); + Assert.assertNotNull(observer.error.get(), "observer.onError must be called"); + Assert.assertTrue(observer.error.get() instanceof ConnectorException); + Assert.assertFalse(observer.completed.get(), + "observer.onCompleted must not be called without a token"); + } +} diff --git a/OpenICF-java-framework/connector-server-grizzly/src/main/java/org/forgerock/openicf/framework/server/grizzly/OpenICFWebSocketApplication.java b/OpenICF-java-framework/connector-server-grizzly/src/main/java/org/forgerock/openicf/framework/server/grizzly/OpenICFWebSocketApplication.java index 3bdbd7a9..dc61bb78 100644 --- a/OpenICF-java-framework/connector-server-grizzly/src/main/java/org/forgerock/openicf/framework/server/grizzly/OpenICFWebSocketApplication.java +++ b/OpenICF-java-framework/connector-server-grizzly/src/main/java/org/forgerock/openicf/framework/server/grizzly/OpenICFWebSocketApplication.java @@ -20,6 +20,8 @@ * with the fields enclosed by brackets [] replaced by * your own identifying information: * "Portions Copyrighted [year] [name of copyright owner]" + * + * Portions Copyrighted 2026 3A Systems, LLC */ package org.forgerock.openicf.framework.server.grizzly; @@ -169,7 +171,9 @@ public static class OpenICFWebSocket extends DefaultWebSocket { new ConcurrentLinkedQueue(); private final ConnectionPrincipal connectionPrincipal; - private RemoteOperationContext context = null; + // Written on the handshake-processing pool thread, read by other message + // threads via getRemoteConnectionContext()/isHandHooked(). + private volatile RemoteOperationContext context = null; private final WebSocketConnectionHolder adapter = new WebSocketConnectionHolder() { diff --git a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/SinglePrincipal.java b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/SinglePrincipal.java index 1eca6ad8..79c968be 100644 --- a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/SinglePrincipal.java +++ b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/SinglePrincipal.java @@ -12,7 +12,7 @@ * information: "Portions copyright [year] [name of copyright owner]". * * Copyright 2015-2016 ForgeRock AS. - * Portions copyright 2025 3A Systems LLC. + * Portions copyright 2025-2026 3A Systems LLC. */ @@ -145,7 +145,9 @@ public void onWebSocketFrame(Frame frame) { private static final Logger logger = Log.getLogger(SinglePrincipal.class); private boolean hasCloseBeenCalled = false; - private RemoteOperationContext context = null; + // Written on the handshake-processing pool thread, read by other message + // threads via getRemoteConnectionContext()/isHandHooked(). + private volatile RemoteOperationContext context = null; private final WebSocketConnectionHolder adapter = new WebSocketConnectionHolder() {