From b2cd1240fec653d129122991506b34de6c49345b Mon Sep 17 00:00:00 2001 From: Florian Habermann Date: Sun, 9 Aug 2026 16:22:53 +0200 Subject: [PATCH 1/3] Make the locking policy of LockedExecutor and StripeLockedExecutor configurable Both executors created their lock via the no-arg `new ReentrantReadWriteLock()`, which is the non-fair (barging) policy, and no overload allowed a caller to choose otherwise. Under that policy an arriving writer always barges ahead of the queue, and an arriving reader only defers when the node at the head of the queue happens to be a writer, so no acquisition order is guaranteed and a waiting thread can be overtaken an unbounded number of times. The Javadoc never mentioned any of this. Fairness is now an explicit choice via `LockedExecutor.New(boolean)` and `StripeLockedExecutor.New(int, boolean)`. The default stays non-fair, so `New()` and `global()` behave exactly as before and no consumer loses throughput. `LockScope` and `StripeLockScope` gained an overridable `createExecutor()` hook so subclasses can select a policy, mirroring the existing `stripeCount()` hook. The Javadoc now documents the policy, its starvation exposure, that reentrancy holds in both modes, and that upgrading a read lock to a write lock deadlocks. Three defects in `StripeLockedExecutor.Default` are fixed along the way. `abs(mutex.hashCode()) % length` produced a negative stripe index for a mutex whose hash code is `Integer.MIN_VALUE`, since `Math.abs` of that value is itself negative; the index is now derived with `Math.floorMod`. The lock array was allocated in the constructor although its field is transient, so a deserialized instance dereferenced null; allocation moved into the lazy accessor, keyed off a non-transient stripe count. The per-element double-checked initialization relied on the array reference's volatility to publish its elements, which it does not do; the array is now populated completely before it is published. The `mutex` parameter is also checked against null, as the Javadoc already promised. The base module gains its first tests, covering the fairness wiring, reentrancy in both policies, lock release on exceptions, mutual exclusion and stripe independence, and each of the three defects. All test library versions move into the root pom's dependencyManagement so base, integration-tests and test-fixtures share them. --- base/pom.xml | 24 ++ .../serializer/concurrency/LockScope.java | 20 +- .../concurrency/LockedExecutor.java | 81 ++++- .../concurrency/StripeLockScope.java | 24 +- .../concurrency/StripeLockedExecutor.java | 133 ++++++-- .../serializer/concurrency/LockScopeTest.java | 121 ++++++++ .../concurrency/LockedExecutorTest.java | 248 +++++++++++++++ .../concurrency/StripeLockScopeTest.java | 123 ++++++++ .../concurrency/StripeLockedExecutorTest.java | 293 ++++++++++++++++++ integration-tests/pom.xml | 20 -- pom.xml | 33 ++ test-fixtures/pom.xml | 1 - 12 files changed, 1060 insertions(+), 61 deletions(-) create mode 100644 base/src/test/java/org/eclipse/serializer/concurrency/LockScopeTest.java create mode 100644 base/src/test/java/org/eclipse/serializer/concurrency/LockedExecutorTest.java create mode 100644 base/src/test/java/org/eclipse/serializer/concurrency/StripeLockScopeTest.java create mode 100644 base/src/test/java/org/eclipse/serializer/concurrency/StripeLockedExecutorTest.java diff --git a/base/pom.xml b/base/pom.xml index a413a2f6..ee2dda3b 100644 --- a/base/pom.xml +++ b/base/pom.xml @@ -24,6 +24,30 @@ org.slf4j slf4j-api + + + + org.junit.jupiter + junit-jupiter + test + + + + + org.apache.maven.plugins + maven-surefire-plugin + + + false + false + + + + + diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/LockScope.java b/base/src/main/java/org/eclipse/serializer/concurrency/LockScope.java index 033ac6f1..d3151552 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/LockScope.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/LockScope.java @@ -48,13 +48,29 @@ private LockedExecutor executor() { if((executor = this.executor) == null) { - executor = this.executor = LockedExecutor.New(); + executor = this.executor = this.createExecutor(); } } } return executor; } - + + /** + * Creates the {@link LockedExecutor} used by this scope. + *

+ * Override to control the executor's configuration, e.g. to request a fair locking policy via + * {@link LockedExecutor#New(boolean)}. + *

+ * This method is called while the monitor of this instance is held. It must therefore return + * quickly and must not acquire any other lock, otherwise it can deadlock the scope. + * + * @return a newly created {@link LockedExecutor} + */ + protected LockedExecutor createExecutor() + { + return LockedExecutor.New(); + } + /** * Executes an operation protected by a read lock. * diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java b/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java index cc13ba48..0fd5f609 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java @@ -26,9 +26,23 @@ /** * Facility to execute operations with a reentrant mutual exclusion. - * + *

+ * Reentrancy: a read operation may be nested inside another read operation, a write operation + * inside another write operation, and a read operation inside a write operation. Nesting a write + * operation inside a read operation is not supported: a {@link ReentrantReadWriteLock} cannot + * upgrade a read lock to a write lock, so such an attempt deadlocks the calling thread. This applies + * to both fairness policies. + *

+ * Fairness: instances created by {@link #New()}, and the instance returned by + * {@link #global()}, use the non-fair policy. Arriving threads may barge ahead of threads + * that are already waiting, so no acquisition order is guaranteed and an individual thread can be + * overtaken an unbounded number of times. This maximizes throughput but permits starvation under + * sustained contention. Use {@link #New(boolean)} with true to obtain a fair instance, + * which serves waiting threads in approximate arrival order. + * * @see ReentrantLock * @see ReadWriteLock + * @see #New(boolean) */ public interface LockedExecutor { @@ -105,35 +119,78 @@ private Static() * Provides a global {@link LockedExecutor} instance. *

* Only a single one exists for the whole VM process, meaning it can be used to create VM-wide locks. - * + *

+ * The shared instance uses the non-fair locking policy. Because it is shared by otherwise + * unrelated parts of the process, its contention - and therefore its exposure to starvation - is + * the sum of all its users. Prefer a dedicated instance from {@link #New(boolean)} where + * acquisition order matters. + * * @return a shared {@link LockedExecutor} instance */ public static LockedExecutor global() { return Static.sharedInstance(); } - - + + /** - * Pseudo-constructor method to create a new {@link LockedExecutor}. - * + * Pseudo-constructor method to create a new {@link LockedExecutor} with the non-fair locking policy. + *

+ * Equivalent to {@link #New(boolean) New(false)}. + * * @return a newly created {@link LockedExecutor} + * @see #New(boolean) */ public static LockedExecutor New() { - return new LockedExecutor.Default(); + return New(false); } - - + + + /** + * Pseudo-constructor method to create a new {@link LockedExecutor} with the given locking policy. + *

+ * A fair executor grants access to the longest-waiting thread and therefore cannot starve any + * thread, but its throughput is considerably lower because every hand-over has to unpark the next + * thread instead of letting an already running one barge in. A non-fair executor lets arriving + * threads barge ahead of waiting ones, which yields a much higher throughput but provides no + * ordering guarantee at all. + *

+ * Reentrancy is unaffected by this choice: a thread that already holds the lock is always let + * through, in both policies. + * + * @param fair true to use a fair locking policy, false for a non-fair one + * @return a newly created {@link LockedExecutor} + */ + public static LockedExecutor New(final boolean fair) + { + return new LockedExecutor.Default(fair); + } + + public static class Default implements LockedExecutor { + private final boolean fair; + private transient volatile ReentrantReadWriteLock reentrantLock; - Default() + Default(final boolean fair) { super(); + + this.fair = fair; } - + + /** + * Tells whether this executor's lock uses the fair policy. + * + * @return true if the locking policy is fair + */ + boolean isFair() + { + return this.fair; + } + private ReentrantReadWriteLock reentrantLock() { /* @@ -148,7 +205,7 @@ private ReentrantReadWriteLock reentrantLock() { if((reentrantLock = this.reentrantLock) == null) { - reentrantLock = this.reentrantLock = new ReentrantReadWriteLock(); + reentrantLock = this.reentrantLock = new ReentrantReadWriteLock(this.fair); } } } diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockScope.java b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockScope.java index c2c91615..a0df18e9 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockScope.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockScope.java @@ -49,23 +49,39 @@ private StripeLockedExecutor executor() { if((executor = this.executor) == null) { - executor = this.executor = StripeLockedExecutor.New(this.stripeCount()); + executor = this.executor = this.createExecutor(); } } } return executor; } - + + /** + * Creates the {@link StripeLockedExecutor} used by this scope. + *

+ * Override to control the executor's configuration, e.g. to request a fair locking policy via + * {@link StripeLockedExecutor#New(int, boolean)}. + *

+ * This method is called while the monitor of this instance is held. It must therefore return + * quickly and must not acquire any other lock, otherwise it can deadlock the scope. + * + * @return a newly created {@link StripeLockedExecutor} + */ + protected StripeLockedExecutor createExecutor() + { + return StripeLockedExecutor.New(this.stripeCount()); + } + /** * Gets the maximum number of stripes used for the {@link StripeLockedExecutor}. - * + * * @return max number of stripes */ protected int stripeCount() { return Runtime.getRuntime().availableProcessors(); } - + /** * Executes an operation protected by a read lock. * diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java index 940c7f0c..b0fff5ad 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java @@ -14,8 +14,9 @@ * #L% */ -import static java.lang.Math.abs; +import static java.lang.Math.floorMod; import static org.eclipse.serializer.math.XMath.positive; +import static org.eclipse.serializer.util.X.notNull; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantLock; @@ -29,9 +30,27 @@ /** * Facility to execute operations with a reentrant mutual exclusion for defined mutexes. - * + *

+ * Mutexes are mapped onto a fixed number of stripes by their {@link Object#hashCode()}. Two distinct + * mutexes that map onto the same stripe share a lock and therefore exclude each other, so the stripe + * count trades memory for the amount of achievable parallelism. + *

+ * Reentrancy: a read operation may be nested inside another read operation, a write operation + * inside another write operation, and a read operation inside a write operation. Nesting a write + * operation inside a read operation for the same stripe is not supported: a + * {@link ReentrantReadWriteLock} cannot upgrade a read lock to a write lock, so such an attempt + * deadlocks the calling thread. This applies to both fairness policies. + *

+ * Fairness: instances created by {@link #New(int)}, and the instance returned by + * {@link #global()}, use the non-fair policy. Arriving threads may barge ahead of threads that + * are already waiting, so no acquisition order is guaranteed and an individual thread can be + * overtaken an unbounded number of times. This maximizes throughput but permits starvation under + * sustained contention. Use {@link #New(int, boolean)} with true to obtain a fair + * instance, which serves waiting threads in approximate arrival order. + * * @see ReentrantLock * @see ReadWriteLock + * @see #New(int, boolean) */ public interface StripeLockedExecutor { @@ -114,62 +133,132 @@ private Static() * Provides a global {@link StripeLockedExecutor} instance. *

* Only a single one exists for the whole VM process, meaning it can be used to create VM-wide locks. - * + *

+ * The shared instance uses the non-fair locking policy. Because it is shared by otherwise + * unrelated parts of the process, its contention - and therefore its exposure to starvation - is + * the sum of all its users. Prefer a dedicated instance from {@link #New(int, boolean)} where + * acquisition order matters. + * * @return a shared {@link StripeLockedExecutor} instance */ public static StripeLockedExecutor global() { return Static.sharedInstance(); } - - - + + + /** - * Pseudo-constructor method to create a new {@link StripeLockedExecutor}. - * - * @param stripeCount maximum number of stripes + * Pseudo-constructor method to create a new {@link StripeLockedExecutor} with the non-fair + * locking policy. + *

+ * Equivalent to {@link #New(int, boolean) New(stripeCount, false)}. + * + * @param stripeCount maximum number of stripes, must be positive * @return a newly created {@link StripeLockedExecutor} + * @see #New(int, boolean) */ public static StripeLockedExecutor New(final int stripeCount) + { + return New(stripeCount, false); + } + + + /** + * Pseudo-constructor method to create a new {@link StripeLockedExecutor} with the given locking + * policy. + *

+ * A fair executor grants access to the longest-waiting thread and therefore cannot starve any + * thread, but its throughput is considerably lower because every hand-over has to unpark the next + * thread instead of letting an already running one barge in. A non-fair executor lets arriving + * threads barge ahead of waiting ones, which yields a much higher throughput but provides no + * ordering guarantee at all. + *

+ * Reentrancy is unaffected by this choice: a thread that already holds a stripe's lock is always + * let through, in both policies. + * + * @param stripeCount maximum number of stripes, must be positive + * @param fair true to use a fair locking policy, false for a non-fair one + * @return a newly created {@link StripeLockedExecutor} + */ + public static StripeLockedExecutor New(final int stripeCount, final boolean fair) { return new StripeLockedExecutor.Default( - positive(stripeCount) + positive(stripeCount), + fair ); } public static class Default implements StripeLockedExecutor { + private final int stripeCount; + private final boolean fair; + private transient volatile ReentrantReadWriteLock[] reentrantLocks; - Default(final int stripeCount) + Default(final int stripeCount, final boolean fair) { super(); - - this.reentrantLocks = new ReentrantReadWriteLock[stripeCount]; + + this.stripeCount = stripeCount; + this.fair = fair; } - - private ReentrantReadWriteLock reentrantLock(final Object mutex) + + /** + * Tells whether this executor's locks use the fair policy. + * + * @return true if the locking policy is fair + */ + boolean isFair() + { + return this.fair; + } + + private ReentrantReadWriteLock[] reentrantLocks() { /* * Double-checked locking to reduce the overhead of acquiring a lock * by testing the locking criterion. * The field (this.reentrantLocks) has to be volatile. + * + * The array is populated completely before it is published via the volatile write, so a + * thread that reads a non-null array is guaranteed to see all of its elements as well. + * Initializing the elements individually would not be safe, since the volatility of the + * field applies to the array reference only, not to its elements. + * + * The array cannot be created in the constructor because the field is transient and is + * therefore null again after the instance has been deserialized. */ - - final int index = abs(mutex.hashCode()) % this.reentrantLocks.length; - ReentrantReadWriteLock reentrantLock = this.reentrantLocks[index]; - if(reentrantLock == null) + ReentrantReadWriteLock[] reentrantLocks = this.reentrantLocks; + if(reentrantLocks == null) { synchronized(this) { - if((reentrantLock = this.reentrantLocks[index]) == null) + if((reentrantLocks = this.reentrantLocks) == null) { - reentrantLock = this.reentrantLocks[index] = new ReentrantReadWriteLock(); + reentrantLocks = new ReentrantReadWriteLock[this.stripeCount]; + for(int i = 0; i < reentrantLocks.length; i++) + { + reentrantLocks[i] = new ReentrantReadWriteLock(this.fair); + } + this.reentrantLocks = reentrantLocks; } } } - return reentrantLock; + return reentrantLocks; + } + + private ReentrantReadWriteLock reentrantLock(final Object mutex) + { + notNull(mutex); + + /* + * floorMod instead of abs(...) % length: abs(Integer.MIN_VALUE) is negative, + * which would yield a negative index for a mutex with that hash code. + */ + final ReentrantReadWriteLock[] reentrantLocks = this.reentrantLocks(); + return reentrantLocks[floorMod(mutex.hashCode(), reentrantLocks.length)]; } @Override diff --git a/base/src/test/java/org/eclipse/serializer/concurrency/LockScopeTest.java b/base/src/test/java/org/eclipse/serializer/concurrency/LockScopeTest.java new file mode 100644 index 00000000..3c051dee --- /dev/null +++ b/base/src/test/java/org/eclipse/serializer/concurrency/LockScopeTest.java @@ -0,0 +1,121 @@ +package org.eclipse.serializer.concurrency; + +/*- + * #%L + * Eclipse Serializer Base + * %% + * Copyright (C) 2026 MicroStream Software + * %% + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * #L% + */ + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Field; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + + +class LockScopeTest +{ + @Test + void createExecutor_notOverridden_producesNonFairExecutor() throws Exception + { + final DefaultScope scope = new DefaultScope(); + scope.touch(); + + assertFalse(((LockedExecutor.Default)executorOf(scope)).isFair()); + } + + @Test + void createExecutor_overridden_isUsed() throws Exception + { + final FairScope scope = new FairScope(); + scope.touch(); + + assertTrue(((LockedExecutor.Default)executorOf(scope)).isFair()); + } + + @Test + void executor_afterTransientFieldCleared_reinitializes() throws Exception + { + final FairScope scope = new FairScope(); + scope.touch(); + + final Field field = LockScope.class.getDeclaredField("executor"); + field.setAccessible(true); + field.set(scope, null); + + assertTrue(scope.touch(), "the scope did not re-create its executor"); + assertNotNull(executorOf(scope)); + } + + @Test + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void write_producer_returnsProducersResult() + { + assertEquals("write", new DefaultScope().writeResult()); + } + + private static LockedExecutor executorOf(final LockScope scope) throws Exception + { + final Field field = LockScope.class.getDeclaredField("executor"); + field.setAccessible(true); + return (LockedExecutor)field.get(scope); + } + + + private static class DefaultScope extends LockScope + { + DefaultScope() + { + super(); + } + + /** + * Runs a trivial read to trigger the lazy executor initialization. + * + * @return whether the action was executed + */ + boolean touch() + { + final AtomicBoolean executed = new AtomicBoolean(); + this.read(() -> executed.set(true)); + return executed.get(); + } + + String writeResult() + { + return this.write(() -> "write"); + } + + } + + + private static final class FairScope extends DefaultScope + { + FairScope() + { + super(); + } + + @Override + protected LockedExecutor createExecutor() + { + return LockedExecutor.New(true); + } + + } + +} diff --git a/base/src/test/java/org/eclipse/serializer/concurrency/LockedExecutorTest.java b/base/src/test/java/org/eclipse/serializer/concurrency/LockedExecutorTest.java new file mode 100644 index 00000000..2072d701 --- /dev/null +++ b/base/src/test/java/org/eclipse/serializer/concurrency/LockedExecutorTest.java @@ -0,0 +1,248 @@ +package org.eclipse.serializer.concurrency; + +/*- + * #%L + * Eclipse Serializer Base + * %% + * Copyright (C) 2026 MicroStream Software + * %% + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * #L% + */ + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; + +import org.eclipse.serializer.functional.Action; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + + +class LockedExecutorTest +{ + @Test + void New_noArguments_isNonFair() + { + assertFalse(((LockedExecutor.Default)LockedExecutor.New()).isFair()); + } + + @Test + void New_fair_isFair() + { + assertTrue(((LockedExecutor.Default)LockedExecutor.New(true)).isFair()); + } + + @Test + void New_nonFair_isNonFair() + { + assertFalse(((LockedExecutor.Default)LockedExecutor.New(false)).isFair()); + } + + @Test + void global_isNonFair() + { + assertFalse(((LockedExecutor.Default)LockedExecutor.global()).isFair()); + } + + @Test + void global_calledTwice_returnsSameInstance() + { + assertEquals(LockedExecutor.global(), LockedExecutor.global()); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void read_nestedRead_doesNotDeadlock(final boolean fair) + { + final LockedExecutor executor = LockedExecutor.New(fair); + final AtomicInteger depth = new AtomicInteger(); + + executor.read(() -> + { + executor.read(() -> + { + depth.incrementAndGet(); + }); + }); + + assertEquals(1, depth.get()); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void write_nestedWrite_doesNotDeadlock(final boolean fair) + { + final LockedExecutor executor = LockedExecutor.New(fair); + final AtomicInteger depth = new AtomicInteger(); + + executor.write(() -> + { + executor.write(() -> + { + depth.incrementAndGet(); + }); + }); + + assertEquals(1, depth.get()); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void write_nestedRead_doesNotDeadlock(final boolean fair) + { + final LockedExecutor executor = LockedExecutor.New(fair); + final AtomicInteger depth = new AtomicInteger(); + + executor.write(() -> + { + executor.read(() -> + { + depth.incrementAndGet(); + }); + }); + + assertEquals(1, depth.get()); + } + + @Test + void read_producer_returnsProducersResult() + { + final String result = LockedExecutor.New().read(() -> "read"); + + assertEquals("read", result); + } + + @Test + void write_producer_returnsProducersResult() + { + final String result = LockedExecutor.New().write(() -> "write"); + + assertEquals("write", result); + } + + @Test + @Timeout(value = 30, unit = TimeUnit.SECONDS) + void write_actionThrows_releasesLock() throws InterruptedException + { + final LockedExecutor executor = LockedExecutor.New(); + + assertThrows(IllegalStateException.class, () -> executor.write((Action)() -> + { + throw new IllegalStateException("expected"); + })); + + assertTrue( + this.acquiresWriteLockWithin(executor, 10, TimeUnit.SECONDS), + "the write lock was not released after the action threw" + ); + } + + @Test + @Timeout(value = 30, unit = TimeUnit.SECONDS) + void read_actionThrows_releasesLock() throws InterruptedException + { + final LockedExecutor executor = LockedExecutor.New(); + + assertThrows(IllegalStateException.class, () -> executor.read((Action)() -> + { + throw new IllegalStateException("expected"); + })); + + assertTrue( + this.acquiresWriteLockWithin(executor, 10, TimeUnit.SECONDS), + "the read lock was not released after the action threw" + ); + } + + /** + * A fair executor must hand the write lock over to a waiting writer even while readers keep + * arriving. With the non-fair policy this is not guaranteed, which is why only the fair + * direction is asserted here. + */ + @Test + @Timeout(value = 120, unit = TimeUnit.SECONDS) + void write_fairExecutorUnderSustainedReadLoad_isNotStarved() throws InterruptedException + { + final int readerCount = Math.max(4, Runtime.getRuntime().availableProcessors()); + + final LockedExecutor executor = LockedExecutor.New(true); + final AtomicBoolean keepReading = new AtomicBoolean(true); + final CountDownLatch readersUp = new CountDownLatch(readerCount); + final CountDownLatch writerDone = new CountDownLatch(1); + + final Thread[] readers = new Thread[readerCount]; + for(int i = 0; i < readerCount; i++) + { + readers[i] = new Thread(() -> + { + readersUp.countDown(); + while(keepReading.get()) + { + executor.read(() -> Thread.onSpinWait()); + } + }, "reader-" + i); + readers[i].setDaemon(true); + readers[i].start(); + } + + final boolean acquired; + try + { + assertTrue(readersUp.await(30, TimeUnit.SECONDS), "the readers did not start"); + + final Thread writer = new Thread( + () -> executor.write(() -> writerDone.countDown()), + "writer" + ); + writer.setDaemon(true); + writer.start(); + + acquired = writerDone.await(30, TimeUnit.SECONDS); + } + finally + { + keepReading.set(false); + for(final Thread reader : readers) + { + reader.join(TimeUnit.SECONDS.toMillis(10)); + } + } + + assertTrue(acquired, "the fair write lock was starved by the sustained read load"); + } + + private boolean acquiresWriteLockWithin( + final LockedExecutor executor, + final long timeout, + final TimeUnit unit + ) + throws InterruptedException + { + final CountDownLatch acquired = new CountDownLatch(1); + + final Thread thread = new Thread(() -> executor.write(() -> acquired.countDown()), "acquirer"); + thread.setDaemon(true); + thread.start(); + + final boolean result = acquired.await(timeout, unit); + thread.join(unit.toMillis(timeout)); + return result; + } + +} diff --git a/base/src/test/java/org/eclipse/serializer/concurrency/StripeLockScopeTest.java b/base/src/test/java/org/eclipse/serializer/concurrency/StripeLockScopeTest.java new file mode 100644 index 00000000..393c3eb8 --- /dev/null +++ b/base/src/test/java/org/eclipse/serializer/concurrency/StripeLockScopeTest.java @@ -0,0 +1,123 @@ +package org.eclipse.serializer.concurrency; + +/*- + * #%L + * Eclipse Serializer Base + * %% + * Copyright (C) 2026 MicroStream Software + * %% + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * #L% + */ + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Field; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + + +class StripeLockScopeTest +{ + @Test + void createExecutor_notOverridden_producesNonFairExecutor() throws Exception + { + final DefaultScope scope = new DefaultScope(); + scope.touch(); + + assertFalse(((StripeLockedExecutor.Default)executorOf(scope)).isFair()); + } + + @Test + void createExecutor_overridden_isUsed() throws Exception + { + final FairScope scope = new FairScope(); + scope.touch(); + + assertTrue(((StripeLockedExecutor.Default)executorOf(scope)).isFair()); + } + + @Test + void executor_afterTransientFieldCleared_reinitializes() throws Exception + { + final FairScope scope = new FairScope(); + scope.touch(); + + final Field field = StripeLockScope.class.getDeclaredField("executor"); + field.setAccessible(true); + field.set(scope, null); + + assertTrue(scope.touch(), "the scope did not re-create its executor"); + assertNotNull(executorOf(scope)); + } + + @Test + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void write_producer_returnsProducersResult() + { + assertEquals("write", new DefaultScope().writeResult()); + } + + private static StripeLockedExecutor executorOf(final StripeLockScope scope) throws Exception + { + final Field field = StripeLockScope.class.getDeclaredField("executor"); + field.setAccessible(true); + return (StripeLockedExecutor)field.get(scope); + } + + + private static class DefaultScope extends StripeLockScope + { + private final Object mutex = new Object(); + + DefaultScope() + { + super(); + } + + /** + * Runs a trivial read to trigger the lazy executor initialization. + * + * @return whether the action was executed + */ + boolean touch() + { + final AtomicBoolean executed = new AtomicBoolean(); + this.read(this.mutex, () -> executed.set(true)); + return executed.get(); + } + + String writeResult() + { + return this.write(this.mutex, () -> "write"); + } + + } + + + private static final class FairScope extends DefaultScope + { + FairScope() + { + super(); + } + + @Override + protected StripeLockedExecutor createExecutor() + { + return StripeLockedExecutor.New(this.stripeCount(), true); + } + + } + +} diff --git a/base/src/test/java/org/eclipse/serializer/concurrency/StripeLockedExecutorTest.java b/base/src/test/java/org/eclipse/serializer/concurrency/StripeLockedExecutorTest.java new file mode 100644 index 00000000..70646ae3 --- /dev/null +++ b/base/src/test/java/org/eclipse/serializer/concurrency/StripeLockedExecutorTest.java @@ -0,0 +1,293 @@ +package org.eclipse.serializer.concurrency; + +/*- + * #%L + * Eclipse Serializer Base + * %% + * Copyright (C) 2026 MicroStream Software + * %% + * This program and the accompanying materials are made + * available under the terms of the Eclipse Public License 2.0 + * which is available at https://www.eclipse.org/legal/epl-2.0/ + * + * SPDX-License-Identifier: EPL-2.0 + * #L% + */ + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Field; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; + +import org.eclipse.serializer.exceptions.NumberRangeException; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + + +class StripeLockedExecutorTest +{ + @Test + void New_stripeCountOnly_isNonFair() + { + assertFalse(((StripeLockedExecutor.Default)StripeLockedExecutor.New(4)).isFair()); + } + + @Test + void New_fair_isFair() + { + assertTrue(((StripeLockedExecutor.Default)StripeLockedExecutor.New(4, true)).isFair()); + } + + @Test + void global_isNonFair() + { + assertFalse(((StripeLockedExecutor.Default)StripeLockedExecutor.global()).isFair()); + } + + @Test + void New_stripeCountIsZero_throwsNumberRangeException() + { + assertThrows(NumberRangeException.class, () -> StripeLockedExecutor.New(0)); + } + + @Test + void New_stripeCountIsNegative_throwsNumberRangeException() + { + assertThrows(NumberRangeException.class, () -> StripeLockedExecutor.New(-1)); + } + + /** + * {@link Math#abs(int)} of {@link Integer#MIN_VALUE} is negative, so deriving the stripe index + * that way produces an {@link ArrayIndexOutOfBoundsException} for such a mutex. + *

+ * The stripe count must not be a divisor of {@link Integer#MIN_VALUE}, since the remainder would + * then be zero and mask the defect. Hence 3 rather than a power of two. + */ + @Test + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void read_mutexHashCodeIsIntegerMinValue_doesNotThrow() + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(3); + final AtomicBoolean executed = new AtomicBoolean(); + + executor.read(new FixedHashCode(Integer.MIN_VALUE), () -> executed.set(true)); + + assertTrue(executed.get()); + } + + /** + * @see #read_mutexHashCodeIsIntegerMinValue_doesNotThrow() + */ + @Test + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void write_mutexHashCodeIsIntegerMinValue_doesNotThrow() + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(3); + final AtomicBoolean executed = new AtomicBoolean(); + + executor.write(new FixedHashCode(Integer.MIN_VALUE), () -> executed.set(true)); + + assertTrue(executed.get()); + } + + @Test + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void read_mutexHashCodeIsNegative_doesNotThrow() + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(3); + final AtomicBoolean executed = new AtomicBoolean(); + + executor.read(new FixedHashCode(-7), () -> executed.set(true)); + + assertTrue(executed.get()); + } + + /** + * Every stripe index has to stay within the array's bounds, for any hash code. + */ + @Test + @Timeout(value = 30, unit = TimeUnit.SECONDS) + void read_extremeHashCodes_stayWithinBounds() + { + final int[] hashCodes = { + Integer.MIN_VALUE, Integer.MIN_VALUE + 1, -1234567, -3, -2, -1, + 0, 1, 2, 3, 1234567, Integer.MAX_VALUE - 1, Integer.MAX_VALUE + }; + + for(int stripeCount = 1; stripeCount <= 9; stripeCount++) + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(stripeCount); + for(final int hashCode : hashCodes) + { + final AtomicBoolean executed = new AtomicBoolean(); + executor.read(new FixedHashCode(hashCode), () -> executed.set(true)); + assertTrue(executed.get(), "hashCode " + hashCode + " / stripeCount " + stripeCount); + } + } + } + + @Test + void read_nullMutex_throwsNullPointerException() + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(4); + + assertThrows(NullPointerException.class, () -> executor.read(null, () -> + { + // must not be reached + })); + } + + @Test + void write_nullMutex_throwsNullPointerException() + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(4); + + assertThrows(NullPointerException.class, () -> executor.write(null, () -> + { + // must not be reached + })); + } + + /** + * The locks are held in a transient field, which is {@code null} again after the instance has + * been deserialized. The executor has to re-create them instead of failing. + */ + @Test + @Timeout(value = 10, unit = TimeUnit.SECONDS) + void read_afterTransientLocksFieldCleared_reinitializesLocks() throws Exception + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(4); + final Object mutex = new Object(); + + executor.read(mutex, () -> + { + // trigger the lazy initialization + }); + + final Field field = StripeLockedExecutor.Default.class.getDeclaredField("reentrantLocks"); + field.setAccessible(true); + field.set(executor, null); + + final AtomicBoolean executed = new AtomicBoolean(); + executor.read(mutex, () -> executed.set(true)); + + assertTrue(executed.get()); + } + + @Test + @Timeout(value = 60, unit = TimeUnit.SECONDS) + void write_distinctStripes_proceedConcurrently() throws InterruptedException + { + final StripeLockedExecutor executor = StripeLockedExecutor.New(4); + + // 0 and 1 map onto distinct stripes for a stripe count of 4 + final Object mutexA = new FixedHashCode(0); + final Object mutexB = new FixedHashCode(1); + + final CyclicBarrier barrier = new CyclicBarrier(2); + final CountDownLatch bothInside = new CountDownLatch(2); + + final Thread threadA = startDaemon("stripe-a", () -> executor.write(mutexA, () -> + { + awaitBarrier(barrier); + bothInside.countDown(); + })); + final Thread threadB = startDaemon("stripe-b", () -> executor.write(mutexB, () -> + { + awaitBarrier(barrier); + bothInside.countDown(); + })); + + final boolean concurrent = bothInside.await(30, TimeUnit.SECONDS); + + threadA.join(TimeUnit.SECONDS.toMillis(10)); + threadB.join(TimeUnit.SECONDS.toMillis(10)); + + assertTrue(concurrent, "writes on distinct stripes did not proceed concurrently"); + } + + @Test + @Timeout(value = 60, unit = TimeUnit.SECONDS) + void write_sameMutex_isMutuallyExclusive() throws InterruptedException + { + final int threadCount = 8; + final int iterations = 500; + + final StripeLockedExecutor executor = StripeLockedExecutor.New(4); + final Object mutex = new Object(); + final AtomicInteger inside = new AtomicInteger(); + final AtomicInteger maxInside = new AtomicInteger(); + final CountDownLatch done = new CountDownLatch(threadCount); + + for(int i = 0; i < threadCount; i++) + { + startDaemon("writer-" + i, () -> + { + for(int n = 0; n < iterations; n++) + { + executor.write(mutex, () -> + { + maxInside.accumulateAndGet(inside.incrementAndGet(), Math::max); + Thread.onSpinWait(); + inside.decrementAndGet(); + }); + } + done.countDown(); + }); + } + + assertTrue(done.await(30, TimeUnit.SECONDS), "the writers did not finish"); + assertEquals(1, maxInside.get(), "more than one thread held the same stripe's write lock"); + } + + private static Thread startDaemon(final String name, final Runnable runnable) + { + final Thread thread = new Thread(runnable, name); + thread.setDaemon(true); + thread.start(); + return thread; + } + + private static void awaitBarrier(final CyclicBarrier barrier) + { + try + { + barrier.await(20, TimeUnit.SECONDS); + } + catch(final Exception e) + { + throw new RuntimeException(e); + } + } + + + /** + * Mutex with a controllable hash code, to target a specific stripe. + */ + private static final class FixedHashCode + { + private final int hashCode; + + FixedHashCode(final int hashCode) + { + super(); + + this.hashCode = hashCode; + } + + @Override + public int hashCode() + { + return this.hashCode; + } + + } + +} diff --git a/integration-tests/pom.xml b/integration-tests/pom.xml index 9172a79c..5afb51f2 100644 --- a/integration-tests/pom.xml +++ b/integration-tests/pom.xml @@ -19,22 +19,6 @@ Serializer-level integration test suite (collections, serialization, communication, layered entities). Test-only, not deployed. https://projects.eclipse.org/projects/technology.serializer - - 6.0.3 - - - - - - org.junit - junit-bom - ${junit-jupiter.version} - pom - import - - - - @@ -83,25 +67,21 @@ commons-io commons-io - 2.21.0 test org.javassist javassist - 3.30.2-GA test org.slf4j slf4j-simple - ${slf4j.version} test net.datafaker datafaker - 2.5.3 test diff --git a/pom.xml b/pom.xml index 802153d8..5f4cc49c 100644 --- a/pom.xml +++ b/pom.xml @@ -47,6 +47,10 @@ false 2 2.0.17 + 6.0.3 + 2.21.0 + 3.30.2-GA + 2.5.3 @@ -81,6 +85,35 @@ slf4j-api ${slf4j.version} + + + + org.junit + junit-bom + ${junit-jupiter.version} + pom + import + + + org.slf4j + slf4j-simple + ${slf4j.version} + + + commons-io + commons-io + ${commons-io.version} + + + org.javassist + javassist + ${javassist.version} + + + net.datafaker + datafaker + ${datafaker.version} + diff --git a/test-fixtures/pom.xml b/test-fixtures/pom.xml index 881b68b2..7a2bc7bb 100644 --- a/test-fixtures/pom.xml +++ b/test-fixtures/pom.xml @@ -30,7 +30,6 @@ org.junit.jupiter junit-jupiter-api - 6.0.3 true From b5089ca1211e87450ba217c815f24395587ff5cc Mon Sep 17 00:00:00 2001 From: Florian Habermann Date: Sun, 9 Aug 2026 16:32:22 +0200 Subject: [PATCH 2/3] Do not claim that a fair executor cannot starve a thread The Javadoc of both `New(boolean)` overloads stated that a fair executor "grants access to the longest-waiting thread and therefore cannot starve any thread". That is stronger than what the underlying lock guarantees: `ReentrantReadWriteLock` documents fair mode as an approximate arrival-order policy, and `ReentrantLock` explicitly notes that the fairness of a lock does not imply fairness of thread scheduling. The wording now says the lock is handed over in approximate arrival order so a waiting thread is not overtaken indefinitely, and adds an explicit note that a fair policy orders the hand-over of the lock itself and cannot influence how the JVM or the operating system schedules the contending threads. --- .../serializer/concurrency/LockedExecutor.java | 14 +++++++++----- .../concurrency/StripeLockedExecutor.java | 14 +++++++++----- 2 files changed, 18 insertions(+), 10 deletions(-) diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java b/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java index 0fd5f609..f54c652d 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/LockedExecutor.java @@ -150,11 +150,15 @@ public static LockedExecutor New() /** * Pseudo-constructor method to create a new {@link LockedExecutor} with the given locking policy. *

- * A fair executor grants access to the longest-waiting thread and therefore cannot starve any - * thread, but its throughput is considerably lower because every hand-over has to unpark the next - * thread instead of letting an already running one barge in. A non-fair executor lets arriving - * threads barge ahead of waiting ones, which yields a much higher throughput but provides no - * ordering guarantee at all. + * A fair executor hands the lock over in approximate arrival order, so a waiting thread is not + * overtaken indefinitely, but its throughput is considerably lower because every hand-over has to + * unpark the next thread instead of letting an already running one barge in. A non-fair executor + * lets arriving threads barge ahead of waiting ones, which yields a much higher throughput but + * provides no ordering guarantee at all. + *

+ * Note that the fairness of a lock does not extend to the scheduling of threads, as documented on + * {@link ReentrantLock}. A fair policy orders the hand-over of the lock itself; it cannot prevent + * the JVM or the operating system from scheduling the contending threads unevenly. *

* Reentrancy is unaffected by this choice: a thread that already holds the lock is always let * through, in both policies. diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java index b0fff5ad..db3e1a70 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java @@ -168,11 +168,15 @@ public static StripeLockedExecutor New(final int stripeCount) * Pseudo-constructor method to create a new {@link StripeLockedExecutor} with the given locking * policy. *

- * A fair executor grants access to the longest-waiting thread and therefore cannot starve any - * thread, but its throughput is considerably lower because every hand-over has to unpark the next - * thread instead of letting an already running one barge in. A non-fair executor lets arriving - * threads barge ahead of waiting ones, which yields a much higher throughput but provides no - * ordering guarantee at all. + * A fair executor hands a stripe's lock over in approximate arrival order, so a waiting thread is + * not overtaken indefinitely, but its throughput is considerably lower because every hand-over has + * to unpark the next thread instead of letting an already running one barge in. A non-fair + * executor lets arriving threads barge ahead of waiting ones, which yields a much higher + * throughput but provides no ordering guarantee at all. + *

+ * Note that the fairness of a lock does not extend to the scheduling of threads, as documented on + * {@link ReentrantLock}. A fair policy orders the hand-over of the lock itself; it cannot prevent + * the JVM or the operating system from scheduling the contending threads unevenly. *

* Reentrancy is unaffected by this choice: a thread that already holds a stripe's lock is always * let through, in both policies. From 8cfec1c23fcc80fbb2dd0a06374d585d33050bc2 Mon Sep 17 00:00:00 2001 From: Florian Habermann Date: Sun, 9 Aug 2026 16:38:27 +0200 Subject: [PATCH 3/3] Document that a stripe executor allocates all of its locks at once The lock array is now populated completely before it is published, so all stripes are created the first time the executor is used rather than one by one on demand. That is a deliberate trade for a correctly published array and a hot path without a CAS, but it does mean an oversized stripe count is paid for up front. Say so on `New(int, boolean)` and point at the degree of parallelism as the yardstick for sizing, so the parameter is not mistaken for the number of mutexes the executor will see. --- .../eclipse/serializer/concurrency/StripeLockedExecutor.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java index db3e1a70..aca35830 100644 --- a/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java +++ b/base/src/main/java/org/eclipse/serializer/concurrency/StripeLockedExecutor.java @@ -180,6 +180,11 @@ public static StripeLockedExecutor New(final int stripeCount) *

* Reentrancy is unaffected by this choice: a thread that already holds a stripe's lock is always * let through, in both policies. + *

+ * The locks of all stripes are created together, the first time the executor is used. The stripe + * count should therefore be sized like a degree of parallelism - the default of + * {@link Runtime#availableProcessors()} used by {@link StripeLockScope} is a good yardstick - and + * not like the number of mutexes the executor will ever see. * * @param stripeCount maximum number of stripes, must be positive * @param fair true to use a fair locking policy, false for a non-fair one