This is an automated email from the ASF dual-hosted git repository.
tabish121 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/qpid-protonj2.git
The following commit(s) were added to refs/heads/main by this push:
new f11d953f PROTON-2963 Allow disposition range to include wrapped
delivery Ids
f11d953f is described below
commit f11d953fb453dce6ec784746a0bb24fc1f59cee1
Author: Timothy Bish <[email protected]>
AuthorDate: Tue Sep 1 18:35:55 2026 -0400
PROTON-2963 Allow disposition range to include wrapped delivery Ids
Handles correctly the cases where the range is an overflow such that the
low ID is large than the high and the tracker needs to scan into the next
set of IDs after an overflow.
---
.../qpid/protonj2/engine/util/UnsettledMap.java | 112 +-
.../protonj2/engine/util/UnsettledMapTest.java | 1887 +++++++++++++-------
2 files changed, 1315 insertions(+), 684 deletions(-)
diff --git
a/protonj2/src/main/java/org/apache/qpid/protonj2/engine/util/UnsettledMap.java
b/protonj2/src/main/java/org/apache/qpid/protonj2/engine/util/UnsettledMap.java
index d7fe4adb..7651baaa 100644
---
a/protonj2/src/main/java/org/apache/qpid/protonj2/engine/util/UnsettledMap.java
+++
b/protonj2/src/main/java/org/apache/qpid/protonj2/engine/util/UnsettledMap.java
@@ -295,22 +295,58 @@ public class UnsettledMap<Delivery> implements
Map<UnsignedInteger, Delivery> {
return;
}
- int readStart = -1;
+ // There are two cases here, either first <= last or first > last which
+ // involves two cycles one for each range before and after the
overflow.
+ if (Integer.compareUnsigned(first, last) <= 0) {
+ forEachInRange(tail, false, first, last, action);
+ } else {
+ final UnsettledBucket<Delivery> nextGeneration =
+ forEachInRange(tail, false, first,
UnsignedInteger.MAX_VALUE.intValue(), action);
+
+ if (nextGeneration != null) {
+ forEachInRange(nextGeneration, true, 0, last, action);
+ }
+ }
+ }
+
+ private UnsettledBucket<Delivery> forEachInRange(UnsettledBucket<Delivery>
tail, boolean continuation, int first, int last, Consumer<Delivery> action) {
boolean foundFirst = false;
boolean foundLast = false;
+ int inGeneration = continuation ? tail.generation : -1;
- for (UnsettledBucket<Delivery> bucket = tail; bucket != null &&
!foundLast; bucket = bucket.next) {
+ UnsettledBucket<Delivery> bucket = tail;
+
+ for (; bucket != null && !foundLast; bucket = bucket.next) {
final int writeOffset = bucket.writeOffset;
- readStart = bucket.readOffset;
+ int readStart = bucket.readOffset;
+
+ // We are locked to a ranged search in a single generation so if
values
+ // jump we can stop and return where the next generation starts in
case
+ // the caller has a second range to check.
+ if (inGeneration >= 0 && bucket.generation != inGeneration) {
+ return bucket;
+ }
+
+ if (!foundFirst) {
+ if (bucket.isCapturedByRange(first, last)) {
+ final int result = bucket.search(first);
+ final int ceiling = result >= 0 ? result : ~result;
- if (!foundFirst && bucket.isCapturedByRange(first, last)) {
- final int result = bucket.search(first);
- final int ceiling = result >= 0 ? result : ~result;
+ // We either found the actual first or we found the
location where it would
+ // be inserted at which is the next logical value we will
use as the start
+ // and proceed with all values from that point until we
find the end or an
+ // element greater than the last value.
- if (ceiling < writeOffset) {
- foundFirst = true;
- readStart = ceiling;
+ if (ceiling < writeOffset) {
+ foundFirst = true;
+ readStart = ceiling;
+ inGeneration = bucket.generation;
+ }
+ } else if (continuation && bucket.lowestDeliveryId > first) {
+ // If continued from a range that include overflow and we
find values ahead of
+ // the lowest point in this sequence we know there can be
no matches so stop now.
+ return null;
}
}
@@ -330,6 +366,8 @@ public class UnsettledMap<Delivery> implements
Map<UnsignedInteger, Delivery> {
}
}
}
+
+ return bucket;
}
/**
@@ -351,23 +389,60 @@ public class UnsettledMap<Delivery> implements
Map<UnsignedInteger, Delivery> {
return;
}
+ // There are two cases here, either first <= last or first > last which
+ // involves two cycles one for each range before and after the
overflow.
+ if (Integer.compareUnsigned(first, last) <= 0) {
+ removeEachInRange(tail, false, first, last, action);
+ } else {
+ final UnsettledBucket<Delivery> nextGeneration =
+ removeEachInRange(tail, false, first,
UnsignedInteger.MAX_VALUE.intValue(), action);
+
+ if (nextGeneration != null) {
+ removeEachInRange(nextGeneration, true, 0, last, action);
+ }
+ }
+ }
+
+ private UnsettledBucket<Delivery>
removeEachInRange(UnsettledBucket<Delivery> tail, boolean continuation, int
first, int last, Consumer<Delivery> action) {
boolean foundFirst = false;
boolean foundLast = false;
int removeStart = 0;
int removeEnd = 0;
+ int inGeneration = continuation ? tail.generation : -1;
- for (UnsettledBucket<Delivery> bucket = tail; bucket != null &&
!foundLast; ) {
+ UnsettledBucket<Delivery> bucket = tail;
+
+ for (; bucket != null && !foundLast; ) {
final int writeOffset = bucket.writeOffset;
removeStart = bucket.readOffset;
- if (!foundFirst && bucket.isCapturedByRange(first, last)) {
- final int result = bucket.search(first);
- final int ceiling = result >= 0 ? result : ~result;
+ // We are locked to a ranged search in a single generation so if
values
+ // jump we can stop and return where the next generation starts in
case
+ // the caller has a second range to check.
+ if (inGeneration >= 0 && bucket.generation != inGeneration) {
+ return bucket;
+ }
+
+ if (!foundFirst) {
+ if (bucket.isCapturedByRange(first, last)) {
+ final int result = bucket.search(first);
+ final int ceiling = result >= 0 ? result : ~result;
- if (ceiling < writeOffset) {
- foundFirst = true;
- removeStart = ceiling;
+ // We either found the actual first or we found the
location where it would
+ // be inserted at which is the next logical value we will
use as the start
+ // and proceed with all values from that point until we
find the end or an
+ // element greater than the last value.
+
+ if (ceiling < writeOffset) {
+ foundFirst = true;
+ removeStart = ceiling;
+ inGeneration = bucket.generation;
+ }
+ } else if (continuation && bucket.lowestDeliveryId > first) {
+ // If continued from a range that include overflow and we
find values ahead of
+ // the lowest point in this sequence we know there can be
no matches so stop now.
+ return null;
}
}
@@ -391,6 +466,8 @@ public class UnsettledMap<Delivery> implements
Map<UnsignedInteger, Delivery> {
bucket = bucket.next;
}
}
+
+ return bucket;
}
@Override
@@ -894,8 +971,7 @@ public class UnsettledMap<Delivery> implements
Map<UnsignedInteger, Delivery> {
/**
* Checks if the given range of delivery IDs potentially captures any
entries in
- * this bucket by checking if the lowest delivery ID in this bucket is
between the
- * given low and high values.
+ * this bucket by checking if the the opposing ends of the ranges
overlap.
*
* @param lowest
* The lowest value that is being searched for.
diff --git
a/protonj2/src/test/java/org/apache/qpid/protonj2/engine/util/UnsettledMapTest.java
b/protonj2/src/test/java/org/apache/qpid/protonj2/engine/util/UnsettledMapTest.java
index 45f3a37c..034df0fd 100644
---
a/protonj2/src/test/java/org/apache/qpid/protonj2/engine/util/UnsettledMapTest.java
+++
b/protonj2/src/test/java/org/apache/qpid/protonj2/engine/util/UnsettledMapTest.java
@@ -795,86 +795,6 @@ public class UnsettledMapTest {
assertEquals(6, tracker.size());
}
- @Test
- public void testForEachDeliveryIteratesOverLargeSeriesOfDeliveries() {
- UnsettledMap<DeliveryType> tracker = createMap();
- assertEquals(0, tracker.size());
-
- final int COUNT = 4080;
-
- for (int i = 0; i < COUNT; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
-
- assertEquals(COUNT, tracker.size());
-
- final AtomicInteger index = new AtomicInteger();
-
- tracker.forEach((delivery) -> index.incrementAndGet());
-
- assertEquals(index.get(), COUNT);
- }
-
- @Test
- public void
testForEachBiConsumerDeliveryIteratesOverLargeSeriesOfDeliveries() {
- UnsettledMap<DeliveryType> tracker = createMap();
- assertEquals(0, tracker.size());
-
- final int COUNT = 4080;
-
- for (int i = 0; i < COUNT; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
-
- assertEquals(COUNT, tracker.size());
-
- final AtomicInteger index = new AtomicInteger();
-
- tracker.forEach((deliveryId, delivery) -> index.incrementAndGet());
-
- assertEquals(index.get(), COUNT);
- }
-
- @Test
- public void testRangedForEachDeliveryIteratesOverSmallSeriesOfDeliveries()
{
- UnsettledMap<DeliveryType> tracker = createMap();
- assertEquals(0, tracker.size());
-
- final int COUNT = 512;
-
- for (int i = 0; i < COUNT; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
-
- assertEquals(COUNT, tracker.size());
-
- final AtomicInteger index = new AtomicInteger();
-
- tracker.forEach(260, 262, (delivery) -> index.incrementAndGet());
-
- assertEquals(index.get(), 3);
- }
-
- @Test
- public void
testRangedForEachDeliveryIteratesSeriesWhenValuesOverflowIntRange() {
- UnsettledMap<DeliveryType> tracker = createMap();
- assertEquals(0, tracker.size());
-
- tracker.put(0, new DeliveryType(0));
- tracker.put(1, new DeliveryType(1));
- tracker.put(Integer.MAX_VALUE, new DeliveryType(Integer.MAX_VALUE));
- tracker.put(Integer.MAX_VALUE + 1, new DeliveryType(Integer.MAX_VALUE
+ 1));
- tracker.put(Integer.MAX_VALUE + 2, new DeliveryType(Integer.MAX_VALUE
+ 2));
- tracker.put(Integer.MAX_VALUE + 3, new DeliveryType(Integer.MAX_VALUE
+ 3));
- tracker.put(Integer.MAX_VALUE + 4, new DeliveryType(Integer.MAX_VALUE
+ 4));
-
- final AtomicInteger index = new AtomicInteger();
-
- tracker.forEach(Integer.MAX_VALUE, Integer.MAX_VALUE + 2, (delivery)
-> index.incrementAndGet());
-
- assertEquals(3, index.get());
- }
-
@SuppressWarnings("unlikely-arg-type")
@Test
public void testValuesCollection() {
@@ -1467,25 +1387,6 @@ public class UnsettledMapTest {
}
}
- @Test
- public void testForEachEntry() {
- UnsettledMap<DeliveryType> tracker = createMap();
-
- final int[] inputValues = {3, 0, -1, 1, -2, 2};
-
- for (int entry : inputValues) {
- tracker.put(entry, new DeliveryType(entry));
- }
-
- final SequenceNumber index = new SequenceNumber(0);
- tracker.forEach((value) -> {
- int i = index.getAndIncrement().intValue();
- assertEquals(new DeliveryType(inputValues[i]), value);
- });
-
- assertEquals(index.intValue(), inputValues.length);
- }
-
@Test
public void testRandomProduceAndConsumeWithBacklog() {
UnsettledMap<DeliveryType> tracker = createMap();
@@ -1935,141 +1836,661 @@ public class UnsettledMapTest {
}
@Test
- public void testRemoveRangeRemovesNoValues() {
- final int afterLast = uintArray.length + 1; // Entries are one based
-
- final AtomicBoolean removed = new AtomicBoolean();
-
- tracker.removeEach(afterLast, afterLast + 10, (delivery) ->
removed.set(true));
-
- assertFalse(removed.get());
- }
+ public void testRemoveFromIteratorFromMiddleBucket() {
+ final int numBuckets = 5;
+ final int bucketSize = 10;
+ final int numEntries = numBuckets * bucketSize;
- @Test
- public void testRemoveRangeRemovesLastValue() {
- final int lastEntry = uintArray.length; // Entries are one based
+ final UnsettledMap<DeliveryType> map = createMap(numBuckets,
bucketSize);
- final AtomicInteger removed = new AtomicInteger();
+ for (int i = 0; i < numEntries; ++i) {
+ map.put(i, new DeliveryType(i));
+ }
- tracker.removeEach(lastEntry, lastEntry, (delivery) ->
removed.incrementAndGet());
+ assertEquals(numEntries, map.size());
- assertEquals(1, removed.get());
- }
+ Iterator<UnsignedInteger> entries = map.keySet().iterator();
- @Test
- public void
testRemoveRangeRemovesLastValueAndRangeOutsideOfActualEntries() {
- final int lastEntry = uintArray.length; // Entries are one based
+ // Move to center of bucket two
+ for (int i = 0; i < bucketSize + (bucketSize / 2); ++i) {
+ entries.next();
+ }
- final AtomicInteger removed = new AtomicInteger();
+ UnsignedInteger lastValue = null;
- tracker.removeEach(lastEntry, lastEntry + 10, (delivery) ->
removed.incrementAndGet());
+ // Remove from center of bucket two into bucket three until a
compaction event should occur.
+ for (int i = 0; i < bucketSize; ++i) {
+ lastValue = entries.next();
+ entries.remove();
+ }
- assertEquals(1, removed.get());
+ assertEquals(lastValue.intValue() + 1, entries.next().intValue());
}
@Test
- public void testRemoveEachWithRangeThatMatchesMany() {
- final UnsettledMap<DeliveryType> map = createMap(5, 50);
+ public void testRemoveFromMiddleBucket() {
+ final int numBuckets = 5;
+ final int bucketSize = 10;
+ final int numEntries = numBuckets * bucketSize;
- final int[] entries = new int[] { 512, 513, 512, 513, 512, 513 };
+ final UnsettledMap<DeliveryType> map = createMap(numBuckets,
bucketSize);
- for(int i : entries) {
+ for (int i = 0; i < numEntries; ++i) {
map.put(i, new DeliveryType(i));
}
- assertEquals(entries.length, map.size());
+ assertEquals(numEntries, map.size());
- final AtomicInteger removed = new AtomicInteger();
+ int position = bucketSize + (bucketSize / 2);
+ int lastValue = 0;
- map.removeEach(512, 513, (delivery) -> removed.incrementAndGet());
+ // Remove from center of bucket two into bucket three until a
compaction event should occur.
+ for (int i = 0; i < bucketSize + 3; ++i) {
+ lastValue = map.get(position).getDeliveryId();
+ map.remove(position++);
+ }
- assertEquals(2, removed.get());
- assertEquals(entries.length - 2, map.size());
+ assertEquals(lastValue + 1, map.get(position).getDeliveryId());
}
@Test
- public void testRemoveAllEntriesFromFirstBucket() {
- doTestRemoveEach(0, 15);
- }
+ public void testRepeatedRemoveOldestHotPath() {
+ UnsettledMap<DeliveryType> tracker = createMap();
- @Test
- public void testRemoveAllEntriesFromMiddleBucket() {
- doTestRemoveEach(16, 31);
- }
+ final int COUNT = 10000;
- @Test
- public void testRemoveAllEntriesFromEndBucket() {
- doTestRemoveEach(32, 47);
- }
+ for (int i = 0; i < COUNT; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
- @Test
- public void testRemoveEntriesSpanningThreeBuckets() {
- doTestRemoveEach(8, 39);
- }
+ for (int i = 0; i < COUNT; i++) {
+ DeliveryType removed = tracker.remove(i);
+ assertNotNull(removed);
+ assertEquals(i, removed.getDeliveryId());
+ }
- @Test
- public void testRemoveAllEntriesWithClosedRange() {
- doTestRemoveEach(0, 47);
+ assertTrue(tracker.isEmpty());
}
@Test
- public void testRemoveAllEntriesWithOpenRange() {
- doTestRemoveEach(0, 64);
- }
-
- public void doTestRemoveEach(int start, int end) {
- final int numBuckets = 3;
- final int bucketSize = 16;
- final int numEntries = numBuckets * bucketSize;
- final int numRemoved = Math.min(end - start + 1, numEntries);
-
- UnsettledMap<DeliveryType> map = createMap(numBuckets, bucketSize);
+ public void testRemoveOldestAcrossBucketBoundaries() {
+ UnsettledMap<DeliveryType> tracker = createMap(3, 4);
- for (int i = 0; i < numEntries; ++i) {
- map.put(i, new DeliveryType(i));
+ for (int i = 0; i < 12; i++) {
+ tracker.put(i, new DeliveryType(i));
}
- assertEquals(numEntries, map.size());
-
- final AtomicInteger removed = new AtomicInteger();
-
- map.removeEach(start, end, (delivery) -> removed.incrementAndGet());
+ for (int i = 0; i < 12; i++) {
+ assertEquals(i, tracker.remove(i).getDeliveryId());
+ }
- assertEquals(numRemoved, removed.get());
- assertEquals(numEntries - numRemoved, map.size());
+ assertEquals(0, tracker.size());
}
@Test
- public void testRemoveEachWithRangeMuchLargerThanContainedEntries() {
- final UnsettledMap<DeliveryType> map = createMap();
- final int NUM_ENTRIES = 100;
+ public void testSlidingWindowPutRemoveOldest() {
+ UnsettledMap<DeliveryType> tracker = createMap();
- map.put(0, new DeliveryType(0));
+ int window = 1024;
- for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
- map.put(j, new DeliveryType(j));
+ for (int i = 0; i < window; i++) {
+ tracker.put(i, new DeliveryType(i));
}
- map.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
-
- assertEquals(NUM_ENTRIES + 2, map.size());
-
- final AtomicInteger removed = new AtomicInteger();
-
- map.removeEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
removed.incrementAndGet());
+ for (int i = 0; i < 5000; i++) {
+ tracker.put(window + i, new DeliveryType(window + i));
+ tracker.remove(i);
+ }
- assertEquals(NUM_ENTRIES + 2, removed.get());
- assertEquals(0, map.size());
+ assertEquals(window, tracker.size());
}
@Test
- public void testRemoveEachWhereLastNotPresentAndNextValuesAreOverflow() {
- final UnsettledMap<DeliveryType> map = createMap();
- final int NUM_ENTRIES = 100;
-
- map.put(0, new DeliveryType(0));
+ public void testMixedTailAndMiddleRemovals() {
+ UnsettledMap<DeliveryType> tracker = createMap();
- for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
+ for (int i = 0; i < 1000; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ tracker.remove(0); // tail
+ tracker.remove(500); // middle
+ tracker.remove(1); // tail again
+
+ assertFalse(tracker.containsKey(0));
+ assertFalse(tracker.containsKey(1));
+ assertFalse(tracker.containsKey(500));
+ assertNotNull(tracker.get(499));
+ }
+
+ @Test
+ public void testDuplicateIdsAfterWrap() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+
+ int nearMax = Integer.MAX_VALUE - 50;
+
+ for (int i = 0; i < 100; i++) {
+ tracker.put(nearMax + i, new DeliveryType(nearMax + i));
+ }
+
+ // Wrap
+ for (int i = 0; i < 100; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ // Validate both ranges exist
+ assertNotNull(tracker.get(nearMax + 10));
+ assertNotNull(tracker.get(10));
+ }
+
+ @Test
+ public void testRemoveDuplicateIdRemovesCorrectInstance() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+
+ tracker.put(1, new DeliveryType(1)); // old
+ // simulate wrap
+ tracker.put(1, new DeliveryType(1)); // new
+
+ tracker.remove(1);
+
+ // One should still remain
+ assertTrue(tracker.containsKey(1));
+ }
+
+ @Test
+ public void testIteratorRemoveAllSequentially() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+
+ for (int i = 0; i < 1000; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ Iterator<?> it = tracker.values().iterator();
+
+ while (it.hasNext()) {
+ it.next();
+ it.remove();
+ }
+
+ assertTrue(tracker.isEmpty());
+ }
+
+ @Test
+ public void testCompactionAtLowWaterMarkBoundary() {
+ UnsettledMap<DeliveryType> tracker = createMap(2, 16);
+
+ for (int i = 0; i < 32; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ for (int i = 0; i < 12; i++) {
+ tracker.remove(i);
+ }
+
+ // Now near low water mark
+ tracker.remove(12);
+
+ // Validate still consistent
+ assertNotNull(tracker.get(20));
+ }
+
+ @Test
+ public void testIteratorRemoveTriggeringCompactionFromHeadBucket() {
+ UnsettledMap<DeliveryType> tracker = createMap(3, 10); // 30 Entries
of capacity
+
+ // Fill the map to capacity but no further should remain at three
buckets
+ for (int i = 0; i < 30; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ Iterator<DeliveryType> it = tracker.values().iterator();
+
+ // Move to the middle before we start to remove entries
+ for (int i = 0; i < 15; i++) {
+ it.next();
+ }
+
+ DeliveryType lastReturned = null;
+
+ // Remove ten elements which should trigger a compaction and a roll
+ // forward of data in the center bucket.
+ for (int i = 0; i < 10; i++) {
+ it.remove();
+ lastReturned = it.next();
+ }
+
+ for (int i = 24; i < 30; i++) {
+ assertEquals(lastReturned.getDeliveryId(), i);
+ if (i < 29) {
+ lastReturned = it.next(); // Don't pass the end.
+ }
+ }
+ }
+
+ @Test
+ public void testIteratorRemoveTriggeringCompactionToHeadBucket() {
+ UnsettledMap<DeliveryType> tracker = createMap(3, 10); // 30 Entries
of capacity
+
+ // Fill the map to capacity but no further should remain at three
buckets
+ for (int i = 0; i < 30; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ Iterator<UnsignedInteger> it = tracker.keySet().iterator();
+
+ // Move to the middle before we start to remove entries
+ for (int i = 0; i < 25; i++) {
+ it.next();
+ }
+
+ // Should now be positioned in the head bucket and we remove the last
five entries
+ for (int i = 0; i < 5; i++) {
+ assertEquals(25 + i, it.next().intValue());
+ it.remove();
+ }
+
+ Iterator<DeliveryType> values = tracker.values().iterator();
+
+ for (int i = 0; i < 10; i++) {
+ values.next();
+ }
+
+ for (int i = 0; i < 5; i++) {
+ values.next();
+ values.remove();
+ }
+ }
+
+ @Test
+ public void testIteratorRemoveWhenHeadBucketShouldHaveWrappedAround() {
+ UnsettledMap<DeliveryType> tracker = createMap(3, 10); // 30 Entries
of capacity
+
+ // Fill the map to capacity but no further should remain at three
buckets
+ for (int i = 0; i < 30; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ // Remove the first 20 so head and tail are now the same.
+ for (int i = 0; i < 20; i++) {
+ tracker.remove(i);
+ }
+
+ // Put another 20 in so that head wraps to the slot behind tail
+ for (int i = 30; i < 50; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(30, tracker.size());
+
+ // Remove five from head so that when we remove five from the middle
bucket it rolls into head
+ tracker.remove(40);
+ tracker.remove(41);
+ tracker.remove(42);
+ tracker.remove(43);
+ tracker.remove(44);
+
+ Iterator<UnsignedInteger> it = tracker.keySet().iterator();
+
+ for (int i = 0; i < 15; ++i) {
+ assertEquals(20 + i, it.next().intValue());
+ }
+
+ for (int i = 0; i < 5; ++i) {
+ assertEquals(35 + i, it.next().intValue());
+ it.remove();
+ }
+
+ for (int i = 0; i < 5; ++i) {
+ assertEquals(45 + i, it.next().intValue());
+ }
+ }
+
+ @Test
+ public void testIteratorRemoveBeyondDefaultBucketSizeDoesNotThrow() {
+ final int bucketSize = 512;
+ final UnsettledMap<DeliveryType> tracker = createMap(2, bucketSize);
// uses test helper
+
+ for (int i = 0; i < bucketSize; ++i) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ final Iterator<DeliveryType> it = tracker.values().iterator();
+
+ DeliveryType removed = null;
+ for (int i = 0; i <= 300; ++i) {
+ removed = it.next();
+ }
+
+ try {
+ it.remove();
+ } catch (IndexOutOfBoundsException ex) {
+ fail("Iterator remove threw IndexOutOfBoundsException for
bucketSize=" + bucketSize +
+ " at index > 256 which is beyond the default bucket capacity
value. ");
+ }
+
+ assertNotNull(removed);
+ assertNull(tracker.get(removed.getDeliveryId()));
+ assertEquals(bucketSize - 1, tracker.size());
+ }
+
+ @Test
+ public void
testRemoveFromMiddleDoesNotLoseTailEntriesWhenCompactionTriggered() {
+ // Use small bucketSize so bucketLowWaterMark is small and compaction
is more likely.
+ final UnsettledMap<DeliveryType> tracker = createMap(6, 10); //
low-water ≈ 3
+
+ // Fill enough to create multiple buckets.
+ for (int i = 0; i < 60; ++i) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ // Make tail bucket sparse (but not empty) so tryCompact(tail) has a
chance to do something.
+ for (int i = 0; i < 8; ++i) {
+ assertNotNull(tracker.remove(i));
+ }
+
+ // Now remove from a middle region (not tail) and ensure the early
tail-adjacent IDs remain.
+ // If removeValue incorrectly compacts tail and corrupts the ring,
these can disappear.
+ assertNotNull(tracker.remove(25)); // removal from a non-tail bucket
+
+ // Sanity: nearby entries should still exist
+ assertNotNull(tracker.get(24));
+ assertNull(tracker.get(25));
+ assertNotNull(tracker.get(26));
+
+ // Tail-adjacent entries (8..15) should still exist
+ for (int i = 8; i < 16; ++i) {
+ assertNotNull(tracker.get(i), "Entry " + i + " missing after
middle removal/compaction");
+ }
+ }
+
+ @Test
+ public void
testRemoveFromNonTailTriggersWrongCompactionAndStillPreservesCorrectness() {
+ final UnsettledMap<DeliveryType> tracker = createMap(6, 10);
+
+ // Fill 3 buckets: 0..29
+ for (int i = 0; i < 30; ++i) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ // Make tail bucket small (remove 0..6 leaves 7..9 in first bucket =>
3 entries == low-water)
+ for (int i = 0; i < 7; ++i) {
+ assertNotNull(tracker.remove(i));
+ }
+
+ // Now remove entries from a *non-tail* bucket to make that bucket
sparse too.
+ // Removing these should make the target bucket hit <= low-water and
trigger the compaction path.
+ assertNotNull(tracker.remove(15));
+ assertNotNull(tracker.remove(16));
+ assertNotNull(tracker.remove(17));
+
+ // If tryCompact(tail) corrupts the map, these will be missing or
inconsistent.
+ for (int i = 7; i < 10; ++i) {
+ assertNotNull(tracker.get(i), "Tail-adjacent entry missing after
non-tail removal triggered compaction");
+ }
+ assertNull(tracker.get(15));
+ assertNull(tracker.get(16));
+ assertNull(tracker.get(17));
+ assertNotNull(tracker.get(18));
+ }
+
+ @Test
+ public void testIterateOverBucketsThatHaveWrapped() {
+ final UnsettledMap<DeliveryType> tracker = createMap();
+
+ tracker.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
+ tracker.put(0, new DeliveryType(0));
+ tracker.put(1, new DeliveryType(1));
+
+ assertEquals(3, tracker.size());
+
+ Iterator<UnsignedInteger> iter = tracker.keySet().iterator();
+
+ final int[] expected = { UnsignedInteger.MAX_VALUE.intValue(), 0, 1 };
+
+ int count = 0;
+
+ while (iter.hasNext()) {
+ assertEquals(expected[count++], iter.next().intValue());
+ }
+
+ assertEquals(3, count);
+ }
+
+ protected void dumpRandomDataSet(int iterations, long seed, boolean
bounded) {
+ final int[] dataSet = new int[iterations];
+
+ random.setSeed(seed);
+
+ for (int i = 0; i < iterations; ++i) {
+ if (bounded) {
+ dataSet[i] = random.nextInt(iterations);
+ } else {
+ dataSet[i] = random.nextInt();
+ }
+ }
+
+ LOG.info("Iterations was {}, Random seed was: {}", iterations , seed);
+ LOG.info("Entries in data set: {}", dataSet);
+ }
+
+ protected static class OutsideEntry<K, V> implements Map.Entry<K, V> {
+
+ private final K key;
+ private V value;
+
+ public OutsideEntry(K key, V value) {
+ this.key = key;
+ this.value = value;
+ }
+
+ @Override
+ public V setValue(V value) {
+ V oldValue = this.value;
+ this.value = value;
+ return oldValue;
+ }
+
+ @Override
+ public V getValue() {
+ return value;
+ }
+
+ @Override
+ public K getKey() {
+ return key;
+ }
+ }
+
+ //----- forEach methods test variations
+
+ @Test
+ public void testForEachEntry() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+
+ final int[] inputValues = {3, 0, -1, 1, -2, 2};
+
+ for (int entry : inputValues) {
+ tracker.put(entry, new DeliveryType(entry));
+ }
+
+ final SequenceNumber index = new SequenceNumber(0);
+ tracker.forEach((value) -> {
+ int i = index.getAndIncrement().intValue();
+ assertEquals(new DeliveryType(inputValues[i]), value);
+ });
+
+ assertEquals(index.intValue(), inputValues.length);
+ }
+
+ @Test
+ public void testForEachDeliveryIteratesOverLargeSeriesOfDeliveries() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+ assertEquals(0, tracker.size());
+
+ final int COUNT = 4080;
+
+ for (int i = 0; i < COUNT; ++i) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(COUNT, tracker.size());
+
+ final AtomicInteger index = new AtomicInteger();
+
+ tracker.forEach((delivery) -> index.incrementAndGet());
+
+ assertEquals(index.get(), COUNT);
+ }
+
+ @Test
+ public void testForEachOnEmptyMap() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+ assertEquals(0, tracker.size());
+
+ final AtomicInteger index = new AtomicInteger();
+
+ tracker.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
index.incrementAndGet());
+
+ assertEquals(index.get(), 0);
+ }
+
+ @Test
+ public void
testForEachBiConsumerDeliveryIteratesOverLargeSeriesOfDeliveries() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+ assertEquals(0, tracker.size());
+
+ final int COUNT = 4080;
+
+ for (int i = 0; i < COUNT; ++i) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(COUNT, tracker.size());
+
+ final AtomicInteger index = new AtomicInteger();
+
+ tracker.forEach((deliveryId, delivery) -> index.incrementAndGet());
+
+ assertEquals(index.get(), COUNT);
+ }
+
+ @Test
+ public void testRangedForEachDeliveryIteratesOverSmallSeriesOfDeliveries()
{
+ UnsettledMap<DeliveryType> tracker = createMap();
+ assertEquals(0, tracker.size());
+
+ final int COUNT = 512;
+
+ for (int i = 0; i < COUNT; ++i) {
+ tracker.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(COUNT, tracker.size());
+
+ final AtomicInteger index = new AtomicInteger();
+
+ tracker.forEach(260, 262, (delivery) -> index.incrementAndGet());
+
+ assertEquals(index.get(), 3);
+ }
+
+ @Test
+ public void
testRangedForEachDeliveryIteratesSeriesWhenValuesOverflowIntRange() {
+ UnsettledMap<DeliveryType> tracker = createMap();
+ assertEquals(0, tracker.size());
+
+ tracker.put(0, new DeliveryType(0));
+ tracker.put(1, new DeliveryType(1));
+ tracker.put(Integer.MAX_VALUE, new DeliveryType(Integer.MAX_VALUE));
+ tracker.put(Integer.MAX_VALUE + 1, new DeliveryType(Integer.MAX_VALUE
+ 1));
+ tracker.put(Integer.MAX_VALUE + 2, new DeliveryType(Integer.MAX_VALUE
+ 2));
+ tracker.put(Integer.MAX_VALUE + 3, new DeliveryType(Integer.MAX_VALUE
+ 3));
+ tracker.put(Integer.MAX_VALUE + 4, new DeliveryType(Integer.MAX_VALUE
+ 4));
+
+ final AtomicInteger index = new AtomicInteger();
+
+ tracker.forEach(Integer.MAX_VALUE, Integer.MAX_VALUE + 2, (delivery)
-> index.incrementAndGet());
+
+ assertEquals(3, index.get());
+ }
+
+ @Test
+ public void testForEachWithRangeMuchLargerThanContainedEntries() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = 100;
+
+ map.put(0, new DeliveryType(0));
+
+ for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
+ map.put(j, new DeliveryType(j));
+ }
+
+ map.put(Integer.MAX_VALUE, new DeliveryType(Integer.MAX_VALUE));
+
+ assertEquals(NUM_ENTRIES + 2, map.size());
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(0, Integer.MAX_VALUE, (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(NUM_ENTRIES + 2, traversed.get());
+ assertEquals(NUM_ENTRIES + 2, map.size());
+ }
+
+ @Test
+ public void testForEachCoversElementsInBetweenGivenRangeInOtherBuckets() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
+ final int NUM_ENTRIES = 100;
+
+ for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
+ map.put(j, new DeliveryType(j));
+ }
+
+ assertEquals(NUM_ENTRIES, map.size());
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(0, Integer.MAX_VALUE, (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(NUM_ENTRIES, traversed.get());
+ assertEquals(NUM_ENTRIES, map.size());
+ }
+
+ @Test
+ public void testForEachWhereLastValueNotPresentButGreaterValuesAre() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = 100;
+
+ map.put(0, new DeliveryType(0));
+
+ for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
+ map.put(j, new DeliveryType(j));
+ }
+
+ map.put(65534, new DeliveryType(65534));
+ map.put(65536, new DeliveryType(65536));
+ map.put(65537, new DeliveryType(65537));
+
+ map.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
+
+ assertEquals(NUM_ENTRIES + 5, map.size());
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(0, 65535, (delivery) -> traversed.incrementAndGet());
+
+ assertEquals(NUM_ENTRIES + 2, traversed.get());
+ }
+
+ @Test
+ public void testForEachWhereLastNotPresentAndNextValuesAreOverflow() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = 100;
+ final int EXPECTED_ENTRIES = 102;
+
+ map.put(0, new DeliveryType(0));
+
+ for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
map.put(j, new DeliveryType(j));
}
@@ -2081,147 +2502,231 @@ public class UnsettledMapTest {
assertEquals(NUM_ENTRIES + 5, map.size());
- final AtomicInteger removed = new AtomicInteger();
+ final AtomicInteger traversed = new AtomicInteger();
- map.removeEach(0, 65535, (delivery) -> removed.incrementAndGet());
+ map.forEach(0, 65535, (delivery) -> traversed.incrementAndGet());
- assertEquals(NUM_ENTRIES + 5, removed.get());
- assertEquals(0, map.size());
+ assertEquals(EXPECTED_ENTRIES, traversed.get());
}
@Test
- public void testRemoveEachEntireDeliveryIdRangeTwoBuckets() {
- final UnsettledMap<DeliveryType> map = createMap(3, 128);
- final int NUM_ENTRIES = UnsignedByte.MAX_VALUE.intValue();
+ public void testForEachWithRangeThatWrapped() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- for (int i = 0; i < NUM_ENTRIES; ++i) {
+ final int[] entries = new int[] { 512, 513, 0, 1, 2, 3,
Integer.MAX_VALUE };
+
+ for(int i : entries) {
map.put(i, new DeliveryType(i));
}
- final AtomicInteger removed = new AtomicInteger();
+ assertEquals(entries.length, map.size());
- map.removeEach(0, NUM_ENTRIES, (delivery) ->
removed.incrementAndGet());
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(0, Integer.MAX_VALUE, (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(2, traversed.get());
+ assertEquals(entries.length, map.size());
+
+ map.remove(512);
+ map.remove(513);
+
+ map.forEach(0, Integer.MAX_VALUE, (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(7, traversed.get());
+ }
+
+ @Test
+ public void testForEachWithRangeThatWrappedAndStartBeyondMaxInt() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
+
+ final int[] entries = new int[] { 512, 513, Integer.MAX_VALUE + 1,
UnsignedInteger.MAX_VALUE.intValue() };
+
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(entries.length, map.size());
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(Integer.MAX_VALUE + 1,
UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(2, traversed.get());
+ assertEquals(entries.length, map.size());
+ }
+
+ @Test
+ public void testForEachFindsNoValues() {
+ final int afterLast = uintArray.length + 1; // Entries are one based
+
+ final AtomicBoolean traversed = new AtomicBoolean();
+
+ tracker.forEach(afterLast, afterLast + 10, (delivery) ->
traversed.set(true));
+
+ assertFalse(traversed.get());
+ }
+
+ @Test
+ public void testForEachRangeFindsOnlyLastValue() {
+ final int lastEntry = uintArray.length; // Entries are one based
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ tracker.forEach(lastEntry, lastEntry, (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(1, traversed.get());
+ }
+
+ @Test
+ public void testForEachRangedLastValueAndRangeOutsideOfActualEntries() {
+ final int lastEntry = uintArray.length; // Entries are one based
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ tracker.forEach(lastEntry, lastEntry + 10, (delivery) ->
traversed.incrementAndGet());
+
+ assertEquals(1, traversed.get());
+ }
+
+ @Test
+ public void testForEachWithRangeThatMatchesMany() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
+
+ final int[] entries = new int[] { 512, 513, 512, 513, 512, 513 };
+
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(entries.length, map.size());
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(512, 513, (delivery) -> traversed.incrementAndGet());
- assertEquals(NUM_ENTRIES, removed.get());
- assertEquals(0, map.size());
+ assertEquals(2, traversed.get());
+ assertEquals(entries.length, map.size());
}
@Test
- public void testRemoveEachEntireDeliveryIdRange() {
- final UnsettledMap<DeliveryType> map = createMap();
- final int NUM_ENTRIES = UnsignedShort.MAX_VALUE.intValue();
+ public void testForEachWithRangeThatFallBetweenManyGenerations() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- for (int i = 0; i < NUM_ENTRIES; ++i) {
+ final int[] entries = new int[] { 512, 513, 512, 513, 512, 513 };
+
+ for(int i : entries) {
map.put(i, new DeliveryType(i));
}
- final AtomicInteger removed = new AtomicInteger();
+ assertEquals(entries.length, map.size());
- map.removeEach(0, NUM_ENTRIES, (delivery) ->
removed.incrementAndGet());
+ final AtomicInteger traversed = new AtomicInteger();
- assertEquals(NUM_ENTRIES, removed.get());
- assertEquals(0, map.size());
+ map.forEach(1, 10, (delivery) -> traversed.incrementAndGet());
+ map.forEach(600, 65535, (delivery) -> traversed.incrementAndGet());
+ map.forEach(510, 511, (delivery) -> traversed.incrementAndGet());
+
+ assertEquals(0, traversed.get());
+ assertEquals(entries.length, map.size());
}
@Test
- public void testForEachWithRangeMuchLargerThanContainedEntries() {
- final UnsettledMap<DeliveryType> map = createMap();
- final int NUM_ENTRIES = 100;
+ public void testForEachWithRangeThatFallAtEndOfLastGeneration() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- map.put(0, new DeliveryType(0));
+ final int[] entries = new int[] { 512, 513, 512, 513, 512, 513, 65530,
65531, 65535 };
- for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
- map.put(j, new DeliveryType(j));
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
}
- map.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
-
- assertEquals(NUM_ENTRIES + 2, map.size());
+ assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> returned = new ArrayList<>();
- map.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
traversed.incrementAndGet());
+ map.forEach(65531, 65580, (delivery) -> {
+ traversed.incrementAndGet();
+ returned.add(delivery.getDeliveryId());
+ });
- assertEquals(NUM_ENTRIES + 2, traversed.get());
- assertEquals(NUM_ENTRIES + 2, map.size());
+ assertEquals(2, traversed.get());
+ assertEquals(entries.length, map.size());
+
+ assertTrue(returned.contains(65531));
+ assertTrue(returned.contains(65535));
}
@Test
- public void testForEachCoversElementsInBetweenGivenRangeInOtherBuckets() {
- final UnsettledMap<DeliveryType> map = createMap(5, 10);
- final int NUM_ENTRIES = 100;
+ public void testForEachBoundedToFirstGeneration() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
- map.put(j, new DeliveryType(j));
+ final int[] entries = new int[] { 512, 513,
UnsignedInteger.MAX_VALUE.intValue(), UnsignedInteger.MAX_VALUE.intValue() };
+
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
}
- assertEquals(NUM_ENTRIES, map.size());
+ assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
- map.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
traversed.incrementAndGet());
+ map.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) -> {
+ assertEquals(entries[traversed.getAndIncrement()],
delivery.getDeliveryId());
+ });
- assertEquals(NUM_ENTRIES, traversed.get());
- assertEquals(NUM_ENTRIES, map.size());
+ assertEquals(3, traversed.get());
+ assertEquals(entries.length, map.size());
}
@Test
- public void testForEachWhereLastValueNotPresentButGreaterValuesAre() {
- final UnsettledMap<DeliveryType> map = createMap();
- final int NUM_ENTRIES = 100;
+ public void testForEachWhenSpanHasValuePresentInSuccessiveGeneration() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- map.put(0, new DeliveryType(0));
+ final int[] entries = new int[] { 0, 1, 65535, Integer.MAX_VALUE,
UnsignedInteger.MAX_VALUE.intValue(), UnsignedInteger.MAX_VALUE.intValue() };
- for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
- map.put(j, new DeliveryType(j));
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
}
- map.put(65534, new DeliveryType(65534));
- map.put(65536, new DeliveryType(65536));
- map.put(65537, new DeliveryType(65537));
-
- map.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
-
- assertEquals(NUM_ENTRIES + 5, map.size());
+ assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
- map.forEach(0, 65535, (delivery) -> traversed.incrementAndGet());
+ map.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) -> {
+ assertEquals(entries[traversed.getAndIncrement()],
delivery.getDeliveryId());
+ });
- assertEquals(NUM_ENTRIES + 2, traversed.get());
+ assertEquals(5, traversed.get());
+ assertEquals(entries.length, map.size());
}
@Test
- public void testForEachWhereLastNotPresentAndNextValuesAreOverflow() {
- final UnsettledMap<DeliveryType> map = createMap();
- final int NUM_ENTRIES = 100;
+ public void testForEachWhereFirstAndLastAreEqualOnlyReturnsOnce() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- map.put(0, new DeliveryType(0));
+ final int[] entries = new int[] {
UnsignedInteger.MAX_VALUE.intValue(), 0, 1,
UnsignedInteger.MAX_VALUE.intValue() };
- for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
- map.put(j, new DeliveryType(j));
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
}
- map.put(65534, new DeliveryType(65534));
-
- map.put(0, new DeliveryType(0));
- map.put(1, new DeliveryType(1));
- map.put(2, new DeliveryType(2));
-
- assertEquals(NUM_ENTRIES + 5, map.size());
+ assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
- map.forEach(0, 65535, (delivery) -> traversed.incrementAndGet());
+ map.forEach(UnsignedInteger.MAX_VALUE.intValue(),
UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
traversed.incrementAndGet());
- assertEquals(NUM_ENTRIES + 5, traversed.get());
+ assertEquals(1, traversed.get());
+ assertEquals(entries.length, map.size());
}
@Test
- public void testForEachWithRangeThatWrapped() {
+ public void testForeachWhenRangeCarriesOverflow() {
final UnsettledMap<DeliveryType> map = createMap(5, 50);
- final int[] entries = new int[] { 512, 513, 0, 1, 2, 3,
UnsignedInteger.MAX_VALUE.intValue() };
+ final int[] entries = new int[] { 65534, 65535, 10, 11, 15, 90 };
for(int i : entries) {
map.put(i, new DeliveryType(i));
@@ -2230,51 +2735,86 @@ public class UnsettledMapTest {
assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> returned = new ArrayList<>();
- map.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), (delivery) ->
traversed.incrementAndGet());
+ map.forEach(65534, 15, (delivery) -> {
+ traversed.incrementAndGet();
+ returned.add(delivery.getDeliveryId());
+ });
- assertEquals(entries.length, traversed.get());
+ assertEquals(5, traversed.get());
assertEquals(entries.length, map.size());
+
+ assertTrue(returned.contains(65534));
+ assertTrue(returned.contains(65535));
+ assertTrue(returned.contains(10));
+ assertTrue(returned.contains(11));
+ assertTrue(returned.contains(15));
}
@Test
- public void testForEachFindsNoValues() {
- final int afterLast = uintArray.length + 1; // Entries are one based
-
- final AtomicBoolean traversed = new AtomicBoolean();
+ public void
testForeachWhenRangeCarriesOverflowButStopsBeforeFirstValueInNextGeneration() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- tracker.forEach(afterLast, afterLast + 10, (delivery) ->
traversed.set(true));
+ final int[] entries = new int[] { 65534, 65535, 10, 11, 15, 90 };
- assertFalse(traversed.get());
- }
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
+ }
- @Test
- public void testForEachRangeFindsOnlyLastValue() {
- final int lastEntry = uintArray.length; // Entries are one based
+ assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> returned = new ArrayList<>();
- tracker.forEach(lastEntry, lastEntry, (delivery) ->
traversed.incrementAndGet());
+ map.forEach(65534, 9, (delivery) -> {
+ traversed.incrementAndGet();
+ returned.add(delivery.getDeliveryId());
+ });
- assertEquals(1, traversed.get());
+ assertEquals(2, traversed.get());
+ assertEquals(entries.length, map.size());
+
+ assertTrue(returned.contains(65534));
+ assertTrue(returned.contains(65535));
}
@Test
- public void testForEachRangedLastValueAndRangeOutsideOfActualEntries() {
- final int lastEntry = uintArray.length; // Entries are one based
+ public void
testForeachWhenRangeCarriesOverflowButStopsBeforeFirstValueInGenerationBeyondNextValue()
{
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
+
+ final int[] entries = new int[] { 65534, 65535, 10, 10, 11, 15, 90 };
+
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
+ }
+
+ assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> returned = new ArrayList<>();
- tracker.forEach(lastEntry, lastEntry + 10, (delivery) ->
traversed.incrementAndGet());
+ assertNotNull(map.remove(10));
- assertEquals(1, traversed.get());
+ map.forEach(65534, 11, (delivery) -> {
+ traversed.incrementAndGet();
+ returned.add(delivery.getDeliveryId());
+ });
+
+ assertEquals(4, traversed.get());
+ assertEquals(entries.length - 1, map.size());
+
+ assertTrue(returned.contains(65534));
+ assertTrue(returned.contains(65535));
+ assertTrue(returned.contains(10));
+ assertTrue(returned.contains(11));
}
@Test
- public void testForEachWithRangeThatMatchesMany() {
+ public void testForeachWhenRangeCarriesOverflowNoOverflowValuesInMap() {
final UnsettledMap<DeliveryType> map = createMap(5, 50);
- final int[] entries = new int[] { 512, 513, 512, 513, 512, 513 };
+ final int[] entries = new int[] { 65534, 65535 };
for(int i : entries) {
map.put(i, new DeliveryType(i));
@@ -2283,444 +2823,452 @@ public class UnsettledMapTest {
assertEquals(entries.length, map.size());
final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> returned = new ArrayList<>();
- map.forEach(512, 513, (delivery) -> traversed.incrementAndGet());
+ map.forEach(65534, 255, (delivery) -> {
+ traversed.incrementAndGet();
+ returned.add(delivery.getDeliveryId());
+ });
assertEquals(2, traversed.get());
assertEquals(entries.length, map.size());
+
+ assertTrue(returned.contains(65534));
+ assertTrue(returned.contains(65535));
}
@Test
- public void testRemoveAllEntriesInSmallChunks() {
- final AtomicInteger removed = new AtomicInteger();
+ public void testForEachInBucketThatHasNoMatchingSequence() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- for (int i = 0; i < uintArray.length; i += 2) {
- tracker.removeEach(uintArray[i].intValue(),
uintArray[i+1].intValue(), (delivery) -> removed.incrementAndGet());
+ final int[] entries = new int[] { 65534, 65535, 65550, 65551 };
+
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
}
- assertEquals(uintArray.length, removed.get());
+ assertEquals(entries.length, map.size());
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ map.forEach(65538, 65549, (delivery) -> {
+ traversed.incrementAndGet();
+ });
+
+ assertEquals(0, traversed.get());
+ assertEquals(entries.length, map.size());
}
@Test
- public void testRemoveFromIteratorFromMiddleBucket() {
- final int numBuckets = 5;
- final int bucketSize = 10;
- final int numEntries = numBuckets * bucketSize;
+ public void testForEachWithRangeThatWrappedAndNonZeroReadOffset() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- final UnsettledMap<DeliveryType> map = createMap(numBuckets,
bucketSize);
+ // Populate items in high range and lower wrapped range
+ final int[] entries = new int[] { -4, -3, -2, -1, 0, 1, 2, 3, 4, 5 };
- for (int i = 0; i < numEntries; ++i) {
+ for (int i : entries) {
map.put(i, new DeliveryType(i));
}
- assertEquals(numEntries, map.size());
-
- Iterator<UnsignedInteger> entries = map.keySet().iterator();
+ assertEquals(entries.length, map.size());
- // Move to center of bucket two
- for (int i = 0; i < bucketSize + (bucketSize / 2); ++i) {
- entries.next();
- }
+ final AtomicInteger traversedCount = new AtomicInteger();
+ final List<Integer> traversed = new ArrayList<>();
- UnsignedInteger lastValue = null;
+ // Starts mid-bucket in the upper segment and ends mid-bucket in the
lower segment.
+ map.forEach(-2, 2, (delivery) -> {
+ traversedCount.incrementAndGet();
+ traversed.add(delivery.getDeliveryId());
+ });
- // Remove from center of bucket two into bucket three until a
compaction event should occur.
- for (int i = 0; i < bucketSize; ++i) {
- lastValue = entries.next();
- entries.remove();
- }
+ assertEquals(5, traversedCount.get());
+ assertEquals(entries.length, map.size());
- assertEquals(lastValue.intValue() + 1, entries.next().intValue());
+ assertTrue(traversed.contains(-2));
+ assertTrue(traversed.contains(-1));
+ assertTrue(traversed.contains(0));
+ assertTrue(traversed.contains(1));
+ assertTrue(traversed.contains(2));
}
@Test
- public void testRemoveFromMiddleBucket() {
- final int numBuckets = 5;
- final int bucketSize = 10;
- final int numEntries = numBuckets * bucketSize;
+ public void testForEachAcrossMultipleGenerations() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- final UnsettledMap<DeliveryType> map = createMap(numBuckets,
bucketSize);
+ // Add sequence that spans two full generations of unsigned integer
wrap
+ map.put(-2, new DeliveryType(-2)); // Gen 0
+ map.put(-1, new DeliveryType(-1)); // Gen 0
+ map.put(0, new DeliveryType(0)); // Gen 1 (overflow)
+ map.put(1, new DeliveryType(1)); // Gen 1
+ map.put(2, new DeliveryType(2)); // Gen 1
- for (int i = 0; i < numEntries; ++i) {
- map.put(i, new DeliveryType(i));
- }
+ final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> visited = new ArrayList<>();
- assertEquals(numEntries, map.size());
+ // Iterate over wrapped range
+ map.forEach(-2, 1, (delivery) -> {
+ traversed.incrementAndGet();
+ visited.add(delivery.getDeliveryId());
+ });
- int position = bucketSize + (bucketSize / 2);
- int lastValue = 0;
+ assertEquals(4, traversed.get());
+ assertEquals(Arrays.asList(-2, -1, 0, 1), visited);
+ }
- // Remove from center of bucket two into bucket three until a
compaction event should occur.
- for (int i = 0; i < bucketSize + 3; ++i) {
- lastValue = map.get(position).getDeliveryId();
- map.remove(position++);
- }
+ @Test
+ public void
testForEachRangedContinuationShortCircuitsWhenBucketLowExceedsLast() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- assertEquals(lastValue + 1, map.get(position).getDeliveryId());
+ // Gen 0: High IDs
+ map.put(-2, new DeliveryType(-2)); // 0xFFFFFFFE
+ map.put(-1, new DeliveryType(-1)); // 0xFFFFFFFF
+
+ // Gen 1: Low IDs that start well past the second bound 'last' (1)
+ map.put(10, new DeliveryType(10));
+ map.put(11, new DeliveryType(11));
+
+ final AtomicInteger traversed = new AtomicInteger();
+
+ // Search range -2 to 1 (wrapped). Gen 0 has -2, -1. Gen 1 starts at
10 (> 1).
+ map.forEach(-2, 1, d -> traversed.incrementAndGet());
+
+ assertEquals(2, traversed.get());
}
@Test
- public void testRepeatedRemoveOldestHotPath() {
- UnsettledMap<DeliveryType> tracker = createMap();
+ public void testForEachRangedContinuationBucketsExceedRangeOfLast() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- final int COUNT = 10000;
+ // Gen 0
+ map.put(-2, new DeliveryType(-2));
+ map.put(-1, new DeliveryType(-1));
- for (int i = 0; i < COUNT; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ // Gen 1
+ map.put(10, new DeliveryType(10));
+ map.put(11, new DeliveryType(11));
- for (int i = 0; i < COUNT; i++) {
- DeliveryType removed = tracker.remove(i);
- assertNotNull(removed);
- assertEquals(i, removed.getDeliveryId());
- }
+ final AtomicInteger traversed = new AtomicInteger();
+ final List<Integer> visited = new ArrayList<>();
- assertTrue(tracker.isEmpty());
+ map.forEach(-2, 1, (delivery) -> {
+ traversed.incrementAndGet();
+ visited.add(delivery.getDeliveryId());
+ });
+
+ assertEquals(2, traversed.get());
+ assertEquals(4, map.size());
+
+ assertTrue(visited.contains(-2));
+ assertTrue(visited.contains(-1));
}
+ //----- removeEach method test variations
+
@Test
- public void testRemoveOldestAcrossBucketBoundaries() {
- UnsettledMap<DeliveryType> tracker = createMap(3, 4);
+ public void testRemoveRangeRemovesNoValues() {
+ final int afterLast = uintArray.length + 1; // Entries are one based
- for (int i = 0; i < 12; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ final AtomicBoolean removed = new AtomicBoolean();
- for (int i = 0; i < 12; i++) {
- assertEquals(i, tracker.remove(i).getDeliveryId());
- }
+ tracker.removeEach(afterLast, afterLast + 10, (delivery) ->
removed.set(true));
- assertEquals(0, tracker.size());
+ assertFalse(removed.get());
}
@Test
- public void testSlidingWindowPutRemoveOldest() {
- UnsettledMap<DeliveryType> tracker = createMap();
-
- int window = 1024;
+ public void testRemoveRangeRemovesLastValue() {
+ final int lastEntry = uintArray.length; // Entries are one based
- for (int i = 0; i < window; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ final AtomicInteger removed = new AtomicInteger();
- for (int i = 0; i < 5000; i++) {
- tracker.put(window + i, new DeliveryType(window + i));
- tracker.remove(i);
- }
+ tracker.removeEach(lastEntry, lastEntry, (delivery) ->
removed.incrementAndGet());
- assertEquals(window, tracker.size());
+ assertEquals(1, removed.get());
}
@Test
- public void testMixedTailAndMiddleRemovals() {
- UnsettledMap<DeliveryType> tracker = createMap();
+ public void
testRemoveRangeRemovesLastValueAndRangeOutsideOfActualEntries() {
+ final int lastEntry = uintArray.length; // Entries are one based
- for (int i = 0; i < 1000; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ final AtomicInteger removed = new AtomicInteger();
- tracker.remove(0); // tail
- tracker.remove(500); // middle
- tracker.remove(1); // tail again
+ tracker.removeEach(lastEntry, lastEntry + 10, (delivery) ->
removed.incrementAndGet());
- assertFalse(tracker.containsKey(0));
- assertFalse(tracker.containsKey(1));
- assertFalse(tracker.containsKey(500));
- assertNotNull(tracker.get(499));
+ assertEquals(1, removed.get());
}
@Test
- public void testDuplicateIdsAfterWrap() {
- UnsettledMap<DeliveryType> tracker = createMap();
+ public void testRemoveEachWithRangeThatMatchesMany() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 50);
- int nearMax = Integer.MAX_VALUE - 50;
+ final int[] entries = new int[] { 512, 513, 512, 513, 512, 513 };
- for (int i = 0; i < 100; i++) {
- tracker.put(nearMax + i, new DeliveryType(nearMax + i));
+ for(int i : entries) {
+ map.put(i, new DeliveryType(i));
}
- // Wrap
- for (int i = 0; i < 100; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ assertEquals(entries.length, map.size());
- // Validate both ranges exist
- assertNotNull(tracker.get(nearMax + 10));
- assertNotNull(tracker.get(10));
+ final AtomicInteger removed = new AtomicInteger();
+
+ map.removeEach(512, 513, (delivery) -> removed.incrementAndGet());
+
+ assertEquals(2, removed.get());
+ assertEquals(entries.length - 2, map.size());
}
@Test
- public void testRemoveDuplicateIdRemovesCorrectInstance() {
- UnsettledMap<DeliveryType> tracker = createMap();
+ public void testRemoveAllEntriesFromFirstBucket() {
+ doTestRemoveEach(0, 15);
+ }
- tracker.put(1, new DeliveryType(1)); // old
- // simulate wrap
- tracker.put(1, new DeliveryType(1)); // new
+ @Test
+ public void testRemoveAllEntriesFromMiddleBucket() {
+ doTestRemoveEach(16, 31);
+ }
- tracker.remove(1);
+ @Test
+ public void testRemoveAllEntriesFromEndBucket() {
+ doTestRemoveEach(32, 47);
+ }
- // One should still remain
- assertTrue(tracker.containsKey(1));
+ @Test
+ public void testRemoveEntriesSpanningThreeBuckets() {
+ doTestRemoveEach(8, 39);
}
@Test
- public void testRemoveEachWithNonZeroReadOffset() {
- UnsettledMap<DeliveryType> tracker = createMap(2, 16);
+ public void testRemoveAllEntriesWithClosedRange() {
+ doTestRemoveEach(0, 47);
+ }
- for (int i = 0; i < 32; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ @Test
+ public void testRemoveAllEntriesWithOpenRange() {
+ doTestRemoveEach(0, 64);
+ }
- for (int i = 0; i < 5; i++) {
- tracker.remove(i);
- }
+ public void doTestRemoveEach(int start, int end) {
+ final int numBuckets = 3;
+ final int bucketSize = 16;
+ final int numEntries = numBuckets * bucketSize;
+ final int numRemoved = Math.min(end - start + 1, numEntries);
- tracker.removeEach(5, 20, d -> {});
+ UnsettledMap<DeliveryType> map = createMap(numBuckets, bucketSize);
- for (int i = 5; i <= 20; i++) {
- assertFalse(tracker.containsKey(i));
+ for (int i = 0; i < numEntries; ++i) {
+ map.put(i, new DeliveryType(i));
}
- }
-
- @Test
- public void testIteratorRemoveAllSequentially() {
- UnsettledMap<DeliveryType> tracker = createMap();
- for (int i = 0; i < 1000; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ assertEquals(numEntries, map.size());
- Iterator<?> it = tracker.values().iterator();
+ final AtomicInteger removed = new AtomicInteger();
- while (it.hasNext()) {
- it.next();
- it.remove();
- }
+ map.removeEach(start, end, (delivery) -> removed.incrementAndGet());
- assertTrue(tracker.isEmpty());
+ assertEquals(numRemoved, removed.get());
+ assertEquals(numEntries - numRemoved, map.size());
}
@Test
- public void testCompactionAtLowWaterMarkBoundary() {
- UnsettledMap<DeliveryType> tracker = createMap(2, 16);
+ public void testRemoveEachWithRangeMuchLargerThanContainedEntries() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = 100;
- for (int i = 0; i < 32; i++) {
- tracker.put(i, new DeliveryType(i));
- }
+ map.put(0, new DeliveryType(0));
- for (int i = 0; i < 12; i++) {
- tracker.remove(i);
+ for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
+ map.put(j, new DeliveryType(j));
}
- // Now near low water mark
- tracker.remove(12);
+ map.put(Integer.MAX_VALUE, new DeliveryType(Integer.MAX_VALUE));
- // Validate still consistent
- assertNotNull(tracker.get(20));
+ assertEquals(NUM_ENTRIES + 2, map.size());
+
+ final AtomicInteger removed = new AtomicInteger();
+
+ map.removeEach(0, Integer.MAX_VALUE, (delivery) ->
removed.incrementAndGet());
+
+ assertEquals(NUM_ENTRIES + 2, removed.get());
+ assertEquals(0, map.size());
}
@Test
- public void testRemoveEachClearsThenReuseMap() {
- UnsettledMap<DeliveryType> tracker = createMap();
+ public void
testRemoveEachWithRangeMuchLargerThanContainedEntriesInRangeAboveMaxInt() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = 100;
- for (int i = 0; i < 100; i++) {
- tracker.put(i, new DeliveryType(i));
+ map.put(Integer.MAX_VALUE, new DeliveryType(Integer.MAX_VALUE));
+
+ for (int i = 0, j = Integer.MAX_VALUE + 512; i < NUM_ENTRIES; ++i, j
+= 25) {
+ map.put(j, new DeliveryType(j));
}
- tracker.removeEach(0, 200, d -> {});
+ map.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
- assertTrue(tracker.isEmpty());
+ assertEquals(NUM_ENTRIES + 2, map.size());
- tracker.put(999, new DeliveryType(999));
+ final AtomicInteger removed = new AtomicInteger();
- assertEquals(1, tracker.size());
- assertNotNull(tracker.get(999));
+ map.removeEach(Integer.MAX_VALUE, UnsignedInteger.MAX_VALUE.intValue()
- 1, (delivery) -> removed.incrementAndGet());
+
+ assertEquals(NUM_ENTRIES + 1, removed.get());
+ assertEquals(1, map.size());
}
@Test
- public void testRecycleBucketWhenTailHasWrapped() {
- final UnsettledMap<DeliveryType> tracker = createMap(5, 2);
+ public void testRemoveEachWhereLastNotPresentAndNextValuesAreOverflow() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = 100;
- // Fill all 5 buckets: bucket0=[0,1], bucket1=[2,3], bucket2=[4,5],
bucket3=[6,7], bucket4=[8,9]
- for (int i = 0; i < 10; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
- assertEquals(10, tracker.size());
+ map.put(0, new DeliveryType(0));
- // Remove 0..5 => drains bucket0, bucket1, bucket2 fully (advances
tail to bucket3)
- for (int i = 0; i < 6; ++i) {
- assertNotNull(tracker.remove(i), "Expected to remove existing id:
" + i);
+ for (int i = 0, j = 512; i < NUM_ENTRIES; ++i, j += 25) {
+ map.put(j, new DeliveryType(j));
}
- assertEquals(4, tracker.size()); // remaining: 6,7,8,9
- // Add 10..13:
- // bucket4 is full so putting 10 advances head 4->0 (wrap) and uses
bucket0 again
- // then 12 advances head 0->1 and uses bucket1 again
- for (int i = 10; i < 14; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
- assertEquals(8, tracker.size()); // now: 6,7,8,9,10,11,12,13
+ map.put(65534, new DeliveryType(65534));
- // Now remove entries that occupy bucket index 0 (10,11) fully.
- // This should recycle bucket index 0 while tail is at 3 and head at 1:
- // tail > head and index(0) < tail(3) => final else branch in
recycleBucket.
- tracker.removeEach(10, 11, d -> {});
+ map.put(0, new DeliveryType(0));
+ map.put(1, new DeliveryType(1));
+ map.put(2, new DeliveryType(2));
- assertEquals(6, tracker.size());
- assertFalse(tracker.containsKey(10));
- assertFalse(tracker.containsKey(11));
- assertTrue(tracker.containsKey(12));
- assertTrue(tracker.containsKey(13));
+ assertEquals(NUM_ENTRIES + 5, map.size());
- // Validate map remains consistent and all remaining values are
removable.
- final int[] remaining = { 6, 7, 8, 9, 12, 13 };
+ final AtomicInteger removed = new AtomicInteger();
- for (int id : remaining) {
- DeliveryType removed = tracker.remove(id);
- assertNotNull(removed, "Expected to remove existing id: " + id);
- assertEquals(id, removed.getDeliveryId());
- }
+ map.removeEach(0, 65535, (delivery) -> removed.incrementAndGet());
- assertTrue(tracker.isEmpty());
+ assertEquals(NUM_ENTRIES + 2, removed.get());
+ assertEquals(3, map.size());
}
@Test
- public void testIteratorRemoveTriggeringCompactionFromHeadBucket() {
- UnsettledMap<DeliveryType> tracker = createMap(3, 10); // 30 Entries
of capacity
-
- // Fill the map to capacity but no further should remain at three
buckets
- for (int i = 0; i < 30; i++) {
- tracker.put(i, new DeliveryType(i));
- }
-
- Iterator<DeliveryType> it = tracker.values().iterator();
+ public void testRemoveEachEntireDeliveryIdRangeTwoBuckets() {
+ final UnsettledMap<DeliveryType> map = createMap(3, 128);
+ final int NUM_ENTRIES = UnsignedByte.MAX_VALUE.intValue();
- // Move to the middle before we start to remove entries
- for (int i = 0; i < 15; i++) {
- it.next();
+ for (int i = 0; i < NUM_ENTRIES; ++i) {
+ map.put(i, new DeliveryType(i));
}
- DeliveryType lastReturned = null;
+ final AtomicInteger removed = new AtomicInteger();
- // Remove ten elements which should trigger a compaction and a roll
- // forward of data in the center bucket.
- for (int i = 0; i < 10; i++) {
- it.remove();
- lastReturned = it.next();
- }
+ map.removeEach(0, NUM_ENTRIES, (delivery) ->
removed.incrementAndGet());
- for (int i = 24; i < 30; i++) {
- assertEquals(lastReturned.getDeliveryId(), i);
- if (i < 29) {
- lastReturned = it.next(); // Don't pass the end.
- }
- }
+ assertEquals(NUM_ENTRIES, removed.get());
+ assertEquals(0, map.size());
}
@Test
- public void testIteratorRemoveTriggeringCompactionToHeadBucket() {
- UnsettledMap<DeliveryType> tracker = createMap(3, 10); // 30 Entries
of capacity
-
- // Fill the map to capacity but no further should remain at three
buckets
- for (int i = 0; i < 30; i++) {
- tracker.put(i, new DeliveryType(i));
- }
-
- Iterator<UnsignedInteger> it = tracker.keySet().iterator();
+ public void testRemoveEachEntireDeliveryIdRange() {
+ final UnsettledMap<DeliveryType> map = createMap();
+ final int NUM_ENTRIES = UnsignedShort.MAX_VALUE.intValue();
- // Move to the middle before we start to remove entries
- for (int i = 0; i < 25; i++) {
- it.next();
+ for (int i = 0; i < NUM_ENTRIES; ++i) {
+ map.put(i, new DeliveryType(i));
}
- // Should now be positioned in the head bucket and we remove the last
five entries
- for (int i = 0; i < 5; i++) {
- assertEquals(25 + i, it.next().intValue());
- it.remove();
- }
+ final AtomicInteger removed = new AtomicInteger();
- Iterator<DeliveryType> values = tracker.values().iterator();
+ map.removeEach(0, NUM_ENTRIES, (delivery) ->
removed.incrementAndGet());
- for (int i = 0; i < 10; i++) {
- values.next();
- }
+ assertEquals(NUM_ENTRIES, removed.get());
+ assertEquals(0, map.size());
+ }
- for (int i = 0; i < 5; i++) {
- values.next();
- values.remove();
+ @Test
+ public void testRemoveAllEntriesInSmallChunks() {
+ final AtomicInteger removed = new AtomicInteger();
+
+ for (int i = 0; i < uintArray.length; i += 2) {
+ tracker.removeEach(uintArray[i].intValue(),
uintArray[i+1].intValue(), (delivery) -> removed.incrementAndGet());
}
+
+ assertEquals(uintArray.length, removed.get());
}
@Test
- public void testIteratorRemoveWhenHeadBucketShouldHaveWrappedAround() {
- UnsettledMap<DeliveryType> tracker = createMap(3, 10); // 30 Entries
of capacity
+ public void testRemoveEachWithNonZeroReadOffset() {
+ UnsettledMap<DeliveryType> tracker = createMap(2, 16);
- // Fill the map to capacity but no further should remain at three
buckets
- for (int i = 0; i < 30; i++) {
+ for (int i = 0; i < 32; i++) {
tracker.put(i, new DeliveryType(i));
}
- // Remove the first 20 so head and tail are now the same.
- for (int i = 0; i < 20; i++) {
+ for (int i = 0; i < 5; i++) {
tracker.remove(i);
}
- // Put another 20 in so that head wraps to the slot behind tail
- for (int i = 30; i < 50; i++) {
- tracker.put(i, new DeliveryType(i));
+ tracker.removeEach(5, 20, d -> {});
+
+ for (int i = 5; i <= 20; i++) {
+ assertFalse(tracker.containsKey(i));
}
+ }
- assertEquals(30, tracker.size());
+ @Test
+ public void testRemoveEachClearsThenReuseMap() {
+ UnsettledMap<DeliveryType> tracker = createMap();
- // Remove five from head so that when we remove five from the middle
bucket it rolls into head
- tracker.remove(40);
- tracker.remove(41);
- tracker.remove(42);
- tracker.remove(43);
- tracker.remove(44);
+ for (int i = 0; i < 100; i++) {
+ tracker.put(i, new DeliveryType(i));
+ }
- Iterator<UnsignedInteger> it = tracker.keySet().iterator();
+ tracker.removeEach(0, 200, d -> {});
- for (int i = 0; i < 15; ++i) {
- assertEquals(20 + i, it.next().intValue());
- }
+ assertTrue(tracker.isEmpty());
- for (int i = 0; i < 5; ++i) {
- assertEquals(35 + i, it.next().intValue());
- it.remove();
- }
+ tracker.put(999, new DeliveryType(999));
- for (int i = 0; i < 5; ++i) {
- assertEquals(45 + i, it.next().intValue());
- }
+ assertEquals(1, tracker.size());
+ assertNotNull(tracker.get(999));
}
@Test
- public void testIteratorRemoveBeyondDefaultBucketSizeDoesNotThrow() {
- final int bucketSize = 512;
- final UnsettledMap<DeliveryType> tracker = createMap(2, bucketSize);
// uses test helper
+ public void testRecycleBucketWhenTailHasWrapped() {
+ final UnsettledMap<DeliveryType> tracker = createMap(5, 2);
- for (int i = 0; i < bucketSize; ++i) {
+ // Fill all 5 buckets: bucket0=[0,1], bucket1=[2,3], bucket2=[4,5],
bucket3=[6,7], bucket4=[8,9]
+ for (int i = 0; i < 10; ++i) {
tracker.put(i, new DeliveryType(i));
}
+ assertEquals(10, tracker.size());
- final Iterator<DeliveryType> it = tracker.values().iterator();
+ // Remove 0..5 => drains bucket0, bucket1, bucket2 fully (advances
tail to bucket3)
+ for (int i = 0; i < 6; ++i) {
+ assertNotNull(tracker.remove(i), "Expected to remove existing id:
" + i);
+ }
+ assertEquals(4, tracker.size()); // remaining: 6,7,8,9
- DeliveryType removed = null;
- for (int i = 0; i <= 300; ++i) {
- removed = it.next();
+ // Add 10..13:
+ // bucket4 is full so putting 10 advances head 4->0 (wrap) and uses
bucket0 again
+ // then 12 advances head 0->1 and uses bucket1 again
+ for (int i = 10; i < 14; ++i) {
+ tracker.put(i, new DeliveryType(i));
}
+ assertEquals(8, tracker.size()); // now: 6,7,8,9,10,11,12,13
- try {
- it.remove();
- } catch (IndexOutOfBoundsException ex) {
- fail("Iterator remove threw IndexOutOfBoundsException for
bucketSize=" + bucketSize +
- " at index > 256 which is beyond the default bucket capacity
value. ");
+ // Now remove entries that occupy bucket index 0 (10,11) fully.
+ // This should recycle bucket index 0 while tail is at 3 and head at 1:
+ // tail > head and index(0) < tail(3) => final else branch in
recycleBucket.
+ tracker.removeEach(10, 11, d -> {});
+
+ assertEquals(6, tracker.size());
+ assertFalse(tracker.containsKey(10));
+ assertFalse(tracker.containsKey(11));
+ assertTrue(tracker.containsKey(12));
+ assertTrue(tracker.containsKey(13));
+
+ // Validate map remains consistent and all remaining values are
removable.
+ final int[] remaining = { 6, 7, 8, 9, 12, 13 };
+
+ for (int id : remaining) {
+ DeliveryType removed = tracker.remove(id);
+ assertNotNull(removed, "Expected to remove existing id: " + id);
+ assertEquals(id, removed.getDeliveryId());
}
- assertNotNull(removed);
- assertNull(tracker.get(removed.getDeliveryId()));
- assertEquals(bucketSize - 1, tracker.size());
+ assertTrue(tracker.isEmpty());
}
@Test
@@ -2793,36 +3341,6 @@ public class UnsettledMapTest {
assertTrue(tracker.isEmpty());
}
- @Test
- public void
testRemoveFromMiddleDoesNotLoseTailEntriesWhenCompactionTriggered() {
- // Use small bucketSize so bucketLowWaterMark is small and compaction
is more likely.
- final UnsettledMap<DeliveryType> tracker = createMap(6, 10); //
low-water ≈ 3
-
- // Fill enough to create multiple buckets.
- for (int i = 0; i < 60; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
-
- // Make tail bucket sparse (but not empty) so tryCompact(tail) has a
chance to do something.
- for (int i = 0; i < 8; ++i) {
- assertNotNull(tracker.remove(i));
- }
-
- // Now remove from a middle region (not tail) and ensure the early
tail-adjacent IDs remain.
- // If removeValue incorrectly compacts tail and corrupts the ring,
these can disappear.
- assertNotNull(tracker.remove(25)); // removal from a non-tail bucket
-
- // Sanity: nearby entries should still exist
- assertNotNull(tracker.get(24));
- assertNull(tracker.get(25));
- assertNotNull(tracker.get(26));
-
- // Tail-adjacent entries (8..15) should still exist
- for (int i = 8; i < 16; ++i) {
- assertNotNull(tracker.get(i), "Entry " + i + " missing after
middle removal/compaction");
- }
- }
-
@Test
public void testRecycleHeadWhenTailNotZeroDoesNotCorruptSpan() {
final UnsettledMap<DeliveryType> tracker = createMap(8, 4);
@@ -2860,36 +3378,6 @@ public class UnsettledMapTest {
assertNotNull(tracker.get(13));
}
- @Test
- public void
testRemoveFromNonTailTriggersWrongCompactionAndStillPreservesCorrectness() {
- final UnsettledMap<DeliveryType> tracker = createMap(6, 10);
-
- // Fill 3 buckets: 0..29
- for (int i = 0; i < 30; ++i) {
- tracker.put(i, new DeliveryType(i));
- }
-
- // Make tail bucket small (remove 0..6 leaves 7..9 in first bucket =>
3 entries == low-water)
- for (int i = 0; i < 7; ++i) {
- assertNotNull(tracker.remove(i));
- }
-
- // Now remove entries from a *non-tail* bucket to make that bucket
sparse too.
- // Removing these should make the target bucket hit <= low-water and
trigger the compaction path.
- assertNotNull(tracker.remove(15));
- assertNotNull(tracker.remove(16));
- assertNotNull(tracker.remove(17));
-
- // If tryCompact(tail) corrupts the map, these will be missing or
inconsistent.
- for (int i = 7; i < 10; ++i) {
- assertNotNull(tracker.get(i), "Tail-adjacent entry missing after
non-tail removal triggered compaction");
- }
- assertNull(tracker.get(15));
- assertNull(tracker.get(16));
- assertNull(tracker.get(17));
- assertNotNull(tracker.get(18));
- }
-
@Test
public void
testRemoveEachSpanningMultipleBucketRecyclesDoesNotSkipBuckets() {
final UnsettledMap<DeliveryType> tracker = createMap(10, 4);
@@ -2919,70 +3407,137 @@ public class UnsettledMapTest {
}
@Test
- public void testIterateOverBucketsThatHaveWrapped() {
- final UnsettledMap<DeliveryType> tracker = createMap();
+ public void testRemoveEachWithRangeThatWrappedAndNonZeroReadOffset() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- tracker.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
- tracker.put(0, new DeliveryType(0));
- tracker.put(1, new DeliveryType(1));
+ // Populate items in high range and lower wrapped range
+ final int[] entries = new int[] { -4, -3, -2, -1, 0, 1, 2, 3, 4, 5 };
- assertEquals(3, tracker.size());
+ for (int i : entries) {
+ map.put(i, new DeliveryType(i));
+ }
- Iterator<UnsignedInteger> iter = tracker.keySet().iterator();
+ assertEquals(entries.length, map.size());
- final int[] expected = { UnsignedInteger.MAX_VALUE.intValue(), 0, 1 };
+ final AtomicInteger removed = new AtomicInteger();
+ final List<Integer> removedIds = new ArrayList<>();
- int count = 0;
+ // Starts mid-bucket in the upper segment and ends mid-bucket in the
lower segment.
+ map.removeEach(-2, 2, (delivery) -> {
+ removed.incrementAndGet();
+ removedIds.add(delivery.getDeliveryId());
+ });
- while (iter.hasNext()) {
- assertEquals(expected[count++], iter.next().intValue());
- }
+ assertEquals(5, removed.get());
+ assertEquals(entries.length - 5, map.size());
- assertEquals(3, count);
+ assertTrue(removedIds.contains(-2));
+ assertTrue(removedIds.contains(-1));
+ assertTrue(removedIds.contains(0));
+ assertTrue(removedIds.contains(1));
+ assertTrue(removedIds.contains(2));
+
+ // Ensure boundary items remain intact
+ assertNotNull(map.get(-4));
+ assertNotNull(map.get(-3));
+ assertNotNull(map.get(3));
+ assertNotNull(map.get(4));
+ assertNotNull(map.get(5));
}
- protected void dumpRandomDataSet(int iterations, long seed, boolean
bounded) {
- final int[] dataSet = new int[iterations];
+ @Test
+ public void testRemoveEachAcrossMultipleGenerations() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- random.setSeed(seed);
+ map.put(-2, new DeliveryType(-2)); // Gen 0
+ map.put(-1, new DeliveryType(-1)); // Gen 0
+ map.put(0, new DeliveryType(0)); // Gen 1
+ map.put(1, new DeliveryType(1)); // Gen 1
+ map.put(2, new DeliveryType(2)); // Gen 1
- for (int i = 0; i < iterations; ++i) {
- if (bounded) {
- dataSet[i] = random.nextInt(iterations);
- } else {
- dataSet[i] = random.nextInt();
- }
- }
+ assertEquals(5, map.size());
- LOG.info("Iterations was {}, Random seed was: {}", iterations , seed);
- LOG.info("Entries in data set: {}", dataSet);
+ final AtomicInteger removed = new AtomicInteger();
+
+ // Remove range wrapping across boundary
+ map.removeEach(-2, 1, (delivery) -> removed.incrementAndGet());
+
+ assertEquals(4, removed.get());
+ assertEquals(1, map.size());
+
+ assertNull(map.get(-2));
+ assertNull(map.get(-1));
+ assertNull(map.get(0));
+ assertNull(map.get(1));
+ assertNotNull(map.get(2));
}
- protected static class OutsideEntry<K, V> implements Map.Entry<K, V> {
+ @Test
+ public void testRemoveEachRangedContinuationBucketsExceedRangeOfLast() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 10);
- private final K key;
- private V value;
+ // Gen 0
+ map.put(-2, new DeliveryType(-2));
+ map.put(-1, new DeliveryType(-1));
- public OutsideEntry(K key, V value) {
- this.key = key;
- this.value = value;
- }
+ // Gen 1
+ map.put(10, new DeliveryType(10));
+ map.put(11, new DeliveryType(11));
- @Override
- public V setValue(V value) {
- V oldValue = this.value;
- this.value = value;
- return oldValue;
- }
+ final AtomicInteger removed = new AtomicInteger();
- @Override
- public V getValue() {
- return value;
+ // Remove range -2 to 1 (wrapped). Should remove -2 and -1, then
short-circuit when seeing 10.
+ map.removeEach(-2, 1, d -> removed.incrementAndGet());
+
+ assertEquals(2, removed.get());
+ assertEquals(2, map.size());
+ assertNull(map.get(-2));
+ assertNull(map.get(-1));
+ assertNotNull(map.get(10));
+ assertNotNull(map.get(11));
+ }
+
+ @Test
+ public void testRangedOperationsTraverseOverEmptyIntermediateBuckets() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 4);
+
+ // Fill 4 buckets
+ for (int i = 0; i < 16; i++) {
+ map.put(i, new DeliveryType(i));
}
- @Override
- public K getKey() {
- return key;
+ // Drain the middle two buckets completely using single removals
+ for (int i = 4; i < 12; i++) {
+ assertNotNull(map.remove(i));
}
+
+ final AtomicInteger forEachCount = new AtomicInteger();
+ map.forEach(0, 15, d -> forEachCount.incrementAndGet());
+ assertEquals(8, forEachCount.get()); // Only 0..3 and 12..15 remain
+
+ final AtomicInteger removeEachCount = new AtomicInteger();
+ map.removeEach(0, 15, d -> removeEachCount.incrementAndGet());
+ assertEquals(8, removeEachCount.get());
+ assertTrue(map.isEmpty());
+ }
+
+ @Test
+ public void testRangedForEachWithDuplicateKeysAcrossGenerations() {
+ final UnsettledMap<DeliveryType> map = createMap(5, 4);
+
+ // Gen 0: ID 5
+ map.put(5, new DeliveryType(5));
+ map.put(UnsignedInteger.MAX_VALUE.intValue(), new
DeliveryType(UnsignedInteger.MAX_VALUE.intValue()));
+
+ // Gen 1: Overflow and re-insert ID 5
+ map.put(5, new DeliveryType(5));
+
+ final AtomicInteger count = new AtomicInteger();
+
+ // Ranged search bounded to Gen 0 range
+ map.forEach(0, UnsignedInteger.MAX_VALUE.intValue(), d ->
count.incrementAndGet());
+
+ // Should stop at end of Gen 0 and not bleed into Gen 1
+ assertEquals(2, count.get());
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]