diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/MessageGroupQueueTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/MessageGroupQueueTests.java index 24b4aba49d..ce74884e82 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/MessageGroupQueueTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/MessageGroupQueueTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2011 the original author or authors. + * Copyright 2002-2015 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -15,47 +15,58 @@ */ package org.springframework.integration.jdbc; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import org.junit.ClassRule; import org.junit.Test; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; -import org.springframework.messaging.Message; -import org.springframework.messaging.support.GenericMessage; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupQueue; import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.store.SimpleMessageStore; +import org.springframework.integration.test.support.LongRunningIntegrationTest; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.GenericMessage; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertNull; -import static org.junit.Assert.assertTrue; /** * @author Oleg Zhurakousky + * @author Artem Bilan */ public class MessageGroupQueueTests { - + @ClassRule + public static LongRunningIntegrationTest longTests = new LongRunningIntegrationTest(); + @Test - public void validateMgqInterruption() throws Exception{ - + public void validateMgqInterruption() throws Exception { + final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 1); - + final AtomicReference exceptionHolder = new AtomicReference(); - + Thread t = new Thread(new Runnable() { - + + @Override public void run() { queue.offer(new GenericMessage("hello")); try { queue.offer(new GenericMessage("hello"), 100, TimeUnit.SECONDS); - } catch (InterruptedException e) { + } + catch (InterruptedException e) { exceptionHolder.set(e); } } + }); t.start(); Thread.sleep(1000); @@ -63,112 +74,144 @@ public class MessageGroupQueueTests { Thread.sleep(1000); assertTrue(exceptionHolder.get() instanceof InterruptedException); } - + @Test - public void testConcurrentReadWrite() throws Exception{ + public void testConcurrentReadWrite() throws Exception { final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 1); final AtomicReference> messageHolder = new AtomicReference>(); - - Thread t1 = new Thread(new Runnable() { + + Thread t1 = new Thread(new Runnable() { + + @Override public void run() { try { messageHolder.set(queue.poll(1000, TimeUnit.SECONDS)); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t2 = new Thread(new Runnable() { + Thread t2 = new Thread(new Runnable() { + + @Override public void run() { try { queue.offer(new GenericMessage("hello"), 1000, TimeUnit.SECONDS); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); t1.start(); t2.start(); Thread.sleep(1000); assertTrue(messageHolder.get() instanceof Message); } - + @Test - public void testConcurrentWriteRead() throws Exception{ + public void testConcurrentWriteRead() throws Exception { final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 1); final AtomicReference> messageHolder = new AtomicReference>(); - + queue.offer(new GenericMessage("hello"), 1000, TimeUnit.SECONDS); - - Thread t1 = new Thread(new Runnable() { + + Thread t1 = new Thread(new Runnable() { + + @Override public void run() { try { queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t2 = new Thread(new Runnable() { + Thread t2 = new Thread(new Runnable() { + + @Override public void run() { try { queue.poll(1000, TimeUnit.SECONDS); messageHolder.set(queue.poll(1000, TimeUnit.SECONDS)); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - + t1.start(); Thread.sleep(1000); t2.start(); Thread.sleep(1000); assertTrue(messageHolder.get().getPayload().equals("Hi")); } - + @Test - public void testConcurrentReadersWithTimeout() throws Exception{ + public void testConcurrentReadersWithTimeout() throws Exception { final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 1); final AtomicReference> messageHolder1 = new AtomicReference>(); final AtomicReference> messageHolder2 = new AtomicReference>(); final AtomicReference> messageHolder3 = new AtomicReference>(); - - Thread t1 = new Thread(new Runnable() { + + Thread t1 = new Thread(new Runnable() { + + @Override public void run() { try { messageHolder1.set(queue.poll(10, TimeUnit.SECONDS)); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t2 = new Thread(new Runnable() { + Thread t2 = new Thread(new Runnable() { + + @Override public void run() { try { messageHolder2.set(queue.poll(10, TimeUnit.SECONDS)); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t3 = new Thread(new Runnable() { + Thread t3 = new Thread(new Runnable() { + + @Override public void run() { try { messageHolder3.set(queue.poll(10, TimeUnit.SECONDS)); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t4 = new Thread(new Runnable() { + Thread t4 = new Thread(new Runnable() { + + @Override public void run() { try { queue.offer(new GenericMessage("Hi"), 10, TimeUnit.SECONDS); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); t1.start(); Thread.sleep(1000); @@ -178,48 +221,61 @@ public class MessageGroupQueueTests { Thread.sleep(1000); t4.start(); Thread.sleep(1000); - assertTrue(messageHolder1.get().getPayload().equals("Hi")); + assertNotNull(messageHolder1.get()); + assertEquals("Hi", messageHolder1.get().getPayload()); Thread.sleep(4000); assertTrue(messageHolder2.get() == null); } - + @Test - public void testConcurrentWritersWithTimeout() throws Exception{ + public void testConcurrentWritersWithTimeout() throws Exception { final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 1); final AtomicReference booleanHolder1 = new AtomicReference(true); final AtomicReference booleanHolder2 = new AtomicReference(true); final AtomicReference booleanHolder3 = new AtomicReference(true); - - Thread t1 = new Thread(new Runnable() { + + Thread t1 = new Thread(new Runnable() { + + @Override public void run() { try { booleanHolder1.set(queue.offer(new GenericMessage("Hi-1"), 2, TimeUnit.SECONDS)); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t2 = new Thread(new Runnable() { + Thread t2 = new Thread(new Runnable() { + + @Override public void run() { try { boolean offered = queue.offer(new GenericMessage("Hi-2"), 2, TimeUnit.SECONDS); System.out.println(offered); booleanHolder2.set(offered); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); - Thread t3 = new Thread(new Runnable() { + Thread t3 = new Thread(new Runnable() { + + @Override public void run() { try { boolean offered = queue.offer(new GenericMessage("Hi-3"), 2, TimeUnit.SECONDS); System.out.println(offered); booleanHolder3.set(offered); - } catch (Exception e) { + } + catch (Exception e) { e.printStackTrace(); } } + }); t1.start(); Thread.sleep(1000); @@ -231,88 +287,105 @@ public class MessageGroupQueueTests { assertFalse(booleanHolder2.get()); assertFalse(booleanHolder3.get()); } + @Test - public void testConcurrentWriteReadMulti() throws Exception{ - final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 4); - final AtomicReference> messageHolder = new AtomicReference>(); - - queue.offer(new GenericMessage("hello"), 1000, TimeUnit.SECONDS); - - Thread t1 = new Thread(new Runnable() { - public void run() { - try { - queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); - queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); - queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); - } catch (Exception e) { - e.printStackTrace(); - } - } - }); - Thread t2 = new Thread(new Runnable() { - public void run() { - try { - queue.poll(1000, TimeUnit.SECONDS); - messageHolder.set(queue.poll(1000, TimeUnit.SECONDS)); - queue.poll(1000, TimeUnit.SECONDS); - queue.poll(1000, TimeUnit.SECONDS); - } catch (Exception e) { - e.printStackTrace(); - } - } - }); - - t1.start(); - Thread.sleep(1000); - t2.start(); - Thread.sleep(1000); - assertTrue(messageHolder.get().getPayload().equals("Hi")); - assertNull(queue.poll(5, TimeUnit.SECONDS)); - } - + public void testConcurrentWriteReadMulti() throws Exception { + final MessageGroupQueue queue = new MessageGroupQueue(new SimpleMessageStore(), 1, 4); + final AtomicReference> messageHolder = new AtomicReference>(); + + queue.offer(new GenericMessage("hello"), 1000, TimeUnit.SECONDS); + + Thread t1 = new Thread(new Runnable() { + + @Override + public void run() { + try { + queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); + queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); + queue.offer(new GenericMessage("Hi"), 1000, TimeUnit.SECONDS); + } + catch (Exception e) { + e.printStackTrace(); + } + } + + }); + Thread t2 = new Thread(new Runnable() { + + @Override + public void run() { + try { + queue.poll(1000, TimeUnit.SECONDS); + messageHolder.set(queue.poll(1000, TimeUnit.SECONDS)); + queue.poll(1000, TimeUnit.SECONDS); + queue.poll(1000, TimeUnit.SECONDS); + } + catch (Exception e) { + e.printStackTrace(); + } + } + + }); + + t1.start(); + Thread.sleep(1000); + t2.start(); + Thread.sleep(1000); + assertTrue(messageHolder.get().getPayload().equals("Hi")); + assertNull(queue.poll(5, TimeUnit.SECONDS)); + } + @Test - public void validateMgqInterruptionStoreLock() throws Exception{ - - MessageGroupStore mgs = Mockito.mock(MessageGroupStore.class); - Mockito.doAnswer(new Answer() { - public MessageGroup answer(InvocationOnMock invocation) - throws Throwable { - Thread.sleep(5000); - return null; - } - }).when(mgs).addMessageToGroup(Mockito.any(Integer.class), Mockito.any(Message.class)); - - MessageGroup mg = Mockito.mock(MessageGroup.class); - Mockito.when(mgs.getMessageGroup(Mockito.any())).thenReturn(mg); - Mockito.when(mg.size()).thenReturn(0); - - final MessageGroupQueue queue = new MessageGroupQueue(mgs, 1, 1); - - final AtomicReference exceptionHolder = new AtomicReference(); - - Thread t1 = new Thread(new Runnable() { - - public void run() { - queue.offer(new GenericMessage("hello")); - } - }); - t1.start(); - Thread.sleep(500); - Thread t2 = new Thread(new Runnable() { - - public void run() { - queue.offer(new GenericMessage("hello")); - try { - queue.offer(new GenericMessage("hello"), 100, TimeUnit.SECONDS); - } catch (InterruptedException e) { - exceptionHolder.set(e); - } - } - }); - t2.start(); - Thread.sleep(1000); - t2.interrupt(); - Thread.sleep(1000); - assertTrue(exceptionHolder.get() instanceof InterruptedException); - } + public void validateMgqInterruptionStoreLock() throws Exception { + + MessageGroupStore mgs = Mockito.mock(MessageGroupStore.class); + Mockito.doAnswer(new Answer() { + + @Override + public MessageGroup answer(InvocationOnMock invocation) throws Throwable { + Thread.sleep(5000); + return null; + } + + }).when(mgs).addMessageToGroup(Mockito.any(Integer.class), Mockito.any(Message.class)); + + MessageGroup mg = Mockito.mock(MessageGroup.class); + Mockito.when(mgs.getMessageGroup(Mockito.any())).thenReturn(mg); + Mockito.when(mg.size()).thenReturn(0); + + final MessageGroupQueue queue = new MessageGroupQueue(mgs, 1, 1); + + final AtomicReference exceptionHolder = new AtomicReference(); + + Thread t1 = new Thread(new Runnable() { + + @Override + public void run() { + queue.offer(new GenericMessage("hello")); + } + + }); + t1.start(); + Thread.sleep(500); + Thread t2 = new Thread(new Runnable() { + + @Override + public void run() { + queue.offer(new GenericMessage("hello")); + try { + queue.offer(new GenericMessage("hello"), 100, TimeUnit.SECONDS); + } + catch (InterruptedException e) { + exceptionHolder.set(e); + } + } + + }); + t2.start(); + Thread.sleep(1000); + t2.interrupt(); + Thread.sleep(1000); + assertTrue(exceptionHolder.get() instanceof InterruptedException); + } + }