Repository: storm
Updated Branches:
  refs/heads/master 1c2ac2eb4 -> aaebc3b23


http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java
----------------------------------------------------------------------
diff --git 
a/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java
 
b/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java
new file mode 100644
index 0000000..8c19ee8
--- /dev/null
+++ 
b/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java
@@ -0,0 +1,233 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ * <p>
+ * http://www.apache.org/licenses/LICENSE-2.0
+ * <p>
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.storm.topology;
+
+import org.apache.storm.utils.Utils;
+import org.apache.storm.windowing.persistence.SimpleWindowPartitionCache;
+import org.apache.storm.windowing.persistence.WindowPartitionCache;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.FutureTask;
+
+/**
+ * Unit tests for {@link SimpleWindowPartitionCache}
+ */
+public class SimpleWindowPartitionCacheTest {
+
+    @Before
+    public void setUp() throws Exception {
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testBuildInvalid1() throws Exception {
+        SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+            .maximumSize(0)
+            .build(null);
+    }
+
+    @Test(expected = IllegalArgumentException.class)
+    public void testBuildInvalid2() throws Exception {
+        SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+            .maximumSize(-1)
+            .build(null);
+    }
+
+    @Test(expected = NullPointerException.class)
+    public void testBuildInvalid3() throws Exception {
+        SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+            .maximumSize(1)
+            .build(null);
+    }
+
+    @Test
+    public void testBuildOk() throws Exception {
+        SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+            .maximumSize(1)
+            .removalListener((key, val, removalCause) -> {
+            })
+            .build(key -> key);
+    }
+
+    @Test
+    public void testGet() throws Exception {
+        List<Integer> removed = new ArrayList<>();
+        List<Integer> loaded = new ArrayList<>();
+        SimpleWindowPartitionCache<Integer, Integer> cache =
+            SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+                .maximumSize(2)
+                .removalListener((key, val, removalCause) -> removed.add(key))
+                .build(key -> {
+                    loaded.add(key);
+                    return key;
+                });
+
+        cache.get(1);
+        cache.get(2);
+        cache.get(3);
+        Assert.assertEquals(Arrays.asList(1, 2, 3), loaded);
+        // since 2 is the largest un-pinned entry before 3 is loaded
+        Assert.assertEquals(Collections.singletonList(2), removed);
+    }
+
+    @Test(expected = NullPointerException.class)
+    public void testGetNull() throws Exception {
+        SimpleWindowPartitionCache<Integer, Integer> cache =
+            SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+                .maximumSize(2)
+                .build(key -> null);
+
+        cache.get(1);
+    }
+
+    @Test
+    public void testEvictNoRemovalListener() throws Exception {
+        SimpleWindowPartitionCache<Integer, Integer> cache =
+            SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+                .maximumSize(1)
+                .build(key -> {
+                    return key;
+                });
+        cache.get(1);
+        cache.get(2);
+        Assert.assertEquals(Collections.singletonMap(2, 2), cache.asMap());
+        cache.invalidate(2);
+        Assert.assertEquals(Collections.emptyMap(), cache.asMap());
+    }
+
+    @Test
+    public void testPinAndGet() throws Exception {
+        List<Integer> removed = new ArrayList<>();
+        List<Integer> loaded = new ArrayList<>();
+        SimpleWindowPartitionCache<Integer, Integer> cache =
+            SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+                .maximumSize(1)
+                .removalListener(new 
WindowPartitionCache.RemovalListener<Integer, Integer>() {
+                    @Override
+                    public void onRemoval(Integer key, Integer val, 
WindowPartitionCache.RemovalCause removalCause) {
+                        removed.add(key);
+                    }
+                })
+                .build(new WindowPartitionCache.CacheLoader<Integer, 
Integer>() {
+                    @Override
+                    public Integer load(Integer key) {
+                        loaded.add(key);
+                        return key;
+                    }
+                });
+
+        cache.get(1);
+        cache.pinAndGet(2);
+        cache.get(3);
+        Assert.assertEquals(Arrays.asList(1, 2, 3), loaded);
+        Assert.assertEquals(Collections.singletonList(1), removed);
+    }
+
+    @Test
+    public void testInvalidate() throws Exception {
+        List<Integer> removed = new ArrayList<>();
+        List<Integer> loaded = new ArrayList<>();
+        SimpleWindowPartitionCache<Integer, Integer> cache =
+            SimpleWindowPartitionCache.<Integer, Integer>newBuilder()
+                .maximumSize(1)
+                .removalListener((key, val, removalCause) -> removed.add(key))
+                .build(key -> {
+                    loaded.add(key);
+                    return key;
+                });
+
+        cache.pinAndGet(1);
+        cache.invalidate(1);
+        Assert.assertEquals(Collections.singletonList(1), loaded);
+        Assert.assertEquals(Collections.emptyList(), removed);
+        Assert.assertEquals(cache.asMap(), Collections.singletonMap(1, 1));
+
+        cache.unpin(1);
+        cache.invalidate(1);
+        Assert.assertTrue(cache.asMap().isEmpty());
+    }
+
+
+    @Test(timeout = 10000)
+    public void testConcurrentGet() throws Exception {
+        List<Integer> loaded = new ArrayList<>();
+        SimpleWindowPartitionCache<Integer, Object> cache =
+            SimpleWindowPartitionCache.<Integer, Object>newBuilder()
+                .maximumSize(1)
+                .build(key -> {
+                    Utils.sleep(1000);
+                    loaded.add(key);
+                    return new Object();
+                });
+
+        FutureTask<Object> ft1 = new FutureTask<>(() -> cache.pinAndGet(1));
+        FutureTask<Object> ft2 = new FutureTask<>(() -> cache.pinAndGet(1));
+        Thread t1 = new Thread(ft1);
+        Thread t2 = new Thread(ft2);
+        t1.start();
+        t2.start();
+        t1.join();
+        t2.join();
+
+        Assert.assertEquals(Collections.singletonList(1), loaded);
+        Assert.assertEquals(ft1.get(), ft2.get());
+    }
+
+    @Test
+    public void testConcurrentUnpin() throws Exception {
+        SimpleWindowPartitionCache<Integer, Object> cache =
+            SimpleWindowPartitionCache.<Integer, Object>newBuilder()
+                .maximumSize(1)
+                .build(key -> new Object());
+
+        cache.pinAndGet(1);
+        FutureTask<Boolean> ft1 = new FutureTask<>(() -> cache.unpin(1));
+        FutureTask<Boolean> ft2 = new FutureTask<>(() -> cache.unpin(1));
+        Thread t1 = new Thread(ft1);
+        Thread t2 = new Thread(ft2);
+        t1.start();
+        t2.start();
+        t1.join();
+        t2.join();
+
+        Assert.assertTrue(ft1.get() || ft2.get());
+        Assert.assertFalse(ft1.get() && ft2.get());
+    }
+
+        @Test
+    public void testEviction() throws Exception {
+        List<Integer> removed = new ArrayList<>();
+        SimpleWindowPartitionCache<Integer, Object> cache =
+            SimpleWindowPartitionCache.<Integer, Object>newBuilder()
+                .maximumSize(1)
+                .removalListener((key, val, removalCause) -> removed.add(key))
+                .build(key -> new Object());
+
+        cache.get(0);
+        cache.pinAndGet(1);
+        Assert.assertEquals(Collections.singletonList(0), removed);
+        cache.get(2);
+        Assert.assertEquals(Collections.singletonList(0), removed);
+    }
+}
\ No newline at end of file

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java
----------------------------------------------------------------------
diff --git 
a/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java 
b/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java
index 178c1bb..d99ecb3 100644
--- a/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java
+++ b/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java
@@ -98,8 +98,8 @@ public class WindowManagerTest {
 
     @Test
     public void testCountBasedWindow() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
CountEvictionPolicy<Integer>(5);
-        TriggerPolicy<Integer> triggerPolicy = new 
CountTriggerPolicy<Integer>(2, windowManager, evictionPolicy);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
CountEvictionPolicy<Integer>(5);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
CountTriggerPolicy<Integer>(2, windowManager, evictionPolicy);
         triggerPolicy.start();
         windowManager.setEvictionPolicy(evictionPolicy);
         windowManager.setTriggerPolicy(triggerPolicy);
@@ -141,7 +141,7 @@ public class WindowManagerTest {
         int threshold = WindowManager.EXPIRE_EVENTS_THRESHOLD;
         int windowLength = 5;
         windowManager.setEvictionPolicy(new CountEvictionPolicy<Integer>(5));
-        TriggerPolicy<Integer> triggerPolicy = new 
TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.HOURS).value, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.HOURS).value, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
         for (int i : seq(1, 5)) {
@@ -203,13 +203,13 @@ public class WindowManagerTest {
 
     @Test
     public void testTimeBasedWindow() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
TimeEvictionPolicy<Integer>(new Duration(1, TimeUnit.SECONDS).value);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
TimeEvictionPolicy<Integer>(new Duration(1, TimeUnit.SECONDS).value);
         windowManager.setEvictionPolicy(evictionPolicy);
         /*
          * Don't wait for Timetrigger to fire since this could lead to timing 
issues in unit tests.
          * Set it to a large value and trigger manually.
           */
-        TriggerPolicy<Integer> triggerPolicy = new 
TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager, 
evictionPolicy);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager, 
evictionPolicy);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
         long now = System.currentTimeMillis();
@@ -265,13 +265,13 @@ public class WindowManagerTest {
 
     @Test
     public void testTimeBasedWindowExpiry() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
TimeEvictionPolicy<Integer>(new Duration(100, TimeUnit.MILLISECONDS).value);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
TimeEvictionPolicy<Integer>(new Duration(100, TimeUnit.MILLISECONDS).value);
         windowManager.setEvictionPolicy(evictionPolicy);
         /*
          * Don't wait for Timetrigger to fire since this could lead to timing 
issues in unit tests.
          * Set it to a large value and trigger manually.
           */
-        TriggerPolicy<Integer> triggerPolicy = new 
TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
         long now = System.currentTimeMillis();
@@ -302,9 +302,9 @@ public class WindowManagerTest {
 
     @Test
     public void testTumblingWindow() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
CountEvictionPolicy<Integer>(3);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
CountEvictionPolicy<Integer>(3);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
CountTriggerPolicy<Integer>(3, windowManager, evictionPolicy);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
CountTriggerPolicy<Integer>(3, windowManager, evictionPolicy);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
         windowManager.add(1);
@@ -332,9 +332,9 @@ public class WindowManagerTest {
 
     @Test
     public void testEventTimeBasedWindow() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkTimeEvictionPolicy<>(20);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkTimeEvictionPolicy<>(20);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 
@@ -398,9 +398,9 @@ public class WindowManagerTest {
 
     @Test
     public void testCountBasedWindowWithEventTs() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkCountEvictionPolicy<>(3);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkCountEvictionPolicy<>(3);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 
@@ -437,9 +437,9 @@ public class WindowManagerTest {
 
     @Test
     public void testCountBasedTriggerWithEventTs() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkTimeEvictionPolicy<Integer>(20);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkTimeEvictionPolicy<Integer>(20);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkCountTriggerPolicy<Integer>(3, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkCountTriggerPolicy<Integer>(3, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 
@@ -477,9 +477,9 @@ public class WindowManagerTest {
 
     @Test
     public void testCountBasedTumblingWithSameEventTs() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkCountEvictionPolicy<>(2);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkCountEvictionPolicy<>(2);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 
@@ -505,9 +505,9 @@ public class WindowManagerTest {
 
     @Test
     public void testCountBasedSlidingWithSameEventTs() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkCountEvictionPolicy<>(5);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkCountEvictionPolicy<>(5);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 
@@ -533,9 +533,9 @@ public class WindowManagerTest {
     }
         @Test
     public void testEventTimeLag() throws Exception {
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkTimeEvictionPolicy<>(20, 5);
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkTimeEvictionPolicy<>(20, 5);
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 
@@ -560,7 +560,7 @@ public class WindowManagerTest {
     @Test
     public void testScanStop() throws Exception {
         final Set<Integer> eventsScanned = new HashSet<>();
-        EvictionPolicy<Integer> evictionPolicy = new 
WatermarkTimeEvictionPolicy<Integer>(20, 5) {
+        EvictionPolicy<Integer, ?> evictionPolicy = new 
WatermarkTimeEvictionPolicy<Integer>(20, 5) {
 
             @Override
             public Action evict(Event<Integer> event) {
@@ -570,7 +570,7 @@ public class WindowManagerTest {
 
         };
         windowManager.setEvictionPolicy(evictionPolicy);
-        TriggerPolicy<Integer> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
+        TriggerPolicy<Integer, ?> triggerPolicy = new 
WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, 
windowManager);
         triggerPolicy.start();
         windowManager.setTriggerPolicy(triggerPolicy);
 

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java
----------------------------------------------------------------------
diff --git 
a/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java
 
b/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java
new file mode 100644
index 0000000..587d874
--- /dev/null
+++ 
b/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java
@@ -0,0 +1,246 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ * <p>
+ * http://www.apache.org/licenses/LICENSE-2.0
+ * <p>
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.storm.windowing.persistence;
+
+import org.apache.storm.state.KeyValueState;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.windowing.Event;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Captor;
+import org.mockito.Mock;
+import org.mockito.Mockito;
+import org.mockito.MockitoAnnotations;
+import org.mockito.invocation.InvocationOnMock;
+import org.mockito.runners.MockitoJUnitRunner;
+import org.mockito.stubbing.Answer;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Deque;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.Consumer;
+import java.util.function.Supplier;
+
+import static org.mockito.AdditionalAnswers.returnsArgAt;
+
+/**
+ * Unit tests for {@link WindowState}
+ */
+@RunWith(MockitoJUnitRunner.class)
+public class WindowStateTest {
+
+    @Mock
+    private KeyValueState<Long, WindowState.WindowPartition<Integer>> 
windowState;
+    @Mock
+    private KeyValueState<String, Deque<Long>> partitionIdsState;
+    @Mock
+    private KeyValueState<String, Optional<?>> systemState;
+    @Mock
+    private Supplier<Map<String, Optional<?>>> supplier;
+    @Captor
+    private ArgumentCaptor<Long> longCaptor;
+    @Captor
+    private ArgumentCaptor<WindowState.WindowPartition<Integer>> 
windowValuesCaptor;
+
+    @Before
+    public void setUp() throws Exception {
+        MockitoAnnotations.initMocks(this);
+    }
+
+    @Test
+    public void testAdd() throws Exception {
+        Mockito.when(partitionIdsState.get(Mockito.any(), 
Mockito.any())).then(returnsArgAt(1));
+        Mockito.when(windowState.get(Mockito.any(), 
Mockito.any())).then(returnsArgAt(1));
+
+        WindowState<Integer> ws = getWindowState(10 * 
WindowState.MAX_PARTITION_EVENTS);
+
+        long partitions = 15;
+        long numEvents = partitions * WindowState.MAX_PARTITION_EVENTS;
+        for (int i = 0; i < numEvents; i++) {
+            ws.add(getEvent(i));
+        }
+        // 5 partitions evicted to window state
+        Mockito.verify(windowState, 
Mockito.times(5)).put(longCaptor.capture(), windowValuesCaptor.capture());
+        Assert.assertEquals(5, longCaptor.getAllValues().size());
+        // each evicted partition has MAX_EVENTS_PER_PARTITION
+        windowValuesCaptor.getAllValues().forEach(wp -> {
+            Assert.assertEquals(WindowState.MAX_PARTITION_EVENTS, wp.size());
+        });
+        // last partition is not evicted
+        Assert.assertFalse(longCaptor.getAllValues().contains(partitions - 1));
+    }
+
+    @Test
+    public void testIterator() throws Exception {
+        Map<Long, WindowState.WindowPartition<Event<Tuple>>> partitionMap = 
new HashMap<>();
+        Mockito.when(partitionIdsState.get(Mockito.any(), 
Mockito.any())).then(returnsArgAt(1));
+        Mockito.when(windowState.get(Mockito.any(), Mockito.any())).then(new 
Answer<Object>() {
+            @Override
+            public Object answer(InvocationOnMock invocation) throws Throwable 
{
+                Object[] args = invocation.getArguments();
+                WindowState.WindowPartition<Event<Tuple>> evicted = 
partitionMap.get(args[0]);
+                return evicted != null ? evicted : args[1];
+            }
+        });
+
+        Mockito.doAnswer(new Answer<Void>() {
+            @Override
+            public Void answer(InvocationOnMock invocation) throws Throwable {
+                Object[] args = invocation.getArguments();
+                partitionMap.put((long)args[0], 
(WindowState.WindowPartition<Event<Tuple>>)args[1]);
+                return null;
+            }
+        }).when(windowState).put(Mockito.any(), Mockito.any());
+
+        Mockito.doAnswer(new Answer<Void>() {
+            @Override
+            public Void answer(InvocationOnMock invocation) throws Throwable {
+                Object[] args = invocation.getArguments();
+                partitionMap.remove(args[0]);
+                return null;
+            }
+        }).when(windowState).delete(Mockito.anyLong());
+
+        Mockito.when(supplier.get()).thenReturn(Collections.emptyMap());
+
+        WindowState<Integer> ws = getWindowState(10 * 
WindowState.MAX_PARTITION_EVENTS);
+
+        long partitions = 15;
+
+        long numEvents = partitions * WindowState.MAX_PARTITION_EVENTS;
+        List<Event<Integer>> expected = new ArrayList<>();
+        for (int i = 0; i < numEvents; i++) {
+            Event<Integer> event = getEvent(i);
+            expected.add(event);
+            ws.add(event);
+        }
+
+        Assert.assertEquals(5, partitionMap.size());
+        Iterator<Event<Integer>> it = ws.iterator();
+        List<Event<Integer>> actual = new ArrayList<>();
+        it.forEachRemaining(actual::add);
+        Assert.assertEquals(expected, actual);
+
+        // iterate again
+        it = ws.iterator();
+        actual.clear();
+        it.forEachRemaining(actual::add);
+        Assert.assertEquals(expected, actual);
+
+        // remove
+        it = ws.iterator();
+        while (it.hasNext()) {
+            it.next();
+            it.remove();
+        }
+
+        it = ws.iterator();
+        actual.clear();
+        it.forEachRemaining(actual::add);
+        Assert.assertEquals(Collections.emptyList(), actual);
+    }
+
+    @Test
+    public void testIteratorPartitionNotEvicted() throws Exception {
+        Map<Long, WindowState.WindowPartition<Event<Tuple>>> partitionMap = 
new HashMap<>();
+        Mockito.when(partitionIdsState.get(Mockito.any(), 
Mockito.any())).then(returnsArgAt(1));
+        Mockito.when(windowState.get(Mockito.any(), Mockito.any())).then(new 
Answer<Object>() {
+            @Override
+            public Object answer(InvocationOnMock invocation) throws Throwable 
{
+                Object[] args = invocation.getArguments();
+                WindowState.WindowPartition<Event<Tuple>> evicted = 
partitionMap.get(args[0]);
+                return evicted != null ? evicted : args[1];
+            }
+        });
+
+        Mockito.doAnswer(new Answer<Void>() {
+            @Override
+            public Void answer(InvocationOnMock invocation) throws Throwable {
+                Object[] args = invocation.getArguments();
+                partitionMap.put((long)args[0], 
(WindowState.WindowPartition<Event<Tuple>>)args[1]);
+                return null;
+            }
+        }).when(windowState).put(Mockito.any(), Mockito.any());
+
+        Mockito.when(supplier.get()).thenReturn(Collections.emptyMap());
+
+        WindowState<Integer> ws = getWindowState(10 * 
WindowState.MAX_PARTITION_EVENTS);
+
+        long partitions = 10;
+
+        long numEvents = partitions * WindowState.MAX_PARTITION_EVENTS;
+        List<Event<Integer>> expected = new ArrayList<>();
+        for (int i = 0; i < numEvents; i++) {
+            Event<Integer> event = getEvent(i);
+            expected.add(event);
+            ws.add(event);
+        }
+
+        // Stop iterating in the middle of the 10th partition
+        Iterator<Event<Integer>> it = ws.iterator();
+        for(int i=0; i<9500; i++) {
+            it.next();
+        }
+
+        for (int i = 0; i < numEvents; i++) {
+            Event<Integer> event = getEvent(i);
+            expected.add(event);
+            ws.add(event);
+        }
+
+        // 10th partition should not have been evicted
+        Assert.assertFalse(partitionMap.containsKey(9L));
+    }
+
+    private Event<Integer> getEvent(int i) {
+        return getEvent(i, 0);
+    }
+
+    private Event<Integer> getEvent(int i, long ts) {
+        return new Event<Integer>() {
+            @Override
+            public long getTimestamp() {
+                return ts;
+            }
+
+            @Override
+            public Integer get() {
+                return i;
+            }
+
+            @Override
+            public boolean isWatermark() {
+                return false;
+            }
+        };
+    }
+
+    private WindowState<Integer> getWindowState(int maxEvents) {
+        return new WindowState<>(windowState, partitionIdsState, systemState,
+            supplier, maxEvents);
+    }
+}
\ No newline at end of file

Reply via email to