atomics: Java Atomics and VarHandle companion code

AtomicLong/LongAdder/VarHandle counters benchmarked against the synchronized
and ReentrantLock baselines from the locks module, a deterministic ABA race
against a hand-rolled Treiber stack plus the AtomicStampedReference fix, and
a VarHandle access-modes demo (plain/opaque/acquire-release/volatile).

Co-Authored-By: Claude Sonnet 5 <[email protected]>
Claude-Session: https://claude.ai/code/session_01FhzLY5p6okFva3qsnsRyvM
This commit is contained in:
2026-09-30 06:23:03 +00:00
co-authored by Claude Sonnet 5
parent 32b8067ecf
commit ff183da163
21 changed files with 845 additions and 0 deletions
@@ -0,0 +1,104 @@
package com.ankurm.atomics;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicStampedReference;
/**
* Two parts. Part 1 forces a real ABA race against {@link TreiberStack} with
* two threads and a latch, so the corruption below is an actual observed
* race outcome, not a described one. Part 2 shows the minimal mechanism
* {@link AtomicStampedReference} uses to detect - not prevent, detect -
* exactly that race.
*/
public final class AbaProblemDemo {
public static void main(String[] args) throws InterruptedException {
part1TreiberStackAba();
System.out.println();
part2StampedReferenceDetectsIt();
}
private static void part1TreiberStackAba() throws InterruptedException {
System.out.println("=== Part 1: a real ABA race against TreiberStack ===");
TreiberStack<String> stack = new TreiberStack<>();
stack.push("C");
stack.push("B");
stack.push("A");
System.out.println("Initial stack (top first): " + stack.contentsSnapshot());
CountDownLatch t1HasRead = new CountDownLatch(1);
CountDownLatch mainHasInterfered = new CountDownLatch(1);
String[] t1Result = new String[1];
Thread t1 = new Thread(() -> {
// Read oldTop=A, newTop=B, then pause right before the CAS - exactly
// where a real thread could be preempted for an arbitrarily long time.
TreiberStack.Node<String>[] read = stack.readForPopForDemo();
TreiberStack.Node<String> oldTop = read[0];
TreiberStack.Node<String> newTop = read[1];
t1HasRead.countDown();
try {
mainHasInterfered.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
boolean success = stack.finishPopForDemo(oldTop, newTop);
t1Result[0] = success ? ("CAS succeeded, pop() would have returned \"" + oldTop.value + "\"")
: "CAS failed (this is what we WANT to see - it did not happen here)";
}, "Thread-1-stale-popper");
t1.start();
t1HasRead.await(); // Thread 1 now holds oldTop=A, newTop=B, has not CAS'd yet.
// Main thread interferes: legitimately pop A, then B (both real pop() calls),
// then push the SAME "A" node object back - simulating a pooled allocator
// that reuses freed nodes instead of always allocating fresh ones.
TreiberStack.Node<String>[] read = stack.readForPopForDemo();
TreiberStack.Node<String> nodeA = read[0];
String popped1 = stack.pop();
String popped2 = stack.pop();
System.out.println("Main thread popped, legitimately: \"" + popped1 + "\", then \"" + popped2 + "\"");
System.out.println("Stack after those two real pops: " + stack.contentsSnapshot());
stack.pushSameNodeForDemo(nodeA); // same object identity as Thread 1's oldTop
System.out.println("Main thread pushed the SAME \"A\" node object back: " + stack.contentsSnapshot());
mainHasInterfered.countDown();
t1.join();
System.out.println("Thread 1's stale CAS result: " + t1Result[0]);
System.out.println("Stack contents after Thread 1's stale CAS: " + stack.contentsSnapshot());
System.out.println("\"B\" is back in the stack even though the main thread already popped it and");
System.out.println("nobody ever pushed it again - Thread 1's CAS matched on reference identity");
System.out.println("alone (top was \"A\" both times it looked) and blindly installed a newTop");
System.out.println("(\"B\") that was computed from a read that happened before two pops and a");
System.out.println("push it never saw. \"A\" was also just handed out twice: once to the main");
System.out.println("thread's first pop(), once to Thread 1's stale one.");
}
private static void part2StampedReferenceDetectsIt() {
System.out.println("=== Part 2: AtomicStampedReference detects the same shape of race ===");
AtomicStampedReference<String> ref = new AtomicStampedReference<>("A", 0);
int[] stampHolder = new int[1];
String staleRef = ref.get(stampHolder);
int staleStamp = stampHolder[0];
System.out.println("Reader captured: ref=\"" + staleRef + "\", stamp=" + staleStamp);
// Simulate the same A -> B -> A round trip, each transition bumping the stamp -
// exactly what a real concurrent writer would do on every successful update.
ref.set("B", staleStamp + 1);
ref.set("A", staleStamp + 2);
System.out.println("After a concurrent A -> B -> A round trip: ref=\"" + ref.getReference()
+ "\", stamp=" + ref.getStamp() + " (reference is back to \"A\", but the stamp moved on)");
boolean plainWouldSucceed = staleRef.equals(ref.getReference()); // what a plain == / equals CAS would see
boolean stampedSucceeds = ref.compareAndSet(staleRef, "Z", staleStamp, staleStamp + 1);
System.out.println("A plain AtomicReference.compareAndSet(\"A\", \"Z\") would see reference == \"A\" and "
+ "succeed: " + plainWouldSucceed);
System.out.println("AtomicStampedReference.compareAndSet(\"A\", \"Z\", " + staleStamp + ", " + (staleStamp + 1)
+ ") actually succeeded: " + stampedSucceeds + " (false is correct - the stamp proves a change "
+ "happened in between, even though the reference alone looks unchanged)");
}
}
@@ -0,0 +1,23 @@
package com.ankurm.atomics;
import java.util.concurrent.atomic.AtomicLong;
/**
* {@link AtomicLong#incrementAndGet()} - a hardware compare-and-swap loop under
* the hood, retried until it succeeds. No lock, no parking, but every thread
* that loses a CAS race spins and retries against the same single contended
* memory location.
*/
public final class AtomicLongCounter implements Counter {
private final AtomicLong count = new AtomicLong();
@Override
public void increment() {
count.incrementAndGet();
}
@Override
public long get() {
return count.get();
}
}
@@ -0,0 +1,7 @@
package com.ankurm.atomics;
/** The same shared counter, protected five different ways. */
public interface Counter {
void increment();
long get();
}
@@ -0,0 +1,50 @@
package com.ankurm.atomics;
import org.openjdk.jmh.annotations.*;
import java.util.concurrent.TimeUnit;
/**
* All five counters, same write-only workload, run at a fixed thread count
* per JVM invocation via {@code -t N}; {@code scripts/run-all.sh} sweeps 1,
* 2, 4, 8, 16, 32 and 64 threads.
*/
@BenchmarkMode(Mode.Throughput)
@OutputTimeUnit(TimeUnit.MILLISECONDS)
@State(Scope.Benchmark)
@Warmup(iterations = 3, time = 1, timeUnit = TimeUnit.SECONDS)
@Measurement(iterations = 5, time = 1, timeUnit = TimeUnit.SECONDS)
@Fork(1)
public class IncrementBenchmark {
private final SynchronizedCounter synchronizedCounter = new SynchronizedCounter();
private final ReentrantLockCounter reentrantLockCounter = new ReentrantLockCounter();
private final AtomicLongCounter atomicLongCounter = new AtomicLongCounter();
private final LongAdderCounter longAdderCounter = new LongAdderCounter();
private final VarHandleCounter varHandleCounter = new VarHandleCounter();
@Benchmark
public void synchronized_() {
synchronizedCounter.increment();
}
@Benchmark
public void reentrantLock() {
reentrantLockCounter.increment();
}
@Benchmark
public void atomicLong() {
atomicLongCounter.increment();
}
@Benchmark
public void longAdder() {
longAdderCounter.increment();
}
@Benchmark
public void varHandleCas() {
varHandleCounter.increment();
}
}
@@ -0,0 +1,26 @@
package com.ankurm.atomics;
import java.util.concurrent.atomic.LongAdder;
/**
* {@link LongAdder} takes the opposite approach to {@link AtomicLongCounter}:
* instead of every thread fighting over one contended CAS location, writes
* are spread across an internal array of per-thread (really, per-probe-hash)
* cells that only get created once contention is actually detected, and
* {@code sum()} adds them all up on read. Writes get cheap; reads get more
* expensive and, crucially, not linearizable with concurrent writes - see
* the README for what that trade-off actually means.
*/
public final class LongAdderCounter implements Counter {
private final LongAdder count = new LongAdder();
@Override
public void increment() {
count.increment();
}
@Override
public long get() {
return count.sum();
}
}
@@ -0,0 +1,29 @@
package com.ankurm.atomics;
import java.util.concurrent.locks.ReentrantLock;
/** The explicit-lock baseline, for comparison against the four lock-free strategies. */
public final class ReentrantLockCounter implements Counter {
private final ReentrantLock lock = new ReentrantLock();
private long count;
@Override
public void increment() {
lock.lock();
try {
count++;
} finally {
lock.unlock();
}
}
@Override
public long get() {
lock.lock();
try {
return count;
} finally {
lock.unlock();
}
}
}
@@ -0,0 +1,16 @@
package com.ankurm.atomics;
/** Baseline: the same lock-based approach benchmarked in the {@code locks} module. */
public final class SynchronizedCounter implements Counter {
private long count;
@Override
public synchronized void increment() {
count++;
}
@Override
public synchronized long get() {
return count;
}
}
@@ -0,0 +1,107 @@
package com.ankurm.atomics;
import java.util.concurrent.atomic.AtomicReference;
/**
* A classic lock-free stack (Treiber, 1986): push and pop both loop on a
* single {@link AtomicReference#compareAndSet} against the top node. This
* textbook implementation is exactly where the textbook ABA problem lives:
* {@code pop()} reads {@code oldTop} and computes {@code newTop} from it,
* and if another thread pops that same node and later pushes the very same
* node object back on - same reference, different {@code next} underneath
* it by then - the CAS below sees the reference it expects and succeeds,
* even though the structure it is about to install ({@code newTop},
* computed from the now-stale read) is no longer correct.
* <p>
* The package-private {@code *ForDemo} methods exist only so
* {@link AbaProblemDemo} can force that exact interleaving deterministically
* - reusing the identical popped {@link Node} object on the way back in,
* which is what a real freelist or pooled-node allocator does and is the
* only way ABA is reproducible on purpose rather than by chance. The public
* {@code push}/{@code pop} API never reuses nodes and is not affected.
*/
public final class TreiberStack<T> {
static final class Node<T> {
final T value;
volatile Node<T> next;
Node(T value, Node<T> next) {
this.value = value;
this.next = next;
}
}
private final AtomicReference<Node<T>> top = new AtomicReference<>();
public void push(T value) {
Node<T> oldTop;
Node<T> newNode = new Node<>(value, null);
do {
oldTop = top.get();
newNode.next = oldTop;
} while (!top.compareAndSet(oldTop, newNode));
}
public T pop() {
Node<T> oldTop;
Node<T> newTop;
do {
oldTop = top.get();
if (oldTop == null) {
return null;
}
newTop = oldTop.next;
} while (!top.compareAndSet(oldTop, newTop));
return oldTop.value;
}
public String contentsSnapshot() {
StringBuilder sb = new StringBuilder("[");
Node<T> n = top.get();
boolean first = true;
int guard = 0;
while (n != null && guard++ < 20) {
if (!first) sb.append(", ");
sb.append(n.value);
first = false;
n = n.next;
}
sb.append("]");
return sb.toString();
}
// --- demo-only access below: never used by push()/pop() above ---
Node<T> topNodeForDemo() {
return top.get();
}
/** Reads oldTop/newTop exactly like pop() does, but stops before the CAS and hands
* both back so the demo can interleave real pop/push calls from another thread
* in between - reproducing the read-side of the race, not simulating it. */
Node<T>[] readForPopForDemo() {
@SuppressWarnings("unchecked")
Node<T>[] result = new Node[2];
result[0] = top.get(); // oldTop
result[1] = result[0] == null ? null : result[0].next; // newTop
return result;
}
/** Completes the CAS a {@link #readForPopForDemo()} call started - this is the
* exact same compareAndSet pop() itself uses, just split in two so the demo can
* inject interference in between. */
boolean finishPopForDemo(Node<T> oldTop, Node<T> newTop) {
return top.compareAndSet(oldTop, newTop);
}
/** Pushes back the SAME node object a previous pop observed, exactly as a pooled
* allocator would - the one operation that makes ABA possible. */
void pushSameNodeForDemo(Node<T> node) {
Node<T> oldTop;
do {
oldTop = top.get();
node.next = oldTop;
} while (!top.compareAndSet(oldTop, node));
}
}
@@ -0,0 +1,69 @@
package com.ankurm.atomics;
import java.lang.invoke.MethodHandles;
import java.lang.invoke.VarHandle;
/**
* {@link VarHandle} exposes four families of access mode on the same field,
* each a different point on the plain-to-volatile ordering spectrum defined
* by {@code VarHandle}'s own class documentation. This demo runs all four
* against one field and prints what each call returns - it does not and
* cannot prove the ordering guarantees themselves on a 2-vCPU single-run
* demo (that would need the kind of large concurrent campaign the jmm
* module ran with jcstress, not a coordination primitive), so treat this
* as "the API surface actually compiles and does what its Javadoc says
* about its own return values," with the ordering claims themselves
* attributed to the Javadoc in the README and the post.
*/
public final class VarHandleAccessModesDemo {
private static final VarHandle FIELD;
static {
try {
FIELD = MethodHandles.lookup()
.findVarHandle(VarHandleAccessModesDemo.class, "value", int.class);
} catch (ReflectiveOperationException e) {
throw new ExceptionInInitializerError(e);
}
}
@SuppressWarnings("unused")
private volatile int value;
public static void main(String[] args) {
VarHandleAccessModesDemo demo = new VarHandleAccessModesDemo();
// Plain: no ordering or visibility guarantee at all - same as a normal field
// read/write. Fastest, and the only mode allowed to be reordered/cached freely.
FIELD.set(demo, 1);
System.out.println("plain set(1) / get() -> " + (int) FIELD.get(demo));
// Opaque: guarantees the write is eventually visible and reads/writes to THIS
// location are not reordered with each other, but gives no happens-before
// relationship with any OTHER variable - "just don't tear or cache forever."
FIELD.setOpaque(demo, 2);
System.out.println("setOpaque(2) / getOpaque() -> " + (int) FIELD.getOpaque(demo));
// Acquire/release: the one-directional half of volatile. setRelease publishes
// everything written before it to a thread that later does a getAcquire on the
// same location - one-way happens-before, cheaper than full volatile on some
// hardware because it doesn't need a full bidirectional fence.
FIELD.setRelease(demo, 3);
System.out.println("setRelease(3) / getAcquire() -> " + (int) FIELD.getAcquire(demo));
// Volatile: full happens-before both ways, same guarantee as a `volatile` field
// or synchronized access to it - what every Counter in this module's benchmark
// that isn't plain/opaque/acquire-release actually relies on.
FIELD.setVolatile(demo, 4);
System.out.println("setVolatile(4) / getVolatile() -> " + (int) FIELD.getVolatile(demo));
// And the compare-and-swap family every lock-free Counter in this module is
// actually built on:
boolean casSucceeded = FIELD.compareAndSet(demo, 4, 5);
boolean casShouldFail = FIELD.compareAndSet(demo, 4, 6); // 4 is stale now, expect false
System.out.println("compareAndSet(4, 5) succeeded -> " + casSucceeded
+ ", second compareAndSet(4, 6) succeeded -> " + casShouldFail
+ " (expected false - value is 5, not 4, by the second call)");
}
}
@@ -0,0 +1,41 @@
package com.ankurm.atomics;
import java.lang.invoke.MethodHandles;
import java.lang.invoke.VarHandle;
/**
* The same compare-and-swap loop {@link AtomicLongCounter} does internally,
* written out by hand against a plain {@code long} field via {@link VarHandle}.
* This is what {@code AtomicLong} is built on, one layer down - no boxing,
* no extra object, just a field and a handle that knows how to fence and
* CAS against it.
*/
public final class VarHandleCounter implements Counter {
private static final VarHandle COUNT;
static {
try {
COUNT = MethodHandles.lookup()
.findVarHandle(VarHandleCounter.class, "count", long.class);
} catch (ReflectiveOperationException e) {
throw new ExceptionInInitializerError(e);
}
}
@SuppressWarnings("unused")
private volatile long count;
@Override
public void increment() {
long current;
do {
current = (long) COUNT.getVolatile(this);
} while (!COUNT.compareAndSet(this, current, current + 1));
}
@Override
public long get() {
return (long) COUNT.getVolatile(this);
}
}
@@ -0,0 +1,70 @@
package com.ankurm.atomics;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import java.util.concurrent.CountDownLatch;
import java.util.stream.IntStream;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Lost-update sanity checks for all five counters - does not test throughput
* or the ABA/memory-ordering claims, only that nothing loses an increment.
*/
class CounterCorrectnessTest {
private static final int THREADS = 8;
private static final int INCREMENTS_PER_THREAD = 50_000;
@Test @Timeout(30)
void synchronizedCounterHasNoLostUpdates() throws InterruptedException {
assertNoLostUpdates(new SynchronizedCounter());
}
@Test @Timeout(30)
void reentrantLockCounterHasNoLostUpdates() throws InterruptedException {
assertNoLostUpdates(new ReentrantLockCounter());
}
@Test @Timeout(30)
void atomicLongCounterHasNoLostUpdates() throws InterruptedException {
assertNoLostUpdates(new AtomicLongCounter());
}
@Test @Timeout(30)
void longAdderCounterHasNoLostUpdates() throws InterruptedException {
assertNoLostUpdates(new LongAdderCounter());
}
@Test @Timeout(30)
void varHandleCounterHasNoLostUpdates() throws InterruptedException {
assertNoLostUpdates(new VarHandleCounter());
}
private void assertNoLostUpdates(Counter counter) throws InterruptedException {
CountDownLatch ready = new CountDownLatch(THREADS);
CountDownLatch start = new CountDownLatch(1);
CountDownLatch done = new CountDownLatch(THREADS);
IntStream.range(0, THREADS).forEach(i -> new Thread(() -> {
ready.countDown();
try {
start.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
for (int j = 0; j < INCREMENTS_PER_THREAD; j++) {
counter.increment();
}
done.countDown();
}).start());
ready.await();
start.countDown();
done.await();
assertEquals((long) THREADS * INCREMENTS_PER_THREAD, counter.get());
}
}