From b8f3172252681242d495691b683da9700a0f136a Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Mon, 5 Oct 2026 15:37:43 +0300 Subject: [PATCH 1/3] [#151] Fix the batch test race and return pooled connectors in UseCase3 tests testBatchUseCase3Failure asserted zero results right after executeBatch(), but the UseCase3 test connector delivers results from its own thread, so the first onNext() could already have run. Drop the transient checks and wait up to 3 s for hasError instead of a fixed 500 ms sleep, as #93 did for testBatchUseCase3. Both UseCase3 tests now close their subscriptions: each one holds a pooled connector until it is closed, so repeated runs exhausted the pool and hung in borrowObject(). Closing the queryBatch subscription of a finished batch threw an NPE, because SubscriptionImpl wrapped the null subscription returned by the connector. Make close() and isUnsubscribed() of the executeBatch and queryBatch wrappers null-safe, as getReturnValue() of queryBatch already was. Fixes #151 --- .../local/operations/SubscriptionImpl.java | 12 +++-- .../api/LocalConnectorInfoManagerTests.java | 54 ++++++++++++------- 2 files changed, 41 insertions(+), 25 deletions(-) diff --git a/OpenICF-java-framework/connector-framework-internal/src/main/java/org/identityconnectors/framework/impl/api/local/operations/SubscriptionImpl.java b/OpenICF-java-framework/connector-framework-internal/src/main/java/org/identityconnectors/framework/impl/api/local/operations/SubscriptionImpl.java index 16da473a0..e47cae9d7 100644 --- a/OpenICF-java-framework/connector-framework-internal/src/main/java/org/identityconnectors/framework/impl/api/local/operations/SubscriptionImpl.java +++ b/OpenICF-java-framework/connector-framework-internal/src/main/java/org/identityconnectors/framework/impl/api/local/operations/SubscriptionImpl.java @@ -163,19 +163,20 @@ public Subscription executeBatch(final List tasks, final Observer batchTasks = batch.build(); - Subscription sub = facade.executeBatch(batchTasks, observer, options); - assertNotNull(sub.getReturnValue()); + final Subscription sub = facade.executeBatch(batchTasks, observer, options); + try { + assertNotNull(sub.getReturnValue()); - final long timeout = System.currentTimeMillis() + 3000; - while (!isComplete.get() && System.currentTimeMillis() < timeout) { - Thread.sleep(100); - } + final long timeout = System.currentTimeMillis() + 3000; + while (!isComplete.get() && System.currentTimeMillis() < timeout) { + Thread.sleep(100); + } - assertEquals(results.size(), batchTasks.size()); - assertTrue(isComplete.get()); - assertFalse(hasError.get()); + assertEquals(results.size(), batchTasks.size()); + assertTrue(isComplete.get()); + assertFalse(hasError.get()); - sub = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); - assertNull(sub.getReturnValue()); + final Subscription query = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); + try { + assertNull(query.getReturnValue()); + } finally { + query.close(); + } + } finally { + // the subscription holds a pooled connector until it is closed + sub.close(); + } } @Test @@ -479,17 +488,22 @@ public void onNext(BatchResult batchResult) { } }; - Subscription sub = facade.executeBatch(batch.build(), observer, options); - assertEquals(results.size(), 0); - assertFalse(isComplete.get()); - assertFalse(hasError.get()); - assertNotNull(sub.getReturnValue()); + final Subscription sub = facade.executeBatch(batch.build(), observer, options); + try { + assertNotNull(sub.getReturnValue()); - Thread.sleep(500); + final long timeout = System.currentTimeMillis() + 3000; + while (!hasError.get() && System.currentTimeMillis() < timeout) { + Thread.sleep(100); + } - assertEquals(results.size(), 2); - assertFalse(isComplete.get()); - assertTrue(hasError.get()); + assertEquals(results.size(), 2); + assertFalse(isComplete.get()); + assertTrue(hasError.get()); + } finally { + // the subscription holds a pooled connector until it is closed + sub.close(); + } } @Test From 8ac6d5712963e2b401bdbb887625724494ee2361 Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Tue, 6 Oct 2026 15:32:06 +0300 Subject: [PATCH 2/3] [#151] Close every batch test subscription and mark UseCase3 complete before onCompleted() Address the review of #153: - TstConnector: BatchUseCase3Processor marks the token complete before it calls observer.onCompleted(), so a test that has seen onCompleted() always gets null from queryBatch(). - LocalConnectorInfoManagerTests: the UseCase0/1, UseCase2 and UseCase0and1Handler tests close their subscriptions in finally, so they no longer hold pooled connectors. - testBatchUseCase3 checks isUnsubscribed() on the closed queryBatch() subscription, which covers the null-subscription guard in SubscriptionImpl. --- .../api/LocalConnectorInfoManagerTests.java | 99 ++++++++++++------- .../testconnector/TstConnector.java | 4 +- 2 files changed, 67 insertions(+), 36 deletions(-) diff --git a/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java b/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java index 62fec0ba0..a6e4eccd1 100644 --- a/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java +++ b/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java @@ -247,11 +247,16 @@ public void onNext(BatchResult batchResult) { } }; - Subscription sub = facade.executeBatch(batch.build(), observer, options); - assertEquals(results.size(), batch.build().size()); - assertTrue(isComplete.get()); - assertFalse(hasError.get()); - assertNull(sub.getReturnValue()); + final Subscription sub = facade.executeBatch(batch.build(), observer, options); + try { + assertEquals(results.size(), batch.build().size()); + assertTrue(isComplete.get()); + assertFalse(hasError.get()); + assertNull(sub.getReturnValue()); + } finally { + // the subscription holds a pooled connector until it is closed + sub.close(); + } } @Test @@ -298,19 +303,27 @@ public void onNext(BatchResult batchResult) { } }; - Subscription sub = facade.executeBatch(batch.build(), observer, options); - assertEquals(results.size(), 0); - assertFalse(isComplete.get()); - assertFalse(hasError.get()); - assertNotNull(sub.getReturnValue()); - - Thread.sleep(500); - sub = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); + final Subscription sub = facade.executeBatch(batch.build(), observer, options); + try { + assertEquals(results.size(), 0); + assertFalse(isComplete.get()); + assertFalse(hasError.get()); + assertNotNull(sub.getReturnValue()); - assertEquals(results.size(), batch.build().size()); - assertTrue(isComplete.get()); - assertFalse(hasError.get()); - assertNull(sub.getReturnValue()); + Thread.sleep(500); + final Subscription query = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); + try { + assertEquals(results.size(), batch.build().size()); + assertTrue(isComplete.get()); + assertFalse(hasError.get()); + assertNull(query.getReturnValue()); + } finally { + query.close(); + } + } finally { + // the subscription holds a pooled connector until it is closed + sub.close(); + } } @Test @@ -358,19 +371,27 @@ public void onNext(BatchResult batchResult) { } }; - Subscription sub = facade.executeBatch(batch.build(), observer, options); - assertEquals(results.size(), 0); - assertFalse(isComplete.get()); - assertFalse(hasError.get()); - assertNotNull(sub.getReturnValue()); - - Thread.sleep(500); - sub = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); + final Subscription sub = facade.executeBatch(batch.build(), observer, options); + try { + assertEquals(results.size(), 0); + assertFalse(isComplete.get()); + assertFalse(hasError.get()); + assertNotNull(sub.getReturnValue()); - assertEquals(results.size(), 2); - assertFalse(isComplete.get()); - assertTrue(hasError.get()); - assertNull(sub.getReturnValue()); + Thread.sleep(500); + final Subscription query = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); + try { + assertEquals(results.size(), 2); + assertFalse(isComplete.get()); + assertTrue(hasError.get()); + assertNull(query.getReturnValue()); + } finally { + query.close(); + } + } finally { + // the subscription holds a pooled connector until it is closed + sub.close(); + } } @Test @@ -437,6 +458,9 @@ public void onNext(BatchResult batchResult) { } finally { query.close(); } + // the connector returned no subscription for the completed batch, and close() + // released the observer, so this checks the wrapper's null-subscription path + assertTrue(query.isUnsubscribed()); } finally { // the subscription holds a pooled connector until it is closed sub.close(); @@ -553,11 +577,16 @@ public void onNext(BatchResult batchResult) { } }; - Subscription sub = facade.executeBatch(batch.build(), observer, options); - assertEquals(results.size(), 2); - //Batch process is complete but handler failed to receive all - assertTrue(hasError.get() ^ isComplete.get()); - assertNull(sub.getReturnValue()); - assertTrue(sub.isUnsubscribed()); + final Subscription sub = facade.executeBatch(batch.build(), observer, options); + try { + assertEquals(results.size(), 2); + //Batch process is complete but handler failed to receive all + assertTrue(hasError.get() ^ isComplete.get()); + assertNull(sub.getReturnValue()); + assertTrue(sub.isUnsubscribed()); + } finally { + // the subscription holds a pooled connector until it is closed + sub.close(); + } } } diff --git a/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java b/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java index ec0e74116..83f15430f 100644 --- a/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java +++ b/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java @@ -488,8 +488,10 @@ public void run() { } } if (complete) { - observer.onCompleted(); + // mark the token complete first: an observer that saw onCompleted() may + // call queryBatch() at once and must find the batch complete BatchRemoteCache.setComplete(token); + observer.onCompleted(); } } catch (Exception e) { // interrupted From d0358352aa46f62e7ba8f2a1c6c9728fb48f0481 Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Tue, 6 Oct 2026 17:59:34 +0300 Subject: [PATCH 3/3] [#151] Drop the UseCase3 processor field and pin the completion order Address the second review of #153: - TstConnector: drop the processorUseCase3 field. Its processor thread set it to null when it finished; now that the tests return the connector to the pool, that late write could hit the next test that borrowed the same instance between its write and read of the field. - testBatchUseCase3 calls queryBatch() from inside onCompleted(), on the processor thread, and asserts it returns null. Reverting the setComplete()-before-onCompleted() order now fails on every run. --- .../api/LocalConnectorInfoManagerTests.java | 36 +++++++++++++++++-- .../testconnector/TstConnector.java | 6 +--- 2 files changed, 35 insertions(+), 7 deletions(-) diff --git a/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java b/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java index a6e4eccd1..f87a05d4a 100644 --- a/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java +++ b/OpenICF-java-framework/connector-framework-internal/src/test/java/org/identityconnectors/framework/impl/api/LocalConnectorInfoManagerTests.java @@ -35,6 +35,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import java.util.jar.Attributes; import java.util.jar.JarEntry; import java.util.jar.JarOutputStream; @@ -405,7 +406,7 @@ public void testBatchUseCase3() throws Exception { api.setProducerBufferSize(0); ConnectorFacadeFactory facf = ConnectorFacadeFactory.getInstance(); - ConnectorFacade facade = facf.newInstance(api); + final ConnectorFacade facade = facf.newInstance(api); OperationOptionsBuilder builder = new OperationOptionsBuilder(); builder.setOption(OperationOptions.OP_FAIL_ON_ERROR, true); @@ -420,10 +421,37 @@ public void testBatchUseCase3() throws Exception { final List results = new ArrayList(); final AtomicBoolean isComplete = new AtomicBoolean(false); final AtomicBoolean hasError = new AtomicBoolean(false); + final AtomicReference lastToken = new AtomicReference(); + final AtomicReference queryOnCompleted = new AtomicReference("not called"); + final AtomicReference queryOnCompletedError = new AtomicReference(); Observer observer = new Observer() { public void onCompleted() { - isComplete.set(true); + // runs on the processor thread inside its onCompleted() call: if the connector + // marked the token complete only after that call, queryBatch() would see it incomplete + try { + final Subscription query = facade.queryBatch(lastToken.get(), new Observer() { + public void onCompleted() { + } + + public void onError(Throwable e) { + } + + public void onNext(BatchResult batchResult) { + } + }, options); + try { + queryOnCompleted.set(query.getReturnValue()); + } finally { + query.close(); + } + } catch (Throwable t) { + // an exception here would reach the processor's catch, whose onError is dropped + // once the observer is released, so record it for the test thread instead + queryOnCompletedError.set(t); + } finally { + isComplete.set(true); + } } public void onError(Throwable e) { @@ -431,6 +459,7 @@ public void onError(Throwable e) { } public void onNext(BatchResult batchResult) { + lastToken.set(batchResult.getToken()); results.add(batchResult); if (batchResult.getError()) { hasError.set(true); @@ -451,6 +480,9 @@ public void onNext(BatchResult batchResult) { assertEquals(results.size(), batchTasks.size()); assertTrue(isComplete.get()); assertFalse(hasError.get()); + assertNull(queryOnCompletedError.get()); + // a queryBatch() made from onCompleted() finds the batch complete + assertNull(queryOnCompleted.get()); final Subscription query = facade.queryBatch((BatchToken) sub.getReturnValue(), observer, options); try { diff --git a/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java b/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java index 83f15430f..e415d8881 100644 --- a/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java +++ b/OpenICF-java-framework/testbundlev1/src/main/java/org/identityconnectors/testconnector/TstConnector.java @@ -239,8 +239,6 @@ public Schema schema() { return builder.build(); } - BatchUseCase3Processor processorUseCase3 = null; - @Override public Subscription executeBatch(final List tasks, final Observer observer, final OperationOptions options) { @@ -268,8 +266,7 @@ public Object getReturnValue() { } }; } else if (options.getOptions().containsKey("TEST_USECASE3")) { - processorUseCase3 = new BatchUseCase3Processor(); - final BatchToken token = processorUseCase3.executeBatch(tasks, options, observer); + final BatchToken token = new BatchUseCase3Processor().executeBatch(tasks, options, observer); return new Subscription() { @Override @@ -500,7 +497,6 @@ public void run() { } } BatchRemoteCache.flushTasks(token); - processorUseCase3 = null; } } }