فهرست منبع

resolve #63 - CQuery iteration is thread-safe, although indexing could not reflect order of values

Ranides Atterwim 4 سال پیش
والد
کامیت
6fe00233f7

+ 3 - 2
assira.core/src/main/java/net/ranides/assira/collection/query/CQueryAbstract.java

@@ -41,6 +41,7 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.Spliterator;
 import java.util.TreeSet;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.BinaryOperator;
 import java.util.function.Consumer;
 import java.util.function.Function;
@@ -522,8 +523,8 @@ public abstract class CQueryAbstract<T> implements CQuery<T>, CQueryFeatures {
     @Override
     public boolean whileEach(Predicates.EachPredicate<? super T> consumer) {
         if(hasFastEach()) {
-            int[] index = {0};
-            return whileEach(v -> consumer.test(index[0]++, v));
+            AtomicInteger index = new AtomicInteger(0);
+            return whileEach(v -> consumer.test(index.getAndIncrement(), v));
         }
         if(hasFastIterator()) {
             return BaseIterable.whileEach(this, consumer);

+ 3 - 2
assira.core/src/main/java/net/ranides/assira/collection/query/derived/CQFilterEach.java

@@ -6,6 +6,7 @@ import net.ranides.assira.collection.query.base.CQAbstractFilter;
 import net.ranides.assira.functional.Predicates;
 
 import java.util.Iterator;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Predicate;
 import java.util.stream.Stream;
 
@@ -20,8 +21,8 @@ public class CQFilterEach<T> extends CQAbstractFilter<T, T> {
 
     @Override
     public Stream<T> stream() {
-        int[] index = {0};
-        return source.stream().filter(v -> p.test(index[0]++, v));
+        AtomicInteger index = new AtomicInteger(0);
+        return source.stream().filter(v -> p.test(index.getAndIncrement(), v));
     }
 
     @Override

+ 4 - 3
assira.core/src/main/java/net/ranides/assira/collection/query/derived/CQSlice.java

@@ -7,6 +7,7 @@ import net.ranides.assira.functional.Predicates;
 import net.ranides.assira.math.MathUtils;
 
 import java.util.Iterator;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Stream;
 
 public class CQSlice<T> extends CQueryAbstract<T> {
@@ -68,7 +69,7 @@ public class CQSlice<T> extends CQueryAbstract<T> {
         if(!hasFastEach()) {
             return super.whileEach(consumer);
         }
-        boolean[] stop = { false };
+        AtomicBoolean stop = new AtomicBoolean(false);
         boolean canceled = source.whileEach((i,v) -> {
             if(i < begin) {
                 return true;
@@ -77,12 +78,12 @@ public class CQSlice<T> extends CQueryAbstract<T> {
                 return false;
             }
             if(i + 1 == end) {
-                stop[0] = true;
+                stop.set(true);
                 return false;
             }
             return true;
         });
-        return canceled || stop[0];
+        return canceled || stop.get();
     }
 
     @Override

+ 3 - 2
assira.core/src/main/java/net/ranides/assira/collection/query/support/BaseCollection.java

@@ -5,6 +5,7 @@ import lombok.experimental.UtilityClass;
 import net.ranides.assira.functional.Consumers;
 
 import java.util.Collection;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Consumer;
 
 @UtilityClass
@@ -15,7 +16,7 @@ public class BaseCollection {
     }
 
     public static <T> void forEach(Collection<T> that, Consumers.EachConsumer<? super T> consumer) {
-        int[] index = {0};
-        that.forEach(v -> consumer.accept(index[0]++, v));
+        AtomicInteger index = new AtomicInteger(0);
+        that.forEach(v -> consumer.accept(index.getAndIncrement(), v));
     }
 }

+ 8 - 6
assira.core/src/main/java/net/ranides/assira/collection/query/support/BaseEach.java

@@ -14,6 +14,8 @@ import java.util.List;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.BiConsumer;
 import java.util.function.BinaryOperator;
 import java.util.function.Function;
@@ -26,18 +28,18 @@ import java.util.stream.Collector;
 public class BaseEach {
 
     public static <T> int size(CQuery<T> that) {
-        int[] out = {0};
-        that.forEach(v -> out[0]++);
-        return out[0];
+        AtomicInteger index = new AtomicInteger(0);
+        that.forEach(v -> index.getAndIncrement());
+        return index.get();
     }
 
     public static <T> boolean isEmpty(CQuery<T> that) {
-        boolean[] out = {true};
+        AtomicBoolean out = new AtomicBoolean(true);
         that.whileEach(v -> {
-            out[0] = false;
+            out.set(false);
             return false;
         });
-        return out[0];
+        return out.get();
     }
 
     public static <T, A, R> R collect(CQuery<T> that, Collector<? super T, A, R> collector) {

+ 3 - 2
assira.core/src/main/java/net/ranides/assira/collection/query/support/BaseStream.java

@@ -11,6 +11,7 @@ import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
 import java.util.Spliterator;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.BinaryOperator;
 import java.util.function.Function;
 import java.util.function.IntFunction;
@@ -79,7 +80,7 @@ public class BaseStream {
     }
 
     public static <T> boolean whileEach(CQuery<T> that, Predicates.EachPredicate<? super T> consumer) {
-        int[] index = {0};
-        return that.stream().allMatch(v -> consumer.test(index[0]++, v));
+        AtomicInteger index = new AtomicInteger(0);
+        return that.stream().allMatch(v -> consumer.test(index.getAndIncrement(), v));
     }
 }